import { parseSSEStream, type ExecEvent } from "@cloudflare/sandbox"; import type { Agent } from "agents"; import * as Data from "effect/Data"; import * as Effect from "effect/Effect"; import type { InstallationConfig } from "../configuration/customer"; import { verifyAssertion } from "../server/bridge"; import { decode, DOMAIN_CONFIGURATION_PURPOSE, DOMAIN_HEALTH_PURPOSE, HEALTH_PURPOSE, opaque, type BridgeClaims, } from "../shared/bridge"; import { customDomainOrigin } from "../shared/domain-origin"; import type { BridgeStore } from "./bridge-store"; class ManagementFailure extends Data.TaggedError("ManagementFailure")<{ message: string; }> {} interface HealthSandbox { execTemporary( command: string, expiresAt: number, id: string, ): Promise>; closeTemporary(): Promise; } function managementCall(operation: () => PromiseLike, message: string) { return Effect.tryPromise({ try: operation, catch: () => new ManagementFailure({ message }), }); } function assertionPayload(assertion: string) { return Effect.try({ try: (): Record => { const parts = assertion.split("."); if (parts.length !== 3) throw new Error(); const claims: unknown = JSON.parse( new TextDecoder().decode(decode(parts[1])), ); if (!claims || typeof claims !== "object" || Array.isArray(claims)) throw new Error(); return claims as Record; }, catch: () => new ManagementFailure({ message: "Invalid management assertion" }), }); } export class InstanceManagement { constructor( private readonly sql: Agent["sql"], private readonly storage: Pick, private readonly installation: () => InstallationConfig, private readonly bridge: BridgeStore, private readonly sandbox: (id: string) => HealthSandbox, private readonly waitUntil: (work: Promise) => void, ) {} recover() { return Effect.gen({ self: this }, function* () { this.bridge.initialize(); for (const id of this.bridge.pending()) { if ( !(yield* managementCall( () => this.sandbox(id).closeTemporary(), "Health cleanup pending", )) ) return yield* Effect.fail( new ManagementFailure({ message: "Health cleanup pending" }), ); this.bridge.closed(id, false); } }); } configuredDomainOrigin() { return ( this.sql<{ origin: string | null }>`SELECT origin FROM flarebot_domain_configuration WHERE singleton = 1`[0]?.origin ?? null ); } configureDomain(assertion: string) { return Effect.gen({ self: this }, function* () { const installation = this.installation(); if (!installation.bridge) return yield* Effect.fail( new ManagementFailure({ message: "Domain configuration unavailable", }), ); const claims = yield* assertionPayload(assertion); const revision = Number(claims.state); const origin = claims.challenge === "" ? null : customDomainOrigin(claims.challenge); if ( !Number.isSafeInteger(revision) || revision < 1 || (claims.challenge && !origin) ) return yield* Effect.fail( new ManagementFailure({ message: "Invalid domain configuration" }), ); yield* verifyAssertion(assertion, installation.bridge, { iss: installation.controlPlaneOrigin, aud: installation.runtimeOrigin, sub: installation.ownerSubject, installationId: installation.installationId, purpose: DOMAIN_CONFIGURATION_PURPOSE, state: String(revision), challenge: origin ?? "", }); yield* Effect.try({ try: () => this.storage.transactionSync(() => { const current = this.sql<{ revision: number; origin: string | null; }>`SELECT revision, origin FROM flarebot_domain_configuration WHERE singleton = 1`[0]; if ( current && (revision < current.revision || (revision === current.revision && origin !== current.origin)) ) throw new ManagementFailure({ message: "Stale domain configuration", }); if (!current) this .sql`INSERT INTO flarebot_domain_configuration VALUES (1, ${revision}, ${origin})`; else if (revision > current.revision) this .sql`UPDATE flarebot_domain_configuration SET revision = ${revision}, origin = ${origin} WHERE singleton = 1`; }), catch: (error) => error instanceof ManagementFailure ? error : new ManagementFailure({ message: "Domain configuration unavailable", }), }); return { revision, origin }; }); } domainHealth(assertion: string, expectedOrigin: string) { return Effect.gen({ self: this }, function* () { const installation = this.installation(); if ( !installation.bridge || this.configuredDomainOrigin() !== expectedOrigin ) return yield* Effect.fail( new ManagementFailure({ message: "Domain unavailable" }), ); const claims = yield* assertionPayload(assertion); if (!opaque(claims.state) || !opaque(claims.challenge)) return yield* Effect.fail( new ManagementFailure({ message: "Invalid domain health" }), ); const { state, challenge } = claims; yield* verifyAssertion(assertion, installation.bridge, { iss: installation.controlPlaneOrigin, aud: expectedOrigin, sub: installation.ownerSubject, installationId: installation.installationId, purpose: DOMAIN_HEALTH_PURPOSE, state, challenge, }); return { installationId: installation.installationId, origin: expectedOrigin, state, challenge, }; }); } bootstrapHealth(claims: BridgeClaims) { return Effect.gen({ self: this }, function* () { const config = this.installation(); const release = config.release; if ( !release || !opaque(claims.state) || !opaque(claims.challenge) || !Number.isSafeInteger(claims.iat) || !Number.isSafeInteger(claims.exp) || claims.iat > Math.floor(Date.now() / 1000) || claims.exp <= claims.iat || claims.exp - claims.iat > 60 || claims.purpose !== HEALTH_PURPOSE || claims.iss !== config.controlPlaneOrigin || claims.aud !== config.runtimeOrigin || claims.sub !== config.ownerSubject || claims.installationId !== config.installationId || claims.operationId !== release.operationId || claims.artifactDigest !== release.artifactDigest || claims.version !== release.version ) return yield* Effect.fail( new ManagementFailure({ message: "Invalid health identity" }), ); const probe = this.bridge.beginHealth( claims.jti, claims.exp * 1000, release.operationId, release.artifactDigest, ); if (!probe.cached) { // The parent owns completion even if the requesting Worker disconnects. const work = Effect.runPromise(this.runProbe(probe.probeId)); this.waitUntil(work.catch(() => {})); yield* managementCall(() => work, "Native health probe unavailable"); } return { identity: "ready", nativeParent: "ready", sandbox: "booted-and-destroyed", } as const; }); } private runProbe(id: string) { return Effect.gen({ self: this }, function* () { const sandbox = this.sandbox(id); let successful = false; const failure = new ManagementFailure({ message: "Native health probe unavailable", }); yield* Effect.acquireUseRelease( Effect.succeed(sandbox), () => Effect.gen(function* () { const stream = yield* managementCall( () => sandbox.execTemporary( "printf 'flarebot-native-health-v1'", Date.now() + 45_000, id, ), failure.message, ); const controller = new AbortController(); const events = parseSSEStream(stream, controller.signal)[ Symbol.asyncIterator ](); yield* Effect.acquireUseRelease( Effect.succeed(events), () => Effect.gen(function* () { let output = ""; while (true) { const next = yield* managementCall( () => events.next(), failure.message, ); if (next.done) break; const event = next.value; if (event.type === "stdout") { output += event.data ?? ""; if (output.length > 128) return yield* Effect.fail(failure); } else if ( event.type === "stderr" || event.type === "error" ) return yield* Effect.fail(failure); else if (event.type === "complete") { successful = (event.exitCode ?? event.result?.exitCode) === 0 && output === "flarebot-native-health-v1"; break; } } }), () => Effect.gen(function* () { controller.abort(); yield* managementCall( () => events.return?.() ?? Promise.resolve({ done: true as const, value: undefined, }), failure.message, ).pipe( Effect.timeoutOrElse({ duration: 5_000, orElse: () => Effect.void, }), Effect.ignore, ); }), ); }).pipe( Effect.timeoutOrElse({ duration: 45_000, orElse: () => Effect.fail(failure), }), ), () => managementCall(() => sandbox.closeTemporary(), failure.message).pipe( Effect.tap((closed) => Effect.sync(() => { if (closed) this.bridge.closed(id, successful); else successful = false; }), ), Effect.timeoutOrElse({ duration: 5_000, orElse: () => Effect.fail(failure), }), Effect.orDie, ), ); if (!successful) return yield* Effect.fail(failure); }); } }