feat(edp): provision Foundry reader grants

This commit is contained in:
Codex
2026-07-19 15:03:48 +03:00
parent 847a08da93
commit def9a24e0d
17 changed files with 891 additions and 7 deletions
@@ -8,6 +8,7 @@ import {
export function readConfig(env = process.env) {
const provisionerApiEnabled = boolean(env.EXTERNAL_DATA_PLANE_PROVISIONING_ENABLED, false);
const managedProvisionerApiEnabled = boolean(env.EXTERNAL_DATA_PLANE_MANAGED_PROVISIONING_ENABLED, false);
const foundryProvisionerApiEnabled = boolean(env.EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONING_ENABLED, false);
const managedProvisionerMaxSkewSeconds = integer(
env.EXTERNAL_DATA_PLANE_MANAGED_PROVISIONER_MAX_SKEW_SECONDS,
60,
@@ -21,6 +22,7 @@ export function readConfig(env = process.env) {
internalAccessToken: optional(env.NODEDC_INTERNAL_ACCESS_TOKEN),
provisionerApiEnabled,
managedProvisionerApiEnabled,
foundryProvisionerApiEnabled,
provisionerAccessToken: provisionerApiEnabled ? secretFile(env.EXTERNAL_DATA_PLANE_PROVISIONER_TOKEN_FILE) : "",
managedProvisionerPublicKey: managedProvisionerApiEnabled
? loadManagedProvisionerPublicKeyFile(env.EXTERNAL_DATA_PLANE_MANAGED_PROVISIONER_PUBLIC_KEY_FILE)
@@ -49,6 +51,38 @@ export function readConfig(env = process.env) {
100,
100_000,
),
foundryProvisionerPublicKey: foundryProvisionerApiEnabled
? loadManagedProvisionerPublicKeyFile(env.EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_PUBLIC_KEY_FILE)
: null,
foundryProvisionerServiceId: foundryProvisionerApiEnabled
? validateManagedProvisionerIdentity(
required(env.EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_SERVICE_ID, "EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_SERVICE_ID"),
"service_id",
)
: "",
foundryProvisionerKeyId: foundryProvisionerApiEnabled
? validateManagedProvisionerIdentity(
required(env.EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_KEY_ID, "EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_KEY_ID"),
"key_id",
)
: "",
foundryProvisionerAudience: foundryProvisionerApiEnabled
? validateManagedProvisionerAudience(
required(env.EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_AUDIENCE, "EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_AUDIENCE"),
)
: "",
foundryProvisionerMaxSkewMs: integer(
env.EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_MAX_SKEW_SECONDS,
60,
5,
300,
) * 1000,
foundryProvisionerReplayCacheMaxEntries: integer(
env.EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_REPLAY_CACHE_MAX_ENTRIES,
10_000,
100,
100_000,
),
rawRetentionDays: integer(env.EXTERNAL_DATA_PLANE_RAW_RETENTION_DAYS, 14, 1, 3650),
maxBatchBytes: integer(env.EXTERNAL_DATA_PLANE_MAX_BATCH_BYTES, 5 * 1024 * 1024, 1024, 50 * 1024 * 1024),
maxFactsPerPublish: integer(env.EXTERNAL_DATA_PLANE_MAX_FACTS_PER_PUBLISH, 5000, 1, 100_000),
@@ -11,6 +11,13 @@ const MANAGED_REQUEST_KEYS = new Set([
"generation",
"capabilityDigest",
]);
const MANAGED_CONSUMER_REQUEST_KEYS = new Set([
"allowedDataProductIds",
"expiresAt",
"generation",
"capabilityDigest",
]);
const MANAGED_CONSUMER_PLAN_KEYS = new Set(["allowedDataProductIds"]);
const SOURCE_KEYS = new Set(["tenantId", "connectionId", "providerId"]);
const SHA256_DIGEST = /^[a-f0-9]{64}$/;
@@ -82,6 +89,38 @@ export function normalizeManagedReaderBindingRequest(value) {
});
}
export function normalizeManagedConsumerReaderPlanRequest(value) {
if (!isPlainObject(value) || containsSecretLikeKey(value) || !hasOnlyKeys(value, MANAGED_CONSUMER_PLAN_KEYS)) {
throw readerError("managed_consumer_reader_plan_request_invalid");
}
const allowedDataProductIds = uniqueIdentifiers(value.allowedDataProductIds);
if (!allowedDataProductIds.length) throw readerError("managed_consumer_reader_scope_invalid");
return Object.freeze({ allowedDataProductIds: Object.freeze([...allowedDataProductIds].sort()) });
}
export function normalizeManagedConsumerReaderBindingRequest(value) {
if (!isPlainObject(value) || containsSecretLikeKey(value) || !hasOnlyKeys(value, MANAGED_CONSUMER_REQUEST_KEYS)) {
throw readerError("managed_consumer_reader_binding_request_invalid");
}
const allowedDataProductIds = uniqueIdentifiers(value.allowedDataProductIds);
if (!allowedDataProductIds.length) throw readerError("managed_consumer_reader_scope_invalid");
if (value.expiresAt !== null) throw readerError("managed_consumer_reader_must_be_durable");
const generation = Number(value.generation);
if (!Number.isSafeInteger(generation) || generation < 1 || generation > 2_147_483_647) {
throw readerError("managed_consumer_reader_generation_invalid");
}
const capabilityDigest = String(value.capabilityDigest || "").toLowerCase();
if (!SHA256_DIGEST.test(capabilityDigest)) {
throw readerError("managed_consumer_reader_capability_digest_invalid");
}
return Object.freeze({
allowedDataProductIds: Object.freeze([...allowedDataProductIds].sort()),
expiresAt: null,
generation,
capabilityDigest,
});
}
export function readerBindingRequestHash(policy) {
const canonical = JSON.stringify({
tenantId: policy.tenantId,
@@ -95,6 +134,15 @@ export function readerBindingRequestHash(policy) {
return createHash("sha256").update(canonical, "utf8").digest("hex");
}
export function consumerReaderBindingRequestHash(policy) {
return createHash("sha256").update(JSON.stringify({
allowedDataProductIds: [...policy.allowedDataProductIds].sort(),
expiresAt: null,
generation: policy.generation,
capabilityDigest: policy.capabilityDigest,
}), "utf8").digest("hex");
}
export function assertReaderProduct(binding, dataProductId, now = new Date()) {
const expired = binding?.expiresAt !== null && new Date(binding?.expiresAt) <= now;
if (!isPlainObject(binding) || binding.active !== true || expired) {
@@ -124,6 +172,21 @@ export function safeReaderBinding(binding) {
};
}
export function safeManagedConsumerReaderBinding(binding) {
return {
id: binding.id,
bindingKey: binding.bindingKey,
generation: Number(binding.generation),
allowedDataProductIds: uniqueIdentifiers(binding.allowedDataProductIds),
active: binding.active === true,
expiresAt: binding.expiresAt === null ? null : new Date(binding.expiresAt).toISOString(),
createdAt: binding.createdAt ? new Date(binding.createdAt).toISOString() : undefined,
rotatedAt: binding.rotatedAt ? new Date(binding.rotatedAt).toISOString() : undefined,
revokedAt: binding.revokedAt ? new Date(binding.revokedAt).toISOString() : undefined,
sourceScope: "resolved-server-side",
};
}
function readerError(code, status = 400) {
return Object.assign(new Error(code), { status, code });
}
@@ -41,6 +41,47 @@ export async function resolveManagedReaderSourceConnection(db, policy) {
throw sourceScopeError("managed_reader_source_scope_ambiguous", 409);
}
export async function resolveManagedReaderScope(db, allowedDataProductIdsValue) {
const allowedDataProductIds = uniqueIdentifiers(allowedDataProductIdsValue);
if (!allowedDataProductIds.length) {
throw sourceScopeError("managed_consumer_reader_scope_request_invalid", 400);
}
const result = await db.query(
`select distinct candidate.tenant_id as "tenantId",
candidate.connection_id as "connectionId", candidate.provider_id as "providerId"
from external_data_plane_writer_bindings as candidate
where candidate.active = true
and (candidate.expires_at is null or candidate.expires_at > now())
and not exists (
select 1
from unnest($1::text[]) as requested(data_product_id)
where not exists (
select 1
from external_data_plane_writer_bindings as coverage
where coverage.tenant_id = candidate.tenant_id
and coverage.connection_id = candidate.connection_id
and coverage.provider_id = candidate.provider_id
and coverage.active = true
and (coverage.expires_at is null or coverage.expires_at > now())
and coverage.allowed_data_product_ids ? requested.data_product_id
)
)
order by candidate.tenant_id asc, candidate.provider_id asc, candidate.connection_id asc
limit 3`,
[allowedDataProductIds],
);
const candidates = result.rows
.map((row) => ({
tenantId: identifier(row.tenantId),
connectionId: identifier(row.connectionId),
providerId: identifier(row.providerId),
}))
.filter((row) => row.tenantId && row.connectionId && row.providerId);
if (candidates.length === 1) return Object.freeze(candidates[0]);
if (!candidates.length) throw sourceScopeError("managed_consumer_reader_source_scope_not_found", 409);
throw sourceScopeError("managed_consumer_reader_source_scope_ambiguous", 409);
}
export async function backfillManagedReaderSourceConnections(db) {
const unresolved = await db.query(
`select id, tenant_id as "tenantId", connection_id as "connectionId",
+168 -1
View File
@@ -25,14 +25,18 @@ import {
} from "./managed-provisioner-auth.mjs";
import {
assertReaderProduct,
consumerReaderBindingRequestHash,
createReaderToken,
hashReaderToken,
normalizeManagedConsumerReaderBindingRequest,
normalizeManagedConsumerReaderPlanRequest,
normalizeManagedReaderBindingRequest,
normalizeReaderBindingRequest,
readerBindingRequestHash,
safeManagedConsumerReaderBinding,
safeReaderBinding,
} from "./reader-binding.mjs";
import { resolveManagedReaderSourceConnection } from "./reader-source-scope.mjs";
import { resolveManagedReaderScope, resolveManagedReaderSourceConnection } from "./reader-source-scope.mjs";
import { migrate } from "./schema.mjs";
import {
createWriterToken,
@@ -60,6 +64,19 @@ const verifyManagedProvisionerRequest = config.managedProvisionerApiEnabled
replayCache: managedProvisionerReplayCache,
})
: null;
const foundryProvisionerReplayCache = config.foundryProvisionerApiEnabled
? new ManagedProvisionerReplayCache({ maxEntries: config.foundryProvisionerReplayCacheMaxEntries })
: null;
const verifyFoundryProvisionerRequest = config.foundryProvisionerApiEnabled
? createManagedProvisionerRequestVerifier({
publicKey: config.foundryProvisionerPublicKey,
serviceId: config.foundryProvisionerServiceId,
keyId: config.foundryProvisionerKeyId,
audience: config.foundryProvisionerAudience,
maxSkewMs: config.foundryProvisionerMaxSkewMs,
replayCache: foundryProvisionerReplayCache,
})
: null;
const pool = new Pool({ connectionString: config.databaseUrl, max: config.databasePoolSize });
const app = express();
const httpServer = createServer(app);
@@ -106,6 +123,8 @@ app.get("/healthz", asyncRoute(async (_req, res) => {
managedWriterBindingLifetime: config.managedProvisionerApiEnabled ? "explicit-revoke" : "disabled",
managedReaderBindingProvisioning: config.managedProvisionerApiEnabled ? "enabled" : "disabled",
managedReaderBindingLifetime: config.managedProvisionerApiEnabled ? "explicit-revoke" : "disabled",
foundryReaderBindingProvisioning: config.foundryProvisionerApiEnabled ? "digest+server-resolved-source" : "disabled",
foundryReaderBindingLifetime: config.foundryProvisionerApiEnabled ? "explicit-revoke" : "disabled",
rawRetentionSweep: {
mode: "server-scheduled",
lastSweepAt: lastRetentionSweepAt,
@@ -378,6 +397,143 @@ app.post("/internal/data-plane/v1/writer-bindings/:bindingId/revoke", requirePro
res.json({ ok: true, writerBinding: safeWriterBinding(result.rows[0]) });
}));
app.post("/internal/data-plane/v1/consumer-reader-bindings/plan", requireFoundryProvisionerApi, asyncRoute(async (req, res) => {
const policy = normalizeManagedConsumerReaderPlanRequest(req.body);
await assertRegisteredProductIds(policy.allowedDataProductIds);
await resolveManagedReaderScope(pool, policy.allowedDataProductIds);
const dataProducts = [];
for (const dataProductId of policy.allowedDataProductIds) {
const definition = await loadDataProductDefinition(pool, dataProductId);
if (!definition || definition.active === false) throw httpError(404, "data_product_not_found");
dataProducts.push(safeDataProductDefinition(definition));
}
res.set("Cache-Control", "no-store, max-age=0");
res.json({ ok: true, sourceScope: "resolved-server-side", dataProducts });
}));
app.put("/internal/data-plane/v1/consumer-reader-bindings/by-key/:bindingKey", requireFoundryProvisionerApi, asyncRoute(async (req, res) => {
const bindingKey = requireIdentifier(req.params.bindingKey, "managed_consumer_reader_binding_key_invalid");
const policy = normalizeManagedConsumerReaderBindingRequest(req.body);
await assertRegisteredProductIds(policy.allowedDataProductIds);
const requestHash = consumerReaderBindingRequestHash(policy);
const client = await pool.connect();
let response;
try {
await client.query("begin");
await client.query("select pg_advisory_xact_lock(hashtextextended($1, 0))", [bindingKey]);
const existing = await client.query(
`select id, binding_key as "bindingKey", request_hash as "requestHash", generation,
tenant_id as "tenantId", connection_id as "connectionId",
source_connection_id as "sourceConnectionId", provider_id as "providerId",
allowed_data_product_ids as "allowedDataProductIds", expires_at as "expiresAt",
active, created_at as "createdAt", rotated_at as "rotatedAt", revoked_at as "revokedAt"
from external_data_plane_reader_bindings
where binding_key = $1 and generation = $2
for update`,
[bindingKey, policy.generation],
);
if (existing.rowCount) {
const binding = existing.rows[0];
if (binding.requestHash !== requestHash) throw httpError(409, "managed_consumer_reader_binding_request_conflict");
if (binding.active !== true || binding.expiresAt !== null) {
throw httpError(409, "managed_consumer_reader_binding_generation_inactive");
}
response = { status: 200, idempotent: true, binding };
} else {
const source = await resolveManagedReaderScope(client, policy.allowedDataProductIds);
const inserted = await client.query(
`insert into external_data_plane_reader_bindings (
id, token_hash, binding_key, request_hash, generation,
tenant_id, connection_id, source_connection_id, provider_id,
allowed_data_product_ids, expires_at
) values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10::jsonb, $11)
returning id, binding_key as "bindingKey", generation,
tenant_id as "tenantId", connection_id as "connectionId",
source_connection_id as "sourceConnectionId", provider_id as "providerId",
allowed_data_product_ids as "allowedDataProductIds", expires_at as "expiresAt",
active, created_at as "createdAt", rotated_at as "rotatedAt", revoked_at as "revokedAt"`,
[
randomUUID(),
policy.capabilityDigest,
bindingKey,
requestHash,
policy.generation,
source.tenantId,
source.connectionId,
source.connectionId,
source.providerId,
JSON.stringify(policy.allowedDataProductIds),
policy.expiresAt,
],
);
response = { status: 201, idempotent: false, binding: inserted.rows[0] };
}
await client.query("commit");
} catch (error) {
await client.query("rollback");
if (error?.code === "23505") throw httpError(409, "managed_consumer_reader_capability_digest_conflict");
throw error;
} finally {
client.release();
}
res.set("Cache-Control", "no-store, max-age=0");
res.status(response.status).json({
ok: true,
idempotent: response.idempotent,
readerBinding: safeManagedConsumerReaderBinding(response.binding),
});
}));
app.post("/internal/data-plane/v1/consumer-reader-bindings/by-key/:bindingKey/generations/:generation/revoke", requireFoundryProvisionerApi, asyncRoute(async (req, res) => {
const bindingKey = requireIdentifier(req.params.bindingKey, "managed_consumer_reader_binding_key_invalid");
const generation = requirePositiveInteger(req.params.generation, "managed_consumer_reader_generation_invalid");
const client = await pool.connect();
let binding;
let idempotent;
try {
await client.query("begin");
await client.query("select pg_advisory_xact_lock(hashtextextended($1, 0))", [bindingKey]);
const existing = await client.query(
`select id, binding_key as "bindingKey", generation,
tenant_id as "tenantId", connection_id as "connectionId",
source_connection_id as "sourceConnectionId", provider_id as "providerId",
allowed_data_product_ids as "allowedDataProductIds", expires_at as "expiresAt",
active, created_at as "createdAt", rotated_at as "rotatedAt", revoked_at as "revokedAt"
from external_data_plane_reader_bindings
where binding_key = $1 and generation = $2
for update`,
[bindingKey, generation],
);
if (!existing.rowCount) throw httpError(404, "managed_consumer_reader_binding_not_found");
if (existing.rows[0].active === true) {
const revoked = await client.query(
`update external_data_plane_reader_bindings
set active = false, revoked_at = now()
where binding_key = $1 and generation = $2 and active = true
returning id, binding_key as "bindingKey", generation,
tenant_id as "tenantId", connection_id as "connectionId",
source_connection_id as "sourceConnectionId", provider_id as "providerId",
allowed_data_product_ids as "allowedDataProductIds", expires_at as "expiresAt",
active, created_at as "createdAt", rotated_at as "rotatedAt", revoked_at as "revokedAt"`,
[bindingKey, generation],
);
binding = revoked.rows[0];
idempotent = false;
} else {
binding = existing.rows[0];
idempotent = true;
}
await client.query("commit");
} catch (error) {
await client.query("rollback");
throw error;
} finally {
client.release();
}
res.set("Cache-Control", "no-store, max-age=0");
res.json({ ok: true, idempotent, readerBinding: safeManagedConsumerReaderBinding(binding) });
}));
app.put("/internal/data-plane/v1/reader-bindings/by-key/:bindingKey", requireManagedProvisionerApi, asyncRoute(async (req, res) => {
const bindingKey = requireIdentifier(req.params.bindingKey, "managed_reader_binding_key_invalid");
const policy = normalizeManagedReaderBindingRequest(req.body);
@@ -918,6 +1074,17 @@ function requireManagedProvisionerApi(req, _res, next) {
}
}
function requireFoundryProvisionerApi(req, _res, next) {
if (!config.foundryProvisionerApiEnabled) return next(httpError(503, "foundry_provisioner_api_disabled"));
if (!verifyFoundryProvisionerRequest) return next(httpError(503, "foundry_provisioner_api_not_configured"));
try {
req.managedProvisionerIdentity = verifyFoundryProvisionerRequest(req);
return next();
} catch (error) {
return next(error);
}
}
function requireProvisionerCredential(req, next) {
if (!config.provisionerAccessToken) return next(httpError(503, "provisioner_api_not_configured"));
const value = bearerToken(req);
@@ -24,12 +24,22 @@ const managedProvisionerKeyId = "engine-edp-managed-provisioner-v1";
const managedProvisionerAudience = "nodedc-external-data-plane.managed-provisioning.v1";
const managedProvisionerPublicKeyPath = join(directory, "engine-public-key.pem");
const { privateKey: managedProvisionerPrivateKey, publicKey: managedProvisionerPublicKey } = generateKeyPairSync("ed25519");
const foundryProvisionerServiceId = "nodedc-module-foundry";
const foundryProvisionerKeyId = "foundry-edp-managed-provisioner-v1";
const foundryProvisionerAudience = "nodedc-external-data-plane.managed-provisioning.v1";
const foundryProvisionerPublicKeyPath = join(directory, "foundry-public-key.pem");
const { privateKey: foundryProvisionerPrivateKey, publicKey: foundryProvisionerPublicKey } = generateKeyPairSync("ed25519");
await writeFile(secretPath, `${provisionerSecret}\n`, { mode: 0o600 });
await writeFile(
managedProvisionerPublicKeyPath,
managedProvisionerPublicKey.export({ type: "spki", format: "pem" }),
{ mode: 0o600 },
);
await writeFile(
foundryProvisionerPublicKeyPath,
foundryProvisionerPublicKey.export({ type: "spki", format: "pem" }),
{ mode: 0o600 },
);
const child = spawn(process.execPath, ["src/server.mjs"], {
cwd: new URL("..", import.meta.url),
@@ -44,6 +54,11 @@ const child = spawn(process.execPath, ["src/server.mjs"], {
EXTERNAL_DATA_PLANE_MANAGED_PROVISIONER_SERVICE_ID: managedProvisionerServiceId,
EXTERNAL_DATA_PLANE_MANAGED_PROVISIONER_KEY_ID: managedProvisionerKeyId,
EXTERNAL_DATA_PLANE_MANAGED_PROVISIONER_AUDIENCE: managedProvisionerAudience,
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONING_ENABLED: "true",
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_PUBLIC_KEY_FILE: foundryProvisionerPublicKeyPath,
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_SERVICE_ID: foundryProvisionerServiceId,
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_KEY_ID: foundryProvisionerKeyId,
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_AUDIENCE: foundryProvisionerAudience,
EXTERNAL_DATA_PLANE_RETENTION_SWEEP_MS: "60000",
EXTERNAL_DATA_PLANE_STREAM_POLL_MS: "250",
EXTERNAL_DATA_PLANE_STREAM_HEARTBEAT_MS: "5000",
@@ -249,6 +264,62 @@ try {
});
assert.equal(rejectedManagedReader.status, 401);
const foundryPlanPath = "/internal/data-plane/v1/consumer-reader-bindings/plan";
const foundryPlanBody = { allowedDataProductIds: [productId] };
const engineCannotProvisionFoundryGrant = await signedFetch(foundryPlanPath, {
method: "POST",
body: foundryPlanBody,
});
assert.equal(engineCannotProvisionFoundryGrant.status, 401);
assert.equal((await engineCannotProvisionFoundryGrant.json()).error, "managed_provisioner_identity_mismatch");
const foundryPlan = await jsonRequest(foundryPlanPath, {
method: "POST",
foundrySignature: true,
body: foundryPlanBody,
});
assert.equal(foundryPlan.sourceScope, "resolved-server-side");
assert.deepEqual(foundryPlan.dataProducts.map((value) => value.id), [productId]);
assert.equal(/tenant|provider|connection/i.test(JSON.stringify(foundryPlan)), false);
const foundryReaderToken = "ndc_edprb_foundry_api_test_0123456789abcdefghijklmnopqrstuvwxyz";
const foundryReaderBody = {
allowedDataProductIds: [productId],
expiresAt: null,
generation: 1,
capabilityDigest: createHash("sha256").update(foundryReaderToken, "utf8").digest("hex"),
};
const foundryReaderPath = "/internal/data-plane/v1/consumer-reader-bindings/by-key/fndrc-api-test-positions";
const foundryReaderResults = await Promise.all([
rawJsonRequest(foundryReaderPath, { method: "PUT", foundrySignature: true, body: foundryReaderBody }),
rawJsonRequest(foundryReaderPath, { method: "PUT", foundrySignature: true, body: foundryReaderBody }),
]);
assert.deepEqual(foundryReaderResults.map(({ value }) => value.idempotent).sort(), [false, true]);
assert.equal(foundryReaderResults[0].value.readerBinding.id, foundryReaderResults[1].value.readerBinding.id);
assert.equal(JSON.stringify(foundryReaderResults[0].value).includes(foundryReaderBody.capabilityDigest), false);
assert.equal(/tenant|provider|connection/i.test(JSON.stringify(foundryReaderResults[0].value)), false);
assert.equal("token" in foundryReaderResults[0].value, false);
const foundrySnapshot = await jsonRequest(
`/internal/data-plane/v1/data-products/${productId}/snapshot`,
{ token: foundryReaderToken },
);
assert.equal(validateDataProductSnapshot(foundrySnapshot).ok, true);
assert.equal(foundrySnapshot.facts.length, 1);
const foundryRevoke = await jsonRequest(`${foundryReaderPath}/generations/1/revoke`, {
method: "POST",
foundrySignature: true,
});
assert.equal(foundryRevoke.idempotent, false);
assert.equal(foundryRevoke.readerBinding.active, false);
const foundryRevokeRetry = await jsonRequest(`${foundryReaderPath}/generations/1/revoke`, {
method: "POST",
foundrySignature: true,
});
assert.equal(foundryRevokeRetry.idempotent, true);
const rejectedFoundryReader = await fetch(`${baseUrl}/internal/data-plane/v1/reader/data-products`, {
headers: { Authorization: `Bearer ${foundryReaderToken}` },
});
assert.equal(rejectedFoundryReader.status, 401);
const secondWriterToken = "ndc_edpwb_second_source_api_test_0123456789abcdefghijklmnopqrstuvwxyz";
await jsonRequest("/internal/data-plane/v1/writer-bindings/by-key/engine.api-connection-2.positions", {
method: "PUT",
@@ -259,6 +330,12 @@ try {
capabilityDigest: createHash("sha256").update(secondWriterToken, "utf8").digest("hex"),
},
});
const ambiguousFoundryPlan = await foundrySignedFetch(foundryPlanPath, {
method: "POST",
body: foundryPlanBody,
});
assert.equal(ambiguousFoundryPlan.status, 409);
assert.equal((await ambiguousFoundryPlan.json()).error, "managed_consumer_reader_source_scope_ambiguous");
const ambiguousReaderToken = "ndc_edprb_ambiguous_api_test_0123456789abcdefghijklmnopqrstuvwxyz";
const ambiguousReader = await signedFetch(
"/internal/data-plane/v1/reader-bindings/by-key/engine.api-ambiguous.positions-reader",
@@ -413,19 +490,20 @@ function publish(runId, observedAt, longitude) {
};
}
async function jsonRequest(path, { method = "GET", token, managedSignature = false, body } = {}) {
const { response, value } = await rawJsonRequest(path, { method, token, managedSignature, body });
async function jsonRequest(path, { method = "GET", token, managedSignature = false, foundrySignature = false, body } = {}) {
const { response, value } = await rawJsonRequest(path, { method, token, managedSignature, foundrySignature, body });
if (!response.ok) throw new Error(`request_failed:${response.status}:${value.error}`);
return value;
}
async function rawJsonRequest(path, { method = "GET", token, managedSignature = false, body } = {}) {
async function rawJsonRequest(path, { method = "GET", token, managedSignature = false, foundrySignature = false, body } = {}) {
const rawBody = body === undefined ? Buffer.alloc(0) : Buffer.from(JSON.stringify(body), "utf8");
const response = await fetch(`${baseUrl}${path}`, {
method,
headers: {
...(token ? { Authorization: `Bearer ${token}` } : {}),
...(managedSignature ? signedManagedHeaders(path, method, rawBody) : {}),
...(foundrySignature ? signedFoundryHeaders(path, method, rawBody) : {}),
...(body === undefined ? {} : { "Content-Type": "application/json" }),
},
body: body === undefined ? undefined : rawBody,
@@ -447,6 +525,18 @@ async function signedFetch(path, { method, body } = {}) {
});
}
async function foundrySignedFetch(path, { method, body } = {}) {
const rawBody = body === undefined ? Buffer.alloc(0) : Buffer.from(JSON.stringify(body), "utf8");
return fetch(`${baseUrl}${path}`, {
method,
headers: {
...signedFoundryHeaders(path, method, rawBody),
...(body === undefined ? {} : { "Content-Type": "application/json" }),
},
body: body === undefined ? undefined : rawBody,
});
}
function signedManagedHeaders(path, method, rawBody) {
const timestamp = new Date().toISOString();
const nonce = randomBytes(24).toString("base64url");
@@ -473,6 +563,44 @@ function signedManagedHeaders(path, method, rawBody) {
};
}
function signedFoundryHeaders(path, method, rawBody) {
return signedHeaders({
path,
method,
rawBody,
audience: foundryProvisionerAudience,
serviceId: foundryProvisionerServiceId,
keyId: foundryProvisionerKeyId,
privateKey: foundryProvisionerPrivateKey,
});
}
function signedHeaders({ path, method, rawBody, audience, serviceId, keyId, privateKey }) {
const timestamp = new Date().toISOString();
const nonce = randomBytes(24).toString("base64url");
const bodySha256 = sha256RawBody(rawBody);
const payload = managedProvisionerSigningPayload({
audience,
serviceId,
keyId,
method,
path,
timestamp,
nonce,
bodySha256,
});
const signature = sign(null, Buffer.from(payload, "utf8"), privateKey).toString("base64url");
return {
[MANAGED_PROVISIONER_HEADERS.serviceId]: serviceId,
[MANAGED_PROVISIONER_HEADERS.keyId]: keyId,
[MANAGED_PROVISIONER_HEADERS.audience]: audience,
[MANAGED_PROVISIONER_HEADERS.timestamp]: timestamp,
[MANAGED_PROVISIONER_HEADERS.nonce]: nonce,
[MANAGED_PROVISIONER_HEADERS.bodySha256]: bodySha256,
[MANAGED_PROVISIONER_HEADERS.signature]: signature,
};
}
function assertOneTimeCapabilityResponse(response) {
assert.equal(response.headers.get("cache-control"), "no-store, max-age=0");
assert.equal(response.headers.get("pragma"), "no-cache");
@@ -66,6 +66,22 @@ try {
assert.equal(managedConfig.managedProvisionerAudience, "nodedc-external-data-plane.managed-provisioning.v1");
assert.equal(managedConfig.managedProvisionerMaxSkewMs, 45_000);
assert.equal(managedConfig.managedProvisionerReplayCacheMaxEntries, 500);
const foundryConfig = readConfig({
...base,
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONING_ENABLED: "true",
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_PUBLIC_KEY_FILE: publicKeyPath,
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_SERVICE_ID: "nodedc-module-foundry",
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_KEY_ID: "foundry-edp-managed-provisioner-v1",
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_AUDIENCE: "nodedc-external-data-plane.managed-provisioning.v1",
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_MAX_SKEW_SECONDS: "30",
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_REPLAY_CACHE_MAX_ENTRIES: "600",
});
assert.equal(foundryConfig.foundryProvisionerApiEnabled, true);
assert.equal(foundryConfig.foundryProvisionerPublicKey.asymmetricKeyType, "ed25519");
assert.equal(foundryConfig.foundryProvisionerServiceId, "nodedc-module-foundry");
assert.equal(foundryConfig.foundryProvisionerKeyId, "foundry-edp-managed-provisioner-v1");
assert.equal(foundryConfig.foundryProvisionerMaxSkewMs, 30_000);
assert.equal(foundryConfig.foundryProvisionerReplayCacheMaxEntries, 600);
assert.throws(() => readConfig({
...base,
EXTERNAL_DATA_PLANE_MANAGED_PROVISIONING_ENABLED: "true",
@@ -88,6 +104,7 @@ try {
EXTERNAL_DATA_PLANE_MANAGED_PROVISIONING_ENABLED: "false",
EXTERNAL_DATA_PLANE_PROVISIONER_TOKEN_FILE: join(directory, "missing"),
EXTERNAL_DATA_PLANE_MANAGED_PROVISIONER_PUBLIC_KEY_FILE: join(directory, "missing"),
EXTERNAL_DATA_PLANE_FOUNDRY_PROVISIONER_PUBLIC_KEY_FILE: join(directory, "missing"),
}).provisionerAccessToken, "");
} finally {
await rm(directory, { recursive: true, force: true });
@@ -1,11 +1,15 @@
import assert from "node:assert/strict";
import {
assertReaderProduct,
consumerReaderBindingRequestHash,
createReaderToken,
hashReaderToken,
normalizeManagedConsumerReaderBindingRequest,
normalizeManagedConsumerReaderPlanRequest,
normalizeManagedReaderBindingRequest,
normalizeReaderBindingRequest,
readerBindingRequestHash,
safeManagedConsumerReaderBinding,
safeReaderBinding,
} from "../src/reader-binding.mjs";
@@ -64,4 +68,31 @@ assert.throws(() => normalizeManagedReaderBindingRequest({
capabilityDigest: managed.capabilityDigest,
}), /managed_reader_binding_must_be_durable/);
const consumerPlan = normalizeManagedConsumerReaderPlanRequest({
allowedDataProductIds: ["fleet.positions.current.v1", "fleet.positions.current.v1"],
});
assert.deepEqual(consumerPlan.allowedDataProductIds, ["fleet.positions.current.v1"]);
const consumerManaged = normalizeManagedConsumerReaderBindingRequest({
allowedDataProductIds: consumerPlan.allowedDataProductIds,
expiresAt: null,
generation: 1,
capabilityDigest: hashReaderToken(token),
});
assert.match(consumerReaderBindingRequestHash(consumerManaged), /^[a-f0-9]{64}$/);
const safeConsumer = safeManagedConsumerReaderBinding({
id: "reader-foundry-01",
bindingKey: "fndrc-reader-managed",
...consumerManaged,
active: true,
});
assert.equal(safeConsumer.sourceScope, "resolved-server-side");
assert.equal(JSON.stringify(safeConsumer).includes(consumerManaged.capabilityDigest), false);
assert.equal("tenantId" in safeConsumer, false);
assert.equal("connectionId" in safeConsumer, false);
assert.equal("providerId" in safeConsumer, false);
assert.throws(() => normalizeManagedConsumerReaderBindingRequest({
...consumerManaged,
capability: "forbidden",
}), /managed_consumer_reader_binding_request_invalid/);
console.log("external-data-plane reader bindings: ok");
@@ -1,5 +1,5 @@
import assert from "node:assert/strict";
import { resolveManagedReaderSourceConnection } from "../src/reader-source-scope.mjs";
import { resolveManagedReaderScope, resolveManagedReaderSourceConnection } from "../src/reader-source-scope.mjs";
const policy = {
tenantId: "tenant-01",
@@ -20,6 +20,24 @@ await assert.rejects(
resolveManagedReaderSourceConnection(db([]), policy),
(error) => error?.status === 409 && error?.code === "managed_reader_source_scope_not_found",
);
assert.deepEqual(
await resolveManagedReaderScope(scopeDb([
{ tenantId: "tenant-01", connectionId: "producer-connection", providerId: "provider-01" },
]), policy.allowedDataProductIds),
{ tenantId: "tenant-01", connectionId: "producer-connection", providerId: "provider-01" },
);
await assert.rejects(
resolveManagedReaderScope(scopeDb([]), policy.allowedDataProductIds),
(error) => error?.status === 409 && error?.code === "managed_consumer_reader_source_scope_not_found",
);
await assert.rejects(
resolveManagedReaderScope(scopeDb([
{ tenantId: "tenant-01", connectionId: "producer-a", providerId: "provider-01" },
{ tenantId: "tenant-01", connectionId: "producer-b", providerId: "provider-01" },
]), policy.allowedDataProductIds),
(error) => error?.status === 409 && error?.code === "managed_consumer_reader_source_scope_ambiguous",
);
await assert.rejects(
resolveManagedReaderSourceConnection(db(["producer-a", "producer-b"]), policy),
(error) => error?.status === 409 && error?.code === "managed_reader_source_scope_ambiguous",
@@ -39,4 +57,14 @@ function db(connectionIds) {
};
}
function scopeDb(rows) {
return {
async query(sql, values) {
assert.match(sql, /external_data_plane_writer_bindings/);
assert.deepEqual(values, [policy.allowedDataProductIds]);
return { rows };
},
};
}
console.log("external-data-plane reader source scope: ok");