session-history-journal.host.spec.ts 36 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923
  1. /** Raw Session journal transport and message-aligned pagination coverage. */
  2. import { describe, expect, it, vi } from 'vitest'
  3. import { Context } from '@deepseek-ai/cordis'
  4. import AgentRegistry, { type Agent, type AssistantStreamFrame } from '@deepseek-ai/dsh-agent'
  5. import SessionStore from '@deepseek-ai/dsh-session'
  6. import { LlmAttemptId, ToolCallId, createMessage, createToolResultMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
  7. import type { Session, SessionEvent, SessionId } from '@deepseek-ai/dsh-session'
  8. import { SessionHistoryController } from '@deepseek-ai/dsh-api-session-controller/src/history.ts'
  9. import type { SessionFollowFrame, SessionPage, SessionWireEvent } from '@deepseek-ai/dsh-api-session-controller/types'
  10. import { createSessionTestRemote, installSessionReadTestServices } from './test-remote.ts'
  11. /** Append a production-shaped human prompt to the session surface. */
  12. function appendUserText(session: Session, text: string): SessionEvent {
  13. return session.append('user/message', createUserMessage({
  14. content: [{ type: 'text', text }], source: { kind: 'user' },
  15. }), { surfaceOp: 'append' })
  16. }
  17. /** Append a production-shaped assistant message to the session surface. */
  18. function appendAssistantText(session: Session, text: string, step: number): SessionEvent {
  19. return session.append('assistant/message', {
  20. turn: 1,
  21. step,
  22. message: createMessage({
  23. role: 'assistant',
  24. content: [{ type: 'text', text }],
  25. source: { kind: 'model', provider: 'p', model: 'm' },
  26. }),
  27. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: [text] }],
  28. }, { surfaceOp: 'append' })
  29. }
  30. /**
  31. * Append a plugin-owned log-only event. The host proxy is projection-only, so it
  32. * declares no compaction vocabulary; the cast writes the real event shape without
  33. * depending on the owning package.
  34. */
  35. function appendExtension(session: Session, type: string, data: unknown): SessionEvent {
  36. return (session.append as unknown as (type: string, data: unknown) => SessionEvent)(type, data)
  37. }
  38. async function harness(): Promise<{ ctx: Context }> {
  39. const ctx = new Context()
  40. await ctx.plugin(SessionStore)
  41. await ctx.plugin(AgentRegistry)
  42. installSessionReadTestServices(ctx)
  43. return { ctx }
  44. }
  45. /** Drain one Session follow until `count` event frames arrive. */
  46. async function collect(
  47. iterable: AsyncIterable<SessionFollowFrame>,
  48. count: number,
  49. abort: AbortController,
  50. ): Promise<SessionFollowFrame[]> {
  51. const frames: SessionFollowFrame[] = []
  52. for await (const frame of iterable) {
  53. frames.push(frame)
  54. if (frames.filter(candidate => candidate.type === 'event').length >= count) abort.abort()
  55. }
  56. return frames
  57. }
  58. /** Open follow and wait until its cursor is fixed before appending fixtures. */
  59. async function openFollow(
  60. history: SessionHistoryController,
  61. sessionId: SessionId,
  62. signal: AbortSignal,
  63. ): Promise<AsyncIterable<SessionFollowFrame>> {
  64. const iterator = history.follow({
  65. address: { kind: 'session', sessionId },
  66. }, signal)[Symbol.asyncIterator]()
  67. await expect(iterator.next()).resolves.toMatchObject({
  68. done: false,
  69. value: { type: 'snapshot' },
  70. })
  71. return { [Symbol.asyncIterator]: () => iterator }
  72. }
  73. /** Abort one follow and await both its iterator and owning Context teardown. */
  74. async function disposeFollow(
  75. ctx: Context,
  76. iterator: AsyncIterator<SessionFollowFrame>,
  77. abort: AbortController,
  78. ): Promise<void> {
  79. abort.abort()
  80. await iterator.return?.()
  81. await ctx.fiber.dispose()
  82. }
  83. /** Read scalar v2 page records for assertions over the logical journal. */
  84. function pageEvents(page: SessionPage): SessionWireEvent[] {
  85. return page.records.map(record => record.event)
  86. }
  87. describe('Session history raw journal', () => {
  88. it('opens an empty opted-in Assistant baseline before any live attempt exists', async () => {
  89. const { ctx } = await harness()
  90. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  91. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  92. const abort = new AbortController()
  93. const iterator = history.follow({
  94. address: { kind: 'session', sessionId: session.id },
  95. assistantStream: true,
  96. }, abort.signal)[Symbol.asyncIterator]()
  97. await expect(iterator.next()).resolves.toMatchObject({
  98. done: false,
  99. value: { type: 'snapshot', assistantStream: { revision: 0 } },
  100. })
  101. abort.abort()
  102. await iterator.next()
  103. await ctx.fiber.dispose()
  104. })
  105. it('filters foreign and opening-baseline frames buffered during the source observation', async () => {
  106. const { ctx } = await harness()
  107. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  108. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  109. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  110. const originalObserve = ctx.sessionQuery.observeSession.bind(ctx.sessionQuery)
  111. const entered = Promise.withResolvers<undefined>()
  112. const release = Promise.withResolvers<undefined>()
  113. const observe = vi.spyOn(ctx.sessionQuery, 'observeSession').mockImplementation(async (...args) => {
  114. entered.resolve(undefined)
  115. await release.promise
  116. return originalObserve(...args)
  117. })
  118. const abort = new AbortController()
  119. const iterator = history.follow({
  120. address: { kind: 'session', sessionId: session.id },
  121. assistantStream: true,
  122. }, abort.signal)[Symbol.asyncIterator]()
  123. const opening = iterator.next()
  124. await entered.promise
  125. const attemptId = LlmAttemptId('buffered-opening-attempt')
  126. ctx.emit('agent/assistant-stream', {
  127. agent,
  128. frame: { type: 'start', attemptId, revision: 1, turn: 1, step: 1 },
  129. })
  130. ctx.emit('agent/assistant-stream', {
  131. agent,
  132. frame: {
  133. type: 'chunk', attemptId, revision: 2, index: 0,
  134. time: 2, chunk: { type: 'text-delta', index: 0, text: 'buffered' },
  135. },
  136. })
  137. const foreign = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  138. ctx.emit('agent/assistant-stream', {
  139. agent: { id: foreign.id, session: foreign, status: 'running', ctx } as Agent,
  140. frame: {
  141. type: 'start', attemptId: LlmAttemptId('foreign-attempt'), revision: 1, turn: 1, step: 1,
  142. },
  143. })
  144. release.resolve(undefined)
  145. await expect(opening).resolves.toMatchObject({
  146. done: false,
  147. value: { type: 'snapshot', assistantStream: { revision: 2 } },
  148. })
  149. const next = iterator.next()
  150. const durable = session.append('turn/start', { turn: 1 })
  151. await expect(next).resolves.toEqual({ done: false, value: { type: 'event', event: durable } })
  152. abort.abort()
  153. await iterator.next()
  154. observe.mockRestore()
  155. await ctx.fiber.dispose()
  156. })
  157. it('opens an opted-in assistant baseline and preserves mixed live FIFO order', async () => {
  158. const { ctx } = await harness()
  159. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  160. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  161. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  162. const attemptId = LlmAttemptId('live-follow-attempt')
  163. const emit = (frame: AssistantStreamFrame): void => {
  164. ctx.emit('agent/assistant-stream', { agent, frame })
  165. }
  166. emit({
  167. type: 'start', attemptId, revision: 1,
  168. turn: 1, step: 1,
  169. })
  170. const firstChunk = { type: 'text-delta', index: 0, text: 'a' } as const
  171. emit({
  172. type: 'chunk', attemptId, revision: 2, index: 0,
  173. time: 1, chunk: firstChunk,
  174. })
  175. const abort = new AbortController()
  176. const iterator = history.follow({
  177. address: { kind: 'session', sessionId: session.id },
  178. assistantStream: true,
  179. }, abort.signal)[Symbol.asyncIterator]()
  180. await expect(iterator.next()).resolves.toMatchObject({
  181. done: false,
  182. value: {
  183. type: 'snapshot',
  184. assistantStream: {
  185. revision: 2,
  186. activeAttempt: {
  187. attemptId,
  188. startedAfterSeq: -1,
  189. turn: 1,
  190. step: 1,
  191. nextIndex: 1,
  192. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['a'] }],
  193. },
  194. },
  195. },
  196. })
  197. const nextFrame: AssistantStreamFrame = {
  198. type: 'chunk', attemptId, revision: 3, index: 1,
  199. time: 2, chunk: { type: 'text-delta', index: 0, text: 'b' },
  200. }
  201. emit(nextFrame)
  202. const message = appendAssistantText(session, 'ab', 1)
  203. const endFrame: AssistantStreamFrame = {
  204. type: 'end', attemptId, revision: 4, index: 2,
  205. outcome: { kind: 'committed', eventType: 'assistant/message', seq: message.seq },
  206. }
  207. emit(endFrame)
  208. await expect(iterator.next()).resolves.toEqual({
  209. done: false, value: { type: 'assistant-stream', frame: nextFrame },
  210. })
  211. await expect(iterator.next()).resolves.toEqual({
  212. done: false, value: { type: 'event', event: message },
  213. })
  214. await expect(iterator.next()).resolves.toEqual({
  215. done: false, value: { type: 'assistant-stream', frame: endFrame },
  216. })
  217. abort.abort()
  218. await iterator.next()
  219. await ctx.fiber.dispose()
  220. })
  221. it('forwards revision one when the attached Agent lifecycle restarts after opening', async () => {
  222. const { ctx } = await harness()
  223. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  224. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  225. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  226. const attemptId = LlmAttemptId(`${session.id}:1`)
  227. ctx.emit('agent/assistant-stream', {
  228. agent,
  229. frame: {
  230. type: 'start', attemptId, revision: 1,
  231. turn: 1, step: 1,
  232. },
  233. })
  234. const oldChunk = { type: 'text-delta', index: 0, text: 'old' } as const
  235. ctx.emit('agent/assistant-stream', {
  236. agent,
  237. frame: {
  238. type: 'chunk', attemptId, revision: 2, index: 0,
  239. time: 101, chunk: oldChunk,
  240. },
  241. })
  242. const abort = new AbortController()
  243. const iterator = history.follow({
  244. address: { kind: 'session', sessionId: session.id },
  245. assistantStream: true,
  246. }, abort.signal)[Symbol.asyncIterator]()
  247. try {
  248. await expect(iterator.next()).resolves.toMatchObject({
  249. done: false,
  250. value: {
  251. type: 'snapshot',
  252. assistantStream: {
  253. revision: 2,
  254. activeAttempt: {
  255. attemptId,
  256. startedAfterSeq: -1,
  257. turn: 1,
  258. step: 1,
  259. nextIndex: 1,
  260. stream: [{ type: 'text-chunks', time0: 101, index: 0, dt: [], texts: ['old'] }],
  261. },
  262. },
  263. },
  264. })
  265. ctx.emit('agent/disposed', { agent })
  266. const replacementAgent = { id: session.id, session, status: 'running', ctx } as Agent
  267. const replacement: AssistantStreamFrame = {
  268. type: 'start', attemptId, revision: 1,
  269. turn: 2, step: 1,
  270. }
  271. ctx.emit('agent/assistant-stream', { agent: replacementAgent, frame: replacement })
  272. await expect(iterator.next()).resolves.toEqual({
  273. done: false,
  274. value: { type: 'assistant-stream', frame: { ...replacement, startedAfterSeq: -1 } },
  275. })
  276. } finally {
  277. await disposeFollow(ctx, iterator, abort)
  278. }
  279. })
  280. it('publishes an empty replacement baseline after an Agent frame revision gap', async () => {
  281. const { ctx } = await harness()
  282. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  283. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  284. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  285. const attemptId = LlmAttemptId('revision-gap-attempt')
  286. ctx.emit('agent/assistant-stream', {
  287. agent,
  288. frame: {
  289. type: 'start', attemptId, revision: 1,
  290. turn: 1, step: 1,
  291. },
  292. })
  293. const chunk = { type: 'text-delta', index: 0, text: 'after gap' } as const
  294. ctx.emit('agent/assistant-stream', {
  295. agent,
  296. frame: {
  297. type: 'chunk', attemptId, revision: 3, index: 0,
  298. time: 101, chunk,
  299. },
  300. })
  301. const abort = new AbortController()
  302. const iterator = history.follow({
  303. address: { kind: 'session', sessionId: session.id },
  304. assistantStream: true,
  305. }, abort.signal)[Symbol.asyncIterator]()
  306. try {
  307. await expect(iterator.next()).resolves.toMatchObject({
  308. done: false,
  309. value: {
  310. type: 'snapshot',
  311. assistantStream: { revision: 3 },
  312. },
  313. })
  314. } finally {
  315. await disposeFollow(ctx, iterator, abort)
  316. }
  317. })
  318. it('drops active attempts when an Agent chunk index is not dense', async () => {
  319. const { ctx } = await harness()
  320. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  321. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  322. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  323. const attemptId = LlmAttemptId('dense-index-attempt')
  324. ctx.emit('agent/assistant-stream', {
  325. agent,
  326. frame: {
  327. type: 'start', attemptId, revision: 1,
  328. turn: 1, step: 1,
  329. },
  330. })
  331. const chunk = { type: 'text-delta', index: 0, text: 'out of order' } as const
  332. ctx.emit('agent/assistant-stream', {
  333. agent,
  334. frame: {
  335. type: 'chunk', attemptId, revision: 2, index: 1,
  336. time: 101, chunk,
  337. },
  338. })
  339. const abort = new AbortController()
  340. const iterator = history.follow({
  341. address: { kind: 'session', sessionId: session.id },
  342. assistantStream: true,
  343. }, abort.signal)[Symbol.asyncIterator]()
  344. try {
  345. await expect(iterator.next()).resolves.toMatchObject({
  346. done: false,
  347. value: {
  348. type: 'snapshot',
  349. assistantStream: { revision: 2 },
  350. },
  351. })
  352. } finally {
  353. await disposeFollow(ctx, iterator, abort)
  354. }
  355. })
  356. it('reuses an unchanged Assistant baseline across follow openings', async () => {
  357. const { ctx } = await harness()
  358. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  359. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  360. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  361. const attemptId = LlmAttemptId('cached-baseline-attempt')
  362. ctx.emit('agent/assistant-stream', {
  363. agent,
  364. frame: {
  365. type: 'start', attemptId, revision: 1,
  366. turn: 1, step: 1,
  367. },
  368. })
  369. const firstAbort = new AbortController()
  370. const firstIterator = history.follow({
  371. address: { kind: 'session', sessionId: session.id },
  372. assistantStream: true,
  373. }, firstAbort.signal)[Symbol.asyncIterator]()
  374. const secondAbort = new AbortController()
  375. const secondIterator = history.follow({
  376. address: { kind: 'session', sessionId: session.id },
  377. assistantStream: true,
  378. }, secondAbort.signal)[Symbol.asyncIterator]()
  379. try {
  380. const first = await firstIterator.next()
  381. if (first.done || first.value.type !== 'snapshot') throw new Error('first follow did not open')
  382. const baseline = first.value.assistantStream
  383. expect(baseline).toMatchObject({ revision: 1, activeAttempt: { attemptId } })
  384. const second = await secondIterator.next()
  385. if (second.done || second.value.type !== 'snapshot') throw new Error('second follow did not open')
  386. expect(second.value.assistantStream).toEqual(baseline)
  387. } finally {
  388. firstAbort.abort()
  389. secondAbort.abort()
  390. await firstIterator.return?.()
  391. await secondIterator.return?.()
  392. await ctx.fiber.dispose()
  393. }
  394. })
  395. it('opens an empty Assistant baseline before the target Agent emits frames', async () => {
  396. const { ctx } = await harness()
  397. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  398. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  399. const abort = new AbortController()
  400. const iterator = history.follow({
  401. address: { kind: 'session', sessionId: session.id },
  402. assistantStream: true,
  403. }, abort.signal)[Symbol.asyncIterator]()
  404. try {
  405. await expect(iterator.next()).resolves.toMatchObject({
  406. done: false,
  407. value: {
  408. type: 'snapshot',
  409. assistantStream: { revision: 0 },
  410. },
  411. })
  412. } finally {
  413. await disposeFollow(ctx, iterator, abort)
  414. }
  415. })
  416. it('filters Assistant frames from another Session out of the target follow', async () => {
  417. const { ctx } = await harness()
  418. const target = ctx.sessions.create(undefined, { meta: { cwd: '/target' } })
  419. const other = ctx.sessions.create(undefined, { meta: { cwd: '/other' } })
  420. const otherAgent = { id: other.id, session: other, status: 'running', ctx } as Agent
  421. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  422. const abort = new AbortController()
  423. const iterator = history.follow({
  424. address: { kind: 'session', sessionId: target.id },
  425. assistantStream: true,
  426. }, abort.signal)[Symbol.asyncIterator]()
  427. try {
  428. await expect(iterator.next()).resolves.toMatchObject({
  429. done: false,
  430. value: { type: 'snapshot' },
  431. })
  432. ctx.emit('agent/assistant-stream', {
  433. agent: otherAgent,
  434. frame: {
  435. type: 'start', attemptId: LlmAttemptId('other-session-attempt'),
  436. revision: 1, turn: 1, step: 1,
  437. },
  438. })
  439. const targetEvent = target.append('turn/start', { turn: 1 })
  440. await expect(iterator.next()).resolves.toEqual({
  441. done: false,
  442. value: { type: 'event', event: targetEvent },
  443. })
  444. } finally {
  445. await disposeFollow(ctx, iterator, abort)
  446. }
  447. })
  448. it('does not replay a buffered Assistant frame already represented by the opening baseline', async () => {
  449. const { ctx } = await harness()
  450. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  451. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  452. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  453. const observationStarted = Promise.withResolvers<undefined>()
  454. const releaseObservation = Promise.withResolvers<undefined>()
  455. const originalObserve = ctx.sessionQuery.observeSession.bind(ctx.sessionQuery)
  456. const observe = vi.spyOn(ctx.sessionQuery, 'observeSession').mockImplementation(async (sessionId, options) => {
  457. observationStarted.resolve(undefined)
  458. await releaseObservation.promise
  459. return await originalObserve(sessionId, options)
  460. })
  461. const abort = new AbortController()
  462. const iterator = history.follow({
  463. address: { kind: 'session', sessionId: session.id },
  464. assistantStream: true,
  465. }, abort.signal)[Symbol.asyncIterator]()
  466. try {
  467. const opening = iterator.next()
  468. await observationStarted.promise
  469. const frame: AssistantStreamFrame = {
  470. type: 'start', attemptId: LlmAttemptId('opening-cut-attempt'),
  471. revision: 1, turn: 1, step: 1,
  472. }
  473. ctx.emit('agent/assistant-stream', { agent, frame })
  474. releaseObservation.resolve(undefined)
  475. await expect(opening).resolves.toMatchObject({
  476. done: false,
  477. value: {
  478. type: 'snapshot',
  479. assistantStream: { revision: 1, activeAttempt: { attemptId: frame.attemptId } },
  480. },
  481. })
  482. const durable = session.append('turn/start', { turn: 1 })
  483. await expect(iterator.next()).resolves.toEqual({
  484. done: false,
  485. value: { type: 'event', event: durable },
  486. })
  487. } finally {
  488. releaseObservation.resolve(undefined)
  489. observe.mockRestore()
  490. abort.abort()
  491. await iterator.return?.()
  492. await ctx.fiber.dispose()
  493. }
  494. })
  495. it('does not release an old-lifecycle frame after the opening baseline resets to revision one', async () => {
  496. const { ctx } = await harness()
  497. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  498. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  499. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  500. const attemptId = LlmAttemptId(`${session.id}:1`)
  501. ctx.emit('agent/assistant-stream', {
  502. agent,
  503. frame: {
  504. type: 'start', attemptId, revision: 1,
  505. turn: 1, step: 1,
  506. },
  507. })
  508. const observationStarted = Promise.withResolvers<undefined>()
  509. const releaseObservation = Promise.withResolvers<undefined>()
  510. const originalObserve = ctx.sessionQuery.observeSession.bind(ctx.sessionQuery)
  511. const observe = vi.spyOn(ctx.sessionQuery, 'observeSession').mockImplementation(async (sessionId, options) => {
  512. observationStarted.resolve(undefined)
  513. await releaseObservation.promise
  514. return await originalObserve(sessionId, options)
  515. })
  516. const abort = new AbortController()
  517. const iterator = history.follow({
  518. address: { kind: 'session', sessionId: session.id },
  519. assistantStream: true,
  520. }, abort.signal)[Symbol.asyncIterator]()
  521. try {
  522. const opening = iterator.next()
  523. await observationStarted.promise
  524. const oldChunk = { type: 'text-delta', index: 0, text: 'old lifecycle' } as const
  525. ctx.emit('agent/assistant-stream', {
  526. agent,
  527. frame: {
  528. type: 'chunk', attemptId, revision: 2, index: 0,
  529. time: 101, chunk: oldChunk,
  530. },
  531. })
  532. ctx.emit('agent/assistant-stream', {
  533. agent,
  534. frame: {
  535. type: 'start', attemptId, revision: 1,
  536. turn: 2, step: 1,
  537. },
  538. })
  539. releaseObservation.resolve(undefined)
  540. await expect(opening).resolves.toMatchObject({
  541. done: false,
  542. value: {
  543. type: 'snapshot',
  544. assistantStream: {
  545. revision: 1,
  546. activeAttempt: {
  547. attemptId, startedAfterSeq: -1,
  548. turn: 2, step: 1, nextIndex: 0, stream: [],
  549. },
  550. },
  551. },
  552. })
  553. const durable = session.append('turn/start', { turn: 2 })
  554. await expect(iterator.next()).resolves.toEqual({
  555. done: false,
  556. value: { type: 'event', event: durable },
  557. })
  558. } finally {
  559. releaseObservation.resolve(undefined)
  560. observe.mockRestore()
  561. abort.abort()
  562. await iterator.return?.()
  563. await ctx.fiber.dispose()
  564. }
  565. })
  566. it('keeps assistant frames out of a durable-only follower', async () => {
  567. const { ctx } = await harness()
  568. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  569. const agent = { id: session.id, session, status: 'running', ctx } as Agent
  570. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  571. const abort = new AbortController()
  572. const iterator = history.follow({
  573. address: { kind: 'session', sessionId: session.id },
  574. }, abort.signal)[Symbol.asyncIterator]()
  575. const opening = await iterator.next()
  576. expect(opening.value).not.toHaveProperty('assistantStream')
  577. const attemptId = LlmAttemptId('durable-only-attempt')
  578. ctx.emit('agent/assistant-stream', {
  579. agent,
  580. frame: {
  581. type: 'start', attemptId, revision: 1,
  582. turn: 1, step: 1,
  583. },
  584. })
  585. const durable = session.append('turn/start', { turn: 1 })
  586. ctx.emit('agent/assistant-stream', {
  587. agent,
  588. frame: {
  589. type: 'end', attemptId, revision: 2, index: 0, outcome: { kind: 'abandoned' },
  590. },
  591. })
  592. const next = session.append('turn/end', {
  593. turn: 1, reason: { kind: 'completed' },
  594. })
  595. await expect(iterator.next()).resolves.toEqual({
  596. done: false, value: { type: 'event', event: durable },
  597. })
  598. await expect(iterator.next()).resolves.toEqual({
  599. done: false, value: { type: 'event', event: next },
  600. })
  601. abort.abort()
  602. await iterator.next()
  603. await ctx.fiber.dispose()
  604. })
  605. it('follows raw tool events and preserves result metadata without a Tools service', async () => {
  606. const { ctx } = await harness()
  607. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  608. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  609. const abort = new AbortController()
  610. const stream = await openFollow(history, session.id, abort.signal)
  611. const collected = collect(stream, 2, abort)
  612. const call = session.append('tool/call', {
  613. turn: 1, step: 1, callId: ToolCallId('raw-call'), name: 'custom', arguments: '{malformed',
  614. })
  615. const result = session.append('tool/result', {
  616. turn: 1, step: 1,
  617. message: createToolResultMessage({
  618. callId: ToolCallId('raw-call'),
  619. content: [{ type: 'text', text: 'raw output' }],
  620. isError: false,
  621. }),
  622. meta: { nested: { count: 2 }, paths: ['a.ts', 'b.ts'] },
  623. }, { surfaceOp: 'append' })
  624. const frames = await collected
  625. expect(frames).toEqual([
  626. { type: 'event', event: call },
  627. { type: 'event', event: result },
  628. ])
  629. expect((frames[1] as Extract<SessionFollowFrame, { type: 'event' }>).event.data)
  630. .toMatchObject({ meta: { nested: { count: 2 }, paths: ['a.ts', 'b.ts'] } })
  631. })
  632. it('follows live results without rescanning Session history', async () => {
  633. const { ctx } = await harness()
  634. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  635. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  636. const abort = new AbortController()
  637. const stream = await openFollow(history, session.id, abort.signal)
  638. const iterator = stream[Symbol.asyncIterator]()
  639. session.append('tool/call', {
  640. turn: 1, step: 1, callId: ToolCallId('live-fast'), name: 'term', arguments: '{"cmd":"pwd"}',
  641. })
  642. await expect(iterator.next()).resolves.toMatchObject({
  643. value: { type: 'event', event: { type: 'tool/call', data: { callId: 'live-fast' } } },
  644. })
  645. const events = vi.spyOn(session, 'snapshotEvents').mockImplementation(() => {
  646. throw new Error('live result rescanned Session history')
  647. })
  648. try {
  649. session.append('tool/result', {
  650. turn: 1, step: 1,
  651. message: createToolResultMessage({
  652. callId: ToolCallId('live-fast'),
  653. content: [{ type: 'text', text: 'ok' }],
  654. isError: false,
  655. }),
  656. }, { surfaceOp: 'append' })
  657. await expect(iterator.next()).resolves.toMatchObject({
  658. value: { type: 'event', event: { type: 'tool/result', data: { message: { source: { callId: 'live-fast' } } } } },
  659. })
  660. } finally {
  661. events.mockRestore()
  662. abort.abort()
  663. await iterator.next()
  664. await ctx.fiber.dispose()
  665. }
  666. })
  667. it('serves raw call and result entries without parsing tool arguments', async () => {
  668. const { ctx } = await harness()
  669. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  670. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  671. const start = session.append('turn/start', { turn: 1 })
  672. const call = session.append('tool/call', {
  673. turn: 1, step: 1, callId: ToolCallId('history-call'), name: 'custom', arguments: '{broken',
  674. })
  675. const result = session.append('tool/result', {
  676. turn: 1, step: 1,
  677. message: createToolResultMessage({
  678. callId: ToolCallId('history-call'),
  679. content: [{ type: 'text', text: 'failed raw output' }],
  680. isError: true,
  681. }),
  682. meta: { persisted: true, count: 3 },
  683. }, { surfaceOp: 'append' })
  684. const response = await remote.page({
  685. address: { kind: 'session', sessionId: session.id },
  686. throughSeq: session.seq - 1,
  687. })
  688. expect(response.ok).toBe(true)
  689. if (!response.ok) throw new Error('unreachable')
  690. expect(response.value.records).toEqual([
  691. { type: 'event', event: start },
  692. { type: 'event', event: call },
  693. { type: 'event', event: result },
  694. ])
  695. })
  696. it('counts only append-origin messages toward maxMessages and keeps each compaction summary with its replacement', async () => {
  697. const { ctx } = await harness()
  698. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  699. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  700. session.append('turn/start', { turn: 1 })
  701. const first = appendUserText(session, 'first prompt')
  702. appendAssistantText(session, 'first reply', 1)
  703. const third = appendUserText(session, 'second prompt')
  704. appendAssistantText(session, 'second reply', 2)
  705. const shadowed = [...session.surface.nodes]
  706. const shadowedStart = shadowed[0]
  707. const shadowedEnd = shadowed.at(-1)
  708. if (shadowedStart === undefined || shadowedEnd === undefined) {
  709. throw new Error('expected a non-empty surface')
  710. }
  711. // A compaction transaction: a log-only summary record immediately followed by the
  712. // replacement that shadows the range.
  713. const summary = appendExtension(session, 'compaction/summary', {
  714. summary: [{ type: 'text', text: 'summary' }],
  715. shadowedRange: { start: shadowed[0], end: shadowed.at(-1) },
  716. shadowedSeqs: shadowed,
  717. shadowedTokenCount: 0,
  718. provider: 'p',
  719. model: 'm',
  720. })
  721. session.append('user/message', createUserMessage({
  722. content: [{ type: 'text', text: '<context_checkpoint>summary</context_checkpoint>' }],
  723. source: { kind: 'plugin', plugin: 'compact' },
  724. }), {
  725. surfaceOp: { op: 'replace', startSeq: shadowedStart, endSeq: shadowedEnd },
  726. sourceEventSeqs: [...shadowed, summary.seq],
  727. })
  728. const response = await remote.page({
  729. address: { kind: 'session', sessionId: session.id },
  730. throughSeq: session.seq - 1,
  731. maxMessages: 2,
  732. })
  733. if (!response.ok) throw new Error('unreachable')
  734. const page = pageEvents(response.value)
  735. // Two append-origin messages fill the page even though a replacement copy of
  736. // the same event type sits in the window: the copy is model-only.
  737. const messages = page.filter(event => event.type === 'user/message' || event.type === 'assistant/message')
  738. expect(messages.map(event => event.seq)).toEqual([third.seq, third.seq + 1, third.seq + 3])
  739. expect(page.some(event => event.seq === first.seq)).toBe(false)
  740. expect(response.value.hasMore).toBe(true)
  741. // The range stays contiguous, so the checkpoint's summary record is readable on
  742. // the same page as the checkpoint itself.
  743. const summaryIndex = page.findIndex(event => event.seq === summary.seq)
  744. expect(summaryIndex).toBeGreaterThan(-1)
  745. expect(page[summaryIndex + 1]?.seq).toBe(summary.seq + 1)
  746. expect(page.map(event => event.seq)).toEqual(page.map((_event, index) => third.seq + index))
  747. })
  748. it('paginates a message with a large embedded stream without expanding physical records', async () => {
  749. const { ctx } = await harness()
  750. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  751. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  752. session.append('turn/start', { turn: 1 })
  753. session.append('step/start', { turn: 1, step: 1 })
  754. const texts = Array.from({ length: 128 }, () => 'x')
  755. const message = session.append('assistant/message', {
  756. turn: 1,
  757. step: 1,
  758. message: createMessage({
  759. role: 'assistant',
  760. content: [{ type: 'text', text: 'x'.repeat(texts.length) }],
  761. source: { kind: 'model', provider: 'p', model: 'm' },
  762. }),
  763. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: texts.slice(1).map(() => 0), texts }],
  764. }, { surfaceOp: 'append' })
  765. const scalarMin = Math.min
  766. const min = vi.spyOn(Math, 'min').mockImplementation((...values) => {
  767. if (values.length > 2) throw new RangeError('variadic minimum rejected by regression harness')
  768. return scalarMin(...values)
  769. })
  770. try {
  771. const response = await remote.page({
  772. address: { kind: 'session', sessionId: session.id },
  773. throughSeq: message.seq,
  774. maxMessages: 1,
  775. })
  776. if (!response.ok) throw new Error('unreachable')
  777. expect(pageEvents(response.value).map(event => event.seq)).toEqual([message.seq])
  778. expect(response.value.records).toEqual([{ type: 'event', event: message }])
  779. expect(response.value.hasMore).toBe(true)
  780. } finally {
  781. min.mockRestore()
  782. }
  783. })
  784. it('keeps an earlier declared source on the same message-aligned page', async () => {
  785. const { ctx } = await harness()
  786. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  787. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  788. const source = session.append('request/context', { provider: 'p', model: 'm' })
  789. const laterSource = session.append('request/context', { provider: 'p', model: 'm' })
  790. const message = session.append('user/message', createUserMessage({
  791. content: [{ type: 'text', text: 'with source' }], source: { kind: 'user' },
  792. }), { sourceEventSeqs: [source.seq, laterSource.seq], surfaceOp: 'append' })
  793. const response = await remote.page({
  794. address: { kind: 'session', sessionId: session.id },
  795. throughSeq: message.seq,
  796. maxMessages: 1,
  797. })
  798. if (!response.ok) throw new Error('unreachable')
  799. expect(pageEvents(response.value).map(event => event.seq)).toEqual([source.seq, laterSource.seq, message.seq])
  800. expect(response.value.hasMore).toBe(false)
  801. await ctx.fiber.dispose()
  802. })
  803. it('keeps compact reasoning and tool-call runs nested in one attempt event', async () => {
  804. const { ctx } = await harness()
  805. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  806. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  807. const callId = ToolCallId('packed-call')
  808. const attempt = session.append('assistant/attempt', {
  809. turn: 1,
  810. step: 1,
  811. stream: [
  812. { type: 'reasoning-chunks', time0: 1, index: 0, dt: [1, 1], texts: ['r0', 'r1', 'r2'] },
  813. { type: 'tool-call-chunks', time0: 4, index: 1, id: callId, dt: [1, 1], args: ['a0', 'a1', 'a2'] },
  814. ],
  815. })
  816. const response = await remote.page({
  817. address: { kind: 'session', sessionId: session.id },
  818. throughSeq: session.seq - 1,
  819. })
  820. if (!response.ok) throw new Error('unreachable')
  821. expect(response.value.records).toEqual([{ type: 'event', event: attempt }])
  822. await ctx.fiber.dispose()
  823. })
  824. it('follows a result after turn/end without reading the addressed Session log', async () => {
  825. const { ctx } = await harness()
  826. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  827. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  828. const abort = new AbortController()
  829. const stream = await openFollow(history, session.id, abort.signal)
  830. const iterator = stream[Symbol.asyncIterator]()
  831. session.append('turn/start', { turn: 1 })
  832. await expect(iterator.next()).resolves.toMatchObject({
  833. value: { type: 'event', event: { type: 'turn/start' } },
  834. })
  835. session.append('tool/call', { turn: 1, step: 1, callId: ToolCallId('c-late'), name: 'term', arguments: '{"cmd":"tail"}' })
  836. await expect(iterator.next()).resolves.toMatchObject({
  837. value: { type: 'event', event: { type: 'tool/call' } },
  838. })
  839. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  840. await expect(iterator.next()).resolves.toMatchObject({
  841. value: { type: 'event', event: { type: 'turn/end' } },
  842. })
  843. const events = vi.spyOn(session, 'snapshotEvents').mockImplementation(() => {
  844. throw new Error('live result rescanned Session history')
  845. })
  846. try {
  847. const result = session.append('tool/result', {
  848. turn: 1, step: 1,
  849. message: createToolResultMessage({
  850. callId: ToolCallId('c-late'),
  851. content: [{ type: 'text', text: 'ok' }],
  852. isError: false,
  853. }),
  854. }, { surfaceOp: 'append' })
  855. await expect(iterator.next()).resolves.toEqual({
  856. done: false,
  857. value: { type: 'event', event: result },
  858. })
  859. } finally {
  860. events.mockRestore()
  861. abort.abort()
  862. await iterator.next()
  863. await ctx.fiber.dispose()
  864. }
  865. })
  866. })