session-projector.test.ts 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827
  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/util/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 { SessionPending } from "@opencode-ai/core/session/pending"
  24. import { Shell } from "@opencode-ai/schema/shell"
  25. import {
  26. InstructionStateTable,
  27. SessionPendingTable,
  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* SessionPending.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* SessionPending.compaction(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(InstructionStateTable)
  200. .values({
  201. session_id: sessionID,
  202. epoch_start: 0,
  203. through_seq: 0,
  204. initial_values: {},
  205. current_values: {},
  206. })
  207. .run()
  208. const events = yield* EventV2.Service
  209. yield* events.publish(SessionEvent.RevertEvent.Staged, {
  210. sessionID,
  211. revert: { messageID: boundary, snapshot: Snapshot.ID.make("tree"), files: [] },
  212. })
  213. expect((yield* db.select({ revert: SessionTable.revert }).from(SessionTable).get())?.revert).toMatchObject({
  214. messageID: boundary,
  215. snapshot: "tree",
  216. files: [],
  217. })
  218. yield* events.publish(SessionEvent.RevertEvent.Cleared, { sessionID })
  219. expect((yield* db.select({ revert: SessionTable.revert }).from(SessionTable).get())?.revert).toBeNull()
  220. yield* events.publish(SessionEvent.RevertEvent.Staged, {
  221. sessionID,
  222. revert: { messageID: boundary, files: [] },
  223. })
  224. yield* events.publish(SessionEvent.RevertEvent.Committed, {
  225. sessionID,
  226. to: boundary,
  227. })
  228. expect(
  229. (yield* db.select({ id: SessionMessageTable.id }).from(SessionMessageTable).all()).map((row) => row.id),
  230. ).toEqual([earlier])
  231. expect(yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get()).toMatchObject({
  232. cost: Money.USD.make(1.25),
  233. tokens_input: 10,
  234. tokens_output: 4,
  235. tokens_reasoning: 2,
  236. tokens_cache_read: 3,
  237. tokens_cache_write: 1,
  238. })
  239. // A committed revert resets the fold cache so the next boundary establishes a new epoch.
  240. expect(yield* db.select().from(InstructionStateTable).get().pipe(Effect.orDie)).toBeUndefined()
  241. }),
  242. )
  243. it.effect("orders projected messages and context by durable aggregate sequence", () =>
  244. Effect.gen(function* () {
  245. const { db } = yield* Database.Service
  246. yield* db
  247. .insert(ProjectTable)
  248. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  249. .run()
  250. .pipe(Effect.orDie)
  251. yield* db
  252. .insert(SessionTable)
  253. .values({
  254. id: sessionID,
  255. project_id: Project.ID.global,
  256. slug: "test",
  257. directory: "/project",
  258. title: "test",
  259. version: "test",
  260. })
  261. .run()
  262. .pipe(Effect.orDie)
  263. const events = yield* EventV2.Service
  264. yield* events.publish(SessionEvent.InputAdmitted, {
  265. sessionID,
  266. inputID: SessionMessage.ID.make("msg_first"),
  267. input: { type: "user", data: { text: "first" }, delivery: "steer" },
  268. })
  269. yield* events.publish(
  270. SessionEvent.InputPromoted,
  271. {
  272. sessionID,
  273. inputID: SessionMessage.ID.make("msg_first"),
  274. },
  275. { id: EventV2.ID.make("evt_z") },
  276. )
  277. yield* events.publish(SessionEvent.InputAdmitted, {
  278. sessionID,
  279. inputID: SessionMessage.ID.make("msg_second"),
  280. input: { type: "user", data: { text: "second" }, delivery: "steer" },
  281. })
  282. yield* events.publish(
  283. SessionEvent.InputPromoted,
  284. {
  285. sessionID,
  286. inputID: SessionMessage.ID.make("msg_second"),
  287. },
  288. { id: EventV2.ID.make("evt_a") },
  289. )
  290. const sessions = yield* SessionV2.Service
  291. const firstPage = yield* sessions.messages({ sessionID, limit: 1, order: "asc" })
  292. expect(firstPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["first"])
  293. const secondPage = yield* sessions.messages({
  294. sessionID,
  295. limit: 1,
  296. order: "asc",
  297. cursor: { id: firstPage[0]!.id, direction: "next" },
  298. })
  299. expect(secondPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["second"])
  300. expect(
  301. (yield* sessions.messages({
  302. sessionID,
  303. limit: 1,
  304. order: "asc",
  305. cursor: { id: secondPage[0]!.id, direction: "previous" },
  306. })).map((message) => (message.type === "user" ? message.text : message.type)),
  307. ).toEqual(["first"])
  308. expect(
  309. (yield* sessions.context(sessionID)).map((message) => (message.type === "user" ? message.text : message.type)),
  310. ).toEqual(["first", "second"])
  311. }).pipe(Effect.provide(sessionsLayer)),
  312. )
  313. it.effect("consumes the pending row and projects the message at promotion", () =>
  314. Effect.gen(function* () {
  315. const { db } = yield* Database.Service
  316. yield* db
  317. .insert(ProjectTable)
  318. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  319. .run()
  320. .pipe(Effect.orDie)
  321. yield* db
  322. .insert(SessionTable)
  323. .values({
  324. id: sessionID,
  325. project_id: Project.ID.global,
  326. slug: "test",
  327. directory: "/project",
  328. title: "test",
  329. version: "test",
  330. })
  331. .run()
  332. .pipe(Effect.orDie)
  333. const events = yield* EventV2.Service
  334. const id = SessionMessage.ID.make("msg_admitted")
  335. const admitted = yield* SessionPending.admit(db, events, {
  336. id,
  337. sessionID,
  338. input: { type: "user", data: { text: "promote me" }, delivery: "steer" },
  339. })
  340. if (!admitted) return yield* Effect.die("Prompt admission failed")
  341. const event = yield* events.publish(SessionEvent.InputPromoted, {
  342. sessionID,
  343. inputID: id,
  344. })
  345. expect(
  346. yield* db.select().from(SessionPendingTable).where(eq(SessionPendingTable.id, id)).get().pipe(Effect.orDie),
  347. ).toBeUndefined()
  348. expect(
  349. yield* db.select().from(SessionMessageTable).where(eq(SessionMessageTable.id, id)).get().pipe(Effect.orDie),
  350. ).toMatchObject({ session_id: sessionID, type: "user", seq: event.durable?.seq })
  351. }),
  352. )
  353. it.effect("projects durable context messages supported by the updater", () =>
  354. Effect.gen(function* () {
  355. const { db } = yield* Database.Service
  356. yield* db
  357. .insert(ProjectTable)
  358. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  359. .run()
  360. .pipe(Effect.orDie)
  361. yield* db
  362. .insert(SessionTable)
  363. .values({
  364. id: sessionID,
  365. project_id: Project.ID.global,
  366. slug: "test",
  367. directory: "/project",
  368. title: "test",
  369. version: "test",
  370. model: previousModel,
  371. })
  372. .run()
  373. .pipe(Effect.orDie)
  374. const events = yield* EventV2.Service
  375. yield* events.publish(SessionEvent.AgentSelected, {
  376. sessionID,
  377. agent: build,
  378. })
  379. yield* events.publish(SessionEvent.ModelSelected, {
  380. sessionID,
  381. model,
  382. })
  383. yield* events.publish(SessionEvent.Synthetic, {
  384. sessionID,
  385. text: "synthetic context",
  386. metadata: { source: "projector-test" },
  387. })
  388. yield* events.publish(SessionEvent.Shell.Started, {
  389. sessionID,
  390. shell: Shell.Info.make({
  391. id: Shell.ID.make("sh_projector"),
  392. status: "running",
  393. command: "pwd",
  394. cwd: "/project",
  395. shell: "/bin/sh",
  396. file: "/tmp/sh_projector.out",
  397. metadata: {},
  398. time: { started: 0 },
  399. }),
  400. })
  401. yield* events.publish(SessionEvent.Shell.Ended, {
  402. sessionID,
  403. shell: Shell.Info.make({
  404. id: Shell.ID.make("sh_projector"),
  405. status: "exited",
  406. command: "pwd",
  407. cwd: "/project",
  408. shell: "/bin/sh",
  409. file: "/tmp/sh_projector.out",
  410. exit: 0,
  411. metadata: {},
  412. time: { started: 0, completed: 1 },
  413. }),
  414. output: { output: "/project", cursor: 8, size: 8, truncated: false },
  415. })
  416. yield* events.publish(SessionEvent.Compaction.Started, {
  417. sessionID,
  418. reason: "manual",
  419. recent: "recent context",
  420. })
  421. yield* events.publish(SessionEvent.Compaction.Delta, {
  422. sessionID,
  423. text: "partial",
  424. })
  425. expect(
  426. yield* db
  427. .select({ id: EventTable.id })
  428. .from(EventTable)
  429. .where(sql`${EventTable.type} like 'session.compaction.delta.%'`)
  430. .all()
  431. .pipe(Effect.orDie),
  432. ).toHaveLength(0)
  433. expect(
  434. yield* db
  435. .select({ data: SessionMessageTable.data })
  436. .from(SessionMessageTable)
  437. .where(eq(SessionMessageTable.type, "compaction"))
  438. .all()
  439. .pipe(Effect.orDie),
  440. ).toEqual([{ data: expect.objectContaining({ status: "running", summary: "", recent: "recent context" }) }])
  441. yield* events.publish(SessionEvent.Compaction.Ended, {
  442. sessionID,
  443. reason: "manual",
  444. text: "summary",
  445. recent: "recent context",
  446. })
  447. const rows = yield* db
  448. .select()
  449. .from(SessionMessageTable)
  450. .where(eq(SessionMessageTable.session_id, sessionID))
  451. .orderBy(asc(SessionMessageTable.seq))
  452. .all()
  453. .pipe(Effect.orDie)
  454. const messages = rows.map((row) =>
  455. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
  456. )
  457. expect(messages.map((message) => message.type)).toEqual([
  458. "agent-switched",
  459. "model-switched",
  460. "synthetic",
  461. "shell",
  462. "compaction",
  463. ])
  464. expect(messages.find((message) => message.type === "synthetic")).toMatchObject({
  465. text: "synthetic context",
  466. metadata: { source: "projector-test" },
  467. })
  468. expect(messages.find((message) => message.type === "model-switched")).toMatchObject({ previous: previousModel })
  469. expect(messages.find((message) => message.type === "shell")).toMatchObject({
  470. command: "pwd",
  471. status: "exited",
  472. exit: 0,
  473. output: { output: "/project", truncated: false },
  474. time: { completed: DateTime.makeUnsafe(0) },
  475. })
  476. expect(messages.find((message) => message.type === "compaction")).toMatchObject({
  477. summary: "summary",
  478. recent: "recent context",
  479. })
  480. expect(
  481. yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie),
  482. ).toMatchObject({
  483. agent: "build",
  484. model,
  485. time_updated: DateTime.toEpochMillis(created),
  486. })
  487. }),
  488. )
  489. it.effect("rejects distinct creator events that reuse one projected message ID", () =>
  490. Effect.gen(function* () {
  491. const { db } = yield* Database.Service
  492. yield* db
  493. .insert(ProjectTable)
  494. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  495. .run()
  496. .pipe(Effect.orDie)
  497. yield* db
  498. .insert(SessionTable)
  499. .values({
  500. id: sessionID,
  501. project_id: Project.ID.global,
  502. slug: "test",
  503. directory: "/project",
  504. title: "test",
  505. version: "test",
  506. })
  507. .run()
  508. .pipe(Effect.orDie)
  509. const events = yield* EventV2.Service
  510. const id = SessionMessage.ID.make("msg_creator_collision")
  511. const { id: _, type, ...data } = encodeMessage({ id, type: "synthetic", text: "existing", time: { created } })
  512. yield* db
  513. .insert(SessionMessageTable)
  514. .values({ id, session_id: sessionID, type, seq: 0, time_created: 0, data })
  515. .run()
  516. const exit = yield* events
  517. .publish(SessionEvent.Step.Started, {
  518. sessionID,
  519. assistantMessageID: id,
  520. agent: build,
  521. model,
  522. })
  523. .pipe(Effect.exit)
  524. expect(exit._tag).toBe("Failure")
  525. expect(
  526. yield* db.select().from(SessionMessageTable).where(eq(SessionMessageTable.id, id)).get().pipe(Effect.orDie),
  527. ).toMatchObject({ type: "synthetic" })
  528. }),
  529. )
  530. it.effect("does not revive a stale incomplete in-memory assistant projection", () =>
  531. Effect.gen(function* () {
  532. const stale = SessionMessage.Assistant.make({
  533. id: SessionMessage.ID.make("msg_assistant_stale"),
  534. type: "assistant",
  535. agent: build,
  536. model,
  537. content: [],
  538. time: { created },
  539. })
  540. const completed = SessionMessage.Assistant.make({
  541. id: SessionMessage.ID.make("msg_assistant_completed"),
  542. type: "assistant",
  543. agent: build,
  544. model,
  545. content: [],
  546. time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
  547. })
  548. expect(
  549. yield* SessionMessageUpdater.memory({ messages: [stale, completed] }).getCurrentAssistant(),
  550. ).toBeUndefined()
  551. }),
  552. )
  553. it.effect("projects retry state and clears it at the next step or execution terminal", () =>
  554. Effect.gen(function* () {
  555. const { db } = yield* Database.Service
  556. yield* db
  557. .insert(ProjectTable)
  558. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  559. .run()
  560. .pipe(Effect.orDie)
  561. yield* db
  562. .insert(SessionTable)
  563. .values({
  564. id: sessionID,
  565. project_id: Project.ID.global,
  566. slug: "test",
  567. directory: "/project",
  568. title: "test",
  569. version: "test",
  570. })
  571. .run()
  572. .pipe(Effect.orDie)
  573. const events = yield* EventV2.Service
  574. const first = SessionMessage.ID.make("msg_retry_first")
  575. const second = SessionMessage.ID.make("msg_retry_second")
  576. yield* events.publish(SessionEvent.Step.Started, { sessionID, assistantMessageID: first, agent: build, model })
  577. yield* events.publish(SessionEvent.RetryScheduled, {
  578. sessionID,
  579. assistantMessageID: first,
  580. attempt: 2,
  581. at: 2_000,
  582. error: { type: "provider.transport", message: "Disconnected" },
  583. })
  584. const decode = (row: typeof SessionMessageTable.$inferSelect) =>
  585. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type })
  586. const firstRow = yield* db
  587. .select()
  588. .from(SessionMessageTable)
  589. .where(eq(SessionMessageTable.id, first))
  590. .get()
  591. .pipe(Effect.orDie)
  592. const projected = firstRow ?? (yield* Effect.die(new Error("Missing retry projection")))
  593. expect(decode(projected)).toMatchObject({
  594. retry: { attempt: 2, at: DateTime.makeUnsafe(2_000), error: { type: "provider.transport" } },
  595. })
  596. yield* events.publish(SessionEvent.Step.Started, { sessionID, assistantMessageID: second, agent: build, model })
  597. yield* events.publish(SessionEvent.RetryScheduled, {
  598. sessionID,
  599. assistantMessageID: second,
  600. attempt: 3,
  601. at: 6_000,
  602. error: { type: "provider.internal", message: "Unavailable" },
  603. })
  604. yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID, reason: "shutdown" })
  605. const rows = yield* db
  606. .select()
  607. .from(SessionMessageTable)
  608. .where(eq(SessionMessageTable.session_id, sessionID))
  609. .orderBy(asc(SessionMessageTable.seq))
  610. .all()
  611. .pipe(Effect.orDie)
  612. expect(decode(rows[0])).not.toHaveProperty("retry")
  613. expect(decode(rows[1])).not.toHaveProperty("retry")
  614. }),
  615. )
  616. it.effect("does not infer restart continuation from lifecycle history", () =>
  617. Effect.gen(function* () {
  618. const { db } = yield* Database.Service
  619. yield* db
  620. .insert(ProjectTable)
  621. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  622. .run()
  623. .pipe(Effect.orDie)
  624. yield* db
  625. .insert(SessionTable)
  626. .values({
  627. id: sessionID,
  628. project_id: Project.ID.global,
  629. slug: "test",
  630. directory: "/project",
  631. title: "test",
  632. version: "test",
  633. })
  634. .run()
  635. .pipe(Effect.orDie)
  636. const events = yield* EventV2.Service
  637. const suspended = () =>
  638. db
  639. .select({ timeSuspended: SessionTable.time_suspended })
  640. .from(SessionTable)
  641. .where(eq(SessionTable.id, sessionID))
  642. .get()
  643. .pipe(Effect.orDie)
  644. yield* events.publish(SessionEvent.Execution.Interrupted, { sessionID, reason: "shutdown" })
  645. expect((yield* suspended())?.timeSuspended).toBeNull()
  646. yield* events.publish(SessionEvent.Execution.Started, { sessionID })
  647. expect((yield* suspended())?.timeSuspended).toBeNull()
  648. }),
  649. )
  650. it.effect("updates only the newest incomplete assistant projection", () =>
  651. Effect.gen(function* () {
  652. const { db } = yield* Database.Service
  653. yield* db
  654. .insert(ProjectTable)
  655. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  656. .run()
  657. .pipe(Effect.orDie)
  658. yield* db
  659. .insert(SessionTable)
  660. .values({
  661. id: sessionID,
  662. project_id: Project.ID.global,
  663. slug: "test",
  664. directory: "/project",
  665. title: "test",
  666. version: "test",
  667. })
  668. .run()
  669. .pipe(Effect.orDie)
  670. yield* db
  671. .insert(SessionMessageTable)
  672. .values([
  673. assistantRow(SessionMessage.ID.make("msg_assistant_1"), 0),
  674. assistantRow(SessionMessage.ID.make("msg_assistant_2"), 1),
  675. ])
  676. .run()
  677. .pipe(Effect.orDie)
  678. const service = yield* EventV2.Service
  679. const usageUpdated = yield* service
  680. .subscribe(SessionEvent.UsageUpdated)
  681. .pipe(Stream.runHead, Effect.forkScoped({ startImmediately: true }))
  682. yield* service.publish(SessionEvent.Step.Ended, {
  683. sessionID,
  684. assistantMessageID: SessionMessage.ID.make("msg_assistant_2"),
  685. finish: "stop",
  686. cost: Money.USD.make(1.25),
  687. tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
  688. })
  689. const rows = yield* db
  690. .select()
  691. .from(SessionMessageTable)
  692. .where(eq(SessionMessageTable.session_id, sessionID))
  693. .orderBy(asc(SessionMessageTable.id))
  694. .all()
  695. .pipe(Effect.orDie)
  696. const messages = rows.map((row) =>
  697. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
  698. )
  699. expect(messages[0]).not.toHaveProperty("time.completed")
  700. expect(messages[1]).toMatchObject({
  701. type: "assistant",
  702. finish: "stop",
  703. cost: Money.USD.make(1.25),
  704. tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
  705. time: { completed: DateTime.makeUnsafe(0) },
  706. })
  707. expect(
  708. yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie),
  709. ).toMatchObject({
  710. cost: 1.25,
  711. tokens_input: 10,
  712. tokens_output: 4,
  713. tokens_reasoning: 2,
  714. tokens_cache_read: 3,
  715. tokens_cache_write: 1,
  716. })
  717. expect(Option.getOrThrow(yield* Fiber.join(usageUpdated)).data).toEqual({
  718. sessionID,
  719. cost: Money.USD.make(1.25),
  720. tokens: { input: 10, output: 4, reasoning: 2, cache: { read: 3, write: 1 } },
  721. })
  722. }),
  723. )
  724. it.effect("does not revive a stale incomplete assistant projection", () =>
  725. Effect.gen(function* () {
  726. const { db } = yield* Database.Service
  727. yield* db
  728. .insert(ProjectTable)
  729. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  730. .run()
  731. .pipe(Effect.orDie)
  732. yield* db
  733. .insert(SessionTable)
  734. .values({
  735. id: sessionID,
  736. project_id: Project.ID.global,
  737. slug: "test",
  738. directory: "/project",
  739. title: "test",
  740. version: "test",
  741. })
  742. .run()
  743. .pipe(Effect.orDie)
  744. yield* db
  745. .insert(SessionMessageTable)
  746. .values([
  747. assistantRow(SessionMessage.ID.make("msg_assistant_stale"), 0),
  748. assistantRow(SessionMessage.ID.make("msg_assistant_completed"), 1, {
  749. created: DateTime.makeUnsafe(1),
  750. completed: DateTime.makeUnsafe(2),
  751. }),
  752. ])
  753. .run()
  754. .pipe(Effect.orDie)
  755. const service = yield* EventV2.Service
  756. yield* service.publish(SessionEvent.Text.Started, {
  757. sessionID,
  758. assistantMessageID: SessionMessage.ID.make("msg_assistant_completed"),
  759. ordinal: 0,
  760. })
  761. const rows = yield* db
  762. .select()
  763. .from(SessionMessageTable)
  764. .where(eq(SessionMessageTable.session_id, sessionID))
  765. .orderBy(asc(SessionMessageTable.id))
  766. .all()
  767. .pipe(Effect.orDie)
  768. const messages = rows.map((row) =>
  769. Schema.decodeUnknownSync(SessionMessage.Info)({ ...row.data, id: row.id, type: row.type }),
  770. )
  771. expect(messages).toEqual([
  772. SessionMessage.Assistant.make({
  773. id: SessionMessage.ID.make("msg_assistant_completed"),
  774. type: "assistant",
  775. agent: build,
  776. model,
  777. content: [SessionMessage.AssistantText.make({ type: "text", text: "" })],
  778. time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
  779. }),
  780. SessionMessage.Assistant.make({
  781. id: SessionMessage.ID.make("msg_assistant_stale"),
  782. type: "assistant",
  783. agent: build,
  784. model,
  785. content: [],
  786. time: { created },
  787. }),
  788. ])
  789. }),
  790. )
  791. })