From b1fa9cfcf4641e5ff025e2f504fc9f9825ecdbaa Mon Sep 17 00:00:00 2001 From: Nathan Beddoe Date: Sun, 13 Sep 2026 04:18:11 +0200 Subject: [PATCH] feat: persist hosted MCP control plane --- .github/workflows/ci.yml | 2 + .../hosted-mcp-registry-authority.ts | 345 +++++++ control-plane/hosted-mcp-registry-codec.ts | 698 ++++++++++++++ control-plane/hosted-mcp-registry-database.ts | 648 +++++++++++++ .../hosted-mcp-registry-projection.ts | 77 ++ control-plane/hosted-mcp-registry-store.ts | 700 ++++++++++++++ .../hosted-mcp-registry-transitions.ts | 683 ++++++++++++++ control-plane/hosted-mcp-registry-types.ts | 284 ++++++ control-plane/installation-registry.ts | 184 +++- package.json | 1 + tests/fixtures/hosted-mcp-registry-worker.ts | 266 ++++++ tests/helpers/hosted-mcp-registry-harness.mjs | 331 +++++++ ...osted-mcp-persistence-adversarial.test.mjs | 621 +++++++++++++ tests/hosted-mcp-persistence.test.mjs | 852 ++++++++++++++++++ tsconfig.worker.json | 1 + 15 files changed, 5692 insertions(+), 1 deletion(-) create mode 100644 control-plane/hosted-mcp-registry-authority.ts create mode 100644 control-plane/hosted-mcp-registry-codec.ts create mode 100644 control-plane/hosted-mcp-registry-database.ts create mode 100644 control-plane/hosted-mcp-registry-projection.ts create mode 100644 control-plane/hosted-mcp-registry-store.ts create mode 100644 control-plane/hosted-mcp-registry-transitions.ts create mode 100644 control-plane/hosted-mcp-registry-types.ts create mode 100644 tests/fixtures/hosted-mcp-registry-worker.ts create mode 100644 tests/helpers/hosted-mcp-registry-harness.mjs create mode 100644 tests/hosted-mcp-persistence-adversarial.test.mjs create mode 100644 tests/hosted-mcp-persistence.test.mjs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index e5c9e45..97a972f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -71,6 +71,8 @@ jobs: - run: pnpm test:oauth - run: pnpm test:bridge-server - run: pnpm test:ownership + - name: Validate native hosted MCP control-plane persistence + run: pnpm test:hosted-mcp-persistence - run: pnpm test:orchestrator - run: pnpm test:installation-status - run: pnpm test:updates diff --git a/control-plane/hosted-mcp-registry-authority.ts b/control-plane/hosted-mcp-registry-authority.ts new file mode 100644 index 0000000..8420615 --- /dev/null +++ b/control-plane/hosted-mcp-registry-authority.ts @@ -0,0 +1,345 @@ +import * as Effect from "effect/Effect"; +import { fingerprint } from "./deployment-values.ts"; +import { + parseHostedMcpCreateInput, + materializeHostedMcpObservation, + parseHostedMcpHealthInput, + parseHostedMcpLifecycleInput, + parseHostedMcpObservationClaim, + parseHostedMcpRevision, + parseHostedMcpScopeId, + parseHostedMcpServerId, + invalidHostedMcpRecord, +} from "./hosted-mcp-registry-codec.ts"; +import { + hostedMcpCreatePayload, + hostedMcpLifecyclePayload, + hostedMcpUpdatePayload, + HostedMcpRegistryStore, +} from "./hosted-mcp-registry-store.ts"; +import { + parseHostedMcpDeploymentIntent, + parseHostedMcpTemplate, +} from "./hosted-mcp-template.ts"; +import type { + HostedMcpInstallationCoordination, + HostedMcpInventoryPage, + HostedMcpOperationState, + HostedMcpPrivilegedInventory, + HostedMcpServerProjection, +} from "./hosted-mcp-registry-types.ts"; + +async function digest(value: unknown) { + return Effect.runPromise(fingerprint(value)); +} + +export class HostedMcpRegistryAuthority { + private readonly store: HostedMcpRegistryStore; + + constructor(storage: DurableObjectStorage) { + this.store = new HostedMcpRegistryStore(storage); + } + + private scope(installationValue: unknown, accountValue: unknown) { + return { + installationId: parseHostedMcpScopeId(installationValue), + accountId: parseHostedMcpScopeId(accountValue), + }; + } + + private async createInput(manifestValue: unknown, intentValue: unknown) { + let manifest; + let intent; + try { + manifest = parseHostedMcpTemplate(manifestValue); + intent = parseHostedMcpDeploymentIntent(manifest, intentValue); + } catch { + return invalidHostedMcpRecord(); + } + return parseHostedMcpCreateInput( + manifest, + intent, + await digest(intent.configuration), + ); + } + + async create( + subject: string, + installationValue: unknown, + accountValue: unknown, + manifestValue: unknown, + intentValue: unknown, + ): Promise { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + const input = await this.createInput(manifestValue, intentValue); + const payload = hostedMcpCreatePayload(installationId, accountId, input); + return this.store.create( + subject, + installationId, + accountId, + input, + await digest(payload), + ); + } + + async update( + subject: string, + installationValue: unknown, + accountValue: unknown, + serverValue: unknown, + expectedRevisionValue: unknown, + manifestValue: unknown, + intentValue: unknown, + ): Promise { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + const serverId = parseHostedMcpServerId(serverValue); + const expectedRevision = parseHostedMcpRevision(expectedRevisionValue); + const input = await this.createInput(manifestValue, intentValue); + const payload = hostedMcpUpdatePayload( + installationId, + accountId, + serverId, + expectedRevision, + input, + ); + return this.store.update( + subject, + installationId, + accountId, + serverId, + expectedRevision, + input, + await digest(payload), + ); + } + + async beginLifecycle( + subject: string, + installationValue: unknown, + accountValue: unknown, + serverValue: unknown, + expectedRevisionValue: unknown, + inputValue: unknown, + ): Promise { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + const serverId = parseHostedMcpServerId(serverValue); + const expectedRevision = parseHostedMcpRevision(expectedRevisionValue); + const input = parseHostedMcpLifecycleInput(inputValue); + const payload = hostedMcpLifecyclePayload( + installationId, + accountId, + serverId, + expectedRevision, + input, + ); + return this.store.beginLifecycle( + subject, + installationId, + accountId, + serverId, + expectedRevision, + input, + await digest(payload), + ); + } + + get( + subject: string, + installationValue: unknown, + accountValue: unknown, + serverValue: unknown, + ): HostedMcpServerProjection | null { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + return this.store.get( + subject, + installationId, + accountId, + parseHostedMcpServerId(serverValue), + ); + } + + list( + subject: string, + installationValue: unknown, + accountValue: unknown, + cursorValue: unknown, + ): HostedMcpInventoryPage { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + const cursor = + cursorValue === null ? null : parseHostedMcpServerId(cursorValue); + return this.store.list(subject, installationId, accountId, cursor); + } +} + +abstract class HostedMcpInternalAuthority { + protected readonly store: HostedMcpRegistryStore; + + constructor(storage: DurableObjectStorage) { + this.store = new HostedMcpRegistryStore(storage); + } + + protected scope(installationValue: unknown, accountValue: unknown) { + return { + installationId: parseHostedMcpScopeId(installationValue), + accountId: parseHostedMcpScopeId(accountValue), + }; + } +} + +// This authority is deliberately not installed as an InstallationRegistry RPC. +// FLA-61 orchestration must put its authenticated grant check in front of this +// capability; owner methods alone can never write Cloudflare resource identity. +export class HostedMcpRegistryReconciler extends HostedMcpInternalAuthority { + async recordObservation( + subject: string, + installationValue: unknown, + accountValue: unknown, + serverValue: unknown, + inputValue: unknown, + ): Promise { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + const serverId = parseHostedMcpServerId(serverValue); + const claim = parseHostedMcpObservationClaim(inputValue); + const context = this.store.reconciliationContext( + subject, + installationId, + accountId, + serverId, + claim.operationId, + ); + const identityDigests = new Map(); + await Promise.all( + claim.resources.map(async ({ key, provider }) => { + if (provider) + identityDigests.set( + key, + await digest({ + serverId, + installationId, + accountId, + key, + provider, + }), + ); + }), + ); + return this.store.recordObservation( + subject, + installationId, + accountId, + serverId, + materializeHostedMcpObservation( + context.server, + claim, + identityDigests, + new Date().toISOString(), + ), + ); + } + + recordHealth( + subject: string, + installationValue: unknown, + accountValue: unknown, + serverValue: unknown, + inputValue: unknown, + ): HostedMcpServerProjection { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + const serverId = parseHostedMcpServerId(serverValue); + const server = this.store.reconciliationServer( + subject, + installationId, + accountId, + serverId, + ); + return this.store.recordHealth( + subject, + installationId, + accountId, + serverId, + parseHostedMcpHealthInput(server, inputValue), + ); + } +} + +// Provider identities are available only to trusted cleanup/orchestration code, +// never from the owner/customer RPC projection. +export class HostedMcpRegistryCoordinator extends HostedMcpInternalAuthority { + coordinateInstallation( + subject: string, + installationValue: unknown, + accountValue: unknown, + action: "upgrade" | "remove", + ): HostedMcpInstallationCoordination { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + if (action !== "upgrade" && action !== "remove") + return invalidHostedMcpRecord(); + return this.store.coordinateInstallation( + subject, + installationId, + accountId, + action, + ); + } + + finishRemoval( + subject: string, + installationValue: unknown, + accountValue: unknown, + ): HostedMcpInstallationCoordination { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + return this.store.finishInstallationRemoval( + subject, + installationId, + accountId, + ); + } + + inventory( + subject: string, + installationValue: unknown, + accountValue: unknown, + serverValue: unknown, + retainedOnly: boolean, + ): HostedMcpPrivilegedInventory { + const { installationId, accountId } = this.scope( + installationValue, + accountValue, + ); + if (typeof retainedOnly !== "boolean") return invalidHostedMcpRecord(); + return this.store.privilegedInventory( + subject, + installationId, + accountId, + serverValue === null ? null : parseHostedMcpServerId(serverValue), + retainedOnly, + ); + } +} diff --git a/control-plane/hosted-mcp-registry-codec.ts b/control-plane/hosted-mcp-registry-codec.ts new file mode 100644 index 0000000..f13e9aa --- /dev/null +++ b/control-plane/hosted-mcp-registry-codec.ts @@ -0,0 +1,698 @@ +import { z } from "zod"; +import { + parseHostedMcpDeploymentIntent, + parseHostedMcpDesiredState, + parseHostedMcpObservedState, + parseHostedMcpTemplate, +} from "./hosted-mcp-template.ts"; +import { + HOSTED_MCP_REGISTRY_SCHEMA_VERSION, + HostedMcpRegistryError, + type HostedMcpCreateInput, + type HostedMcpDeploymentProvider, + type HostedMcpHealthInput, + type HostedMcpLifecycleInput, + type HostedMcpObservationClaim, + type HostedMcpObservationInput, + type HostedMcpOperation, + type HostedMcpOwnedResource, + type HostedMcpProviderIdentity, + type HostedMcpResourcePlan, + type HostedMcpServer, +} from "./hosted-mcp-registry-types.ts"; +import type { + HostedMcpDesiredState, + HostedMcpObservedState, + HostedMcpReleaseIdentity, + HostedMcpTemplateManifest, +} from "../shared/hosted-mcp.ts"; +import { safeJsonSnapshot } from "../shared/safe-json.ts"; + +const identifier = z.string().regex(/^[a-f0-9]{32}$/); +const serverId = z.string().regex(/^hosted-mcp-[a-f0-9]{32}$/); +const revision = z.number().int().positive().max(Number.MAX_SAFE_INTEGER); +const sequence = z.number().int().nonnegative().max(Number.MAX_SAFE_INTEGER); +const timestamp = z.iso.datetime(); +const digest = z.string().regex(/^[a-f0-9]{64}$/); +const providerId = z + .string() + .regex( + /^(?:[a-f0-9]{32}|[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12})$/, + ); +const providerName = z + .string() + .min(1) + .max(63) + .regex(/^[a-z](?:[a-z0-9-]*[a-z0-9])?$/); +const className = z + .string() + .min(1) + .max(64) + .regex(/^[A-Z][A-Za-z0-9]*$/); +const workerScriptName = z.string().regex(/^flarebot-hosted-mcp-[a-f0-9]{32}$/); +const workersDevHostname = z + .string() + .max(253) + .regex( + /^flarebot-hosted-mcp-[a-f0-9]{32}\.[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.workers\.dev$/, + ); +const workersAccountSubdomain = z + .string() + .min(1) + .max(63) + .regex(/^[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?$/); +const nullableEndpoint = z.string().max(2048).nullable(); +const ownerSubject = z + .string() + .min(1) + .max(512) + .refine((value) => !!value.trim() && !/[\u0000-\u001f\u007f]/.test(value)); + +const operationKind = z.enum([ + "create", + "update", + "suspend", + "resume", + "delete", + "recover", +]); + +const operationStatus = z.enum([ + "pending", + "succeeded", + "failed", + "recovery-required", +]); + +const deploymentState = z.enum([ + "absent", + "deploying", + "provisioned", + "updating", + "suspending", + "suspended", + "deleting", + "deleted", + "failed", + "recovery-required", +]); + +const providerIdentity = z.discriminatedUnion("kind", [ + z.strictObject({ + kind: z.literal("r2-bucket"), + accountId: identifier, + bucketName: providerName, + }), + z.strictObject({ + kind: z.literal("durable-object"), + accountId: identifier, + namespaceId: providerId, + namespaceName: providerName, + className, + }), +]); + +const deploymentProvider = z.discriminatedUnion("kind", [ + z.strictObject({ + kind: z.literal("worker"), + accountId: identifier, + scriptName: workerScriptName, + versionId: providerId, + hostname: workersDevHostname, + }), + z.strictObject({ + kind: z.literal("container"), + accountId: identifier, + scriptName: workerScriptName, + versionId: providerId, + hostname: workersDevHostname, + applicationId: providerId, + applicationName: providerName, + imageDigest: digest, + namespaceId: providerId, + namespaceName: providerName, + className, + }), +]); + +const resourcePlanSchema = z.strictObject({ + inventoryId: identifier, + key: z + .string() + .min(1) + .max(64) + .regex(/^[a-z][A-Za-z0-9]*$/), + kind: z.enum(["r2-bucket", "durable-object"]), + providerName, + className: className.nullable(), + retention: z.enum(["delete", "owner-choice"]), + action: z.enum(["preserve", "delete", "retain"]), + generation: revision, +}); + +const ownedResourceSchema = z.strictObject({ + schemaVersion: z.literal(HOSTED_MCP_REGISTRY_SCHEMA_VERSION), + inventoryId: identifier, + serverId, + key: z + .string() + .min(1) + .max(64) + .regex(/^[a-z][A-Za-z0-9]*$/), + kind: z.enum(["r2-bucket", "durable-object"]), + state: z.enum([ + "pending", + "ready", + "retained", + "deleting", + "deleted", + "unknown", + ]), + identityDigest: digest.nullable(), + cleanup: z.enum([ + "not-requested", + "pending", + "retained", + "complete", + "failed", + ]), + provider: providerIdentity.nullable(), + updatedAt: timestamp, +}); + +const storedOperationSchema = z.strictObject({ + schemaVersion: z.literal(HOSTED_MCP_REGISTRY_SCHEMA_VERSION), + operationId: identifier, + requestId: identifier, + serverId, + installationId: identifier, + ownerSubject, + accountId: identifier, + kind: operationKind, + status: operationStatus, + startRevision: revision, + observationSequence: sequence, + cutover: z.boolean(), + phase: deploymentState, + failure: z.unknown().nullable(), + payloadDigest: digest, + payload: z.unknown(), + createdAt: timestamp, + updatedAt: timestamp, +}); + +const storedServerSchema = z.strictObject({ + schemaVersion: z.literal(HOSTED_MCP_REGISTRY_SCHEMA_VERSION), + serverId, + installationId: identifier, + ownerSubject, + accountId: identifier, + accountSubdomain: workersAccountSubdomain, + name: z.string().regex(/^flarebot-hosted-mcp-[a-f0-9]{32}$/), + runtime: z.enum(["worker", "container"]), + revision, + configurationRevision: revision, + manifest: z.unknown(), + activeManifest: z.unknown().nullable(), + deploymentIntent: z.unknown(), + desiredState: z.unknown(), + observedState: z.unknown(), + deploymentProvider: z.unknown().nullable(), + resourcePlan: z.unknown(), + endpoint: nullableEndpoint, + currentOperationId: identifier.nullable(), + tombstonedAt: timestamp.nullable(), + createdAt: timestamp, + updatedAt: timestamp, +}); + +const lifecycleInputSchema = z.strictObject({ + schemaVersion: z.literal(1), + requestId: identifier, + action: z.enum(["suspend", "resume", "delete", "recover"]), + resourceDisposition: z + .array( + z.strictObject({ + key: z + .string() + .min(1) + .max(64) + .regex(/^[a-z][A-Za-z0-9]*$/), + action: z.enum(["delete", "retain"]), + }), + ) + .max(32), +}); + +const observationClaimSchema = z.strictObject({ + operationId: identifier, + expectedRevision: revision, + sequence: revision, + outcome: z.enum(["progress", "succeeded", "failed", "recovery-required"]), + phase: deploymentState, + failure: z.unknown().nullable(), + deploymentProvider: deploymentProvider.nullable(), + resources: z + .array( + z.strictObject({ + key: z + .string() + .min(1) + .max(64) + .regex(/^[a-z][A-Za-z0-9]*$/), + state: ownedResourceSchema.shape.state, + cleanup: ownedResourceSchema.shape.cleanup, + provider: providerIdentity.nullable(), + }), + ) + .max(32), +}); + +const healthInputSchema = z.strictObject({ + expectedRevision: revision, + readiness: z.unknown(), + credential: z.unknown(), + connection: z.unknown().nullable(), + failure: z.unknown().nullable(), + updatedAt: timestamp, +}); + +export function invalidHostedMcpRecord(): never { + throw new HostedMcpRegistryError("hosted_mcp_invalid"); +} + +function snapshot(value: unknown): unknown { + try { + return safeJsonSnapshot(value); + } catch { + return invalidHostedMcpRecord(); + } +} + +export function parseStoredJson(value: string): unknown { + try { + return JSON.parse(value); + } catch { + return invalidHostedMcpRecord(); + } +} + +function parseHostedMcpEndpoint( + value: unknown, + expectedPath: string, +): string | null { + const parsed = nullableEndpoint.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + if (parsed.data === null) return null; + let url: URL; + try { + url = new URL(parsed.data); + } catch { + return invalidHostedMcpRecord(); + } + if ( + url.protocol !== "https:" || + url.username || + url.password || + url.search || + url.hash || + url.pathname !== expectedPath || + url.href !== parsed.data || + url.href.length > 2048 + ) + return invalidHostedMcpRecord(); + return url.href; +} + +export function parseHostedMcpOwnedResources( + value: unknown, + expectedServerId: string, +): HostedMcpOwnedResource[] { + const parsed = z + .array(ownedResourceSchema) + .max(64) + .safeParse(snapshot(value)); + if ( + !parsed.success || + new Set(parsed.data.map(({ inventoryId }) => inventoryId)).size !== + parsed.data.length || + parsed.data.some( + (resource) => + resource.serverId !== expectedServerId || + (resource.provider !== null && + resource.provider.kind !== resource.kind) || + (resource.state === "ready" && + (!resource.identityDigest || !resource.provider)) || + (resource.state === "retained" && resource.cleanup !== "retained") || + (resource.state === "deleted" && resource.cleanup !== "complete"), + ) + ) + return invalidHostedMcpRecord(); + return parsed.data as HostedMcpOwnedResource[]; +} + +function parseHostedMcpFailure( + value: unknown, +): HostedMcpObservedState["failure"] { + if (value === null) return null; + try { + const snapshotValue = snapshot(value); + const uncertain = + typeof snapshotValue === "object" && + snapshotValue !== null && + "code" in snapshotValue && + snapshotValue.code === "outcome_uncertain"; + return parseHostedMcpObservedState({ + schemaVersion: 1, + serverId: `hosted-mcp-${"0".repeat(32)}`, + revision: 1, + runtime: "worker", + deployment: { + state: uncertain ? "recovery-required" : "failed", + activeRelease: null, + }, + readiness: { state: "unknown", checkedAt: null }, + resources: [], + credential: { + state: "missing", + revision: 1, + updatedAt: "2000-01-01T00:00:00.000Z", + }, + connection: null, + failure: snapshotValue, + updatedAt: "2000-01-01T00:00:00.000Z", + }).failure; + } catch { + return invalidHostedMcpRecord(); + } +} + +export function parseHostedMcpOperation(value: unknown): HostedMcpOperation { + const parsed = storedOperationSchema.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + const failure = parseHostedMcpFailure(parsed.data.failure); + if ( + (parsed.data.status === "pending" || parsed.data.status === "succeeded") !== + (failure === null) || + (parsed.data.cutover && + parsed.data.kind !== "update" && + parsed.data.kind !== "recover") + ) + return invalidHostedMcpRecord(); + return { ...parsed.data, failure } as HostedMcpOperation; +} + +const sameRelease = ( + left: HostedMcpReleaseIdentity | null, + right: HostedMcpReleaseIdentity | null, +) => JSON.stringify(left) === JSON.stringify(right); + +export function parseHostedMcpServer(value: unknown): HostedMcpServer { + const parsed = storedServerSchema.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + const raw = parsed.data; + let manifest: HostedMcpTemplateManifest; + let activeManifest: HostedMcpTemplateManifest | null; + let deploymentIntent; + let desiredState: HostedMcpDesiredState; + let observedState: HostedMcpObservedState; + let parsedDeploymentProvider: HostedMcpDeploymentProvider | null; + let resourcePlan: HostedMcpResourcePlan[]; + try { + manifest = parseHostedMcpTemplate(raw.manifest); + activeManifest = raw.activeManifest + ? parseHostedMcpTemplate(raw.activeManifest) + : null; + deploymentIntent = parseHostedMcpDeploymentIntent( + manifest, + raw.deploymentIntent, + ); + desiredState = parseHostedMcpDesiredState(manifest, raw.desiredState); + observedState = parseHostedMcpObservedState(raw.observedState); + parsedDeploymentProvider = raw.deploymentProvider + ? (deploymentProvider.parse( + raw.deploymentProvider, + ) as HostedMcpDeploymentProvider) + : null; + resourcePlan = z.array(resourcePlanSchema).max(32).parse(raw.resourcePlan); + } catch { + return invalidHostedMcpRecord(); + } + if ( + raw.serverId !== desiredState.serverId || + raw.serverId !== observedState.serverId || + raw.runtime !== manifest.runtime.type || + raw.runtime !== observedState.runtime || + raw.revision < desiredState.revision || + raw.revision < observedState.revision || + raw.configurationRevision > raw.revision || + !sameRelease( + activeManifest?.identity ?? null, + observedState.deployment.activeRelease, + ) || + new Set(resourcePlan.map(({ inventoryId }) => inventoryId)).size !== + resourcePlan.length || + new Set(resourcePlan.map(({ key }) => key)).size !== resourcePlan.length || + resourcePlan.some( + (resource) => + (resource.kind === "durable-object") !== + (resource.className !== null) || + (resource.action === "retain" && resource.retention !== "owner-choice"), + ) || + (parsedDeploymentProvider !== null && + deriveHostedMcpEndpoint( + { ...raw, manifest } as HostedMcpServer, + parsedDeploymentProvider, + ) !== raw.endpoint) || + (parsedDeploymentProvider === null && raw.endpoint !== null) || + (raw.tombstonedAt !== null && + (desiredState.lifecycle !== "deleted" || + observedState.deployment.state !== "deleted" || + raw.currentOperationId !== null || + raw.endpoint !== null)) + ) + return invalidHostedMcpRecord(); + return { + ...raw, + manifest, + activeManifest, + deploymentIntent, + desiredState, + observedState, + deploymentProvider: parsedDeploymentProvider, + resourcePlan, + } as HostedMcpServer; +} + +export function deriveHostedMcpEndpoint( + server: Pick< + HostedMcpServer, + "accountId" | "accountSubdomain" | "name" | "runtime" | "manifest" + >, + provider: HostedMcpDeploymentProvider, +): string { + if ( + provider.accountId !== server.accountId || + provider.scriptName !== server.name || + provider.hostname !== + `${server.name}.${server.accountSubdomain}.workers.dev` + ) + return invalidHostedMcpRecord(); + if (server.runtime === "worker") { + if (server.manifest.runtime.type !== "worker" || provider.kind !== "worker") + return invalidHostedMcpRecord(); + } else { + if ( + server.manifest.runtime.type !== "container" || + provider.kind !== "container" + ) + return invalidHostedMcpRecord(); + const containerBinding = server.manifest.runtime.gateway.bindings.find( + (binding) => binding.kind === "container", + ); + if ( + provider.applicationName !== `${server.name}-app` || + provider.namespaceName !== `${server.name}-ns` || + provider.className !== containerBinding?.className || + provider.imageDigest !== server.manifest.runtime.image.digest + ) + return invalidHostedMcpRecord(); + } + const path = server.manifest.protocol.mcpPath; + return parseHostedMcpEndpoint(`https://${provider.hostname}${path}`, path)!; +} + +export function parseHostedMcpAccountSubdomain( + runtimeOrigin: string | null, + installationId: string, +): string { + if (runtimeOrigin === null) return invalidHostedMcpRecord(); + let url: URL; + try { + url = new URL(runtimeOrigin); + } catch { + return invalidHostedMcpRecord(); + } + const prefix = `flarebot-${installationId}.`; + const suffix = ".workers.dev"; + if ( + url.protocol !== "https:" || + url.username || + url.password || + url.port || + url.pathname !== "/" || + url.search || + url.hash || + !url.hostname.startsWith(prefix) || + !url.hostname.endsWith(suffix) + ) + return invalidHostedMcpRecord(); + const subdomain = url.hostname.slice(prefix.length, -suffix.length); + const parsed = workersAccountSubdomain.safeParse(subdomain); + if (!parsed.success) return invalidHostedMcpRecord(); + return parsed.data; +} + +function assertResourceProvider( + server: HostedMcpServer, + plan: HostedMcpResourcePlan, + provider: HostedMcpProviderIdentity | null, +) { + if (provider === null) return; + if ( + provider.accountId !== server.accountId || + provider.kind !== plan.kind || + (provider.kind === "r2-bucket" && + provider.bucketName !== plan.providerName) || + (provider.kind === "durable-object" && + (provider.namespaceName !== plan.providerName || + provider.className !== plan.className)) + ) + return invalidHostedMcpRecord(); +} + +export function parseHostedMcpObservationClaim( + value: unknown, +): HostedMcpObservationClaim { + const parsed = observationClaimSchema.safeParse(snapshot(value)); + if ( + !parsed.success || + new Set(parsed.data.resources.map(({ key }) => key)).size !== + parsed.data.resources.length + ) + return invalidHostedMcpRecord(); + const failure = parseHostedMcpFailure(parsed.data.failure); + if ( + (parsed.data.outcome === "progress" || + parsed.data.outcome === "succeeded") !== + (failure === null) + ) + return invalidHostedMcpRecord(); + return { ...parsed.data, failure } as HostedMcpObservationClaim; +} + +export function materializeHostedMcpObservation( + server: HostedMcpServer, + claim: HostedMcpObservationClaim, + identityDigests: ReadonlyMap, + now: string, +): HostedMcpObservationInput { + if (claim.deploymentProvider) + deriveHostedMcpEndpoint(server, claim.deploymentProvider); + const resources = claim.resources.map((resource) => { + const plan = server.resourcePlan.find(({ key }) => key === resource.key); + if (!plan) return invalidHostedMcpRecord(); + assertResourceProvider(server, plan, resource.provider); + return parseHostedMcpOwnedResources( + [ + { + schemaVersion: 1, + inventoryId: plan.inventoryId, + serverId: server.serverId, + key: plan.key, + kind: plan.kind, + state: resource.state, + identityDigest: resource.provider + ? (identityDigests.get(resource.key) ?? invalidHostedMcpRecord()) + : null, + cleanup: resource.cleanup, + provider: resource.provider, + updatedAt: now, + }, + ], + server.serverId, + )[0]; + }); + return { ...claim, resources }; +} + +export function parseHostedMcpHealthInput( + server: HostedMcpServer, + value: unknown, +): HostedMcpHealthInput { + const parsed = healthInputSchema.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + try { + const observed = parseHostedMcpObservedState({ + ...server.observedState, + revision: server.revision + 1, + readiness: parsed.data.readiness, + credential: parsed.data.credential, + connection: parsed.data.connection, + failure: parsed.data.failure, + updatedAt: parsed.data.updatedAt, + }); + return { + expectedRevision: parsed.data.expectedRevision, + readiness: observed.readiness, + credential: observed.credential, + connection: observed.connection, + failure: observed.failure, + updatedAt: observed.updatedAt, + }; + } catch { + return invalidHostedMcpRecord(); + } +} + +export function parseHostedMcpCreateInput( + manifestValue: unknown, + intentValue: unknown, + configurationDigest: string, +): HostedMcpCreateInput { + try { + const manifest = parseHostedMcpTemplate(manifestValue); + const intent = parseHostedMcpDeploymentIntent(manifest, intentValue); + if (!digest.safeParse(configurationDigest).success) + return invalidHostedMcpRecord(); + return { manifest, intent, configurationDigest }; + } catch { + return invalidHostedMcpRecord(); + } +} + +export function parseHostedMcpLifecycleInput( + value: unknown, +): HostedMcpLifecycleInput { + const parsed = lifecycleInputSchema.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + return parsed.data; +} + +export function parseHostedMcpScopeId(value: unknown): string { + const parsed = identifier.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + return parsed.data; +} + +export function parseHostedMcpServerId(value: unknown): string { + const parsed = serverId.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + return parsed.data; +} + +export function parseHostedMcpRevision(value: unknown): number { + const parsed = revision.safeParse(snapshot(value)); + if (!parsed.success) return invalidHostedMcpRecord(); + return parsed.data; +} diff --git a/control-plane/hosted-mcp-registry-database.ts b/control-plane/hosted-mcp-registry-database.ts new file mode 100644 index 0000000..fe49a3e --- /dev/null +++ b/control-plane/hosted-mcp-registry-database.ts @@ -0,0 +1,648 @@ +import { + invalidHostedMcpRecord, + parseHostedMcpOperation, + parseHostedMcpOwnedResources, + parseHostedMcpServer, + parseStoredJson, +} from "./hosted-mcp-registry-codec.ts"; +import { + projectHostedMcpOperation, + projectHostedMcpServer, +} from "./hosted-mcp-registry-projection.ts"; +import { parseInstallation } from "./installation-metadata.ts"; +import { + HOSTED_MCP_REGISTRY_SCHEMA_VERSION, + HostedMcpRegistryError, + type HostedMcpOperation, + type HostedMcpOperationState, + type HostedMcpOwnedResource, + type HostedMcpServer, + type HostedMcpServerProjection, +} from "./hosted-mcp-registry-types.ts"; + +type ServerRow = Record & { + server_id: string; + installation_id: string; + account_id: string; + revision: number; + record_json: string; +}; + +type OperationRow = Record & { + operation_id: string; + request_id: string; + server_id: string; + installation_id: string; + account_id: string; + status: string; + observation_sequence: number; + record_json: string; +}; + +type ResourceRow = Record & { + inventory_id: string; + server_id: string; + state: string; + cleanup: string; + record_json: string; +}; + +const json = (value: unknown) => JSON.stringify(value); + +const resourceStateTransitions: Record< + HostedMcpOwnedResource["state"], + HostedMcpOwnedResource["state"][] +> = { + unknown: ["unknown", "pending", "ready", "deleting", "deleted", "retained"], + pending: ["pending", "ready", "deleting", "deleted"], + ready: ["ready", "deleting", "deleted", "retained"], + deleting: ["deleting", "deleted"], + deleted: ["deleted"], + retained: ["retained"], +}; + +const cleanupTransitions: Record< + HostedMcpOwnedResource["cleanup"], + HostedMcpOwnedResource["cleanup"][] +> = { + "not-requested": ["not-requested", "pending", "retained", "complete"], + pending: ["pending", "failed", "retained", "complete"], + failed: ["failed", "pending", "retained", "complete"], + retained: ["retained"], + complete: ["complete"], +}; + +function conflict(): never { + throw new HostedMcpRegistryError("hosted_mcp_conflict"); +} + +function missing(): never { + throw new HostedMcpRegistryError("hosted_mcp_not_found"); +} + +export class HostedMcpRegistryDatabase { + constructor(private readonly storage: DurableObjectStorage) {} + + transact(operation: () => T): T { + this.ensureSchema(); + return this.storage.transactionSync(operation); + } + + private ensureSchema() { + this.storage.transactionSync(() => { + this.storage.sql + .exec( + ` + CREATE TABLE IF NOT EXISTS flarebot_hosted_mcp_schema ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + version INTEGER NOT NULL + ) + `, + ) + .toArray(); + const versions = this.storage.sql + .exec<{ version: number }>( + "SELECT version FROM flarebot_hosted_mcp_schema WHERE singleton = 1", + ) + .toArray(); + if ( + versions.length > 1 || + (versions.length === 1 && + versions[0].version !== HOSTED_MCP_REGISTRY_SCHEMA_VERSION) + ) + throw new HostedMcpRegistryError("hosted_mcp_future_schema"); + + this.storage.sql + .exec( + ` + CREATE TABLE IF NOT EXISTS hosted_mcp_servers ( + server_id TEXT PRIMARY KEY, + installation_id TEXT NOT NULL, + account_id TEXT NOT NULL, + revision INTEGER NOT NULL, + record_json TEXT NOT NULL + ); + CREATE INDEX IF NOT EXISTS hosted_mcp_servers_installation + ON hosted_mcp_servers (installation_id, account_id, server_id); + CREATE TABLE IF NOT EXISTS hosted_mcp_operations ( + operation_id TEXT PRIMARY KEY, + request_id TEXT NOT NULL, + server_id TEXT NOT NULL, + installation_id TEXT NOT NULL, + account_id TEXT NOT NULL, + status TEXT NOT NULL, + observation_sequence INTEGER NOT NULL CHECK (observation_sequence >= 0), + record_json TEXT NOT NULL, + UNIQUE (installation_id, request_id) + ); + CREATE INDEX IF NOT EXISTS hosted_mcp_operations_server + ON hosted_mcp_operations (server_id, status, operation_id); + CREATE TABLE IF NOT EXISTS hosted_mcp_resources ( + inventory_id TEXT PRIMARY KEY, + server_id TEXT NOT NULL, + state TEXT NOT NULL, + cleanup TEXT NOT NULL, + record_json TEXT NOT NULL + ); + CREATE INDEX IF NOT EXISTS hosted_mcp_resources_server + ON hosted_mcp_resources (server_id, state, cleanup, inventory_id) + ; + CREATE TABLE IF NOT EXISTS hosted_mcp_installation_fences ( + installation_id TEXT PRIMARY KEY, + account_id TEXT NOT NULL, + owner_subject TEXT NOT NULL, + state TEXT NOT NULL CHECK (state IN ('removing', 'removed')), + updated_at TEXT NOT NULL + ) + `, + ) + .toArray(); + if (versions.length === 0) + this.storage.sql + .exec( + "INSERT INTO flarebot_hosted_mcp_schema (singleton, version) VALUES (1, ?)", + HOSTED_MCP_REGISTRY_SCHEMA_VERSION, + ) + .toArray(); + }); + } + + assertParent( + subject: string, + installationId: string, + accountId: string, + requireReady: boolean, + ) { + let installation; + try { + installation = parseInstallation( + this.storage.kv.get(`installation:${installationId}`), + ); + } catch { + return missing(); + } + if ( + installation.ownerSubject !== subject || + installation.installationId !== installationId || + installation.accountId !== accountId + ) + return missing(); + if (requireReady && installation.status !== "ready") + throw new HostedMcpRegistryError("hosted_mcp_installation_not_ready"); + return installation; + } + + assertMutationAllowed( + subject: string, + installationId: string, + accountId: string, + allowDelete: boolean, + ) { + const rows = this.storage.sql + .exec<{ + account_id: string; + owner_subject: string; + state: string; + }>( + `SELECT account_id, owner_subject, state + FROM hosted_mcp_installation_fences WHERE installation_id = ?`, + installationId, + ) + .toArray(); + if (rows.length === 0) return; + if ( + rows.length !== 1 || + rows[0].account_id !== accountId || + rows[0].owner_subject !== subject || + rows[0].state === "removed" || + !allowDelete + ) + return conflict(); + } + + readServer( + subject: string, + installationId: string, + accountId: string, + serverId: string, + ): HostedMcpServer { + const rows = this.storage.sql + .exec( + `SELECT server_id, installation_id, account_id, revision, record_json + FROM hosted_mcp_servers WHERE server_id = ?`, + serverId, + ) + .toArray(); + if (rows.length !== 1) return missing(); + const row = rows[0]; + const record = parseHostedMcpServer(parseStoredJson(row.record_json)); + if ( + record.ownerSubject !== subject || + record.installationId !== installationId || + record.accountId !== accountId || + row.server_id !== record.serverId || + row.installation_id !== record.installationId || + row.account_id !== record.accountId || + row.revision !== record.revision + ) + return missing(); + return record; + } + + readOperation(operationId: string): HostedMcpOperation { + const rows = this.storage.sql + .exec( + `SELECT operation_id, request_id, server_id, installation_id, account_id, status, observation_sequence, record_json + FROM hosted_mcp_operations WHERE operation_id = ?`, + operationId, + ) + .toArray(); + if (rows.length !== 1) return missing(); + const row = rows[0]; + const record = parseHostedMcpOperation(parseStoredJson(row.record_json)); + if ( + row.operation_id !== record.operationId || + row.request_id !== record.requestId || + row.server_id !== record.serverId || + row.installation_id !== record.installationId || + row.account_id !== record.accountId || + row.status !== record.status || + row.observation_sequence !== record.observationSequence + ) + return invalidHostedMcpRecord(); + return record; + } + + currentOperation(server: HostedMcpServer) { + if (!server.currentOperationId) return null; + const operation = this.readOperation(server.currentOperationId); + if (operation.serverId !== server.serverId) return invalidHostedMcpRecord(); + return operation; + } + + projection(server: HostedMcpServer) { + return projectHostedMcpServer(server, this.currentOperation(server)); + } + + operationForRequest(installationId: string, requestId: string) { + const rows = this.storage.sql + .exec( + `SELECT operation_id, request_id, server_id, installation_id, account_id, status, observation_sequence, record_json + FROM hosted_mcp_operations + WHERE installation_id = ? AND request_id = ?`, + installationId, + requestId, + ) + .toArray(); + if (rows.length === 0) return null; + if (rows.length !== 1) return invalidHostedMcpRecord(); + return this.readOperation(rows[0].operation_id); + } + + replay( + subject: string, + installationId: string, + accountId: string, + requestId: string, + payloadDigest: string, + payload: unknown, + ): HostedMcpOperationState | null { + const operation = this.operationForRequest(installationId, requestId); + if (!operation) return null; + const server = this.readServer( + subject, + installationId, + accountId, + operation.serverId, + ); + if ( + operation.payloadDigest !== payloadDigest || + operation.installationId !== installationId || + operation.accountId !== accountId || + operation.ownerSubject !== subject || + json(operation.payload) !== json(payload) + ) + return conflict(); + return { + server: this.projection(server), + operation: projectHostedMcpOperation(operation), + }; + } + + insertServer(server: HostedMcpServer) { + const value = parseHostedMcpServer(server); + const rows = this.storage.sql + .exec<{ server_id: string }>( + `INSERT INTO hosted_mcp_servers + (server_id, installation_id, account_id, revision, record_json) + VALUES (?, ?, ?, ?, ?) RETURNING server_id`, + value.serverId, + value.installationId, + value.accountId, + value.revision, + json(value), + ) + .toArray(); + if (rows.length !== 1 || rows[0].server_id !== value.serverId) + return conflict(); + } + + updateServer(server: HostedMcpServer, expectedRevision: number) { + const value = parseHostedMcpServer(server); + const rows = this.storage.sql + .exec<{ server_id: string }>( + `UPDATE hosted_mcp_servers SET revision = ?, record_json = ? + WHERE server_id = ? AND revision = ? RETURNING server_id`, + value.revision, + json(value), + value.serverId, + expectedRevision, + ) + .toArray(); + if (rows.length !== 1 || rows[0].server_id !== value.serverId) + return conflict(); + } + + insertOperation(operation: HostedMcpOperation) { + const value = parseHostedMcpOperation(operation); + const rows = this.storage.sql + .exec<{ operation_id: string }>( + `INSERT INTO hosted_mcp_operations + (operation_id, request_id, server_id, installation_id, account_id, status, observation_sequence, record_json) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) RETURNING operation_id`, + value.operationId, + value.requestId, + value.serverId, + value.installationId, + value.accountId, + value.status, + value.observationSequence, + json(value), + ) + .toArray(); + if (rows.length !== 1 || rows[0].operation_id !== value.operationId) + return conflict(); + } + + updatePendingOperation( + operation: HostedMcpOperation, + expectedObservationSequence: number, + ) { + const value = parseHostedMcpOperation(operation); + const rows = this.storage.sql + .exec<{ operation_id: string }>( + `UPDATE hosted_mcp_operations + SET status = ?, observation_sequence = ?, record_json = ? + WHERE operation_id = ? AND status = 'pending' + AND observation_sequence = ? + RETURNING operation_id`, + value.status, + value.observationSequence, + json(value), + value.operationId, + expectedObservationSequence, + ) + .toArray(); + if (rows.length !== 1 || rows[0].operation_id !== value.operationId) + return conflict(); + } + + resources(serverId: string): HostedMcpOwnedResource[] { + return this.storage.sql + .exec( + `SELECT inventory_id, server_id, state, cleanup, record_json + FROM hosted_mcp_resources WHERE server_id = ? ORDER BY inventory_id`, + serverId, + ) + .toArray() + .map((row) => { + const resource = parseHostedMcpOwnedResources( + [parseStoredJson(row.record_json)], + serverId, + )[0]; + if ( + row.inventory_id !== resource.inventoryId || + row.server_id !== resource.serverId || + row.state !== resource.state || + row.cleanup !== resource.cleanup + ) + return invalidHostedMcpRecord(); + return resource; + }); + } + + mergeResources( + server: HostedMcpServer, + operation: HostedMcpOperation, + incoming: HostedMcpOwnedResource[], + ) { + const prior = new Map( + this.resources(server.serverId).map((resource) => [ + resource.inventoryId, + resource, + ]), + ); + for (const resource of incoming) { + const plan = server.resourcePlan.find( + ({ inventoryId }) => inventoryId === resource.inventoryId, + ); + const existing = prior.get(resource.inventoryId); + const destructiveBeforeCutover = + existing !== undefined && + plan?.action === "delete" && + (operation.kind === "update" || operation.kind === "recover") && + !operation.cutover && + (existing.state !== resource.state || + existing.cleanup !== resource.cleanup); + if ( + !plan || + destructiveBeforeCutover || + plan.key !== resource.key || + plan.kind !== resource.kind || + (resource.provider !== null && + resource.provider.accountId !== server.accountId) || + (plan.action === "preserve" && + (resource.state === "retained" || + resource.state === "deleting" || + resource.state === "deleted" || + resource.cleanup !== "not-requested")) || + (plan.action === "delete" && + operation.kind !== "update" && + operation.kind !== "delete" && + operation.kind !== "recover") || + (plan.action === "retain" && + (plan.retention !== "owner-choice" || + operation.kind !== "delete" || + resource.state !== "retained" || + resource.cleanup !== "retained")) || + (plan.action !== "retain" && + (resource.state === "retained" || resource.cleanup === "retained")) || + (existing && + (existing.serverId !== resource.serverId || + existing.key !== resource.key || + existing.kind !== resource.kind || + (existing.identityDigest !== null && + existing.identityDigest !== resource.identityDigest) || + (existing.provider !== null && + json(existing.provider) !== json(resource.provider)) || + !resourceStateTransitions[existing.state].includes( + resource.state, + ) || + !cleanupTransitions[existing.cleanup].includes(resource.cleanup) || + resource.updatedAt < existing.updatedAt)) + ) + return conflict(); + const value = existing + ? { + ...resource, + identityDigest: existing.identityDigest ?? resource.identityDigest, + provider: existing.provider ?? resource.provider, + } + : resource; + const rows = this.storage.sql + .exec<{ inventory_id: string }>( + `INSERT INTO hosted_mcp_resources + (inventory_id, server_id, state, cleanup, record_json) + VALUES (?, ?, ?, ?, ?) + ON CONFLICT(inventory_id) DO UPDATE SET + state = excluded.state, + cleanup = excluded.cleanup, + record_json = excluded.record_json + WHERE hosted_mcp_resources.server_id = excluded.server_id + RETURNING inventory_id`, + value.inventoryId, + value.serverId, + value.state, + value.cleanup, + json(value), + ) + .toArray(); + if (rows.length !== 1 || rows[0].inventory_id !== value.inventoryId) + return conflict(); + prior.set(value.inventoryId, value); + } + return [...prior.values()]; + } + + serverExists(serverId: string) { + return ( + this.storage.sql + .exec<{ count: number }>( + "SELECT COUNT(*) AS count FROM hosted_mcp_servers WHERE server_id = ?", + serverId, + ) + .one().count !== 0 + ); + } + + serverIds( + installationId: string, + accountId: string, + after = "", + limit?: number, + ) { + const suffix = limit === undefined ? "" : " LIMIT ?"; + const params: SqlStorageValue[] = [installationId, accountId, after]; + if (limit !== undefined) params.push(limit); + return this.storage.sql + .exec<{ server_id: string }>( + `SELECT server_id FROM hosted_mcp_servers + WHERE installation_id = ? AND account_id = ? AND server_id > ? + ORDER BY server_id${suffix}`, + ...params, + ) + .toArray() + .map(({ server_id }) => server_id); + } + + assertNoPendingOperations(installationId: string, accountId: string) { + const count = this.pendingOperationCount(installationId, accountId); + if (count !== 0) return conflict(); + } + + pendingOperationCount(installationId: string, accountId: string) { + return this.storage.sql + .exec<{ count: number }>( + `SELECT COUNT(*) AS count + FROM hosted_mcp_operations AS operations + JOIN hosted_mcp_servers AS servers + ON servers.server_id = operations.server_id + WHERE servers.installation_id = ? + AND servers.account_id = ? + AND operations.status = 'pending'`, + installationId, + accountId, + ) + .one().count; + } + + beginRemoval( + subject: string, + installationId: string, + accountId: string, + now: string, + ) { + const rows = this.storage.sql + .exec<{ installation_id: string }>( + `INSERT INTO hosted_mcp_installation_fences + (installation_id, account_id, owner_subject, state, updated_at) + VALUES (?, ?, ?, 'removing', ?) + ON CONFLICT(installation_id) DO UPDATE SET updated_at = excluded.updated_at + WHERE hosted_mcp_installation_fences.account_id = excluded.account_id + AND hosted_mcp_installation_fences.owner_subject = excluded.owner_subject + AND hosted_mcp_installation_fences.state = 'removing' + RETURNING installation_id`, + installationId, + accountId, + subject, + now, + ) + .toArray(); + if (rows.length !== 1 || rows[0].installation_id !== installationId) + return conflict(); + } + + finishRemoval( + subject: string, + installationId: string, + accountId: string, + now: string, + ) { + const rows = this.storage.sql + .exec<{ installation_id: string }>( + `UPDATE hosted_mcp_installation_fences + SET state = 'removed', updated_at = ? + WHERE installation_id = ? AND account_id = ? AND owner_subject = ? + AND state = 'removing' + RETURNING installation_id`, + now, + installationId, + accountId, + subject, + ) + .toArray(); + if (rows.length !== 1 || rows[0].installation_id !== installationId) + return conflict(); + } + + assertRemovalFence( + subject: string, + installationId: string, + accountId: string, + ) { + const rows = this.storage.sql + .exec<{ account_id: string; owner_subject: string; state: string }>( + `SELECT account_id, owner_subject, state + FROM hosted_mcp_installation_fences WHERE installation_id = ?`, + installationId, + ) + .toArray(); + if ( + rows.length !== 1 || + rows[0].account_id !== accountId || + rows[0].owner_subject !== subject || + !["removing", "removed"].includes(rows[0].state) + ) + return missing(); + } +} diff --git a/control-plane/hosted-mcp-registry-projection.ts b/control-plane/hosted-mcp-registry-projection.ts new file mode 100644 index 0000000..898aca0 --- /dev/null +++ b/control-plane/hosted-mcp-registry-projection.ts @@ -0,0 +1,77 @@ +import type { + HostedMcpOperation, + HostedMcpOperationProjection, + HostedMcpServer, + HostedMcpServerProjection, +} from "./hosted-mcp-registry-types.ts"; + +// The customer DTO deliberately keeps FLA-58's immutable release identity, +// including source/artifact hashes, while removing configuration, catalog and +// provider-derived fingerprints. Privileged cleanup DTOs never use this seam. +export function projectHostedMcpServer( + server: HostedMcpServer, + operation: HostedMcpOperation | null, +): HostedMcpServerProjection { + return { + schemaVersion: 1, + serverId: server.serverId, + installationId: server.installationId, + accountId: server.accountId, + name: server.name, + runtime: server.runtime, + revision: server.revision, + configurationRevision: server.configurationRevision, + desired: { + schemaVersion: 1, + serverId: server.desiredState.serverId, + revision: server.desiredState.revision, + lifecycle: server.desiredState.lifecycle, + release: server.desiredState.release, + configuration: { revision: server.desiredState.configuration.revision }, + resourceDisposition: server.desiredState.resourceDisposition, + }, + observed: { + schemaVersion: 1, + serverId: server.observedState.serverId, + revision: server.observedState.revision, + runtime: server.observedState.runtime, + deployment: server.observedState.deployment, + readiness: server.observedState.readiness, + resources: server.observedState.resources.map(({ key, kind, state }) => ({ + key, + kind, + state, + })), + credential: server.observedState.credential, + connection: server.observedState.connection + ? { + id: server.observedState.connection.id, + revision: server.observedState.connection.revision, + enabled: server.observedState.connection.enabled, + state: server.observedState.connection.state, + } + : null, + failure: server.observedState.failure, + updatedAt: server.observedState.updatedAt, + }, + activeRelease: server.observedState.deployment.activeRelease, + endpoint: server.endpoint, + operation: operation ? projectHostedMcpOperation(operation) : null, + tombstonedAt: server.tombstonedAt, + createdAt: server.createdAt, + updatedAt: server.updatedAt, + }; +} + +export function projectHostedMcpOperation( + operation: HostedMcpOperation, +): HostedMcpOperationProjection { + return { + operationId: operation.operationId, + requestId: operation.requestId, + kind: operation.kind, + status: operation.status, + startRevision: operation.startRevision, + failure: operation.failure, + }; +} diff --git a/control-plane/hosted-mcp-registry-store.ts b/control-plane/hosted-mcp-registry-store.ts new file mode 100644 index 0000000..a622182 --- /dev/null +++ b/control-plane/hosted-mcp-registry-store.ts @@ -0,0 +1,700 @@ +import { + invalidHostedMcpRecord, + parseHostedMcpAccountSubdomain, + parseHostedMcpOperation, + parseHostedMcpServer, +} from "./hosted-mcp-registry-codec.ts"; +import { projectHostedMcpOperation } from "./hosted-mcp-registry-projection.ts"; +import { HostedMcpRegistryDatabase } from "./hosted-mcp-registry-database.ts"; +import { + desiredForLifecycle, + emptyObservedState, + completeObservation, + initialResourcePlan, + initialDesiredState, + lifecycleTransitionAllowed, + prepareObservation, + resourcePlanForLifecycle, + transitionHealth, + updateResourcePlan, +} from "./hosted-mcp-registry-transitions.ts"; +import { + HostedMcpRegistryError, + type HostedMcpCreateInput, + type HostedMcpInstallationCoordination, + type HostedMcpHealthInput, + type HostedMcpInventoryPage, + type HostedMcpLifecycleInput, + type HostedMcpObservationInput, + type HostedMcpOperationState, + type HostedMcpOwnedResource, + type HostedMcpPrivilegedInventory, + type HostedMcpReconciliationContext, + type HostedMcpServerProjection, +} from "./hosted-mcp-registry-types.ts"; +import { parseHostedMcpDesiredState } from "./hosted-mcp-template.ts"; + +const json = (value: unknown) => JSON.stringify(value); +const same = (left: unknown, right: unknown) => json(left) === json(right); +const nowIso = () => new Date().toISOString(); +const identifier = () => crypto.randomUUID().replaceAll("-", ""); + +function conflict(): never { + throw new HostedMcpRegistryError("hosted_mcp_conflict"); +} + +export function hostedMcpCreatePayload( + installationId: string, + accountId: string, + input: HostedMcpCreateInput, +) { + return { + schemaVersion: 1, + kind: "create" as const, + installationId, + accountId, + manifest: input.manifest, + intent: input.intent, + configurationDigest: input.configurationDigest, + }; +} + +export function hostedMcpUpdatePayload( + installationId: string, + accountId: string, + serverId: string, + expectedRevision: number, + input: HostedMcpCreateInput, +) { + return { + ...hostedMcpCreatePayload(installationId, accountId, input), + kind: "update" as const, + serverId, + expectedRevision, + }; +} + +export function hostedMcpLifecyclePayload( + installationId: string, + accountId: string, + serverId: string, + expectedRevision: number, + input: HostedMcpLifecycleInput, +) { + return { ...input, installationId, accountId, serverId, expectedRevision }; +} + +export class HostedMcpRegistryStore { + private readonly database: HostedMcpRegistryDatabase; + + constructor(storage: DurableObjectStorage) { + this.database = new HostedMcpRegistryDatabase(storage); + } + + create( + subject: string, + installationId: string, + accountId: string, + input: HostedMcpCreateInput, + payloadDigest: string, + ): HostedMcpOperationState { + return this.database.transact(() => { + const payload = hostedMcpCreatePayload(installationId, accountId, input); + const prior = this.database.replay( + subject, + installationId, + accountId, + input.intent.requestId, + payloadDigest, + payload, + ); + if (prior) return prior; + const parent = this.database.assertParent( + subject, + installationId, + accountId, + true, + ); + this.database.assertMutationAllowed( + subject, + installationId, + accountId, + false, + ); + + const suffix = identifier(); + const serverId = `hosted-mcp-${suffix}`; + const operationId = identifier(); + const now = nowIso(); + const server = parseHostedMcpServer({ + schemaVersion: 1, + serverId, + installationId, + ownerSubject: subject, + accountId, + accountSubdomain: parseHostedMcpAccountSubdomain( + parent.resources.runtimeOrigin, + installationId, + ), + name: `flarebot-hosted-mcp-${suffix}`, + runtime: input.manifest.runtime.type, + revision: 1, + configurationRevision: 1, + manifest: input.manifest, + activeManifest: null, + deploymentIntent: input.intent, + desiredState: initialDesiredState( + input.manifest, + serverId, + input.configurationDigest, + ), + observedState: emptyObservedState( + serverId, + input.manifest.runtime.type, + now, + ), + deploymentProvider: null, + resourcePlan: initialResourcePlan(input.manifest, serverId), + endpoint: null, + currentOperationId: operationId, + tombstonedAt: null, + createdAt: now, + updatedAt: now, + }); + const operation = parseHostedMcpOperation({ + schemaVersion: 1, + operationId, + requestId: input.intent.requestId, + serverId, + installationId, + ownerSubject: subject, + accountId, + kind: "create", + status: "pending", + startRevision: 1, + observationSequence: 0, + cutover: false, + phase: "absent", + failure: null, + payloadDigest, + payload, + createdAt: now, + updatedAt: now, + }); + this.database.insertServer(server); + this.database.insertOperation(operation); + return { + server: this.database.projection(server), + operation: projectHostedMcpOperation(operation), + }; + }); + } + + update( + subject: string, + installationId: string, + accountId: string, + serverId: string, + expectedRevision: number, + input: HostedMcpCreateInput, + payloadDigest: string, + ): HostedMcpOperationState { + return this.database.transact(() => { + const payload = hostedMcpUpdatePayload( + installationId, + accountId, + serverId, + expectedRevision, + input, + ); + const prior = this.database.replay( + subject, + installationId, + accountId, + input.intent.requestId, + payloadDigest, + payload, + ); + if (prior) return prior; + this.database.assertParent(subject, installationId, accountId, true); + this.database.assertMutationAllowed( + subject, + installationId, + accountId, + false, + ); + + const current = this.database.readServer( + subject, + installationId, + accountId, + serverId, + ); + if ( + current.revision !== expectedRevision || + current.tombstonedAt || + this.database.currentOperation(current)?.status === "pending" || + current.desiredState.lifecycle !== "running" || + !current.activeManifest + ) + return conflict(); + + const nextRevision = current.revision + 1; + const resourcePlan = updateResourcePlan( + current, + input.manifest, + nextRevision, + this.database.resources(serverId), + ); + const disposition = input.manifest.resources.map(({ key, retention }) => { + const previous = current.desiredState.resourceDisposition.find( + (resource) => resource.key === key, + ); + return { + key, + action: + retention === "owner-choice" && previous?.action === "retain" + ? ("retain" as const) + : ("delete" as const), + }; + }); + let desired; + try { + desired = parseHostedMcpDesiredState(input.manifest, { + schemaVersion: 1, + serverId, + revision: nextRevision, + lifecycle: "running", + release: input.manifest.identity, + configuration: { + revision: current.configurationRevision + 1, + digest: input.configurationDigest, + }, + resourceDisposition: disposition, + }); + } catch { + return invalidHostedMcpRecord(); + } + const operationId = identifier(); + const now = nowIso(); + const server = parseHostedMcpServer({ + ...current, + revision: nextRevision, + configurationRevision: current.configurationRevision + 1, + manifest: input.manifest, + deploymentIntent: input.intent, + desiredState: desired, + resourcePlan, + currentOperationId: operationId, + updatedAt: now, + }); + const operation = parseHostedMcpOperation({ + schemaVersion: 1, + operationId, + requestId: input.intent.requestId, + serverId, + installationId, + ownerSubject: subject, + accountId, + kind: "update", + status: "pending", + startRevision: nextRevision, + observationSequence: 0, + cutover: false, + phase: current.observedState.deployment.state, + failure: null, + payloadDigest, + payload, + createdAt: now, + updatedAt: now, + }); + this.database.updateServer(server, current.revision); + this.database.insertOperation(operation); + return { + server: this.database.projection(server), + operation: projectHostedMcpOperation(operation), + }; + }); + } + + beginLifecycle( + subject: string, + installationId: string, + accountId: string, + serverId: string, + expectedRevision: number, + input: HostedMcpLifecycleInput, + payloadDigest: string, + ): HostedMcpOperationState { + return this.database.transact(() => { + const payload = hostedMcpLifecyclePayload( + installationId, + accountId, + serverId, + expectedRevision, + input, + ); + const prior = this.database.replay( + subject, + installationId, + accountId, + input.requestId, + payloadDigest, + payload, + ); + if (prior) return prior; + this.database.assertParent( + subject, + installationId, + accountId, + input.action !== "delete", + ); + this.database.assertMutationAllowed( + subject, + installationId, + accountId, + input.action === "delete", + ); + + const current = this.database.readServer( + subject, + installationId, + accountId, + serverId, + ); + const previous = this.database.currentOperation(current); + if ( + current.revision !== expectedRevision || + current.tombstonedAt || + previous?.status === "pending" || + !lifecycleTransitionAllowed(current, previous, input) + ) + return conflict(); + if ( + input.action !== "delete" && + !same( + input.resourceDisposition, + current.desiredState.resourceDisposition, + ) + ) + return conflict(); + + const nextRevision = current.revision + 1; + const operationId = identifier(); + const now = nowIso(); + const server = parseHostedMcpServer({ + ...current, + revision: nextRevision, + desiredState: desiredForLifecycle(current, nextRevision, input), + resourcePlan: resourcePlanForLifecycle(current, input), + currentOperationId: operationId, + updatedAt: now, + }); + const operation = parseHostedMcpOperation({ + schemaVersion: 1, + operationId, + requestId: input.requestId, + serverId, + installationId, + ownerSubject: subject, + accountId, + kind: input.action, + status: "pending", + startRevision: nextRevision, + observationSequence: 0, + cutover: false, + phase: current.observedState.deployment.state, + failure: null, + payloadDigest, + payload, + createdAt: now, + updatedAt: now, + }); + this.database.updateServer(server, current.revision); + this.database.insertOperation(operation); + return { + server: this.database.projection(server), + operation: projectHostedMcpOperation(operation), + }; + }); + } + + recordObservation( + subject: string, + installationId: string, + accountId: string, + serverId: string, + input: HostedMcpObservationInput, + ): HostedMcpOperationState { + return this.database.transact(() => { + this.database.assertParent(subject, installationId, accountId, false); + const current = this.database.readServer( + subject, + installationId, + accountId, + serverId, + ); + const operation = this.database.readOperation(input.operationId); + const prepared = prepareObservation(current, operation, input); + const inventory = this.database.mergeResources( + current, + operation, + prepared.resources, + ); + const next = completeObservation(prepared, inventory, nowIso()); + this.database.updatePendingOperation( + next.operation, + operation.observationSequence, + ); + this.database.updateServer(next.server, current.revision); + return { + server: this.database.projection(next.server), + operation: projectHostedMcpOperation(next.operation), + }; + }); + } + + recordHealth( + subject: string, + installationId: string, + accountId: string, + serverId: string, + input: HostedMcpHealthInput, + ): HostedMcpServerProjection { + return this.database.transact(() => { + this.database.assertParent(subject, installationId, accountId, false); + const current = this.database.readServer( + subject, + installationId, + accountId, + serverId, + ); + const server = transitionHealth(current, input, nowIso()); + this.database.updateServer(server, current.revision); + return this.database.projection(server); + }); + } + + reconciliationContext( + subject: string, + installationId: string, + accountId: string, + serverId: string, + operationId: string, + ): HostedMcpReconciliationContext { + return this.database.transact(() => { + this.database.assertParent(subject, installationId, accountId, false); + const server = this.database.readServer( + subject, + installationId, + accountId, + serverId, + ); + const operation = this.database.readOperation(operationId); + if ( + operation.serverId !== server.serverId || + operation.installationId !== installationId || + operation.accountId !== accountId || + operation.ownerSubject !== subject + ) + return conflict(); + return { server, operation }; + }); + } + + reconciliationServer( + subject: string, + installationId: string, + accountId: string, + serverId: string, + ) { + return this.database.transact(() => { + this.database.assertParent(subject, installationId, accountId, false); + return this.database.readServer( + subject, + installationId, + accountId, + serverId, + ); + }); + } + + get( + subject: string, + installationId: string, + accountId: string, + serverId: string, + ): HostedMcpServerProjection | null { + return this.database.transact(() => { + this.database.assertParent(subject, installationId, accountId, false); + return this.database.serverExists(serverId) + ? this.database.projection( + this.database.readServer( + subject, + installationId, + accountId, + serverId, + ), + ) + : null; + }); + } + + list( + subject: string, + installationId: string, + accountId: string, + cursor: string | null, + ): HostedMcpInventoryPage { + return this.database.transact(() => { + this.database.assertParent(subject, installationId, accountId, false); + const ids = this.database.serverIds( + installationId, + accountId, + cursor ?? "", + 51, + ); + const servers = ids + .slice(0, 50) + .map((serverId) => + this.database.projection( + this.database.readServer( + subject, + installationId, + accountId, + serverId, + ), + ), + ); + return { + servers, + nextCursor: ids.length > 50 ? (servers.at(-1)?.serverId ?? null) : null, + }; + }); + } + + privilegedInventory( + subject: string, + installationId: string, + accountId: string, + serverId: string | null, + retainedOnly: boolean, + ): HostedMcpPrivilegedInventory { + return this.database.transact(() => { + this.database.assertRemovalFence(subject, installationId, accountId); + if (serverId) + this.database.readServer(subject, installationId, accountId, serverId); + const ids = serverId + ? [serverId] + : this.database.serverIds(installationId, accountId); + return { + installationId, + accountId, + resources: ids.flatMap((id) => + this.database + .resources(id) + .filter( + (resource) => !retainedOnly || resource.cleanup === "retained", + ), + ), + }; + }); + } + + coordinateInstallation( + subject: string, + installationId: string, + accountId: string, + action: "upgrade" | "remove", + ): HostedMcpInstallationCoordination { + return this.database.transact(() => { + this.database.assertParent( + subject, + installationId, + accountId, + action === "upgrade", + ); + if (action === "remove") + this.database.beginRemoval( + subject, + installationId, + accountId, + nowIso(), + ); + if (action === "upgrade") + this.database.assertNoPendingOperations(installationId, accountId); + const servers = this.database + .serverIds(installationId, accountId) + .map((serverId) => + this.database.readServer( + subject, + installationId, + accountId, + serverId, + ), + ); + const activeServerCount = servers.filter( + (server) => !server.tombstonedAt, + ).length; + const pendingOperationCount = this.database.pendingOperationCount( + installationId, + accountId, + ); + const retainedResourceCount = servers.flatMap((server) => + this.database + .resources(server.serverId) + .filter(({ cleanup }) => cleanup === "retained"), + ).length; + return { + installationId, + accountId, + action, + serverCount: servers.length, + activeServerCount, + pendingOperationCount, + retainedResourceCount, + ready: + pendingOperationCount === 0 && + (action === "upgrade" || activeServerCount === 0), + }; + }); + } + + finishInstallationRemoval( + subject: string, + installationId: string, + accountId: string, + ): HostedMcpInstallationCoordination { + return this.database.transact(() => { + this.database.assertRemovalFence(subject, installationId, accountId); + this.database.assertNoPendingOperations(installationId, accountId); + const servers = this.database + .serverIds(installationId, accountId) + .map((serverId) => + this.database.readServer( + subject, + installationId, + accountId, + serverId, + ), + ); + if (servers.some(({ tombstonedAt }) => !tombstonedAt)) return conflict(); + const retainedResourceCount = servers.flatMap((server) => + this.database + .resources(server.serverId) + .filter(({ cleanup }) => cleanup === "retained"), + ).length; + this.database.finishRemoval(subject, installationId, accountId, nowIso()); + return { + installationId, + accountId, + action: "remove", + serverCount: servers.length, + activeServerCount: 0, + pendingOperationCount: 0, + retainedResourceCount, + ready: true, + }; + }); + } +} diff --git a/control-plane/hosted-mcp-registry-transitions.ts b/control-plane/hosted-mcp-registry-transitions.ts new file mode 100644 index 0000000..2aff36a --- /dev/null +++ b/control-plane/hosted-mcp-registry-transitions.ts @@ -0,0 +1,683 @@ +import { + parseHostedMcpDesiredState, + parseHostedMcpObservedState, +} from "./hosted-mcp-template.ts"; +import { + deriveHostedMcpEndpoint, + invalidHostedMcpRecord, + parseHostedMcpOperation, + parseHostedMcpServer, +} from "./hosted-mcp-registry-codec.ts"; +import type { + HostedMcpDeploymentProvider, + HostedMcpHealthInput, + HostedMcpLifecycleInput, + HostedMcpObservationInput, + HostedMcpOperation, + HostedMcpOperationKind, + HostedMcpOwnedResource, + HostedMcpResourcePlan, + HostedMcpServer, +} from "./hosted-mcp-registry-types.ts"; +import { HostedMcpRegistryError } from "./hosted-mcp-registry-types.ts"; +import type { + HostedMcpDesiredState, + HostedMcpObservedState, + HostedMcpTemplateManifest, +} from "../shared/hosted-mcp.ts"; +import { assertHostedMcpRollbackCompatible } from "./hosted-mcp-template.ts"; + +const json = (value: unknown) => JSON.stringify(value); +const same = (left: unknown, right: unknown) => json(left) === json(right); + +function conflict(): never { + throw new HostedMcpRegistryError("hosted_mcp_conflict"); +} + +export function initialDesiredState( + manifest: HostedMcpTemplateManifest, + serverId: string, + configurationDigest: string, +): HostedMcpDesiredState { + return parseHostedMcpDesiredState(manifest, { + schemaVersion: 1, + serverId, + revision: 1, + lifecycle: "running", + release: manifest.identity, + configuration: { revision: 1, digest: configurationDigest }, + resourceDisposition: manifest.resources.map(({ key }) => ({ + key, + action: "delete", + })), + }); +} + +const resourceProviderName = (serverId: string, suffix: string | number) => + `flarebot-${serverId.slice("hosted-mcp-".length)}-${suffix}`; + +export function initialResourcePlan( + manifest: HostedMcpTemplateManifest, + serverId: string, +): HostedMcpResourcePlan[] { + return manifest.resources.map((resource, index) => ({ + inventoryId: crypto.randomUUID().replaceAll("-", ""), + key: resource.key, + kind: resource.kind, + providerName: resourceProviderName(serverId, index + 1), + className: resource.kind === "durable-object" ? resource.className : null, + retention: resource.retention, + action: "preserve", + generation: 1, + })); +} + +function sameManifest( + left: HostedMcpTemplateManifest, + right: HostedMcpTemplateManifest, +) { + return json(left) === json(right); +} + +export function updateResourcePlan( + current: HostedMcpServer, + manifest: HostedMcpTemplateManifest, + nextRevision: number, + inventory: HostedMcpOwnedResource[], +): HostedMcpResourcePlan[] { + const active = current.activeManifest; + if (!active || active.identity.templateId !== manifest.identity.templateId) + return conflict(); + if (same(active.identity, manifest.identity)) { + if (!sameManifest(active, manifest)) return conflict(); + } else { + let compatible = false; + for (const [from, target] of [ + [manifest, active], + [active, manifest], + ] as const) { + try { + assertHostedMcpRollbackCompatible(from, target); + compatible = true; + } catch {} + } + if (!compatible) return conflict(); + } + + const next = new Map(current.resourcePlan.map((entry) => [entry.key, entry])); + for (const entry of current.resourcePlan) + next.set(entry.key, { ...entry, action: "delete" }); + for (const [index, resource] of manifest.resources.entries()) { + const prior = next.get(resource.key); + const priorInventory = prior + ? inventory.find(({ inventoryId }) => inventoryId === prior.inventoryId) + : null; + const requiresNewGeneration = + priorInventory?.state === "deleted" || + priorInventory?.state === "retained"; + const className = + resource.kind === "durable-object" ? resource.className : null; + if ( + prior && + (prior.kind !== resource.kind || + prior.className !== className || + prior.retention !== resource.retention) + ) + return conflict(); + next.set( + resource.key, + prior && !requiresNewGeneration + ? { ...prior, action: "preserve" } + : { + inventoryId: crypto.randomUUID().replaceAll("-", ""), + key: resource.key, + kind: resource.kind, + providerName: resourceProviderName( + current.serverId, + `${String(nextRevision)}-${String(index + 1)}`, + ), + className, + retention: resource.retention, + action: "preserve", + generation: nextRevision, + }, + ); + } + return [...next.values()]; +} + +export function resourcePlanForLifecycle( + server: HostedMcpServer, + input: HostedMcpLifecycleInput, +): HostedMcpResourcePlan[] { + if (input.action !== "delete") return server.resourcePlan; + const disposition = new Map( + input.resourceDisposition.map(({ key, action }) => [key, action]), + ); + return server.resourcePlan.map((entry) => { + const action = disposition.get(entry.key) ?? "delete"; + if (action === "retain" && entry.retention !== "owner-choice") + return conflict(); + return { ...entry, action }; + }); +} + +export function emptyObservedState( + serverId: string, + runtime: "worker" | "container", + now: string, +): HostedMcpObservedState { + return parseHostedMcpObservedState({ + schemaVersion: 1, + serverId, + revision: 1, + runtime, + deployment: { state: "absent", activeRelease: null }, + readiness: { state: "unknown", checkedAt: null }, + resources: [], + credential: { state: "missing", revision: 1, updatedAt: now }, + connection: null, + failure: null, + updatedAt: now, + }); +} + +export function desiredForLifecycle( + server: HostedMcpServer, + nextRevision: number, + input: HostedMcpLifecycleInput, +): HostedMcpDesiredState { + const lifecycle = + input.action === "delete" + ? "deleted" + : input.action === "suspend" + ? "suspended" + : "running"; + try { + return parseHostedMcpDesiredState(server.manifest, { + ...server.desiredState, + revision: nextRevision, + lifecycle, + resourceDisposition: input.resourceDisposition, + }); + } catch { + return invalidHostedMcpRecord(); + } +} + +export function lifecycleTransitionAllowed( + server: HostedMcpServer, + previous: HostedMcpOperation | null, + input: HostedMcpLifecycleInput, +) { + const observed = server.observedState.deployment.state; + return ( + (input.action === "suspend" && + server.desiredState.lifecycle === "running" && + observed === "provisioned") || + (input.action === "resume" && + server.desiredState.lifecycle === "suspended" && + observed === "suspended") || + (input.action === "delete" && observed !== "deleted") || + (input.action === "recover" && + (previous?.status === "failed" || + previous?.status === "recovery-required") && + previous.kind !== "delete" && + (server.activeManifest !== null || + observed === "failed" || + observed === "recovery-required")) + ); +} + +export function progressStateAllowed( + kind: HostedMcpOperationKind, + from: HostedMcpObservedState["deployment"]["state"], + to: HostedMcpObservedState["deployment"]["state"], + cutover: boolean, +) { + if (cutover) return from === "provisioned" && to === "provisioned"; + const allowed = (() => { + switch (kind) { + case "create": + return from === "absent" + ? ["absent", "deploying"] + : from === "deploying" + ? ["deploying"] + : []; + case "update": + return from === "provisioned" + ? ["provisioned", "updating"] + : from === "updating" + ? ["updating", "provisioned"] + : []; + case "suspend": + return from === "provisioned" + ? ["provisioned", "suspending"] + : from === "suspending" + ? ["suspending"] + : []; + case "resume": + return from === "suspended" + ? ["suspended", "updating"] + : from === "updating" + ? ["updating"] + : []; + case "delete": + return from === "deleting" + ? ["deleting"] + : ["absent", "provisioned", "suspended"].includes(from) + ? [from, "deleting"] + : []; + case "recover": + if ( + ["absent", "deploying", "failed", "recovery-required"].includes(from) + ) + return [from, "deploying"]; + if (["provisioned", "updating"].includes(from)) + return ["provisioned", "updating"]; + if (["suspended", "suspending"].includes(from)) + return [from, "updating"]; + return []; + } + })(); + return allowed.includes(to); +} + +export interface PreparedHostedMcpObservation { + current: HostedMcpServer; + operation: HostedMcpOperation; + input: HostedMcpObservationInput; + endpoint: string | null; + deploymentProvider: HostedMcpDeploymentProvider | null; + cutover: boolean; + targetState: "provisioned" | "suspended" | "deleted"; + resources: HostedMcpOwnedResource[]; +} + +function stableDeploymentIdentity(provider: HostedMcpDeploymentProvider) { + const common = { + kind: provider.kind, + accountId: provider.accountId, + scriptName: provider.scriptName, + hostname: provider.hostname, + }; + return provider.kind === "worker" + ? common + : { + ...common, + applicationId: provider.applicationId, + applicationName: provider.applicationName, + namespaceId: provider.namespaceId, + namespaceName: provider.namespaceName, + className: provider.className, + }; +} + +function deploymentProviderForObservation( + current: HostedMcpServer, + operation: HostedMcpOperation, + input: HostedMcpObservationInput, + promotes: boolean, +) { + if (!promotes) { + if (input.deploymentProvider !== null) return conflict(); + return current.deploymentProvider; + } + const provider = input.deploymentProvider ?? current.deploymentProvider; + if (!provider) return conflict(); + deriveHostedMcpEndpoint(current, provider); + const previous = current.deploymentProvider; + if (!previous) return provider; + if ( + !same( + stableDeploymentIdentity(previous), + stableDeploymentIdentity(provider), + ) + ) + return conflict(); + const identityMayAdvance = + !operation.cutover && + (operation.kind === "update" || operation.kind === "recover"); + if (!identityMayAdvance && !same(previous, provider)) return conflict(); + return provider; +} + +export function prepareObservation( + current: HostedMcpServer, + operation: HostedMcpOperation, + input: HostedMcpObservationInput, +): PreparedHostedMcpObservation { + if ( + current.tombstonedAt || + current.currentOperationId !== input.operationId || + operation.serverId !== current.serverId || + operation.status !== "pending" || + operation.startRevision !== input.expectedRevision || + input.sequence !== operation.observationSequence + 1 || + current.revision < input.expectedRevision + ) + return conflict(); + + const targetState = + current.desiredState.lifecycle === "deleted" + ? "deleted" + : current.desiredState.lifecycle === "suspended" + ? "suspended" + : "provisioned"; + if ( + (input.outcome === "succeeded" && input.phase !== targetState) || + ((input.outcome === "failed" || input.outcome === "recovery-required") && + input.phase !== input.outcome) || + (input.outcome === "progress" && + !progressStateAllowed( + operation.kind, + operation.phase, + input.phase, + operation.cutover, + )) + ) + return conflict(); + + const cutover = + input.outcome === "progress" && + input.phase === "provisioned" && + !operation.cutover && + (operation.kind === "update" || operation.kind === "recover") && + current.activeManifest !== null; + if ( + cutover && + current.activeManifest !== null && + !same(current.activeManifest.identity, current.manifest.identity) && + input.deploymentProvider === null + ) + return conflict(); + if ( + input.outcome === "succeeded" && + targetState === "provisioned" && + current.resourcePlan.some(({ action }) => action === "delete") && + !operation.cutover + ) + return conflict(); + const promotes = + cutover || (input.outcome === "succeeded" && targetState === "provisioned"); + const deploymentProvider = deploymentProviderForObservation( + current, + operation, + input, + promotes, + ); + + const endpoint = promotes + ? deploymentProvider + ? deriveHostedMcpEndpoint(current, deploymentProvider) + : null + : targetState === "deleted" && input.outcome === "succeeded" + ? null + : current.endpoint; + if ( + (promotes && endpoint === null) || + (!promotes && input.deploymentProvider !== null) + ) + return conflict(); + + return { + current, + operation, + input, + endpoint, + deploymentProvider, + cutover, + targetState, + resources: input.resources, + }; +} + +export function completeObservation( + prepared: PreparedHostedMcpObservation, + inventory: HostedMcpOwnedResource[], + now: string, +) { + const { + current, + operation, + input, + endpoint, + deploymentProvider, + cutover, + targetState, + } = prepared; + const authorityTime = monotonicTimestamp( + now, + current.updatedAt, + current.observedState.updatedAt, + operation.updatedAt, + ); + const currentInventory = inventory.filter((resource) => + current.resourcePlan.some( + ({ inventoryId }) => inventoryId === resource.inventoryId, + ), + ); + const liveResources = currentInventory + .filter(({ state }) => state !== "retained" && state !== "deleted") + .map(({ key, kind, state, identityDigest }) => ({ + key, + kind, + state, + identityDigest, + })); + const activeResourceKeys = new Set( + current.manifest.resources.map(({ key }) => key), + ); + const provisioning = + cutover || (input.outcome === "succeeded" && targetState === "provisioned"); + if ( + provisioning && + (currentInventory.some((resource) => { + const plan = current.resourcePlan.find( + ({ inventoryId }) => inventoryId === resource.inventoryId, + ); + return ( + !plan || + (plan.action === "preserve" && resource.state !== "ready") || + (!cutover && + plan.action === "delete" && + (resource.state !== "deleted" || resource.cleanup !== "complete")) || + !activeResourceKeys.has(resource.key) !== (plan.action === "delete") + ); + }) || + current.resourcePlan.some( + ({ inventoryId }) => + !currentInventory.some( + (resource) => resource.inventoryId === inventoryId, + ), + )) + ) + return conflict(); + if ( + input.outcome === "succeeded" && + targetState === "deleted" && + currentInventory.some( + ({ cleanup }) => cleanup !== "retained" && cleanup !== "complete", + ) + ) + return conflict(); + + const nextOperation = parseHostedMcpOperation({ + ...operation, + status: input.outcome === "progress" ? "pending" : input.outcome, + observationSequence: input.sequence, + cutover: operation.cutover || cutover, + phase: input.phase, + failure: input.failure, + updatedAt: authorityTime, + }); + const activeManifest = provisioning + ? current.manifest + : targetState === "deleted" && input.outcome === "succeeded" + ? null + : current.activeManifest; + const invalidatesHealth = + cutover || + (!operation.cutover && + input.outcome === "succeeded" && + targetState === "provisioned" && + current.activeManifest !== null && + (operation.kind === "update" || operation.kind === "recover")); + let observed = current.observedState; + if (input.outcome === "succeeded" || cutover) { + const deleting = targetState === "deleted"; + const running = targetState === "provisioned" || cutover; + try { + observed = parseHostedMcpObservedState({ + ...current.observedState, + revision: current.revision + 1, + deployment: { + state: cutover ? "provisioned" : targetState, + activeRelease: deleting ? null : current.desiredState.release, + }, + readiness: + running && !invalidatesHealth + ? current.observedState.readiness + : { state: "unknown", checkedAt: null }, + resources: deleting ? [] : liveResources, + credential: current.observedState.credential, + connection: + running && invalidatesHealth + ? invalidateConnectionHealth(current.observedState.connection) + : running + ? current.observedState.connection + : null, + failure: null, + updatedAt: authorityTime, + }); + } catch { + return conflict(); + } + } else if (operation.cutover) { + try { + observed = parseHostedMcpObservedState({ + ...current.observedState, + revision: current.revision + 1, + resources: liveResources, + updatedAt: authorityTime, + }); + } catch { + return conflict(); + } + } else if (!current.activeManifest) { + try { + observed = parseHostedMcpObservedState({ + ...current.observedState, + revision: current.revision + 1, + deployment: { state: input.phase, activeRelease: null }, + readiness: { state: "unknown", checkedAt: null }, + resources: [], + connection: null, + failure: input.failure, + updatedAt: authorityTime, + }); + } catch { + return conflict(); + } + } + const server = parseHostedMcpServer({ + ...current, + revision: current.revision + 1, + activeManifest, + observedState: observed, + deploymentProvider: provisioning + ? deploymentProvider + : targetState === "deleted" && input.outcome === "succeeded" + ? null + : current.deploymentProvider, + endpoint, + currentOperationId: + input.outcome === "succeeded" && targetState === "deleted" + ? null + : current.currentOperationId, + tombstonedAt: + input.outcome === "succeeded" && targetState === "deleted" + ? authorityTime + : null, + updatedAt: authorityTime, + }); + return { server, operation: nextOperation }; +} + +function invalidateConnectionHealth( + connection: HostedMcpObservedState["connection"], +): HostedMcpObservedState["connection"] { + return connection + ? { + ...connection, + state: connection.enabled ? "connecting" : "disabled", + catalogDigest: null, + } + : null; +} + +function monotonicTimestamp(now: string, ...prior: string[]) { + const previous = Math.max(...prior.map((value) => Date.parse(value))); + const candidate = Date.parse(now); + return new Date(Math.max(candidate, previous + 1)).toISOString(); +} + +export function transitionHealth( + current: HostedMcpServer, + input: HostedMcpHealthInput, + now: string, +) { + const previous = current.observedState; + const authorityTime = monotonicTimestamp( + now, + current.updatedAt, + previous.updatedAt, + previous.credential.updatedAt, + ); + const sameCredentialRevision = + input.credential.revision === previous.credential.revision; + const previousConnection = previous.connection; + const connection = input.connection; + if ( + current.tombstonedAt || + current.revision !== input.expectedRevision || + current.desiredState.lifecycle !== "running" || + input.credential.revision < previous.credential.revision || + (sameCredentialRevision && + input.credential.state !== previous.credential.state) || + (previousConnection !== null && connection === null) || + (previousConnection !== null && + connection !== null && + (connection.id !== previousConnection.id || + connection.revision < previousConnection.revision || + (connection.revision === previousConnection.revision && + !same(connection, previousConnection)))) + ) + return conflict(); + const readiness = + input.readiness.state === "unknown" + ? { state: "unknown" as const, checkedAt: null } + : { state: input.readiness.state, checkedAt: authorityTime }; + const credential = sameCredentialRevision + ? previous.credential + : { + state: input.credential.state, + revision: input.credential.revision, + updatedAt: authorityTime, + }; + const observed = parseHostedMcpObservedState({ + ...previous, + revision: current.revision + 1, + readiness, + credential, + connection, + failure: input.failure, + updatedAt: authorityTime, + }); + return parseHostedMcpServer({ + ...current, + revision: current.revision + 1, + observedState: observed, + updatedAt: authorityTime, + }); +} diff --git a/control-plane/hosted-mcp-registry-types.ts b/control-plane/hosted-mcp-registry-types.ts new file mode 100644 index 0000000..ac97a84 --- /dev/null +++ b/control-plane/hosted-mcp-registry-types.ts @@ -0,0 +1,284 @@ +import type { HostedMcpDeploymentIntent } from "./hosted-mcp-template.ts"; +import type { + HostedMcpDesiredState, + HostedMcpObservedState, + HostedMcpReleaseIdentity, + HostedMcpTemplateManifest, +} from "../shared/hosted-mcp.ts"; + +export const HOSTED_MCP_REGISTRY_SCHEMA_VERSION = 1; + +export type HostedMcpOperationKind = + "create" | "update" | "suspend" | "resume" | "delete" | "recover"; + +export type HostedMcpOperationStatus = + "pending" | "succeeded" | "failed" | "recovery-required"; + +export type HostedMcpDeploymentProvider = + | { + kind: "worker"; + accountId: string; + scriptName: string; + versionId: string; + hostname: string; + } + | { + kind: "container"; + accountId: string; + scriptName: string; + versionId: string; + hostname: string; + applicationId: string; + applicationName: string; + imageDigest: string; + namespaceId: string; + namespaceName: string; + className: string; + }; + +export type HostedMcpProviderIdentity = + | { kind: "r2-bucket"; accountId: string; bucketName: string } + | { + kind: "durable-object"; + accountId: string; + namespaceId: string; + namespaceName: string; + className: string; + }; + +export interface HostedMcpResourcePlan { + inventoryId: string; + key: string; + kind: "r2-bucket" | "durable-object"; + providerName: string; + className: string | null; + retention: "delete" | "owner-choice"; + action: "preserve" | "delete" | "retain"; + generation: number; +} + +export interface HostedMcpOwnedResource { + schemaVersion: 1; + inventoryId: string; + serverId: string; + key: string; + kind: "r2-bucket" | "durable-object"; + state: "pending" | "ready" | "retained" | "deleting" | "deleted" | "unknown"; + identityDigest: string | null; + cleanup: "not-requested" | "pending" | "retained" | "complete" | "failed"; + provider: HostedMcpProviderIdentity | null; + updatedAt: string; +} + +export interface HostedMcpOperation { + schemaVersion: 1; + operationId: string; + requestId: string; + serverId: string; + installationId: string; + ownerSubject: string; + accountId: string; + kind: HostedMcpOperationKind; + status: HostedMcpOperationStatus; + startRevision: number; + observationSequence: number; + cutover: boolean; + phase: HostedMcpObservedState["deployment"]["state"]; + failure: HostedMcpObservedState["failure"]; + payloadDigest: string; + payload: unknown; + createdAt: string; + updatedAt: string; +} + +export interface HostedMcpServer { + schemaVersion: 1; + serverId: string; + installationId: string; + ownerSubject: string; + accountId: string; + accountSubdomain: string; + name: string; + runtime: "worker" | "container"; + revision: number; + configurationRevision: number; + manifest: HostedMcpTemplateManifest; + activeManifest: HostedMcpTemplateManifest | null; + deploymentIntent: HostedMcpDeploymentIntent; + desiredState: HostedMcpDesiredState; + observedState: HostedMcpObservedState; + deploymentProvider: HostedMcpDeploymentProvider | null; + resourcePlan: HostedMcpResourcePlan[]; + endpoint: string | null; + currentOperationId: string | null; + tombstonedAt: string | null; + createdAt: string; + updatedAt: string; +} + +export interface HostedMcpServerProjection { + schemaVersion: 1; + serverId: string; + installationId: string; + accountId: string; + name: string; + runtime: "worker" | "container"; + revision: number; + configurationRevision: number; + desired: HostedMcpDesiredProjection; + observed: HostedMcpObservedProjection; + activeRelease: HostedMcpReleaseIdentity | null; + endpoint: string | null; + operation: HostedMcpOperationProjection | null; + tombstonedAt: string | null; + createdAt: string; + updatedAt: string; +} + +export interface HostedMcpDesiredProjection { + schemaVersion: 1; + serverId: string; + revision: number; + lifecycle: HostedMcpDesiredState["lifecycle"]; + // Immutable artifact/source hashes are reviewed FLA-58 public release + // identity, not a configuration- or credential-derived fingerprint. + release: HostedMcpReleaseIdentity; + configuration: { revision: number }; + resourceDisposition: HostedMcpDesiredState["resourceDisposition"]; +} + +export interface HostedMcpObservedProjection { + schemaVersion: 1; + serverId: string; + revision: number; + runtime: HostedMcpObservedState["runtime"]; + deployment: HostedMcpObservedState["deployment"]; + readiness: HostedMcpObservedState["readiness"]; + resources: Array< + Pick + >; + credential: HostedMcpObservedState["credential"]; + connection: null | Pick< + NonNullable, + "id" | "revision" | "enabled" | "state" + >; + failure: HostedMcpObservedState["failure"]; + updatedAt: string; +} + +export type HostedMcpOperationProjection = Pick< + HostedMcpOperation, + "operationId" | "requestId" | "kind" | "status" | "startRevision" | "failure" +>; + +export interface HostedMcpInventoryPage { + servers: HostedMcpServerProjection[]; + nextCursor: string | null; +} + +export interface HostedMcpCreateInput { + manifest: HostedMcpTemplateManifest; + intent: HostedMcpDeploymentIntent; + configurationDigest: string; +} + +export interface HostedMcpLifecycleInput { + schemaVersion: 1; + requestId: string; + action: "suspend" | "resume" | "delete" | "recover"; + resourceDisposition: Array<{ key: string; action: "delete" | "retain" }>; +} + +export interface HostedMcpOperationState { + server: HostedMcpServerProjection; + operation: HostedMcpOperationProjection; +} + +export interface HostedMcpObservationInput { + operationId: string; + expectedRevision: number; + sequence: number; + outcome: "progress" | "succeeded" | "failed" | "recovery-required"; + phase: HostedMcpObservedState["deployment"]["state"]; + failure: HostedMcpObservedState["failure"]; + deploymentProvider: HostedMcpDeploymentProvider | null; + resources: HostedMcpOwnedResource[]; +} + +export interface HostedMcpResourceObservationClaim { + key: string; + state: HostedMcpOwnedResource["state"]; + cleanup: HostedMcpOwnedResource["cleanup"]; + provider: HostedMcpProviderIdentity | null; +} + +export interface HostedMcpObservationClaim { + operationId: string; + expectedRevision: number; + sequence: number; + outcome: HostedMcpObservationInput["outcome"]; + phase: HostedMcpObservationInput["phase"]; + failure: HostedMcpObservedState["failure"]; + deploymentProvider: HostedMcpDeploymentProvider | null; + resources: HostedMcpResourceObservationClaim[]; +} + +export interface HostedMcpHealthInput { + expectedRevision: number; + readiness: HostedMcpObservedState["readiness"]; + credential: HostedMcpObservedState["credential"]; + connection: HostedMcpObservedState["connection"]; + failure: HostedMcpObservedState["failure"]; + updatedAt: string; +} + +export interface HostedMcpReconciliationContext { + server: HostedMcpServer; + operation: HostedMcpOperation; +} + +export interface HostedMcpInstallationCoordination { + installationId: string; + accountId: string; + action: "upgrade" | "remove"; + serverCount: number; + activeServerCount: number; + pendingOperationCount: number; + retainedResourceCount: number; + ready: boolean; +} + +export interface HostedMcpPrivilegedInventory { + installationId: string; + accountId: string; + resources: HostedMcpOwnedResource[]; +} + +export type HostedMcpRegistryErrorCode = + | "hosted_mcp_invalid" + | "hosted_mcp_not_found" + | "hosted_mcp_conflict" + | "hosted_mcp_installation_not_ready" + | "hosted_mcp_future_schema"; + +export class HostedMcpRegistryError extends Error { + constructor(readonly code: HostedMcpRegistryErrorCode) { + super(code); + this.name = "HostedMcpRegistryError"; + } +} + +export type HostedMcpRegistryResult = + { ok: true; value: T } | { ok: false; error: HostedMcpRegistryErrorCode }; + +export function hostedMcpRegistryResult( + operation: () => T, +): HostedMcpRegistryResult { + try { + return { ok: true, value: operation() }; + } catch (error) { + if (error instanceof HostedMcpRegistryError) + return { ok: false, error: error.code }; + throw error; + } +} diff --git a/control-plane/installation-registry.ts b/control-plane/installation-registry.ts index e0262ba..3728900 100644 --- a/control-plane/installation-registry.ts +++ b/control-plane/installation-registry.ts @@ -2,6 +2,18 @@ import { Agent } from "agents"; import * as Effect from "effect/Effect"; import { z } from "zod"; import type { Env } from "./config.ts"; +import { + HostedMcpRegistryAuthority, + HostedMcpRegistryCoordinator, +} from "./hosted-mcp-registry-authority.ts"; +import { + HostedMcpRegistryError, + type HostedMcpInstallationCoordination, + type HostedMcpInventoryPage, + type HostedMcpOperationState, + type HostedMcpRegistryResult, + type HostedMcpServerProjection, +} from "./hosted-mcp-registry-types.ts"; import { beginDomain, changeDomain, @@ -87,6 +99,32 @@ export class InstallationRegistry extends Agent { ) throw new InstallationNotFound(); } + #hostedAuthority() { + return new HostedMcpRegistryAuthority(this.ctx.storage); + } + #hostedCoordinator() { + return new HostedMcpRegistryCoordinator(this.ctx.storage); + } + async #executeHosted( + subject: string, + operation: (authority: HostedMcpRegistryAuthority) => A | Promise, + publish = false, + ): Promise> { + try { + try { + this.#owner(subject); + } catch { + throw new HostedMcpRegistryError("hosted_mcp_not_found"); + } + const value = await operation(this.#hostedAuthority()); + if (publish) this.setState({ revision: this.state.revision + 1 }); + return { ok: true, value }; + } catch (error) { + if (error instanceof HostedMcpRegistryError) + return { ok: false, error: error.code }; + throw error; + } + } #record(value: unknown, subject: string, id: string) { const record = parseInstallation(value); if (record.ownerSubject !== subject || record.installationId !== id) @@ -393,6 +431,24 @@ export class InstallationRegistry extends Agent { .max(Date.now() + 30 * 60_000), deadline, ); + const replayKey = `operation-request:${requestId}`; + if (upgrade && this.ctx.storage.kv.get(replayKey) === undefined) { + try { + // No await separates this synchronous hosted-server fence from + // opening the parent storage transaction below. A new parent + // upgrade therefore cannot race an unfinished hosted mutation. + this.#hostedCoordinator().coordinateInstallation( + subject, + id, + accountId, + "upgrade", + ); + } catch (error) { + if (error instanceof HostedMcpRegistryError) + return yield* new InstallationConflict(); + throw error; + } + } return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const current = yield* this.#readRecord(tx, subject, id); @@ -406,7 +462,6 @@ export class InstallationRegistry extends Agent { ) ) return yield* new InstallationConflict(); - const replayKey = `operation-request:${requestId}`; const replayValue = yield* Effect.promise(() => tx.get(replayKey)); const replay = replayValue === undefined @@ -625,4 +680,131 @@ export class InstallationRegistry extends Agent { true, ); } + + createHostedMcpServer( + subject: string, + installationId: unknown, + accountId: unknown, + manifest: unknown, + deploymentIntent: unknown, + ): Promise> { + return this.#executeHosted( + subject, + (authority) => + authority.create( + subject, + installationId, + accountId, + manifest, + deploymentIntent, + ), + true, + ); + } + + updateHostedMcpServer( + subject: string, + installationId: unknown, + accountId: unknown, + serverId: unknown, + expectedRevision: unknown, + manifest: unknown, + deploymentIntent: unknown, + ): Promise> { + return this.#executeHosted( + subject, + (authority) => + authority.update( + subject, + installationId, + accountId, + serverId, + expectedRevision, + manifest, + deploymentIntent, + ), + true, + ); + } + + beginHostedMcpLifecycle( + subject: string, + installationId: unknown, + accountId: unknown, + serverId: unknown, + expectedRevision: unknown, + input: unknown, + ): Promise> { + return this.#executeHosted( + subject, + (authority) => + authority.beginLifecycle( + subject, + installationId, + accountId, + serverId, + expectedRevision, + input, + ), + true, + ); + } + + getHostedMcpServer( + subject: string, + installationId: unknown, + accountId: unknown, + serverId: unknown, + ): Promise> { + return this.#executeHosted(subject, (authority) => + authority.get(subject, installationId, accountId, serverId), + ); + } + + listHostedMcpServers( + subject: string, + installationId: unknown, + accountId: unknown, + cursor: unknown = null, + ): Promise> { + return this.#executeHosted(subject, (authority) => + authority.list(subject, installationId, accountId, cursor), + ); + } + + coordinateHostedMcpInstallation( + subject: string, + installationId: unknown, + accountId: unknown, + action: "upgrade" | "remove", + ): Promise> { + return this.#executeHosted( + subject, + () => + this.#hostedCoordinator().coordinateInstallation( + subject, + installationId, + accountId, + action, + ), + true, + ); + } + + finishHostedMcpInstallationRemoval( + subject: string, + installationId: unknown, + accountId: unknown, + ): Promise> { + return this.#executeHosted( + subject, + () => + this.#hostedCoordinator().finishRemoval( + subject, + installationId, + accountId, + ), + true, + ); + } } diff --git a/package.json b/package.json index d3731bb..8d00156 100644 --- a/package.json +++ b/package.json @@ -71,6 +71,7 @@ "test:hosted-mcp-artifacts": "node --test tests/hosted-mcp-artifact.test.mjs tests/hosted-mcp-docker-image.test.mjs tests/hosted-mcp-native-cleanup.test.mjs tests/hosted-mcp-output.test.mjs tests/hosted-mcp-publication-store.test.mjs", "test:hosted-mcp-artifacts-production": "node --test tests/hosted-mcp-production-output.test.mjs", "test:hosted-mcp-artifacts-native": "node --test --test-concurrency=1 tests/hosted-mcp-artifact-native.test.mjs", + "test:hosted-mcp-persistence": "node --test --test-concurrency=1 --test-reporter=tap tests/hosted-mcp-persistence.test.mjs tests/hosted-mcp-persistence-adversarial.test.mjs", "test:skill-package": "node --test tests/skill-package.test.mjs", "test:bundled-skills": "node --test tests/bundled-skills.test.mjs", "test:installed-skills": "node --test tests/installed-skills.test.mjs", diff --git a/tests/fixtures/hosted-mcp-registry-worker.ts b/tests/fixtures/hosted-mcp-registry-worker.ts new file mode 100644 index 0000000..a959d8a --- /dev/null +++ b/tests/fixtures/hosted-mcp-registry-worker.ts @@ -0,0 +1,266 @@ +import { InstallationRegistry as ProductionRegistry } from "../../control-plane/installation-registry.ts"; +import { + HostedMcpRegistryCoordinator, + HostedMcpRegistryReconciler, +} from "../../control-plane/hosted-mcp-registry-authority.ts"; +import { hostedMcpRegistryResult } from "../../control-plane/hosted-mcp-registry-types.ts"; +import { registryName } from "../../control-plane/installation-metadata.ts"; +import { parseInstallation } from "../../control-plane/installation-metadata.ts"; +import type { Env } from "../../control-plane/config.ts"; + +export class InstallationRegistry extends ProductionRegistry { + fixtureProductionReconcilerSurface() { + const registry = this as unknown as Record; + return { + observation: typeof registry.recordHostedMcpObservation === "function", + health: typeof registry.recordHostedMcpHealth === "function", + inventory: typeof registry.hostedMcpResourceInventory === "function", + }; + } + + async fixtureRecordHostedMcpObservation( + subject: string, + installationId: unknown, + accountId: unknown, + serverId: unknown, + input: unknown, + ) { + try { + return { + ok: true as const, + value: await new HostedMcpRegistryReconciler( + this.ctx.storage, + ).recordObservation( + subject, + installationId, + accountId, + serverId, + input, + ), + }; + } catch (error) { + return hostedMcpRegistryResult(() => { + throw error; + }); + } + } + + fixtureRecordHostedMcpHealth( + subject: string, + installationId: unknown, + accountId: unknown, + serverId: unknown, + input: unknown, + ) { + return hostedMcpRegistryResult(() => + new HostedMcpRegistryReconciler(this.ctx.storage).recordHealth( + subject, + installationId, + accountId, + serverId, + input, + ), + ); + } + + fixturePrivilegedHostedMcpInventory( + subject: string, + installationId: unknown, + accountId: unknown, + serverId: unknown, + retainedOnly: boolean, + ) { + return hostedMcpRegistryResult(() => + new HostedMcpRegistryCoordinator(this.ctx.storage).inventory( + subject, + installationId, + accountId, + serverId, + retainedOnly, + ), + ); + } + + fixtureRemoveParentInstallation(installationId: string) { + this.ctx.storage.kv.delete(`installation:${installationId}`); + } + + async fixtureSetParentStatus( + installationId: string, + status: "ready" | "failed", + ) { + const key = `installation:${installationId}`; + const current = parseInstallation(this.ctx.storage.kv.get(key)); + await this.ctx.storage.put( + key, + parseInstallation({ + ...current, + status, + progress: status === "ready" ? "complete" : "deploying", + errorCode: status === "ready" ? null : "deployment_failed", + revision: current.revision + 1, + updatedAt: Math.max(Date.now(), current.updatedAt), + }), + ); + } + + async fixtureSeedInstallation( + subject: string, + installationId: string, + accountId: string, + ) { + const now = Date.now(); + const release = { + version: "1.0.0", + sourceRevision: "1".repeat(40), + artifactDigest: "2".repeat(64), + }; + await this.ctx.storage.put(`installation:${installationId}`, { + schemaVersion: 1, + installationId, + ownerSubject: subject, + accountId, + createdAt: now, + updatedAt: now, + revision: 1, + resources: { + workerName: `flarebot-${installationId}`, + sandboxApplicationName: `flarebot-shell-${installationId}`, + personalAgentNamespaceId: "3".repeat(32), + sandboxNamespaceId: "4".repeat(32), + sandboxApplicationId: "5".repeat(32), + runtimeOrigin: `https://flarebot-${installationId}.fixture.workers.dev`, + }, + desiredRelease: release, + installedRelease: { ...release, installedAt: now }, + status: "ready", + progress: "complete", + errorCode: null, + operationId: "6".repeat(32), + }); + } + + fixtureHostedRows() { + const count = (table: string) => { + try { + return this.ctx.storage.sql + .exec<{ count: number }>(`SELECT COUNT(*) AS count FROM ${table}`) + .one().count; + } catch { + return 0; + } + }; + let schemaVersion: number | null = null; + try { + schemaVersion = this.ctx.storage.sql + .exec<{ version: number }>( + "SELECT version FROM flarebot_hosted_mcp_schema WHERE singleton = 1", + ) + .one().version; + } catch {} + return { + schemaVersion, + servers: count("hosted_mcp_servers"), + operations: count("hosted_mcp_operations"), + resources: count("hosted_mcp_resources"), + }; + } + + fixtureSetHostedSchemaVersion(version: number) { + this.ctx.storage.transactionSync(() => { + this.ctx.storage.sql + .exec( + ` + CREATE TABLE IF NOT EXISTS flarebot_hosted_mcp_schema ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + version INTEGER NOT NULL + ) + `, + ) + .toArray(); + this.ctx.storage.sql + .exec( + `INSERT INTO flarebot_hosted_mcp_schema (singleton, version) + VALUES (1, ?) ON CONFLICT(singleton) DO UPDATE SET version = excluded.version`, + version, + ) + .toArray(); + }); + } +} + +type FixtureStub = DurableObjectStub; + +export default { + async fetch(request, env) { + if (request.method !== "POST") return new Response(null, { status: 405 }); + const body = (await request.json()) as Record; + const owner = body.registryOwner as string; + const subject = (body.subject as string | undefined) ?? owner; + const stub = env.INSTALLATIONS.get( + env.INSTALLATIONS.idFromName(registryName(owner)), + ) as FixtureStub; + const args = (body.args as unknown[] | undefined) ?? []; + const invoke = (method: unknown) => + (method as (...values: unknown[]) => Promise)(subject, ...args); + switch (new URL(request.url).pathname) { + case "/seed": + await stub.fixtureSeedInstallation( + owner, + body.installationId as string, + body.accountId as string, + ); + return Response.json({ ok: true }); + case "/rows": + return Response.json(await stub.fixtureHostedRows()); + case "/production-surface": + return Response.json(await stub.fixtureProductionReconcilerSurface()); + case "/future": + await stub.fixtureSetHostedSchemaVersion(body.version as number); + return Response.json({ ok: true }); + case "/parent-status": + await stub.fixtureSetParentStatus( + body.installationId as string, + body.status as "ready" | "failed", + ); + return Response.json({ ok: true }); + case "/start": + return Response.json(await invoke(stub.start)); + case "/create": + return Response.json(await invoke(stub.createHostedMcpServer)); + case "/update": + return Response.json(await invoke(stub.updateHostedMcpServer)); + case "/lifecycle": + return Response.json(await invoke(stub.beginHostedMcpLifecycle)); + case "/observe": + return Response.json( + await invoke(stub.fixtureRecordHostedMcpObservation), + ); + case "/health": + return Response.json(await invoke(stub.fixtureRecordHostedMcpHealth)); + case "/get": + return Response.json(await invoke(stub.getHostedMcpServer)); + case "/list": + return Response.json(await invoke(stub.listHostedMcpServers)); + case "/inventory": + return Response.json( + await invoke(stub.fixturePrivilegedHostedMcpInventory), + ); + case "/coordinate": + return Response.json( + await invoke(stub.coordinateHostedMcpInstallation), + ); + case "/finish-removal": + return Response.json( + await invoke(stub.finishHostedMcpInstallationRemoval), + ); + case "/remove-parent": + await stub.fixtureRemoveParentInstallation( + body.installationId as string, + ); + return Response.json({ ok: true }); + default: + return new Response(null, { status: 404 }); + } + }, +} satisfies ExportedHandler; diff --git a/tests/helpers/hosted-mcp-registry-harness.mjs b/tests/helpers/hosted-mcp-registry-harness.mjs new file mode 100644 index 0000000..bd110c7 --- /dev/null +++ b/tests/helpers/hosted-mcp-registry-harness.mjs @@ -0,0 +1,331 @@ +import assert from "node:assert/strict"; +import { randomBytes } from "node:crypto"; +import { mkdtemp, mkdir, readFile, rm, writeFile } from "node:fs/promises"; +import { createServer } from "node:net"; +import { tmpdir } from "node:os"; +import { join, resolve } from "node:path"; +import { unstable_dev } from "wrangler"; + +export const id = () => randomBytes(16).toString("hex"); +const deadline = (promise, milliseconds, label) => + Promise.race([ + promise, + new Promise((_, reject) => { + const timer = setTimeout( + () => reject(new Error(`${label} exceeded ${milliseconds}ms`)), + milliseconds, + ); + timer.unref(); + }), + ]); + +async function freePort() { + const server = createServer(); + await new Promise((done) => server.listen(0, "127.0.0.1", done)); + const { port } = server.address(); + await new Promise((done) => server.close(done)); + return port; +} + +export const intentFor = (manifest, requestId, greeting = "hello") => ({ + schemaVersion: 1, + requestId, + template: manifest.identity, + configuration: manifest.configuration.fields.some( + ({ key }) => key === "greeting", + ) + ? { greeting } + : {}, +}); + +export const providerName = (server, suffix = 1) => + `flarebot-${server.serverId.slice("hosted-mcp-".length)}-${suffix}`; + +export const deploymentProvider = (server, account, overrides = {}) => { + const common = { + accountId: account, + scriptName: server.name, + versionId: "1".repeat(32), + hostname: `${server.name}.fixture.workers.dev`, + }; + return server.runtime === "worker" + ? { kind: "worker", ...common, ...overrides } + : { + kind: "container", + ...common, + applicationId: "2".repeat(32), + applicationName: `${server.name}-app`, + imageDigest: "8".repeat(64), + namespaceId: "3".repeat(32), + namespaceName: `${server.name}-ns`, + className: "HostedMcpContainer", + ...overrides, + }; +}; + +export const resourceClaim = (server, account, overrides = {}) => + server.runtime === "worker" + ? { + key: "notes", + state: "ready", + cleanup: "not-requested", + provider: { + kind: "durable-object", + accountId: account, + namespaceId: "9".repeat(32), + namespaceName: providerName(server), + className: "NotesStore", + }, + ...overrides, + } + : { + key: "archive", + state: "ready", + cleanup: "not-requested", + provider: { + kind: "r2-bucket", + accountId: account, + bucketName: providerName(server), + }, + ...overrides, + }; + +export const observation = ( + server, + operation, + outcome, + phase, + overrides = {}, +) => ({ + operationId: operation.operationId, + expectedRevision: operation.startRevision, + sequence: 1, + outcome, + phase, + failure: + outcome === "failed" + ? { + stage: "deployment", + code: "deployment_failed", + outcome: "known", + replay: "allowed", + } + : outcome === "recovery-required" + ? { + stage: "mcp", + code: "outcome_uncertain", + outcome: "uncertain", + replay: "forbidden", + } + : null, + deploymentProvider: null, + resources: [], + ...overrides, +}); + +export const health = (server, overrides = {}) => { + const now = new Date().toISOString(); + return { + expectedRevision: server.revision, + readiness: { state: "ready", checkedAt: now }, + credential: { + state: "configured", + revision: server.observed.credential.revision + 1, + updatedAt: now, + }, + connection: { + id: "mcp-00000000-0000-0000-0000-000000000000", + revision: 1, + enabled: true, + state: "ready", + catalogDigest: "a".repeat(64), + }, + failure: null, + updatedAt: now, + ...overrides, + }; +}; + +export function version(manifest, value, fill) { + const next = structuredClone(manifest); + next.identity = { + ...manifest.identity, + version: value, + sourceRevision: fill.repeat(40), + artifactDigest: fill.repeat(64), + }; + next.compatibility.rollbackTargets = [ + { + identity: structuredClone(manifest.identity), + dataSchemaVersion: manifest.data.schemaVersion, + }, + ]; + return next; +} + +export async function registryHarness(t) { + const lifetime = new AbortController(); + let temporary = null; + let worker = null; + let acquisition = null; + + t.after(async () => { + lifetime.abort(new Error("test lifetime ended")); + let acquired = worker; + if (!acquired && acquisition) { + try { + acquired = await deadline( + acquisition, + 5_000, + "late worker acquisition", + ); + } catch {} + } + if (acquired) { + try { + await deadline(acquired.stop(), 10_000, "worker cleanup"); + } catch (error) { + t.diagnostic(`worker cleanup: ${error}`); + } + } + if (temporary) + await deadline( + rm(temporary, { recursive: true, force: true }), + 10_000, + "fixture cleanup", + ); + }); + + temporary = await mkdtemp(join(tmpdir(), "flarebot-hosted-mcp-")); + const assets = join(temporary, "assets"); + await mkdir(assets); + const port = await freePort(); + const origin = `http://127.0.0.1:${port}`; + const configPath = join(temporary, "wrangler.json"); + await writeFile( + configPath, + JSON.stringify({ + name: `flarebot-hosted-mcp-${id()}`, + main: resolve("tests/fixtures/hosted-mcp-registry-worker.ts"), + compatibility_date: "2026-09-04", + compatibility_flags: ["nodejs_compat", "global_fetch_strictly_public"], + assets: { directory: assets, binding: "ASSETS" }, + durable_objects: { + bindings: [ + { name: "INSTALLATIONS", class_name: "InstallationRegistry" }, + ], + }, + exports: { + InstallationRegistry: { type: "durable-object", storage: "sqlite" }, + }, + observability: { enabled: false }, + }), + ); + + const start = async () => { + acquisition = unstable_dev("tests/fixtures/hosted-mcp-registry-worker.ts", { + config: configPath, + local: true, + ip: "127.0.0.1", + port, + inspectorPort: 0, + persist: true, + persistTo: temporary, + logLevel: "error", + experimental: { disableExperimentalWarning: true, watch: false }, + }); + acquisition.then( + (lateWorker) => { + if (lifetime.signal.aborted) void lateWorker.stop(); + }, + () => {}, + ); + worker = await deadline(acquisition, 30_000, "worker acquisition"); + return worker; + }; + await start(); + + const call = async (path, registryOwner, args = [], extra = {}) => { + const response = await fetch(`${origin}/${path}`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ registryOwner, args, ...extra }), + signal: AbortSignal.any([lifetime.signal, AbortSignal.timeout(8_000)]), + }); + assert.equal(response.status, 200, `${path} HTTP status`); + return response.json(); + }; + return { + call, + restart: async () => { + await deadline(worker.stop(), 10_000, "worker restart stop"); + worker = null; + await start(); + }, + seed: (owner, installationId, accountId) => + call("seed", owner, [], { installationId, accountId }), + rows: (owner) => call("rows", owner), + }; +} + +export const fixtureManifest = (runtime = "worker") => + readFile( + new URL(`../fixtures/hosted-mcp/${runtime}.json`, import.meta.url), + ).then(JSON.parse); + +export async function createReady( + registry, + owner, + installation, + account, + manifest, +) { + let result = await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, id()), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + const { server, operation } = result.value; + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + return result.value; +} + +const deniedKeys = new Set([ + "ownerSubject", + "manifest", + "activeManifest", + "deploymentIntent", + "deploymentProvider", + "resourcePlan", + "provider", + "identityDigest", + "catalogDigest", + "configurationDigest", + "payloadDigest", + "payload", +]); + +export function assertSafeProjection(value, path = "result") { + if (Array.isArray(value)) { + value.forEach((entry, index) => + assertSafeProjection(entry, `${path}[${index}]`), + ); + return; + } + if (!value || typeof value !== "object") return; + for (const [key, nested] of Object.entries(value)) { + assert.equal(deniedKeys.has(key), false, `${path}.${key}`); + assertSafeProjection(nested, `${path}.${key}`); + } +} diff --git a/tests/hosted-mcp-persistence-adversarial.test.mjs b/tests/hosted-mcp-persistence-adversarial.test.mjs new file mode 100644 index 0000000..593b730 --- /dev/null +++ b/tests/hosted-mcp-persistence-adversarial.test.mjs @@ -0,0 +1,621 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { + assertSafeProjection, + createReady, + deploymentProvider, + fixtureManifest, + health, + id, + intentFor, + observation, + registryHarness, + resourceClaim, + version, +} from "./helpers/hosted-mcp-registry-harness.mjs"; + +test( + "native observation sequencing and authority time prevent stale regressions", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const owner = "hosted-owner-sequence"; + const installation = "7".repeat(32); + const account = "8".repeat(32); + await registry.seed(owner, installation, account); + + let result = await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, id()), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + let { server, operation } = result.value; + + t.diagnostic("phase: delayed progress cannot regress an operation phase"); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "deploying", { sequence: 1 }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + assert.equal(server.observed.deployment.state, "deploying"); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "absent", { sequence: 2 }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + sequence: 2, + deploymentProvider: deploymentProvider(server, account), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + + t.diagnostic("phase: caller timestamps cannot wedge authority time"); + const farFuture = "9999-12-31T23:59:59.999Z"; + result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + readiness: { state: "ready", checkedAt: farFuture }, + credential: { + state: "configured", + revision: server.observed.credential.revision + 1, + updatedAt: farFuture, + }, + updatedAt: farFuture, + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + assert.notEqual(server.observed.updatedAt, farFuture); + const firstAuthorityTime = server.observed.readiness.checkedAt; + + result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + readiness: { state: "unhealthy", checkedAt: firstAuthorityTime }, + credential: server.observed.credential, + connection: { + ...server.observed.connection, + revision: server.observed.connection.revision + 1, + state: "error", + catalogDigest: null, + }, + updatedAt: firstAuthorityTime, + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + assert.equal(server.observed.readiness.state, "unhealthy"); + assert.ok(server.observed.readiness.checkedAt > firstAuthorityTime); + await registry.restart(); + result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + readiness: { state: "ready", checkedAt: firstAuthorityTime }, + credential: server.observed.credential, + connection: { + ...server.observed.connection, + revision: server.observed.connection.revision + 1, + state: "ready", + catalogDigest: "c".repeat(64), + }, + updatedAt: firstAuthorityTime, + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.ok( + result.value.observed.readiness.checkedAt > + server.observed.readiness.checkedAt, + ); + }, +); + +test( + "native cutover invalidates old health through failed cleanup and restart", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const manifestV2 = version(manifest, "2.0.0", "d"); + manifestV2.resources = []; + manifestV2.data.resources = []; + manifestV2.runtime.bindings = manifestV2.runtime.bindings.filter( + ({ kind }) => kind !== "durable-object", + ); + const owner = "hosted-owner-cutover-health"; + const installation = "d".repeat(32); + const account = "e".repeat(32); + await registry.seed(owner, installation, account); + + let { server } = await createReady( + registry, + owner, + installation, + account, + manifest, + ); + let result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + + result = await registry.call("update", owner, [ + installation, + account, + server.serverId, + server.revision, + manifestV2, + intentFor(manifestV2, id()), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + let { operation } = result.value; + server = result.value.server; + + t.diagnostic("phase: V1 health remains writable before V2 cutover"); + result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + assert.equal(server.observed.readiness.state, "ready"); + const staleV1Health = health(server, { + connection: { + ...server.observed.connection, + revision: server.observed.connection.revision + 1, + state: "ready", + catalogDigest: "d".repeat(64), + }, + }); + + t.diagnostic("phase: V2 cutover invalidates V1 readiness and transport"); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "provisioned", { + sequence: 1, + deploymentProvider: deploymentProvider(server, account, { + versionId: "2".repeat(32), + }), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + assert.deepEqual(server.activeRelease, manifestV2.identity); + assert.deepEqual(server.observed.readiness, { + state: "unknown", + checkedAt: null, + }); + assert.equal(server.observed.connection.state, "connecting"); + assert.deepEqual( + await registry.call("health", owner, [ + installation, + account, + server.serverId, + staleV1Health, + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + + await registry.restart(); + result = await registry.call("get", owner, [ + installation, + account, + server.serverId, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + operation = server.operation; + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "updating", { + sequence: 2, + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "provisioned", { + sequence: 2, + resources: [ + resourceClaim(server, account, { + state: "deleting", + cleanup: "pending", + }), + ], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "failed", "failed", { + sequence: 3, + resources: [ + resourceClaim(server, account, { + state: "deleting", + cleanup: "pending", + }), + ], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + + t.diagnostic("phase: failed cleanup restart cannot revive V1 health"); + await registry.restart(); + result = await registry.call("get", owner, [ + installation, + account, + server.serverId, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + assert.deepEqual(server.activeRelease, manifestV2.identity); + assert.equal(server.observed.resources[0].state, "deleting"); + assert.deepEqual(server.observed.readiness, { + state: "unknown", + checkedAt: null, + }); + assert.equal(server.observed.connection.state, "connecting"); + + result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + connection: { + ...server.observed.connection, + revision: server.observed.connection.revision + 1, + state: "ready", + catalogDigest: "e".repeat(64), + }, + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.value.observed.readiness.state, "ready"); + assert.equal(result.value.observed.connection.state, "ready"); + }, +); + +test( + "native container authority freezes account and provider identity", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest("container"); + const owner = "hosted-owner-container"; + const installation = "9".repeat(32); + const account = "a".repeat(32); + await registry.seed(owner, installation, account); + + let result = await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, id()), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + let { server, operation } = result.value; + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account, { + hostname: `${server.name}.attacker.workers.dev`, + }), + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_invalid" }, + ); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assertSafeProjection(result); + server = result.value.server; + + t.diagnostic("phase: resume cannot substitute stable container identity"); + result = await registry.call("lifecycle", owner, [ + installation, + account, + server.serverId, + server.revision, + { + schemaVersion: 1, + requestId: id(), + action: "suspend", + resourceDisposition: server.desired.resourceDisposition, + }, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + ({ server, operation } = result.value); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "suspended", { + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + result = await registry.call("lifecycle", owner, [ + installation, + account, + server.serverId, + server.revision, + { + schemaVersion: 1, + requestId: id(), + action: "resume", + resourceDisposition: server.desired.resourceDisposition, + }, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + ({ server, operation } = result.value); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account, { + applicationId: "4".repeat(32), + }), + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + + t.diagnostic( + "phase: update may advance version but not application identity", + ); + const manifestV2 = version(manifest, "2.0.0", "b"); + result = await registry.call("update", owner, [ + installation, + account, + server.serverId, + server.revision, + manifestV2, + intentFor(manifestV2, id()), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + ({ server, operation } = result.value); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account, { + versionId: "5".repeat(32), + namespaceId: "6".repeat(32), + }), + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account, { + versionId: "5".repeat(32), + }), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.deepEqual(result.value.server.activeRelease, manifestV2.identity); + }, +); + +test( + "native failed delete rejects incoherent recovery for delete and retain plans", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const owner = "hosted-owner-delete-recovery"; + const installation = "b".repeat(32); + const account = "c".repeat(32); + await registry.seed(owner, installation, account); + + for (const action of ["delete", "retain"]) { + let { server } = await createReady( + registry, + owner, + installation, + account, + manifest, + ); + const disposition = [{ key: "notes", action }]; + let result = await registry.call("lifecycle", owner, [ + installation, + account, + server.serverId, + server.revision, + { + schemaVersion: 1, + requestId: id(), + action: "delete", + resourceDisposition: disposition, + }, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + let operation; + ({ server, operation } = result.value); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "failed", "failed"), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + assert.deepEqual( + await registry.call("lifecycle", owner, [ + installation, + account, + server.serverId, + server.revision, + { + schemaVersion: 1, + requestId: id(), + action: "recover", + resourceDisposition: disposition, + }, + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + const retry = await registry.call("lifecycle", owner, [ + installation, + account, + server.serverId, + server.revision, + { + schemaVersion: 1, + requestId: id(), + action: "delete", + resourceDisposition: disposition, + }, + ]); + assert.equal(retry.ok, true, JSON.stringify(retry)); + } + }, +); + +test( + "native production installation upgrade waits for hosted operations", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const owner = "hosted-owner-parent-upgrade"; + const installation = "d".repeat(32); + const account = "e".repeat(32); + await registry.seed(owner, installation, account); + + let hosted = await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, id()), + ]); + assert.equal(hosted.ok, true, JSON.stringify(hosted)); + let { server, operation } = hosted.value; + const fromRelease = { + version: "1.0.0", + sourceRevision: "1".repeat(40), + artifactDigest: "2".repeat(64), + }; + const toRelease = { + version: "2.0.0", + sourceRevision: "7".repeat(40), + artifactDigest: "8".repeat(64), + }; + const upgrade = { + fromRelease, + toRelease, + deployedOperationId: "3".repeat(32), + configDigest: "4".repeat(64), + versionId: "5".repeat(32), + deploymentId: "6".repeat(32), + sourceCodeHash: "7".repeat(64), + fingerprint: "8".repeat(64), + containerFingerprint: "9".repeat(64), + targetContainerFingerprint: "a".repeat(64), + }; + const start = () => + registry.call("start", owner, [ + installation, + account, + id(), + toRelease, + Date.now() + 60_000, + null, + upgrade, + ]); + + assert.deepEqual(await start(), { + ok: false, + error: "installation_conflict", + }); + hosted = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(hosted.ok, true, JSON.stringify(hosted)); + const started = await start(); + assert.equal(started.ok, true, JSON.stringify(started)); + assert.equal(started.value.installation.status, "updating"); + }, +); diff --git a/tests/hosted-mcp-persistence.test.mjs b/tests/hosted-mcp-persistence.test.mjs new file mode 100644 index 0000000..962e167 --- /dev/null +++ b/tests/hosted-mcp-persistence.test.mjs @@ -0,0 +1,852 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { + assertSafeProjection, + createReady, + deploymentProvider, + fixtureManifest, + health, + id, + intentFor, + observation, + providerName, + registryHarness, + resourceClaim, + version, +} from "./helpers/hosted-mcp-registry-harness.mjs"; + +test( + "native registry atomically scopes replay and rejects future schema", + { timeout: 90_000 }, + async (t) => { + t.diagnostic("phase: acquire isolated persisted registry"); + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const owner = "hosted-owner-replay"; + const installation = "a".repeat(32); + const secondInstallation = "b".repeat(32); + const account = "c".repeat(32); + const secondAccount = "d".repeat(32); + await registry.seed(owner, installation, account); + await registry.seed(owner, secondInstallation, secondAccount); + + t.diagnostic("phase: concurrent exact replay"); + const requestId = id(); + const args = [ + installation, + account, + manifest, + intentFor(manifest, requestId), + ]; + const results = await Promise.all( + Array.from({ length: 8 }, () => registry.call("create", owner, args)), + ); + results.forEach((result) => assert.deepEqual(result, results[0])); + assert.equal(results[0].ok, true); + assert.equal((await registry.rows(owner)).operations, 1); + + t.diagnostic("phase: exact replay precedes current parent readiness"); + await registry.call("parent-status", owner, [], { + installationId: installation, + status: "failed", + }); + assert.equal( + (await registry.call("create", owner, args)).value.operation.operationId, + results[0].value.operation.operationId, + ); + assert.deepEqual( + await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, id()), + ]), + { ok: false, error: "hosted_mcp_installation_not_ready" }, + ); + await registry.call("parent-status", owner, [], { + installationId: installation, + status: "ready", + }); + assert.deepEqual( + await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, requestId, "changed"), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + assert.equal((await registry.rows(owner)).operations, 1); + + t.diagnostic("phase: identical request IDs are scoped per installation"); + const second = await registry.call("create", owner, [ + secondInstallation, + secondAccount, + manifest, + intentFor(manifest, requestId), + ]); + assert.equal(second.ok, true); + assert.equal((await registry.rows(owner)).operations, 2); + assert.equal( + ( + await registry.call("coordinate", owner, [ + installation, + account, + "remove", + ]) + ).ok, + true, + ); + await registry.restart(); + assert.equal( + (await registry.call("create", owner, args)).value.operation.operationId, + results[0].value.operation.operationId, + ); + assert.deepEqual( + await registry.call("get", owner, [ + installation, + secondAccount, + results[0].value.server.serverId, + ]), + { ok: false, error: "hosted_mcp_not_found" }, + ); + + t.diagnostic("phase: future schema refuses before hosted writes"); + const futureOwner = "hosted-owner-future"; + const futureInstallation = "e".repeat(32); + const futureAccount = "f".repeat(32); + await registry.seed(futureOwner, futureInstallation, futureAccount); + await registry.call("future", futureOwner, [], { version: 2 }); + const before = await registry.rows(futureOwner); + assert.deepEqual( + await registry.call("list", futureOwner, [ + futureInstallation, + futureAccount, + null, + ]), + { ok: false, error: "hosted_mcp_future_schema" }, + ); + assert.deepEqual(await registry.rows(futureOwner), before); + }, +); + +test( + "native owner RPCs are recursively safe and provider authority is closed", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const owner = "hosted-owner-projection"; + const otherOwner = "hosted-owner-other"; + const installation = "1".repeat(32); + const account = "2".repeat(32); + await registry.seed(owner, installation, account); + + t.diagnostic( + "phase: trusted reconciler is absent from the owner RPC surface", + ); + assert.deepEqual(await registry.call("production-surface", owner), { + observation: false, + health: false, + inventory: false, + }); + + t.diagnostic("phase: public mutation result is recursively safe"); + const created = await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, id()), + ]); + assert.equal(created.ok, true); + assertSafeProjection(created); + const { server, operation } = created.value; + assert.equal(server.desired.configuration.revision, 1); + assert.deepEqual(server.desired.release, manifest.identity); + + t.diagnostic( + "phase: arbitrary endpoint and provider identities fail closed", + ); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + endpoint: "https://evil.example/mcp", + deploymentProvider: deploymentProvider(server, account), + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_invalid" }, + ); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account, { + hostname: `${server.name}.attacker.workers.dev`, + }), + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_invalid" }, + ); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account, { + hostname: "foreign.fixture.workers.dev", + }), + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_invalid" }, + ); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account), + resources: [ + resourceClaim(server, account, { + provider: { + kind: "durable-object", + accountId: account, + namespaceId: "9".repeat(32), + namespaceName: "unrelated-resource", + className: "NotesStore", + }, + }), + ], + }), + ]), + { ok: false, error: "hosted_mcp_invalid" }, + ); + + const ready = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(ready.ok, true, JSON.stringify(ready)); + assertSafeProjection(ready); + assertSafeProjection( + await registry.call("get", owner, [ + installation, + account, + server.serverId, + ]), + ); + assertSafeProjection( + await registry.call("list", owner, [installation, account, null]), + ); + assert.deepEqual( + await registry.call( + "get", + owner, + [installation, account, server.serverId], + { + subject: otherOwner, + }, + ), + { ok: false, error: "hosted_mcp_not_found" }, + ); + }, +); + +test( + "native update preserves active health and fences semantic regressions", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const owner = "hosted-owner-update"; + const installation = "3".repeat(32); + const account = "4".repeat(32); + await registry.seed(owner, installation, account); + let { server } = await createReady( + registry, + owner, + installation, + account, + manifest, + ); + + t.diagnostic("phase: establish independently versioned health axes"); + let result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + assert.equal(server.observed.readiness.state, "ready"); + assertSafeProjection(result); + + const manifestV2 = version(manifest, "2.0.0", "b"); + const updateArgsV2 = [ + installation, + account, + server.serverId, + server.revision, + manifestV2, + intentFor(manifestV2, id(), "v2"), + ]; + result = await registry.call("update", owner, updateArgsV2); + assert.equal(result.ok, true, JSON.stringify(result)); + assertSafeProjection(result); + let operation; + ({ server, operation } = result.value); + const operationStartRevision = operation.startRevision; + await registry.call("parent-status", owner, [], { + installationId: installation, + status: "failed", + }); + await registry.restart(); + const updateReplay = await registry.call("update", owner, updateArgsV2); + assert.equal(updateReplay.ok, true, JSON.stringify(updateReplay)); + assert.equal( + updateReplay.value.operation.operationId, + operation.operationId, + ); + await registry.call("parent-status", owner, [], { + installationId: installation, + status: "ready", + }); + + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "updating", { + resources: [ + resourceClaim(server, account, { + state: "retained", + cleanup: "retained", + }), + ], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + + t.diagnostic("phase: health remains writable while update is pending"); + result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + credential: { + ...server.observed.credential, + revision: server.observed.credential.revision + 1, + updatedAt: new Date().toISOString(), + }, + connection: { + ...server.observed.connection, + revision: server.observed.connection.revision + 1, + catalogDigest: "b".repeat(64), + }, + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + assert.ok(server.revision > operationStartRevision); + + t.diagnostic( + "phase: failed staging preserves active release and readiness", + ); + const endpoint = server.endpoint; + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "failed", "failed", { + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assertSafeProjection(result); + ({ server, operation } = result.value); + assert.equal(operation.status, "failed"); + assert.equal(operation.failure.code, "deployment_failed"); + assert.deepEqual(server.activeRelease, manifest.identity); + assert.equal(server.endpoint, endpoint); + assert.equal(server.observed.deployment.state, "provisioned"); + assert.equal(server.observed.readiness.state, "ready"); + + t.diagnostic("phase: credential and connection generations cannot regress"); + assert.deepEqual( + await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + credential: { ...server.observed.credential, state: "revoked" }, + connection: { + ...server.observed.connection, + state: "disconnected", + catalogDigest: "b".repeat(64), + }, + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + assert.deepEqual( + await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + connection: { + ...server.observed.connection, + id: "mcp-11111111-1111-1111-1111-111111111111", + revision: server.observed.connection.revision + 1, + catalogDigest: "b".repeat(64), + }, + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + + t.diagnostic("phase: unrelated template updates are refused"); + const unrelated = version(manifestV2, "3.0.0", "c"); + unrelated.identity.templateId = "unrelated-template"; + unrelated.compatibility.rollbackTargets[0].identity.templateId = + "unrelated-template"; + assert.deepEqual( + await registry.call("update", owner, [ + installation, + account, + server.serverId, + server.revision, + unrelated, + intentFor(unrelated, id(), "foreign"), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + + t.diagnostic("phase: recovery promotes the compatible staged release"); + result = await registry.call("lifecycle", owner, [ + installation, + account, + server.serverId, + server.revision, + { + schemaVersion: 1, + requestId: id(), + action: "recover", + resourceDisposition: server.desired.resourceDisposition, + }, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + ({ server, operation } = result.value); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "provisioned", { + sequence: 1, + deploymentProvider: deploymentProvider(server, account, { + versionId: "2".repeat(32), + }), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assertSafeProjection(result); + assert.deepEqual(result.value.server.activeRelease, manifestV2.identity); + assert.equal(result.value.server.observed.readiness.state, "unknown"); + assert.equal(result.value.server.observed.connection.state, "connecting"); + + server = result.value.server; + await registry.restart(); + result = await registry.call("get", owner, [ + installation, + account, + server.serverId, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + operation = server.operation; + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "updating", { + sequence: 2, + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + sequence: 2, + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + result = await registry.call("health", owner, [ + installation, + account, + server.serverId, + health(server, { + connection: { + ...server.observed.connection, + revision: server.observed.connection.revision + 1, + state: "ready", + catalogDigest: "c".repeat(64), + }, + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.value.observed.readiness.state, "ready"); + + t.diagnostic("phase: removed resources follow a monotonic cleanup plan"); + server = result.value; + const manifestV3 = version(manifestV2, "3.0.0", "d"); + manifestV3.resources = []; + manifestV3.data.resources = []; + manifestV3.runtime.bindings = manifestV3.runtime.bindings.filter( + ({ kind }) => kind !== "durable-object", + ); + result = await registry.call("update", owner, [ + installation, + account, + server.serverId, + server.revision, + manifestV3, + intentFor(manifestV3, id(), "v3"), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assertSafeProjection(result); + ({ server, operation } = result.value); + + t.diagnostic("phase: destructive cleanup is refused before cutover"); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "updating", { + resources: [ + resourceClaim(server, account, { + state: "deleting", + cleanup: "pending", + }), + ], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "failed", "failed", { + resources: [ + resourceClaim(server, account, { + state: "deleted", + cleanup: "complete", + }), + ], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + await registry.restart(); + result = await registry.call("get", owner, [ + installation, + account, + server.serverId, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value; + operation = server.operation; + assert.deepEqual(server.activeRelease, manifestV2.identity); + assert.equal(server.observed.readiness.state, "ready"); + assert.equal(server.observed.resources[0].state, "ready"); + + t.diagnostic("phase: cutover precedes destructive cleanup"); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "provisioned", { + sequence: 1, + deploymentProvider: deploymentProvider(server, account, { + versionId: "3".repeat(32), + }), + resources: [resourceClaim(server, account)], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + assert.deepEqual(server.activeRelease, manifestV3.identity); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "provisioned", { + sequence: 2, + resources: [ + resourceClaim(server, account, { + state: "deleting", + cleanup: "pending", + }), + ], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + assert.equal(server.observed.resources[0].state, "deleting"); + assert.deepEqual( + await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "progress", "provisioned", { + sequence: 3, + resources: [resourceClaim(server, account)], + }), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + sequence: 3, + resources: [ + resourceClaim(server, account, { + state: "deleted", + cleanup: "complete", + }), + ], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assertSafeProjection(result); + assert.deepEqual(result.value.server.observed.resources, []); + + t.diagnostic("phase: a re-added logical key receives a fresh generation"); + server = result.value.server; + await registry.restart(); + const manifestV4 = version(manifestV3, "4.0.0", "e"); + manifestV4.resources = structuredClone(manifest.resources); + manifestV4.data.resources = structuredClone(manifest.data.resources); + manifestV4.runtime.bindings = structuredClone(manifest.runtime.bindings); + const generation = server.revision + 1; + result = await registry.call("update", owner, [ + installation, + account, + server.serverId, + server.revision, + manifestV4, + intentFor(manifestV4, id(), "v4"), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + ({ server, operation } = result.value); + const newProviderName = providerName(server, `${generation}-1`); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "provisioned", { + deploymentProvider: deploymentProvider(server, account, { + versionId: "4".repeat(32), + }), + resources: [ + resourceClaim(server, account, { + provider: { + kind: "durable-object", + accountId: account, + namespaceId: "8".repeat(32), + namespaceName: newProviderName, + className: "NotesStore", + }, + }), + ], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + server = result.value.server; + assert.equal(server.observed.resources[0].state, "ready"); + assert.equal( + ( + await registry.call("coordinate", owner, [ + installation, + account, + "remove", + ]) + ).ok, + true, + ); + const history = await registry.call("inventory", owner, [ + installation, + account, + server.serverId, + false, + ]); + assert.equal(history.ok, true, JSON.stringify(history)); + assert.equal(history.value.resources.length, 2); + assert.equal( + new Set(history.value.resources.map(({ inventoryId }) => inventoryId)) + .size, + 2, + ); + assert.deepEqual( + new Set(history.value.resources.map(({ state }) => state)), + new Set(["deleted", "ready"]), + ); + }, +); + +test( + "native removal fence survives restart and retains privileged cleanup inventory", + { timeout: 90_000 }, + async (t) => { + const registry = await registryHarness(t); + const manifest = await fixtureManifest(); + const owner = "hosted-owner-removal"; + const otherOwner = "hosted-owner-removal-other"; + const installation = "5".repeat(32); + const account = "6".repeat(32); + await registry.seed(owner, installation, account); + let { server } = await createReady( + registry, + owner, + installation, + account, + manifest, + ); + + t.diagnostic("phase: parent removal durably fences new hosted mutation"); + let result = await registry.call("coordinate", owner, [ + installation, + account, + "remove", + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.value.ready, false); + assert.equal(result.value.activeServerCount, 1); + assertSafeProjection(result); + await registry.restart(); + assert.deepEqual( + await registry.call("create", owner, [ + installation, + account, + manifest, + intentFor(manifest, id()), + ]), + { ok: false, error: "hosted_mcp_conflict" }, + ); + + t.diagnostic("phase: fence still admits required server deletion"); + const deleteArgs = [ + installation, + account, + server.serverId, + server.revision, + { + schemaVersion: 1, + requestId: id(), + action: "delete", + resourceDisposition: [{ key: "notes", action: "retain" }], + }, + ]; + result = await registry.call("lifecycle", owner, deleteArgs); + assert.equal(result.ok, true, JSON.stringify(result)); + assertSafeProjection(result); + const deleteOperationId = result.value.operation.operationId; + let operation; + ({ server, operation } = result.value); + result = await registry.call("observe", owner, [ + installation, + account, + server.serverId, + observation(server, operation, "succeeded", "deleted", { + resources: [ + resourceClaim(server, account, { + state: "retained", + cleanup: "retained", + }), + ], + }), + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.ok(result.value.server.tombstonedAt); + result = await registry.call("finish-removal", owner, [ + installation, + account, + ]); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.value.ready, true); + assert.equal(result.value.retainedResourceCount, 1); + assertSafeProjection(result); + + t.diagnostic( + "phase: privileged retained inventory survives parent deletion", + ); + await registry.call("remove-parent", owner, [], { + installationId: installation, + }); + await registry.restart(); + const deleteReplay = await registry.call("lifecycle", owner, deleteArgs); + assert.equal(deleteReplay.ok, true, JSON.stringify(deleteReplay)); + assert.equal(deleteReplay.value.operation.operationId, deleteOperationId); + assert.equal(deleteReplay.value.operation.status, "succeeded"); + const inventory = await registry.call("inventory", owner, [ + installation, + account, + null, + true, + ]); + assert.equal(inventory.ok, true, JSON.stringify(inventory)); + assert.equal(inventory.value.resources.length, 1); + assert.equal( + inventory.value.resources[0].provider.namespaceName, + providerName(server), + ); + assert.deepEqual( + await registry.call( + "inventory", + owner, + [installation, account, null, true], + { + subject: otherOwner, + }, + ), + { ok: false, error: "hosted_mcp_not_found" }, + ); + assert.deepEqual( + await registry.call("get", owner, [ + installation, + account, + server.serverId, + ]), + { ok: false, error: "hosted_mcp_not_found" }, + ); + }, +); diff --git a/tsconfig.worker.json b/tsconfig.worker.json index 5d3ca60..73d6639 100644 --- a/tsconfig.worker.json +++ b/tsconfig.worker.json @@ -25,6 +25,7 @@ "tests/fixtures/effect-boundaries-worker.ts", "tests/fixtures/oauth-worker.ts", "tests/fixtures/ownership-worker.ts", + "tests/fixtures/hosted-mcp-registry-worker.ts", "tests/fixtures/orchestrator-worker.ts", "tests/fixtures/bridge-control-worker.ts", "tests/fixtures/bridge-customer-worker.ts", -- 2.51.2