server-session.test.ts 45 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127
  1. import { describe, expect, test } from "bun:test"
  2. import type { retry } from "@opencode-ai/core/util/retry"
  3. import type { Message, OpencodeClient, Part, Session } from "@opencode-ai/sdk/v2/client"
  4. import { createServerSession } from "./server-session"
  5. const session = (id: string, parentID?: string): Session => ({
  6. id,
  7. slug: id,
  8. projectID: "project",
  9. directory: "/repo",
  10. title: id,
  11. version: "1",
  12. parentID,
  13. time: { created: 1, updated: 1 },
  14. })
  15. type UserMessage = Extract<Message, { role: "user" }>
  16. type TextPart = Extract<Part, { type: "text" }>
  17. type MessageResponse = {
  18. data: { info: Message; parts: Part[] }[]
  19. response: { headers: Headers }
  20. }
  21. const userMessage = (id: string, input: Partial<UserMessage> = {}): UserMessage => ({
  22. id,
  23. sessionID: "child",
  24. role: "user",
  25. time: { created: 1 },
  26. agent: "build",
  27. model: { providerID: "provider", modelID: "model" },
  28. ...input,
  29. })
  30. const textPart = (messageID: string, input: Partial<TextPart> = {}): TextPart => ({
  31. id: "part",
  32. sessionID: "child",
  33. messageID,
  34. type: "text",
  35. text: "text",
  36. ...input,
  37. })
  38. const response = (data: MessageResponse["data"] = [], cursor?: string): MessageResponse => ({
  39. data,
  40. response: { headers: new Headers(cursor ? { "x-next-cursor": cursor } : undefined) },
  41. })
  42. const deferredResponse = () => Promise.withResolvers<MessageResponse>()
  43. function messageClient(...responses: Array<MessageResponse | Promise<MessageResponse>>) {
  44. let index = 0
  45. const requests: unknown[] = []
  46. const waiting = new Map<number, () => void>()
  47. const client = {
  48. session: {
  49. get: async () => ({ data: session("child", "root") }),
  50. messages: (input: unknown) => {
  51. requests.push(input)
  52. waiting.get(requests.length)?.()
  53. waiting.delete(requests.length)
  54. return responses[index++]
  55. },
  56. },
  57. } as unknown as OpencodeClient
  58. return Object.assign(client, {
  59. requests,
  60. requested(count: number) {
  61. if (requests.length >= count) return Promise.resolve()
  62. return new Promise<void>((resolve) => waiting.set(count, resolve))
  63. },
  64. })
  65. }
  66. const retryImmediately: typeof retry = async (task, options = {}) => {
  67. const attempts = options.attempts ?? 3
  68. for (let attempt = 0; ; attempt++) {
  69. try {
  70. return await task()
  71. } catch (error) {
  72. if (attempt === attempts - 1) throw error
  73. }
  74. }
  75. }
  76. function setup(sessions: Record<string, Session>) {
  77. const get: unknown[] = []
  78. const messages: unknown[] = []
  79. const client = {
  80. session: {
  81. get: async (input: unknown) => {
  82. get.push(input)
  83. const id = (input as { sessionID: string }).sessionID
  84. return { data: sessions[id] }
  85. },
  86. messages: async (input: unknown) => {
  87. messages.push(input)
  88. return response()
  89. },
  90. diff: async () => ({ data: [] }),
  91. todo: async () => ({ data: [] }),
  92. },
  93. } as unknown as OpencodeClient
  94. return { get, messages, store: createServerSession(client) }
  95. }
  96. describe("server session", () => {
  97. test("resolves lineage by session ID without directory", async () => {
  98. const ctx = setup({ child: session("child", "root"), root: session("root") })
  99. const result = await ctx.store.lineage.resolve("child")
  100. expect(result.root.id).toBe("root")
  101. expect(ctx.get).toEqual([{ sessionID: "child" }, { sessionID: "root" }])
  102. expect(ctx.store.lineage.peek("child")).toEqual(result)
  103. })
  104. test("loads session content through the server client", async () => {
  105. const ctx = setup({ root: session("root") })
  106. await ctx.store.sync("root")
  107. expect(ctx.get).toEqual([{ sessionID: "root" }])
  108. expect(ctx.messages).toEqual([{ sessionID: "root", limit: 2, before: undefined }])
  109. expect(ctx.store.data.message.root).toEqual([])
  110. })
  111. test("merges live events into the initial page", async () => {
  112. const pending = deferredResponse()
  113. const user = userMessage("message-1")
  114. const live = userMessage("message-2", { time: { created: 2 } })
  115. const livePart = textPart(live.id, { text: "live" })
  116. const store = createServerSession(messageClient(pending.promise))
  117. const loading = store.sync("child")
  118. store.apply({ type: "message.updated", properties: { info: live } })
  119. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: livePart, time: 2 } })
  120. pending.resolve(response([{ info: user, parts: [] }]))
  121. await loading
  122. expect(store.data.message.child).toEqual([user, live])
  123. expect(store.data.part[live.id]).toEqual([livePart])
  124. })
  125. test("preserves same-ID live updates over the initial page", async () => {
  126. const pending = deferredResponse()
  127. const fetched = userMessage("message")
  128. const fetchedPart = textPart(fetched.id, { text: "fetched" })
  129. const live = { ...fetched, time: { created: 2 } }
  130. const livePart = { ...fetchedPart, text: "live" }
  131. const store = createServerSession(messageClient(pending.promise))
  132. const loading = store.sync("child")
  133. store.apply({ type: "message.updated", properties: { info: live } })
  134. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: livePart, time: 2 } })
  135. pending.resolve(response([{ info: fetched, parts: [fetchedPart] }]))
  136. await loading
  137. expect(store.data.message.child).toEqual([live])
  138. expect(store.data.part[live.id]).toEqual([livePart])
  139. })
  140. test("preserves removals received during the initial load", async () => {
  141. const pending = deferredResponse()
  142. const removed = userMessage("message-1")
  143. const kept = { ...removed, id: "message-2" }
  144. const part = textPart(kept.id, { text: "removed" })
  145. const store = createServerSession(messageClient(pending.promise))
  146. const loading = store.sync("child")
  147. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: removed.id } })
  148. store.apply({
  149. type: "message.part.removed",
  150. properties: { sessionID: "child", messageID: kept.id, partID: part.id },
  151. })
  152. pending.resolve(
  153. response([
  154. { info: removed, parts: [] },
  155. { info: kept, parts: [part] },
  156. ]),
  157. )
  158. await loading
  159. expect(store.data.message.child).toEqual([kept])
  160. expect(store.data.part[kept.id]).toBeUndefined()
  161. })
  162. test("keeps removal tracking isolated across load generations", async () => {
  163. const firstResponse = deferredResponse()
  164. const secondResponse = deferredResponse()
  165. const message = userMessage("message")
  166. const store = createServerSession(messageClient(firstResponse.promise, secondResponse.promise))
  167. const first = store.sync("child")
  168. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  169. store.apply({
  170. type: "session.deleted",
  171. properties: { sessionID: "child", info: session("child", "root") },
  172. })
  173. const second = store.sync("child")
  174. firstResponse.resolve(response())
  175. await first
  176. secondResponse.resolve(response([{ info: message, parts: [] }]))
  177. await second
  178. expect(store.data.message.child).toEqual([message])
  179. })
  180. test("tracks removals in a replacement load generation", async () => {
  181. const firstResponse = deferredResponse()
  182. const secondResponse = deferredResponse()
  183. const message = userMessage("message")
  184. const store = createServerSession(messageClient(firstResponse.promise, secondResponse.promise))
  185. const first = store.sync("child")
  186. store.apply({
  187. type: "session.deleted",
  188. properties: { sessionID: "child", info: session("child", "root") },
  189. })
  190. const second = store.sync("child")
  191. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  192. firstResponse.resolve(response())
  193. await first
  194. secondResponse.resolve(response([{ info: message, parts: [] }]))
  195. await second
  196. expect(store.data.message.child).toEqual([])
  197. })
  198. test("preserves remove then re-add when a refresh omits the message", async () => {
  199. const pending = deferredResponse()
  200. const message = userMessage("message")
  201. const store = createServerSession(messageClient(response([{ info: message, parts: [] }]), pending.promise))
  202. await store.sync("child")
  203. const refreshing = store.sync("child", { force: true })
  204. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  205. store.apply({ type: "message.updated", properties: { info: message } })
  206. pending.resolve(response())
  207. await refreshing
  208. expect(store.data.message.child).toEqual([message])
  209. })
  210. test("preserves a re-added message without restoring removed parts", async () => {
  211. const pending = deferredResponse()
  212. const message = userMessage("message")
  213. const part = textPart(message.id, { text: "stale" })
  214. const store = createServerSession(messageClient(response([{ info: message, parts: [] }]), pending.promise))
  215. await store.sync("child")
  216. const refreshing = store.sync("child", { force: true })
  217. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  218. store.apply({ type: "message.updated", properties: { info: message } })
  219. pending.resolve(response([{ info: message, parts: [part] }]))
  220. await refreshing
  221. expect(store.data.message.child).toEqual([message])
  222. expect(store.data.part[message.id]).toBeUndefined()
  223. })
  224. test("preserves optimistic parts re-added after removal during a refresh", async () => {
  225. const pending = deferredResponse()
  226. const message = userMessage("message")
  227. const stale = textPart(message.id, { id: "stale", text: "stale" })
  228. const part = textPart(message.id, { id: "optimistic", text: "optimistic" })
  229. const store = createServerSession(
  230. messageClient(response([{ info: message, parts: [] }]), pending.promise, response()),
  231. )
  232. await store.sync("child")
  233. const refreshing = store.sync("child", { force: true })
  234. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  235. store.optimistic.add({ sessionID: "child", message, parts: [part] })
  236. pending.resolve(response([{ info: message, parts: [stale] }]))
  237. await refreshing
  238. expect(store.data.message.child).toEqual([message])
  239. expect(store.data.part[message.id]).toEqual([part])
  240. await store.sync("child", { force: true })
  241. expect(store.data.message.child).toEqual([message])
  242. expect(store.data.part[message.id]).toEqual([part])
  243. })
  244. test("drops stale event content omitted by a complete initial page", async () => {
  245. const stale = userMessage("stale")
  246. const store = createServerSession(messageClient(response()))
  247. store.apply({ type: "message.updated", properties: { info: stale } })
  248. await store.sync("child")
  249. expect(store.data.message.child).toEqual([])
  250. })
  251. test("preserves event content outside an incomplete initial page", async () => {
  252. const live = userMessage("message-1")
  253. const fetched = userMessage("message-2", { time: { created: 2 } })
  254. const store = createServerSession(messageClient(response([{ info: fetched, parts: [] }], "older")))
  255. store.apply({ type: "message.updated", properties: { info: live } })
  256. await store.sync("child")
  257. expect(store.data.message.child).toEqual([live, fetched])
  258. })
  259. test("does not restore removed optimistic content on refresh", async () => {
  260. const message = userMessage("message")
  261. const part = textPart(message.id, { text: "removed" })
  262. const kept = { ...message, id: "kept" }
  263. const keptPart = { ...part, id: "kept-part", messageID: kept.id }
  264. const store = createServerSession(messageClient(response([{ info: kept, parts: [] }])))
  265. store.optimistic.add({ sessionID: "child", message, parts: [part] })
  266. store.optimistic.add({ sessionID: "child", message: kept, parts: [keptPart] })
  267. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  268. store.apply({
  269. type: "message.part.removed",
  270. properties: { sessionID: "child", messageID: kept.id, partID: keptPart.id },
  271. })
  272. await store.sync("child", { force: true })
  273. expect(store.data.message.child).toEqual([kept])
  274. expect(store.data.part[message.id]).toBeUndefined()
  275. expect(store.data.part[kept.id]).toBeUndefined()
  276. })
  277. test("replaces confirmed optimistic content with the initial page", async () => {
  278. const optimistic = userMessage("message")
  279. const fetched = { ...optimistic, time: { created: 2 } }
  280. const store = createServerSession(messageClient(response([{ info: fetched, parts: [] }])))
  281. store.optimistic.add({ sessionID: "child", message: optimistic, parts: [] })
  282. await store.sync("child")
  283. expect(store.data.message.child).toEqual([fetched])
  284. })
  285. test("replaces a confirmed optimistic part with fetched content", async () => {
  286. const pending = deferredResponse()
  287. const message = userMessage("message")
  288. const optimistic = textPart(message.id, { text: "optimistic" })
  289. const fetched = { ...optimistic, text: "fetched" }
  290. const store = createServerSession(messageClient(pending.promise))
  291. const loading = store.sync("child")
  292. store.optimistic.add({ sessionID: "child", message, parts: [optimistic] })
  293. pending.resolve(response([{ info: message, parts: [fetched] }]))
  294. await loading
  295. expect(store.data.part[message.id]).toEqual([fetched])
  296. })
  297. test("rolls back only unconfirmed optimistic parts", async () => {
  298. const pending = deferredResponse()
  299. const message = userMessage("message")
  300. const confirmed = textPart(message.id, { id: "confirmed", text: "confirmed" })
  301. const pendingPart = textPart(message.id, { id: "pending", text: "pending" })
  302. const store = createServerSession(messageClient(pending.promise))
  303. const loading = store.sync("child")
  304. store.optimistic.add({ sessionID: "child", message, parts: [confirmed, pendingPart] })
  305. pending.resolve(response([{ info: message, parts: [confirmed] }]))
  306. await loading
  307. store.optimistic.remove({ sessionID: "child", messageID: message.id })
  308. expect(store.data.message.child).toEqual([message])
  309. expect(store.data.part[message.id]).toEqual([confirmed])
  310. })
  311. test("updates confirmed optimistic parts from later pages", async () => {
  312. const message = userMessage("message")
  313. const confirmed = textPart(message.id, { id: "confirmed", text: "first" })
  314. const updated = { ...confirmed, text: "updated" }
  315. const pendingPart = textPart(message.id, { id: "pending", text: "pending" })
  316. const store = createServerSession(
  317. messageClient(response([{ info: message, parts: [confirmed] }]), response([{ info: message, parts: [updated] }])),
  318. )
  319. store.optimistic.add({ sessionID: "child", message, parts: [confirmed, pendingPart] })
  320. await store.sync("child")
  321. await store.sync("child", { force: true })
  322. store.optimistic.remove({ sessionID: "child", messageID: message.id })
  323. expect(store.data.part[message.id]).toEqual([updated])
  324. })
  325. test("does not restore a confirmed optimistic part after its removal event", async () => {
  326. const message = userMessage("message")
  327. const confirmed = textPart(message.id, { id: "confirmed", text: "confirmed" })
  328. const pendingPart = textPart(message.id, { id: "pending", text: "pending" })
  329. const store = createServerSession(
  330. messageClient(response([{ info: message, parts: [confirmed] }]), response([{ info: message, parts: [] }])),
  331. )
  332. store.optimistic.add({ sessionID: "child", message, parts: [confirmed, pendingPart] })
  333. await store.sync("child")
  334. store.apply({
  335. type: "message.part.removed",
  336. properties: { sessionID: "child", messageID: message.id, partID: confirmed.id },
  337. })
  338. await store.sync("child", { force: true })
  339. expect(store.data.part[message.id]).toEqual([pendingPart])
  340. })
  341. test("clears delta buffers when removing optimistic content", () => {
  342. const message = userMessage("message")
  343. const part = textPart(message.id, { text: "optimistic" })
  344. const store = setup({ child: session("child") }).store
  345. store.optimistic.add({ sessionID: "child", message, parts: [part] })
  346. store.apply({
  347. type: "message.part.delta",
  348. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: " delta" },
  349. })
  350. store.optimistic.remove({ sessionID: "child", messageID: message.id })
  351. expect(store.data.part[message.id]).toBeUndefined()
  352. expect(store.data.part_text_accum_delta[part.id]).toBeUndefined()
  353. })
  354. test("does not remove content confirmed by a message event", () => {
  355. const message = userMessage("message")
  356. const part = textPart(message.id)
  357. const store = setup({ child: session("child") }).store
  358. store.optimistic.add({ sessionID: "child", message, parts: [part] })
  359. store.apply({ type: "message.updated", properties: { sessionID: "child", info: message } })
  360. store.optimistic.remove({ sessionID: "child", messageID: message.id })
  361. expect(store.data.message.child).toEqual([message])
  362. expect(store.data.part[message.id]).toBeUndefined()
  363. })
  364. test("does not remove parts confirmed by part events", () => {
  365. const message = userMessage("message")
  366. const part = textPart(message.id)
  367. const store = setup({ child: session("child") }).store
  368. store.optimistic.add({ sessionID: "child", message, parts: [part] })
  369. store.apply({ type: "message.updated", properties: { sessionID: "child", info: message } })
  370. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } })
  371. store.optimistic.remove({ sessionID: "child", messageID: message.id })
  372. expect(store.data.message.child).toEqual([message])
  373. expect(store.data.part[message.id]).toEqual([part])
  374. })
  375. test("treats a part event as confirmation when it precedes the message event", () => {
  376. const message = userMessage("message")
  377. const part = textPart(message.id)
  378. const store = setup({ child: session("child") }).store
  379. store.optimistic.add({ sessionID: "child", message, parts: [part] })
  380. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } })
  381. store.optimistic.remove({ sessionID: "child", messageID: message.id })
  382. expect(store.data.message.child).toEqual([message])
  383. expect(store.data.part[message.id]).toEqual([part])
  384. })
  385. test("clears stale parts when the initial page has none", async () => {
  386. const pending = deferredResponse()
  387. const message = userMessage("message")
  388. const part = textPart(message.id, { text: "stale" })
  389. const store = createServerSession(messageClient(pending.promise))
  390. store.apply({ type: "message.updated", properties: { info: message } })
  391. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 1 } })
  392. const loading = store.sync("child")
  393. pending.resolve(response([{ info: message, parts: [] }]))
  394. await loading
  395. expect(store.data.part[message.id]).toBeUndefined()
  396. })
  397. test("clears delta buffers for parts omitted by the initial page", async () => {
  398. const pending = deferredResponse()
  399. const message = userMessage("message")
  400. const kept = textPart(message.id, { id: "part-1", text: "kept" })
  401. const removed: Part = { ...kept, id: "part-2", text: "removed" }
  402. const store = createServerSession(messageClient(pending.promise))
  403. store.apply({ type: "message.updated", properties: { info: message } })
  404. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: kept, time: 1 } })
  405. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: removed, time: 1 } })
  406. store.apply({
  407. type: "message.part.delta",
  408. properties: { sessionID: "child", messageID: message.id, partID: removed.id, field: "text", delta: " delta" },
  409. })
  410. const loading = store.sync("child")
  411. pending.resolve(response([{ info: message, parts: [kept] }]))
  412. await loading
  413. expect(store.data.part[message.id]).toEqual([kept])
  414. expect(store.data.part_text_accum_delta[removed.id]).toBeUndefined()
  415. })
  416. test("clears a stale delta buffer when a refresh replaces its part", async () => {
  417. const message = userMessage("message")
  418. const stale = textPart(message.id, { text: "stale" })
  419. const fetched = { ...stale, text: "fetched" }
  420. const store = createServerSession(
  421. messageClient(response([{ info: message, parts: [stale] }]), response([{ info: message, parts: [fetched] }])),
  422. )
  423. await store.sync("child")
  424. store.apply({
  425. type: "message.part.delta",
  426. properties: { sessionID: "child", messageID: message.id, partID: stale.id, field: "text", delta: " delta" },
  427. })
  428. await store.sync("child", { force: true })
  429. expect(store.data.part[message.id]).toEqual([fetched])
  430. expect(store.data.part_text_accum_delta[stale.id]).toBeUndefined()
  431. })
  432. test("preserves a non-durable delta received before refresh", async () => {
  433. const message = userMessage("message")
  434. const part = textPart(message.id, { text: "stale" })
  435. const store = createServerSession(
  436. messageClient(response([{ info: message, parts: [part] }]), response([{ info: message, parts: [{ ...part }] }])),
  437. )
  438. await store.sync("child")
  439. store.apply({
  440. type: "message.part.delta",
  441. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: " delta" },
  442. })
  443. await store.sync("child", { force: true })
  444. expect(store.data.part[message.id]).toEqual([{ ...part, text: "stale delta" }])
  445. expect(store.data.part_text_accum_delta[part.id]).toBe("stale delta")
  446. })
  447. test("accepts fetched text that intentionally replaces an accumulated prefix", async () => {
  448. const message = userMessage("message")
  449. const part = textPart(message.id, { text: "abc" })
  450. const fetched = { ...part, text: "ab" }
  451. const store = createServerSession(
  452. messageClient(response([{ info: message, parts: [part] }]), response([{ info: message, parts: [fetched] }])),
  453. )
  454. await store.sync("child")
  455. store.apply({
  456. type: "message.part.delta",
  457. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: "def" },
  458. })
  459. await store.sync("child", { force: true })
  460. expect(store.data.part[message.id]).toEqual([fetched])
  461. expect(store.data.part_text_accum_delta[part.id]).toBeUndefined()
  462. })
  463. test("preserves an unpersisted delta suffix after partial server catch-up", async () => {
  464. const message = userMessage("message")
  465. const part = textPart(message.id, { text: "a" })
  466. const fetched = { ...part, text: "ab" }
  467. const store = createServerSession(
  468. messageClient(response([{ info: message, parts: [part] }]), response([{ info: message, parts: [fetched] }])),
  469. )
  470. await store.sync("child")
  471. store.apply({
  472. type: "message.part.delta",
  473. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: "bc" },
  474. })
  475. await store.sync("child", { force: true })
  476. expect(store.data.part[message.id]).toEqual([{ ...part, text: "abc" }])
  477. expect(store.data.part_text_accum_delta[part.id]).toBe("abc")
  478. })
  479. test("clears delta state after exact server catch-up", async () => {
  480. const message = userMessage("message")
  481. const part = textPart(message.id, { text: "a" })
  482. const fetched = { ...part, text: "ab" }
  483. const store = createServerSession(
  484. messageClient(response([{ info: message, parts: [part] }]), response([{ info: message, parts: [fetched] }])),
  485. )
  486. await store.sync("child")
  487. store.apply({
  488. type: "message.part.delta",
  489. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: "b" },
  490. })
  491. await store.sync("child", { force: true })
  492. expect(store.data.part[message.id]).toEqual([fetched])
  493. expect(store.data.part_text_accum_delta[part.id]).toBeUndefined()
  494. })
  495. test("uses the successful retry response over events from a failed attempt", async () => {
  496. const failed = Promise.withResolvers<MessageResponse>()
  497. const retried = Promise.withResolvers<MessageResponse>()
  498. const message = userMessage("message")
  499. const stale = textPart(message.id, { text: "stale" })
  500. const intermediate = { ...stale, text: "intermediate" }
  501. const fetched = { ...stale, text: "fetched" }
  502. const client = messageClient(failed.promise, retried.promise)
  503. const store = createServerSession(client, { retry: retryImmediately })
  504. store.apply({ type: "message.updated", properties: { info: message } })
  505. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: stale, time: 1 } })
  506. const loading = store.sync("child")
  507. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: intermediate, time: 2 } })
  508. failed.reject(new Error("failed to fetch"))
  509. await client.requested(2)
  510. retried.resolve(response([{ info: message, parts: [fetched] }]))
  511. await loading
  512. expect(store.data.part[message.id]).toEqual([fetched])
  513. })
  514. test("preserves non-durable deltas across message retries", async () => {
  515. const failed = Promise.withResolvers<MessageResponse>()
  516. const retried = Promise.withResolvers<MessageResponse>()
  517. const message = userMessage("message")
  518. const part = textPart(message.id, { text: "stale" })
  519. const client = messageClient(failed.promise, retried.promise)
  520. const store = createServerSession(client, { retry: retryImmediately })
  521. store.apply({ type: "message.updated", properties: { info: message } })
  522. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 1 } })
  523. const loading = store.sync("child")
  524. store.apply({
  525. type: "message.part.delta",
  526. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: " delta" },
  527. })
  528. failed.reject(new Error("failed to fetch"))
  529. await client.requested(2)
  530. retried.resolve(response([{ info: message, parts: [part] }]))
  531. await loading
  532. expect(store.data.part[message.id]).toEqual([{ ...part, text: "stale delta" }])
  533. })
  534. test("preserves part removals across message retries", async () => {
  535. const failed = Promise.withResolvers<MessageResponse>()
  536. const retried = Promise.withResolvers<MessageResponse>()
  537. const message = userMessage("message")
  538. const part = textPart(message.id)
  539. const client = messageClient(response([{ info: message, parts: [part] }]), failed.promise, retried.promise)
  540. const store = createServerSession(client, { retry: retryImmediately })
  541. await store.sync("child")
  542. const loading = store.sync("child", { force: true })
  543. store.apply({
  544. type: "message.part.removed",
  545. properties: { sessionID: "child", messageID: message.id, partID: part.id },
  546. })
  547. failed.reject(new Error("failed to fetch"))
  548. await client.requested(3)
  549. retried.resolve(response([{ info: message, parts: [part] }]))
  550. await loading
  551. expect(store.data.part[message.id]).toBeUndefined()
  552. })
  553. test("preserves message removals across message retries", async () => {
  554. const failed = Promise.withResolvers<MessageResponse>()
  555. const retried = Promise.withResolvers<MessageResponse>()
  556. const message = userMessage("message")
  557. const part = textPart(message.id)
  558. const client = messageClient(response([{ info: message, parts: [part] }]), failed.promise, retried.promise)
  559. const store = createServerSession(client, { retry: retryImmediately })
  560. await store.sync("child")
  561. const loading = store.sync("child", { force: true })
  562. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  563. failed.reject(new Error("failed to fetch"))
  564. await client.requested(3)
  565. retried.resolve(response([{ info: message, parts: [part] }]))
  566. await loading
  567. expect(store.data.message.child).toEqual([])
  568. expect(store.data.part[message.id]).toBeUndefined()
  569. })
  570. test("preserves optimistic re-adds across message retries", async () => {
  571. const failed = Promise.withResolvers<MessageResponse>()
  572. const retried = Promise.withResolvers<MessageResponse>()
  573. const message = userMessage("message")
  574. const stale = textPart(message.id, { id: "stale", text: "stale" })
  575. const optimistic = textPart(message.id, { id: "optimistic", text: "optimistic" })
  576. const client = messageClient(response([{ info: message, parts: [stale] }]), failed.promise, retried.promise)
  577. const store = createServerSession(client, { retry: retryImmediately })
  578. await store.sync("child")
  579. const loading = store.sync("child", { force: true })
  580. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  581. store.optimistic.add({ sessionID: "child", message, parts: [optimistic] })
  582. failed.reject(new Error("failed to fetch"))
  583. await client.requested(3)
  584. retried.resolve(response([{ info: message, parts: [stale] }]))
  585. await loading
  586. expect(store.data.message.child).toEqual([message])
  587. expect(store.data.part[message.id]).toEqual([optimistic])
  588. })
  589. test("accepts part omission from a successful retry after an earlier delta", async () => {
  590. const failed = Promise.withResolvers<MessageResponse>()
  591. const retried = Promise.withResolvers<MessageResponse>()
  592. const message = userMessage("message")
  593. const part = textPart(message.id)
  594. const client = messageClient(response([{ info: message, parts: [part] }]), failed.promise, retried.promise)
  595. const store = createServerSession(client, { retry: retryImmediately })
  596. await store.sync("child")
  597. const loading = store.sync("child", { force: true })
  598. store.apply({
  599. type: "message.part.delta",
  600. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: " delta" },
  601. })
  602. failed.reject(new Error("failed to fetch"))
  603. await client.requested(3)
  604. retried.resolve(response([{ info: message, parts: [] }]))
  605. await loading
  606. expect(store.data.part[message.id]).toBeUndefined()
  607. expect(store.data.part_text_accum_delta[part.id]).toBeUndefined()
  608. })
  609. test("clears load-owned orphan parts when all retries fail", async () => {
  610. const first = Promise.withResolvers<MessageResponse>()
  611. const second = Promise.withResolvers<MessageResponse>()
  612. const third = Promise.withResolvers<MessageResponse>()
  613. const message = userMessage("message")
  614. const part = textPart(message.id)
  615. const client = messageClient(first.promise, second.promise, third.promise)
  616. const store = createServerSession(client, { retry: retryImmediately })
  617. const loading = store.sync("child").catch((error) => error)
  618. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } })
  619. first.reject(new Error("failed to fetch"))
  620. await client.requested(2)
  621. second.reject(new Error("failed to fetch"))
  622. await client.requested(3)
  623. third.reject(new Error("failed to fetch"))
  624. await loading
  625. expect(store.data.part[message.id]).toBeUndefined()
  626. })
  627. test("preserves live updates during a forced refresh", async () => {
  628. const pending = deferredResponse()
  629. const stale = userMessage("message")
  630. const stalePart = textPart(stale.id, { text: "stale" })
  631. const store = createServerSession(messageClient(response([{ info: stale, parts: [stalePart] }]), pending.promise))
  632. await store.sync("child")
  633. const refreshing = store.sync("child", { force: true })
  634. const live = { ...stale, time: { created: 2 } }
  635. store.apply({ type: "message.updated", properties: { info: live } })
  636. store.apply({
  637. type: "message.part.delta",
  638. properties: { sessionID: "child", messageID: stale.id, partID: stalePart.id, field: "text", delta: " live" },
  639. })
  640. pending.resolve(response([{ info: stale, parts: [stalePart] }]))
  641. await refreshing
  642. expect(store.data.message.child).toEqual([live])
  643. expect(store.data.part[stale.id]).toEqual([{ ...stalePart, text: "stale live" }])
  644. })
  645. test("keeps fetched message metadata when only a part changes", async () => {
  646. const pending = deferredResponse()
  647. const stale = userMessage("message")
  648. const fetched = { ...stale, time: { created: 2 } }
  649. const part = textPart(stale.id, { text: "stale" })
  650. const store = createServerSession(messageClient(response([{ info: stale, parts: [part] }]), pending.promise))
  651. await store.sync("child")
  652. const refreshing = store.sync("child", { force: true })
  653. store.apply({
  654. type: "message.part.delta",
  655. properties: { sessionID: "child", messageID: stale.id, partID: part.id, field: "text", delta: " live" },
  656. })
  657. pending.resolve(response([{ info: fetched, parts: [part] }]))
  658. await refreshing
  659. expect(store.data.message.child).toEqual([fetched])
  660. expect(store.data.part[stale.id]).toEqual([{ ...part, text: "stale live" }])
  661. })
  662. test("preserves a part update when a forced refresh omits its message", async () => {
  663. const pending = deferredResponse()
  664. const message = userMessage("message")
  665. const stale = textPart(message.id, { text: "stale" })
  666. const live = { ...stale, text: "live" }
  667. const store = createServerSession(messageClient(response([{ info: message, parts: [stale] }]), pending.promise))
  668. await store.sync("child")
  669. const refreshing = store.sync("child", { force: true })
  670. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: live, time: 2 } })
  671. pending.resolve(response())
  672. await refreshing
  673. expect(store.data.message.child).toEqual([message])
  674. expect(store.data.part[message.id]).toEqual([live])
  675. })
  676. test("ignores a late part update after its message is removed", async () => {
  677. const pending = deferredResponse()
  678. const message = userMessage("message")
  679. const part = textPart(message.id)
  680. const store = createServerSession(messageClient(pending.promise))
  681. const loading = store.sync("child")
  682. store.apply({ type: "message.updated", properties: { info: message } })
  683. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  684. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } })
  685. pending.resolve(response([{ info: message, parts: [part] }]))
  686. await loading
  687. expect(store.data.message.child).toEqual([])
  688. expect(store.data.part[message.id]).toBeUndefined()
  689. })
  690. test("ignores a late part update after a completed message removal", () => {
  691. const message = userMessage("message")
  692. const part = textPart(message.id)
  693. const store = setup({ child: session("child") }).store
  694. store.apply({ type: "message.updated", properties: { info: message } })
  695. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  696. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } })
  697. expect(store.data.part[message.id]).toBeUndefined()
  698. })
  699. test("does not restore a completed message removal from a stale refresh", async () => {
  700. const message = userMessage("message")
  701. const part = textPart(message.id)
  702. const store = createServerSession(
  703. messageClient(response([{ info: message, parts: [part] }]), response([{ info: message, parts: [part] }])),
  704. )
  705. await store.sync("child")
  706. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } })
  707. await store.sync("child", { force: true })
  708. expect(store.data.message.child).toEqual([])
  709. expect(store.data.part[message.id]).toBeUndefined()
  710. })
  711. test("does not restore a completed part removal from a stale refresh", async () => {
  712. const message = userMessage("message")
  713. const part = textPart(message.id)
  714. const store = createServerSession(
  715. messageClient(response([{ info: message, parts: [part] }]), response([{ info: message, parts: [part] }])),
  716. )
  717. await store.sync("child")
  718. store.apply({
  719. type: "message.part.removed",
  720. properties: { sessionID: "child", messageID: message.id, partID: part.id },
  721. })
  722. await store.sync("child", { force: true })
  723. expect(store.data.part[message.id]).toBeUndefined()
  724. })
  725. test("does not cache skipped optimistic parts", () => {
  726. const message = userMessage("message")
  727. const part = { id: "part", sessionID: "child", messageID: message.id, type: "step-start" as const }
  728. const store = setup({ child: session("child") }).store
  729. store.optimistic.add({ sessionID: "child", message, parts: [part] })
  730. expect(store.data.part[message.id]).toEqual([])
  731. })
  732. test("clears stale delta buffers when replacing optimistic parts", () => {
  733. const message = userMessage("message")
  734. const stale = textPart(message.id, { id: "stale", text: "stale" })
  735. const optimistic = textPart(message.id, { id: "optimistic", text: "optimistic" })
  736. const store = setup({ child: session("child") }).store
  737. store.optimistic.add({ sessionID: "child", message, parts: [stale] })
  738. store.apply({
  739. type: "message.part.delta",
  740. properties: { sessionID: "child", messageID: message.id, partID: stale.id, field: "text", delta: " delta" },
  741. })
  742. store.optimistic.add({ sessionID: "child", message, parts: [optimistic] })
  743. expect(store.data.part_text_accum_delta[stale.id]).toBeUndefined()
  744. expect(store.data.part_text_accum_delta[optimistic.id]).toBeUndefined()
  745. })
  746. test("preserves removals during history prepend", async () => {
  747. const pending = deferredResponse()
  748. const latest = userMessage("message-2", { time: { created: 2 } })
  749. const older = { ...latest, id: "message-1", time: { created: 1 } }
  750. const store = createServerSession(messageClient(response([{ info: latest, parts: [] }], "older"), pending.promise))
  751. await store.sync("child")
  752. const loading = store.history.loadMore("child")
  753. store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: older.id } })
  754. pending.resolve(response([{ info: older, parts: [] }]))
  755. await loading
  756. expect(store.data.message.child).toEqual([latest])
  757. })
  758. test("preserves loaded history during an incomplete refresh", async () => {
  759. const older = userMessage("message-1")
  760. const latest = userMessage("message-2", { time: { created: 2 } })
  761. const fresh = userMessage("message-3", { time: { created: 3 } })
  762. const store = createServerSession(
  763. messageClient(
  764. response(
  765. [
  766. { info: older, parts: [] },
  767. { info: latest, parts: [] },
  768. ],
  769. "older",
  770. ),
  771. response(
  772. [
  773. { info: latest, parts: [] },
  774. { info: fresh, parts: [] },
  775. ],
  776. "older",
  777. ),
  778. ),
  779. )
  780. await store.sync("child")
  781. await store.sync("child", { force: true })
  782. expect(store.data.message.child).toEqual([older, latest, fresh])
  783. })
  784. test("drops stale recent messages omitted by an incomplete refresh", async () => {
  785. const third = userMessage("message-3", { time: { created: 3 } })
  786. const fourth = userMessage("message-4", { time: { created: 4 } })
  787. const stale = userMessage("message-5", { time: { created: 5 } })
  788. const store = createServerSession(
  789. messageClient(
  790. response(
  791. [
  792. { info: fourth, parts: [] },
  793. { info: stale, parts: [] },
  794. ],
  795. "older",
  796. ),
  797. response(
  798. [
  799. { info: third, parts: [] },
  800. { info: fourth, parts: [] },
  801. ],
  802. "older",
  803. ),
  804. ),
  805. )
  806. await store.sync("child")
  807. await store.sync("child", { force: true })
  808. expect(store.data.message.child).toEqual([third, fourth])
  809. })
  810. test("uses message creation time for incomplete refresh boundaries", async () => {
  811. const older = userMessage("msg_z", { time: { created: 1 } })
  812. const boundary = userMessage("msg_m", { time: { created: 2 } })
  813. const stale = userMessage("msg_a", { time: { created: 3 } })
  814. const store = createServerSession(
  815. messageClient(
  816. response(
  817. [
  818. { info: older, parts: [] },
  819. { info: stale, parts: [] },
  820. ],
  821. "older",
  822. ),
  823. response([{ info: boundary, parts: [] }], "older"),
  824. ),
  825. )
  826. await store.sync("child")
  827. await store.sync("child", { force: true })
  828. expect(store.data.message.child).toEqual([boundary, older])
  829. })
  830. test("preserves a part update for a message being loaded from history", async () => {
  831. const pending = deferredResponse()
  832. const latest = userMessage("message-2", { time: { created: 2 } })
  833. const older = userMessage("message-1")
  834. const stale = textPart(older.id, { text: "stale" })
  835. const live = { ...stale, text: "live" }
  836. const store = createServerSession(messageClient(response([{ info: latest, parts: [] }], "older"), pending.promise))
  837. await store.sync("child")
  838. const loading = store.history.loadMore("child")
  839. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part: live, time: 2 } })
  840. pending.resolve(response([{ info: older, parts: [stale] }]))
  841. await loading
  842. expect(store.data.part[older.id]).toEqual([live])
  843. })
  844. test("does not clear newer orphan parts after terminal history prepend", async () => {
  845. const pending = deferredResponse()
  846. const latest = userMessage("message-2", { time: { created: 2 } })
  847. const older = userMessage("message-1")
  848. const newer = userMessage("message-3", { time: { created: 3 } })
  849. const part = textPart(newer.id, { text: "live" })
  850. const store = createServerSession(messageClient(response([{ info: latest, parts: [] }], "older"), pending.promise))
  851. await store.sync("child")
  852. const loading = store.history.loadMore("child")
  853. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 3 } })
  854. pending.resolve(response([{ info: older, parts: [] }]))
  855. await loading
  856. store.apply({ type: "message.updated", properties: { sessionID: "child", info: newer } })
  857. expect(store.data.part[newer.id]).toEqual([part])
  858. })
  859. test("accepts an authoritative history part after an earlier unknown-parent update", async () => {
  860. const pending = deferredResponse()
  861. const history = deferredResponse()
  862. const latest = userMessage("message-2", { time: { created: 2 } })
  863. const older = userMessage("message-1")
  864. const part = textPart(older.id, { text: "live" })
  865. const store = createServerSession(messageClient(pending.promise, history.promise))
  866. const loading = store.sync("child")
  867. store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } })
  868. pending.resolve(response([{ info: latest, parts: [] }], "older"))
  869. await loading
  870. expect(store.data.part[older.id]).toEqual([part])
  871. const loadingHistory = store.history.loadMore("child")
  872. history.resolve(response([{ info: older, parts: [{ ...part, text: "stale" }] }]))
  873. await loadingHistory
  874. expect(store.data.part[older.id]).toEqual([{ ...part, text: "stale" }])
  875. })
  876. test("preserves an unknown-parent part removal across pages", async () => {
  877. const initial = deferredResponse()
  878. const history = deferredResponse()
  879. const latest = userMessage("message-2", { time: { created: 2 } })
  880. const older = userMessage("message-1")
  881. const part = textPart(older.id)
  882. const store = createServerSession(messageClient(initial.promise, history.promise))
  883. const loading = store.sync("child")
  884. store.apply({
  885. type: "message.part.removed",
  886. properties: { sessionID: "child", messageID: older.id, partID: part.id },
  887. })
  888. initial.resolve(response([{ info: latest, parts: [] }], "older"))
  889. await loading
  890. const loadingHistory = store.history.loadMore("child")
  891. history.resolve(response([{ info: older, parts: [part] }]))
  892. await loadingHistory
  893. expect(store.data.part[older.id]).toBeUndefined()
  894. })
  895. test("clears orphaned parts when a refresh drops a message", async () => {
  896. const message = userMessage("message")
  897. const part = textPart(message.id, { text: "stale" })
  898. const store = createServerSession(messageClient(response([{ info: message, parts: [part] }]), response()))
  899. await store.sync("child")
  900. store.apply({
  901. type: "message.part.delta",
  902. properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: " delta" },
  903. })
  904. await store.sync("child", { force: true })
  905. expect(store.data.message.child).toEqual([])
  906. expect(store.data.part[message.id]).toBeUndefined()
  907. expect(store.data.part_text_accum_delta[part.id]).toBeUndefined()
  908. })
  909. test("applies events without a directory store", () => {
  910. const ctx = setup({})
  911. ctx.store.apply({ type: "session.created", properties: { sessionID: "root", info: session("root") } })
  912. ctx.store.apply({ type: "session.status", properties: { sessionID: "root", status: { type: "busy" } } })
  913. expect(ctx.store.get("root")?.directory).toBe("/repo")
  914. expect(ctx.store.data.session_working("root")).toBe(true)
  915. expect(ctx.get).toEqual([])
  916. })
  917. test("preserves pinned session content under server-wide cache pressure", () => {
  918. const ctx = setup({})
  919. ctx.store.pin("active")
  920. ctx.store.optimistic.add({
  921. sessionID: "active",
  922. message: {
  923. id: "message",
  924. sessionID: "active",
  925. role: "assistant",
  926. time: { created: 1 },
  927. parentID: "parent",
  928. modelID: "model",
  929. providerID: "provider",
  930. mode: "build",
  931. agent: "agent",
  932. path: { cwd: "/repo", root: "/repo" },
  933. cost: 0,
  934. tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
  935. },
  936. parts: [],
  937. })
  938. for (let index = 0; index < 50; index++) {
  939. ctx.store.remember(session(`session-${index}`))
  940. ctx.store.apply({
  941. type: "session.status",
  942. properties: { sessionID: `session-${index}`, status: { type: "idle" } },
  943. })
  944. }
  945. expect(ctx.store.data.message.active?.map((message) => message.id)).toEqual(["message"])
  946. expect(ctx.store.data.session_status["session-0"]).toBeUndefined()
  947. })
  948. })