diff --git a/apps/ai/migrations/chat-session/0000_baseline.sql b/apps/ai/migrations/chat-session/0000_baseline.sql new file mode 100644 index 000000000..7555de8ba --- /dev/null +++ b/apps/ai/migrations/chat-session/0000_baseline.sql @@ -0,0 +1,17 @@ +-- The ChatSession Durable Object's SQLite schema. Idempotent: objects that predate these +-- migration files already have the tables and adopt this as their first applied file. +-- `payload` is the encoded `ChatEvent` minus its `seq`, which is the key. +CREATE TABLE IF NOT EXISTS events ( + seq INTEGER PRIMARY KEY AUTOINCREMENT, + created_at INTEGER NOT NULL, + payload TEXT NOT NULL +); +CREATE TABLE IF NOT EXISTS session ( + id INTEGER PRIMARY KEY CHECK (id = 1), + running INTEGER NOT NULL DEFAULT 0, + running_since INTEGER, + running_message_id TEXT, + running_input TEXT, + running_resumes INTEGER +); +INSERT OR IGNORE INTO session (id, running) VALUES (1, 0); diff --git a/apps/ai/src/chat/ChatSession.ts b/apps/ai/src/chat/ChatSession.ts index cf6234008..058c970ec 100644 --- a/apps/ai/src/chat/ChatSession.ts +++ b/apps/ai/src/chat/ChatSession.ts @@ -34,7 +34,7 @@ * for the route graph. */ import * as Cloudflare from "alchemy/Cloudflare" -import { Effect, Option, Schema } from "effect" +import { Effect, FileSystem, Option, Path, Schema } from "effect" import { ChatTurnOrigin as ChatTurnOriginSchema, ChatTurnTenant, @@ -83,23 +83,40 @@ interface SessionRow extends Record { readonly running_resumes: number | null } -/** Rows the DO writes. `payload` is the encoded `ChatEvent` minus its `seq`, which is the key. */ -const SCHEMA = ` -CREATE TABLE IF NOT EXISTS events ( - seq INTEGER PRIMARY KEY AUTOINCREMENT, - created_at INTEGER NOT NULL, - payload TEXT NOT NULL -); -CREATE TABLE IF NOT EXISTS session ( - id INTEGER PRIMARY KEY CHECK (id = 1), - running INTEGER NOT NULL DEFAULT 0, - running_since INTEGER, - running_message_id TEXT, - running_input TEXT, - running_resumes INTEGER -); -INSERT OR IGNORE INTO session (id, running) VALUES (1, 0); -` +/** + * The schema's migration files, relative to the repo root where alchemy runs. Alchemy embeds them + * in the bundle at deploy; each activation applies the pending ones before the first call. + */ +export const CHAT_SESSION_MIGRATIONS = "apps/ai/migrations/chat-session" + +/** Columns added after the first live objects were created, with their declared types. */ +const LATE_SESSION_COLUMNS = { + running_since: "INTEGER", + running_message_id: "TEXT", + running_input: "TEXT", + running_resumes: "INTEGER", +} as const + +interface ColumnRow extends Record { + readonly name: string +} + +/** + * Bring a `session` table that predates the migration files up to the baseline's shape, which + * `CREATE TABLE IF NOT EXISTS` alone cannot do. A fresh object has no table, so nothing runs. + */ +export const addMissingSessionColumns = (sql: SqlStorage): void => { + const present = new Set( + sql + .exec("SELECT name FROM pragma_table_info('session')") + .toArray() + .map((row) => row.name), + ) + if (present.size === 0) return + for (const [column, type] of Object.entries(LATE_SESSION_COLUMNS)) { + if (!present.has(column)) sql.exec(`ALTER TABLE session ADD COLUMN ${column} ${type}`) + } +} /** * How long a subscription sits silent before the connection is recycled. Cloudflare caps a request @@ -267,21 +284,8 @@ export class ChatSession { private readonly applier: ProposalApplier = applyThroughWorker, private readonly runner: TurnRunner = runThroughWorker, ) { + // The schema is the activation's to apply (`makeChatSessionActivation`), before this runs. this.sql = ctx.storage.sql - this.sql.exec(SCHEMA) - // Sessions created before the watchdog columns existed (local dev only — the class has - // never been deployed) would otherwise fail every read against `session`. - for (const column of [ - "running_since INTEGER", - "running_message_id TEXT", - "running_input TEXT", - "running_resumes INTEGER", - ]) { - // A failure means the column is already present. - Effect.runSync( - Effect.ignore(Effect.try(() => this.sql.exec(`ALTER TABLE session ADD COLUMN ${column}`))), - ) - } } /** Highest assigned seq, i.e. the cursor a client that has read everything holds. */ @@ -827,42 +831,63 @@ const isChatSessionNamespace = ( scope: Cloudflare.DurableObject, ): scope is Cloudflare.DurableObject => scope.name === "ChatSession" +/** The part of alchemy's `SqlMigrations` the activation uses; a test applies the files itself. */ +export interface ChatSessionSchema { + readonly apply: () => Effect.Effect +} + /** - * One activation, in alchemy's two phases: the outer Effect resolves the state, env and namespace - * (it also runs at plan time, against a mock state, so it must not touch storage), the inner one - * builds the session inside the object's `blockConcurrencyWhile`, so the schema statements have - * run before the first call reaches it, hibernation wakes included. The turn and apply graphs reach - * other sessions through the namespace (`ChatSessions`). + * One activation, in alchemy's two phases: the outer Effect resolves the state, env, namespace and + * migration files (it also runs at plan time, against a mock state, so it must not touch storage), + * the inner one migrates and builds the session inside the object's `blockConcurrencyWhile`, so the + * schema is current before the first call reaches it, hibernation wakes included. The turn and + * apply graphs reach other sessions through the namespace (`ChatSessions`). */ -export const activateChatSession = Effect.map( - Effect.all([ - Cloudflare.DurableObjectState, - Cloudflare.WorkerEnvironment, - WorkersAiGateway, - Effect.serviceOption(Cloudflare.DurableObjectScope), - ]), - ([state, env, workersAi, scope]) => { - // Always there in the isolate; absent only where a test builds the activation by hand. - const chatSessions: ChatSessionNamespace | undefined = Option.getOrUndefined( - Option.filter(scope, isChatSessionNamespace), - ) - return Effect.sync(() => - chatSessionRpc( - new ChatSession( - state.raw, - env, - (input) => applyThroughWorker({ ...input, chatSessions }), - (input) => runThroughWorker({ ...input, workersAi, chatSessions }), - ), - ), - ) - }, +export const makeChatSessionActivation = ( + migrations: Effect.Effect, never, R>, +) => + Effect.map( + Effect.all([ + Cloudflare.DurableObjectState, + Cloudflare.WorkerEnvironment, + WorkersAiGateway, + Effect.serviceOption(Cloudflare.DurableObjectScope), + migrations, + ]), + ([state, env, workersAi, scope, schema]) => { + // Always there in the isolate; absent only where a test builds the activation by hand. + const chatSessions: ChatSessionNamespace | undefined = Option.getOrUndefined( + Option.filter(scope, isChatSessionNamespace), + ) + return Effect.gen(function* () { + yield* Effect.sync(() => addMissingSessionColumns(state.raw.storage.sql)) + // An object whose schema cannot be brought current cannot serve a single call. + yield* schema.apply().pipe(Effect.orDie) + return chatSessionRpc( + new ChatSession( + state.raw, + env, + (input) => applyThroughWorker({ ...input, chatSessions }), + (input) => runThroughWorker({ ...input, workersAi, chatSessions }), + ), + ) + }) + }, + ) + +export const activateChatSession = makeChatSessionActivation( + Cloudflare.SqlMigrations(CHAT_SESSION_MIGRATIONS), ) /** The activation, as the layer the host Worker provides. */ // The activation's requirements are named rather than inferred: `.make` discharges -// `DurableObjectServices` (both of these) through its own `Exclude`, while inference would widen +// `DurableObjectServices` (all but the gateway) through its own `Exclude`, while inference would widen // them into the layer's requirements and surface them all the way up in `alchemy.run.ts`. export const ChatSessionLive = ChatSessionObject.make< - Cloudflare.DurableObjectState | Cloudflare.WorkerEnvironment | WorkersAiGateway + | Cloudflare.DurableObjectState + | Cloudflare.WorkerEnvironment + | Cloudflare.Worker + | WorkersAiGateway + | FileSystem.FileSystem + | Path.Path >(activateChatSession) diff --git a/apps/ai/src/chat/ChatSessionObject.test.ts b/apps/ai/src/chat/ChatSessionObject.test.ts index be25a0458..f289c3e42 100644 --- a/apps/ai/src/chat/ChatSessionObject.test.ts +++ b/apps/ai/src/chat/ChatSessionObject.test.ts @@ -2,9 +2,9 @@ import { assert, describe, it } from "@effect/vitest" import { encodeChatTurnTenant, type ChatTurnTenant } from "@maple/domain/chat-session" import * as Cloudflare from "alchemy/Cloudflare" import { Effect, Layer, Option } from "effect" -import { makeFakeDurableObjectState } from "../../test/chat/fake-do-state" +import { applyChatSessionMigrations, makeFakeDurableObjectState } from "../../test/chat/fake-do-state" import { WorkersAiGateway } from "../platform/WorkersAiHttpClient" -import { activateChatSession } from "./ChatSession" +import { makeChatSessionActivation } from "./ChatSession" const TENANT = encodeChatTurnTenant({ orgId: "org_test" as ChatTurnTenant["orgId"], @@ -13,22 +13,34 @@ const TENANT = encodeChatTurnTenant({ authMode: "self_hosted", }) +/** Alchemy's `SqlMigrations` captures the files at deploy; a test applies the same files directly. */ +const activateChatSession = Effect.flatMap(Cloudflare.DurableObjectState, (state) => + makeChatSessionActivation( + Effect.succeed({ + apply: () => Effect.sync(() => applyChatSessionMigrations(state.raw.storage.sql)), + }), + ), +) + /** One activation the way alchemy's Durable Object bridge performs it: the outer phase under the state, then the inner. */ -const activate = Effect.gen(function* () { - const state = makeFakeDurableObjectState() - // SAFETY: the fake carries the `storage.sql` and `waitUntil` the class reads, and nothing else. - const raw = state as unknown as import("@cloudflare/workers-types").DurableObjectState - const build = yield* activateChatSession.pipe( - Effect.provide( - Layer.mergeAll( - Layer.succeed(Cloudflare.DurableObjectState, Cloudflare.fromDurableObjectState(raw)), - Layer.succeed(Cloudflare.WorkerEnvironment, {}), - Layer.succeed(WorkersAiGateway, Option.none()), +const activateOn = (state: ReturnType) => + Effect.gen(function* () { + // SAFETY: the fake carries the `storage.sql` and `waitUntil` the class reads, and nothing else. + const raw = state as unknown as import("@cloudflare/workers-types").DurableObjectState + const build = yield* activateChatSession.pipe( + Effect.provide( + Layer.mergeAll( + Layer.succeed(Cloudflare.DurableObjectState, Cloudflare.fromDurableObjectState(raw)), + Layer.succeed(Cloudflare.WorkerEnvironment, {}), + Layer.succeed(WorkersAiGateway, Option.none()), + ), ), - ), - ) - return { rpc: yield* build, state } -}) + ) + return { rpc: yield* build, state } + }) + +/** A fresh object per test, so no test reads another's events. */ +const activate = () => activateOn(makeFakeDurableObjectState({ migrated: false })) describe("the ChatSession Durable Object on alchemy's form", () => { it.effect("the outer phase touches no storage, so it can run against alchemy's plan-time mock", () => @@ -51,7 +63,7 @@ describe("the ChatSession Durable Object on alchemy's form", () => { it.effect("exposes the stub's surface over the session it built", () => Effect.gen(function* () { - const { rpc } = yield* activate + const { rpc } = yield* activate() assert.strictEqual(yield* rpc.cursor(), 0) const seq = yield* rpc.append({ type: "user-message", id: "u1", text: "hello" }) assert.strictEqual(seq, 1) @@ -66,7 +78,7 @@ describe("the ChatSession Durable Object on alchemy's form", () => { it.effect("begins a turn over RPC and lets the class own it", () => Effect.gen(function* () { - const { rpc, state } = yield* activate + const { rpc, state } = yield* activate() const begun = yield* rpc.beginTurn({ sessionId: "org_test:tab", messageId: "m1", @@ -84,4 +96,27 @@ describe("the ChatSession Durable Object on alchemy's form", () => { assert.strictEqual(yield* rpc.running(), false) }), ) + + it.effect("brings a session table from before the late columns up to the baseline", () => + Effect.gen(function* () { + const state = makeFakeDurableObjectState({ migrated: false }) + state.storage.sql.exec( + "CREATE TABLE session (id INTEGER PRIMARY KEY CHECK (id = 1), running INTEGER NOT NULL DEFAULT 0)", + ) + state.storage.sql.exec("INSERT INTO session (id, running) VALUES (1, 0)") + const { rpc } = yield* activateOn(state) + + assert.strictEqual(yield* rpc.running(), false) + const columns = state.storage.sql + .exec("SELECT name FROM pragma_table_info('session')") + .toArray() + .map((row) => row.name) + assert.includeMembers(columns, [ + "running_since", + "running_message_id", + "running_input", + "running_resumes", + ]) + }), + ) }) diff --git a/apps/ai/test/chat/fake-do-state.ts b/apps/ai/test/chat/fake-do-state.ts index 813e9a41d..bbd39729f 100644 --- a/apps/ai/test/chat/fake-do-state.ts +++ b/apps/ai/test/chat/fake-do-state.ts @@ -9,7 +9,21 @@ * `node:sqlite` is the same engine Durable Object storage exposes, and the `SqlStorage` surface the * class touches is small: `exec(sql, ...bindings)` returning a cursor with `one()` and `toArray()`. */ +import { readdirSync, readFileSync } from "node:fs" +import path from "node:path" import { DatabaseSync } from "node:sqlite" +import { CHAT_SESSION_MIGRATIONS } from "../../src/chat/ChatSession" + +const migrationsDir = path.resolve(import.meta.dirname, "../../../..", CHAT_SESSION_MIGRATIONS) + +/** Every migration file, in order, as alchemy's `apply` runs them (minus its history table). */ +export const applyChatSessionMigrations = (sql: SqlStorage): void => { + for (const file of readdirSync(migrationsDir) + .filter((name) => name.endsWith(".sql")) + .toSorted()) { + sql.exec(readFileSync(path.join(migrationsDir, file), "utf8")) + } +} /** Everything the DO under test reads off its state, and nothing more. */ export interface FakeDurableObjectState { @@ -24,14 +38,15 @@ export interface FakeDurableObjectState { readonly pending: Array> } -export const makeFakeDurableObjectState = (): FakeDurableObjectState => { +/** `migrated: false` leaves the database empty, for a test that migrates through the activation. */ +export const makeFakeDurableObjectState = ({ migrated = true } = {}): FakeDurableObjectState => { const db = new DatabaseSync(":memory:") const pending: Array> = [] const alarms: Array = [] const sql = { exec: (statement: string, ...bindings: ReadonlyArray) => { - // `ChatSession`'s CREATE TABLE block is several statements in one string; `node:sqlite` + // A migration file is several statements in one string; `node:sqlite` // splits those only through `exec`, while parameterised statements need `prepare`. if (bindings.length === 0 && /;\s*\S/.test(statement.trim())) { db.exec(statement) @@ -48,6 +63,8 @@ export const makeFakeDurableObjectState = (): FakeDurableObjectState => { }, } + if (migrated) applyChatSessionMigrations(sql as SqlStorage) + return { storage: { sql: sql as SqlStorage, diff --git a/apps/chat-bot/src/relay/ConnectorRelay.test.ts b/apps/chat-bot/src/relay/ConnectorRelay.test.ts index b0c182c65..8930aa44f 100644 --- a/apps/chat-bot/src/relay/ConnectorRelay.test.ts +++ b/apps/chat-bot/src/relay/ConnectorRelay.test.ts @@ -13,7 +13,12 @@ import { chatConnectorId } from "@maple/chat-platform" import { ChatConversationKey } from "@maple/domain/chat-session" import { Effect, Option, Schema } from "effect" import { describe, expect, it } from "vitest" -import { ConnectorRelay, type ConnectorRelayRuntime } from "./ConnectorRelay.ts" +import { + ConnectorRelay, + type ConnectorRelayLedger, + type ConnectorRelayRuntime, + ConnectorRelayStorageError, +} from "./ConnectorRelay.ts" import { decodeRelayTurnCheckpoint, type RelayTurnCheckpoint, type SettleOutcome } from "./settle.ts" import type { RelayHost } from "./run.ts" @@ -63,37 +68,54 @@ const conversation = (channelId: string, key = channelId): ChatConversation => ( opened: true, }) -/** One object's storage, and the alarm it arms — enough of the platform to drive the class. */ +/** + * One object's storage and its pending jobs — enough of the platform to drive the class. The + * ledger keeps checkpoints in the same store, as the isolate's does; a test runs a job by calling + * `turnDue` / `keepAliveDue` itself. + */ const objectState = () => { const stored = new Map() const pending: Array> = [] - const alarms: Array = [] - /** When the alarm is due; `null` once it has fired, which the tests do by calling `alarm()`. */ - const scheduled = { at: null as number | null } + /** Pending turn jobs by checkpoint key, and how often each was scheduled. */ + const turnJobs = new Map() + const keepAlive = { scheduled: 0 } + const schedule = (key: string) => turnJobs.set(key, (turnJobs.get(key) ?? 0) + 1) + const ledger: ConnectorRelayLedger = { + record: (key, checkpoint) => + Effect.sync(() => { + stored.set(key, checkpoint) + schedule(key) + }), + forget: (key) => + Effect.sync(() => { + stored.delete(key) + turnJobs.delete(key) + }), + read: (key) => Effect.sync(() => stored.get(key)), + revisit: (key) => Effect.sync(() => void (stored.has(key) && schedule(key))), + keepAlive: Effect.sync(() => void (keepAlive.scheduled += 1)), + } return { stored, pending, - alarms, - scheduled, + turnJobs, + keepAlive, + ledger, waitUntil: (promise: Promise) => void pending.push(promise), storage: { - getAlarm: () => Promise.resolve(scheduled.at), - setAlarm: (at: number) => { - alarms.push(at) - return Promise.resolve() - }, get: (key: string) => Promise.resolve(stored.get(key) as A | undefined), - put: (key: string, value: boolean | RelayTurnCheckpoint) => { + put: (key: string, value: boolean) => { stored.set(key, value) return Promise.resolve() }, - delete: (key: string) => Promise.resolve(stored.delete(key)), - list: ({ prefix }: { prefix: string }) => - Promise.resolve(new Map([...stored].filter(([key]) => key.startsWith(prefix)))), }, } } +/** A relay over `state`, with the heavy half `load` stands in for. */ +const relayOn = (state: ReturnType, load?: () => Promise) => + new ConnectorRelay(state, {}, state.ledger, load) + /** * A deployment of relay objects, addressed the way the Worker env addresses them. * @@ -109,7 +131,7 @@ const deployment = () => { if (existing !== undefined) return existing const state = objectState() objects.set(name, state) - const relay = new ConnectorRelay(state, env) + const relay = new ConnectorRelay(state, env, state.ledger) relays.set(name, relay) return relay } @@ -170,7 +192,7 @@ describe("remembering the conversations the bot opened", () => { // answer mentions here and nothing else, so it has to reach the turn to be logged there. const broken = objectState() broken.storage.put = () => Promise.reject(new Error("storage unavailable")) - const relay = new ConnectorRelay(broken, {}) + const relay = relayOn(broken) const here = { ...message, channelId: "channel-1" } const error = yield* Effect.flip( @@ -185,7 +207,7 @@ describe("remembering the conversations the bot opened", () => { Effect.gen(function* () { // Nothing to write to, and the turn this rides on is somebody's question: the answer must // still be given, at the cost of the follow-ups after it. - const relay = new ConnectorRelay(objectState(), {}) + const relay = relayOn(objectState()) const ports = relay.relayPorts(message) yield* ports.rememberConversation(conversation("thread_7")) @@ -209,18 +231,18 @@ const checkpoint = (): RelayTurnCheckpoint => ) const TURN_KEY = "turn:org_1:bot-testchat-thread_7:a1" -describe("waking after an eviction", () => { - it("does nothing on an alarm with no turn recorded", async () => { +describe("turn jobs", () => { + it("does nothing for a job whose turn was already forgotten", async () => { const state = objectState() state.stored.set("opened:thread_7", true) - await Effect.runPromise(new ConnectorRelay(state, {}).alarm()) + await Effect.runPromise(relayOn(state).turnDue(TURN_KEY)) expect(state.pending).toEqual([]) - expect(state.alarms).toEqual([]) + expect([...state.turnJobs]).toEqual([]) expect([...state.stored]).toEqual([["opened:thread_7", true]]) }) - it("leaves a turn this activation is still relaying to it, and clears it when the turn ends", async () => { + it("only revisits a turn this activation is still relaying, and drops it and its job when the turn ends", async () => { const state = objectState() let finish = () => {} const finished = new Promise((resolve) => { @@ -237,69 +259,186 @@ describe("waking after an eviction", () => { await finished }, }) - const relay = new ConnectorRelay(state, {}, run.load) + const relay = relayOn(state, run.load) await Effect.runPromise(relay.deliver(message)) await recorded expect([...state.stored.keys()]).toEqual([TURN_KEY]) - // The keep-alive lands mid-turn: the checkpoint is this activation's own, not an orphan. - await Effect.runPromise(relay.alarm()) + // The job lands mid-turn: the checkpoint is this activation's own, not an orphan. + await Effect.runPromise(relay.turnDue(TURN_KEY)) + expect(state.turnJobs.get(TURN_KEY)).toBe(2) finish() await Promise.all(state.pending) expect(run.settled).toEqual([]) expect([...state.stored]).toEqual([]) + expect([...state.turnJobs]).toEqual([]) }) - it("settles a turn it finds recorded, and clears it once settled", async () => { + it("settles an evicted turn, and drops it and its job once settled", async () => { const state = objectState() const run = heavy({ settle: "done" }) state.stored.set(TURN_KEY, checkpoint()) - await Effect.runPromise(new ConnectorRelay(state, {}, run.load).alarm()) + await Effect.runPromise(relayOn(state, run.load).turnDue(TURN_KEY)) - await Promise.all(state.pending) - // Kept resident while it works, like any turn this object relays. - expect(state.alarms).toHaveLength(1) expect(run.settled).toEqual([checkpoint()]) expect([...state.stored]).toEqual([]) + expect([...state.turnJobs]).toEqual([]) }) it("keeps a turn the session is still running, and comes back for it", async () => { const state = objectState() const run = heavy({ settle: "pending" }) state.stored.set(TURN_KEY, checkpoint()) - const relay = new ConnectorRelay(state, {}, run.load) + const relay = relayOn(state, run.load) - await Effect.runPromise(relay.alarm()) - await Promise.all(state.pending) - await Effect.runPromise(relay.alarm()) - await Promise.all(state.pending) + await Effect.runPromise(relay.turnDue(TURN_KEY)) + await Effect.runPromise(relay.turnDue(TURN_KEY)) expect(run.settled).toEqual([checkpoint(), checkpoint()]) expect([...state.stored.keys()]).toEqual([TURN_KEY]) - expect(state.alarms).toHaveLength(2) + expect(state.turnJobs.get(TURN_KEY)).toBe(2) + }) + + it("drops a turn checkpoint it can no longer read", async () => { + // One an older build wrote: dropped rather than thrown on. + const state = objectState() + state.stored.set(TURN_KEY, { sessionId: "org_1:bot-testchat-thread_7" }) + await Effect.runPromise(relayOn(state).turnDue(TURN_KEY)) + + expect([...state.stored]).toEqual([]) }) +}) - it("does not push back an alarm that is already due, however busy the conversation", async () => { - // Events under 30s apart would otherwise postpone the alarm — and any settle — indefinitely. +describe("the keep-alive", () => { + it("is not pushed back, however busy the conversation", async () => { + // Events under 30s apart would otherwise postpone it indefinitely. const state = objectState() - state.scheduled.at = 1 - const relay = new ConnectorRelay(state, {}, heavy({ event: () => Promise.resolve() }).load) + const relay = relayOn(state, heavy({ event: () => Promise.resolve() }).load) await Effect.runPromise(relay.deliver(message)) await Effect.runPromise(relay.deliver(message)) await Promise.all(state.pending) - expect(state.alarms).toEqual([]) + expect(state.keepAlive.scheduled).toBe(1) }) - it("drops a turn checkpoint it can no longer read", async () => { - // One an older build wrote: dropped rather than thrown on. + it("re-arms while a turn runs, lapses once none does, and arms again for the next event", async () => { const state = objectState() - state.stored.set(TURN_KEY, { sessionId: "org_1:bot-testchat-thread_7" }) - await Effect.runPromise(new ConnectorRelay(state, {}).alarm()) + let finish = () => {} + const finished = new Promise((resolve) => { + finish = resolve + }) + const relay = relayOn(state, heavy({ event: () => finished }).load) + + await Effect.runPromise(relay.deliver(message)) + await Effect.runPromise(relay.keepAliveDue()) + expect(state.keepAlive.scheduled).toBe(2) + + finish() + await Promise.all(state.pending) + await Effect.runPromise(relay.keepAliveDue()) + expect(state.keepAlive.scheduled).toBe(2) + + await Effect.runPromise(relay.deliver(message)) + expect(state.keepAlive.scheduled).toBe(3) + }) + + it("arms again on the next event after a schedule that failed", async () => { + const state = objectState() + const relay = relayOn(state, heavy({ event: () => Promise.resolve() }).load) + const working = state.ledger.keepAlive + Object.assign(state.ledger, { + keepAlive: Effect.fail( + new ConnectorRelayStorageError({ operation: "keepAlive", message: "down", cause: undefined }), + ), + }) + + await Effect.runPromise(relay.deliver(message)) + await Promise.all(state.pending) + Object.assign(state.ledger, { keepAlive: working }) + await Effect.runPromise(relay.deliver(message)) + await Promise.all(state.pending) + + expect(state.keepAlive.scheduled).toBe(1) + }) +}) + +describe("a job racing the end of its turn", () => { + it("does not settle a turn that finished after the job read its checkpoint", async () => { + const state = objectState() + let finish = () => {} + const finished = new Promise((resolve) => { + finish = resolve + }) + let recordedTurn = () => {} + const recorded = new Promise((resolve) => { + recordedTurn = resolve + }) + const run = heavy({ + event: async (host) => { + await Effect.runPromise(host.recordTurn(checkpoint())) + recordedTurn() + await finished + }, + }) + const relay = relayOn(state, run.load) + await Effect.runPromise(relay.deliver(message)) + await recorded + + // The job's read lands, then the turn ends before it decides. + const read = state.ledger.read + Object.assign(state.ledger, { + read: (key: string) => + read(key).pipe( + Effect.tap(() => + Effect.promise(async () => { + finish() + await Promise.all(state.pending) + }), + ), + ), + }) + await Effect.runPromise(relay.turnDue(TURN_KEY)) + expect(run.settled).toEqual([]) + expect([...state.stored]).toEqual([]) + expect([...state.turnJobs]).toEqual([]) + }) + + it("settles a turn whose forget failed, rather than revisiting it forever", async () => { + const state = objectState() + let recordedTurn = () => {} + const recorded = new Promise((resolve) => { + recordedTurn = resolve + }) + const run = heavy({ + event: async (host) => { + await Effect.runPromise(host.recordTurn(checkpoint())) + recordedTurn() + }, + settle: "done", + }) + const forget = state.ledger.forget + Object.assign(state.ledger, { + forget: () => + Effect.fail( + new ConnectorRelayStorageError({ + operation: "forget", + message: "down", + cause: undefined, + }), + ), + }) + const relay = relayOn(state, run.load) + await Effect.runPromise(relay.deliver(message)) + await recorded await Promise.all(state.pending) + Object.assign(state.ledger, { forget }) + + await Effect.runPromise(relay.turnDue(TURN_KEY)) + + expect(run.settled).toEqual([checkpoint()]) expect([...state.stored]).toEqual([]) }) }) diff --git a/apps/chat-bot/src/relay/ConnectorRelay.ts b/apps/chat-bot/src/relay/ConnectorRelay.ts index 096e25119..2d0ea93cc 100644 --- a/apps/chat-bot/src/relay/ConnectorRelay.ts +++ b/apps/chat-bot/src/relay/ConnectorRelay.ts @@ -24,25 +24,42 @@ import { BoundChatSessions, type ChatSessionNamespace } from "@maple/backend/platform/chat-sessions" import type { ChatConversation, InboundEvent } from "@maple/chat-platform" import { summarizeCause } from "@maple/backend/platform/describe-cause" +import * as Alchemy from "alchemy" import * as Cloudflare from "alchemy/Cloudflare" -import { Effect, Schema } from "effect" +import { Context, Effect, Schema } from "effect" import { ConversationNotRecorded } from "./conversation.ts" import type { RelayTurnCheckpoint, SettleOutcome } from "./settle.ts" import { connectorConversationRelayName, connectorRelayByName } from "./stub.ts" -/** What this object reads off its Durable Object state. */ +/** What this object reads off its Durable Object state. Turn checkpoints go through the ledger. */ interface ConnectorRelayState { readonly storage: { - getAlarm(): Promise - setAlarm(scheduledTime: number): Promise get(key: string): Promise - put(key: string, value: boolean | RelayTurnCheckpoint): Promise - delete(key: string): Promise - list(options: { prefix: string }): Promise> + put(key: string, value: boolean): Promise } waitUntil(promise: Promise): void } +/** + * Turn checkpoints and the durable jobs that wake this object, each write committed together with + * the job it implies. Alchemy callbacks in the isolate (`durableLedger`); a test stands in for them. + */ +export interface ConnectorRelayLedger { + /** Store the checkpoint and schedule its turn job, in one transaction. */ + readonly record: ( + key: string, + checkpoint: RelayTurnCheckpoint, + ) => Effect.Effect + /** Drop the checkpoint and cancel its turn job, in one transaction. */ + readonly forget: (key: string) => Effect.Effect + /** The stored checkpoint, or `undefined` once the turn was forgotten. */ + readonly read: (key: string) => Effect.Effect + /** Run the turn job again after {@link KEEP_ALIVE_MS}, if the checkpoint is still stored. */ + readonly revisit: (key: string) => Effect.Effect + /** Run the keep-alive job after {@link KEEP_ALIVE_MS}. */ + readonly keepAlive: Effect.Effect +} + /** * Where the one durable fact lives, under the conversation it is about. * @@ -53,7 +70,7 @@ interface ConnectorRelayState { const openedKey = (conversationKey: string): string => `opened:${conversationKey}` /** One checkpoint per turn: a channel whose threads are conversations can relay several at once. */ -const TURN_PREFIX = "turn:" +export const TURN_PREFIX = "turn:" const turnKey = (checkpoint: RelayTurnCheckpoint): string => `${TURN_PREFIX}${checkpoint.sessionId}:${checkpoint.turnMessageId}` @@ -74,14 +91,14 @@ export interface ConnectorRelayPorts { } /** - * How often a relaying object re-arms its alarm. + * How often a relaying object wakes itself. * * An outbound fetch never keeps a Durable Object alive and an object with no incoming event is * evicted inside a couple of minutes, so a turn nobody is streaming FROM this object would be cut - * off mid-answer. The alarm is that incoming event, and it is the same 30 seconds the chat session + * off mid-answer. A due job is that incoming event, and it is the same 30 seconds the chat session * itself uses for the same reason. */ -const KEEP_ALIVE_MS = 30 * 1000 +export const KEEP_ALIVE_MS = 30 * 1000 /** The heavy half, loaded on first use (see `run`); a parameter so a test can stand in for it. */ export type ConnectorRelayRuntime = Pick @@ -131,15 +148,18 @@ const runtimeCall = (operation: string, call: () => Promise) => }) export class ConnectorRelay { - /** How many events this activation is still working on. Zero means the alarm may stop. */ + /** How many events this activation is still working on. Zero lets the keep-alive lapse. */ private live = 0 - /** Checkpoints this activation is relaying; one in storage but not here is an evicted turn's. */ + /** Whether a keep-alive job is pending; set once per lapse so a busy channel never postpones it. */ + private keepingAlive = false + /** Checkpoints this activation is relaying; one stored but not here is an evicted turn's. */ private readonly relaying = new Set() private unlinkedNoticeAt: number | undefined constructor( private readonly ctx: ConnectorRelayState, private readonly env: Record, + private readonly ledger: ConnectorRelayLedger, private readonly runtime: () => Promise = loadRuntime, /** maple-ai's `ChatSession` namespace, which the Worker binds cross-script. */ private readonly chatSessions?: ChatSessionNamespace, @@ -155,32 +175,37 @@ export class ConnectorRelay { deliver(event: InboundEvent): Effect.Effect { return Effect.sync(() => { this.live += 1 - this.armKeepAlive() + if (!this.keepingAlive) { + this.keepingAlive = true + this.ctx.waitUntil(Effect.runPromise(this.armKeepAlive)) + } this.ctx.waitUntil(Effect.runPromise(this.run(event))) }) } + /** The keep-alive job: re-armed while this activation has work, left to lapse once it has none. */ + keepAliveDue(): Effect.Effect { + return Effect.suspend(() => { + if (this.live > 0) return this.armKeepAlive + this.keepingAlive = false + return Effect.void + }) + } + /** - * The keep-alive, and the settle of turns an evicted activation left behind. It keeps firing while - * any checkpoint is left, which is how a turn the session is still running gets settled later. + * A turn's job. A turn this activation relays is only looked at again later; one it does not is + * an evicted activation's, and is settled here. A forgotten turn's job has nothing left to do. */ - alarm(): Effect.Effect { - return storageCall("list", () => this.ctx.storage.list({ prefix: TURN_PREFIX })).pipe( - Effect.matchEffect({ - // Unlistable this time: re-armed, so the checkpoints are looked at again on the next tick. - onFailure: (error) => - logStorageFailure(error).pipe(Effect.andThen(Effect.sync(() => this.armKeepAlive()))), - onSuccess: (recorded) => - Effect.sync(() => { - for (const [key, checkpoint] of recorded) { - if (this.relaying.has(key)) continue - this.live += 1 - this.relaying.add(key) - this.ctx.waitUntil(Effect.runPromise(this.settle(key, checkpoint))) - } - if (this.live > 0) this.armKeepAlive() - }), - }), + turnDue(key: string): Effect.Effect { + return this.ledger.read(key).pipe( + Effect.flatMap((checkpoint) => + checkpoint === undefined + ? Effect.void + : this.relaying.has(key) + ? this.ledger.revisit(key) + : this.settle(key, checkpoint), + ), + Effect.catchTag(STORAGE_ERROR, logStorageFailure), ) } @@ -253,23 +278,14 @@ export class ConnectorRelay { ) } - /** A pending alarm is kept, not pushed back: a busy conversation would postpone it forever. */ - private armKeepAlive(): void { - this.ctx.waitUntil( - Effect.runPromise( - storageCall("getAlarm", () => this.ctx.storage.getAlarm()).pipe( - Effect.flatMap((pending) => - pending === null - ? storageCall("setAlarm", () => - this.ctx.storage.setAlarm(Date.now() + KEEP_ALIVE_MS), - ) - : Effect.void, - ), - Effect.catchTag(STORAGE_ERROR, logStorageFailure), - ), + /** A failed schedule leaves no job to clear the flag, so it falls here and the next event re-arms. */ + private readonly armKeepAlive: Effect.Effect = Effect.suspend(() => this.ledger.keepAlive).pipe( + Effect.catchTag(STORAGE_ERROR, (error) => + logStorageFailure(error).pipe( + Effect.andThen(Effect.sync(() => void (this.keepingAlive = false))), ), - ) - } + ), + ) /** * Everything below this line is behind a dynamic import: it reaches the connector registry, the @@ -306,7 +322,7 @@ export class ConnectorRelay { } /** Kept only while the session is still running the turn; cleared however else it ends. */ - private settle(key: string, checkpoint: unknown): Effect.Effect { + private settle(key: string, checkpoint: unknown): Effect.Effect { return runtimeCall("settle", async () => { const { settleInboundTurn } = await this.runtime() return settleInboundTurn( @@ -326,30 +342,37 @@ export class ConnectorRelay { ), ), Effect.flatMap((outcome) => - outcome === "done" ? this.forgetTurn(key) : Effect.sync(() => void this.relaying.delete(key)), + outcome === "done" ? this.ledger.forget(key) : this.ledger.revisit(key), ), - Effect.ensuring(Effect.sync(() => (this.live -= 1))), + // A settle that re-records the turn marks it live; it is not, once the settle returns. + Effect.ensuring(Effect.sync(() => void this.relaying.delete(key))), ) } - /** Marked before the write, so an alarm meanwhile does not take a live turn for an evicted one. */ + /** Marked before the write, so a job meanwhile does not take a live turn for an evicted one. */ private recordTurn(checkpoint: RelayTurnCheckpoint): Effect.Effect { const key = turnKey(checkpoint) return Effect.suspend(() => { this.relaying.add(key) - return storageCall("put", () => this.ctx.storage.put(key, checkpoint)) + return this.ledger.record(key, checkpoint) }).pipe(Effect.catchTag(STORAGE_ERROR, logStorageFailure)) } /** - * Left in `relaying`: an alarm that listed the key before the delete must not settle a turn - * that already finished. Keys are per turn, so none is ever reused. + * The checkpoint and its job go together. The key stays in `relaying`: a job that read the + * checkpoint before this delete must still see a live turn. A forget that failed rolled both + * back, so the key leaves `relaying` and the surviving job settles the turn. Keys are per turn. */ private forgetTurn(key: string): Effect.Effect { - return storageCall("delete", () => this.ctx.storage.delete(key)).pipe( - Effect.asVoid, - Effect.catchTag(STORAGE_ERROR, logStorageFailure), - ) + return this.ledger + .forget(key) + .pipe( + Effect.catchTag(STORAGE_ERROR, (error) => + logStorageFailure(error).pipe( + Effect.andThen(Effect.sync(() => void this.relaying.delete(key))), + ), + ), + ) } /** @@ -374,24 +397,99 @@ export class ConnectorRelay { export interface ConnectorRelayApi { readonly deliver: (event: InboundEvent) => Effect.Effect - readonly alarm: () => Effect.Effect readonly remember: (conversationKey: string) => Effect.Effect } +const storageFailed = (operation: string) => (cause: unknown) => + new ConnectorRelayStorageError({ + operation, + message: `The relay object's storage ${operation} failed`, + cause, + }) + +/** + * The ledger on alchemy's callbacks: a checkpoint and its job commit in one storage transaction, and + * the job store keeps the native alarm on the earliest one. `runtime` is the instance's own context, + * captured at activation, because the turns run detached on `waitUntil`. + */ +const durableLedger = ( + state: Cloudflare.DurableObjectState["Service"], + turnJob: Alchemy.Callback, + keepAliveJob: Alchemy.Callback, + runtime: Context.Context, +): ConnectorRelayLedger => { + const inContext = (operation: string, effect: Effect.Effect) => + effect.pipe(Effect.mapError(storageFailed(operation)), Effect.provideContext(runtime)) + const again = { after: KEEP_ALIVE_MS } + return { + record: (key, checkpoint) => + inContext( + "record", + state.storage.transaction( + state.storage + .put(key, checkpoint) + .pipe(Effect.andThen(turnJob.schedule(key, { ...again, payload: key }))), + ), + ), + forget: (key) => + inContext( + "forget", + state.storage.transaction( + state.storage.delete(key).pipe(Effect.andThen(turnJob.cancel(key))), + ), + ), + read: (key) => inContext("get", state.storage.get(key)), + revisit: (key) => + inContext( + "revisit", + state.storage.transaction( + state.storage + .get(key) + .pipe( + Effect.flatMap((stored) => + stored === undefined + ? Effect.void + : turnJob.schedule(key, { ...again, payload: key }), + ), + ), + ), + ), + keepAlive: inContext("keepAlive", keepAliveJob.schedule("keep-alive", { ...again, payload: null })), + } +} + /** * One activation, in alchemy's two phases: the outer Effect resolves the state, env and the chat * session namespace (the Worker provides it, `worker.ts`) — it also runs at plan time against a - * mock state, so it must not touch storage — and the inner one returns the object's methods as - * Effects, which alchemy's bridge runs per RPC call. + * mock state, so it must not touch storage — and the inner one registers the jobs and returns the + * object's methods as Effects, which alchemy's bridge runs per RPC call. */ export const activateConnectorRelay = Effect.map( Effect.all([Cloudflare.DurableObjectState, Cloudflare.WorkerEnvironment, BoundChatSessions]), ([state, env, chatSessions]) => - Effect.sync(() => { - const relay = new ConnectorRelay(state.raw, env, loadRuntime, chatSessions) + Effect.gen(function* () { + // The handlers reach `relay` only when a job runs, after it is built below. + const turnJob = yield* Alchemy.makeCallback("relay-turn", (key: string) => + Effect.suspend(() => relay.turnDue(key)), + ) + const keepAliveJob = yield* Alchemy.makeCallback("relay-keep-alive", (_: null) => + Effect.suspend(() => relay.keepAliveDue()), + ) + const runtime = yield* Effect.context() + const ledger = durableLedger(state, turnJob, keepAliveJob, runtime) + const relay = new ConnectorRelay(state.raw, env, ledger, loadRuntime, chatSessions) + // Every checkpoint stored at activation is a dead activation's. Scheduling each one covers + // the ones written before checkpoints had jobs; for the rest it only brings the job forward. + const stored = yield* state.storage.list({ prefix: TURN_PREFIX }) + yield* Effect.forEach(stored.keys(), (key) => turnJob.schedule(key, { after: 0, payload: key }), { + discard: true, + }).pipe( + Effect.catchTag("CallbackError", (error) => + logStorageFailure(storageFailed("schedule")(error)), + ), + ) return { deliver: (event) => relay.deliver(event), - alarm: () => relay.alarm(), remember: (conversationKey) => relay.remember(conversationKey), } satisfies ConnectorRelayApi }), diff --git a/apps/chat-bot/src/relay/settle.ts b/apps/chat-bot/src/relay/settle.ts index a8e3ac727..bc0905440 100644 --- a/apps/chat-bot/src/relay/settle.ts +++ b/apps/chat-bot/src/relay/settle.ts @@ -38,7 +38,7 @@ export const decodeRelayTurnCheckpoint = Schema.decodeUnknownOption(RelayTurnChe /** The session's `TURN_STALE_MS` (25 min, when it expires a turn nobody ended) plus a margin. */ const CHECKPOINT_TTL_MS = 30 * 60 * 1000 -/** `"pending"`: the session is still running the turn, so the checkpoint waits for the next alarm. */ +/** `"pending"`: the session is still running the turn, so the checkpoint waits for its next job. */ export type SettleOutcome = "pending" | "done" /** Anything that goes wrong is logged once, the cause summarized, and the checkpoint dropped. */ diff --git a/bun.lock b/bun.lock index b5226deb4..c3b140e3a 100644 --- a/bun.lock +++ b/bun.lock @@ -9,6 +9,7 @@ }, "devDependencies": { "@cloudflare/workers-types": "catalog:alchemy", + "@effect/platform-node": "catalog:effect", "@effect/tsgo": "0.48.0", "@effect/vitest": "catalog:effect", "@maple/alchemy-portless": "workspace:*", @@ -949,10 +950,10 @@ }, }, "patchedDependencies": { - "alchemy@2.0.0-beta.80": "patches/alchemy@2.0.0-beta.80.patch", + "@effect/ai-openrouter@4.0.0": "patches/@effect%2Fai-openrouter@4.0.0.patch", "@effect/vitest@4.0.0": "patches/@effect%2Fvitest@4.0.0.patch", + "alchemy@2.0.0-beta.81": "patches/alchemy@2.0.0-beta.81.patch", "effect@4.0.0": "patches/effect@4.0.0.patch", - "@effect/ai-openrouter@4.0.0": "patches/@effect%2Fai-openrouter@4.0.0.patch", }, "overrides": { "@effect/sql-d1": "4.0.0", @@ -968,7 +969,7 @@ "catalogs": { "alchemy": { "@cloudflare/workers-types": "4.20260603.1", - "alchemy": "2.0.0-beta.80", + "alchemy": "2.0.0-beta.81", }, "charts": { "@tanstack/charts": "0.18.0", @@ -980,6 +981,7 @@ "@effect/atom-react": "4.0.0", "@effect/language-service": "^0.87.3", "@effect/platform-bun": "4.0.0", + "@effect/platform-node": "4.0.0", "@effect/vitest": "4.0.0", "effect": "4.0.0", }, @@ -1020,11 +1022,11 @@ "@alcalzone/ansi-tokenize": ["@alcalzone/ansi-tokenize@0.2.5", "", { "dependencies": { "ansi-styles": "^6.2.1", "is-fullwidth-code-point": "^5.0.0" } }, "sha512-3NX/MpTdroi0aKz134A6RC2Gb2iXVECN4QaAXnvCIxxIm3C3AVB1mkUe8NaaiyvOpDfsrqWhYtj+Q6a62RrTsw=="], - "@alchemy.run/cloudflare-runtime": ["@alchemy.run/cloudflare-runtime@2.0.0-beta.80", "", { "dependencies": { "@alchemy.run/node-utils": "2.0.0-beta.80", "@cloudflare/unenv-preset": "^2.16.1", "@puppeteer/browsers": "^3.2.0", "capnp-es": "^0.0.16", "magic-string": "^0.30.21", "mime": "^4.0.7", "sharp": "^0.35.3", "unenv": "^2.0.0-rc.24", "workerd": "1.20260918.1", "yauzl": "^3.4.0" }, "peerDependencies": { "@distilled.cloud/cloudflare": "1.0.0-rc.13", "@effect/platform-bun": "^4.0.0", "@effect/platform-node": "^4.0.0", "effect": "^4.0.0", "rolldown": "1.2.5", "vite": "^7.0.0 || ^8.0.0" }, "optionalPeers": ["@effect/platform-bun", "@effect/platform-node", "rolldown", "vite"] }, "sha512-vrMDVAQ+LlzFGkgfIbi/65iwskdSxxWYhJX03sFL72mvkGxvb8RQSKrbxLGHhMz0QSZh796dOff1z9zOwxsIbA=="], + "@alchemy.run/cloudflare-runtime": ["@alchemy.run/cloudflare-runtime@2.0.0-beta.81", "", { "dependencies": { "@alchemy.run/node-utils": "2.0.0-beta.81", "@cloudflare/unenv-preset": "^2.16.1", "@puppeteer/browsers": "^3.2.0", "capnp-es": "^0.0.16", "magic-string": "^0.30.21", "mime": "^4.0.7", "sharp": "^0.35.3", "unenv": "^2.0.0-rc.24", "workerd": "1.20260918.1", "yauzl": "^3.4.0" }, "peerDependencies": { "@distilled.cloud/cloudflare": "1.0.0-rc.13", "@effect/platform-bun": "^4.0.0", "@effect/platform-node": "^4.0.0", "effect": "^4.0.0", "rolldown": "1.2.5", "vite": "^7.0.0 || ^8.0.0" }, "optionalPeers": ["@effect/platform-bun", "@effect/platform-node", "rolldown", "vite"] }, "sha512-gGmshJZSFbeyK94UMRgRcrE0prwJrp3H3kQ8aZmGOuWzwJMMxaoCuf4Nulf8nGsvHcjMEAfFi45AXRRQBgbgQQ=="], - "@alchemy.run/floci": ["@alchemy.run/floci@2.0.0-beta.80", "", { "peerDependencies": { "effect": "^4.0.0" } }, "sha512-gtUXFitVY9pSQQYCDtSVyjfd5YOKY3Wleg3WCG5yWno7/l7xBo4kfOMbM63tzWnmka9LewVSITbMfwnC4D/73A=="], + "@alchemy.run/floci": ["@alchemy.run/floci@2.0.0-beta.81", "", { "peerDependencies": { "effect": "^4.0.0" } }, "sha512-0mYZqFrQLlRwDUsoJlI/ArnEaAaFid2ZPsJWvUyC5llAScWxlST7pdjQd1r/t1YIFFkhOV9JQO1wa7FI1nY6EQ=="], - "@alchemy.run/node-utils": ["@alchemy.run/node-utils@2.0.0-beta.80", "", { "dependencies": { "chokidar": "^5.0.0", "rolldown": "1.2.5" } }, "sha512-uOhVDQZTjwh1AU4mrk0wB5rQaziTwSrVMfqwFP8wALfU//nveu4goOWZimHtqytXiybENh+iEbALH25EFAkc7A=="], + "@alchemy.run/node-utils": ["@alchemy.run/node-utils@2.0.0-beta.81", "", { "dependencies": { "chokidar": "^5.0.0", "rolldown": "1.2.5" } }, "sha512-AcxxGMFbbdiKww7OUbBd/fsSDyRJjoICAIqGwgkJxrIQkkaNt6+9SObrznmnOx70OonA2GH87RGnEaTtORHffw=="], "@alchemy.run/sigil": ["@alchemy.run/sigil@0.1.0-alpha.1", "", { "peerDependencies": { "@types/react": ">=19.2.0", "react-devtools-core": ">=6.1.2" }, "optionalPeers": ["@types/react", "react-devtools-core"] }, "sha512-RtHELG2w+GluVewRsuKm66idtPzC+4STptMr/61X1Vy5uBkDzcCzjzV8G6jUa4iBv6Il0hHV7UL6wvv//9nD/Q=="], @@ -2504,7 +2506,7 @@ "ajv-keywords": ["ajv-keywords@5.1.0", "", { "dependencies": { "fast-deep-equal": "^3.1.3" }, "peerDependencies": { "ajv": "^8.8.2" } }, "sha512-YCS/JNFAUyr5vAuhk1DWm1CBxRHW9LbJ2ozWeemrIqpbsqKjHVxYPyi5GC0rjZIT5JxJ3virVTS8wk4i/Z+krw=="], - "alchemy": ["alchemy@2.0.0-beta.80", "", { "dependencies": { "@alchemy.run/cloudflare-runtime": "2.0.0-beta.80", "@alchemy.run/floci": "2.0.0-beta.80", "@alchemy.run/node-utils": "2.0.0-beta.80", "@alchemy.run/sigil": "0.1.0-alpha.1", "@distilled.cloud/acme": "1.0.0-rc.13", "@distilled.cloud/aws": "1.0.0-rc.13", "@distilled.cloud/axiom": "1.0.0-rc.13", "@distilled.cloud/cloudflare": "1.0.0-rc.13", "@distilled.cloud/core": "1.0.0-rc.13", "@distilled.cloud/doppler": "1.0.0-rc.13", "@distilled.cloud/fly-io": "1.0.0-rc.13", "@distilled.cloud/gcp": "1.0.0-rc.13", "@distilled.cloud/hetzner": "1.0.0-rc.13", "@distilled.cloud/infisical": "1.0.0-rc.13", "@distilled.cloud/neon": "1.0.0-rc.13", "@distilled.cloud/planetscale": "1.0.0-rc.13", "@distilled.cloud/prisma": "1.0.0-rc.13", "@distilled.cloud/railway": "1.0.0-rc.13", "@distilled.cloud/stripe": "1.0.0-rc.13", "@distilled.cloud/zerossl": "1.0.0-rc.13", "@effect/sql-d1": "^4.0.0", "@effect/sql-sqlite-do": "^4.0.0", "@libsql/client": "^0.17.0", "@neon/functions": "0.11.0", "@octokit/rest": "^22.0.1", "@octokit/webhooks": "^14.2.0", "@prisma/dev": "^0.20.0", "@types/aws-lambda": "^8.10.152", "capnweb": "^0.12.0", "fflate": "^0.8.3", "libsodium-wrappers": "^0.8.3", "pathe": "^2.0.3", "picomatch": "^4.0.4", "rolldown": "1.2.5", "tinyglobby": "^0.2.17", "yaml": "^2.9.0" }, "peerDependencies": { "@alchemy.run/frontend-frameworks": "2.0.0-beta.80", "@aws/durable-execution-sdk-js": "^2.1.0", "@effect/platform-bun": "^4.0.0", "@effect/platform-node": "^4.0.0", "@effect/sql-mysql2": "^4.0.0", "@effect/sql-pg": "^4.0.0", "@effect/vitest": "^4.0.0", "@prisma/orm-postgres": "8.0.0-rc.11", "@vercel/nft": "^1.10.2", "drizzle-kit": "1.0.0-rc.5-ab785fc", "drizzle-orm": "1.0.0-rc.5-ab785fc", "effect": "^4.0.0", "mongodb": "^6.10.0", "mysql2": "^3.24.2", "pg": "^8.22.0", "prisma": "8.0.0-rc.15", "vite": "^8.0.7", "ws": "^8.21.0" }, "optionalPeers": ["@alchemy.run/frontend-frameworks", "@aws/durable-execution-sdk-js", "@effect/platform-bun", "@effect/platform-node", "@effect/sql-mysql2", "@effect/sql-pg", "@effect/vitest", "@prisma/orm-postgres", "@vercel/nft", "drizzle-kit", "drizzle-orm", "mongodb", "mysql2", "pg", "prisma", "vite", "ws"], "bin": { "alchemy": "./bin/cli.js" } }, "sha512-2y7mswypPiApxuEdp8NGMfzCGB5+9W6OmQAuePyZ2xBkgoCpEqtKAnASuvKGUZPSYNESENbamXDuI5428ewu9w=="], + "alchemy": ["alchemy@2.0.0-beta.81", "", { "dependencies": { "@alchemy.run/cloudflare-runtime": "2.0.0-beta.81", "@alchemy.run/floci": "2.0.0-beta.81", "@alchemy.run/node-utils": "2.0.0-beta.81", "@alchemy.run/sigil": "0.1.0-alpha.1", "@distilled.cloud/acme": "1.0.0-rc.13", "@distilled.cloud/aws": "1.0.0-rc.13", "@distilled.cloud/axiom": "1.0.0-rc.13", "@distilled.cloud/cloudflare": "1.0.0-rc.13", "@distilled.cloud/core": "1.0.0-rc.13", "@distilled.cloud/doppler": "1.0.0-rc.13", "@distilled.cloud/fly-io": "1.0.0-rc.13", "@distilled.cloud/gcp": "1.0.0-rc.13", "@distilled.cloud/hetzner": "1.0.0-rc.13", "@distilled.cloud/infisical": "1.0.0-rc.13", "@distilled.cloud/neon": "1.0.0-rc.13", "@distilled.cloud/planetscale": "1.0.0-rc.13", "@distilled.cloud/prisma": "1.0.0-rc.13", "@distilled.cloud/railway": "1.0.0-rc.13", "@distilled.cloud/stripe": "1.0.0-rc.13", "@distilled.cloud/zerossl": "1.0.0-rc.13", "@effect/sql-d1": "^4.0.0", "@effect/sql-sqlite-do": "^4.0.0", "@libsql/client": "^0.17.0", "@neon/functions": "0.11.0", "@octokit/rest": "^22.0.1", "@octokit/webhooks": "^14.2.0", "@prisma/dev": "^0.20.0", "@types/aws-lambda": "^8.10.152", "capnweb": "^0.12.0", "fflate": "^0.8.3", "libsodium-wrappers": "^0.8.3", "pathe": "^2.0.3", "picomatch": "^4.0.4", "rolldown": "1.2.5", "tinyglobby": "^0.2.17", "yaml": "^2.9.0" }, "peerDependencies": { "@alchemy.run/frontend-frameworks": "2.0.0-beta.81", "@aws/durable-execution-sdk-js": "^2.1.0", "@effect/platform-bun": "^4.0.0", "@effect/platform-node": "^4.0.0", "@effect/sql-mysql2": "^4.0.0", "@effect/sql-pg": "^4.0.0", "@effect/vitest": "^4.0.0", "@prisma/orm-postgres": "8.0.0-rc.11", "@vercel/nft": "^1.10.2", "drizzle-kit": "1.0.0-rc.5-ab785fc", "drizzle-orm": "1.0.0-rc.5-ab785fc", "effect": "^4.0.0", "mongodb": "^6.10.0", "mysql2": "^3.24.2", "pg": "^8.22.0", "prisma": "8.0.0-rc.15", "vite": "^8.0.7", "ws": "^8.21.0" }, "optionalPeers": ["@alchemy.run/frontend-frameworks", "@aws/durable-execution-sdk-js", "@effect/platform-bun", "@effect/platform-node", "@effect/sql-mysql2", "@effect/sql-pg", "@effect/vitest", "@prisma/orm-postgres", "@vercel/nft", "drizzle-kit", "drizzle-orm", "mongodb", "mysql2", "pg", "prisma", "vite", "ws"], "bin": { "alchemy": "./bin/cli.js" } }, "sha512-AxEUQSYSln9NgZMuYnSz9Q7ERR8BrKrGpx9vOwJtWity7Ze8/7vI4k62u2EAoJRpBbaOTnI2Bp67Au1y9J7FTA=="], "am-i-vibing": ["am-i-vibing@0.4.0", "", { "dependencies": { "process-ancestry": "^0.1.0" }, "bin": { "am-i-vibing": "dist/cli.mjs" } }, "sha512-MxT4XZL7pzLHpuvhDKdMaQHMGGkJDLluKBLsbstn+8wv9sWcFT6h+0ve9qkml95amVTZtZV83gQe2hY+ojgHLg=="], diff --git a/package.json b/package.json index 760a867b5..f4b872fcb 100644 --- a/package.json +++ b/package.json @@ -49,6 +49,7 @@ }, "devDependencies": { "@cloudflare/workers-types": "catalog:alchemy", + "@effect/platform-node": "catalog:effect", "@effect/tsgo": "0.48.0", "@effect/vitest": "catalog:effect", "@maple/alchemy-portless": "workspace:*", @@ -81,7 +82,7 @@ "catalogs": { "alchemy": { "@cloudflare/workers-types": "4.20260603.1", - "alchemy": "2.0.0-beta.80" + "alchemy": "2.0.0-beta.81" }, "charts": { "@tanstack/charts": "0.18.0", @@ -93,6 +94,7 @@ "@effect/atom-react": "4.0.0", "@effect/language-service": "^0.87.3", "@effect/platform-bun": "4.0.0", + "@effect/platform-node": "4.0.0", "@effect/vitest": "4.0.0", "effect": "4.0.0" }, @@ -126,6 +128,6 @@ "@effect/ai-openrouter@4.0.0": "patches/@effect%2Fai-openrouter@4.0.0.patch", "@effect/vitest@4.0.0": "patches/@effect%2Fvitest@4.0.0.patch", "effect@4.0.0": "patches/effect@4.0.0.patch", - "alchemy@2.0.0-beta.80": "patches/alchemy@2.0.0-beta.80.patch" + "alchemy@2.0.0-beta.81": "patches/alchemy@2.0.0-beta.81.patch" } } diff --git a/patches/alchemy@2.0.0-beta.80.patch b/patches/alchemy@2.0.0-beta.81.patch similarity index 83% rename from patches/alchemy@2.0.0-beta.80.patch rename to patches/alchemy@2.0.0-beta.81.patch index 60ab783df..a3cf9ae5d 100644 --- a/patches/alchemy@2.0.0-beta.80.patch +++ b/patches/alchemy@2.0.0-beta.81.patch @@ -1,12 +1,12 @@ diff --git a/lib/AWS/AutoScaling/AutoScalingGroup.js b/lib/AWS/AutoScaling/AutoScalingGroup.js -index f925594..827bb86 100644 +index a0fcddc..bea6b3d 100644 --- a/lib/AWS/AutoScaling/AutoScalingGroup.js +++ b/lib/AWS/AutoScaling/AutoScalingGroup.js -@@ -246,7 +246,17 @@ export const AutoScalingGroupProvider = () => Provider.effect(AutoScalingGroup, +@@ -245,7 +245,17 @@ export const AutoScalingGroupProvider = () => Provider.effect(AutoScalingGroup, TerminationPolicies: news.terminationPolicies, Tags: toTags(autoScalingGroupName, desiredTags), }) -- .pipe(Effect.catch((error) => error?._tag === "AlreadyExistsFault" +- .pipe(Effect.catch((error) => error?._tag === "AlreadyExistsFault" ? Effect.void : Effect.fail(error))); + .pipe( + // maple patch: a just-created instance profile takes a few seconds to + // reach EC2, and the launch template is rejected until it does. @@ -17,11 +17,11 @@ index f925594..827bb86 100644 + Schedule.recurs(10), + Schedule.exponential("1 second"), + ]), -+ }), Effect.catch((error) => error?._tag === "AlreadyExistsFault" - ? Effect.void - : Effect.fail(error))); ++ }), Effect.catch((error) => error?._tag === "AlreadyExistsFault" ? Effect.void : Effect.fail(error))); existing = yield* describeGroup(autoScalingGroupName).pipe(Effect.filterOrFail(Boolean, () => new Error(`Auto Scaling Group '${autoScalingGroupName}' was not readable after create`)), Effect.retry({ -@@ -265,7 +275,9 @@ export const AutoScalingGroupProvider = () => Provider.effect(AutoScalingGroup, + while: () => true, + schedule: Schedule.max([Schedule.recurs(8), Schedule.exponential("250 millis")]), +@@ -259,7 +269,9 @@ export const AutoScalingGroupProvider = () => Provider.effect(AutoScalingGroup, AutoScalingGroupName: autoScalingGroupName, MinSize: news.minSize, MaxSize: news.maxSize, @@ -33,7 +33,7 @@ index f925594..827bb86 100644 VPCZoneIdentifier: news.subnetIds.join(","), HealthCheckType: healthCheckType, diff --git a/lib/AWS/ECS/CapacityProvider.js b/lib/AWS/ECS/CapacityProvider.js -index c274c91..51984f5 100644 +index 401c3eb..c060ef5 100644 --- a/lib/AWS/ECS/CapacityProvider.js +++ b/lib/AWS/ECS/CapacityProvider.js @@ -91,7 +91,7 @@ export const CapacityProviderProvider = () => Provider.effect(CapacityProvider, @@ -46,19 +46,19 @@ index c274c91..51984f5 100644 } const internalTags = yield* createInternalTags(id); diff --git a/lib/AWS/ECS/Service.js b/lib/AWS/ECS/Service.js -index b4045f4..a8da70c 100644 +index fff172b..8b857da 100644 --- a/lib/AWS/ECS/Service.js +++ b/lib/AWS/ECS/Service.js -@@ -675,6 +675,8 @@ const transformServiceProps = (id, props) => Effect.gen(function* () { +@@ -663,6 +663,8 @@ const transformServiceProps = (id, props) => Effect.gen(function* () { return next; }).pipe(Namespace.push(id)); }); +/** maple patch: whether the service's tasks get their own ENI. */ +const usesAwsvpc = (props) => (props?.networkMode ?? "awsvpc") === "awsvpc"; const composeManagedIngress = (id, props, lbProp) => Effect.gen(function* () { - const config = lbProp === true - ? {} -@@ -978,7 +980,7 @@ const composeManagedIngress = (id, props, lbProp) => Effect.gen(function* () { + const config = lbProp === true ? {} : isELBv2Listener(lbProp) ? { listener: lbProp } : lbProp; + const rules = config.rules ?? (config.listener !== undefined ? [{}] : []); +@@ -955,7 +957,7 @@ const composeManagedIngress = (id, props, lbProp) => Effect.gen(function* () { vpcId: network.vpcId, port: spec.port, protocol: spec.protocol, @@ -67,7 +67,7 @@ index b4045f4..a8da70c 100644 healthCheckPath: isNetworkTg ? wantsHttpCheck ? (health?.path ?? "/") -@@ -1844,7 +1846,9 @@ export const ServiceProvider = () => Provider.effect(Service, Effect.gen(functio +@@ -1790,7 +1792,9 @@ export const ServiceProvider = () => Provider.effect(Service, Effect.gen(functio platformVersion: news.platformVersion, deploymentConfiguration: news.deploymentConfiguration, healthCheckGracePeriodSeconds: toWireSeconds(news.healthCheckGracePeriod), @@ -79,10 +79,10 @@ index b4045f4..a8da70c 100644 placementConstraints: news.placementConstraints, placementStrategy: news.placementStrategy, diff --git a/lib/Cloudflare/Workers/DurableObject.js b/lib/Cloudflare/Workers/DurableObject.js -index 14dca618c2..e90d8b1c61 100644 +index cd388b9..b5c5f1c 100644 --- a/lib/Cloudflare/Workers/DurableObject.js +++ b/lib/Cloudflare/Workers/DurableObject.js -@@ -971,22 +971,17 @@ export const DurableObject = taggedFunction(DurableObjectScope, function (...arg +@@ -968,22 +968,17 @@ export const DurableObject = taggedFunction(DurableObjectScope, function (...arg return Effect.die(new Error(`DurableObject '${namespace}' is not a DurableObject`)); } })); @@ -113,10 +113,10 @@ index 14dca618c2..e90d8b1c61 100644 // Class-form declarations (`DurableObject()("Name", props?)`) can // carry `transferredFrom` so the host's local binding drives a diff --git a/src/AWS/AutoScaling/AutoScalingGroup.ts b/src/AWS/AutoScaling/AutoScalingGroup.ts -index 5b40c65..2eb9162 100644 +index 25dac6f..93bab9d 100644 --- a/src/AWS/AutoScaling/AutoScalingGroup.ts +++ b/src/AWS/AutoScaling/AutoScalingGroup.ts -@@ -470,6 +470,17 @@ export const AutoScalingGroupProvider = () => +@@ -453,6 +453,17 @@ export const AutoScalingGroupProvider = () => Tags: toTags(autoScalingGroupName, desiredTags), } as any) .pipe( @@ -132,9 +132,9 @@ index 5b40c65..2eb9162 100644 + ]), + }), Effect.catch((error: any) => - error?._tag === "AlreadyExistsFault" - ? Effect.void -@@ -503,7 +514,9 @@ export const AutoScalingGroupProvider = () => + error?._tag === "AlreadyExistsFault" ? Effect.void : Effect.fail(error), + ), +@@ -481,7 +492,9 @@ export const AutoScalingGroupProvider = () => AutoScalingGroupName: autoScalingGroupName, MinSize: news.minSize, MaxSize: news.maxSize, @@ -146,10 +146,10 @@ index 5b40c65..2eb9162 100644 VPCZoneIdentifier: (news.subnetIds as string[]).join(","), HealthCheckType: healthCheckType, diff --git a/src/AWS/ECS/CapacityProvider.ts b/src/AWS/ECS/CapacityProvider.ts -index c1e65c2..c78d9c2 100644 +index cb68fdb..2857e7b 100644 --- a/src/AWS/ECS/CapacityProvider.ts +++ b/src/AWS/ECS/CapacityProvider.ts -@@ -188,7 +188,13 @@ export const CapacityProviderProvider = () => +@@ -180,7 +180,13 @@ export const CapacityProviderProvider = () => read: Effect.fn(function* ({ id, olds, output }) { const name = output?.name ?? (yield* toName(id, olds ?? {})); const found = yield* describe(name); @@ -165,10 +165,10 @@ index c1e65c2..c78d9c2 100644 } const internalTags = yield* createInternalTags(id); diff --git a/src/AWS/ECS/Service.ts b/src/AWS/ECS/Service.ts -index eac86f0..745df88 100644 +index 48a26ed..1a933af 100644 --- a/src/AWS/ECS/Service.ts +++ b/src/AWS/ECS/Service.ts -@@ -1658,6 +1658,11 @@ const transformServiceProps = ( +@@ -1609,6 +1609,11 @@ const transformServiceProps = ( }).pipe(Namespace.push(id)); }); @@ -180,7 +180,7 @@ index eac86f0..745df88 100644 const composeManagedIngress = ( id: string, props: ServiceProps, -@@ -2071,7 +2076,9 @@ const composeManagedIngress = ( +@@ -1998,7 +2003,9 @@ const composeManagedIngress = ( vpcId: network.vpcId as string, port: spec.port as number, protocol: spec.protocol, @@ -191,10 +191,10 @@ index eac86f0..745df88 100644 healthCheckPath: isNetworkTg ? wantsHttpCheck ? (health?.path ?? "/") -@@ -3249,7 +3256,10 @@ export const ServiceProvider = () => - healthCheckGracePeriodSeconds: toWireSeconds( - news.healthCheckGracePeriod, - ), +@@ -3099,7 +3106,10 @@ export const ServiceProvider = () => + platformVersion: news.platformVersion, + deploymentConfiguration: news.deploymentConfiguration, + healthCheckGracePeriodSeconds: toWireSeconds(news.healthCheckGracePeriod), - networkConfiguration: networkConfigurationOf(network, securityGroups), + // ECS rejects an awsvpcConfiguration for host/bridge tasks (maple patch). + networkConfiguration: usesAwsvpc(news) @@ -204,10 +204,10 @@ index eac86f0..745df88 100644 placementConstraints: news.placementConstraints, placementStrategy: news.placementStrategy, diff --git a/src/Cloudflare/Workers/DurableObject.ts b/src/Cloudflare/Workers/DurableObject.ts -index f849cca331..ebbebadd1f 100644 +index 07cea49..50e068f 100644 --- a/src/Cloudflare/Workers/DurableObject.ts +++ b/src/Cloudflare/Workers/DurableObject.ts -@@ -1303,7 +1303,9 @@ export const DurableObject: DurableObjectClass = taggedFunction( +@@ -1238,7 +1238,9 @@ export const DurableObject: DurableObjectClass = taggedFunction( }), ); @@ -218,11 +218,11 @@ index f849cca331..ebbebadd1f 100644 Type: TypeId, LogicalId: namespace, name: namespace, -@@ -1315,17 +1317,11 @@ export const DurableObject: DurableObjectClass = taggedFunction( - getByName: ( - name: string, - options?: DurableObjectGetDurableObjectOptions, -- ) => makeRpcStub(binding.getByName(name, options), { errors }), +@@ -1246,17 +1248,11 @@ export const DurableObject: DurableObjectClass = taggedFunction( + Output.map((durableObjectNamespaces) => durableObjectNamespaces?.[namespace]), + ), + getByName: (name: string, options?: DurableObjectGetDurableObjectOptions) => +- makeRpcStub(binding.getByName(name, options), { errors }), - // newUniqueId: () => use((ns) => ns.newUniqueId()), - // idFromName: (name: string) => use((ns) => ns.idFromName(name)), - // idFromString: (id: string) => use((ns) => ns.idFromString(id)), @@ -233,7 +233,7 @@ index f849cca331..ebbebadd1f 100644 - // jurisdiction: (jurisdiction: cf.DurableObjectJurisdiction) => - // use((ns) => ns.jurisdiction(jurisdiction) as any), - }; -+ ) => makeRpcStub(resolve().getByName(name, options), { errors }), ++ makeRpcStub(resolve().getByName(name, options), { errors }), + jurisdiction: (jurisdiction: DurableObjectJurisdiction) => + client(() => resolve().jurisdiction(jurisdiction)), + });