703 lines
26 KiB
TypeScript
703 lines
26 KiB
TypeScript
export * as SessionProjector from "./projector"
|
|
|
|
import { and, asc, desc, eq, gt, inArray, lt, or, sql } from "drizzle-orm"
|
|
import { DateTime, Effect, Layer, Schema } from "effect"
|
|
import { Database } from "../database/database"
|
|
import { EventV2 } from "../event"
|
|
import { makeGlobalNode } from "../effect/app-node"
|
|
import { SessionEvent } from "./event"
|
|
import { SessionV1 } from "../v1/session"
|
|
import { WorkspaceTable } from "../control-plane/workspace.sql"
|
|
import { SessionMessage } from "./message"
|
|
import { SessionMessageUpdater } from "./message-updater"
|
|
import { SessionInput } from "./input"
|
|
import { WorkspaceV2 } from "../workspace"
|
|
import { SessionContextCheckpoint } from "./context-checkpoint"
|
|
import {
|
|
MessageTable,
|
|
PartTable,
|
|
SessionContextCheckpointTable,
|
|
SessionInputTable,
|
|
SessionMessageTable,
|
|
SessionTable,
|
|
} from "./sql"
|
|
import type { DeepMutable } from "../schema"
|
|
import { Slug } from "../util/slug"
|
|
|
|
type DatabaseService = Database.Interface["db"]
|
|
type MessageEvent = Exclude<SessionEvent.Event, typeof SessionEvent.Forked.Type>
|
|
|
|
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Message)
|
|
const encodeMessage = Schema.encodeSync(SessionMessage.Message)
|
|
|
|
export class SessionAlreadyProjected extends Error {}
|
|
|
|
type Usage = {
|
|
cost: number
|
|
tokens: {
|
|
input: number
|
|
output: number
|
|
reasoning: number
|
|
cache: { read: number; write: number }
|
|
}
|
|
}
|
|
|
|
const ForkBatchSize = 500
|
|
|
|
const emptyUsage = (): Usage => ({
|
|
cost: 0,
|
|
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
|
|
})
|
|
|
|
const forkTitle = (value: string) => {
|
|
const match = value.match(/^(.+) \(fork #(\d+)\)$/)
|
|
if (match) return `${match[1]} (fork #${Number.parseInt(match[2], 10) + 1})`
|
|
return `${value} (fork #1)`
|
|
}
|
|
|
|
function usage(part: (typeof SessionV1.Event.PartUpdated.Type)["data"]["part"] | unknown): Usage | undefined {
|
|
if (typeof part !== "object" || part === null) return undefined
|
|
const value = part as Record<string, unknown>
|
|
if (value.type !== "step-finish") return undefined
|
|
if (!("cost" in value) || !("tokens" in value)) return undefined
|
|
return { cost: value.cost as Usage["cost"], tokens: value.tokens as Usage["tokens"] }
|
|
}
|
|
|
|
function addUsage(target: Usage, value: Usage) {
|
|
target.cost += value.cost
|
|
target.tokens.input += value.tokens.input
|
|
target.tokens.output += value.tokens.output
|
|
target.tokens.reasoning += value.tokens.reasoning
|
|
target.tokens.cache.read += value.tokens.cache.read
|
|
target.tokens.cache.write += value.tokens.cache.write
|
|
}
|
|
|
|
function messageUsage(row: typeof SessionMessageTable.$inferSelect): Usage | undefined {
|
|
if (row.type !== "assistant") return undefined
|
|
const message = decodeMessage({ ...row.data, id: row.id, type: row.type })
|
|
if (message.type !== "assistant" || message.cost === undefined || message.tokens === undefined) return undefined
|
|
return { cost: message.cost, tokens: message.tokens }
|
|
}
|
|
|
|
function sessionRow(info: SessionV1.SessionInfo): typeof SessionTable.$inferInsert {
|
|
return {
|
|
id: info.id,
|
|
project_id: info.projectID,
|
|
workspace_id: info.workspaceID ?? null,
|
|
parent_id: info.parentID,
|
|
slug: info.slug,
|
|
directory: info.directory,
|
|
path: info.path,
|
|
title: info.title,
|
|
agent: info.agent,
|
|
model: info.model,
|
|
version: info.version,
|
|
share_url: info.share?.url,
|
|
summary_additions: info.summary?.additions,
|
|
summary_deletions: info.summary?.deletions,
|
|
summary_files: info.summary?.files,
|
|
summary_diffs: info.summary?.diffs ? [...info.summary.diffs] : undefined,
|
|
metadata: info.metadata,
|
|
cost: info.cost ?? 0,
|
|
tokens_input: (info.tokens ?? { input: 0 }).input,
|
|
tokens_output: (info.tokens ?? { output: 0 }).output,
|
|
tokens_reasoning: (info.tokens ?? { reasoning: 0 }).reasoning,
|
|
tokens_cache_read: (info.tokens ?? { cache: { read: 0 } }).cache.read,
|
|
tokens_cache_write: (info.tokens ?? { cache: { write: 0 } }).cache.write,
|
|
revert: info.revert ? { ...info.revert, messageID: SessionMessage.ID.make(info.revert.messageID) } : null,
|
|
permission: info.permission ? [...info.permission] : undefined,
|
|
time_created: info.time.created,
|
|
time_updated: info.time.updated,
|
|
time_compacting: info.time.compacting,
|
|
time_archived: info.time.archived,
|
|
}
|
|
}
|
|
|
|
function messageData(
|
|
info: (typeof SessionV1.Event.MessageUpdated.Type)["data"]["info"],
|
|
): typeof MessageTable.$inferInsert.data {
|
|
const { id: _, sessionID: __, ...rest } = info
|
|
return rest as DeepMutable<typeof rest>
|
|
}
|
|
|
|
function partData(part: (typeof SessionV1.Event.PartUpdated.Type)["data"]["part"]): typeof PartTable.$inferInsert.data {
|
|
const { id: _, messageID: __, sessionID: ___, ...rest } = part
|
|
return rest as DeepMutable<typeof rest>
|
|
}
|
|
|
|
function applyUsage(
|
|
db: DatabaseService,
|
|
sessionID: (typeof SessionV1.Event.MessageUpdated.Type)["data"]["sessionID"],
|
|
value: Usage,
|
|
sign = 1,
|
|
) {
|
|
return db
|
|
.update(SessionTable)
|
|
.set({
|
|
cost: sql`${SessionTable.cost} + ${value.cost * sign}`,
|
|
tokens_input: sql`${SessionTable.tokens_input} + ${value.tokens.input * sign}`,
|
|
tokens_output: sql`${SessionTable.tokens_output} + ${value.tokens.output * sign}`,
|
|
tokens_reasoning: sql`${SessionTable.tokens_reasoning} + ${value.tokens.reasoning * sign}`,
|
|
tokens_cache_read: sql`${SessionTable.tokens_cache_read} + ${value.tokens.cache.read * sign}`,
|
|
tokens_cache_write: sql`${SessionTable.tokens_cache_write} + ${value.tokens.cache.write * sign}`,
|
|
time_updated: sql`${SessionTable.time_updated}`,
|
|
})
|
|
.where(eq(SessionTable.id, sessionID))
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}
|
|
|
|
const projectFork = Effect.fn("SessionProjector.projectFork")(function* (
|
|
db: DatabaseService,
|
|
event: typeof SessionEvent.Forked.Type,
|
|
) {
|
|
const parent = yield* db
|
|
.select()
|
|
.from(SessionTable)
|
|
.where(eq(SessionTable.id, event.data.parentID))
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
if (!parent) return yield* Effect.die(`Fork parent session not found: ${event.data.parentID}`)
|
|
const boundary = event.data.messageID
|
|
? yield* db
|
|
.select({ seq: SessionMessageTable.seq })
|
|
.from(SessionMessageTable)
|
|
.where(
|
|
and(
|
|
eq(SessionMessageTable.session_id, event.data.parentID),
|
|
eq(SessionMessageTable.id, event.data.messageID),
|
|
),
|
|
)
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
: undefined
|
|
if (event.data.messageID && !boundary)
|
|
return yield* Effect.die(`Fork boundary message not found: ${event.data.messageID}`)
|
|
const copied = yield* db
|
|
.select({ seq: SessionMessageTable.seq })
|
|
.from(SessionMessageTable)
|
|
.where(
|
|
and(
|
|
eq(SessionMessageTable.session_id, event.data.parentID),
|
|
boundary === undefined ? undefined : lt(SessionMessageTable.seq, boundary.seq),
|
|
),
|
|
)
|
|
.orderBy(desc(SessionMessageTable.seq))
|
|
.limit(1)
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
const copiedSeq = copied?.seq ?? 0
|
|
|
|
const stored = yield* db
|
|
.insert(SessionTable)
|
|
.values({
|
|
id: event.data.sessionID,
|
|
parent_id: event.data.parentID,
|
|
project_id: parent.project_id,
|
|
workspace_id: parent.workspace_id,
|
|
slug: Slug.create(),
|
|
directory: parent.directory,
|
|
path: parent.path,
|
|
title: forkTitle(parent.title),
|
|
agent: parent.agent,
|
|
model: parent.model,
|
|
version: parent.version,
|
|
cost: 0,
|
|
tokens_input: 0,
|
|
tokens_output: 0,
|
|
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),
|
|
})
|
|
.onConflictDoNothing()
|
|
.returning({ sessionID: SessionTable.id })
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
if (!stored) return yield* Effect.die(new SessionAlreadyProjected())
|
|
|
|
// The fork inherits the parent's transcript, so it inherits the context
|
|
// checkpoint that transcript was built against: copied message seqs keep
|
|
// folding at the same baseline horizon.
|
|
const checkpoint = yield* db
|
|
.select()
|
|
.from(SessionContextCheckpointTable)
|
|
.where(eq(SessionContextCheckpointTable.session_id, event.data.parentID))
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
if (checkpoint) {
|
|
yield* db
|
|
.insert(SessionContextCheckpointTable)
|
|
.values({ ...checkpoint, session_id: event.data.sessionID })
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}
|
|
|
|
const usage = emptyUsage()
|
|
let cursor = -1
|
|
while (true) {
|
|
const rows = yield* db
|
|
.select()
|
|
.from(SessionMessageTable)
|
|
.where(
|
|
and(
|
|
eq(SessionMessageTable.session_id, event.data.parentID),
|
|
gt(SessionMessageTable.seq, cursor),
|
|
copiedSeq === 0 ? undefined : lt(SessionMessageTable.seq, copiedSeq + 1),
|
|
),
|
|
)
|
|
.orderBy(asc(SessionMessageTable.seq))
|
|
.limit(ForkBatchSize)
|
|
.all()
|
|
.pipe(Effect.orDie)
|
|
if (rows.length === 0) break
|
|
|
|
const idMap = new Map(rows.map((row) => [row.id, SessionMessage.ID.create()]))
|
|
yield* db
|
|
.insert(SessionMessageTable)
|
|
.values(
|
|
rows.map((row) => {
|
|
const id = idMap.get(row.id)
|
|
if (!id) throw new Error(`Fork message ID mapping missing: ${row.id}`)
|
|
return {
|
|
id,
|
|
session_id: event.data.sessionID,
|
|
type: row.type,
|
|
seq: row.seq,
|
|
time_created: row.time_created,
|
|
time_updated: row.time_updated,
|
|
data: row.type === "synthetic" ? { ...row.data, sessionID: event.data.sessionID } : row.data,
|
|
}
|
|
}),
|
|
)
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
|
|
const inputRows = yield* db
|
|
.select()
|
|
.from(SessionInputTable)
|
|
.where(
|
|
and(
|
|
eq(SessionInputTable.session_id, event.data.parentID),
|
|
inArray(
|
|
SessionInputTable.id,
|
|
rows.map((row) => row.id),
|
|
),
|
|
),
|
|
)
|
|
.all()
|
|
.pipe(Effect.orDie)
|
|
if (inputRows.length > 0) {
|
|
yield* db
|
|
.insert(SessionInputTable)
|
|
.values(
|
|
inputRows.flatMap((row) => {
|
|
const id = idMap.get(row.id)
|
|
return id
|
|
? [
|
|
{
|
|
id,
|
|
session_id: event.data.sessionID,
|
|
prompt: row.prompt,
|
|
delivery: row.delivery,
|
|
admitted_seq: row.admitted_seq,
|
|
promoted_seq: row.promoted_seq,
|
|
time_created: row.time_created,
|
|
},
|
|
]
|
|
: []
|
|
}),
|
|
)
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}
|
|
|
|
for (const row of rows) {
|
|
const value = messageUsage(row)
|
|
if (value) addUsage(usage, value)
|
|
}
|
|
cursor = rows.at(-1)!.seq
|
|
}
|
|
|
|
yield* db
|
|
.update(SessionTable)
|
|
.set({
|
|
cost: usage.cost,
|
|
tokens_input: usage.tokens.input,
|
|
tokens_output: usage.tokens.output,
|
|
tokens_reasoning: usage.tokens.reasoning,
|
|
tokens_cache_read: usage.tokens.cache.read,
|
|
tokens_cache_write: usage.tokens.cache.write,
|
|
})
|
|
.where(eq(SessionTable.id, event.data.sessionID))
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
if (copiedSeq > 0) yield* EventV2.reserveSequence(db, event.data.sessionID, copiedSeq)
|
|
})
|
|
|
|
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.Message) => {
|
|
if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence")
|
|
const encoded = encodeMessage(message)
|
|
const { id, type, ...data } = encoded
|
|
return db
|
|
.update(SessionMessageTable)
|
|
.set({ type, time_created: DateTime.toEpochMillis(message.time.created), data })
|
|
.where(
|
|
and(
|
|
eq(SessionMessageTable.id, SessionMessage.ID.make(id)),
|
|
eq(SessionMessageTable.session_id, event.data.sessionID),
|
|
),
|
|
)
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}
|
|
const appendMessage = (message: SessionMessage.Message) => insertMessage(db, event, message)
|
|
const adapter: SessionMessageUpdater.Adapter = {
|
|
getCurrentAssistant() {
|
|
return Effect.gen(function* () {
|
|
// A newer turn supersedes stale incomplete rows; never resume an older assistant projection.
|
|
const row = yield* db
|
|
.select()
|
|
.from(SessionMessageTable)
|
|
.where(
|
|
and(eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.type, "assistant")),
|
|
)
|
|
.orderBy(desc(SessionMessageTable.seq))
|
|
.limit(1)
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
if (!row) return
|
|
const message = decodeRow(row)
|
|
return message.type === "assistant" && !message.time.completed ? message : undefined
|
|
})
|
|
},
|
|
getAssistant(messageID) {
|
|
return Effect.gen(function* () {
|
|
const row = yield* db
|
|
.select()
|
|
.from(SessionMessageTable)
|
|
.where(
|
|
and(
|
|
eq(SessionMessageTable.id, messageID),
|
|
eq(SessionMessageTable.session_id, event.data.sessionID),
|
|
eq(SessionMessageTable.type, "assistant"),
|
|
),
|
|
)
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
if (!row) return
|
|
const message = decodeRow(row)
|
|
return message.type === "assistant" ? message : undefined
|
|
})
|
|
},
|
|
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")))
|
|
.orderBy(desc(SessionMessageTable.seq))
|
|
.all()
|
|
.pipe(Effect.orDie)
|
|
return rows
|
|
.map(decodeRow)
|
|
.find((message): message is SessionMessage.Shell => message.type === "shell" && message.callID === callID)
|
|
})
|
|
},
|
|
updateAssistant: updateMessage,
|
|
updateShell: updateMessage,
|
|
appendMessage,
|
|
}
|
|
yield* SessionMessageUpdater.update(adapter, event)
|
|
})
|
|
}
|
|
|
|
function insertMessage(db: DatabaseService, event: SessionEvent.Event, message: SessionMessage.Message) {
|
|
if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence")
|
|
const encoded = encodeMessage(message)
|
|
const { id, type, ...data } = encoded
|
|
return db
|
|
.insert(SessionMessageTable)
|
|
.values({
|
|
id: SessionMessage.ID.make(id),
|
|
session_id: event.data.sessionID,
|
|
type,
|
|
seq: event.durable.seq,
|
|
time_created: DateTime.toEpochMillis(message.time.created),
|
|
data,
|
|
})
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}
|
|
|
|
const layer = Layer.effectDiscard(
|
|
Effect.gen(function* () {
|
|
const events = yield* EventV2.Service
|
|
const { db } = yield* Database.Service
|
|
yield* events.project(SessionV1.Event.Created, (event) =>
|
|
Effect.gen(function* () {
|
|
const stored = yield* db
|
|
.insert(SessionTable)
|
|
.values(sessionRow(event.data.info))
|
|
.onConflictDoNothing()
|
|
.returning({ sessionID: SessionTable.id })
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
if (!stored) return yield* Effect.die(new SessionAlreadyProjected())
|
|
if (event.data.info.workspaceID) {
|
|
yield* db
|
|
.update(WorkspaceTable)
|
|
.set({ time_used: Date.now() })
|
|
.where(eq(WorkspaceTable.id, event.data.info.workspaceID))
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}
|
|
}),
|
|
)
|
|
yield* events.project(SessionV1.Event.Updated, (event) =>
|
|
db
|
|
.update(SessionTable)
|
|
.set(sessionRow(event.data.info))
|
|
.where(eq(SessionTable.id, event.data.sessionID))
|
|
.run()
|
|
.pipe(Effect.orDie),
|
|
)
|
|
yield* events.project(SessionEvent.Moved, (event) =>
|
|
Effect.gen(function* () {
|
|
yield* db
|
|
.update(SessionTable)
|
|
.set({
|
|
directory: event.data.location.directory,
|
|
path: event.data.subdirectory,
|
|
workspace_id: event.data.location.workspaceID ? WorkspaceV2.ID.make(event.data.location.workspaceID) : null,
|
|
time_updated: DateTime.toEpochMillis(event.data.timestamp),
|
|
})
|
|
.where(eq(SessionTable.id, event.data.sessionID))
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
yield* SessionContextCheckpoint.reset(db, event.data.sessionID)
|
|
}),
|
|
)
|
|
yield* events.project(SessionV1.Event.Deleted, (event) =>
|
|
db.delete(SessionTable).where(eq(SessionTable.id, event.data.sessionID)).run().pipe(Effect.orDie),
|
|
)
|
|
yield* events.project(SessionV1.Event.MessageUpdated, (event) =>
|
|
Effect.gen(function* () {
|
|
const time_created = event.data.info.time.created
|
|
const id = event.data.info.id
|
|
const sessionID = event.data.info.sessionID
|
|
const data = messageData(event.data.info)
|
|
yield* db
|
|
.insert(MessageTable)
|
|
.values({ id, session_id: sessionID, time_created, data })
|
|
.onConflictDoUpdate({ target: MessageTable.id, set: { data } })
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}),
|
|
)
|
|
yield* events.project(SessionV1.Event.MessageRemoved, (event) =>
|
|
Effect.gen(function* () {
|
|
const rows = yield* db
|
|
.select()
|
|
.from(PartTable)
|
|
.where(and(eq(PartTable.message_id, event.data.messageID), eq(PartTable.session_id, event.data.sessionID)))
|
|
.all()
|
|
.pipe(Effect.orDie)
|
|
for (const row of rows) {
|
|
const previous = usage(row.data)
|
|
if (previous) yield* applyUsage(db, event.data.sessionID, previous, -1)
|
|
}
|
|
yield* db
|
|
.delete(MessageTable)
|
|
.where(and(eq(MessageTable.id, event.data.messageID), eq(MessageTable.session_id, event.data.sessionID)))
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}),
|
|
)
|
|
yield* events.project(SessionV1.Event.PartRemoved, (event) =>
|
|
Effect.gen(function* () {
|
|
const row = yield* db
|
|
.select()
|
|
.from(PartTable)
|
|
.where(and(eq(PartTable.id, event.data.partID), eq(PartTable.session_id, event.data.sessionID)))
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
const previous = row && usage(row.data)
|
|
if (previous) yield* applyUsage(db, event.data.sessionID, previous, -1)
|
|
yield* db
|
|
.delete(PartTable)
|
|
.where(and(eq(PartTable.id, event.data.partID), eq(PartTable.session_id, event.data.sessionID)))
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
}),
|
|
)
|
|
yield* events.project(SessionV1.Event.PartUpdated, (event) =>
|
|
Effect.gen(function* () {
|
|
const id = event.data.part.id
|
|
const messageID = event.data.part.messageID
|
|
const sessionID = event.data.part.sessionID
|
|
const data = partData(event.data.part)
|
|
const row = yield* db.select().from(PartTable).where(eq(PartTable.id, id)).get().pipe(Effect.orDie)
|
|
yield* db
|
|
.insert(PartTable)
|
|
.values({ id, message_id: messageID, session_id: sessionID, time_created: event.data.time, data })
|
|
.onConflictDoUpdate({ target: PartTable.id, set: { data } })
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
const previous = row && usage(row.data)
|
|
const next = usage(event.data.part)
|
|
if (previous) yield* applyUsage(db, row.session_id, previous, -1)
|
|
if (next) yield* applyUsage(db, sessionID, next)
|
|
}),
|
|
)
|
|
yield* events.project(SessionEvent.AgentSwitched, (event) =>
|
|
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, Effect.andThen(run(db, event))),
|
|
)
|
|
yield* events.project(SessionEvent.ModelSwitched, (event) =>
|
|
Effect.gen(function* () {
|
|
yield* 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* run(db, event)
|
|
}),
|
|
)
|
|
yield* events.project(SessionEvent.Renamed, (event) =>
|
|
db
|
|
.update(SessionTable)
|
|
.set({ title: event.data.title, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
|
|
.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) =>
|
|
Effect.gen(function* () {
|
|
if (event.durable === undefined) return yield* Effect.die("Durable Session event is missing aggregate sequence")
|
|
yield* SessionInput.projectPrompted(db, {
|
|
id: event.data.messageID,
|
|
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* events.project(SessionEvent.PromptAdmitted, (event) =>
|
|
Effect.gen(function* () {
|
|
if (event.durable === undefined) return yield* Effect.die("Durable Session event is missing aggregate sequence")
|
|
yield* SessionInput.projectAdmitted(db, {
|
|
admittedSeq: event.durable.seq,
|
|
id: event.data.messageID,
|
|
sessionID: event.data.sessionID,
|
|
prompt: event.data.prompt,
|
|
delivery: event.data.delivery,
|
|
timeCreated: event.data.timestamp,
|
|
})
|
|
}),
|
|
)
|
|
yield* events.project(SessionEvent.ContextUpdated, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Synthetic, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Skill.Activated, (event) =>
|
|
insertMessage(db, event, {
|
|
id: event.data.messageID,
|
|
type: "skill",
|
|
name: event.data.name,
|
|
text: event.data.text,
|
|
time: { created: event.data.timestamp },
|
|
}),
|
|
)
|
|
yield* events.project(SessionEvent.Shell.Started, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Shell.Ended, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Step.Started, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Step.Ended, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Step.Failed, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Text.Started, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Text.Ended, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Tool.Input.Started, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Tool.Input.Ended, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Tool.Called, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Tool.Progress, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Tool.Success, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Tool.Failed, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Reasoning.Started, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Reasoning.Ended, (event) => run(db, event))
|
|
// yield* events.project(SessionEvent.Retried, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.Compaction.Ended, (event) => run(db, event))
|
|
yield* events.project(SessionEvent.RevertEvent.Staged, (event) =>
|
|
db
|
|
.update(SessionTable)
|
|
.set({
|
|
revert: { ...event.data.revert, files: event.data.revert.files ? [...event.data.revert.files] : undefined },
|
|
time_updated: DateTime.toEpochMillis(event.data.timestamp),
|
|
})
|
|
.where(eq(SessionTable.id, event.data.sessionID))
|
|
.run()
|
|
.pipe(Effect.orDie, Effect.asVoid),
|
|
)
|
|
yield* events.project(SessionEvent.RevertEvent.Cleared, (event) =>
|
|
db
|
|
.update(SessionTable)
|
|
.set({ revert: null, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
|
|
.where(eq(SessionTable.id, event.data.sessionID))
|
|
.run()
|
|
.pipe(Effect.orDie, Effect.asVoid),
|
|
)
|
|
yield* events.project(SessionEvent.RevertEvent.Committed, (event) =>
|
|
Effect.gen(function* () {
|
|
const boundary = yield* db
|
|
.select({ seq: SessionMessageTable.seq })
|
|
.from(SessionMessageTable)
|
|
.where(
|
|
and(
|
|
eq(SessionMessageTable.session_id, event.data.sessionID),
|
|
eq(SessionMessageTable.id, event.data.messageID),
|
|
),
|
|
)
|
|
.get()
|
|
.pipe(Effect.orDie)
|
|
if (!boundary) return yield* Effect.die(`Revert boundary message not found: ${event.data.messageID}`)
|
|
yield* db
|
|
.delete(SessionMessageTable)
|
|
.where(
|
|
and(eq(SessionMessageTable.session_id, event.data.sessionID), gt(SessionMessageTable.seq, boundary.seq)),
|
|
)
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
yield* db
|
|
.delete(SessionInputTable)
|
|
.where(
|
|
and(
|
|
eq(SessionInputTable.session_id, event.data.sessionID),
|
|
or(gt(SessionInputTable.admitted_seq, boundary.seq), gt(SessionInputTable.promoted_seq, boundary.seq)),
|
|
),
|
|
)
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
yield* db
|
|
.update(SessionTable)
|
|
.set({ revert: null, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
|
|
.where(eq(SessionTable.id, event.data.sessionID))
|
|
.run()
|
|
.pipe(Effect.orDie)
|
|
yield* SessionContextCheckpoint.reset(db, event.data.sessionID)
|
|
}),
|
|
)
|
|
}),
|
|
)
|
|
|
|
export const node = makeGlobalNode({ name: "session-projector", layer, deps: [EventV2.node, Database.node] })
|