Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions apps/ai/migrations/chat-session/0000_baseline.sql
Original file line number Diff line number Diff line change
@@ -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);
149 changes: 87 additions & 62 deletions apps/ai/src/chat/ChatSession.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -83,23 +83,40 @@ interface SessionRow extends Record<string, SqlStorageValue> {
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<string, SqlStorageValue> {
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<ColumnRow>("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
Expand Down Expand Up @@ -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. */
Expand Down Expand Up @@ -827,42 +831,63 @@ const isChatSessionNamespace = (
scope: Cloudflare.DurableObject,
): scope is Cloudflare.DurableObject<ChatSessionObject> => scope.name === "ChatSession"

/** The part of alchemy's `SqlMigrations` the activation uses; a test applies the files itself. */
export interface ChatSessionSchema<E, R> {
readonly apply: () => Effect.Effect<void, E, R>
}

/**
* 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 = <SE, SR, R>(
migrations: Effect.Effect<ChatSessionSchema<SE, SR>, 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)
71 changes: 53 additions & 18 deletions apps/ai/src/chat/ChatSessionObject.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"],
Expand All @@ -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<typeof makeFakeDurableObjectState>) =>
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", () =>
Expand All @@ -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)
Expand All @@ -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",
Expand All @@ -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",
])
}),
)
})
21 changes: 19 additions & 2 deletions apps/ai/test/chat/fake-do-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -24,14 +38,15 @@ export interface FakeDurableObjectState {
readonly pending: Array<Promise<unknown>>
}

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<Promise<unknown>> = []
const alarms: Array<number> = []

const sql = {
exec: (statement: string, ...bindings: ReadonlyArray<unknown>) => {
// `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)
Expand All @@ -48,6 +63,8 @@ export const makeFakeDurableObjectState = (): FakeDurableObjectState => {
},
}

if (migrated) applyChatSessionMigrations(sql as SqlStorage)

return {
storage: {
sql: sql as SqlStorage,
Expand Down
Loading
Loading