session-projector.test.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458
  1. import { describe, expect } from "bun:test"
  2. import { DateTime, Effect, Layer, Schema } from "effect"
  3. import { asc, eq } from "drizzle-orm"
  4. import { Database } from "@opencode-ai/core/database/database"
  5. import { EventV2 } from "@opencode-ai/core/event"
  6. import { ModelV2 } from "@opencode-ai/core/model"
  7. import { Project } from "@opencode-ai/core/project"
  8. import { ProjectTable } from "@opencode-ai/core/project/sql"
  9. import { ProviderV2 } from "@opencode-ai/core/provider"
  10. import { AbsolutePath } from "@opencode-ai/core/schema"
  11. import { SessionV2 } from "@opencode-ai/core/session"
  12. import { SessionEvent } from "@opencode-ai/core/session/event"
  13. import { SessionMessage } from "@opencode-ai/core/session/message"
  14. import { Prompt } from "@opencode-ai/core/session/prompt"
  15. import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater"
  16. import { SessionProjector } from "@opencode-ai/core/session/projector"
  17. import { SessionExecution } from "@opencode-ai/core/session/execution"
  18. import { SessionInput } from "@opencode-ai/core/session/input"
  19. import { SessionStore } from "@opencode-ai/core/session/store"
  20. import { SessionInputTable, SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql"
  21. import { testEffect } from "./lib/effect"
  22. const database = Database.layerFromPath(":memory:")
  23. const events = EventV2.layer.pipe(Layer.provide(database))
  24. const projector = SessionProjector.layer.pipe(Layer.provide(events), Layer.provide(database))
  25. const it = testEffect(Layer.mergeAll(database, events, projector))
  26. const sessionID = SessionV2.ID.make("ses_projector_test")
  27. const created = DateTime.makeUnsafe(0)
  28. const model = { id: ModelV2.ID.make("model"), providerID: ProviderV2.ID.make("provider") }
  29. const encodeMessage = Schema.encodeSync(SessionMessage.Message)
  30. const assistantRow = (
  31. id: SessionMessage.ID,
  32. seq: number,
  33. time: { created: DateTime.Utc; completed?: DateTime.Utc } = { created },
  34. ) => {
  35. const {
  36. id: _,
  37. type,
  38. ...data
  39. } = encodeMessage(new SessionMessage.Assistant({ id, type: "assistant", agent: "build", model, content: [], time }))
  40. return { id, session_id: sessionID, type, seq, time_created: DateTime.toEpochMillis(time.created), data }
  41. }
  42. describe("SessionProjector", () => {
  43. it.effect("orders projected messages and context by durable aggregate sequence", () =>
  44. Effect.gen(function* () {
  45. const { db } = yield* Database.Service
  46. yield* db
  47. .insert(ProjectTable)
  48. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  49. .run()
  50. .pipe(Effect.orDie)
  51. yield* db
  52. .insert(SessionTable)
  53. .values({
  54. id: sessionID,
  55. project_id: Project.ID.global,
  56. slug: "test",
  57. directory: "/project",
  58. title: "test",
  59. version: "test",
  60. })
  61. .run()
  62. .pipe(Effect.orDie)
  63. const events = yield* EventV2.Service
  64. yield* events.publish(
  65. SessionEvent.Prompted,
  66. { sessionID, timestamp: created, prompt: new Prompt({ text: "first" }), delivery: "steer" },
  67. { id: SessionMessage.ID.make("evt_z") },
  68. )
  69. yield* events.publish(
  70. SessionEvent.Prompted,
  71. { sessionID, timestamp: created, prompt: new Prompt({ text: "second" }), delivery: "steer" },
  72. { id: SessionMessage.ID.make("evt_a") },
  73. )
  74. const sessions = yield* SessionV2.Service
  75. const firstPage = yield* sessions.messages({ sessionID, limit: 1, order: "asc" })
  76. expect(firstPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["first"])
  77. const secondPage = yield* sessions.messages({
  78. sessionID,
  79. limit: 1,
  80. order: "asc",
  81. cursor: { id: firstPage[0]!.id, direction: "next" },
  82. })
  83. expect(secondPage.map((message) => (message.type === "user" ? message.text : message.type))).toEqual(["second"])
  84. expect(
  85. (yield* sessions.messages({
  86. sessionID,
  87. limit: 1,
  88. order: "asc",
  89. cursor: { id: secondPage[0]!.id, direction: "previous" },
  90. })).map((message) => (message.type === "user" ? message.text : message.type)),
  91. ).toEqual(["first"])
  92. expect(
  93. (yield* sessions.context(sessionID)).map((message) => (message.type === "user" ? message.text : message.type)),
  94. ).toEqual(["first", "second"])
  95. }).pipe(
  96. Effect.provide(
  97. SessionV2.layer.pipe(
  98. Layer.provide(events),
  99. Layer.provide(database),
  100. Layer.provide(Project.defaultLayer),
  101. Layer.provide(SessionStore.layer.pipe(Layer.provide(database))),
  102. Layer.provide(SessionExecution.noopLayer),
  103. ),
  104. ),
  105. ),
  106. )
  107. it.effect("marks an admitted inbox row promoted with the Prompted event sequence", () =>
  108. Effect.gen(function* () {
  109. const { db } = yield* Database.Service
  110. yield* db
  111. .insert(ProjectTable)
  112. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  113. .run()
  114. .pipe(Effect.orDie)
  115. yield* db
  116. .insert(SessionTable)
  117. .values({
  118. id: sessionID,
  119. project_id: Project.ID.global,
  120. slug: "test",
  121. directory: "/project",
  122. title: "test",
  123. version: "test",
  124. })
  125. .run()
  126. .pipe(Effect.orDie)
  127. const events = yield* EventV2.Service
  128. const id = SessionMessage.ID.make("evt_admitted")
  129. yield* SessionInput.admit(db, { id, sessionID, prompt: new Prompt({ text: "promote me" }), delivery: "steer" })
  130. const event = yield* events.publish(
  131. SessionEvent.Prompted,
  132. { sessionID, timestamp: created, prompt: new Prompt({ text: "promote me" }), delivery: "steer" },
  133. { id },
  134. )
  135. expect(
  136. yield* db.select().from(SessionInputTable).where(eq(SessionInputTable.id, id)).get().pipe(Effect.orDie),
  137. ).toMatchObject({ promoted_seq: event.seq })
  138. }),
  139. )
  140. it.effect("projects durable context messages supported by the updater", () =>
  141. Effect.gen(function* () {
  142. const { db } = yield* Database.Service
  143. yield* db
  144. .insert(ProjectTable)
  145. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  146. .run()
  147. .pipe(Effect.orDie)
  148. yield* db
  149. .insert(SessionTable)
  150. .values({
  151. id: sessionID,
  152. project_id: Project.ID.global,
  153. slug: "test",
  154. directory: "/project",
  155. title: "test",
  156. version: "test",
  157. })
  158. .run()
  159. .pipe(Effect.orDie)
  160. const events = yield* EventV2.Service
  161. yield* events.publish(SessionEvent.AgentSwitched, { sessionID, timestamp: created, agent: "build" })
  162. yield* events.publish(SessionEvent.ModelSwitched, { sessionID, timestamp: created, model })
  163. yield* events.publish(SessionEvent.Synthetic, { sessionID, timestamp: created, text: "synthetic context" })
  164. yield* events.publish(SessionEvent.Shell.Started, {
  165. sessionID,
  166. timestamp: created,
  167. callID: "shell-1",
  168. command: "pwd",
  169. })
  170. yield* events.publish(SessionEvent.Shell.Ended, {
  171. sessionID,
  172. timestamp: DateTime.makeUnsafe(1),
  173. callID: "shell-1",
  174. output: "/project",
  175. })
  176. yield* events.publish(SessionEvent.Compaction.Started, { sessionID, timestamp: created, reason: "manual" })
  177. yield* events.publish(SessionEvent.Compaction.Delta, { sessionID, timestamp: created, text: "partial" })
  178. yield* events.publish(SessionEvent.Compaction.Ended, {
  179. sessionID,
  180. timestamp: DateTime.makeUnsafe(1),
  181. text: "summary",
  182. include: "msg-1",
  183. })
  184. const rows = yield* db
  185. .select()
  186. .from(SessionMessageTable)
  187. .where(eq(SessionMessageTable.session_id, sessionID))
  188. .orderBy(asc(SessionMessageTable.seq))
  189. .all()
  190. .pipe(Effect.orDie)
  191. const messages = rows.map((row) =>
  192. Schema.decodeUnknownSync(SessionMessage.Message)({ ...row.data, id: row.id, type: row.type }),
  193. )
  194. expect(messages.map((message) => message.type)).toEqual([
  195. "agent-switched",
  196. "model-switched",
  197. "synthetic",
  198. "shell",
  199. "compaction",
  200. ])
  201. expect(messages.find((message) => message.type === "shell")).toMatchObject({
  202. output: "/project",
  203. time: { completed: DateTime.makeUnsafe(1) },
  204. })
  205. expect(messages.find((message) => message.type === "compaction")).toMatchObject({
  206. summary: "summary",
  207. include: "msg-1",
  208. })
  209. expect(
  210. yield* db.select().from(SessionTable).where(eq(SessionTable.id, sessionID)).get().pipe(Effect.orDie),
  211. ).toMatchObject({
  212. agent: "build",
  213. model,
  214. time_updated: DateTime.toEpochMillis(created),
  215. })
  216. }),
  217. )
  218. it.effect("rejects a Prompted event that conflicts with an admitted inbox row", () =>
  219. Effect.gen(function* () {
  220. const { db } = yield* Database.Service
  221. yield* db
  222. .insert(ProjectTable)
  223. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  224. .run()
  225. .pipe(Effect.orDie)
  226. yield* db
  227. .insert(SessionTable)
  228. .values({
  229. id: sessionID,
  230. project_id: Project.ID.global,
  231. slug: "test",
  232. directory: "/project",
  233. title: "test",
  234. version: "test",
  235. })
  236. .run()
  237. .pipe(Effect.orDie)
  238. const events = yield* EventV2.Service
  239. const id = SessionMessage.ID.make("evt_conflict")
  240. yield* SessionInput.admit(db, { id, sessionID, prompt: new Prompt({ text: "admitted" }), delivery: "steer" })
  241. const exit = yield* events
  242. .publish(
  243. SessionEvent.Prompted,
  244. { sessionID, timestamp: created, prompt: new Prompt({ text: "different" }), delivery: "steer" },
  245. { id },
  246. )
  247. .pipe(Effect.exit)
  248. expect(String(exit)).toContain("Prompt projection conflicts with admitted input")
  249. expect(
  250. yield* db.select().from(SessionInputTable).where(eq(SessionInputTable.id, id)).get().pipe(Effect.orDie),
  251. ).toMatchObject({ promoted_seq: null })
  252. }),
  253. )
  254. it.effect("rejects a Prompted delivery mode that conflicts with an admitted inbox row", () =>
  255. Effect.gen(function* () {
  256. const { db } = yield* Database.Service
  257. yield* db
  258. .insert(ProjectTable)
  259. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  260. .run()
  261. .pipe(Effect.orDie)
  262. yield* db
  263. .insert(SessionTable)
  264. .values({
  265. id: sessionID,
  266. project_id: Project.ID.global,
  267. slug: "test",
  268. directory: "/project",
  269. title: "test",
  270. version: "test",
  271. })
  272. .run()
  273. .pipe(Effect.orDie)
  274. const events = yield* EventV2.Service
  275. const id = SessionMessage.ID.make("evt_delivery_conflict")
  276. const prompt = new Prompt({ text: "admitted" })
  277. yield* SessionInput.admit(db, { id, sessionID, prompt, delivery: "queue" })
  278. const exit = yield* events
  279. .publish(SessionEvent.Prompted, { sessionID, timestamp: created, prompt, delivery: "steer" }, { id })
  280. .pipe(Effect.exit)
  281. expect(String(exit)).toContain("Prompt projection conflicts with admitted input")
  282. expect(
  283. yield* db.select().from(SessionInputTable).where(eq(SessionInputTable.id, id)).get().pipe(Effect.orDie),
  284. ).toMatchObject({ delivery: "queue", promoted_seq: null })
  285. }),
  286. )
  287. it.effect("does not revive a stale incomplete in-memory assistant projection", () =>
  288. Effect.gen(function* () {
  289. const stale = new SessionMessage.Assistant({
  290. id: SessionMessage.ID.make("evt_assistant_stale"),
  291. type: "assistant",
  292. agent: "build",
  293. model,
  294. content: [],
  295. time: { created },
  296. })
  297. const completed = new SessionMessage.Assistant({
  298. id: SessionMessage.ID.make("evt_assistant_completed"),
  299. type: "assistant",
  300. agent: "build",
  301. model,
  302. content: [],
  303. time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
  304. })
  305. expect(
  306. yield* SessionMessageUpdater.memory({ messages: [stale, completed] }).getCurrentAssistant(),
  307. ).toBeUndefined()
  308. }),
  309. )
  310. it.effect("updates only the newest incomplete assistant projection", () =>
  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. yield* db
  331. .insert(SessionMessageTable)
  332. .values([
  333. assistantRow(SessionMessage.ID.make("evt_assistant_1"), 0),
  334. assistantRow(SessionMessage.ID.make("evt_assistant_2"), 1),
  335. ])
  336. .run()
  337. .pipe(Effect.orDie)
  338. const service = yield* EventV2.Service
  339. yield* service.publish(SessionEvent.Step.Ended, {
  340. sessionID,
  341. timestamp: DateTime.makeUnsafe(1),
  342. assistantMessageID: SessionMessage.ID.make("evt_assistant_2"),
  343. finish: "stop",
  344. cost: 0,
  345. tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
  346. })
  347. const rows = yield* db
  348. .select()
  349. .from(SessionMessageTable)
  350. .where(eq(SessionMessageTable.session_id, sessionID))
  351. .orderBy(asc(SessionMessageTable.id))
  352. .all()
  353. .pipe(Effect.orDie)
  354. const messages = rows.map((row) =>
  355. Schema.decodeUnknownSync(SessionMessage.Message)({ ...row.data, id: row.id, type: row.type }),
  356. )
  357. expect(messages[0]).not.toHaveProperty("time.completed")
  358. expect(messages[1]).toMatchObject({
  359. type: "assistant",
  360. finish: "stop",
  361. time: { completed: DateTime.makeUnsafe(1) },
  362. })
  363. }),
  364. )
  365. it.effect("does not revive a stale incomplete assistant projection", () =>
  366. Effect.gen(function* () {
  367. const { db } = yield* Database.Service
  368. yield* db
  369. .insert(ProjectTable)
  370. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  371. .run()
  372. .pipe(Effect.orDie)
  373. yield* db
  374. .insert(SessionTable)
  375. .values({
  376. id: sessionID,
  377. project_id: Project.ID.global,
  378. slug: "test",
  379. directory: "/project",
  380. title: "test",
  381. version: "test",
  382. })
  383. .run()
  384. .pipe(Effect.orDie)
  385. yield* db
  386. .insert(SessionMessageTable)
  387. .values([
  388. assistantRow(SessionMessage.ID.make("evt_assistant_stale"), 0),
  389. assistantRow(SessionMessage.ID.make("evt_assistant_completed"), 1, {
  390. created: DateTime.makeUnsafe(1),
  391. completed: DateTime.makeUnsafe(2),
  392. }),
  393. ])
  394. .run()
  395. .pipe(Effect.orDie)
  396. const service = yield* EventV2.Service
  397. yield* service.publish(SessionEvent.Text.Started, {
  398. sessionID,
  399. timestamp: DateTime.makeUnsafe(3),
  400. textID: "text-stale",
  401. })
  402. const rows = yield* db
  403. .select()
  404. .from(SessionMessageTable)
  405. .where(eq(SessionMessageTable.session_id, sessionID))
  406. .orderBy(asc(SessionMessageTable.id))
  407. .all()
  408. .pipe(Effect.orDie)
  409. const messages = rows.map((row) =>
  410. Schema.decodeUnknownSync(SessionMessage.Message)({ ...row.data, id: row.id, type: row.type }),
  411. )
  412. expect(messages).toEqual([
  413. new SessionMessage.Assistant({
  414. id: SessionMessage.ID.make("evt_assistant_completed"),
  415. type: "assistant",
  416. agent: "build",
  417. model,
  418. content: [],
  419. time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) },
  420. }),
  421. new SessionMessage.Assistant({
  422. id: SessionMessage.ID.make("evt_assistant_stale"),
  423. type: "assistant",
  424. agent: "build",
  425. model,
  426. content: [],
  427. time: { created },
  428. }),
  429. ])
  430. }),
  431. )
  432. })