session-projector.test.ts 27 KB

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