From ffaaad89a44fd4069c8f5f232c9cfaea1163dc6b Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Tue, 11 Aug 2026 15:32:51 -0400 Subject: [PATCH] refactor(core): centralize session message rows --- packages/core/src/session/history.ts | 7 ++-- packages/core/src/session/message-row.ts | 20 +++++++++++ packages/core/src/session/pending.ts | 4 +-- packages/core/src/session/projector.ts | 28 +++++----------- packages/core/src/session/revert.ts | 4 +-- packages/core/src/session/store.ts | 7 ++-- packages/core/src/session/transfer.ts | 10 +++--- .../core/test/session-message-row.test.ts | 33 +++++++++++++++++++ 8 files changed, 76 insertions(+), 37 deletions(-) create mode 100644 packages/core/src/session/message-row.ts create mode 100644 packages/core/test/session-message-row.test.ts diff --git a/packages/core/src/session/history.ts b/packages/core/src/session/history.ts index d2b024c276..0a04f857fa 100644 --- a/packages/core/src/session/history.ts +++ b/packages/core/src/session/history.ts @@ -1,8 +1,9 @@ import { and, asc, desc, eq, gte, sql } from "drizzle-orm" -import { Effect, Schema } from "effect" +import { Effect } from "effect" import { Database } from "../database/database" import { MessageDecodeError } from "./error" import { SessionMessage } from "./message" +import { SessionMessageRow } from "./message-row" import { SessionSchema } from "./schema" import { Instructions } from "../instructions/index" import { InstructionState } from "./instruction-state" @@ -10,8 +11,6 @@ import { SessionMessageTable } from "./sql" type DatabaseService = Database.Interface["db"] -const decode = Schema.decodeUnknownEffect(SessionMessage.Info) - export const latestCompaction = Effect.fnUntraced(function* (db: DatabaseService, sessionID: SessionSchema.ID) { return yield* db .select({ seq: SessionMessageTable.seq }) @@ -50,7 +49,7 @@ const messageRows = Effect.fnUntraced(function* ( }) const decodeMessageRow = (row: typeof SessionMessageTable.$inferSelect) => - decode({ ...row.data, id: row.id, type: row.type }).pipe( + SessionMessageRow.decode(row).pipe( Effect.mapError( () => new MessageDecodeError({ diff --git a/packages/core/src/session/message-row.ts b/packages/core/src/session/message-row.ts new file mode 100644 index 0000000000..2ee4f5f225 --- /dev/null +++ b/packages/core/src/session/message-row.ts @@ -0,0 +1,20 @@ +export * as SessionMessageRow from "./message-row" + +import { Schema } from "effect" +import { SessionMessage } from "./message" +import type { SessionMessageTable } from "./sql" + +export type Representation = Pick + +const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Info) +const decodeMessageSync = Schema.decodeUnknownSync(SessionMessage.Info) +const encodeMessage = Schema.encodeSync(SessionMessage.Info) + +export const decode = (row: Representation) => decodeMessage({ ...row.data, id: row.id, type: row.type }) + +export const decodeSync = (row: Representation) => decodeMessageSync({ ...row.data, id: row.id, type: row.type }) + +export function encode(message: SessionMessage.Info): Representation { + const { id, type, ...data } = encodeMessage(message) + return { id: SessionMessage.ID.make(id), type, data } +} diff --git a/packages/core/src/session/pending.ts b/packages/core/src/session/pending.ts index 9e6dad94d1..3b51c61c3c 100644 --- a/packages/core/src/session/pending.ts +++ b/packages/core/src/session/pending.ts @@ -17,6 +17,7 @@ import { Bus } from "../bus" import { KeyedMutex } from "../effect/keyed-mutex" import { SessionEvent } from "./event" import { SessionMessage } from "./message" +import { SessionMessageRow } from "./message-row" import { SessionSchema } from "./schema" import { SessionMessageTable, SessionPendingTable } from "./sql" @@ -35,7 +36,6 @@ const decodeUser = Schema.decodeUnknownSync(UserData) const encodeUser = Schema.encodeSync(UserData) const decodeSynthetic = Schema.decodeUnknownSync(SyntheticData) const encodeSynthetic = Schema.encodeSync(SyntheticData) -const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info) const inboxLocks = KeyedMutex.makeUnsafe() type PendingRef = { readonly id: SessionMessage.ID; readonly sessionID: SessionSchema.ID } @@ -113,7 +113,7 @@ const promotedFromMessage = Effect.fn("SessionPending.promotedFromMessage")(func if (row === undefined) return undefined if (row.session_id !== sessionID || (row.type !== "user" && row.type !== "synthetic")) return yield* Effect.die(new LifecycleConflict({ id })) - const message = decodeMessage({ ...row.data, id: row.id, type: row.type }) + const message = SessionMessageRow.decodeSync(row) const base = { id, sessionID, timeCreated: message.time.created, delivery } if (message.type === "user") return User.make({ diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index 3588dde6ef..d8ab0aef68 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -9,6 +9,7 @@ import { Agent } from "../agent" import { Model } from "../model" import { SessionEvent } from "./event" import { SessionMessage } from "./message" +import { SessionMessageRow } from "./message-row" import { SessionMessageUpdater } from "./message-updater" import { SessionPending } from "./pending" import { Workspace } from "../workspace" @@ -22,9 +23,6 @@ type DatabaseService = Database.Interface["db"] type CurrentDurableEvent = Extract type MessageEvent = Exclude -const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info) -const encodeMessage = Schema.encodeSync(SessionMessage.Info) - export class SessionAlreadyProjected extends Error {} type Usage = { @@ -210,22 +208,15 @@ const projectFork = Effect.fn("SessionProjector.projectFork")(function* ( function run(db: DatabaseService, event: MessageEvent) { return Effect.gen(function* () { - const decodeRow = (row: typeof SessionMessageTable.$inferSelect) => - decodeMessage({ ...row.data, id: row.id, type: row.type }) + const decodeRow = (row: typeof SessionMessageTable.$inferSelect) => SessionMessageRow.decodeSync(row) const updateMessage = (message: SessionMessage.Info) => { if (event.durable === undefined) return Effect.die(new Error("Durable Session event is missing aggregate sequence")) - const encoded = encodeMessage(message) - const { id, type, ...data } = encoded + const row = SessionMessageRow.encode(message) return db .update(SessionMessageTable) - .set({ type, time_created: DateTime.toEpochMillis(message.time.created), data }) - .where( - and( - eq(SessionMessageTable.id, SessionMessage.ID.make(id)), - eq(SessionMessageTable.session_id, event.data.sessionID), - ), - ) + .set({ type: row.type, time_created: DateTime.toEpochMillis(message.time.created), data: row.data }) + .where(and(eq(SessionMessageTable.id, row.id), eq(SessionMessageTable.session_id, event.data.sessionID))) .run() .pipe(Effect.orDie) } @@ -343,17 +334,16 @@ function run(db: DatabaseService, event: MessageEvent) { function insertMessage(db: DatabaseService, event: SessionEvent.DurableEvent, message: SessionMessage.Info) { if (event.durable === undefined) return Effect.die(new Error("Durable Session event is missing aggregate sequence")) - const encoded = encodeMessage(message) - const { id, type, ...data } = encoded + const row = SessionMessageRow.encode(message) return db .insert(SessionMessageTable) .values({ - id: SessionMessage.ID.make(id), + id: row.id, session_id: event.data.sessionID, - type, + type: row.type, seq: event.durable.seq, time_created: DateTime.toEpochMillis(message.time.created), - data, + data: row.data, }) .run() .pipe(Effect.orDie) diff --git a/packages/core/src/session/revert.ts b/packages/core/src/session/revert.ts index 429170e987..1f415d5ba7 100644 --- a/packages/core/src/session/revert.ts +++ b/packages/core/src/session/revert.ts @@ -8,6 +8,7 @@ import { RelativePath } from "../schema" import { Snapshot } from "../snapshot" import { SessionEvent } from "./event" import { SessionMessage } from "./message" +import { SessionMessageRow } from "./message-row" import { SessionSchema } from "./schema" import { SessionMessageTable } from "./sql" @@ -46,10 +47,9 @@ const plan = Effect.fn("SessionRevert.plan")(function* (input: BoundaryInput) { .orderBy(asc(SessionMessageTable.seq)) .all() .pipe(Effect.orDie) - const decode = Schema.decodeUnknownEffect(SessionMessage.Info) const files = new Map() for (const row of rows) { - const message = yield* decode({ ...row.data, id: row.id, type: row.type }).pipe(Effect.orDie) + const message = yield* SessionMessageRow.decode(row).pipe(Effect.orDie) if (message.type !== "assistant" || !message.snapshot?.start) continue for (const file of message.snapshot.files ?? []) if (!files.has(file)) files.set(file, Snapshot.ID.make(message.snapshot.start)) diff --git a/packages/core/src/session/store.ts b/packages/core/src/session/store.ts index 1101b950ad..7805d3b407 100644 --- a/packages/core/src/session/store.ts +++ b/packages/core/src/session/store.ts @@ -1,12 +1,13 @@ export * as SessionStore from "./store" import { and, eq, isNotNull, isNull, sql } from "drizzle-orm" -import { Context, Effect, Layer, Schema } from "effect" +import { Context, Effect, Layer } from "effect" import { Database } from "../database/database" import { makeGlobalNode } from "@opencode-ai/util/effect/app-node" import { SessionHistory } from "./history" import { MessageDecodeError } from "./error" import { SessionMessage } from "./message" +import { SessionMessageRow } from "./message-row" import { Session } from "@opencode-ai/schema/session" import { SessionMessageTable, SessionTable } from "./sql" import { fromRow } from "./info" @@ -51,8 +52,6 @@ const layer = Layer.effect( Service, Effect.gen(function* () { const { db } = yield* Database.Service - const decodeMessage = Schema.decodeUnknownEffect(SessionMessage.Info) - return Service.of({ get: Effect.fn("SessionStore.get")(function* (sessionID) { const row = yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie) @@ -71,7 +70,7 @@ const layer = Layer.effect( return row ? { sessionID: Session.ID.make(row.session_id), - message: yield* decodeMessage({ ...row.data, id: row.id, type: row.type }).pipe(Effect.orDie), + message: yield* SessionMessageRow.decode(row).pipe(Effect.orDie), } : undefined }), diff --git a/packages/core/src/session/transfer.ts b/packages/core/src/session/transfer.ts index e9919deadb..1a877c940f 100644 --- a/packages/core/src/session/transfer.ts +++ b/packages/core/src/session/transfer.ts @@ -18,6 +18,7 @@ import { Session } from "../session" import { Slug } from "../util/slug" import { SessionEvent } from "./event" import { SessionMessage } from "./message" +import { SessionMessageRow } from "./message-row" import { SessionProjector } from "./projector" import { SessionMessageTable, SessionTable } from "./sql" @@ -47,8 +48,6 @@ const layer = Layer.effect( const { db } = yield* Database.Service const projects = yield* Project.Service const sessions = yield* Session.Service - const encodeMessage = Schema.encodeSync(SessionMessage.Info) - const persistProject = (project: Project.Resolved) => upsertProject(db, project).pipe(Effect.orDie) return Service.of({ @@ -71,15 +70,14 @@ const layer = Layer.effect( const project = yield* projects.resolve(input.location.directory) yield* persistProject(project) const messages = input.data.messages.map((message, index) => { - const encoded = encodeMessage(message) - const { id: _, type, ...data } = encoded + const row = SessionMessageRow.encode(message) return { id: message.id, session_id: sessionID, - type, + type: row.type, seq: index + 1, time_created: DateTime.toEpochMillis(message.time.created), - data, + data: row.data, } }) yield* bus diff --git a/packages/core/test/session-message-row.test.ts b/packages/core/test/session-message-row.test.ts new file mode 100644 index 0000000000..678c27f414 --- /dev/null +++ b/packages/core/test/session-message-row.test.ts @@ -0,0 +1,33 @@ +import { expect, test } from "bun:test" +import { DateTime, Effect } from "effect" +import { SessionMessage } from "@opencode-ai/core/session/message" +import { SessionMessageRow } from "@opencode-ai/core/session/message-row" + +const message = SessionMessage.Synthetic.make({ + id: SessionMessage.ID.make("msg_row"), + type: "synthetic", + text: "hello", + time: { created: DateTime.makeUnsafe(1_000) }, +}) + +test("round trips the persisted message representation", async () => { + const row = SessionMessageRow.encode(message) + + expect(row.id).toBe(message.id) + expect(row.type).toBe(message.type) + expect(row.data).toHaveProperty("text", message.text) + expect(row.data).toHaveProperty("time.created", 1_000) + expect(await Effect.runPromise(SessionMessageRow.decode(row))).toEqual(message) + expect(SessionMessageRow.decodeSync(row)).toEqual(message) +}) + +test("canonical columns override stale values in message data", () => { + const row = SessionMessageRow.encode(message) + const data = { + ...row.data, + id: SessionMessage.ID.make("msg_stale"), + type: "system" as const, + } + + expect(SessionMessageRow.decodeSync({ ...row, data })).toEqual(message) +})