Buckets:
| export * as SessionProjector from "./projector" | |
| import { and, desc, eq, sql } from "drizzle-orm" | |
| import { DateTime, Effect, Layer, Schema } from "effect" | |
| import { Database } from "../database/database" | |
| import { EventV2 } from "../event" | |
| import { LayerNode } from "../effect/layer-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 { SessionContextEpoch } from "./context-epoch" | |
| import { MessageTable, PartTable, SessionMessageTable, SessionTable } from "./sql" | |
| import type { DeepMutable } from "../schema" | |
| type DatabaseService = Database.Interface["db"] | |
| const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Message) | |
| const encodeMessage = Schema.encodeSync(SessionMessage.Message) | |
| class PromptAlreadyProjected extends Error {} | |
| export class SessionAlreadyProjected extends Error {} | |
| type Usage = { | |
| cost: number | |
| tokens: { | |
| input: number | |
| output: number | |
| reasoning: number | |
| cache: { read: number; write: number } | |
| } | |
| } | |
| 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 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 ?? 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) | |
| } | |
| function run(db: DatabaseService, event: SessionEvent.Event) { | |
| 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.seq === undefined) return Effect.die("Synchronized 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.seq === undefined) return Effect.die("Synchronized 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.seq, | |
| time_created: DateTime.toEpochMillis(message.time.created), | |
| data, | |
| }) | |
| .run() | |
| .pipe(Effect.orDie) | |
| } | |
| export const layer = Layer.effectDiscard( | |
| Effect.gen(function* () { | |
| const events = yield* EventV2.Service | |
| const { db } = yield* Database.Service | |
| yield* events.beforeCommit((event) => SessionInput.guardReservedID(db, event)) | |
| 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* SessionContextEpoch.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) => { | |
| if (event.seq === undefined) return Effect.die("Synchronized Session event is missing aggregate sequence") | |
| return 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)), | |
| Effect.andThen(SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.seq)), | |
| ) | |
| }) | |
| 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) | |
| if (event.seq === undefined) | |
| return yield* Effect.die("Synchronized Session event is missing aggregate sequence") | |
| yield* SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.seq) | |
| }), | |
| ) | |
| yield* events.project(SessionEvent.Prompted, (event) => | |
| Effect.gen(function* () { | |
| const messageID = event.data.messageID | |
| const existing = yield* db | |
| .select({ id: SessionMessageTable.id }) | |
| .from(SessionMessageTable) | |
| .where(eq(SessionMessageTable.id, messageID)) | |
| .get() | |
| .pipe(Effect.orDie) | |
| if (existing) return yield* Effect.die(new PromptAlreadyProjected()) | |
| yield* run(db, event) | |
| if (event.seq === undefined) | |
| return yield* Effect.die("Synchronized Session event is missing aggregate sequence") | |
| yield* SessionInput.projectLegacyPrompted(db, { | |
| id: messageID, | |
| sessionID: event.data.sessionID, | |
| prompt: event.data.prompt, | |
| delivery: event.data.delivery, | |
| timeCreated: event.data.timestamp, | |
| promotedSeq: event.seq, | |
| }) | |
| }), | |
| ) | |
| yield* events.project(SessionEvent.PromptLifecycle.Admitted, (event) => | |
| Effect.gen(function* () { | |
| if (event.seq === undefined) | |
| return yield* Effect.die("Synchronized Session event is missing aggregate sequence") | |
| yield* SessionInput.projectAdmitted(db, { | |
| admittedSeq: event.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.PromptLifecycle.Promoted, (event) => | |
| Effect.gen(function* () { | |
| if (event.seq === undefined) | |
| return yield* Effect.die("Synchronized Session event is missing aggregate sequence") | |
| yield* insertMessage( | |
| db, | |
| event, | |
| yield* SessionInput.projectPromoted(db, { | |
| id: event.data.messageID, | |
| sessionID: event.data.sessionID, | |
| prompt: event.data.prompt, | |
| timeCreated: event.data.timeCreated, | |
| promotedSeq: event.seq, | |
| }), | |
| ) | |
| }), | |
| ) | |
| yield* events.project(SessionEvent.InterruptRequested, () => Effect.void) | |
| yield* events.project(SessionEvent.ContextUpdated, (event) => { | |
| if (!event.replay || event.seq === undefined) return run(db, event) | |
| return run(db, event).pipe( | |
| Effect.andThen(SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.seq)), | |
| ) | |
| }) | |
| yield* events.project(SessionEvent.Synthetic, (event) => run(db, event)) | |
| 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) => { | |
| if (event.version === 1) return Effect.void | |
| const seq = event.seq | |
| if (seq === undefined) return Effect.die("Synchronized Session event is missing aggregate sequence") | |
| return Effect.gen(function* () { | |
| yield* run(db, event) | |
| yield* SessionContextEpoch.requestReplacement(db, event.data.sessionID, seq) | |
| }) | |
| }) | |
| }), | |
| ) | |
| export const defaultLayer = layer.pipe(Layer.provide(EventV2.defaultLayer), Layer.provide(Database.defaultLayer)) | |
| export const node = LayerNode.make(layer, [EventV2.node, Database.node]) | |
Xet Storage Details
- Size:
- 18.6 kB
- Xet hash:
- 8289124c71df86a5d3489726122e74296b072fb9eb9e4a3e6306e4ef45a0c2b9
·
Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.