| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827 |
- import { describe, expect } from "bun:test"
- import { DateTime, Effect, Fiber, Option, Schema, Stream } from "effect"
- import { asc, eq, sql } from "drizzle-orm"
- import { Database } from "@opencode-ai/core/database/database"
- import { AgentV2 } from "@opencode-ai/core/agent"
- import { LayerNode } from "@opencode-ai/util/effect/layer-node"
- import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
- import { EventV2 } from "@opencode-ai/core/event"
- import { EventTable } from "@opencode-ai/core/event/sql"
- import { ModelV2 } from "@opencode-ai/core/model"
- import { Project } from "@opencode-ai/core/project"
- import { ProjectTable } from "@opencode-ai/core/project/sql"
- import { ProviderV2 } from "@opencode-ai/core/provider"
- import { AbsolutePath } from "@opencode-ai/core/schema"
- import { SessionV2 } from "@opencode-ai/core/session"
- import { SessionEvent } from "@opencode-ai/core/session/event"
- import { SessionMessage } from "@opencode-ai/core/session/message"
- import { Money } from "@opencode-ai/schema/money"
- import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater"
- import { SessionProjector } from "@opencode-ai/core/session/projector"
- import { SessionExecution } from "@opencode-ai/core/session/execution"
- import { fromRow } from "@opencode-ai/core/session/info"
- import { SessionPending } from "@opencode-ai/core/session/pending"
- import { Shell } from "@opencode-ai/schema/shell"
- import {
- InstructionStateTable,
- SessionPendingTable,
- SessionMessageTable,
- SessionTable,
- } from "@opencode-ai/core/session/sql"
- import { testEffect } from "./lib/effect"
- import { Snapshot } from "@opencode-ai/core/snapshot"
- const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node, SessionProjector.node])))
- const sessionsLayer = AppNodeBuilder.build(SessionV2.node, [[SessionExecution.node, SessionExecution.noopLayer]])
- const sessionID = SessionV2.ID.make("ses_projector_test")
- const created = DateTime.makeUnsafe(0)
- const model = { id: ModelV2.ID.make("model"), providerID: ProviderV2.ID.make("provider") }
- const previousModel = { ...model, variant: ModelV2.VariantID.make("medium") }
- const encodeMessage = Schema.encodeSync(SessionMessage.Info)
- const build = AgentV2.defaultID
- const assistantRow = (
- id: SessionMessage.ID,
- seq: number,
- time: { created: DateTime.Utc; completed?: DateTime.Utc } = { created },
- usage?: Pick<SessionMessage.Assistant, "cost" | "tokens">,
- ) => {
- const {
- id: _,
- type,
- ...data
- } = encodeMessage(
- SessionMessage.Assistant.make({ id, type: "assistant", agent: build, model, content: [], time, ...usage }),
- )
- return { id, session_id: sessionID, type, seq, time_created: DateTime.toEpochMillis(time.created), data }
- }
- describe("SessionProjector", () => {
- it.effect("does not settle a pending manual compaction on an auto failure", () =>
- Effect.gen(function* () {
- const db = (yield* Database.Service).db
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- const events = yield* EventV2.Service
- const inputID = SessionMessage.ID.make("msg_manual_compaction")
- yield* SessionPending.admitCompaction(db, events, { id: inputID, sessionID })
- yield* events.publish(SessionEvent.Compaction.Failed, {
- sessionID,
- reason: "auto",
- error: { type: "compaction.failed", message: "Auto compaction failed" },
- })
- expect(yield* SessionPending.compaction(db, sessionID)).toMatchObject({ id: inputID })
- }),
- )
- it.effect("loads legacy revert storage into canonical state", () =>
- Effect.gen(function* () {
- const db = (yield* Database.Service).db
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- const legacy = JSON.stringify({
- messageID: "msg_boundary",
- snapshot: "tree",
- diff: "legacy patch",
- files: [{ path: "src/old.ts", status: "modified", additions: 1, deletions: 0, patch: "@@" }],
- })
- yield* db.run(sql`update session set revert = ${legacy} where id = ${sessionID}`)
- const stored = yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get()
- if (!stored) return yield* Effect.die("Session row missing")
- const storedRevert = fromRow(stored).revert
- expect(String(storedRevert?.messageID)).toBe("msg_boundary")
- expect(String(storedRevert?.snapshot)).toBe("tree")
- expect(storedRevert?.files).toEqual([
- { file: "src/old.ts", status: "modified", additions: 1, deletions: 0, patch: "@@" },
- ])
- }),
- )
- it.effect("folds live compaction deltas into running memory state", () =>
- Effect.gen(function* () {
- const state = {
- messages: [
- SessionMessage.CompactionRunning.make({
- id: SessionMessage.ID.make("msg_compaction"),
- type: "compaction",
- status: "running",
- reason: "manual",
- summary: "partial ",
- recent: "recent",
- time: { created },
- }),
- ],
- }
- yield* SessionMessageUpdater.update(
- SessionMessageUpdater.memory(state),
- SessionEvent.Compaction.Delta.make({
- id: EventV2.ID.make("evt_delta"),
- type: "session.compaction.delta",
- created,
- data: { sessionID, text: "summary" },
- }),
- )
- expect(state.messages[0]).toMatchObject({ status: "running", summary: "partial summary", recent: "recent" })
- }),
- )
- it.effect("projects staged, cleared, and committed reverts", () =>
- Effect.gen(function* () {
- const db = (yield* Database.Service).db
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- cost: 1.25,
- tokens_input: 10,
- tokens_output: 4,
- tokens_reasoning: 2,
- tokens_cache_read: 3,
- tokens_cache_write: 1,
- })
- .run()
- const boundary = SessionMessage.ID.make("msg_boundary")
- const earlier = SessionMessage.ID.make("msg_earlier")
- yield* db
- .insert(SessionMessageTable)
- .values([
- assistantRow(earlier, 0),
- assistantRow(
- boundary,
- 1,
- { created },
- {
- cost: Money.USD.make(0.5),
- tokens: { input: 4, output: 1, reasoning: 1, cache: { read: 1, write: 0 } },
- },
- ),
- assistantRow(
- SessionMessage.ID.make("msg_later"),
- 2,
- { created },
- {
- cost: Money.USD.make(0.75),
- tokens: { input: 6, output: 3, reasoning: 1, cache: { read: 2, write: 1 } },
- },
- ),
- ])
- .run()
- yield* db
- .insert(InstructionStateTable)
- .values({
- session_id: sessionID,
- epoch_start: 0,
- through_seq: 0,
- initial_values: {},
- current_values: {},
- })
- .run()
- const events = yield* EventV2.Service
- yield* events.publish(SessionEvent.RevertEvent.Staged, {
- sessionID,
- revert: { messageID: boundary, snapshot: Snapshot.ID.make("tree"), files: [] },
- })
- expect((yield* db.select({ revert: SessionTable.revert }).from(SessionTable).get())?.revert).toMatchObject({
- messageID: boundary,
- snapshot: "tree",
- files: [],
- })
- yield* events.publish(SessionEvent.RevertEvent.Cleared, { sessionID })
- expect((yield* db.select({ revert: SessionTable.revert }).from(SessionTable).get())?.revert).toBeNull()
- yield* events.publish(SessionEvent.RevertEvent.Staged, {
- sessionID,
- revert: { messageID: boundary, files: [] },
- })
- yield* events.publish(SessionEvent.RevertEvent.Committed, {
- sessionID,
- to: boundary,
- })
- expect(
- (yield* db.select({ id: SessionMessageTable.id }).from(SessionMessageTable).all()).map((row) => row.id),
- ).toEqual([earlier])
- expect(yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get()).toMatchObject({
- cost: Money.USD.make(1.25),
- tokens_input: 10,
- tokens_output: 4,
- tokens_reasoning: 2,
- tokens_cache_read: 3,
- tokens_cache_write: 1,
- })
- // A committed revert resets the fold cache so the next boundary establishes a new epoch.
- expect(yield* db.select().from(InstructionStateTable).get().pipe(Effect.orDie)).toBeUndefined()
- }),
- )
- it.effect("orders projected messages and context by durable aggregate sequence", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- .pipe(Effect.orDie)
- const events = yield* EventV2.Service
- yield* events.publish(SessionEvent.InputAdmitted, {
- sessionID,
- inputID: SessionMessage.ID.make("msg_first"),
- input: { type: "user", data: { text: "first" }, delivery: "steer" },
- })
- yield* events.publish(
- SessionEvent.InputPromoted,
- {
- sessionID,
- inputID: SessionMessage.ID.make("msg_first"),
- },
- { id: EventV2.ID.make("evt_z") },
- )
- yield* events.publish(SessionEvent.InputAdmitted, {
- sessionID,
- inputID: SessionMessage.ID.make("msg_second"),
- input: { type: "user", data: { text: "second" }, delivery: "steer" },
- })
- yield* events.publish(
- SessionEvent.InputPromoted,
- {
- sessionID,
- inputID: SessionMessage.ID.make("msg_second"),
- },
- { id: EventV2.ID.make("evt_a") },
- )
- const sessions = yield* SessionV2.Service
- const firstPage = yield* sessions.messages({ sessionID, limit: 1, order: "asc" })
- expect(firstPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["first"])
- const secondPage = yield* sessions.messages({
- sessionID,
- limit: 1,
- order: "asc",
- cursor: { id: firstPage[0]!.id, direction: "next" },
- })
- expect(secondPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["second"])
- expect(
- (yield* sessions.messages({
- sessionID,
- limit: 1,
- order: "asc",
- cursor: { id: secondPage[0]!.id, direction: "previous" },
- })).map((message) => (message.type === "user" ? message.text : message.type)),
- ).toEqual(["first"])
- expect(
- (yield* sessions.context(sessionID)).map((message) => (message.type === "user" ? message.text : message.type)),
- ).toEqual(["first", "second"])
- }).pipe(Effect.provide(sessionsLayer)),
- )
- it.effect("consumes the pending row and projects the message at promotion", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- .pipe(Effect.orDie)
- const events = yield* EventV2.Service
- const id = SessionMessage.ID.make("msg_admitted")
- const admitted = yield* SessionPending.admit(db, events, {
- id,
- sessionID,
- input: { type: "user", data: { text: "promote me" }, delivery: "steer" },
- })
- if (!admitted) return yield* Effect.die("Prompt admission failed")
- const event = yield* events.publish(SessionEvent.InputPromoted, {
- sessionID,
- inputID: id,
- })
- expect(
- yield* db.select().from(SessionPendingTable).where(eq(SessionPendingTable.id, id)).get().pipe(Effect.orDie),
- ).toBeUndefined()
- expect(
- yield* db.select().from(SessionMessageTable).where(eq(SessionMessageTable.id, id)).get().pipe(Effect.orDie),
- ).toMatchObject({ session_id: sessionID, type: "user", seq: event.durable?.seq })
- }),
- )
- it.effect("projects durable context messages supported by the updater", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- model: previousModel,
- })
- .run()
- .pipe(Effect.orDie)
- const events = yield* EventV2.Service
- yield* events.publish(SessionEvent.AgentSelected, {
- sessionID,
- agent: build,
- })
- yield* events.publish(SessionEvent.ModelSelected, {
- sessionID,
- model,
- })
- yield* events.publish(SessionEvent.Synthetic, {
- sessionID,
- text: "synthetic context",
- metadata: { source: "projector-test" },
- })
- yield* events.publish(SessionEvent.Shell.Started, {
- sessionID,
- shell: Shell.Info.make({
- id: Shell.ID.make("sh_projector"),
- status: "running",
- command: "pwd",
- cwd: "/project",
- shell: "/bin/sh",
- file: "/tmp/sh_projector.out",
- metadata: {},
- time: { started: 0 },
- }),
- })
- yield* events.publish(SessionEvent.Shell.Ended, {
- sessionID,
- shell: Shell.Info.make({
- id: Shell.ID.make("sh_projector"),
- status: "exited",
- command: "pwd",
- cwd: "/project",
- shell: "/bin/sh",
- file: "/tmp/sh_projector.out",
- exit: 0,
- metadata: {},
- time: { started: 0, completed: 1 },
- }),
- output: { output: "/project", cursor: 8, size: 8, truncated: false },
- })
- yield* events.publish(SessionEvent.Compaction.Started, {
- sessionID,
- reason: "manual",
- recent: "recent context",
- })
- yield* events.publish(SessionEvent.Compaction.Delta, {
- sessionID,
- text: "partial",
- })
- expect(
- yield* db
- .select({ id: EventTable.id })
- .from(EventTable)
- .where(sql`${EventTable.type} like 'session.compaction.delta.%'`)
- .all()
- .pipe(Effect.orDie),
- ).toHaveLength(0)
- expect(
- yield* db
- .select({ data: SessionMessageTable.data })
- .from(SessionMessageTable)
- .where(eq(SessionMessageTable.type, "compaction"))
- .all()
- .pipe(Effect.orDie),
- ).toEqual([{ data: expect.objectContaining({ status: "running", summary: "", recent: "recent context" }) }])
- yield* events.publish(SessionEvent.Compaction.Ended, {
- sessionID,
- reason: "manual",
- text: "summary",
- recent: "recent context",
- })
- const rows = yield* db
- .select()
- .from(SessionMessageTable)
- .where(eq(SessionMessageTable.session_id, sessionID))
- .orderBy(asc(SessionMessageTable.seq))
- .all()
- .pipe(Effect.orDie)
- const messages = rows.map((row) =>
- Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
- )
- expect(messages.map((message) => message.type)).toEqual([
- "agent-switched",
- "model-switched",
- "synthetic",
- "shell",
- "compaction",
- ])
- expect(messages.find((message) => message.type === "synthetic")).toMatchObject({
- text: "synthetic context",
- metadata: { source: "projector-test" },
- })
- expect(messages.find((message) => message.type === "model-switched")).toMatchObject({ previous: previousModel })
- expect(messages.find((message) => message.type === "shell")).toMatchObject({
- command: "pwd",
- status: "exited",
- exit: 0,
- output: { output: "/project", truncated: false },
- time: { completed: DateTime.makeUnsafe(0) },
- })
- expect(messages.find((message) => message.type === "compaction")).toMatchObject({
- summary: "summary",
- recent: "recent context",
- })
- expect(
- yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie),
- ).toMatchObject({
- agent: "build",
- model,
- time_updated: DateTime.toEpochMillis(created),
- })
- }),
- )
- it.effect("rejects distinct creator events that reuse one projected message ID", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- .pipe(Effect.orDie)
- const events = yield* EventV2.Service
- const id = SessionMessage.ID.make("msg_creator_collision")
- const { id: _, type, ...data } = encodeMessage({ id, type: "synthetic", text: "existing", time: { created } })
- yield* db
- .insert(SessionMessageTable)
- .values({ id, session_id: sessionID, type, seq: 0, time_created: 0, data })
- .run()
- const exit = yield* events
- .publish(SessionEvent.Step.Started, {
- sessionID,
- assistantMessageID: id,
- agent: build,
- model,
- })
- .pipe(Effect.exit)
- expect(exit._tag).toBe("Failure")
- expect(
- yield* db.select().from(SessionMessageTable).where(eq(SessionMessageTable.id, id)).get().pipe(Effect.orDie),
- ).toMatchObject({ type: "synthetic" })
- }),
- )
- it.effect("does not revive a stale incomplete in-memory assistant projection", () =>
- Effect.gen(function* () {
- const stale = SessionMessage.Assistant.make({
- id: SessionMessage.ID.make("msg_assistant_stale"),
- type: "assistant",
- agent: build,
- model,
- content: [],
- time: { created },
- })
- const completed = SessionMessage.Assistant.make({
- id: SessionMessage.ID.make("msg_assistant_completed"),
- type: "assistant",
- agent: build,
- model,
- content: [],
- time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
- })
- expect(
- yield* SessionMessageUpdater.memory({ messages: [stale, completed] }).getCurrentAssistant(),
- ).toBeUndefined()
- }),
- )
- it.effect("projects retry state and clears it at the next step or execution terminal", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- .pipe(Effect.orDie)
- const events = yield* EventV2.Service
- const first = SessionMessage.ID.make("msg_retry_first")
- const second = SessionMessage.ID.make("msg_retry_second")
- yield* events.publish(SessionEvent.Step.Started, { sessionID, assistantMessageID: first, agent: build, model })
- yield* events.publish(SessionEvent.RetryScheduled, {
- sessionID,
- assistantMessageID: first,
- attempt: 2,
- at: 2_000,
- error: { type: "provider.transport", message: "Disconnected" },
- })
- const decode = (row: typeof SessionMessageTable.$inferSelect) =>
- Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type })
- const firstRow = yield* db
- .select()
- .from(SessionMessageTable)
- .where(eq(SessionMessageTable.id, first))
- .get()
- .pipe(Effect.orDie)
- const projected = firstRow ?? (yield* Effect.die(new Error("Missing retry projection")))
- expect(decode(projected)).toMatchObject({
- retry: { attempt: 2, at: DateTime.makeUnsafe(2_000), error: { type: "provider.transport" } },
- })
- yield* events.publish(SessionEvent.Step.Started, { sessionID, assistantMessageID: second, agent: build, model })
- yield* events.publish(SessionEvent.RetryScheduled, {
- sessionID,
- assistantMessageID: second,
- attempt: 3,
- at: 6_000,
- error: { type: "provider.internal", message: "Unavailable" },
- })
- yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID, reason: "shutdown" })
- const rows = yield* db
- .select()
- .from(SessionMessageTable)
- .where(eq(SessionMessageTable.session_id, sessionID))
- .orderBy(asc(SessionMessageTable.seq))
- .all()
- .pipe(Effect.orDie)
- expect(decode(rows[0])).not.toHaveProperty("retry")
- expect(decode(rows[1])).not.toHaveProperty("retry")
- }),
- )
- it.effect("does not infer restart continuation from lifecycle history", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- .pipe(Effect.orDie)
- const events = yield* EventV2.Service
- const suspended = () =>
- db
- .select({ timeSuspended: SessionTable.time_suspended })
- .from(SessionTable)
- .where(eq(SessionTable.id, sessionID))
- .get()
- .pipe(Effect.orDie)
- yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID, reason: "shutdown" })
- expect((yield* suspended())?.timeSuspended).toBeNull()
- yield* events.publish(SessionEvent.Execution.Started, { sessionID })
- expect((yield* suspended())?.timeSuspended).toBeNull()
- }),
- )
- it.effect("updates only the newest incomplete assistant projection", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionMessageTable)
- .values([
- assistantRow(SessionMessage.ID.make("msg_assistant_1"), 0),
- assistantRow(SessionMessage.ID.make("msg_assistant_2"), 1),
- ])
- .run()
- .pipe(Effect.orDie)
- const service = yield* EventV2.Service
- const usageUpdated = yield* service
- .subscribe(SessionEvent.UsageUpdated)
- .pipe(Stream.runHead, Effect.forkScoped({ startImmediately: true }))
- yield* service.publish(SessionEvent.Step.Ended, {
- sessionID,
- assistantMessageID: SessionMessage.ID.make("msg_assistant_2"),
- finish: "stop",
- cost: Money.USD.make(1.25),
- tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
- })
- const rows = yield* db
- .select()
- .from(SessionMessageTable)
- .where(eq(SessionMessageTable.session_id, sessionID))
- .orderBy(asc(SessionMessageTable.id))
- .all()
- .pipe(Effect.orDie)
- const messages = rows.map((row) =>
- Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
- )
- expect(messages[0]).not.toHaveProperty("time.completed")
- expect(messages[1]).toMatchObject({
- type: "assistant",
- finish: "stop",
- cost: Money.USD.make(1.25),
- tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
- time: { completed: DateTime.makeUnsafe(0) },
- })
- expect(
- yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie),
- ).toMatchObject({
- cost: 1.25,
- tokens_input: 10,
- tokens_output: 4,
- tokens_reasoning: 2,
- tokens_cache_read: 3,
- tokens_cache_write: 1,
- })
- expect(Option.getOrThrow(yield* Fiber.join(usageUpdated)).data).toEqual({
- sessionID,
- cost: Money.USD.make(1.25),
- tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
- })
- }),
- )
- it.effect("does not revive a stale incomplete assistant projection", () =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionTable)
- .values({
- id: sessionID,
- project_id: Project.ID.global,
- slug: "test",
- directory: "/project",
- title: "test",
- version: "test",
- })
- .run()
- .pipe(Effect.orDie)
- yield* db
- .insert(SessionMessageTable)
- .values([
- assistantRow(SessionMessage.ID.make("msg_assistant_stale"), 0),
- assistantRow(SessionMessage.ID.make("msg_assistant_completed"), 1, {
- created: DateTime.makeUnsafe(1),
- completed: DateTime.makeUnsafe(2),
- }),
- ])
- .run()
- .pipe(Effect.orDie)
- const service = yield* EventV2.Service
- yield* service.publish(SessionEvent.Text.Started, {
- sessionID,
- assistantMessageID: SessionMessage.ID.make("msg_assistant_completed"),
- ordinal: 0,
- })
- const rows = yield* db
- .select()
- .from(SessionMessageTable)
- .where(eq(SessionMessageTable.session_id, sessionID))
- .orderBy(asc(SessionMessageTable.id))
- .all()
- .pipe(Effect.orDie)
- const messages = rows.map((row) =>
- Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
- )
- expect(messages).toEqual([
- SessionMessage.Assistant.make({
- id: SessionMessage.ID.make("msg_assistant_completed"),
- type: "assistant",
- agent: build,
- model,
- content: [SessionMessage.AssistantText.make({ type: "text", text: "" })],
- time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
- }),
- SessionMessage.Assistant.make({
- id: SessionMessage.ID.make("msg_assistant_stale"),
- type: "assistant",
- agent: build,
- model,
- content: [],
- time: { created },
- }),
- ])
- }),
- )
- })
|