Compare commits

...
Author SHA1 Message Date
Kit Langton 41c0e8e472 refactor(core): share session message row codec 2026-09-27 19:10:32 -03:00
10 changed files with 65 additions and 49 deletions

No files matched your search

+1 -2
View File
@@ -156,8 +156,7 @@ const layer = Layer.effect(
})
const configured = Effect.fnUntraced(function* (sessionID: SessionSchema.ID, agentID?: Agent.ID) {
const session = yield* sessions.get(sessionID)
if (!session) return yield* new SessionErrors.NotFoundError({ sessionID })
const session = yield* sessions.require(sessionID)
const agent = yield* agents.resolve(agentID ?? session.agent)
return merge(agent?.permissions ?? missingAgentPermissions, session.permissions ?? [])
})
+3 -3
View File
@@ -1,8 +1,9 @@
import { and, asc, desc, eq, gte, or, sql } from "drizzle-orm"
import { Effect, Schema } from "effect"
import { Effect } from "effect"
import { Database } from "../database/database.js"
import { MessageDecodeError } from "./error.js"
import { SessionMessage } from "./message.js"
import { SessionMessageRow } from "./message-row.js"
import { SessionSchema } from "./schema.js"
import { Instructions } from "../instructions/index.js"
import { InstructionState } from "./instruction-state.js"
@@ -11,7 +12,6 @@ import { SessionMessageTable } from "./sql.js"
type DatabaseService = Database.Interface["db"]
const decode = Schema.decodeUnknownEffect(SessionMessage.Info)
/**
* Which completed compactions bound a history read. Local summaries always do. Native
@@ -60,7 +60,7 @@ export const latestCompaction = Effect.fnUntraced(function* (
})
export const decodeMessageRow = (row: typeof SessionMessageTable.$inferSelect) =>
decode({ ...row.data, id: row.id, type: row.type }).pipe(
SessionMessageRow.decodeEffect(row).pipe(
Effect.tap((message) =>
SessionProviderContext.isCheckpoint(message)
? SessionProviderContext.validate(message.providerContext)
+2 -2
View File
@@ -21,6 +21,7 @@ import { Bus } from "../bus.js"
import { KeyedMutex } from "../effect/keyed-mutex.js"
import { SessionEvent } from "./event.js"
import { SessionMessage } from "./message.js"
import { SessionMessageRow } from "./message-row.js"
import { SessionSchema } from "./schema.js"
import { SessionInboxTable, SessionMessageTable } from "./sql.js"
@@ -55,7 +56,6 @@ const decodeCompaction = Schema.decodeUnknownSync(CompactionPayload)
const encodeCompaction = Schema.encodeSync(CompactionPayload)
const decodeMove = Schema.decodeUnknownSync(MovePayload)
const encodeMove = Schema.encodeSync(MovePayload)
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info)
const inboxLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>()
type PendingRef = { readonly id: SessionMessage.ID; readonly sessionID: SessionSchema.ID }
@@ -125,7 +125,7 @@ const promotedFromMessage = Effect.fn("SessionInbox.promotedFromMessage")(functi
if (row === undefined) return undefined
if (row.session_id !== sessionID || (row.type !== "user" && row.type !== "synthetic"))
return yield* new LifecycleConflict({ id })
const message = decodeMessage({ ...row.data, id: row.id, type: row.type })
const message = SessionMessageRow.decode(row)
const base = { id, sessionID, time: { created: message.time.created }, delivery }
if (message.type === "user")
return User.make({
+25
View File
@@ -0,0 +1,25 @@
export * as SessionMessageRow from "./message-row.js"
import { Schema } from "effect"
import { SessionMessage } from "./message.js"
import type { SessionMessageTable } from "./sql.js"
type Row = Pick<typeof SessionMessageTable.$inferSelect, "id" | "type" | "data">
const decodeSync = Schema.decodeUnknownSync(SessionMessage.Info)
const decodeUnknown = Schema.decodeUnknownEffect(SessionMessage.Info)
const encodeSync = Schema.encodeSync(SessionMessage.Info)
const fields = (row: Row) => ({ ...row.data, id: row.id, type: row.type })
/** Decodes a stored message row, throwing when the stored data is invalid. */
export const decode = (row: Row) => decodeSync(fields(row))
export const decodeEffect = (row: Row) => decodeUnknown(fields(row))
/** Splits a message into its row identity, type, and JSON data columns. */
export const encode = (message: SessionMessage.Info) => {
const encoded = encodeSync(message)
const { id, type, ...data } = encoded
return { id: SessionMessage.ID.make(id), type, data }
}
+1 -5
View File
@@ -65,11 +65,7 @@ const layer = Layer.effect(
const database = yield* Database.Service
const bus = yield* Bus.Service
const get = Effect.fn("SessionMove.get")(function* (sessionID: Session.ID) {
const session = yield* store.get(sessionID)
if (!session) return yield* new NotFoundError({ sessionID })
return session
})
const get = store.require
const resolveDestination = Effect.fn("SessionMove.resolveDestination")(function* (
session: Session.Info,
+12 -18
View File
@@ -10,6 +10,7 @@ import { Agent } from "@opencode/schema/agent"
import { Model } from "@opencode/schema/model"
import { SessionEvent } from "./event.js"
import { SessionMessage } from "./message.js"
import { SessionMessageRow } from "./message-row.js"
import { SessionMessageUpdater } from "./message-updater.js"
import { SessionInbox } from "./inbox.js"
import { Workspace } from "@opencode/schema/workspace"
@@ -31,9 +32,6 @@ type MessageEvent = Exclude<
typeof SessionEvent.Forked.Type | typeof SessionEvent.Deleted.Type
>
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info)
const encodeMessage = Schema.encodeSync(SessionMessage.Info)
export class SessionAlreadyProjected extends Error {}
type Usage = {
@@ -227,17 +225,14 @@ 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 updateMessage = (message: SessionMessage.Info) => {
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 })
.set({ type: row.type, time_created: DateTime.toEpochMillis(message.time.created), data: row.data })
.where(
and(
eq(SessionMessageTable.id, SessionMessage.ID.make(id)),
eq(SessionMessageTable.id, row.id),
eq(SessionMessageTable.session_id, event.data.sessionID),
),
)
@@ -309,7 +304,7 @@ function run(db: DatabaseService, event: MessageEvent) {
.get()
.pipe(Effect.orDie)
if (!row) return
const message = decodeRow(row)
const message = SessionMessageRow.decode(row)
return message.type === "assistant" && !message.time.completed ? message : undefined
})
},
@@ -328,7 +323,7 @@ function run(db: DatabaseService, event: MessageEvent) {
.get()
.pipe(Effect.orDie)
if (!row) return
const message = decodeRow(row)
const message = SessionMessageRow.decode(row)
return message.type === "assistant" ? message : undefined
})
},
@@ -349,7 +344,7 @@ function run(db: DatabaseService, event: MessageEvent) {
.get()
.pipe(Effect.orDie)
if (!row) return
const message = decodeRow(row)
const message = SessionMessageRow.decode(row)
return message.type === "shell" ? message : undefined
})
},
@@ -370,7 +365,7 @@ function run(db: DatabaseService, event: MessageEvent) {
.get()
.pipe(Effect.orDie)
if (!row) return
const message = decodeRow(row)
const message = SessionMessageRow.decode(row)
return message.type === "compaction" ? message : undefined
})
},
@@ -384,17 +379,16 @@ function run(db: DatabaseService, event: MessageEvent) {
}
function insertMessage(db: DatabaseService, event: SessionEvent.DurableEvent, message: SessionMessage.Info) {
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)
+3 -3
View File
@@ -1,7 +1,7 @@
export * as SessionRevert from "./revert.js"
import { and, asc, eq, gt } from "drizzle-orm"
import { Effect, Schema } from "effect"
import { Effect } from "effect"
import { Database } from "../database/database.js"
import { Bus } from "../bus.js"
import { Instance } from "../instance/service.js"
@@ -10,6 +10,7 @@ import { Snapshot } from "../snapshot.js"
import { SessionEvent } from "./event.js"
import { MessageNotFoundError } from "./error.js"
import { SessionMessage } from "./message.js"
import { SessionMessageRow } from "./message-row.js"
import { SessionSchema } from "./schema.js"
import { SessionMessageTable } from "./sql.js"
@@ -104,10 +105,9 @@ const plan = Effect.fn("SessionRevert.plan")(function* (db: Database.Interface["
.orderBy(asc(SessionMessageTable.seq))
.all()
.pipe(Effect.orDie)
const decode = Schema.decodeUnknownEffect(SessionMessage.Info)
const files = new Map<RelativePath, Snapshot.ID>()
for (const row of rows) {
const message = yield* decode({ ...row.data, id: row.id, type: row.type }).pipe(Effect.orDie)
const message = yield* SessionMessageRow.decodeEffect(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))
+1 -6
View File
@@ -16,7 +16,6 @@ import {
CompactionConflictError,
InboxConflictError,
MessageNotFoundError,
NotFoundError,
PromptConflictError,
SyntheticConflictError,
} from "./error.js"
@@ -50,11 +49,7 @@ export const make = Effect.fn("Session.make")(function* () {
const fs = yield* FSUtil.Service
const scope = yield* Scope.Scope
const get = Effect.fn("Session.get")(function* (sessionID: SessionSchema.ID) {
const session = yield* store.get(sessionID)
if (!session) return yield* new NotFoundError({ sessionID })
return session
})
const get = store.require
const message = Effect.fn("Session.message")(function* (sessionID: SessionSchema.ID, messageID: SessionMessage.ID) {
const stored = yield* store.message(messageID)
return stored?.sessionID === sessionID ? stored.message : undefined
+12 -4
View File
@@ -8,7 +8,7 @@ import { AbsolutePath, PositiveInt, RelativePath } from "@opencode/schema/schema
import { Database } from "../database/database.js"
import { makeGlobalNode } from "@opencode/util/effect/app-node"
import { SessionHistory } from "./history.js"
import { MessageDecodeError } from "./error.js"
import { MessageDecodeError, NotFoundError } from "./error.js"
import { SessionMessage } from "./message.js"
import { Session } from "@opencode/schema/session"
import { SessionMessageTable, SessionTable } from "./sql.js"
@@ -52,6 +52,7 @@ export type MessagesInput = {
export interface Interface {
readonly get: (sessionID: Session.ID) => Effect.Effect<Session.Info | undefined>
readonly require: (sessionID: Session.ID) => Effect.Effect<Session.Info, NotFoundError>
readonly list: (input?: ListInput) => Effect.Effect<Session.Info[]>
readonly messages: (input: MessagesInput) => Effect.Effect<SessionMessage.Info[], MessageDecodeError>
readonly context: (sessionID: Session.ID) => Effect.Effect<SessionMessage.Info[], MessageDecodeError>
@@ -91,10 +92,17 @@ const layer = Layer.effect(
Effect.gen(function* () {
const { db } = yield* Database.Service
const get = Effect.fnUntraced(function* (sessionID: Session.ID) {
const row = yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie)
return row ? fromRow(row) : undefined
})
return Service.of({
get: Effect.fnUntraced(function* (sessionID) {
const row = yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie)
return row ? fromRow(row) : undefined
get,
require: Effect.fn("SessionStore.require")(function* (sessionID) {
const session = yield* get(sessionID)
if (!session) return yield* new NotFoundError({ sessionID })
return session
}),
list: Effect.fn("SessionStore.list")(function* (input = {}) {
const direction = input.anchor?.direction ?? "next"
+5 -6
View File
@@ -19,6 +19,7 @@ import { Session } from "../session.js"
import { Slug } from "../util/slug.js"
import { SessionEvent } from "./event.js"
import { SessionMessage } from "./message.js"
import { SessionMessageRow } from "./message-row.js"
import { SessionProjector } from "./projector.js"
import { SessionMessageTable, SessionTable } from "./sql.js"
@@ -51,7 +52,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)
return Service.of({
export: Effect.fn("SessionTransfer.export")(function* (input) {
@@ -75,15 +75,14 @@ const layer = Layer.effect(
yield* upsertProject(db, project).pipe(Effect.orDie)
const importedAt = yield* Clock.currentTimeMillis
const messages = input.data.messages.filter(isSettled).map((message, index) => {
const encoded = encodeMessage(message)
const { id: _, type, ...data } = encoded
const row = SessionMessageRow.encode(message)
return {
id: message.id,
id: row.id,
session_id: sessionID,
type,
type: row.type,
seq: index + 1,
time_created: DateTime.toEpochMillis(message.time.created),
data,
data: row.data,
}
})
yield* bus