session-projector.test.ts 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785
  1. import { describe, expect } from "bun:test"
  2. import { DateTime, Effect, Fiber, Option, Schema, Stream } from "effect"
  3. import { asc, eq, sql } from "drizzle-orm"
  4. import { Database } from "@opencode-ai/core/database/database"
  5. import { AgentV2 } from "@opencode-ai/core/agent"
  6. import { LayerNode } from "@opencode-ai/core/effect/layer-node"
  7. import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
  8. import { EventV2 } from "@opencode-ai/core/event"
  9. import { EventTable } from "@opencode-ai/core/event/sql"
  10. import { ModelV2 } from "@opencode-ai/core/model"
  11. import { Project } from "@opencode-ai/core/project"
  12. import { ProjectTable } from "@opencode-ai/core/project/sql"
  13. import { ProviderV2 } from "@opencode-ai/core/provider"
  14. import { AbsolutePath } from "@opencode-ai/core/schema"
  15. import { SessionV2 } from "@opencode-ai/core/session"
  16. import { SessionEvent } from "@opencode-ai/core/session/event"
  17. import { SessionMessage } from "@opencode-ai/core/session/message"
  18. import { Prompt } from "@opencode-ai/schema/prompt"
  19. import { Money } from "@opencode-ai/schema/money"
  20. import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater"
  21. import { SessionProjector } from "@opencode-ai/core/session/projector"
  22. import { SessionExecution } from "@opencode-ai/core/session/execution"
  23. import { fromRow } from "@opencode-ai/core/session/info"
  24. import { SessionInput } from "@opencode-ai/core/session/input"
  25. import { Shell } from "@opencode-ai/schema/shell"
  26. import {
  27. InstructionCheckpointTable,
  28. SessionInputTable,
  29. SessionMessageTable,
  30. SessionTable,
  31. } from "@opencode-ai/core/session/sql"
  32. import { testEffect } from "./lib/effect"
  33. import { Snapshot } from "@opencode-ai/core/snapshot"
  34. const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node, SessionProjector.node])))
  35. const sessionsLayer = AppNodeBuilder.build(SessionV2.node, [[SessionExecution.node, SessionExecution.noopLayer]])
  36. const sessionID = SessionV2.ID.make("ses_projector_test")
  37. const created = DateTime.makeUnsafe(0)
  38. const model = { id: ModelV2.ID.make("model"), providerID: ProviderV2.ID.make("provider") }
  39. const previousModel = { ...model, variant: ModelV2.VariantID.make("medium") }
  40. const encodeMessage = Schema.encodeSync(SessionMessage.Info)
  41. const build = AgentV2.defaultID
  42. const assistantRow = (
  43. id: SessionMessage.ID,
  44. seq: number,
  45. time: { created: DateTime.Utc; completed?: DateTime.Utc } = { created },
  46. usage?: Pick<SessionMessage.Assistant, "cost" | "tokens">,
  47. ) => {
  48. const {
  49. id: _,
  50. type,
  51. ...data
  52. } = encodeMessage(
  53. SessionMessage.Assistant.make({ id, type: "assistant", agent: build, model, content: [], time, ...usage }),
  54. )
  55. return { id, session_id: sessionID, type, seq, time_created: DateTime.toEpochMillis(time.created), data }
  56. }
  57. describe("SessionProjector", () => {
  58. it.effect("does not settle a pending manual compaction on an auto failure", () =>
  59. Effect.gen(function* () {
  60. const db = (yield* Database.Service).db
  61. yield* db
  62. .insert(ProjectTable)
  63. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  64. .run()
  65. yield* db
  66. .insert(SessionTable)
  67. .values({
  68. id: sessionID,
  69. project_id: Project.ID.global,
  70. slug: "test",
  71. directory: "/project",
  72. title: "test",
  73. version: "test",
  74. })
  75. .run()
  76. const events = yield* EventV2.Service
  77. const inputID = SessionMessage.ID.make("msg_manual_compaction")
  78. yield* SessionInput.admitCompaction(db, events, { id: inputID, sessionID })
  79. yield* events.publish(SessionEvent.Compaction.Failed, {
  80. sessionID,
  81. reason: "auto",
  82. error: { type: "compaction.failed", message: "Auto compaction failed" },
  83. })
  84. expect(yield* SessionInput.pendingCompaction(db, sessionID)).toMatchObject({ id: inputID })
  85. }),
  86. )
  87. it.effect("loads legacy revert storage into canonical state", () =>
  88. Effect.gen(function* () {
  89. const db = (yield* Database.Service).db
  90. yield* db
  91. .insert(ProjectTable)
  92. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  93. .run()
  94. yield* db
  95. .insert(SessionTable)
  96. .values({
  97. id: sessionID,
  98. project_id: Project.ID.global,
  99. slug: "test",
  100. directory: "/project",
  101. title: "test",
  102. version: "test",
  103. })
  104. .run()
  105. const legacy = JSON.stringify({
  106. messageID: "msg_boundary",
  107. snapshot: "tree",
  108. diff: "legacy patch",
  109. files: [{ path: "src/old.ts", status: "modified", additions: 1, deletions: 0, patch: "@@" }],
  110. })
  111. yield* db.run(sql`update session set revert = ${legacy} where id = ${sessionID}`)
  112. const stored = yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get()
  113. if (!stored) return yield* Effect.die("Session row missing")
  114. const storedRevert = fromRow(stored).revert
  115. expect(String(storedRevert?.messageID)).toBe("msg_boundary")
  116. expect(String(storedRevert?.snapshot)).toBe("tree")
  117. expect(storedRevert?.files).toEqual([
  118. { file: "src/old.ts", status: "modified", additions: 1, deletions: 0, patch: "@@" },
  119. ])
  120. }),
  121. )
  122. it.effect("folds live compaction deltas into running memory state", () =>
  123. Effect.gen(function* () {
  124. const state = {
  125. messages: [
  126. SessionMessage.CompactionRunning.make({
  127. id: SessionMessage.ID.make("msg_compaction"),
  128. type: "compaction",
  129. status: "running",
  130. reason: "manual",
  131. summary: "partial ",
  132. recent: "recent",
  133. time: { created },
  134. }),
  135. ],
  136. }
  137. yield* SessionMessageUpdater.update(
  138. SessionMessageUpdater.memory(state),
  139. SessionEvent.Compaction.Delta.make({
  140. id: EventV2.ID.make("evt_delta"),
  141. type: "session.compaction.delta",
  142. created,
  143. data: { sessionID, text: "summary" },
  144. }),
  145. )
  146. expect(state.messages[0]).toMatchObject({ status: "running", summary: "partial summary", recent: "recent" })
  147. }),
  148. )
  149. it.effect("projects staged, cleared, and committed reverts", () =>
  150. Effect.gen(function* () {
  151. const db = (yield* Database.Service).db
  152. yield* db
  153. .insert(ProjectTable)
  154. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  155. .run()
  156. yield* db
  157. .insert(SessionTable)
  158. .values({
  159. id: sessionID,
  160. project_id: Project.ID.global,
  161. slug: "test",
  162. directory: "/project",
  163. title: "test",
  164. version: "test",
  165. cost: 1.25,
  166. tokens_input: 10,
  167. tokens_output: 4,
  168. tokens_reasoning: 2,
  169. tokens_cache_read: 3,
  170. tokens_cache_write: 1,
  171. })
  172. .run()
  173. const boundary = SessionMessage.ID.make("msg_boundary")
  174. const earlier = SessionMessage.ID.make("msg_earlier")
  175. yield* db
  176. .insert(SessionMessageTable)
  177. .values([
  178. assistantRow(earlier, 0),
  179. assistantRow(
  180. boundary,
  181. 1,
  182. { created },
  183. {
  184. cost: Money.USD.make(0.5),
  185. tokens: { input: 4, output: 1, reasoning: 1, cache: { read: 1, write: 0 } },
  186. },
  187. ),
  188. assistantRow(
  189. SessionMessage.ID.make("msg_later"),
  190. 2,
  191. { created },
  192. {
  193. cost: Money.USD.make(0.75),
  194. tokens: { input: 6, output: 3, reasoning: 1, cache: { read: 2, write: 1 } },
  195. },
  196. ),
  197. ])
  198. .run()
  199. yield* db
  200. .insert(InstructionCheckpointTable)
  201. .values({ session_id: sessionID, baseline: "baseline", snapshot: {}, baseline_seq: 0 })
  202. .run()
  203. const events = yield* EventV2.Service
  204. yield* events.publish(SessionEvent.RevertEvent.Staged, {
  205. sessionID,
  206. revert: { messageID: boundary, snapshot: Snapshot.ID.make("tree"), files: [] },
  207. })
  208. expect((yield* db.select({ revert: SessionTable.revert }).from(SessionTable).get())?.revert).toMatchObject({
  209. messageID: boundary,
  210. snapshot: "tree",
  211. files: [],
  212. })
  213. yield* events.publish(SessionEvent.RevertEvent.Cleared, { sessionID })
  214. expect((yield* db.select({ revert: SessionTable.revert }).from(SessionTable).get())?.revert).toBeNull()
  215. yield* events.publish(SessionEvent.RevertEvent.Staged, {
  216. sessionID,
  217. revert: { messageID: boundary, files: [] },
  218. })
  219. yield* events.publish(SessionEvent.RevertEvent.Committed, {
  220. sessionID,
  221. to: boundary,
  222. })
  223. expect(
  224. (yield* db.select({ id: SessionMessageTable.id }).from(SessionMessageTable).all()).map((row) => row.id),
  225. ).toEqual([earlier])
  226. expect(yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get()).toMatchObject({
  227. cost: Money.USD.make(1.25),
  228. tokens_input: 10,
  229. tokens_output: 4,
  230. tokens_reasoning: 2,
  231. tokens_cache_read: 3,
  232. tokens_cache_write: 1,
  233. })
  234. // A committed revert resets the context checkpoint so the next turn re-initializes.
  235. expect(yield* db.select().from(InstructionCheckpointTable).get().pipe(Effect.orDie)).toBeUndefined()
  236. }),
  237. )
  238. it.effect("orders projected messages and context by durable aggregate sequence", () =>
  239. Effect.gen(function* () {
  240. const { db } = yield* Database.Service
  241. yield* db
  242. .insert(ProjectTable)
  243. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  244. .run()
  245. .pipe(Effect.orDie)
  246. yield* db
  247. .insert(SessionTable)
  248. .values({
  249. id: sessionID,
  250. project_id: Project.ID.global,
  251. slug: "test",
  252. directory: "/project",
  253. title: "test",
  254. version: "test",
  255. })
  256. .run()
  257. .pipe(Effect.orDie)
  258. const events = yield* EventV2.Service
  259. yield* events.publish(SessionEvent.PromptAdmitted, {
  260. sessionID,
  261. inputID: SessionMessage.ID.make("msg_first"),
  262. prompt: Prompt.make({ text: "first" }),
  263. delivery: "steer",
  264. })
  265. yield* events.publish(
  266. SessionEvent.PromptPromoted,
  267. {
  268. sessionID,
  269. inputID: SessionMessage.ID.make("msg_first"),
  270. },
  271. { id: EventV2.ID.make("evt_z") },
  272. )
  273. yield* events.publish(SessionEvent.PromptAdmitted, {
  274. sessionID,
  275. inputID: SessionMessage.ID.make("msg_second"),
  276. prompt: Prompt.make({ text: "second" }),
  277. delivery: "steer",
  278. })
  279. yield* events.publish(
  280. SessionEvent.PromptPromoted,
  281. {
  282. sessionID,
  283. inputID: SessionMessage.ID.make("msg_second"),
  284. },
  285. { id: EventV2.ID.make("evt_a") },
  286. )
  287. const sessions = yield* SessionV2.Service
  288. const firstPage = yield* sessions.messages({ sessionID, limit: 1, order: "asc" })
  289. expect(firstPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["first"])
  290. const secondPage = yield* sessions.messages({
  291. sessionID,
  292. limit: 1,
  293. order: "asc",
  294. cursor: { id: firstPage[0]!.id, direction: "next" },
  295. })
  296. expect(secondPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["second"])
  297. expect(
  298. (yield* sessions.messages({
  299. sessionID,
  300. limit: 1,
  301. order: "asc",
  302. cursor: { id: secondPage[0]!.id, direction: "previous" },
  303. })).map((message) => (message.type === "user" ? message.text : message.type)),
  304. ).toEqual(["first"])
  305. expect(
  306. (yield* sessions.context(sessionID)).map((message) => (message.type === "user" ? message.text : message.type)),
  307. ).toEqual(["first", "second"])
  308. }).pipe(Effect.provide(sessionsLayer)),
  309. )
  310. it.effect("marks an inbox row promoted with the PromptPromoted event sequence", () =>
  311. Effect.gen(function* () {
  312. const { db } = yield* Database.Service
  313. yield* db
  314. .insert(ProjectTable)
  315. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  316. .run()
  317. .pipe(Effect.orDie)
  318. yield* db
  319. .insert(SessionTable)
  320. .values({
  321. id: sessionID,
  322. project_id: Project.ID.global,
  323. slug: "test",
  324. directory: "/project",
  325. title: "test",
  326. version: "test",
  327. })
  328. .run()
  329. .pipe(Effect.orDie)
  330. const events = yield* EventV2.Service
  331. const id = SessionMessage.ID.make("msg_admitted")
  332. const admitted = yield* SessionInput.admit(db, events, {
  333. id,
  334. sessionID,
  335. prompt: Prompt.make({ text: "promote me" }),
  336. delivery: "steer",
  337. })
  338. if (!admitted) return yield* Effect.die("Prompt admission failed")
  339. const event = yield* events.publish(SessionEvent.PromptPromoted, {
  340. sessionID,
  341. inputID: id,
  342. })
  343. expect(
  344. yield* db.select().from(SessionInputTable).where(eq(SessionInputTable.id, id)).get().pipe(Effect.orDie),
  345. ).toMatchObject({ promoted_seq: event.durable?.seq })
  346. }),
  347. )
  348. it.effect("projects durable context messages supported by the updater", () =>
  349. Effect.gen(function* () {
  350. const { db } = yield* Database.Service
  351. yield* db
  352. .insert(ProjectTable)
  353. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  354. .run()
  355. .pipe(Effect.orDie)
  356. yield* db
  357. .insert(SessionTable)
  358. .values({
  359. id: sessionID,
  360. project_id: Project.ID.global,
  361. slug: "test",
  362. directory: "/project",
  363. title: "test",
  364. version: "test",
  365. model: previousModel,
  366. })
  367. .run()
  368. .pipe(Effect.orDie)
  369. const events = yield* EventV2.Service
  370. yield* events.publish(SessionEvent.AgentSelected, {
  371. sessionID,
  372. agent: build,
  373. })
  374. yield* events.publish(SessionEvent.ModelSelected, {
  375. sessionID,
  376. model,
  377. })
  378. yield* events.publish(SessionEvent.Synthetic, {
  379. sessionID,
  380. text: "synthetic context",
  381. metadata: { source: "projector-test" },
  382. })
  383. yield* events.publish(SessionEvent.Shell.Started, {
  384. sessionID,
  385. shell: Shell.Info.make({
  386. id: Shell.ID.make("sh_projector"),
  387. status: "running",
  388. command: "pwd",
  389. cwd: "/project",
  390. shell: "/bin/sh",
  391. file: "/tmp/sh_projector.out",
  392. metadata: {},
  393. time: { started: 0 },
  394. }),
  395. })
  396. yield* events.publish(SessionEvent.Shell.Ended, {
  397. sessionID,
  398. shell: Shell.Info.make({
  399. id: Shell.ID.make("sh_projector"),
  400. status: "exited",
  401. command: "pwd",
  402. cwd: "/project",
  403. shell: "/bin/sh",
  404. file: "/tmp/sh_projector.out",
  405. exit: 0,
  406. metadata: {},
  407. time: { started: 0, completed: 1 },
  408. }),
  409. output: { output: "/project", cursor: 8, size: 8, truncated: false },
  410. })
  411. yield* events.publish(SessionEvent.Compaction.Started, {
  412. sessionID,
  413. reason: "manual",
  414. recent: "recent context",
  415. })
  416. yield* events.publish(SessionEvent.Compaction.Delta, {
  417. sessionID,
  418. text: "partial",
  419. })
  420. expect(
  421. yield* db
  422. .select({ id: EventTable.id })
  423. .from(EventTable)
  424. .where(sql`${EventTable.type} like 'session.compaction.delta.%'`)
  425. .all()
  426. .pipe(Effect.orDie),
  427. ).toHaveLength(0)
  428. expect(
  429. yield* db
  430. .select({ data: SessionMessageTable.data })
  431. .from(SessionMessageTable)
  432. .where(eq(SessionMessageTable.type, "compaction"))
  433. .all()
  434. .pipe(Effect.orDie),
  435. ).toEqual([{ data: expect.objectContaining({ status: "running", summary: "", recent: "recent context" }) }])
  436. yield* events.publish(SessionEvent.Compaction.Ended, {
  437. sessionID,
  438. reason: "manual",
  439. text: "summary",
  440. recent: "recent context",
  441. })
  442. const rows = yield* db
  443. .select()
  444. .from(SessionMessageTable)
  445. .where(eq(SessionMessageTable.session_id, sessionID))
  446. .orderBy(asc(SessionMessageTable.seq))
  447. .all()
  448. .pipe(Effect.orDie)
  449. const messages = rows.map((row) =>
  450. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
  451. )
  452. expect(messages.map((message) => message.type)).toEqual([
  453. "agent-switched",
  454. "model-switched",
  455. "synthetic",
  456. "shell",
  457. "compaction",
  458. ])
  459. expect(messages.find((message) => message.type === "synthetic")).toMatchObject({
  460. text: "synthetic context",
  461. metadata: { source: "projector-test" },
  462. })
  463. expect(messages.find((message) => message.type === "model-switched")).toMatchObject({ previous: previousModel })
  464. expect(messages.find((message) => message.type === "shell")).toMatchObject({
  465. command: "pwd",
  466. status: "exited",
  467. exit: 0,
  468. output: { output: "/project", truncated: false },
  469. time: { completed: DateTime.makeUnsafe(0) },
  470. })
  471. expect(messages.find((message) => message.type === "compaction")).toMatchObject({
  472. summary: "summary",
  473. recent: "recent context",
  474. })
  475. expect(
  476. yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie),
  477. ).toMatchObject({
  478. agent: "build",
  479. model,
  480. time_updated: DateTime.toEpochMillis(created),
  481. })
  482. }),
  483. )
  484. it.effect("rejects distinct creator events that reuse one projected message ID", () =>
  485. Effect.gen(function* () {
  486. const { db } = yield* Database.Service
  487. yield* db
  488. .insert(ProjectTable)
  489. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  490. .run()
  491. .pipe(Effect.orDie)
  492. yield* db
  493. .insert(SessionTable)
  494. .values({
  495. id: sessionID,
  496. project_id: Project.ID.global,
  497. slug: "test",
  498. directory: "/project",
  499. title: "test",
  500. version: "test",
  501. })
  502. .run()
  503. .pipe(Effect.orDie)
  504. const events = yield* EventV2.Service
  505. const id = SessionMessage.ID.make("msg_creator_collision")
  506. const { id: _, type, ...data } = encodeMessage({ id, type: "synthetic", text: "existing", time: { created } })
  507. yield* db
  508. .insert(SessionMessageTable)
  509. .values({ id, session_id: sessionID, type, seq: 0, time_created: 0, data })
  510. .run()
  511. const exit = yield* events
  512. .publish(SessionEvent.Step.Started, {
  513. sessionID,
  514. assistantMessageID: id,
  515. agent: build,
  516. model,
  517. })
  518. .pipe(Effect.exit)
  519. expect(exit._tag).toBe("Failure")
  520. expect(
  521. yield* db.select().from(SessionMessageTable).where(eq(SessionMessageTable.id, id)).get().pipe(Effect.orDie),
  522. ).toMatchObject({ type: "synthetic" })
  523. }),
  524. )
  525. it.effect("does not revive a stale incomplete in-memory assistant projection", () =>
  526. Effect.gen(function* () {
  527. const stale = SessionMessage.Assistant.make({
  528. id: SessionMessage.ID.make("msg_assistant_stale"),
  529. type: "assistant",
  530. agent: build,
  531. model,
  532. content: [],
  533. time: { created },
  534. })
  535. const completed = SessionMessage.Assistant.make({
  536. id: SessionMessage.ID.make("msg_assistant_completed"),
  537. type: "assistant",
  538. agent: build,
  539. model,
  540. content: [],
  541. time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
  542. })
  543. expect(
  544. yield* SessionMessageUpdater.memory({ messages: [stale, completed] }).getCurrentAssistant(),
  545. ).toBeUndefined()
  546. }),
  547. )
  548. it.effect("projects retry state and clears it at the next step or execution terminal", () =>
  549. Effect.gen(function* () {
  550. const { db } = yield* Database.Service
  551. yield* db
  552. .insert(ProjectTable)
  553. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  554. .run()
  555. .pipe(Effect.orDie)
  556. yield* db
  557. .insert(SessionTable)
  558. .values({
  559. id: sessionID,
  560. project_id: Project.ID.global,
  561. slug: "test",
  562. directory: "/project",
  563. title: "test",
  564. version: "test",
  565. })
  566. .run()
  567. .pipe(Effect.orDie)
  568. const events = yield* EventV2.Service
  569. const first = SessionMessage.ID.make("msg_retry_first")
  570. const second = SessionMessage.ID.make("msg_retry_second")
  571. yield* events.publish(SessionEvent.Step.Started, { sessionID, assistantMessageID: first, agent: build, model })
  572. yield* events.publish(SessionEvent.RetryScheduled, {
  573. sessionID,
  574. assistantMessageID: first,
  575. attempt: 2,
  576. at: 2_000,
  577. error: { type: "provider.transport", message: "Disconnected" },
  578. })
  579. const decode = (row: typeof SessionMessageTable.$inferSelect) =>
  580. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type })
  581. const firstRow = yield* db
  582. .select()
  583. .from(SessionMessageTable)
  584. .where(eq(SessionMessageTable.id, first))
  585. .get()
  586. .pipe(Effect.orDie)
  587. const projected = firstRow ?? (yield* Effect.die(new Error("Missing retry projection")))
  588. expect(decode(projected)).toMatchObject({
  589. retry: { attempt: 2, at: DateTime.makeUnsafe(2_000), error: { type: "provider.transport" } },
  590. })
  591. yield* events.publish(SessionEvent.Step.Started, { sessionID, assistantMessageID: second, agent: build, model })
  592. yield* events.publish(SessionEvent.RetryScheduled, {
  593. sessionID,
  594. assistantMessageID: second,
  595. attempt: 3,
  596. at: 6_000,
  597. error: { type: "provider.internal", message: "Unavailable" },
  598. })
  599. yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID, reason: "shutdown" })
  600. const rows = yield* db
  601. .select()
  602. .from(SessionMessageTable)
  603. .where(eq(SessionMessageTable.session_id, sessionID))
  604. .orderBy(asc(SessionMessageTable.seq))
  605. .all()
  606. .pipe(Effect.orDie)
  607. expect(decode(rows[0])).not.toHaveProperty("retry")
  608. expect(decode(rows[1])).not.toHaveProperty("retry")
  609. }),
  610. )
  611. it.effect("updates only the newest incomplete assistant projection", () =>
  612. Effect.gen(function* () {
  613. const { db } = yield* Database.Service
  614. yield* db
  615. .insert(ProjectTable)
  616. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  617. .run()
  618. .pipe(Effect.orDie)
  619. yield* db
  620. .insert(SessionTable)
  621. .values({
  622. id: sessionID,
  623. project_id: Project.ID.global,
  624. slug: "test",
  625. directory: "/project",
  626. title: "test",
  627. version: "test",
  628. })
  629. .run()
  630. .pipe(Effect.orDie)
  631. yield* db
  632. .insert(SessionMessageTable)
  633. .values([
  634. assistantRow(SessionMessage.ID.make("msg_assistant_1"), 0),
  635. assistantRow(SessionMessage.ID.make("msg_assistant_2"), 1),
  636. ])
  637. .run()
  638. .pipe(Effect.orDie)
  639. const service = yield* EventV2.Service
  640. const usageUpdated = yield* service
  641. .subscribe(SessionEvent.UsageUpdated)
  642. .pipe(Stream.runHead, Effect.forkScoped({ startImmediately: true }))
  643. yield* service.publish(SessionEvent.Step.Ended, {
  644. sessionID,
  645. assistantMessageID: SessionMessage.ID.make("msg_assistant_2"),
  646. finish: "stop",
  647. cost: Money.USD.make(1.25),
  648. tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
  649. })
  650. const rows = yield* db
  651. .select()
  652. .from(SessionMessageTable)
  653. .where(eq(SessionMessageTable.session_id, sessionID))
  654. .orderBy(asc(SessionMessageTable.id))
  655. .all()
  656. .pipe(Effect.orDie)
  657. const messages = rows.map((row) =>
  658. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
  659. )
  660. expect(messages[0]).not.toHaveProperty("time.completed")
  661. expect(messages[1]).toMatchObject({
  662. type: "assistant",
  663. finish: "stop",
  664. cost: Money.USD.make(1.25),
  665. tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
  666. time: { completed: DateTime.makeUnsafe(0) },
  667. })
  668. expect(
  669. yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie),
  670. ).toMatchObject({
  671. cost: 1.25,
  672. tokens_input: 10,
  673. tokens_output: 4,
  674. tokens_reasoning: 2,
  675. tokens_cache_read: 3,
  676. tokens_cache_write: 1,
  677. })
  678. expect(Option.getOrThrow(yield* Fiber.join(usageUpdated)).data).toEqual({
  679. sessionID,
  680. cost: Money.USD.make(1.25),
  681. tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
  682. })
  683. }),
  684. )
  685. it.effect("does not revive a stale incomplete assistant projection", () =>
  686. Effect.gen(function* () {
  687. const { db } = yield* Database.Service
  688. yield* db
  689. .insert(ProjectTable)
  690. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  691. .run()
  692. .pipe(Effect.orDie)
  693. yield* db
  694. .insert(SessionTable)
  695. .values({
  696. id: sessionID,
  697. project_id: Project.ID.global,
  698. slug: "test",
  699. directory: "/project",
  700. title: "test",
  701. version: "test",
  702. })
  703. .run()
  704. .pipe(Effect.orDie)
  705. yield* db
  706. .insert(SessionMessageTable)
  707. .values([
  708. assistantRow(SessionMessage.ID.make("msg_assistant_stale"), 0),
  709. assistantRow(SessionMessage.ID.make("msg_assistant_completed"), 1, {
  710. created: DateTime.makeUnsafe(1),
  711. completed: DateTime.makeUnsafe(2),
  712. }),
  713. ])
  714. .run()
  715. .pipe(Effect.orDie)
  716. const service = yield* EventV2.Service
  717. yield* service.publish(SessionEvent.Text.Started, {
  718. sessionID,
  719. assistantMessageID: SessionMessage.ID.make("msg_assistant_completed"),
  720. ordinal: 0,
  721. })
  722. const rows = yield* db
  723. .select()
  724. .from(SessionMessageTable)
  725. .where(eq(SessionMessageTable.session_id, sessionID))
  726. .orderBy(asc(SessionMessageTable.id))
  727. .all()
  728. .pipe(Effect.orDie)
  729. const messages = rows.map((row) =>
  730. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
  731. )
  732. expect(messages).toEqual([
  733. SessionMessage.Assistant.make({
  734. id: SessionMessage.ID.make("msg_assistant_completed"),
  735. type: "assistant",
  736. agent: build,
  737. model,
  738. content: [SessionMessage.AssistantText.make({ type: "text", text: "" })],
  739. time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
  740. }),
  741. SessionMessage.Assistant.make({
  742. id: SessionMessage.ID.make("msg_assistant_stale"),
  743. type: "assistant",
  744. agent: build,
  745. model,
  746. content: [],
  747. time: { created },
  748. }),
  749. ])
  750. }),
  751. )
  752. })