Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
20 kB · 613 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614import { 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<Env> { #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<A>( subject: string, program: Effect.Effect<A, RegistryFailure>, ): Promise<RegistryResult<A>> { 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<RegistryResult<Installation>> { 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<RegistryResult<Installation | null>> { 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<RegistryResult<DomainRecord | null>> { 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<RegistryResult<string | null>> { 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<RegistryResult<DomainRecord>> { 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<RegistryResult<DomainRecord>> { 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<RegistryResult<InstallationPage>> { 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<RegistryResult<InstallationState>> { 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<RegistryResult<Installation | null>> { 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<RegistryResult<InstallationOperation | null>> { 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<RegistryResult<InstallationState>> { 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<RegistryResult<InstallationOperation>> { 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<RegistryResult<Installation>> { 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; }), ); }), ); }}