conversation-assembler.spec.ts 36 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959
  1. import { describe, expect, it, vi } from 'vitest'
  2. import type { SessionEvent } from '@deepseek-ai/dsh-session/types'
  3. import { ConversationNodeAssembler } from '../src/client/sessions/conversation-assembler.ts'
  4. import type {
  5. ConversationEventInput, ConversationMatch, ConversationNodeContext,
  6. ConversationNodeDefinition, ConversationViewDefinition, ConversationViewNode,
  7. } from '../src/client/contract/conversation.ts'
  8. interface ScopeProbeStepData {
  9. readonly value: number
  10. }
  11. interface ScopeProbeTurnData {
  12. readonly valueSeenFromStep: number
  13. }
  14. declare module '../src/client/contract/conversation.ts' {
  15. interface ConversationStepDataMap {
  16. 'scope-probe': ScopeProbeStepData
  17. }
  18. interface ConversationTurnDataMap {
  19. 'scope-probe': ScopeProbeTurnData
  20. }
  21. }
  22. interface TestSnapshot {
  23. readonly order: readonly string[]
  24. readonly nodes: ReadonlyMap<string, ConversationViewNode>
  25. }
  26. class TestEventDefinitions {
  27. constructor(
  28. readonly definitions: readonly ConversationNodeDefinition[],
  29. readonly fallback?: ConversationNodeDefinition,
  30. ) {}
  31. entries(): readonly ConversationNodeDefinition[] {
  32. return this.definitions
  33. }
  34. fallbackEntry(): ConversationNodeDefinition | undefined {
  35. return this.fallback
  36. }
  37. }
  38. class TestViewDefinitions {
  39. constructor(readonly definitions: readonly ConversationViewDefinition[]) {}
  40. entries(): readonly ConversationViewDefinition[] {
  41. return this.definitions
  42. }
  43. }
  44. function testView(
  45. apply = vi.fn(),
  46. ): ConversationViewDefinition<ConversationViewNode, TestSnapshot> {
  47. return {
  48. target: 'chat',
  49. create: () => {
  50. let current: TestSnapshot = { order: [], nodes: new Map() }
  51. return {
  52. empty: current,
  53. replace: ({ nodes }) => {
  54. current = { order: nodes.map(node => node.key), nodes: new Map(nodes.map(node => [node.key, node])) }
  55. return current
  56. },
  57. apply: ({ upserts }) => {
  58. apply(upserts)
  59. const nodes = new Map(current.nodes)
  60. const order = [...current.order]
  61. for (const node of upserts) {
  62. if (!nodes.has(node.key)) order.push(node.key)
  63. nodes.set(node.key, node)
  64. }
  65. current = { order, nodes }
  66. return current
  67. },
  68. }
  69. },
  70. }
  71. }
  72. function at(seq: number, type: string, data: unknown): SessionEvent {
  73. return { seq, time: 1_700_000_000_000 + seq, type, data } as SessionEvent
  74. }
  75. function input(event: SessionEvent): ConversationEventInput {
  76. return { event, view: undefined }
  77. }
  78. function chatSnapshot(assembler: ConversationNodeAssembler): TestSnapshot | undefined {
  79. return assembler.snapshot('chat') as TestSnapshot | undefined
  80. }
  81. function node(context: Parameters<ConversationNodeDefinition['buildViewNode']>[0], data: unknown): ConversationViewNode {
  82. return {
  83. key: context.key,
  84. kind: context.kind,
  85. id: context.id,
  86. target: 'chat',
  87. data,
  88. }
  89. }
  90. describe('ConversationNodeAssembler', () => {
  91. it('appends through an exact business-id Context without replaying unrelated Contexts', () => {
  92. const starts = vi.fn((
  93. _context: ConversationNodeContext<{ callSeq: number; results: number }>,
  94. match: ConversationMatch,
  95. ) => ({ callSeq: match.event.seq, results: 0 }))
  96. const updates = vi.fn((context: { state: { callSeq: number; results: number } }) => ({
  97. ...context.state,
  98. results: context.state.results + 1,
  99. }))
  100. const definition: ConversationNodeDefinition<{ callSeq: number; results: number }> = {
  101. kind: 'tool',
  102. match: (event) => {
  103. if (event.type === 'tool/call') return { id: String(event.data.callId), role: 'start' }
  104. if (event.type === 'tool/result') return { id: String(event.data.message.source.callId), role: 'update' }
  105. return null
  106. },
  107. start: starts,
  108. update: updates,
  109. buildViewNode: context => node(context, context.state),
  110. }
  111. const assembler = new ConversationNodeAssembler(
  112. new TestEventDefinitions([definition]),
  113. new TestViewDefinitions([testView()]),
  114. )
  115. assembler.replaceWindow([
  116. input(at(1, 'tool/call', { turn: 1, step: 1, callId: 'a', name: 'x', arguments: '{}' })),
  117. input(at(2, 'tool/call', { turn: 1, step: 1, callId: 'b', name: 'x', arguments: '{}' })),
  118. ], false)
  119. assembler.flush()
  120. starts.mockClear()
  121. assembler.append(input(at(3, 'tool/result', {
  122. turn: 1,
  123. step: 1,
  124. message: { source: { type: 'tool-result', callId: 'a' }, content: [], isError: false },
  125. })))
  126. assembler.flush()
  127. expect(starts).not.toHaveBeenCalled()
  128. expect(updates).toHaveBeenCalledOnce()
  129. const snapshot = chatSnapshot(assembler)
  130. expect([...snapshot?.nodes.values() ?? []].map(value => value.data)).toEqual([
  131. { callSeq: 1, results: 1 },
  132. { callSeq: 2, results: 0 },
  133. ])
  134. })
  135. it('keeps one Match collection while a long Context appends without replay', () => {
  136. const starts = vi.fn(() => 0)
  137. const updates = vi.fn((context: ConversationNodeContext<number> & { readonly state: number }) => (
  138. context.state + 1
  139. ))
  140. const matchCollections = new Set<readonly ConversationMatch[]>()
  141. const definition: ConversationNodeDefinition<number> = {
  142. kind: 'append-linear',
  143. match: (event) => {
  144. const type: string = event.type
  145. if (type === 'linear/start') return { id: 'one', role: 'start' }
  146. if (type === 'linear/update') return { id: 'one', role: 'update' }
  147. return null
  148. },
  149. start: (context) => {
  150. matchCollections.add(context.matches)
  151. return starts()
  152. },
  153. update: (context) => {
  154. matchCollections.add(context.matches)
  155. return updates(context)
  156. },
  157. buildViewNode: context => node(context, context.state),
  158. }
  159. const assembler = new ConversationNodeAssembler(
  160. new TestEventDefinitions([definition]),
  161. new TestViewDefinitions([testView()]),
  162. )
  163. assembler.replaceWindow([input(at(1, 'linear/start', {}))], false)
  164. starts.mockClear()
  165. for (let seq = 2; seq <= 1_001; seq++) {
  166. assembler.append(input(at(seq, 'linear/update', {})))
  167. }
  168. assembler.flush()
  169. expect(starts).not.toHaveBeenCalled()
  170. expect(updates).toHaveBeenCalledTimes(1_000)
  171. expect(matchCollections.size).toBe(1)
  172. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(1_000)
  173. })
  174. it('merges an older page and replays its affected Context once', () => {
  175. const starts = vi.fn(() => 0)
  176. const updates = vi.fn((context: ConversationNodeContext<number> & { readonly state: number }) => (
  177. context.state + 1
  178. ))
  179. const definition: ConversationNodeDefinition<number> = {
  180. kind: 'prepend-linear',
  181. match: (event) => {
  182. const type: string = event.type
  183. if (type === 'linear/start') return { id: 'one', role: 'start' }
  184. if (type === 'linear/update') return { id: 'one', role: 'update' }
  185. return null
  186. },
  187. start: starts,
  188. update: updates,
  189. buildViewNode: context => node(context, context.state),
  190. }
  191. const assembler = new ConversationNodeAssembler(
  192. new TestEventDefinitions([definition]),
  193. new TestViewDefinitions([testView()]),
  194. )
  195. const current = Array.from({ length: 100 }, (_, index) => (
  196. input(at(index + 102, 'linear/update', {}))
  197. ))
  198. assembler.replaceWindow(current, true)
  199. assembler.flush()
  200. expect(starts).not.toHaveBeenCalled()
  201. expect(updates).not.toHaveBeenCalled()
  202. const older = [
  203. input(at(1, 'linear/start', {})),
  204. ...Array.from({ length: 100 }, (_, index) => (
  205. input(at(index + 2, 'linear/update', {}))
  206. )),
  207. ]
  208. assembler.prepend(older, false)
  209. assembler.flush()
  210. expect(starts).toHaveBeenCalledOnce()
  211. expect(updates).toHaveBeenCalledTimes(200)
  212. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(200)
  213. })
  214. it('collects an update before its start and replays it once prepend supplies the start', () => {
  215. const updates = vi.fn((context: { state: { settled: boolean } }) => ({ ...context.state, settled: true }))
  216. const definition: ConversationNodeDefinition<{ settled: boolean }> = {
  217. kind: 'tool',
  218. match: (event) => {
  219. if (event.type === 'tool/call') return { id: String(event.data.callId), role: 'start' }
  220. if (event.type === 'tool/result') return { id: String(event.data.message.source.callId), role: 'update' }
  221. return null
  222. },
  223. start: () => ({ settled: false }),
  224. update: updates,
  225. buildViewNode: context => node(context, context.state ?? { pendingStart: true }),
  226. }
  227. const assembler = new ConversationNodeAssembler(
  228. new TestEventDefinitions([definition]),
  229. new TestViewDefinitions([testView()]),
  230. )
  231. assembler.replaceWindow([input(at(10, 'tool/result', {
  232. turn: 1,
  233. step: 1,
  234. message: { source: { type: 'tool-result', callId: 'a' }, content: [], isError: false },
  235. }))], true)
  236. assembler.flush()
  237. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data)
  238. .toEqual({ pendingStart: true })
  239. assembler.prepend([input(at(5, 'tool/call', {
  240. turn: 1, step: 1, callId: 'a', name: 'x', arguments: '{}',
  241. }))], false)
  242. assembler.flush()
  243. expect(updates).toHaveBeenCalledOnce()
  244. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data)
  245. .toEqual({ settled: true })
  246. })
  247. it('rejects a Definition whose declared start follows an update in log order', () => {
  248. const definition: ConversationNodeDefinition<null> = {
  249. kind: 'invalid-lifecycle',
  250. match: event => event.type === 'turn/end'
  251. ? { id: 'one', role: 'start' }
  252. : event.type === 'turn/start' ? { id: 'one', role: 'update' } : null,
  253. start: () => null,
  254. update: context => context.state,
  255. buildViewNode: () => null,
  256. }
  257. const assembler = new ConversationNodeAssembler(
  258. new TestEventDefinitions([definition]),
  259. new TestViewDefinitions([testView()]),
  260. )
  261. expect(() => assembler.replaceWindow([
  262. input(at(1, 'turn/start', { turn: 1 })),
  263. input(at(2, 'turn/end', { turn: 1, reason: { kind: 'completed' } })),
  264. ], false)).toThrow('received an update before its start Match')
  265. })
  266. it('replays a window-gap reader when prepend supplies a nearer predecessor', () => {
  267. const source: ConversationNodeDefinition<number> = {
  268. kind: 'source',
  269. match: event => event.type === 'user/message'
  270. ? { id: String(event.data.id), role: 'start' }
  271. : null,
  272. start: (_context, match) => Number((match.event.data as { value?: unknown }).value ?? 0),
  273. update: context => context.state,
  274. buildViewNode: () => null,
  275. }
  276. const consumerStart = vi.fn((
  277. _context: Parameters<ConversationNodeDefinition<number>['start']>[0],
  278. _match: Parameters<ConversationNodeDefinition<number>['start']>[1],
  279. reader: Parameters<ConversationNodeDefinition<number>['start']>[2],
  280. ) => reader.previous<number>('source')?.state ?? -1)
  281. const consumer: ConversationNodeDefinition<number> = {
  282. kind: 'consumer',
  283. match: event => event.type === 'assistant/message'
  284. ? { id: `${event.data.turn}:${event.data.step}`, role: 'start' }
  285. : null,
  286. start: consumerStart,
  287. update: context => context.state,
  288. buildViewNode: context => node(context, context.state),
  289. }
  290. const assembler = new ConversationNodeAssembler(
  291. new TestEventDefinitions([source, consumer]),
  292. new TestViewDefinitions([testView()]),
  293. )
  294. assembler.replaceWindow([input(at(10, 'assistant/message', {
  295. turn: 2, step: 1, message: { role: 'assistant', content: [] },
  296. }))], true)
  297. assembler.flush()
  298. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(-1)
  299. assembler.prepend([input(at(5, 'user/message', {
  300. id: 'm1', value: 7, content: [], source: { kind: 'user' },
  301. }))], false)
  302. assembler.flush()
  303. expect(consumerStart).toHaveBeenCalledTimes(2)
  304. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(7)
  305. })
  306. it('keeps the predecessor index ordered across prepend and append', () => {
  307. const source: ConversationNodeDefinition<number> = {
  308. kind: 'source',
  309. match: event => event.type === 'user/message'
  310. ? { id: String(event.data.id), role: 'start' }
  311. : null,
  312. start: (_context, match) => match.event.seq,
  313. update: context => context.state,
  314. buildViewNode: () => null,
  315. }
  316. const consumer: ConversationNodeDefinition<number> = {
  317. kind: 'consumer',
  318. match: event => event.type === 'assistant/message'
  319. ? { id: `${event.data.turn}:${event.data.step}`, role: 'start' }
  320. : null,
  321. start: (_context, _match, reader) => reader.previous<number>('source')?.state ?? -1,
  322. update: context => context.state,
  323. buildViewNode: context => node(context, context.state),
  324. }
  325. const assembler = new ConversationNodeAssembler(
  326. new TestEventDefinitions([source, consumer]),
  327. new TestViewDefinitions([testView()]),
  328. )
  329. assembler.replaceWindow([
  330. input(at(40, 'user/message', { id: 'm40', content: [], source: { kind: 'user' } })),
  331. input(at(50, 'assistant/message', {
  332. turn: 1, step: 1, message: { role: 'assistant', content: [] },
  333. })),
  334. ], true)
  335. assembler.flush()
  336. assembler.prepend([
  337. input(at(10, 'user/message', { id: 'm10', content: [], source: { kind: 'user' } })),
  338. input(at(30, 'user/message', { id: 'm30', content: [], source: { kind: 'user' } })),
  339. ], false)
  340. assembler.flush()
  341. assembler.append(input(at(60, 'user/message', {
  342. id: 'm60', content: [], source: { kind: 'user' },
  343. })))
  344. assembler.append(input(at(70, 'assistant/message', {
  345. turn: 2, step: 1, message: { role: 'assistant', content: [] },
  346. })))
  347. assembler.flush()
  348. expect([...chatSnapshot(assembler)?.nodes.values() ?? []].map(value => value.data))
  349. .toEqual([40, 60])
  350. })
  351. it('replays a window-gap reader when an empty prepend closes the unknown prefix', () => {
  352. const consumerStart = vi.fn((
  353. _context: Parameters<ConversationNodeDefinition<number>['start']>[0],
  354. _match: Parameters<ConversationNodeDefinition<number>['start']>[1],
  355. reader: Parameters<ConversationNodeDefinition<number>['start']>[2],
  356. ) => reader.previous<number>('source')?.state ?? -1)
  357. const consumer: ConversationNodeDefinition<number> = {
  358. kind: 'consumer',
  359. match: event => event.type === 'assistant/message'
  360. ? { id: `${event.data.turn}:${event.data.step}`, role: 'start' }
  361. : null,
  362. start: consumerStart,
  363. update: context => context.state,
  364. buildViewNode: context => node(context, context.state),
  365. }
  366. const assembler = new ConversationNodeAssembler(
  367. new TestEventDefinitions([consumer]),
  368. new TestViewDefinitions([testView()]),
  369. )
  370. assembler.replaceWindow([input(at(10, 'assistant/message', {
  371. turn: 2, step: 1, message: { role: 'assistant', content: [] },
  372. }))], true)
  373. assembler.flush()
  374. expect(assembler.prepend([], false)).toBe('immediate')
  375. assembler.flush()
  376. expect(consumerStart).toHaveBeenCalledTimes(2)
  377. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(-1)
  378. })
  379. it('replays direct dependents when an append revises their predecessor Context', () => {
  380. const source: ConversationNodeDefinition<number> = {
  381. kind: 'source',
  382. match: (event) => {
  383. if (event.type === 'user/message') return { id: 'one', role: 'start' }
  384. if ((event.type as string) === 'source/update') return { id: 'one', role: 'update' }
  385. return null
  386. },
  387. start: () => 1,
  388. update: (_context, match) => (match.event.data as unknown as { value: number }).value,
  389. buildViewNode: () => null,
  390. }
  391. const consumerStart = vi.fn((
  392. _context: Parameters<ConversationNodeDefinition<number>['start']>[0],
  393. _match: Parameters<ConversationNodeDefinition<number>['start']>[1],
  394. reader: Parameters<ConversationNodeDefinition<number>['start']>[2],
  395. ) => reader.previous<number>('source')?.state ?? -1)
  396. const consumer: ConversationNodeDefinition<number> = {
  397. kind: 'consumer',
  398. match: event => event.type === 'assistant/message'
  399. ? { id: 'one', role: 'start' }
  400. : null,
  401. start: consumerStart,
  402. update: context => context.state,
  403. buildViewNode: context => node(context, context.state),
  404. }
  405. const assembler = new ConversationNodeAssembler(
  406. new TestEventDefinitions([source, consumer]),
  407. new TestViewDefinitions([testView()]),
  408. )
  409. assembler.replaceWindow([
  410. input(at(1, 'user/message', { id: 'source', content: [], source: { kind: 'user' } })),
  411. input(at(2, 'assistant/message', { turn: 1, step: 1, message: { role: 'assistant', content: [] } })),
  412. ], false)
  413. assembler.flush()
  414. expect(assembler.append(input(at(3, 'source/update', { value: 2 })))).toBe('immediate')
  415. assembler.flush()
  416. expect(consumerStart).toHaveBeenCalledTimes(2)
  417. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(2)
  418. })
  419. it('replays a transitive dependency closure in start order', () => {
  420. const sourceA: ConversationNodeDefinition<number> = {
  421. kind: 'diamond-a',
  422. match: (event) => {
  423. if (event.type === 'user/message') return { id: 'one', role: 'start' }
  424. if ((event.type as string) === 'diamond/a') return { id: 'one', role: 'update' }
  425. return null
  426. },
  427. start: () => 1,
  428. update: (_context, match) => (match.event.data as unknown as { value: number }).value,
  429. buildViewNode: () => null,
  430. }
  431. const sourceX: ConversationNodeDefinition<number> = {
  432. kind: 'diamond-x',
  433. match: (event) => {
  434. if (event.type === 'turn/start') return { id: 'one', role: 'start' }
  435. if ((event.type as string) === 'diamond/x') return { id: 'one', role: 'update' }
  436. return null
  437. },
  438. start: () => 10,
  439. update: (_context, match) => (match.event.data as unknown as { value: number }).value,
  440. buildViewNode: () => null,
  441. }
  442. const middle: ConversationNodeDefinition<number> = {
  443. kind: 'diamond-b',
  444. match: event => event.type === 'assistant/message'
  445. ? { id: 'one', role: 'start' }
  446. : null,
  447. start: (_context, _match, reader) => (
  448. (reader.previous<number>('diamond-a')?.state ?? 0)
  449. + (reader.previous<number>('diamond-x')?.state ?? 0)
  450. ),
  451. update: context => context.state,
  452. buildViewNode: context => node(context, context.state),
  453. }
  454. const consumer: ConversationNodeDefinition<number> = {
  455. kind: 'diamond-c',
  456. match: event => event.type === 'tool/call'
  457. ? { id: 'one', role: 'start' }
  458. : null,
  459. start: (_context, _match, reader) => (
  460. (reader.previous<number>('diamond-a')?.state ?? 0) * 100
  461. + (reader.previous<number>('diamond-b')?.state ?? 0)
  462. ),
  463. update: context => context.state,
  464. buildViewNode: context => node(context, context.state),
  465. }
  466. const assembler = new ConversationNodeAssembler(
  467. new TestEventDefinitions([sourceA, sourceX, middle, consumer]),
  468. new TestViewDefinitions([testView()]),
  469. )
  470. assembler.replaceWindow([
  471. input(at(1, 'user/message', { id: 'source', content: [], source: { kind: 'user' } })),
  472. input(at(2, 'turn/start', { turn: 1 })),
  473. input(at(3, 'assistant/message', { turn: 1, step: 1, message: { role: 'assistant', content: [] } })),
  474. input(at(4, 'tool/call', { turn: 1, step: 1, callId: 'call', name: 'x', arguments: '{}' })),
  475. ], false)
  476. assembler.append(input(at(5, 'diamond/x', { value: 20 })))
  477. assembler.append(input(at(6, 'diamond/a', { value: 2 })))
  478. assembler.flush()
  479. const value = [...chatSnapshot(assembler)?.nodes.values() ?? []]
  480. .find(candidate => candidate.kind === 'diamond-c')
  481. expect(value?.data).toBe(222)
  482. })
  483. it('replays Location-derived State and rebuilds only owned Nodes when a step closes', () => {
  484. const apply = vi.fn()
  485. const starts = vi.fn((
  486. _context: Parameters<ConversationNodeDefinition<string>['start']>[0],
  487. match: Parameters<ConversationNodeDefinition<string>['start']>[1],
  488. ) => match.location.kind === 'step' ? match.location.step.status : 'missing')
  489. const definition: ConversationNodeDefinition<string> = {
  490. kind: 'step',
  491. match: event => event.type === 'step/start'
  492. ? { id: `${event.data.turn}:${event.data.step}`, role: 'start' }
  493. : null,
  494. start: starts,
  495. update: context => context.state,
  496. buildViewNode: context => node(context, context.state),
  497. }
  498. const assembler = new ConversationNodeAssembler(
  499. new TestEventDefinitions([definition]),
  500. new TestViewDefinitions([testView(apply)]),
  501. )
  502. assembler.replaceWindow([
  503. input(at(1, 'turn/start', { turn: 1 })),
  504. input(at(2, 'step/start', { turn: 1, step: 1 })),
  505. ], false)
  506. assembler.flush()
  507. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe('open')
  508. assembler.append(input(at(3, 'step/end', { turn: 1, step: 1 })))
  509. assembler.flush()
  510. expect(starts).toHaveBeenCalledTimes(2)
  511. expect(apply).toHaveBeenCalledOnce()
  512. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe('closed')
  513. })
  514. it('lets one Context publish Step and Turn data in phase order', () => {
  515. interface State {
  516. readonly turn: number
  517. readonly step: number
  518. readonly value: number
  519. }
  520. const definition: ConversationNodeDefinition<State> = {
  521. kind: 'scope-probe',
  522. match: (event) => {
  523. if (event.type === 'step/start') {
  524. return { id: `${event.data.turn}:${event.data.step}`, role: 'start' }
  525. }
  526. if ((event.type as string) === 'scope-probe/update') {
  527. return { id: '1:1', role: 'update' }
  528. }
  529. return null
  530. },
  531. start: (_context, match) => {
  532. if (match.event.type !== 'step/start') throw new Error('scope probe requires step/start')
  533. return { turn: match.event.data.turn, step: match.event.data.step, value: 1 }
  534. },
  535. update: (_context, match) => ({
  536. turn: 1,
  537. step: 1,
  538. value: (match.event.data as unknown as { value: number }).value,
  539. }),
  540. buildLocationData: (context, scope) => {
  541. const state = context.state
  542. if (state === undefined) return null
  543. if (scope === 'step') {
  544. return {
  545. kind: 'step',
  546. turn: state.turn,
  547. step: state.step,
  548. key: 'scope-probe',
  549. value: { value: state.value },
  550. }
  551. }
  552. const location = context.start?.location
  553. const stepValue = location?.kind === 'step'
  554. ? location.step.data.get('scope-probe')?.value
  555. : undefined
  556. return {
  557. kind: 'turn',
  558. turn: state.turn,
  559. key: 'scope-probe',
  560. value: { valueSeenFromStep: stepValue ?? -1 },
  561. }
  562. },
  563. buildViewNode: (context) => {
  564. const location = context.start?.location
  565. if (location?.kind !== 'step') return null
  566. return node(context, {
  567. step: location.step.data.get('scope-probe')?.value,
  568. turn: location.turn.data.get('scope-probe')?.valueSeenFromStep,
  569. })
  570. },
  571. }
  572. const assembler = new ConversationNodeAssembler(
  573. new TestEventDefinitions([definition]),
  574. new TestViewDefinitions([testView()]),
  575. )
  576. assembler.replaceWindow([
  577. input(at(1, 'turn/start', { turn: 1 })),
  578. input(at(2, 'step/start', { turn: 1, step: 1 })),
  579. ], false)
  580. assembler.flush()
  581. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data)
  582. .toEqual({ step: 1, turn: 1 })
  583. assembler.append(input(at(3, 'scope-probe/update', { turn: 1, step: 1, value: 2 })))
  584. assembler.flush()
  585. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data)
  586. .toEqual({ step: 2, turn: 2 })
  587. })
  588. it('updates existing turn Locations when their Step membership changes', () => {
  589. const apply = vi.fn()
  590. const definition: ConversationNodeDefinition<null> = {
  591. kind: 'turn-probe',
  592. match: event => event.type === 'turn/start'
  593. ? { id: String(event.data.turn), role: 'start' }
  594. : null,
  595. start: () => null,
  596. update: context => context.state,
  597. buildViewNode: context => node(context, context.start?.location.kind === 'turn'
  598. ? context.start.location.turn.steps.length
  599. : -1),
  600. }
  601. const assembler = new ConversationNodeAssembler(
  602. new TestEventDefinitions([definition]),
  603. new TestViewDefinitions([testView(apply)]),
  604. )
  605. assembler.replaceWindow([input(at(1, 'turn/start', { turn: 1 }))], false)
  606. assembler.flush()
  607. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(0)
  608. assembler.append(input(at(2, 'step/start', { turn: 1, step: 1 })))
  609. assembler.flush()
  610. expect(apply).toHaveBeenCalledOnce()
  611. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(1)
  612. })
  613. it('publishes a changed timeline even when no business Definition claims the boundary', () => {
  614. const apply = vi.fn()
  615. const assembler = new ConversationNodeAssembler(
  616. new TestEventDefinitions([]),
  617. new TestViewDefinitions([testView(apply)]),
  618. )
  619. assembler.replaceWindow([], false)
  620. assembler.flush()
  621. assembler.append(input(at(1, 'turn/start', { turn: 1 })))
  622. assembler.flush()
  623. expect(apply).toHaveBeenCalledOnce()
  624. expect(chatSnapshot(assembler)?.order).toEqual([])
  625. })
  626. it('clears the prior Step at a new Turn and honors explicit session ownership', () => {
  627. const definition: ConversationNodeDefinition<null> = {
  628. kind: 'location-probe',
  629. match: (event) => {
  630. if ((event.type as string) === 'command/run') {
  631. return {
  632. id: (event.data as unknown as { commandId: string }).commandId,
  633. role: 'start',
  634. }
  635. }
  636. if ((event.type as string) === 'compact/start') {
  637. return {
  638. id: (event.data as unknown as { compactionId: string }).compactionId,
  639. role: 'start',
  640. }
  641. }
  642. return null
  643. },
  644. start: () => null,
  645. update: context => context.state,
  646. buildViewNode: (context) => {
  647. const location = context.start?.location
  648. const data = location?.kind === 'step'
  649. ? `step:${location.turn.turn}:${location.step.step}`
  650. : location?.kind === 'turn' ? `turn:${location.turn.turn}` : location?.kind
  651. return node(context, data)
  652. },
  653. }
  654. const assembler = new ConversationNodeAssembler(
  655. new TestEventDefinitions([definition]),
  656. new TestViewDefinitions([testView()]),
  657. )
  658. assembler.replaceWindow([
  659. input(at(1, 'turn/start', { turn: 1 })),
  660. input(at(2, 'step/start', { turn: 1, step: 1 })),
  661. input(at(3, 'turn/start', { turn: 2 })),
  662. input(at(4, 'command/run', { commandId: 'command', name: 'x' })),
  663. input(at(5, 'compact/start', { compactionId: 'compact', turn: null })),
  664. ], false)
  665. assembler.flush()
  666. expect([...chatSnapshot(assembler)?.nodes.values() ?? []].map(value => value.data))
  667. .toEqual(['turn:2', 'session'])
  668. })
  669. it('assigns turn boundaries to the Turn even when a Step remains open', () => {
  670. const definition: ConversationNodeDefinition<null> = {
  671. kind: 'turn-boundary-probe',
  672. match: event => event.type === 'turn/end'
  673. ? { id: String(event.data.turn), role: 'start' }
  674. : null,
  675. start: () => null,
  676. update: context => context.state,
  677. buildViewNode: context => node(context, context.start?.location.kind),
  678. }
  679. const assembler = new ConversationNodeAssembler(
  680. new TestEventDefinitions([definition]),
  681. new TestViewDefinitions([testView()]),
  682. )
  683. assembler.replaceWindow([
  684. input(at(1, 'turn/start', { turn: 1 })),
  685. input(at(2, 'step/start', { turn: 1, step: 1 })),
  686. ], false)
  687. assembler.flush()
  688. assembler.append(input(at(3, 'turn/end', { turn: 1, reason: { kind: 'aborted' } })))
  689. assembler.flush()
  690. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe('turn')
  691. })
  692. it('carries explicit coordinates across coordinate-free events in a partial window and live tail', () => {
  693. const definition: ConversationNodeDefinition<null> = {
  694. kind: 'location-probe',
  695. match: event => (event.type as string) === 'tool/code-dispatch-start'
  696. ? { id: String(event.seq), role: 'start' }
  697. : null,
  698. start: () => null,
  699. update: context => context.state,
  700. buildViewNode: (context) => {
  701. const location = context.start?.location
  702. return node(context, location?.kind === 'step'
  703. ? `${location.turn.turn}:${location.step.step}`
  704. : location?.kind)
  705. },
  706. }
  707. const assembler = new ConversationNodeAssembler(
  708. new TestEventDefinitions([definition]),
  709. new TestViewDefinitions([testView()]),
  710. )
  711. assembler.replaceWindow([
  712. input(at(10, 'tool/call', { turn: 2, step: 3, callId: 'root', name: 'x', arguments: '{}' })),
  713. input(at(11, 'tool/code-dispatch-start', { rootCallId: 'root', subCallId: 'a' })),
  714. ], true)
  715. assembler.flush()
  716. assembler.append(input(at(12, 'tool/code-dispatch-start', { rootCallId: 'root', subCallId: 'b' })))
  717. assembler.flush()
  718. expect([...chatSnapshot(assembler)?.nodes.values() ?? []].map(value => value.data))
  719. .toEqual(['2:3', '2:3'])
  720. })
  721. it('treats loaded end boundaries as closed when their starts precede the window', () => {
  722. const definition: ConversationNodeDefinition<null> = {
  723. kind: 'location-probe',
  724. match: event => event.type === 'tool/call'
  725. ? { id: String(event.data.callId), role: 'start' }
  726. : null,
  727. start: () => null,
  728. update: context => context.state,
  729. buildViewNode: (context) => {
  730. const location = context.start?.location
  731. return node(context, location?.kind === 'step'
  732. ? `${location.turn.status}:${location.step.status}`
  733. : location?.kind)
  734. },
  735. }
  736. const assembler = new ConversationNodeAssembler(
  737. new TestEventDefinitions([definition]),
  738. new TestViewDefinitions([testView()]),
  739. )
  740. assembler.replaceWindow([
  741. input(at(10, 'tool/call', { turn: 2, step: 3, callId: 'root', name: 'x', arguments: '{}' })),
  742. input(at(11, 'step/end', { turn: 2, step: 3 })),
  743. input(at(12, 'turn/end', { turn: 2, reason: { kind: 'completed' } })),
  744. ], true)
  745. assembler.flush()
  746. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data)
  747. .toBe('closed:closed')
  748. })
  749. it('restarts State creation from undefined when Location changes replay a Context', () => {
  750. const seen = vi.fn((context: Parameters<ConversationNodeDefinition<number>['start']>[0]) => {
  751. expect(context.state).toBeUndefined()
  752. return 1
  753. })
  754. const definition: ConversationNodeDefinition<number> = {
  755. kind: 'replay-probe',
  756. match: event => event.type === 'step/start'
  757. ? { id: `${event.data.turn}:${event.data.step}`, role: 'start' }
  758. : null,
  759. start: seen,
  760. update: context => context.state,
  761. buildViewNode: context => node(context, context.state),
  762. }
  763. const assembler = new ConversationNodeAssembler(
  764. new TestEventDefinitions([definition]),
  765. new TestViewDefinitions([testView()]),
  766. )
  767. assembler.replaceWindow([input(at(1, 'step/start', { turn: 1, step: 1 }))], false)
  768. assembler.flush()
  769. assembler.append(input(at(2, 'step/end', { turn: 1, step: 1 })))
  770. assembler.flush()
  771. expect(seen).toHaveBeenCalledTimes(2)
  772. })
  773. it('does not invoke the fallback when an ordinary non-rendering Definition claims an event', () => {
  774. const fallbackStart = vi.fn(() => 'fallback')
  775. const claimed: ConversationNodeDefinition<null> = {
  776. kind: 'claimed',
  777. match: event => (event.type as string) === 'command/run' ? { id: 'claimed', role: 'start' } : null,
  778. start: () => null,
  779. update: context => context.state,
  780. buildViewNode: () => null,
  781. }
  782. const fallback: ConversationNodeDefinition<string> = {
  783. kind: 'fallback',
  784. match: event => ({ id: String(event.seq), role: 'start' }),
  785. start: fallbackStart,
  786. update: context => context.state,
  787. buildViewNode: context => node(context, context.state),
  788. }
  789. const assembler = new ConversationNodeAssembler(
  790. new TestEventDefinitions([claimed], fallback),
  791. new TestViewDefinitions([testView()]),
  792. )
  793. assembler.replaceWindow([input(at(1, 'command/run', { commandId: 'one', name: 'x' }))], false)
  794. assembler.flush()
  795. expect(fallbackStart).not.toHaveBeenCalled()
  796. expect(chatSnapshot(assembler)?.order).toEqual([])
  797. })
  798. it('rejects withdrawing a previously materialized Node during an incremental update', () => {
  799. const definition: ConversationNodeDefinition<boolean> = {
  800. kind: 'toggle',
  801. match: (event) => {
  802. if ((event.type as string) === 'command/run') return { id: 'one', role: 'start' }
  803. if ((event.type as string) === 'toggle/hide') return { id: 'one', role: 'update' }
  804. return null
  805. },
  806. start: () => true,
  807. update: () => false,
  808. buildViewNode: context => context.state === true ? node(context, true) : null,
  809. }
  810. const assembler = new ConversationNodeAssembler(
  811. new TestEventDefinitions([definition]),
  812. new TestViewDefinitions([testView()]),
  813. )
  814. assembler.replaceWindow([input(at(1, 'command/run', { commandId: 'one', name: 'x' }))], false)
  815. assembler.flush()
  816. expect(chatSnapshot(assembler)?.order).toHaveLength(1)
  817. assembler.append(input(at(2, 'toggle/hide', {})))
  818. expect(() => assembler.flush()).toThrow(/withdrew materialized target "chat"/)
  819. expect(chatSnapshot(assembler)?.order).toHaveLength(1)
  820. })
  821. it('fails loud when a Definition returns undefined State', () => {
  822. const startUndefined: ConversationNodeDefinition = {
  823. kind: 'undefined-start',
  824. match: event => (event.type as string) === 'command/run' ? { id: 'one', role: 'start' } : null,
  825. start: () => undefined,
  826. update: context => context.state,
  827. buildViewNode: () => null,
  828. }
  829. const startAssembler = new ConversationNodeAssembler(
  830. new TestEventDefinitions([startUndefined]),
  831. new TestViewDefinitions([testView()]),
  832. )
  833. expect(() => startAssembler.replaceWindow([
  834. input(at(1, 'command/run', { commandId: 'one', name: 'x' })),
  835. ], false)).toThrow(/Definition "undefined-start" returned undefined from start/)
  836. const updateUndefined: ConversationNodeDefinition<boolean> = {
  837. kind: 'undefined-update',
  838. match: (event) => {
  839. if ((event.type as string) === 'command/run') return { id: 'one', role: 'start' }
  840. if ((event.type as string) === 'command/done') return { id: 'one', role: 'update' }
  841. return null
  842. },
  843. start: () => true,
  844. update: () => undefined as never,
  845. buildViewNode: context => node(context, context.state),
  846. }
  847. const updateAssembler = new ConversationNodeAssembler(
  848. new TestEventDefinitions([updateUndefined]),
  849. new TestViewDefinitions([testView()]),
  850. )
  851. updateAssembler.replaceWindow([
  852. input(at(1, 'command/run', { commandId: 'one', name: 'x' })),
  853. ], false)
  854. expect(() => updateAssembler.append(
  855. input(at(2, 'command/done', { commandId: 'one', kind: 'success' })),
  856. )).toThrow(/Definition "undefined-update" returned undefined from update/)
  857. })
  858. it('rejects a duplicate start before mutating the existing Context', () => {
  859. const definition: ConversationNodeDefinition<number> = {
  860. kind: 'single-start',
  861. match: event => (event.type as string) === 'command/run' ? { id: 'one', role: 'start' } : null,
  862. start: (_context, match) => match.event.seq,
  863. update: context => context.state,
  864. buildViewNode: context => node(context, context.state),
  865. }
  866. const assembler = new ConversationNodeAssembler(
  867. new TestEventDefinitions([definition]),
  868. new TestViewDefinitions([testView()]),
  869. )
  870. assembler.replaceWindow([
  871. input(at(1, 'command/run', { commandId: 'one', name: 'x' })),
  872. ], false)
  873. assembler.flush()
  874. expect(() => assembler.append(
  875. input(at(2, 'command/run', { commandId: 'two', name: 'x' })),
  876. )).toThrow(/received more than one start Match/)
  877. assembler.flush()
  878. expect([...chatSnapshot(assembler)?.nodes.values() ?? []][0]?.data).toBe(1)
  879. })
  880. })