messages-pagination.test.ts 34 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057
  1. import { describe, expect, test } from "bun:test"
  2. import { LayerNode } from "@kirincode-ai/core/effect/layer-node"
  3. import { SessionV1 } from "@kirincode-ai/core/v1/session"
  4. import { SessionProjector } from "@kirincode-ai/core/session/projector"
  5. import { Effect, Option } from "effect"
  6. import { Session as SessionNs } from "@/session/session"
  7. import { MessageV2 } from "../../src/session/message-v2"
  8. import { MessageID, PartID, type SessionID } from "../../src/session/schema"
  9. import { NotFoundError } from "@/storage/storage"
  10. import { testEffect } from "../lib/effect"
  11. import { ProviderV2 } from "@kirincode-ai/core/provider"
  12. import { ModelV2 } from "@kirincode-ai/core/model"
  13. const it = testEffect(LayerNode.compile(LayerNode.group([SessionNs.node, MessageV2.node, SessionProjector.node])))
  14. const withSession = <A, E, R>(
  15. fn: (input: { session: SessionNs.Interface; sessionID: SessionID }) => Effect.Effect<A, E, R>,
  16. ) =>
  17. Effect.acquireUseRelease(
  18. Effect.gen(function* () {
  19. const session = yield* SessionNs.Service
  20. const created = yield* session.create({})
  21. return { session, sessionID: created.id }
  22. }),
  23. fn,
  24. (input) => input.session.remove(input.sessionID).pipe(Effect.ignore),
  25. )
  26. // Helper functions using Effect.gen
  27. const fill = Effect.fn("Test.fill")(function* (
  28. sessionID: SessionID,
  29. count: number,
  30. time = (i: number) => Date.now() + i,
  31. ) {
  32. const session = yield* SessionNs.Service
  33. const ids = [] as MessageID[]
  34. for (let i = 0; i < count; i++) {
  35. const id = MessageID.ascending()
  36. ids.push(id)
  37. yield* session.updateMessage({
  38. id,
  39. sessionID,
  40. role: "user",
  41. time: { created: time(i) },
  42. agent: "test",
  43. model: { providerID: "test", modelID: "test" },
  44. tools: {},
  45. mode: "",
  46. } as unknown as SessionV1.Info)
  47. yield* session.updatePart({
  48. id: PartID.ascending(),
  49. sessionID,
  50. messageID: id,
  51. type: "text",
  52. text: `m${i}`,
  53. })
  54. }
  55. return ids
  56. })
  57. const addUser = Effect.fn("Test.addUser")(function* (sessionID: SessionID, text?: string) {
  58. const session = yield* SessionNs.Service
  59. const id = MessageID.ascending()
  60. yield* session.updateMessage({
  61. id,
  62. sessionID,
  63. role: "user",
  64. time: { created: Date.now() },
  65. agent: "test",
  66. model: { providerID: "test", modelID: "test" },
  67. tools: {},
  68. mode: "",
  69. } as unknown as SessionV1.Info)
  70. if (text) {
  71. yield* session.updatePart({
  72. id: PartID.ascending(),
  73. sessionID,
  74. messageID: id,
  75. type: "text",
  76. text,
  77. })
  78. }
  79. return id
  80. })
  81. const addAssistant = Effect.fn("Test.addAssistant")(function* (
  82. sessionID: SessionID,
  83. parentID: MessageID,
  84. opts?: { summary?: boolean; finish?: string; error?: SessionV1.Assistant["error"] },
  85. ) {
  86. const session = yield* SessionNs.Service
  87. const id = MessageID.ascending()
  88. yield* session.updateMessage({
  89. id,
  90. sessionID,
  91. role: "assistant",
  92. time: { created: Date.now() },
  93. parentID,
  94. modelID: ModelV2.ID.make("test"),
  95. providerID: ProviderV2.ID.make("test"),
  96. mode: "",
  97. agent: "default",
  98. path: { cwd: "/", root: "/" },
  99. cost: 0,
  100. tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
  101. summary: opts?.summary,
  102. finish: opts?.finish,
  103. error: opts?.error,
  104. } as unknown as SessionV1.Info)
  105. return id
  106. })
  107. const addCompactionPart = Effect.fn("Test.addCompactionPart")(function* (
  108. sessionID: SessionID,
  109. messageID: MessageID,
  110. tailStartID?: MessageID,
  111. ) {
  112. const session = yield* SessionNs.Service
  113. yield* session.updatePart({
  114. id: PartID.ascending(),
  115. sessionID,
  116. messageID,
  117. type: "compaction",
  118. auto: true,
  119. tail_start_id: tailStartID,
  120. } as any)
  121. })
  122. describe("MessageV2.page", () => {
  123. it.instance("returns page result", () =>
  124. withSession(({ sessionID }) =>
  125. Effect.gen(function* () {
  126. yield* fill(sessionID, 2)
  127. const result = yield* MessageV2.page({ sessionID, limit: 10 })
  128. expect(result).toBeDefined()
  129. expect(result.items).toBeArray()
  130. }),
  131. ),
  132. )
  133. it.instance("pages backward with opaque cursors", () =>
  134. withSession(({ sessionID }) =>
  135. Effect.gen(function* () {
  136. const ids = yield* fill(sessionID, 6)
  137. const a = yield* MessageV2.page({ sessionID, limit: 2 })
  138. expect(a.items.map((item) => item.info.id)).toEqual(ids.slice(-2))
  139. expect(a.items.every((item) => item.parts.length === 1)).toBe(true)
  140. expect(a.more).toBe(true)
  141. expect(a.cursor).toBeTruthy()
  142. const b = yield* MessageV2.page({ sessionID, limit: 2, before: a.cursor! })
  143. expect(b.items.map((item) => item.info.id)).toEqual(ids.slice(-4, -2))
  144. expect(b.more).toBe(true)
  145. expect(b.cursor).toBeTruthy()
  146. const c = yield* MessageV2.page({ sessionID, limit: 2, before: b.cursor! })
  147. expect(c.items.map((item) => item.info.id)).toEqual(ids.slice(0, 2))
  148. expect(c.more).toBe(false)
  149. expect(c.cursor).toBeUndefined()
  150. }),
  151. ),
  152. )
  153. it.instance("returns items in chronological order within a page", () =>
  154. withSession(({ sessionID }) =>
  155. Effect.gen(function* () {
  156. const ids = yield* fill(sessionID, 4)
  157. const result = yield* MessageV2.page({ sessionID, limit: 4 })
  158. expect(result.items.map((item) => item.info.id)).toEqual(ids)
  159. }),
  160. ),
  161. )
  162. it.instance("returns empty items for session with no messages", () =>
  163. withSession(({ sessionID }) =>
  164. Effect.gen(function* () {
  165. const result = yield* MessageV2.page({ sessionID, limit: 10 })
  166. expect(result.items).toEqual([])
  167. expect(result.more).toBe(false)
  168. expect(result.cursor).toBeUndefined()
  169. }),
  170. ),
  171. )
  172. it.instance("fails with NotFoundError for non-existent session", () =>
  173. Effect.gen(function* () {
  174. const fake = "non-existent-session" as SessionID
  175. const error = yield* Effect.flip(MessageV2.page({ sessionID: fake, limit: 10 }))
  176. expect(error).toBeInstanceOf(NotFoundError)
  177. expect(error.message).toBe(`Session not found: ${fake}`)
  178. }),
  179. )
  180. it.instance("handles exact limit boundary", () =>
  181. withSession(({ sessionID }) =>
  182. Effect.gen(function* () {
  183. const ids = yield* fill(sessionID, 3)
  184. const result = yield* MessageV2.page({ sessionID, limit: 3 })
  185. expect(result.items.map((item) => item.info.id)).toEqual(ids)
  186. expect(result.more).toBe(false)
  187. expect(result.cursor).toBeUndefined()
  188. }),
  189. ),
  190. )
  191. it.instance("limit of 1 returns single newest message", () =>
  192. withSession(({ sessionID }) =>
  193. Effect.gen(function* () {
  194. const ids = yield* fill(sessionID, 5)
  195. const result = yield* MessageV2.page({ sessionID, limit: 1 })
  196. expect(result.items).toHaveLength(1)
  197. expect(result.items[0].info.id).toBe(ids[ids.length - 1])
  198. expect(result.more).toBe(true)
  199. }),
  200. ),
  201. )
  202. it.instance("hydrates multiple parts per message", () =>
  203. withSession(({ session, sessionID }) =>
  204. Effect.gen(function* () {
  205. const [id] = yield* fill(sessionID, 1)
  206. yield* session.updatePart({
  207. id: PartID.ascending(),
  208. sessionID,
  209. messageID: id,
  210. type: "text",
  211. text: "extra",
  212. })
  213. const result = yield* MessageV2.page({ sessionID, limit: 10 })
  214. expect(result.items).toHaveLength(1)
  215. expect(result.items[0].parts).toHaveLength(2)
  216. }),
  217. ),
  218. )
  219. it.instance("accepts cursors from fractional timestamps", () =>
  220. withSession(({ sessionID }) =>
  221. Effect.gen(function* () {
  222. const ids = yield* fill(sessionID, 4, (i: number) => 1000.5 + i)
  223. const a = yield* MessageV2.page({ sessionID, limit: 2 })
  224. const b = yield* MessageV2.page({ sessionID, limit: 2, before: a.cursor! })
  225. expect(a.items.map((item) => item.info.id)).toEqual(ids.slice(-2))
  226. expect(b.items.map((item) => item.info.id)).toEqual(ids.slice(0, 2))
  227. }),
  228. ),
  229. )
  230. it.instance("messages with same timestamp are ordered by id", () =>
  231. withSession(({ sessionID }) =>
  232. Effect.gen(function* () {
  233. const ids = yield* fill(sessionID, 4, () => 1000)
  234. const a = yield* MessageV2.page({ sessionID, limit: 2 })
  235. expect(a.items.map((item) => item.info.id)).toEqual(ids.slice(-2))
  236. expect(a.more).toBe(true)
  237. const b = yield* MessageV2.page({ sessionID, limit: 2, before: a.cursor! })
  238. expect(b.items.map((item) => item.info.id)).toEqual(ids.slice(0, 2))
  239. expect(b.more).toBe(false)
  240. }),
  241. ),
  242. )
  243. it.instance("does not return messages from other sessions", () =>
  244. Effect.gen(function* () {
  245. const session = yield* SessionNs.Service
  246. const a = yield* session.create({})
  247. const b = yield* session.create({})
  248. yield* fill(a.id, 3)
  249. yield* fill(b.id, 2)
  250. const resultA = yield* MessageV2.page({ sessionID: a.id, limit: 10 })
  251. const resultB = yield* MessageV2.page({ sessionID: b.id, limit: 10 })
  252. expect(resultA.items).toHaveLength(3)
  253. expect(resultB.items).toHaveLength(2)
  254. expect(resultA.items.every((item) => item.info.sessionID === a.id)).toBe(true)
  255. expect(resultB.items.every((item) => item.info.sessionID === b.id)).toBe(true)
  256. yield* session.remove(a.id)
  257. yield* session.remove(b.id)
  258. }),
  259. )
  260. it.instance("large limit returns all messages without cursor", () =>
  261. withSession(({ sessionID }) =>
  262. Effect.gen(function* () {
  263. const ids = yield* fill(sessionID, 10)
  264. const result = yield* MessageV2.page({ sessionID, limit: 100 })
  265. expect(result.items).toHaveLength(10)
  266. expect(result.items.map((item) => item.info.id)).toEqual(ids)
  267. expect(result.more).toBe(false)
  268. expect(result.cursor).toBeUndefined()
  269. }),
  270. ),
  271. )
  272. })
  273. describe("MessageV2.stream", () => {
  274. it.instance("yields items newest first", () =>
  275. withSession(({ sessionID }) =>
  276. Effect.gen(function* () {
  277. const ids = yield* fill(sessionID, 5)
  278. const items = yield* MessageV2.stream(sessionID)
  279. expect(items.map((item) => item.info.id)).toEqual(ids.slice().reverse())
  280. }),
  281. ),
  282. )
  283. it.instance("yields nothing for empty session", () =>
  284. withSession(({ sessionID }) =>
  285. Effect.gen(function* () {
  286. const items = yield* MessageV2.stream(sessionID)
  287. expect(items).toHaveLength(0)
  288. }),
  289. ),
  290. )
  291. it.instance("yields single message", () =>
  292. withSession(({ sessionID }) =>
  293. Effect.gen(function* () {
  294. const ids = yield* fill(sessionID, 1)
  295. const items = yield* MessageV2.stream(sessionID)
  296. expect(items).toHaveLength(1)
  297. expect(items[0].info.id).toBe(ids[0])
  298. }),
  299. ),
  300. )
  301. it.instance("hydrates parts for each yielded message", () =>
  302. withSession(({ sessionID }) =>
  303. Effect.gen(function* () {
  304. yield* fill(sessionID, 3)
  305. const items = yield* MessageV2.stream(sessionID)
  306. for (const item of items) {
  307. expect(item.parts).toHaveLength(1)
  308. expect(item.parts[0].type).toBe("text")
  309. }
  310. }),
  311. ),
  312. )
  313. it.instance("handles sets exceeding internal page size", () =>
  314. withSession(({ sessionID }) =>
  315. Effect.gen(function* () {
  316. const ids = yield* fill(sessionID, 60)
  317. const items = yield* MessageV2.stream(sessionID)
  318. expect(items).toHaveLength(60)
  319. expect(items[0].info.id).toBe(ids[ids.length - 1])
  320. expect(items[59].info.id).toBe(ids[0])
  321. }),
  322. ),
  323. )
  324. it.instance("returns an Effect", () =>
  325. withSession(({ sessionID }) =>
  326. Effect.gen(function* () {
  327. yield* fill(sessionID, 1)
  328. const result = yield* MessageV2.stream(sessionID)
  329. expect(result).toHaveLength(1)
  330. }),
  331. ),
  332. )
  333. })
  334. describe("MessageV2.parts", () => {
  335. it.instance("returns parts for a message", () =>
  336. withSession(({ sessionID }) =>
  337. Effect.gen(function* () {
  338. const [id] = yield* fill(sessionID, 1)
  339. const result = yield* MessageV2.parts(id)
  340. expect(result).toHaveLength(1)
  341. expect(result[0].type).toBe("text")
  342. expect((result[0] as SessionV1.TextPart).text).toBe("m0")
  343. }),
  344. ),
  345. )
  346. it.instance("returns empty array for message with no parts", () =>
  347. withSession(({ sessionID }) =>
  348. Effect.gen(function* () {
  349. const id = yield* addUser(sessionID)
  350. const result = yield* MessageV2.parts(id)
  351. expect(result).toEqual([])
  352. }),
  353. ),
  354. )
  355. it.instance("returns multiple parts in order", () =>
  356. withSession(({ session, sessionID }) =>
  357. Effect.gen(function* () {
  358. const [id] = yield* fill(sessionID, 1)
  359. yield* session.updatePart({
  360. id: PartID.ascending(),
  361. sessionID,
  362. messageID: id,
  363. type: "text",
  364. text: "second",
  365. })
  366. yield* session.updatePart({
  367. id: PartID.ascending(),
  368. sessionID,
  369. messageID: id,
  370. type: "text",
  371. text: "third",
  372. })
  373. const result = yield* MessageV2.parts(id)
  374. expect(result).toHaveLength(3)
  375. expect((result[0] as SessionV1.TextPart).text).toBe("m0")
  376. expect((result[1] as SessionV1.TextPart).text).toBe("second")
  377. expect((result[2] as SessionV1.TextPart).text).toBe("third")
  378. }),
  379. ),
  380. )
  381. it.instance("returns empty for non-existent message id", () =>
  382. Effect.gen(function* () {
  383. yield* SessionNs.Service
  384. const result = yield* MessageV2.parts(MessageID.ascending())
  385. expect(result).toEqual([])
  386. }),
  387. )
  388. it.instance("parts contain sessionID and messageID", () =>
  389. withSession(({ sessionID }) =>
  390. Effect.gen(function* () {
  391. const [id] = yield* fill(sessionID, 1)
  392. const result = yield* MessageV2.parts(id)
  393. expect(result[0].sessionID).toBe(sessionID)
  394. expect(result[0].messageID).toBe(id)
  395. }),
  396. ),
  397. )
  398. })
  399. describe("MessageV2.get", () => {
  400. it.instance("returns message with hydrated parts", () =>
  401. withSession(({ sessionID }) =>
  402. Effect.gen(function* () {
  403. const [id] = yield* fill(sessionID, 1)
  404. const result = yield* MessageV2.get({ sessionID, messageID: id })
  405. expect(result.info.id).toBe(id)
  406. expect(result.info.sessionID).toBe(sessionID)
  407. expect(result.info.role).toBe("user")
  408. expect(result.parts).toHaveLength(1)
  409. expect((result.parts[0] as SessionV1.TextPart).text).toBe("m0")
  410. }),
  411. ),
  412. )
  413. it.instance("fails with NotFoundError for non-existent message", () =>
  414. withSession(({ sessionID }) =>
  415. Effect.gen(function* () {
  416. const messageID = MessageID.ascending()
  417. const error = yield* Effect.flip(MessageV2.get({ sessionID, messageID }))
  418. expect(error).toBeInstanceOf(NotFoundError)
  419. expect(error.message).toBe(`Message not found: ${messageID}`)
  420. }),
  421. ),
  422. )
  423. it.instance("scopes by session id", () =>
  424. Effect.gen(function* () {
  425. const session = yield* SessionNs.Service
  426. const a = yield* session.create({})
  427. const b = yield* session.create({})
  428. const [id] = yield* fill(a.id, 1)
  429. const error = yield* Effect.flip(MessageV2.get({ sessionID: b.id, messageID: id }))
  430. expect(error).toBeInstanceOf(NotFoundError)
  431. expect(error.message).toBe(`Message not found: ${id}`)
  432. const result = yield* MessageV2.get({ sessionID: a.id, messageID: id })
  433. expect(result.info.id).toBe(id)
  434. yield* session.remove(a.id)
  435. yield* session.remove(b.id)
  436. }),
  437. )
  438. it.instance("returns message with multiple parts", () =>
  439. withSession(({ session, sessionID }) =>
  440. Effect.gen(function* () {
  441. const [id] = yield* fill(sessionID, 1)
  442. yield* session.updatePart({
  443. id: PartID.ascending(),
  444. sessionID,
  445. messageID: id,
  446. type: "text",
  447. text: "extra",
  448. })
  449. const result = yield* MessageV2.get({ sessionID, messageID: id })
  450. expect(result.parts).toHaveLength(2)
  451. }),
  452. ),
  453. )
  454. it.instance("returns assistant message with correct role", () =>
  455. withSession(({ session, sessionID }) =>
  456. Effect.gen(function* () {
  457. const uid = yield* addUser(sessionID, "hello")
  458. const aid = yield* addAssistant(sessionID, uid)
  459. yield* session.updatePart({
  460. id: PartID.ascending(),
  461. sessionID,
  462. messageID: aid,
  463. type: "text",
  464. text: "response",
  465. })
  466. const result = yield* MessageV2.get({ sessionID, messageID: aid })
  467. expect(result.info.role).toBe("assistant")
  468. expect(result.parts).toHaveLength(1)
  469. expect((result.parts[0] as SessionV1.TextPart).text).toBe("response")
  470. }),
  471. ),
  472. )
  473. it.instance("returns message with zero parts", () =>
  474. withSession(({ sessionID }) =>
  475. Effect.gen(function* () {
  476. const id = yield* addUser(sessionID)
  477. const result = yield* MessageV2.get({ sessionID, messageID: id })
  478. expect(result.info.id).toBe(id)
  479. expect(result.parts).toEqual([])
  480. }),
  481. ),
  482. )
  483. })
  484. describe("Session.messages", () => {
  485. it.instance("returns all messages in chronological order across pages", () =>
  486. withSession(({ session, sessionID }) =>
  487. Effect.gen(function* () {
  488. const ids = yield* fill(sessionID, 55)
  489. const result = yield* session.messages({ sessionID })
  490. expect(result.map((item) => item.info.id)).toEqual(ids)
  491. }),
  492. ),
  493. )
  494. it.instance("fails with NotFoundError for non-existent session", () =>
  495. Effect.gen(function* () {
  496. const session = yield* SessionNs.Service
  497. const fake = "non-existent-session" as SessionID
  498. const error = yield* Effect.flip(session.messages({ sessionID: fake }))
  499. expect(error).toBeInstanceOf(NotFoundError)
  500. expect(error.message).toBe(`Session not found: ${fake}`)
  501. }),
  502. )
  503. })
  504. describe("Session.findMessage", () => {
  505. it.instance("searches newest-first", () =>
  506. withSession(({ session, sessionID }) =>
  507. Effect.gen(function* () {
  508. const ids = yield* fill(sessionID, 3)
  509. const result = yield* session.findMessage(sessionID, () => true)
  510. expect(Option.isSome(result) ? result.value.info.id : undefined).toBe(ids.at(-1))
  511. }),
  512. ),
  513. )
  514. it.instance("fails with NotFoundError for non-existent session", () =>
  515. Effect.gen(function* () {
  516. const session = yield* SessionNs.Service
  517. const fake = "non-existent-session" as SessionID
  518. const error = yield* Effect.flip(session.findMessage(fake, () => true))
  519. expect(error).toBeInstanceOf(NotFoundError)
  520. expect(error.message).toBe(`Session not found: ${fake}`)
  521. }),
  522. )
  523. })
  524. describe("MessageV2.filterCompacted", () => {
  525. it.instance("returns all messages when no compaction", () =>
  526. withSession(({ sessionID }) =>
  527. Effect.gen(function* () {
  528. const ids = yield* fill(sessionID, 5)
  529. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  530. expect(result).toHaveLength(5)
  531. // reversed from newest-first to chronological
  532. expect(result.map((item) => item.info.id)).toEqual(ids)
  533. }),
  534. ),
  535. )
  536. it.instance("stops at compaction boundary and returns chronological order", () =>
  537. withSession(({ session, sessionID }) =>
  538. Effect.gen(function* () {
  539. // Chronological: u1(+compaction part), a1(summary, parentID=u1), u2, a2
  540. // Stream (newest first): a2, u2, a1(adds u1 to completed), u1(in completed + compaction) -> break
  541. const u1 = yield* addUser(sessionID, "first question")
  542. const a1 = yield* addAssistant(sessionID, u1, { summary: true, finish: "end_turn" })
  543. yield* session.updatePart({
  544. id: PartID.ascending(),
  545. sessionID,
  546. messageID: a1,
  547. type: "text",
  548. text: "summary",
  549. })
  550. yield* addCompactionPart(sessionID, u1)
  551. const u2 = yield* addUser(sessionID, "new question")
  552. const a2 = yield* addAssistant(sessionID, u2)
  553. yield* session.updatePart({
  554. id: PartID.ascending(),
  555. sessionID,
  556. messageID: a2,
  557. type: "text",
  558. text: "new response",
  559. })
  560. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  561. // Includes compaction boundary: u1, a1, u2, a2
  562. expect(result[0].info.id).toBe(u1)
  563. expect(result.length).toBe(4)
  564. }),
  565. ),
  566. )
  567. it.live("handles empty iterable", () =>
  568. Effect.sync(() => {
  569. const result = MessageV2.filterCompacted([])
  570. expect(result).toEqual([])
  571. }),
  572. )
  573. it.instance("does not break on compaction part without matching summary", () =>
  574. withSession(({ sessionID }) =>
  575. Effect.gen(function* () {
  576. const u1 = yield* addUser(sessionID, "hello")
  577. yield* addCompactionPart(sessionID, u1)
  578. yield* addUser(sessionID, "world")
  579. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  580. expect(result).toHaveLength(2)
  581. }),
  582. ),
  583. )
  584. it.instance("skips assistant with error even if marked as summary", () =>
  585. withSession(({ sessionID }) =>
  586. Effect.gen(function* () {
  587. const u1 = yield* addUser(sessionID, "hello")
  588. yield* addCompactionPart(sessionID, u1)
  589. const error = new SessionV1.APIError({
  590. message: "boom",
  591. isRetryable: true,
  592. }).toObject() as SessionV1.Assistant["error"]
  593. yield* addAssistant(sessionID, u1, { summary: true, finish: "end_turn", error })
  594. yield* addUser(sessionID, "retry")
  595. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  596. // Error assistant doesn't add to completed, so compaction boundary never triggers
  597. expect(result).toHaveLength(3)
  598. }),
  599. ),
  600. )
  601. it.instance("skips assistant without finish even if marked as summary", () =>
  602. withSession(({ sessionID }) =>
  603. Effect.gen(function* () {
  604. const u1 = yield* addUser(sessionID, "hello")
  605. yield* addCompactionPart(sessionID, u1)
  606. // summary=true but no finish
  607. yield* addAssistant(sessionID, u1, { summary: true })
  608. yield* addUser(sessionID, "next")
  609. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  610. expect(result).toHaveLength(3)
  611. }),
  612. ),
  613. )
  614. it.instance("retains original tail when compaction stores tail_start_id", () =>
  615. withSession(({ session, sessionID }) =>
  616. Effect.gen(function* () {
  617. const u1 = yield* addUser(sessionID, "first")
  618. const a1 = yield* addAssistant(sessionID, u1, { finish: "end_turn" })
  619. yield* session.updatePart({
  620. id: PartID.ascending(),
  621. sessionID,
  622. messageID: a1,
  623. type: "text",
  624. text: "first reply",
  625. })
  626. const u2 = yield* addUser(sessionID, "second")
  627. const a2 = yield* addAssistant(sessionID, u2, { finish: "end_turn" })
  628. yield* session.updatePart({
  629. id: PartID.ascending(),
  630. sessionID,
  631. messageID: a2,
  632. type: "text",
  633. text: "second reply",
  634. })
  635. const c1 = yield* addUser(sessionID)
  636. yield* addCompactionPart(sessionID, c1, u2)
  637. const s1 = yield* addAssistant(sessionID, c1, { summary: true, finish: "end_turn" })
  638. yield* session.updatePart({
  639. id: PartID.ascending(),
  640. sessionID,
  641. messageID: s1,
  642. type: "text",
  643. text: "summary",
  644. })
  645. const u3 = yield* addUser(sessionID, "third")
  646. const a3 = yield* addAssistant(sessionID, u3, { finish: "end_turn" })
  647. yield* session.updatePart({
  648. id: PartID.ascending(),
  649. sessionID,
  650. messageID: a3,
  651. type: "text",
  652. text: "third reply",
  653. })
  654. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  655. expect(result.map((item) => item.info.id)).toEqual([c1, s1, u2, a2, u3, a3])
  656. }),
  657. ),
  658. )
  659. it.instance("fork remaps compaction tail_start_id for filterCompacted", () =>
  660. Effect.gen(function* () {
  661. const session = yield* SessionNs.Service
  662. const created = yield* session.create({})
  663. const u1 = yield* addUser(created.id, "first")
  664. const a1 = yield* addAssistant(created.id, u1, { finish: "end_turn" })
  665. yield* session.updatePart({
  666. id: PartID.ascending(),
  667. sessionID: created.id,
  668. messageID: a1,
  669. type: "text",
  670. text: "first reply",
  671. })
  672. const u2 = yield* addUser(created.id, "second")
  673. const a2 = yield* addAssistant(created.id, u2, { finish: "end_turn" })
  674. yield* session.updatePart({
  675. id: PartID.ascending(),
  676. sessionID: created.id,
  677. messageID: a2,
  678. type: "text",
  679. text: "second reply",
  680. })
  681. const c1 = yield* addUser(created.id)
  682. yield* addCompactionPart(created.id, c1, u2)
  683. const s1 = yield* addAssistant(created.id, c1, { summary: true, finish: "end_turn" })
  684. yield* session.updatePart({
  685. id: PartID.ascending(),
  686. sessionID: created.id,
  687. messageID: s1,
  688. type: "text",
  689. text: "summary",
  690. })
  691. const u3 = yield* addUser(created.id, "third")
  692. const a3 = yield* addAssistant(created.id, u3, { finish: "end_turn" })
  693. yield* session.updatePart({
  694. id: PartID.ascending(),
  695. sessionID: created.id,
  696. messageID: a3,
  697. type: "text",
  698. text: "third reply",
  699. })
  700. const parentFiltered = MessageV2.filterCompacted(yield* MessageV2.stream(created.id))
  701. expect(parentFiltered.map((item) => item.info.id)).toEqual([c1, s1, u2, a2, u3, a3])
  702. const forked = yield* session.fork({ sessionID: created.id })
  703. const childFiltered = MessageV2.filterCompacted(yield* MessageV2.stream(forked.id))
  704. expect(childFiltered).toHaveLength(parentFiltered.length)
  705. const tailPart = childFiltered.flatMap((m) => m.parts).find((p) => p.type === "compaction")
  706. expect(tailPart?.type).toBe("compaction")
  707. if (!tailPart || tailPart.type !== "compaction") throw new Error("Expected forked compaction part")
  708. expect(tailPart.tail_start_id).toBeDefined()
  709. expect(childFiltered.some((m) => m.info.id === tailPart.tail_start_id)).toBe(true)
  710. yield* session.remove(forked.id)
  711. yield* session.remove(created.id)
  712. }),
  713. )
  714. it.instance("retains an assistant tail when compaction starts inside a turn", () =>
  715. withSession(({ session, sessionID }) =>
  716. Effect.gen(function* () {
  717. const u1 = yield* addUser(sessionID, "first")
  718. const a1 = yield* addAssistant(sessionID, u1, { finish: "end_turn" })
  719. yield* session.updatePart({
  720. id: PartID.ascending(),
  721. sessionID,
  722. messageID: a1,
  723. type: "text",
  724. text: "first reply",
  725. })
  726. const u2 = yield* addUser(sessionID, "second")
  727. const a2 = yield* addAssistant(sessionID, u2, { finish: "end_turn" })
  728. yield* session.updatePart({
  729. id: PartID.ascending(),
  730. sessionID,
  731. messageID: a2,
  732. type: "text",
  733. text: "second reply",
  734. })
  735. const a3 = yield* addAssistant(sessionID, u2, { finish: "end_turn" })
  736. yield* session.updatePart({
  737. id: PartID.ascending(),
  738. sessionID,
  739. messageID: a3,
  740. type: "text",
  741. text: "tail reply",
  742. })
  743. const c1 = yield* addUser(sessionID)
  744. yield* addCompactionPart(sessionID, c1, a3)
  745. const s1 = yield* addAssistant(sessionID, c1, { summary: true, finish: "end_turn" })
  746. yield* session.updatePart({
  747. id: PartID.ascending(),
  748. sessionID,
  749. messageID: s1,
  750. type: "text",
  751. text: "summary",
  752. })
  753. const u3 = yield* addUser(sessionID, "third")
  754. const a4 = yield* addAssistant(sessionID, u3, { finish: "end_turn" })
  755. yield* session.updatePart({
  756. id: PartID.ascending(),
  757. sessionID,
  758. messageID: a4,
  759. type: "text",
  760. text: "third reply",
  761. })
  762. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  763. expect(result.map((item) => item.info.id)).toEqual([c1, s1, a3, u3, a4])
  764. }),
  765. ),
  766. )
  767. it.instance("prefers latest compaction boundary when repeated compactions exist", () =>
  768. withSession(({ session, sessionID }) =>
  769. Effect.gen(function* () {
  770. const u1 = yield* addUser(sessionID, "first")
  771. const a1 = yield* addAssistant(sessionID, u1, { finish: "end_turn" })
  772. yield* session.updatePart({
  773. id: PartID.ascending(),
  774. sessionID,
  775. messageID: a1,
  776. type: "text",
  777. text: "first reply",
  778. })
  779. const u2 = yield* addUser(sessionID, "second")
  780. const a2 = yield* addAssistant(sessionID, u2, { finish: "end_turn" })
  781. yield* session.updatePart({
  782. id: PartID.ascending(),
  783. sessionID,
  784. messageID: a2,
  785. type: "text",
  786. text: "second reply",
  787. })
  788. const c1 = yield* addUser(sessionID)
  789. yield* addCompactionPart(sessionID, c1, u2)
  790. const s1 = yield* addAssistant(sessionID, c1, { summary: true, finish: "end_turn" })
  791. yield* session.updatePart({
  792. id: PartID.ascending(),
  793. sessionID,
  794. messageID: s1,
  795. type: "text",
  796. text: "summary one",
  797. })
  798. const u3 = yield* addUser(sessionID, "third")
  799. const a3 = yield* addAssistant(sessionID, u3, { finish: "end_turn" })
  800. yield* session.updatePart({
  801. id: PartID.ascending(),
  802. sessionID,
  803. messageID: a3,
  804. type: "text",
  805. text: "third reply",
  806. })
  807. const c2 = yield* addUser(sessionID)
  808. yield* addCompactionPart(sessionID, c2, u3)
  809. const s2 = yield* addAssistant(sessionID, c2, { summary: true, finish: "end_turn" })
  810. yield* session.updatePart({
  811. id: PartID.ascending(),
  812. sessionID,
  813. messageID: s2,
  814. type: "text",
  815. text: "summary two",
  816. })
  817. const u4 = yield* addUser(sessionID, "fourth")
  818. const a4 = yield* addAssistant(sessionID, u4, { finish: "end_turn" })
  819. yield* session.updatePart({
  820. id: PartID.ascending(),
  821. sessionID,
  822. messageID: a4,
  823. type: "text",
  824. text: "fourth reply",
  825. })
  826. const result = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
  827. expect(result.map((item) => item.info.id)).toEqual([c2, s2, u3, a3, u4, a4])
  828. }),
  829. ),
  830. )
  831. test("works with array input", () => {
  832. // filterCompacted accepts any Iterable, not just generators
  833. const id = MessageID.ascending()
  834. const items: SessionV1.WithParts[] = [
  835. {
  836. info: {
  837. id,
  838. sessionID: "s1",
  839. role: "user",
  840. time: { created: 1 },
  841. agent: "test",
  842. model: { providerID: "test", modelID: "test" },
  843. } as unknown as SessionV1.Info,
  844. parts: [{ type: "text", text: "hello" }] as unknown as SessionV1.Part[],
  845. },
  846. ]
  847. const result = MessageV2.filterCompacted(items)
  848. expect(result).toHaveLength(1)
  849. expect(result[0].info.id).toBe(id)
  850. })
  851. })
  852. describe("MessageV2.cursor", () => {
  853. test("encode/decode roundtrip", () => {
  854. const input = { id: MessageID.ascending(), time: 1234567890 }
  855. const encoded = MessageV2.cursor.encode(input)
  856. const decoded = MessageV2.cursor.decode(encoded)
  857. expect(decoded.id).toBe(input.id)
  858. expect(decoded.time).toBe(input.time)
  859. })
  860. test("encode/decode with fractional time", () => {
  861. const input = { id: MessageID.ascending(), time: 1234567890.5 }
  862. const encoded = MessageV2.cursor.encode(input)
  863. const decoded = MessageV2.cursor.decode(encoded)
  864. expect(decoded.time).toBe(1234567890.5)
  865. })
  866. test("encoded cursor is base64url", () => {
  867. const encoded = MessageV2.cursor.encode({ id: MessageID.ascending(), time: 0 })
  868. expect(encoded).toMatch(/^[A-Za-z0-9_-]+$/)
  869. })
  870. })
  871. describe("MessageV2 consistency", () => {
  872. it.instance("page hydration matches get for each message", () =>
  873. withSession(({ sessionID }) =>
  874. Effect.gen(function* () {
  875. yield* fill(sessionID, 3)
  876. const paged = yield* MessageV2.page({ sessionID, limit: 10 })
  877. for (const item of paged.items) {
  878. const got = yield* MessageV2.get({ sessionID, messageID: item.info.id as MessageID })
  879. expect(got.info).toEqual(item.info)
  880. expect(got.parts).toEqual(item.parts)
  881. }
  882. }),
  883. ),
  884. )
  885. it.instance("parts from get match standalone parts call", () =>
  886. withSession(({ sessionID }) =>
  887. Effect.gen(function* () {
  888. const [id] = yield* fill(sessionID, 1)
  889. const got = yield* MessageV2.get({ sessionID, messageID: id })
  890. const standalone = yield* MessageV2.parts(id)
  891. expect(got.parts).toEqual(standalone)
  892. }),
  893. ),
  894. )
  895. it.instance("stream collects same messages as exhaustive page iteration", () =>
  896. withSession(({ sessionID }) =>
  897. Effect.gen(function* () {
  898. yield* fill(sessionID, 7)
  899. const streamed = yield* MessageV2.stream(sessionID)
  900. const paged = [] as SessionV1.WithParts[]
  901. let cursor: string | undefined
  902. while (true) {
  903. const result = yield* MessageV2.page({ sessionID, limit: 3, before: cursor })
  904. for (let i = result.items.length - 1; i >= 0; i--) {
  905. paged.push(result.items[i])
  906. }
  907. if (!result.more || !result.cursor) break
  908. cursor = result.cursor
  909. }
  910. expect(streamed.map((m) => m.info.id)).toEqual(paged.map((m) => m.info.id))
  911. }),
  912. ),
  913. )
  914. it.instance("filterCompacted of full stream returns same as Array.from when no compaction", () =>
  915. withSession(({ sessionID }) =>
  916. Effect.gen(function* () {
  917. yield* fill(sessionID, 4)
  918. const stream = yield* MessageV2.stream(sessionID)
  919. const filtered = MessageV2.filterCompacted(stream)
  920. const all = stream.toReversed()
  921. expect(filtered.map((m) => m.info.id)).toEqual(all.map((m) => m.info.id))
  922. }),
  923. ),
  924. )
  925. })