import { DurableObject } from "cloudflare:workers"; import * as Effect from "effect/Effect"; import { z } from "zod"; import type { Env } from "./config.ts"; import { beginDomain, changeDomain, type DomainChanges, } from "./domain-lifecycle.ts"; import { domainRecord, domainStart, type DomainRecord, type DomainStart, } from "./domain-metadata.ts"; import { InstallationConflict, InstallationNotFound, InvalidMetadata, type RegistryFailure, } from "./installation-errors.ts"; import { beginInstallation, changeInstallation, } from "./installation-lifecycle.ts"; import { installationChanges, installationId, ownerSubject, parse, parseInstallation, registryName, releaseIdentity, sameRelease, type Installation, type InstallationChanges, type ReleaseIdentity, } from "./installation-metadata.ts"; import { registryResult, type RegistryResult } from "./registry-result.ts"; import { transaction } from "./storage.ts"; import { installationRecovery, operationChanges, operationSchema, upgradeBaseline, type InstallationOperation, type InstallationRecovery, type UpgradeBaseline, } from "./operation.ts"; const rowKey = (id: string) => `installation:${id}`; const replaySchema = z.strictObject({ installationId, accountId: installationId, }); const operationReplaySchema = z.strictObject({ installationId, operationId: installationId, release: releaseIdentity, }); export interface InstallationPage { installations: Installation[]; nextCursor: string | null; } export interface InstallationState { installation: Installation; operation: InstallationOperation; } // Private control-plane binding only. One owner object keeps reservation replay, // records and owner listing in the same native atomic storage boundary. export class InstallationRegistry extends DurableObject { #owner(subject: string) { parse(ownerSubject, subject); if ( !this.ctx.id.equals( this.env.INSTALLATIONS.idFromName(registryName(subject)), ) ) throw new InstallationNotFound(); } #record(value: unknown, subject: string, id: string) { const record = parseInstallation(value); if (record.ownerSubject !== subject || record.installationId !== id) throw new InstallationNotFound(); return record; } #execute( subject: string, program: Effect.Effect, ): Promise> { return registryResult( Effect.suspend(() => { this.#owner(subject); return program; }), ); } #readRecord( storage: DurableObjectStorage | DurableObjectTransaction, subject: string, id: string, ) { return Effect.promise(() => storage.get(rowKey(id))).pipe( Effect.map((value) => this.#record(value, subject, id)), ); } reserve( subject: string, accountId: string, requestId: string, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, accountId); parse(installationId, requestId); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const replayKey = `request:${requestId}`; const saved = yield* Effect.promise(() => tx.get(replayKey)); if (saved !== undefined) { const replay = parse(replaySchema, saved); if (replay.accountId !== accountId) return yield* new InstallationConflict(); const record = yield* this.#readRecord( tx, subject, replay.installationId, ); if (record.accountId !== replay.accountId) return yield* new InvalidMetadata(); return record; } const id = crypto.randomUUID().replaceAll("-", ""); if (yield* Effect.promise(() => tx.get(rowKey(id)))) return yield* new InstallationConflict(); const now = Date.now(); const record = parseInstallation({ schemaVersion: 1, installationId: id, ownerSubject: subject, accountId, createdAt: now, updatedAt: now, revision: 1, resources: { workerName: `flarebot-${id}`, sandboxApplicationName: `flarebot-shell-${id}`, personalAgentNamespaceId: null, sandboxNamespaceId: null, sandboxApplicationId: null, runtimeOrigin: null, }, desiredRelease: null, installedRelease: null, status: "reserved", errorCode: null, operationId: null, }); yield* Effect.promise(() => tx.put({ [rowKey(id)]: record, [replayKey]: parse(replaySchema, { installationId: id, accountId, }), }), ); return record; }), ); }), ); } get( subject: string, id: string, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); const value = yield* Effect.promise(() => this.ctx.storage.get(rowKey(id)), ); return value === undefined ? null : this.#record(value, subject, id); }), ); } getDomain( subject: string, id: string, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); yield* this.#readRecord(this.ctx.storage, subject, id); const value = yield* Effect.promise(() => this.ctx.storage.get(`domain:${id}`), ); return value === undefined ? null : parse(domainRecord, value); }), ); } domainReplay( subject: string, id: string, requestId: string, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse(installationId, requestId); yield* this.#readRecord(this.ctx.storage, subject, id); const value = yield* Effect.promise(() => this.ctx.storage.get(`domain-request:${id}:${requestId}`), ); return value === undefined ? null : parse( z.strictObject({ intent: domainStart, operationId: installationId, }), value, ).operationId; }), ); } startDomain( subject: string, id: string, requestId: string, input: DomainStart, deadline = Date.now() + 20 * 60_000, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse(installationId, requestId); const intent = parse(domainStart, input); parse( z .number() .int() .min(Date.now() + 30_000) .max(Date.now() + 30 * 60_000), deadline, ); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const installation = yield* this.#readRecord(tx, subject, id); if ( !installation.installedRelease || installation.status !== "ready" ) return yield* new InstallationConflict(); const key = `domain:${id}`; const saved = yield* Effect.promise(() => tx.get(key)); const current = saved === undefined ? null : parse(domainRecord, saved); const replayKey = `domain-request:${id}:${requestId}`; const replay = yield* Effect.promise(() => tx.get<{ intent: DomainStart; operationId: string; }>(replayKey), ); if (replay) { if ( JSON.stringify(replay.intent) !== JSON.stringify(intent) || current?.operationId !== replay.operationId ) return yield* new InstallationConflict(); return current!; } const value = beginDomain( current, intent, crypto.randomUUID().replaceAll("-", ""), deadline, ); yield* Effect.promise(() => tx.put({ [key]: value, [replayKey]: { intent, operationId: value.operationId }, }), ); return value; }), ); }), ); } updateDomain( subject: string, id: string, operationId: string, revision: number, changes: DomainChanges, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse(installationId, operationId); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { yield* this.#readRecord(tx, subject, id); const key = `domain:${id}`; const current = parse( domainRecord, yield* Effect.promise(() => tx.get(key)), ); const next = changeDomain(current, operationId, revision, changes); yield* Effect.promise(() => tx.put(key, next)); return next; }), ); }), ); } list( subject: string, cursor: string | null = null, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { if (cursor !== null) parse(installationId, cursor); const rows = yield* Effect.promise(() => this.ctx.storage.list({ prefix: "installation:", limit: 51, ...(cursor === null ? {} : { startAfter: rowKey(cursor) }), }), ); const records = [...rows].map(([key, value]) => this.#record(value, subject, key.slice("installation:".length)), ); const installations = records.slice(0, 50); return { installations, nextCursor: records.length > 50 ? installations.at(-1)!.installationId : null, }; }), ); } start( subject: string, id: string, accountId: string, requestId: string, release: ReleaseIdentity, deadline: number, recovery: InstallationRecovery | null = null, upgrade: UpgradeBaseline | null = null, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse(installationId, accountId); parse(installationId, requestId); const desired = parse(releaseIdentity, release); if (upgrade) upgrade = parse(upgradeBaseline, upgrade); if (recovery) recovery = parse(installationRecovery, recovery); parse( z .number() .int() .min(Date.now() + 30_000) .max(Date.now() + 30 * 60_000), deadline, ); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const current = yield* this.#readRecord(tx, subject, id); if (current.accountId !== accountId) return yield* new InstallationConflict(); const domain = yield* Effect.promise(() => tx.get(`domain:${id}`)); if ( domain && ["connecting", "removing"].includes( parse(domainRecord, domain).status, ) ) return yield* new InstallationConflict(); const replayKey = `operation-request:${requestId}`; const replayValue = yield* Effect.promise(() => tx.get(replayKey)); const replay = replayValue === undefined ? null : parse(operationReplaySchema, replayValue); const previousValue = yield* Effect.promise(() => tx.get(`operation:${id}`), ); const previous = previousValue === undefined ? null : parse(operationSchema, previousValue); if (replay) { if ( replay.installationId !== id || replay.operationId !== current.operationId || !sameRelease(replay.release, desired) || !previous ) return yield* new InstallationConflict(); return { installation: current, operation: previous, }; } const operationId = crypto.randomUUID().replaceAll("-", ""); const { installation, operation } = beginInstallation( current, previous, { desired, deadline, recovery, upgrade, operationId, now: Math.max(Date.now(), current.updatedAt), }, ); yield* Effect.promise(() => tx.put({ [rowKey(id)]: installation, [`operation:${id}`]: operation, [replayKey]: parse(operationReplaySchema, { installationId: id, operationId, release: desired, }), }), ); return { installation, operation }; }), ); }), ); } replay( subject: string, id: string, requestId: string, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse(installationId, requestId); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const saved = yield* Effect.promise(() => tx.get(`operation-request:${requestId}`), ); if (saved === undefined) return null; const replay = parse(operationReplaySchema, saved); const current = yield* this.#readRecord(tx, subject, id); if ( replay.installationId !== id || replay.operationId !== current.operationId || !sameRelease(replay.release, current.desiredRelease) ) return yield* new InstallationConflict(); return current; }), ); }), ); } operation( subject: string, id: string, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); const value = yield* Effect.promise(() => this.ctx.storage.get(`operation:${id}`), ); return value === undefined ? null : parse(operationSchema, value); }), ); } active( subject: string, id: string, operationId: string, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse(installationId, operationId); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const installation = yield* this.#readRecord(tx, subject, id); const operation = parse( operationSchema, yield* Effect.promise(() => tx.get(`operation:${id}`)), ); if ( installation.operationId !== operationId || operation.operationId !== operationId || !["installing", "updating"].includes(installation.status) ) return yield* new InstallationConflict(); return { installation, operation }; }), ); }), ); } intent( subject: string, id: string, operationId: string, expectedRevision: number, changes: Partial< Pick< InstallationOperation, | "workerIntent" | "containerIntent" | "containerUpdate" | "rolloutIntent" | "rolloutId" > >, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse(installationId, operationId); const input = parse(operationChanges, changes); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const installation = yield* this.#readRecord(tx, subject, id); const operation = parse( operationSchema, yield* Effect.promise(() => tx.get(`operation:${id}`)), ); if ( installation.operationId !== operationId || operation.operationId !== operationId || installation.revision !== expectedRevision || !["installing", "updating"].includes(installation.status) ) return yield* new InstallationConflict(); const next = parse(operationSchema, { ...operation, ...input }); // Intent cannot be cleared or rewritten to conceal an ambiguous write. for (const field of [ "workerIntent", "containerIntent", "containerUpdate", "rolloutIntent", "rolloutId", ] as const) if (operation[field] && field in input) return yield* new InstallationConflict(); yield* Effect.promise(() => tx.put(`operation:${id}`, next)); return next; }), ); }), ); } update( subject: string, id: string, expectedRevision: number, input: InstallationChanges, ): Promise> { return this.#execute( subject, Effect.gen({ self: this }, function* () { parse(installationId, id); parse( z.number().int().positive().max(Number.MAX_SAFE_INTEGER), expectedRevision, ); const changes = parse(installationChanges, input); return yield* transaction(this.ctx.storage, (tx) => Effect.gen({ self: this }, function* () { const value = yield* Effect.promise(() => tx.get(rowKey(id))); if (value === undefined) return yield* new InstallationNotFound(); const current = this.#record(value, subject, id); const next = changeInstallation( current, expectedRevision, changes, Math.max(Date.now(), current.updatedAt), ); yield* Effect.promise(() => tx.put(rowKey(id), next)); return next; }), ); }), ); } }