refactor(core): project session events in core

This commit is contained in:
Dax Raad
2026-05-24 21:45:01 -04:00
parent b18ff8fd68
commit dedcb9ba91
8 changed files with 475 additions and 362 deletions
+12 -14
View File
@@ -29,7 +29,8 @@ export type Payload<D extends Definition = Definition> = {
readonly metadata?: Record<string, unknown>
}
export type Sync = (event: Payload) => Effect.Effect<void>
export type Projector<D extends Definition = Definition> = (event: Payload<D>) => Effect.Effect<void>
type AnyProjector = (event: Payload) => Effect.Effect<void>
export const registry = new Map<string, Definition>()
@@ -69,8 +70,6 @@ export interface PublishOptions {
readonly metadata?: Record<string, unknown>
}
export type Unsubscribe = Effect.Effect<void>
export interface Interface {
readonly publish: <D extends Definition>(
definition: D,
@@ -80,7 +79,7 @@ export interface Interface {
readonly publishEvent: <D extends Definition>(event: Payload<D>) => Effect.Effect<Payload<D>>
readonly subscribe: <D extends Definition>(definition: D) => Stream.Stream<Payload<D>>
readonly all: () => Stream.Stream<Payload>
readonly sync: (handler: Sync) => Effect.Effect<Unsubscribe>
readonly project: <D extends Definition>(definition: D, projector: Projector<D>) => Effect.Effect<void>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/Event") {}
@@ -90,7 +89,7 @@ export const layer = Layer.effect(
Effect.gen(function* () {
const all = yield* PubSub.unbounded<Payload>()
const typed = new Map<string, PubSub.PubSub<Payload>>()
const syncHandlers = new Array<Sync>()
const projectors = new Map<string, AnyProjector[]>()
const getOrCreate = (definition: Definition) =>
Effect.gen(function* () {
@@ -110,8 +109,8 @@ export const layer = Layer.effect(
function publishEvent<D extends Definition>(event: Payload<D>) {
return Effect.gen(function* () {
for (const sync of syncHandlers) {
yield* sync(event as Payload)
for (const projector of projectors.get(event.type) ?? []) {
yield* projector(event as Payload)
}
const pubsub = typed.get(event.type)
if (pubsub) yield* PubSub.publish(pubsub, event as Payload)
@@ -141,16 +140,15 @@ export const layer = Layer.effect(
)
const streamAll = (): Stream.Stream<Payload> => Stream.fromPubSub(all)
const sync = (handler: Sync): Effect.Effect<Unsubscribe> =>
const project = <D extends Definition>(definition: D, projector: Projector<D>): Effect.Effect<void> =>
Effect.sync(() => {
syncHandlers.push(handler)
return Effect.sync(() => {
const index = syncHandlers.indexOf(handler)
if (index >= 0) syncHandlers.splice(index, 1)
})
const list = projectors.get(definition.type) ?? []
list.push((event) => projector(event as Payload<D>))
projectors.set(definition.type, list)
})
return Service.of({ publish, publishEvent, subscribe, all: streamAll, sync })
return Service.of({ publish, publishEvent, subscribe, all: streamAll, project })
}),
)
+151 -90
View File
@@ -1,4 +1,5 @@
import { produce, type WritableDraft } from "immer"
import { Effect } from "effect"
import { SessionEvent } from "./event"
import { SessionMessage } from "./message"
@@ -6,18 +7,17 @@ export type MemoryState = {
messages: SessionMessage.Message[]
}
export interface Adapter<Result> {
readonly getCurrentAssistant: () => SessionMessage.Assistant | undefined
readonly getCurrentCompaction: () => SessionMessage.Compaction | undefined
readonly getCurrentShell: (callID: string) => SessionMessage.Shell | undefined
readonly updateAssistant: (assistant: SessionMessage.Assistant) => void
readonly updateCompaction: (compaction: SessionMessage.Compaction) => void
readonly updateShell: (shell: SessionMessage.Shell) => void
readonly appendMessage: (message: SessionMessage.Message) => void
readonly finish: () => Result
export interface Adapter {
readonly getCurrentAssistant: () => Effect.Effect<SessionMessage.Assistant | undefined>
readonly getCurrentCompaction: () => Effect.Effect<SessionMessage.Compaction | undefined>
readonly getCurrentShell: (callID: string) => Effect.Effect<SessionMessage.Shell | undefined>
readonly updateAssistant: (assistant: SessionMessage.Assistant) => Effect.Effect<void>
readonly updateCompaction: (compaction: SessionMessage.Compaction) => Effect.Effect<void>
readonly updateShell: (shell: SessionMessage.Shell) => Effect.Effect<void>
readonly appendMessage: (message: SessionMessage.Message) => Effect.Effect<void>
}
export function memory(state: MemoryState): Adapter<MemoryState> {
export function memory(state: MemoryState): Adapter {
const activeAssistantIndex = () =>
state.messages.findLastIndex((message) => message.type === "assistant" && !message.time.completed)
const activeCompactionIndex = () => state.messages.findLastIndex((message) => message.type === "compaction")
@@ -26,55 +26,65 @@ export function memory(state: MemoryState): Adapter<MemoryState> {
return {
getCurrentAssistant() {
const index = activeAssistantIndex()
if (index < 0) return
const assistant = state.messages[index]
return assistant?.type === "assistant" ? assistant : undefined
return Effect.sync(() => {
const index = activeAssistantIndex()
if (index < 0) return
const assistant = state.messages[index]
return assistant?.type === "assistant" ? assistant : undefined
})
},
getCurrentCompaction() {
const index = activeCompactionIndex()
if (index < 0) return
const compaction = state.messages[index]
return compaction?.type === "compaction" ? compaction : undefined
return Effect.sync(() => {
const index = activeCompactionIndex()
if (index < 0) return
const compaction = state.messages[index]
return compaction?.type === "compaction" ? compaction : undefined
})
},
getCurrentShell(callID) {
const index = activeShellIndex(callID)
if (index < 0) return
const shell = state.messages[index]
return shell?.type === "shell" ? shell : undefined
return Effect.sync(() => {
const index = activeShellIndex(callID)
if (index < 0) return
const shell = state.messages[index]
return shell?.type === "shell" ? shell : undefined
})
},
updateAssistant(assistant) {
const index = activeAssistantIndex()
if (index < 0) return
const current = state.messages[index]
if (current?.type !== "assistant") return
state.messages[index] = assistant
return Effect.sync(() => {
const index = activeAssistantIndex()
if (index < 0) return
const current = state.messages[index]
if (current?.type !== "assistant") return
state.messages[index] = assistant
})
},
updateCompaction(compaction) {
const index = activeCompactionIndex()
if (index < 0) return
const current = state.messages[index]
if (current?.type !== "compaction") return
state.messages[index] = compaction
return Effect.sync(() => {
const index = activeCompactionIndex()
if (index < 0) return
const current = state.messages[index]
if (current?.type !== "compaction") return
state.messages[index] = compaction
})
},
updateShell(shell) {
const index = activeShellIndex(shell.callID)
if (index < 0) return
const current = state.messages[index]
if (current?.type !== "shell") return
state.messages[index] = shell
return Effect.sync(() => {
const index = activeShellIndex(shell.callID)
if (index < 0) return
const current = state.messages[index]
if (current?.type !== "shell") return
state.messages[index] = shell
})
},
appendMessage(message) {
state.messages.push(message)
},
finish() {
return state
return Effect.sync(() => {
state.messages.push(message)
})
},
}
}
export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Event): Result {
const currentAssistant = adapter.getCurrentAssistant()
export function update(adapter: Adapter, event: SessionEvent.Event) {
type DraftAssistant = WritableDraft<SessionMessage.Assistant>
type DraftTool = WritableDraft<SessionMessage.AssistantTool>
type DraftText = WritableDraft<SessionMessage.AssistantText>
@@ -91,9 +101,10 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
const latestReasoning = (assistant: DraftAssistant | undefined, reasoningID: string) =>
assistant?.content.findLast((item): item is DraftReasoning => item.type === "reasoning" && item.id === reasoningID)
SessionEvent.All.match(event, {
return Effect.gen(function* () {
yield* SessionEvent.All.match(event, {
"session.next.agent.switched": (event) => {
adapter.appendMessage(
return adapter.appendMessage(
new SessionMessage.AgentSwitched({
id: event.id,
type: "agent-switched",
@@ -104,7 +115,7 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
)
},
"session.next.model.switched": (event) => {
adapter.appendMessage(
return adapter.appendMessage(
new SessionMessage.ModelSwitched({
id: event.id,
type: "model-switched",
@@ -115,7 +126,7 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
)
},
"session.next.prompted": (event) => {
adapter.appendMessage(
return adapter.appendMessage(
new SessionMessage.User({
id: event.id,
type: "user",
@@ -129,7 +140,7 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
)
},
"session.next.synthetic": (event) => {
adapter.appendMessage(
return adapter.appendMessage(
new SessionMessage.Synthetic({
sessionID: event.data.sessionID,
text: event.data.text,
@@ -140,7 +151,7 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
)
},
"session.next.shell.started": (event) => {
adapter.appendMessage(
return adapter.appendMessage(
new SessionMessage.Shell({
id: event.id,
type: "shell",
@@ -153,39 +164,46 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
)
},
"session.next.shell.ended": (event) => {
const currentShell = adapter.getCurrentShell(event.data.callID)
return Effect.gen(function* () {
const currentShell = yield* adapter.getCurrentShell(event.data.callID)
if (currentShell) {
adapter.updateShell(
yield* adapter.updateShell(
produce(currentShell, (draft) => {
draft.output = event.data.output
draft.time.completed = event.data.timestamp
}),
)
}
})
},
"session.next.step.started": (event) => {
if (currentAssistant) {
adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.time.completed = event.data.timestamp
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.time.completed = event.data.timestamp
}),
)
}
yield* adapter.appendMessage(
new SessionMessage.Assistant({
id: event.id,
type: "assistant",
agent: event.data.agent,
model: event.data.model,
time: { created: event.data.timestamp },
content: [],
snapshot: event.data.snapshot ? { start: event.data.snapshot } : undefined,
}),
)
}
adapter.appendMessage(
new SessionMessage.Assistant({
id: event.id,
type: "assistant",
agent: event.data.agent,
model: event.data.model,
time: { created: event.data.timestamp },
content: [],
snapshot: event.data.snapshot ? { start: event.data.snapshot } : undefined,
}),
)
})
},
"session.next.step.ended": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.time.completed = event.data.timestamp
draft.finish = event.data.finish
@@ -195,10 +213,13 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.step.failed": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.time.completed = event.data.timestamp
draft.finish = "error"
@@ -206,10 +227,13 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.text.started": () => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.content.push({
type: "text",
@@ -218,30 +242,39 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.text.delta": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestText(draft)
if (match) match.text += event.data.delta
}),
)
}
})
},
"session.next.text.ended": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestText(draft)
if (match) match.text = event.data.text
}),
)
}
})
},
"session.next.tool.input.started": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.content.push({
type: "tool",
@@ -258,10 +291,13 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.tool.input.delta": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestTool(draft, event.data.callID)
// oxlint-disable-next-line no-base-to-string -- event.delta is a Schema.String (runtime string)
@@ -269,11 +305,14 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.tool.input.ended": () => {},
"session.next.tool.input.ended": () => Effect.void,
"session.next.tool.called": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match) {
@@ -289,10 +328,13 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.tool.progress": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "running") {
@@ -302,10 +344,13 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.tool.success": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "running") {
@@ -321,10 +366,13 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.tool.failed": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestTool(draft, event.data.callID)
if (match && match.state.status === "running") {
@@ -341,10 +389,13 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.reasoning.started": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
draft.content.push({
type: "reasoning",
@@ -354,30 +405,37 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
}),
)
}
})
},
"session.next.reasoning.delta": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestReasoning(draft, event.data.reasoningID)
if (match) match.text += event.data.delta
}),
)
}
})
},
"session.next.reasoning.ended": (event) => {
return Effect.gen(function* () {
const currentAssistant = yield* adapter.getCurrentAssistant()
if (currentAssistant) {
adapter.updateAssistant(
yield* adapter.updateAssistant(
produce(currentAssistant, (draft) => {
const match = latestReasoning(draft, event.data.reasoningID)
if (match) match.text = event.data.text
}),
)
}
})
},
"session.next.retried": () => {},
"session.next.retried": () => Effect.void,
"session.next.compaction.started": (event) => {
adapter.appendMessage(
return adapter.appendMessage(
new SessionMessage.Compaction({
id: event.id,
type: "compaction",
@@ -389,29 +447,32 @@ export function update<Result>(adapter: Adapter<Result>, event: SessionEvent.Eve
)
},
"session.next.compaction.delta": (event) => {
const currentCompaction = adapter.getCurrentCompaction()
return Effect.gen(function* () {
const currentCompaction = yield* adapter.getCurrentCompaction()
if (currentCompaction) {
adapter.updateCompaction(
yield* adapter.updateCompaction(
produce(currentCompaction, (draft) => {
draft.summary += event.data.text
}),
)
}
})
},
"session.next.compaction.ended": (event) => {
const currentCompaction = adapter.getCurrentCompaction()
return Effect.gen(function* () {
const currentCompaction = yield* adapter.getCurrentCompaction()
if (currentCompaction) {
adapter.updateCompaction(
yield* adapter.updateCompaction(
produce(currentCompaction, (draft) => {
draft.summary = event.data.text
draft.include = event.data.include
}),
)
}
})
},
})
return adapter.finish()
})
}
export * as SessionMessageUpdater from "./message-updater"
+266
View File
@@ -0,0 +1,266 @@
export * as SessionProjector from "./projector"
import { and, eq } from "drizzle-orm"
import { DateTime, Effect, Layer, Schema } from "effect"
import { Database } from "../database/database"
import { EventV2 } from "../event"
import { SessionEvent } from "./event"
import { SessionMessage } from "./message"
import { SessionMessageUpdater } from "./message-updater"
import { SessionMessageTable, SessionTable } from "./sql"
type DatabaseService = Database.Interface["db"]
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Message)
const encodeMessage = Schema.encodeSync(SessionMessage.Message)
function run(db: DatabaseService, event: SessionEvent.Event) {
return Effect.gen(function* () {
const adapter: SessionMessageUpdater.Adapter = {
getCurrentAssistant() {
return Effect.gen(function* () {
const rows = yield* db
.select()
.from(SessionMessageTable)
.where(and(eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.type, "assistant")))
.all()
.pipe(Effect.orDie)
return rows
.map((row) => decodeMessage({ ...row.data, id: row.id, type: row.type }))
.find(
(message): message is SessionMessage.Assistant => message.type === "assistant" && !message.time.completed,
)
})
},
getCurrentCompaction() {
return Effect.gen(function* () {
const rows = yield* db
.select()
.from(SessionMessageTable)
.where(and(eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.type, "compaction")))
.all()
.pipe(Effect.orDie)
return rows
.map((row) => decodeMessage({ ...row.data, id: row.id, type: row.type }))
.find((message): message is SessionMessage.Compaction => message.type === "compaction")
})
},
getCurrentShell(callID) {
return Effect.gen(function* () {
const rows = yield* db
.select()
.from(SessionMessageTable)
.where(and(eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.type, "shell")))
.all()
.pipe(Effect.orDie)
return rows
.map((row) => decodeMessage({ ...row.data, id: row.id, type: row.type }))
.find((message): message is SessionMessage.Shell => message.type === "shell" && message.callID === callID)
})
},
updateAssistant(message) {
return Effect.gen(function* () {
const encoded = encodeMessage(message)
const { id, type, ...data } = encoded
yield* db
.insert(SessionMessageTable)
.values([
{
id: SessionMessage.ID.make(id),
session_id: event.data.sessionID,
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
])
.onConflictDoUpdate({
target: SessionMessageTable.id,
set: {
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
})
.run()
.pipe(Effect.orDie)
})
},
updateCompaction(message) {
return Effect.gen(function* () {
const encoded = encodeMessage(message)
const { id, type, ...data } = encoded
yield* db
.insert(SessionMessageTable)
.values([
{
id: SessionMessage.ID.make(id),
session_id: event.data.sessionID,
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
])
.onConflictDoUpdate({
target: SessionMessageTable.id,
set: {
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
})
.run()
.pipe(Effect.orDie)
})
},
updateShell(message) {
return Effect.gen(function* () {
const encoded = encodeMessage(message)
const { id, type, ...data } = encoded
yield* db
.insert(SessionMessageTable)
.values([
{
id: SessionMessage.ID.make(id),
session_id: event.data.sessionID,
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
])
.onConflictDoUpdate({
target: SessionMessageTable.id,
set: {
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
})
.run()
.pipe(Effect.orDie)
})
},
appendMessage(message) {
return Effect.gen(function* () {
const encoded = encodeMessage(message)
const { id, type, ...data } = encoded
yield* db
.insert(SessionMessageTable)
.values([
{
id: SessionMessage.ID.make(id),
session_id: event.data.sessionID,
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
])
.onConflictDoUpdate({
target: SessionMessageTable.id,
set: {
type,
time_created: DateTime.toEpochMillis(message.time.created),
data,
},
})
.run()
.pipe(Effect.orDie)
})
},
}
yield* SessionMessageUpdater.update(adapter, event)
})
}
export const layer = Layer.effectDiscard(
Effect.gen(function* () {
const events = yield* EventV2.Service
const database = yield* Database.Service
yield* events.project(SessionEvent.AgentSwitched, (event) =>
Effect.gen(function* () {
const message = Schema.encodeSync(SessionMessage.AgentSwitched)(
new SessionMessage.AgentSwitched({
id: event.id,
type: "agent-switched",
metadata: event.metadata,
agent: event.data.agent,
time: { created: event.data.timestamp },
}),
)
const data = { metadata: message.metadata, agent: message.agent, time: message.time }
yield* database.db
.update(SessionTable)
.set({ agent: event.data.agent, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie)
yield* database.db
.insert(SessionMessageTable)
.values([
{
id: SessionMessage.ID.make(event.id),
session_id: event.data.sessionID,
type: "agent-switched",
time_created: DateTime.toEpochMillis(event.data.timestamp),
data,
},
])
.run()
.pipe(Effect.orDie)
}),
)
yield* events.project(SessionEvent.ModelSwitched, (event) =>
Effect.gen(function* () {
const message = Schema.encodeSync(SessionMessage.ModelSwitched)(
new SessionMessage.ModelSwitched({
id: event.id,
type: "model-switched",
metadata: event.metadata,
model: event.data.model,
time: { created: event.data.timestamp },
}),
)
const data = { metadata: message.metadata, model: message.model, time: message.time }
yield* database.db
.update(SessionTable)
.set({ model: event.data.model, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
.where(eq(SessionTable.id, event.data.sessionID))
.run()
.pipe(Effect.orDie)
yield* database.db
.insert(SessionMessageTable)
.values([
{
id: SessionMessage.ID.make(event.id),
session_id: event.data.sessionID,
type: "model-switched",
time_created: DateTime.toEpochMillis(event.data.timestamp),
data,
},
])
.run()
.pipe(Effect.orDie)
}),
)
yield* events.project(SessionEvent.Prompted, (event) => run(database.db, event))
yield* events.project(SessionEvent.Synthetic, (event) => run(database.db, event))
yield* events.project(SessionEvent.Shell.Started, (event) => run(database.db, event))
yield* events.project(SessionEvent.Shell.Ended, (event) => run(database.db, event))
yield* events.project(SessionEvent.Step.Started, (event) => run(database.db, event))
yield* events.project(SessionEvent.Step.Ended, (event) => run(database.db, event))
yield* events.project(SessionEvent.Step.Failed, (event) => run(database.db, event))
yield* events.project(SessionEvent.Text.Started, (event) => run(database.db, event))
yield* events.project(SessionEvent.Text.Ended, (event) => run(database.db, event))
yield* events.project(SessionEvent.Tool.Input.Started, (event) => run(database.db, event))
yield* events.project(SessionEvent.Tool.Input.Ended, (event) => run(database.db, event))
yield* events.project(SessionEvent.Tool.Called, (event) => run(database.db, event))
yield* events.project(SessionEvent.Tool.Success, (event) => run(database.db, event))
yield* events.project(SessionEvent.Tool.Failed, (event) => run(database.db, event))
yield* events.project(SessionEvent.Reasoning.Started, (event) => run(database.db, event))
yield* events.project(SessionEvent.Reasoning.Ended, (event) => run(database.db, event))
yield* events.project(SessionEvent.Retried, (event) => run(database.db, event))
yield* events.project(SessionEvent.Compaction.Started, (event) => run(database.db, event))
yield* events.project(SessionEvent.Compaction.Ended, (event) => run(database.db, event))
}),
)
export const defaultLayer = layer.pipe(Layer.provide(EventV2.defaultLayer), Layer.provide(Database.defaultLayer))
+5 -6
View File
@@ -90,25 +90,24 @@ describe("EventV2", () => {
}),
)
it.effect("runs sync handlers inline", () =>
it.effect("runs projectors inline", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
const unsubscribe = yield* events.sync((event) =>
yield* events.project(Message, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
const event = yield* events.publish(Message, { text: "hello" })
yield* unsubscribe
yield* events.publish(Message, { text: "after unsubscribe" })
expect(received).toEqual([event])
expect(received).toEqual([event, expect.objectContaining({ data: { text: "after unsubscribe" } })])
}),
)
it.effect("runs sync handlers before publishing to streams", () =>
it.effect("runs projectors before publishing to streams", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<string>()
@@ -117,7 +116,7 @@ describe("EventV2", () => {
Stream.runForEach(() => Effect.sync(() => received.push("stream"))),
Effect.forkScoped,
)
yield* events.sync((event) =>
yield* events.project(Message, (event) =>
Effect.sync(() => {
received.push(event.type)
}),
+14 -19
View File
@@ -7,10 +7,11 @@ import { InstanceRef, WorkspaceRef } from "@/effect/instance-ref"
import { InstanceStore } from "@/project/instance-store"
import { SyncEvent } from "@/sync"
import { EventV2 } from "@opencode-ai/core/event"
import { SessionProjector } from "@opencode-ai/core/session/projector"
import "@opencode-ai/core/account"
import "@opencode-ai/core/catalog"
import "@opencode-ai/core/session/event"
import { Context, Effect, Layer, Option } from "effect"
import { Context, Effect, Layer, Option, Stream } from "effect"
export function toSyncDefinition<D extends EventV2.Definition>(definition: D) {
const result = {
@@ -30,7 +31,6 @@ export const layer = Layer.effect(
Effect.gen(function* () {
const events = yield* EventV2.Service
const bus = yield* ProjectBus.Service
const sync = yield* SyncEvent.Service
const publishGlobal = (event: EventV2.Payload) =>
Effect.sync(() => {
@@ -60,28 +60,23 @@ export const layer = Layer.effect(
})
}
const unsubscribe = yield* events.sync((event) => {
const definition = EventV2.registry.get(event.type)
if (!definition) return Effect.void
const aggregateID = definition.aggregate
? (event.data as Record<string, unknown>)[definition.aggregate]
: undefined
if (definition.version !== undefined && typeof aggregateID === "string") {
return provideEventLocation(event, sync.run(toSyncDefinition(definition), event.data))
}
return provideEventLocation(
event,
bus.publish({ type: definition.type, properties: definition.data }, event.data, { id: event.id }),
)
})
yield* Effect.addFinalizer(() => unsubscribe)
yield* events.all().pipe(
Stream.runForEach((event) => {
const definition = EventV2.registry.get(event.type)
if (!definition) return Effect.void
return provideEventLocation(
event,
bus.publish({ type: definition.type, properties: definition.data }, event.data, { id: event.id }),
)
}),
Effect.forkScoped,
)
return Service.of(events)
}),
)
export const defaultLayer = layer.pipe(
Layer.provideMerge(SessionProjector.defaultLayer),
Layer.provide(EventV2.defaultLayer),
Layer.provide(SyncEvent.defaultLayer),
Layer.provide(ProjectBus.defaultLayer),
@@ -1,204 +0,0 @@
import { and, desc, eq } from "@/storage/db"
import type { Database } from "@/storage/db"
import { SessionMessage } from "@opencode-ai/core/session/message"
import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater"
import { SessionEvent } from "@opencode-ai/core/session/event"
import * as DateTime from "effect/DateTime"
import { SyncEvent } from "@/sync"
import { EventV2Bridge } from "@/event-v2-bridge"
import { SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql"
import type { SessionID } from "./schema"
import { Schema } from "effect"
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Message)
type SessionMessageData = NonNullable<(typeof SessionMessageTable.$inferInsert)["data"]>
function encodeDateTimes(value: unknown): unknown {
if (DateTime.isDateTime(value)) return DateTime.toEpochMillis(value)
if (Array.isArray(value)) return value.map(encodeDateTimes)
if (typeof value === "object" && value !== null) {
return Object.fromEntries(Object.entries(value).map(([key, item]) => [key, encodeDateTimes(item)]))
}
return value
}
function encodeMessageData(value: unknown): SessionMessageData {
return encodeDateTimes(value) as SessionMessageData
}
function sqlite(db: Database.TxOrDb, sessionID: SessionID): SessionMessageUpdater.Adapter<void> {
return {
getCurrentAssistant() {
return db
.select()
.from(SessionMessageTable)
.where(and(eq(SessionMessageTable.session_id, sessionID), eq(SessionMessageTable.type, "assistant")))
.orderBy(desc(SessionMessageTable.id))
.all()
.map((row) => decodeMessage({ ...(row.data as object), id: row.id, type: row.type }))
.find((message): message is SessionMessage.Assistant => message.type === "assistant" && !message.time.completed)
},
getCurrentCompaction() {
return db
.select()
.from(SessionMessageTable)
.where(and(eq(SessionMessageTable.session_id, sessionID), eq(SessionMessageTable.type, "compaction")))
.orderBy(desc(SessionMessageTable.id))
.all()
.map((row) => decodeMessage({ ...(row.data as object), id: row.id, type: row.type }))
.find((message): message is SessionMessage.Compaction => message.type === "compaction")
},
getCurrentShell(callID) {
return db
.select()
.from(SessionMessageTable)
.where(and(eq(SessionMessageTable.session_id, sessionID), eq(SessionMessageTable.type, "shell")))
.orderBy(desc(SessionMessageTable.id))
.all()
.map((row) => decodeMessage({ ...(row.data as object), id: row.id, type: row.type }))
.find((message): message is SessionMessage.Shell => message.type === "shell" && message.callID === callID)
},
updateAssistant(assistant) {
const { id, type, ...data } = assistant
db.update(SessionMessageTable)
.set({ data: encodeMessageData(data) })
.where(
and(
eq(SessionMessageTable.id, id),
eq(SessionMessageTable.session_id, sessionID),
eq(SessionMessageTable.type, type),
),
)
.run()
},
updateCompaction(compaction) {
const { id, type, ...data } = compaction
db.update(SessionMessageTable)
.set({ data: encodeMessageData(data) })
.where(
and(
eq(SessionMessageTable.id, id),
eq(SessionMessageTable.session_id, sessionID),
eq(SessionMessageTable.type, type),
),
)
.run()
},
updateShell(shell) {
const { id, type, ...data } = shell
db.update(SessionMessageTable)
.set({ data: encodeMessageData(data) })
.where(
and(
eq(SessionMessageTable.id, id),
eq(SessionMessageTable.session_id, sessionID),
eq(SessionMessageTable.type, type),
),
)
.run()
},
appendMessage(message) {
const { id, type, ...data } = message
db.insert(SessionMessageTable)
.values([
{
id,
session_id: sessionID,
type,
time_created: DateTime.toEpochMillis(message.time.created),
data: encodeMessageData(data),
},
])
.run()
},
finish() {},
}
}
function update(db: Database.TxOrDb, event: SessionEvent.Event) {
SessionMessageUpdater.update(sqlite(db, event.data.sessionID), event)
}
export default [
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.AgentSwitched), (db, data, event) => {
db.update(SessionTable)
.set({
agent: data.agent,
time_updated: DateTime.toEpochMillis(data.timestamp),
})
.where(eq(SessionTable.id, data.sessionID))
.run()
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.agent.switched", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.ModelSwitched), (db, data, event) => {
db.update(SessionTable)
.set({
model: data.model,
time_updated: DateTime.toEpochMillis(data.timestamp),
})
.where(eq(SessionTable.id, data.sessionID))
.run()
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.model.switched", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Prompted), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.prompted", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Synthetic), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.synthetic", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Shell.Started), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.shell.started", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Shell.Ended), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.shell.ended", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Step.Started), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.step.started", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Step.Ended), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.step.ended", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Step.Failed), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.step.failed", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Text.Started), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.text.started", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Text.Delta), () => {}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Text.Ended), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.text.ended", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Tool.Input.Started), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.tool.input.started", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Tool.Input.Delta), () => {}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Tool.Input.Ended), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.tool.input.ended", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Tool.Called), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.tool.called", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Tool.Success), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.tool.success", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Tool.Failed), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.tool.failed", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Reasoning.Started), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.reasoning.started", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Reasoning.Delta), () => {}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Reasoning.Ended), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.reasoning.ended", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Retried), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.retried", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Compaction.Started), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.compaction.started", data })
}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Compaction.Delta), () => {}),
SyncEvent.project(EventV2Bridge.toSyncDefinition(SessionEvent.Compaction.Ended), (db, data, event) => {
update(db, { id: SessionMessage.ID.make(event.id), type: "session.next.compaction.ended", data })
}),
]
@@ -10,7 +10,6 @@ import { MessageV2 } from "./message-v2"
import { SessionTable, MessageTable, PartTable } from "@opencode-ai/core/session/sql"
import { WorkspaceTable } from "@opencode-ai/core/control-plane/workspace.sql"
import { Log } from "@opencode-ai/core/util/log"
import nextProjectors from "./projectors-next"
const log = Log.create({ service: "session.projector" })
@@ -197,6 +196,4 @@ export default [
log.warn("ignored late part update", { partID: id, messageID, sessionID })
}
}),
...nextProjectors,
]
@@ -1,4 +1,5 @@
import { expect, test } from "bun:test"
import { Effect } from "effect"
import * as DateTime from "effect/DateTime"
import { SessionID } from "../../src/session/schema"
import { EventV2 } from "@opencode-ai/core/event"
@@ -11,7 +12,7 @@ test("step snapshots carry over to assistant messages", () => {
const state: SessionMessageUpdater.MemoryState = { messages: [] }
const sessionID = SessionID.make("session")
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.step.started",
data: {
@@ -25,9 +26,9 @@ test("step snapshots carry over to assistant messages", () => {
},
snapshot: "before",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.step.ended",
data: {
@@ -43,7 +44,7 @@ test("step snapshots carry over to assistant messages", () => {
},
snapshot: "after",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
expect(state.messages[0]?.type).toBe("assistant")
if (state.messages[0]?.type !== "assistant") return
@@ -55,7 +56,7 @@ test("text ended populates assistant text content", () => {
const state: SessionMessageUpdater.MemoryState = { messages: [] }
const sessionID = SessionID.make("session")
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.step.started",
data: {
@@ -68,18 +69,18 @@ test("text ended populates assistant text content", () => {
variant: ModelV2.VariantID.make("default"),
},
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.text.started",
data: {
sessionID,
timestamp: DateTime.makeUnsafe(2),
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.text.ended",
data: {
@@ -87,7 +88,7 @@ test("text ended populates assistant text content", () => {
timestamp: DateTime.makeUnsafe(3),
text: "hello assistant",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
expect(state.messages[0]?.type).toBe("assistant")
if (state.messages[0]?.type !== "assistant") return
@@ -99,7 +100,7 @@ test("tool completion stores completed timestamp", () => {
const sessionID = SessionID.make("session")
const callID = "call"
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.step.started",
data: {
@@ -112,9 +113,9 @@ test("tool completion stores completed timestamp", () => {
variant: ModelV2.VariantID.make("default"),
},
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.tool.input.started",
data: {
@@ -123,9 +124,9 @@ test("tool completion stores completed timestamp", () => {
callID,
name: "bash",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.tool.called",
data: {
@@ -136,9 +137,9 @@ test("tool completion stores completed timestamp", () => {
input: { command: "pwd" },
provider: { executed: true, metadata: { source: "provider" } },
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.tool.success",
data: {
@@ -149,7 +150,7 @@ test("tool completion stores completed timestamp", () => {
content: [{ type: "text", text: "/tmp" }],
provider: { executed: true, metadata: { status: "done" } },
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
expect(state.messages[0]?.type).toBe("assistant")
if (state.messages[0]?.type !== "assistant") return
@@ -164,7 +165,7 @@ test("compaction events reduce to compaction message", () => {
const sessionID = SessionID.make("session")
const id = EventV2.ID.create()
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id,
type: "session.next.compaction.started",
data: {
@@ -172,9 +173,9 @@ test("compaction events reduce to compaction message", () => {
timestamp: DateTime.makeUnsafe(1),
reason: "auto",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.compaction.delta",
data: {
@@ -182,9 +183,9 @@ test("compaction events reduce to compaction message", () => {
timestamp: DateTime.makeUnsafe(2),
text: "hello ",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.compaction.delta",
data: {
@@ -192,9 +193,9 @@ test("compaction events reduce to compaction message", () => {
timestamp: DateTime.makeUnsafe(3),
text: "summary",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
Effect.runSync(SessionMessageUpdater.update(SessionMessageUpdater.memory(state), {
id: EventV2.ID.create(),
type: "session.next.compaction.ended",
data: {
@@ -203,7 +204,7 @@ test("compaction events reduce to compaction message", () => {
text: "final summary",
include: "recent context",
},
} satisfies SessionEvent.Event)
} satisfies SessionEvent.Event))
expect(state.messages).toHaveLength(1)
expect(state.messages[0]).toMatchObject({