refactor(schema): rename V2 session events and normalize payloads (#35217)

This commit is contained in:
Kit Langton 2026-07-03 14:25:59 -04:00 committed by GitHub
commit 394e0b9045
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
63 changed files with 1888 additions and 1989 deletions

View file

@ -106,8 +106,7 @@ const layer = Layer.effect(
yield* events.publish(SessionEvent.Moved, {
sessionID: input.sessionID,
location: Location.Ref.make({ directory }),
subdirectory: RelativePath.make(path.relative(destination.directory, directory).replaceAll("\\", "/")),
timestamp: yield* DateTime.now,
subpath: RelativePath.make(path.relative(destination.directory, directory).replaceAll("\\", "/")),
})
if (patch) {

View file

@ -41,5 +41,7 @@ export const migrations = (
import("./migration/20260622170816_reset_v2_session_state"),
import("./migration/20260622202450_simplify_session_input"),
import("./migration/20260702134641_add_session_context_entry"),
import("./migration/20260703090000_reset_v2_event_rename_sweep"),
import("./migration/20260703181610_event_created_column"),
])
).map((module) => module.default) satisfies DatabaseMigration.Migration[]

View file

@ -0,0 +1,17 @@
import { Effect } from "effect"
import type { DatabaseMigration } from "../migration"
export default {
id: "20260703090000_reset_v2_event_rename_sweep",
up(tx) {
return Effect.gen(function* () {
yield* tx.run(`DELETE FROM \`session_input\`;`)
yield* tx.run(`DELETE FROM \`session_message\`;`)
yield* tx.run(`DELETE FROM \`event\`;`)
yield* tx.run(`DELETE FROM \`event_sequence\`;`)
// `created` column is added by the generated 20260703181610_event_created_column
// migration, which runs after this wipe (NOT NULL without default is safe on the
// emptied table).
})
},
} satisfies DatabaseMigration.Migration

View file

@ -0,0 +1,11 @@
import { Effect } from "effect"
import type { DatabaseMigration } from "../migration"
export default {
id: "20260703181610_event_created_column",
up(tx) {
return Effect.gen(function* () {
yield* tx.run(`ALTER TABLE \`event\` ADD \`created\` integer NOT NULL;`)
})
},
} satisfies DatabaseMigration.Migration

View file

@ -81,6 +81,7 @@ export default {
\`id\` text PRIMARY KEY,
\`aggregate_id\` text NOT NULL,
\`seq\` integer NOT NULL,
\`created\` integer NOT NULL,
\`type\` text NOT NULL,
\`data\` text NOT NULL,
CONSTRAINT \`fk_event_aggregate_id_event_sequence_aggregate_id_fk\` FOREIGN KEY (\`aggregate_id\`) REFERENCES \`event_sequence\`(\`aggregate_id\`) ON DELETE CASCADE

View file

@ -1,6 +1,6 @@
export * as EventV2 from "./event"
import { Cause, Context, Effect, Layer, Option, PubSub, Queue, Schema, Stream } from "effect"
import { Cause, Context, DateTime, Effect, Layer, Option, PubSub, Queue, Schema, Stream } from "effect"
import { Event } from "@opencode-ai/schema/event"
import type { Data, Definition, Payload } from "@opencode-ai/schema/event"
import type { EventLog } from "@opencode-ai/schema/event-log"
@ -55,6 +55,7 @@ export const reserveSequence = Effect.fn("EventV2.reserveSequence")(function* (
export type SerializedEvent = {
readonly id: ID
readonly type: string
readonly created?: DateTime.Utc
readonly seq: number
readonly aggregateID: string
readonly data: Record<string, unknown>
@ -81,6 +82,7 @@ const decodeSerializedEvent = (event: SerializedEvent): Payload => {
}
return {
id: event.id,
created: event.created ?? DateTime.makeUnsafe(0),
type: definition.type,
durable: envelope(event.aggregateID, event.seq, definition.durable.version),
data: Schema.decodeUnknownSync(definition.data)(event.data),
@ -295,6 +297,7 @@ export const layerWith = (options?: LayerOptions) =>
if (
stored?.id === event.id &&
stored.type === versionedType(definition.type, durable.version) &&
stored.created === DateTime.toEpochMillis(event.created ?? DateTime.makeUnsafe(0)) &&
isDeepStrictEqual(stored.data, encoded)
) {
if (input.ownerID && row?.ownerID == null) {
@ -366,6 +369,7 @@ export const layerWith = (options?: LayerOptions) =>
id: event.id,
aggregate_id: aggregateID,
seq,
created: DateTime.toEpochMillis(event.created ?? DateTime.makeUnsafe(0)),
type: versionedType(definition.type, durable.version),
data: encoded,
},
@ -471,6 +475,7 @@ export const layerWith = (options?: LayerOptions) =>
definition,
{
id: options?.id ?? ID.create(),
created: yield* DateTime.now,
...(options?.metadata ? { metadata: options.metadata } : {}),
type: definition.type,
...(location ? { location } : {}),
@ -494,6 +499,7 @@ export const layerWith = (options?: LayerOptions) =>
} else {
const payload = {
id: event.id,
created: event.created ?? DateTime.makeUnsafe(0),
type: definition.type,
data: Schema.decodeUnknownSync(definition.data)(event.data),
} as Payload
@ -610,6 +616,7 @@ export const layerWith = (options?: LayerOptions) =>
return [
decodeSerializedEvent({
id: event.id,
created: DateTime.makeUnsafe(event.created),
aggregateID: event.aggregate_id,
seq: event.seq,
type: event.type,

View file

@ -15,6 +15,7 @@ export const EventTable = sqliteTable(
.notNull()
.references(() => EventSequenceTable.aggregate_id, { onDelete: "cascade" }),
seq: integer().notNull(),
created: integer().notNull(),
type: text().notNull(),
data: text({ mode: "json" }).$type<Record<string, unknown>>().notNull(),
},

View file

@ -363,8 +363,7 @@ const layer = Layer.effect(
yield* events.publish(SessionEvent.Forked, {
sessionID,
parentID: parent.id,
messageID: input.messageID,
timestamp: yield* DateTime.now,
from: input.messageID,
})
return yield* result.get(sessionID).pipe(Effect.orDie)
}),
@ -551,16 +550,13 @@ const layer = Layer.effect(
Effect.gen(function* () {
activeShells.add(input.sessionID)
if ((yield* execution.active).has(input.sessionID)) yield* execution.awaitIdle(input.sessionID)
const messageID = SessionMessage.ID.create()
const callID = Identifier.ascending()
yield* events.publish(
SessionEvent.Shell.Started,
{
sessionID: input.sessionID,
messageID,
callID,
command: input.command,
timestamp: yield* DateTime.now,
},
{ id: input.id },
)
@ -571,7 +567,6 @@ const layer = Layer.effect(
sessionID: input.sessionID,
callID,
output,
timestamp: yield* DateTime.now,
})
}).pipe(
Effect.ensuring(
@ -590,8 +585,6 @@ const layer = Layer.effect(
if (!skill) return yield* new SkillNotFoundError({ skill: input.skill })
yield* events.publish(SessionEvent.Skill.Activated, {
sessionID: input.sessionID,
messageID: input.id ?? SessionMessage.ID.create(),
timestamp: yield* DateTime.now,
name: skill.name,
text: skill.content,
})
@ -602,10 +595,8 @@ const layer = Layer.effect(
}),
switchAgent: Effect.fn("V2Session.switchAgent")(function* (input) {
yield* result.get(input.sessionID)
yield* events.publish(SessionEvent.AgentSwitched, {
yield* events.publish(SessionEvent.AgentSelected, {
sessionID: input.sessionID,
messageID: SessionMessage.ID.create(),
timestamp: yield* DateTime.now,
agent: input.agent,
})
}),
@ -617,10 +608,8 @@ const layer = Layer.effect(
(session.model.variant ?? "default") === (input.model.variant ?? "default")
)
return
yield* events.publish(SessionEvent.ModelSwitched, {
yield* events.publish(SessionEvent.ModelSelected, {
sessionID: input.sessionID,
messageID: SessionMessage.ID.create(),
timestamp: yield* DateTime.now,
model: input.model,
})
}),
@ -628,7 +617,6 @@ const layer = Layer.effect(
yield* result.get(input.sessionID)
yield* events.publish(SessionEvent.Renamed, {
sessionID: input.sessionID,
timestamp: yield* DateTime.now,
title: input.title,
})
}),
@ -676,8 +664,6 @@ const layer = Layer.effect(
yield* result.get(input.sessionID)
yield* events.publish(SessionEvent.Synthetic, {
sessionID: input.sessionID,
messageID: SessionMessage.ID.create(),
timestamp: yield* DateTime.now,
text: input.text,
description: input.description,
metadata: input.metadata,

View file

@ -7,7 +7,7 @@ import { EventV2 } from "../event"
import { makeLocationNode } from "../effect/app-node"
import { llmClient } from "../effect/app-node-platform"
import { SessionEvent } from "./event"
import { SessionMessage } from "./message"
import type { SessionMessage } from "./message"
import { SessionRunnerModel } from "./runner/model"
import { SessionSchema } from "./schema"
import { Token } from "../util/token"
@ -206,11 +206,8 @@ const make = (dependencies: Dependencies) => {
const summaryPrompt = buildPrompt({ previousSummary: input.previousSummary, context: input.context })
const summaryOutput = Math.min(output || SUMMARY_OUTPUT_TOKENS, SUMMARY_OUTPUT_TOKENS)
if (Token.estimate(summaryPrompt) > context - summaryOutput) return false
const messageID = SessionMessage.ID.create()
yield* dependencies.events.publish(SessionEvent.Compaction.Started, {
sessionID: input.sessionID,
messageID,
timestamp: yield* DateTime.now,
reason: input.reason,
})
@ -238,8 +235,6 @@ const make = (dependencies: Dependencies) => {
if (!summarized || failed || !summary.trim()) return false
yield* dependencies.events.publish(SessionEvent.Compaction.Ended, {
sessionID: input.sessionID,
messageID,
timestamp: yield* DateTime.now,
reason: input.reason,
text: summary,
recent: input.recent,

View file

@ -1,13 +1,12 @@
export * as SessionContextCheckpoint from "./context-checkpoint"
import { eq } from "drizzle-orm"
import { DateTime, Effect, Option, Schema } from "effect"
import { Effect, Option, Schema } from "effect"
import type { Database } from "../database/database"
import { EventV2 } from "../event"
import { SystemContext } from "../system-context/index"
import { SessionEvent } from "./event"
import { SessionHistory } from "./history"
import { SessionMessage } from "./message"
import { SessionSchema } from "./schema"
import { SessionContextCheckpointTable } from "./sql"
@ -50,7 +49,7 @@ export const prepare = Effect.fn("SessionContextCheckpoint.prepare")(function* (
yield* events.publish(
SessionEvent.ContextUpdated,
{ sessionID, messageID: SessionMessage.ID.create(), timestamp: yield* DateTime.now, text: result.text },
{ sessionID, text: result.text },
{ commit: () => advance(db, sessionID, result.applied).pipe(Effect.orDie) },
)
return { baseline: stored.baseline, baselineSeq: stored.baseline_seq }

View file

@ -36,7 +36,6 @@ const layer = Layer.effect(
Exit.isFailure(exit) && !Cause.hasInterrupts(exit.cause) ? Cause.squash(exit.cause) : undefined
yield* events.publish(SessionEvent.ExecutionSettled, {
sessionID,
timestamp: yield* DateTime.now,
outcome: Exit.isSuccess(exit) ? "success" : Cause.hasInterrupts(exit.cause) ? "interrupted" : "failure",
error:
failure !== undefined

View file

@ -50,12 +50,10 @@ export const admit = Effect.fn("SessionInput.admit")(function* (
) {
const existing = yield* find(db, input.id)
if (existing !== undefined) return existing
const timestamp = yield* DateTime.now
return yield* events
.publish(SessionEvent.PromptAdmitted, {
messageID: input.id,
inputID: input.id,
sessionID: input.sessionID,
timestamp,
prompt: input.prompt,
delivery: input.delivery,
})
@ -70,7 +68,7 @@ export const admit = Effect.fn("SessionInput.admit")(function* (
sessionID: input.sessionID,
prompt: input.prompt,
delivery: input.delivery,
timeCreated: timestamp,
timeCreated: event.created,
}),
),
),
@ -115,14 +113,11 @@ export const projectAdmitted = Effect.fn("SessionInput.projectAdmitted")(functio
if (!stored) return yield* Effect.die(new LifecycleConflict({ id: input.id }))
})
export const projectPrompted = Effect.fn("SessionInput.projectPrompted")(function* (
export const projectPromptPromoted = Effect.fn("SessionInput.projectPromptPromoted")(function* (
db: DatabaseService,
input: {
readonly id: SessionMessage.ID
readonly sessionID: SessionSchema.ID
readonly prompt: Prompt
readonly delivery: Delivery
readonly timeCreated: DateTime.Utc
readonly promotedSeq: number
},
) {
@ -141,15 +136,16 @@ export const projectPrompted = Effect.fn("SessionInput.projectPrompted")(functio
.pipe(Effect.orDie)
if (updated) {
const stored = fromRow(updated)
if (!matchesProjection(stored, input)) return yield* Effect.die(new LifecycleConflict({ id: input.id }))
return
if (stored.sessionID !== input.sessionID) return yield* Effect.die(new LifecycleConflict({ id: input.id }))
return stored
}
// Every Prompted event is published from an admitted inbox row, so a missing or
// Every PromptPromoted event is published from an admitted inbox row, so a missing or
// divergent row on replay is an invariant violation.
const stored = yield* find(db, input.id)
if (!stored || !matchesProjection(stored, input) || stored.promotedSeq !== input.promotedSeq)
if (!stored || stored.sessionID !== input.sessionID || stored.promotedSeq !== input.promotedSeq)
return yield* Effect.die(new LifecycleConflict({ id: input.id }))
return stored
})
export const hasPending = Effect.fn("SessionInput.hasPending")(function* (
@ -206,12 +202,9 @@ const publish = Effect.fn("SessionInput.publish")(function* (
for (const row of rows) {
const id = SessionMessage.ID.make(row.id)
yield* events
.publish(SessionEvent.Prompted, {
.publish(SessionEvent.PromptPromoted, {
sessionID,
timestamp: DateTime.makeUnsafe(row.time_created),
messageID: id,
prompt: decodePrompt(row.prompt),
delivery: row.delivery,
inputID: id,
})
.pipe(
Effect.catchDefect((defect) =>

View file

@ -60,9 +60,9 @@ const layer = Layer.effect(
const files = yield* Effect.forEach(
toInject,
(path) =>
fs.readFileStringSafe(path).pipe(
Effect.map((content) => (content === undefined ? undefined : { path, content })),
),
fs
.readFileStringSafe(path)
.pipe(Effect.map((content) => (content === undefined ? undefined : { path, content }))),
{ concurrency: "unbounded" },
)
const readable = files.filter((file): file is { path: string; content: string } => file !== undefined)
@ -74,8 +74,6 @@ const layer = Layer.effect(
// metadata so it survives across Location layer restarts.
yield* events.publish(SessionEvent.Synthetic, {
sessionID: input.sessionID,
messageID: SessionMessage.ID.create(),
timestamp: yield* DateTime.now,
text: readable.map((file) => `Instructions from: ${file.path}\n${file.content}`).join("\n\n"),
description: `Loaded ${readable.map((file) => describePath(root, file.path)).join(", ")}`,
metadata: { instruction: { paths: readable.map((file) => file.path) } },

View file

@ -100,112 +100,100 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
return Effect.gen(function* () {
yield* SessionEvent.All.match(event, {
"session.next.agent.switched": (event) => {
"agent.selected": (event) => {
return adapter.appendMessage(
SessionMessage.AgentSwitched.make({
id: event.data.messageID,
SessionMessage.AgentSelected.make({
id: SessionMessage.ID.fromEvent(event.id),
type: "agent-switched",
metadata: event.metadata,
agent: event.data.agent,
time: { created: event.data.timestamp },
time: { created: event.created },
}),
)
},
"session.next.model.switched": (event) => {
"model.selected": (event) => {
return adapter.appendMessage(
SessionMessage.ModelSwitched.make({
id: event.data.messageID,
SessionMessage.ModelSelected.make({
id: SessionMessage.ID.fromEvent(event.id),
type: "model-switched",
metadata: event.metadata,
model: event.data.model,
time: { created: event.data.timestamp },
time: { created: event.created },
}),
)
},
"session.next.moved": () => Effect.void,
"session.next.renamed": () => Effect.void,
"session.next.forked": () => Effect.void,
"session.next.prompted": (event) => {
return adapter.appendMessage(
SessionMessage.User.make({
id: event.data.messageID,
type: "user",
metadata: event.metadata,
text: event.data.prompt.text,
files: event.data.prompt.files,
agents: event.data.prompt.agents,
time: { created: event.data.timestamp },
}),
)
},
"session.next.prompt.admitted": () => Effect.void,
"session.next.execution.settled": () => Effect.void,
"session.next.context.updated": (event) =>
"session.moved": () => Effect.void,
renamed: () => Effect.void,
forked: () => Effect.void,
"prompt.promoted": () => Effect.void,
"prompt.admitted": () => Effect.void,
"execution.settled": () => Effect.void,
"session.context.updated": (event) =>
adapter.appendMessage(
SessionMessage.System.make({
id: event.data.messageID,
id: SessionMessage.ID.fromEvent(event.id),
type: "system",
text: event.data.text,
time: { created: event.data.timestamp },
time: { created: event.created },
}),
),
"session.next.synthetic": (event) => {
synthetic: (event) => {
return adapter.appendMessage(
SessionMessage.Synthetic.make({
sessionID: event.data.sessionID,
text: event.data.text,
description: event.data.description,
metadata: event.data.metadata,
id: event.data.messageID,
id: SessionMessage.ID.fromEvent(event.id),
type: "synthetic",
time: { created: event.data.timestamp },
time: { created: event.created },
}),
)
},
"session.next.skill.activated": (event) => {
"skill.activated": (event) => {
return adapter.appendMessage(
SessionMessage.Skill.make({
id: event.data.messageID,
id: SessionMessage.ID.fromEvent(event.id),
type: "skill",
name: event.data.name,
text: event.data.text,
time: { created: event.data.timestamp },
time: { created: event.created },
}),
)
},
"session.next.shell.started": (event) => {
"shell.started": (event) => {
return adapter.appendMessage(
SessionMessage.Shell.make({
id: event.data.messageID,
id: SessionMessage.ID.fromEvent(event.id),
type: "shell",
metadata: event.metadata,
callID: event.data.callID,
command: event.data.command,
output: "",
time: { created: event.data.timestamp },
time: { created: event.created },
}),
)
},
"session.next.shell.ended": (event) => {
"shell.ended": (event) => {
return Effect.gen(function* () {
const currentShell = yield* adapter.getCurrentShell(event.data.callID)
if (currentShell) {
yield* adapter.updateShell(
produce(currentShell, (draft) => {
draft.output = event.data.output
draft.time.completed = event.data.timestamp
draft.time.completed = event.created
}),
)
}
})
},
"session.next.step.started": (event) => {
"step.started": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.time.completed = event.data.timestamp
draft.time.completed = event.created
}),
)
}
@ -215,16 +203,16 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
type: "assistant",
agent: event.data.agent,
model: event.data.model,
time: { created: event.data.timestamp },
time: { created: event.created },
content: [],
snapshot: event.data.snapshot ? { start: event.data.snapshot } : undefined,
}),
)
})
},
"session.next.step.ended": (event) => {
"step.ended": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
draft.time.completed = event.data.timestamp
draft.time.completed = event.created
draft.finish = event.data.finish
draft.cost = event.data.cost
draft.tokens = event.data.tokens
@ -236,33 +224,33 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
}
})
},
"session.next.step.failed": (event) => {
"step.failed": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
draft.time.completed = event.data.timestamp
draft.time.completed = event.created
draft.finish = "error"
draft.error = event.data.error
})
},
"session.next.text.started": (event) => {
"text.started": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
draft.content.push(
castDraft(SessionMessage.AssistantText.make({ type: "text", id: event.data.textID, text: "" })),
)
})
},
"session.next.text.delta": (event) => {
"text.delta": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestText(draft, event.data.textID)
if (match) match.text += event.data.delta
})
},
"session.next.text.ended": (event) => {
"text.ended": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestText(draft, event.data.textID)
if (match) match.text = event.data.text
})
},
"session.next.tool.input.started": (event) => {
"tool.input.started": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
draft.content.push(
castDraft(
@ -270,26 +258,26 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
type: "tool",
id: event.data.callID,
name: event.data.name,
time: { created: event.data.timestamp },
time: { created: event.created },
state: SessionMessage.ToolStatePending.make({ status: "pending", input: "" }),
}),
),
)
})
},
"session.next.tool.input.delta": () => Effect.void,
"session.next.tool.input.ended": (event) => {
"tool.input.delta": () => Effect.void,
"tool.input.ended": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "pending") match.state.input = event.data.text
})
},
"session.next.tool.called": (event) => {
"tool.called": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match) {
match.provider = event.data.provider
match.time.ran = event.data.timestamp
match.time.ran = event.created
match.state = castDraft(
SessionMessage.ToolStateRunning.make({
status: "running",
@ -301,7 +289,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
}
})
},
"session.next.tool.progress": (event) => {
"tool.progress": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "running") {
@ -310,7 +298,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
}
})
},
"session.next.tool.success": (event) => {
"tool.success": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "running") {
@ -319,7 +307,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
metadata: match.provider?.metadata,
resultMetadata: event.data.provider.metadata,
}
match.time.completed = event.data.timestamp
match.time.completed = event.created
match.state = castDraft(
SessionMessage.ToolStateCompleted.make({
status: "completed",
@ -333,7 +321,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
}
})
},
"session.next.tool.failed": (event) => {
"tool.failed": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && (match.state.status === "pending" || match.state.status === "running")) {
@ -342,7 +330,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
metadata: match.provider?.metadata,
resultMetadata: event.data.provider.metadata,
}
match.time.completed = event.data.timestamp
match.time.completed = event.created
match.state = castDraft(
SessionMessage.ToolStateError.make({
status: "error",
@ -356,7 +344,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
}
})
},
"session.next.reasoning.started": (event) => {
"reasoning.started": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
draft.content.push(
castDraft(
@ -365,47 +353,47 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
id: event.data.reasoningID,
text: "",
providerMetadata: event.data.providerMetadata,
time: { created: event.data.timestamp },
time: { created: event.created },
}),
),
)
})
},
"session.next.reasoning.delta": (event) => {
"reasoning.delta": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestReasoning(draft, event.data.reasoningID)
if (match) match.text += event.data.delta
})
},
"session.next.reasoning.ended": (event) => {
"reasoning.ended": (event) => {
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
const match = latestReasoning(draft, event.data.reasoningID)
if (match) {
match.text = event.data.text
match.time = { created: match.time?.created ?? event.data.timestamp, completed: event.data.timestamp }
match.time = { created: match.time?.created ?? event.created, completed: event.created }
if (event.data.providerMetadata !== undefined) match.providerMetadata = event.data.providerMetadata
}
})
},
"session.next.retried": () => Effect.void,
"session.next.compaction.started": () => Effect.void,
"session.next.compaction.delta": () => Effect.void,
"session.next.compaction.ended": (event) => {
retried: () => Effect.void,
"compaction.started": () => Effect.void,
"compaction.delta": () => Effect.void,
"compaction.ended": (event) => {
return adapter.appendMessage(
SessionMessage.Compaction.make({
id: event.data.messageID,
id: SessionMessage.ID.fromEvent(event.id),
type: "compaction",
metadata: event.metadata,
reason: event.data.reason,
summary: event.data.text,
recent: event.data.recent,
time: { created: event.data.timestamp },
time: { created: event.created },
}),
)
},
"session.next.revert.staged": () => Effect.void,
"session.next.revert.cleared": () => Effect.void,
"session.next.revert.committed": () => Effect.void,
"revert.staged": () => Effect.void,
"revert.cleared": () => Effect.void,
"revert.committed": () => Effect.void,
})
})
}

View file

@ -158,21 +158,17 @@ const projectFork = Effect.fn("SessionProjector.projectFork")(function* (
.get()
.pipe(Effect.orDie)
if (!parent) return yield* Effect.die(new Error(`Fork parent session not found: ${event.data.parentID}`))
const boundary = event.data.messageID
const boundary = event.data.from
? yield* db
.select({ seq: SessionMessageTable.seq })
.from(SessionMessageTable)
.where(
and(
eq(SessionMessageTable.session_id, event.data.parentID),
eq(SessionMessageTable.id, event.data.messageID),
),
and(eq(SessionMessageTable.session_id, event.data.parentID), eq(SessionMessageTable.id, event.data.from)),
)
.get()
.pipe(Effect.orDie)
: undefined
if (event.data.messageID && !boundary)
return yield* Effect.die(new Error(`Fork boundary message not found: ${event.data.messageID}`))
if (event.data.from && !boundary) return yield* Effect.die(new Error(`Fork boundary message not found: ${event.data.from}`))
const copied = yield* db
.select({ seq: SessionMessageTable.seq })
.from(SessionMessageTable)
@ -208,8 +204,8 @@ const projectFork = Effect.fn("SessionProjector.projectFork")(function* (
tokens_reasoning: 0,
tokens_cache_read: 0,
tokens_cache_write: 0,
time_created: DateTime.toEpochMillis(event.data.timestamp),
time_updated: DateTime.toEpochMillis(event.data.timestamp),
time_created: DateTime.toEpochMillis(event.created),
time_updated: DateTime.toEpochMillis(event.created),
})
.onConflictDoNothing()
.returning({ sessionID: SessionTable.id })
@ -474,9 +470,9 @@ const layer = Layer.effectDiscard(
.update(SessionTable)
.set({
directory: event.data.location.directory,
path: event.data.subdirectory,
path: event.data.subpath,
workspace_id: event.data.location.workspaceID ? WorkspaceV2.ID.make(event.data.location.workspaceID) : null,
time_updated: DateTime.toEpochMillis(event.data.timestamp),
time_updated: DateTime.toEpochMillis(event.created),
})
.where(eq(SessionTable.id, event.data.sessionID))
.run()
@ -556,19 +552,19 @@ const layer = Layer.effectDiscard(
if (next) yield* applyUsage(db, sessionID, next)
}),
)
yield* events.project(SessionEvent.AgentSwitched, (event) =>
yield* events.project(SessionEvent.AgentSelected, (event) =>
db
.update(SessionTable)
.set({ agent: event.data.agent, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
.set({ agent: event.data.agent, time_updated: DateTime.toEpochMillis(event.created) })
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie, Effect.andThen(run(db, event))),
)
yield* events.project(SessionEvent.ModelSwitched, (event) =>
yield* events.project(SessionEvent.ModelSelected, (event) =>
Effect.gen(function* () {
yield* db
.update(SessionTable)
.set({ model: event.data.model, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
.set({ model: event.data.model, time_updated: DateTime.toEpochMillis(event.created) })
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie)
@ -578,25 +574,30 @@ const layer = Layer.effectDiscard(
yield* events.project(SessionEvent.Renamed, (event) =>
db
.update(SessionTable)
.set({ title: event.data.title, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
.set({ title: event.data.title, time_updated: DateTime.toEpochMillis(event.created) })
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie),
)
yield* events.project(SessionEvent.Forked, (event) => projectFork(db, event))
yield* events.project(SessionEvent.Prompted, (event) =>
yield* events.project(SessionEvent.PromptPromoted, (event) =>
Effect.gen(function* () {
if (event.durable === undefined)
return yield* Effect.die(new Error("Durable Session event is missing aggregate sequence"))
yield* SessionInput.projectPrompted(db, {
id: event.data.messageID,
const input = yield* SessionInput.projectPromptPromoted(db, {
id: event.data.inputID,
sessionID: event.data.sessionID,
prompt: event.data.prompt,
delivery: event.data.delivery,
timeCreated: event.data.timestamp,
promotedSeq: event.durable.seq,
})
yield* run(db, event)
yield* insertMessage(db, event, {
id: input.id,
type: "user",
metadata: event.metadata,
text: input.prompt.text,
files: input.prompt.files,
agents: input.prompt.agents,
time: { created: event.created },
})
}),
)
yield* events.project(SessionEvent.PromptAdmitted, (event) =>
@ -605,11 +606,11 @@ const layer = Layer.effectDiscard(
return yield* Effect.die(new Error("Durable Session event is missing aggregate sequence"))
yield* SessionInput.projectAdmitted(db, {
admittedSeq: event.durable.seq,
id: event.data.messageID,
id: event.data.inputID,
sessionID: event.data.sessionID,
prompt: event.data.prompt,
delivery: event.data.delivery,
timeCreated: event.data.timestamp,
timeCreated: event.created,
})
}),
)
@ -617,11 +618,11 @@ const layer = Layer.effectDiscard(
yield* events.project(SessionEvent.Synthetic, (event) => run(db, event))
yield* events.project(SessionEvent.Skill.Activated, (event) =>
insertMessage(db, event, {
id: event.data.messageID,
id: SessionMessage.ID.fromEvent(event.id),
type: "skill",
name: event.data.name,
text: event.data.text,
time: { created: event.data.timestamp },
time: { created: event.created },
}),
)
yield* events.project(SessionEvent.Shell.Started, (event) => run(db, event))
@ -646,7 +647,7 @@ const layer = Layer.effectDiscard(
.update(SessionTable)
.set({
revert: { ...event.data.revert, files: event.data.revert.files ? [...event.data.revert.files] : undefined },
time_updated: DateTime.toEpochMillis(event.data.timestamp),
time_updated: DateTime.toEpochMillis(event.created),
})
.where(eq(SessionTable.id, event.data.sessionID))
.run()
@ -655,7 +656,7 @@ const layer = Layer.effectDiscard(
yield* events.project(SessionEvent.RevertEvent.Cleared, (event) =>
db
.update(SessionTable)
.set({ revert: null, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
.set({ revert: null, time_updated: DateTime.toEpochMillis(event.created) })
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie, Effect.asVoid),
@ -693,7 +694,7 @@ const layer = Layer.effectDiscard(
.pipe(Effect.orDie)
yield* db
.update(SessionTable)
.set({ revert: null, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
.set({ revert: null, time_updated: DateTime.toEpochMillis(event.created) })
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie)

View file

@ -89,7 +89,6 @@ export const stage = Effect.fn("SessionRevert.stage")(function* (input: {
} satisfies SessionSchema.Info["revert"]
yield* events.publish(SessionEvent.RevertEvent.Staged, {
sessionID: input.session.id,
timestamp: yield* DateTime.now,
revert,
})
return revert
@ -106,7 +105,6 @@ export const clear = Effect.fn("SessionRevert.clear")(function* (session: Sessio
const events = yield* EventV2.Service
yield* events.publish(SessionEvent.RevertEvent.Cleared, {
sessionID: session.id,
timestamp: yield* DateTime.now,
})
})
@ -116,6 +114,5 @@ export const commit = Effect.fn("SessionRevert.commit")(function* (session: Sess
yield* events.publish(SessionEvent.RevertEvent.Committed, {
sessionID: session.id,
messageID: session.revert.messageID,
timestamp: yield* DateTime.now,
})
})

View file

@ -134,7 +134,6 @@ const layer = Layer.effect(
if (tool.type !== "tool" || (tool.state.status !== "pending" && tool.state.status !== "running")) continue
yield* events.publish(SessionEvent.Tool.Failed, {
sessionID,
timestamp: yield* DateTime.now,
assistantMessageID: message.id,
callID: tool.id,
error: { type: "unknown", message: "Tool execution interrupted" },
@ -297,7 +296,6 @@ const layer = Layer.effect(
yield* serialized(
events.publish(SessionEvent.Step.Ended, {
sessionID: session.id,
timestamp: yield* DateTime.now,
assistantMessageID: yield* publisher.startAssistant(),
finish: settlement.finish,
cost: 0,

View file

@ -78,7 +78,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* events.publish(SessionEvent.Step.Started, {
...input,
assistantMessageID,
timestamp: yield* timestamp,
snapshot: input.snapshot,
})
return assistantMessageID
@ -123,7 +122,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* events.publish(SessionEvent.Text.Ended, {
sessionID: input.sessionID,
assistantMessageID: yield* currentAssistantMessageID(),
timestamp: yield* timestamp,
textID,
text: value,
})
@ -134,7 +132,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* events.publish(SessionEvent.Reasoning.Ended, {
sessionID: input.sessionID,
assistantMessageID: yield* currentAssistantMessageID(),
timestamp: yield* timestamp,
reasoningID,
text: value,
providerMetadata,
@ -147,7 +144,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
if (!tool) return yield* Effect.die(new Error(`Tool input end before start: ${callID}`))
yield* events.publish(SessionEvent.Tool.Input.Ended, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID: tool.assistantMessageID,
callID,
text: value,
@ -176,7 +172,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* toolInput.start(event.id)
yield* events.publish(SessionEvent.Tool.Input.Started, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID,
callID: event.id,
name: event.name,
@ -204,7 +199,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
assistantFailed = true
yield* events.publish(SessionEvent.Step.Failed, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID,
error: { type: "unknown", message },
})
@ -219,7 +213,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
tool.settled = true
yield* events.publish(SessionEvent.Tool.Failed, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID: tool.assistantMessageID,
callID,
error: { type: "unknown", message },
@ -248,7 +241,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* events.publish(SessionEvent.Text.Started, {
sessionID: input.sessionID,
assistantMessageID: yield* startAssistant(),
timestamp: yield* timestamp,
textID: event.id,
})
return
@ -257,7 +249,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* events.publish(SessionEvent.Text.Delta, {
sessionID: input.sessionID,
assistantMessageID: yield* currentAssistantMessageID(),
timestamp: yield* timestamp,
textID: event.id,
delta: event.text,
})
@ -270,7 +261,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* events.publish(SessionEvent.Reasoning.Started, {
sessionID: input.sessionID,
assistantMessageID: yield* startAssistant(),
timestamp: yield* timestamp,
reasoningID: event.id,
providerMetadata: event.providerMetadata,
})
@ -280,7 +270,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* events.publish(SessionEvent.Reasoning.Delta, {
sessionID: input.sessionID,
assistantMessageID: yield* currentAssistantMessageID(),
timestamp: yield* timestamp,
reasoningID: event.id,
delta: event.text,
})
@ -300,7 +289,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* toolInput.append(event.id, event.text)
yield* events.publish(SessionEvent.Tool.Input.Delta, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID: tool.assistantMessageID,
callID: event.id,
delta: event.text,
@ -322,7 +310,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
tool.providerMetadata = event.providerMetadata
yield* events.publish(SessionEvent.Tool.Called, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID: tool.assistantMessageID,
callID: event.id,
tool: event.name,
@ -352,7 +339,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
if ("error" in result) {
yield* events.publish(SessionEvent.Tool.Failed, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID: tool.assistantMessageID,
callID: event.id,
error: result.error,
@ -363,7 +349,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
}
yield* events.publish(SessionEvent.Tool.Success, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID: tool.assistantMessageID,
callID: event.id,
...result,
@ -382,7 +367,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
tool.settled = true
yield* events.publish(SessionEvent.Tool.Failed, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID: tool.assistantMessageID,
callID: event.id,
error: { type: "unknown", message: event.message },

View file

@ -42,9 +42,10 @@ const make = (dependencies: Dependencies) => {
if (!firstUser) return
const agent = yield* dependencies.agents.get(AgentV2.ID.make("title"))
if (!agent) return
const resolved = yield* (agent.model
? dependencies.models.resolve({ ...session, model: agent.model })
: dependencies.models.resolve(session)
const resolved = yield* (
agent.model
? dependencies.models.resolve({ ...session, model: agent.model })
: dependencies.models.resolve(session)
).pipe(Effect.catch(() => Effect.succeed(undefined)))
if (!resolved) return
const chunks: string[] = []
@ -76,7 +77,6 @@ const make = (dependencies: Dependencies) => {
if (!title) return
yield* dependencies.events.publish(SessionEvent.Renamed, {
sessionID: session.id,
timestamp: yield* DateTime.now,
title: truncate(title),
})
})