transport.client.spec.ts 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717
  1. import { describe, expect, it, vi } from 'vitest'
  2. import {
  3. RemoteStream,
  4. RemoteStreamCarrierError,
  5. type RemoteStreamOptions,
  6. } from '@deepseek-ai/dsh-api-gateway/client'
  7. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  8. import { LlmAttemptId } from '@deepseek-ai/dsh-llm'
  9. import { SESSION_FORMAT_VERSION } from '@deepseek-ai/dsh-session/types'
  10. import type { RemoteResult } from '@deepseek-ai/dsh-typert-protocol'
  11. import {
  12. createSessionControlStream,
  13. SessionEventStream,
  14. type SessionJournalChange,
  15. type SessionRemote,
  16. } from '../src/client/index.ts'
  17. import type { SessionRemotes } from '../src/client/sessions/remotes.ts'
  18. import type {
  19. SessionAddress,
  20. SessionAssistantStreamBaseline,
  21. SessionAssistantStreamFrame,
  22. SessionControlFrame,
  23. SessionEventEntry,
  24. SessionFollowFrame,
  25. SessionFollowRequest,
  26. SessionHistoryRecord,
  27. SessionPage,
  28. SessionPageRequest,
  29. SessionWireEvent,
  30. } from '../src/types.ts'
  31. type SessionTransportRemote = Pick<SessionRemote, 'control' | 'follow' | 'page'>
  32. const ADDRESS: SessionAddress = { kind: 'session', sessionId: 'session-1' as never }
  33. const AVAILABLE_CONNECTION = {
  34. generation: {
  35. getSnapshot: () => ({ id: 1, host: { home: '/home/fixture' } }),
  36. subscribe: () => () => {},
  37. },
  38. }
  39. function entry(seq: number): SessionEventEntry {
  40. return { type: 'event', event: { type: 'turn/start', seq, time: seq, data: { turn: seq } } }
  41. }
  42. function page(records: readonly SessionHistoryRecord[], hasMore = false): SessionPage {
  43. return { records, hasMore }
  44. }
  45. function snapshot(
  46. cursor: number,
  47. records: readonly SessionHistoryRecord[],
  48. hasMore = false,
  49. assistantStream: SessionAssistantStreamBaseline = { revision: 0 },
  50. ): SessionFollowFrame {
  51. return {
  52. type: 'snapshot',
  53. header: {
  54. version: SESSION_FORMAT_VERSION,
  55. id: ADDRESS.kind === 'session' ? ADDRESS.sessionId : ADDRESS.childSessionId,
  56. createdAt: 0,
  57. isSeeded: false,
  58. },
  59. cursor,
  60. records,
  61. hasMore,
  62. projections: { asOfSeq: cursor, values: {} },
  63. assistantStream,
  64. }
  65. }
  66. function assistantFrame(frame: SessionAssistantStreamFrame): SessionFollowFrame {
  67. return { type: 'assistant-stream', frame }
  68. }
  69. function sessionClient(remote: SessionTransportRemote): SessionRemotes {
  70. return {
  71. session: remote as SessionRemote,
  72. $stream: <Item>(options: RemoteStreamOptions<Item>) => (
  73. new RemoteStream(AVAILABLE_CONNECTION, options)
  74. ),
  75. commands: { execute: () => Promise.reject(new Error('stream tests never run commands')) },
  76. subagents: {
  77. list: () => Promise.reject(new Error('stream tests never read the subagent catalog')),
  78. prompt: () => Promise.reject(new Error('stream tests never prompt a subagent')),
  79. interruptByParent: () => Promise.reject(new Error('stream tests never interrupt a subagent')),
  80. },
  81. }
  82. }
  83. interface FollowGeneration {
  84. readonly frames: readonly SessionFollowFrame[]
  85. readonly terminal?: Error
  86. readonly hold?: boolean
  87. readonly waitAfterFrames?: Promise<void>
  88. }
  89. class ScriptedSessionRemote implements SessionTransportRemote {
  90. readonly followRequests: SessionFollowRequest[] = []
  91. readonly pageRequests: SessionPageRequest[] = []
  92. readonly signals: AbortSignal[] = []
  93. constructor(
  94. private readonly generations: FollowGeneration[],
  95. private readonly pages: RemoteResult<SessionPage>[],
  96. private readonly controlFrames: readonly SessionControlFrame[] = [],
  97. private readonly holdControl = true,
  98. ) {}
  99. async *follow(request: SessionFollowRequest, signal = new AbortController().signal): AsyncIterable<SessionFollowFrame> {
  100. const generation = this.generations.shift()
  101. if (generation === undefined) throw new Error('no scripted Session generation')
  102. this.followRequests.push(request)
  103. this.signals.push(signal)
  104. for (const frame of generation.frames) yield frame
  105. await generation.waitAfterFrames
  106. if (generation.terminal !== undefined) throw generation.terminal
  107. if (generation.hold === true && !signal.aborted) {
  108. await new Promise<void>((resolve) => {
  109. signal.addEventListener('abort', () => { resolve() }, { once: true })
  110. })
  111. }
  112. }
  113. page(request: SessionPageRequest): Promise<RemoteResult<SessionPage>> {
  114. this.pageRequests.push(request)
  115. const result = this.pages.shift()
  116. if (result === undefined) throw new Error('no scripted Session page')
  117. return Promise.resolve(result)
  118. }
  119. async *control(signal = new AbortController().signal): AsyncIterable<SessionControlFrame> {
  120. for (const frame of this.controlFrames) yield frame
  121. if (this.holdControl && !signal.aborted) {
  122. await new Promise<void>((resolve) => {
  123. signal.addEventListener('abort', () => { resolve() }, { once: true })
  124. })
  125. }
  126. }
  127. }
  128. function surfaceEvent(type = 'user/message'): SessionWireEvent {
  129. return { type, seq: 10, time: 10, data: {}, surfaceOp: 'append' }
  130. }
  131. const invalidWireEvents: [string, unknown][] = [
  132. ['null event', null],
  133. ['array event', []],
  134. ['extra envelope key', { ...surfaceEvent(), obsolete: true }],
  135. ['unknown ignorable extra envelope key', { type: 'extension/event', seq: 10, time: 10, data: {}, ignorable: true, obsolete: true }],
  136. ['missing data', { type: 'turn/start', seq: 10, time: 10 }],
  137. ['invalid type', { ...surfaceEvent(), type: null }],
  138. ['fractional sequence', { ...surfaceEvent(), seq: 1.5 }],
  139. ['negative sequence', { ...surfaceEvent(), seq: -1 }],
  140. ['negative zero sequence', { ...surfaceEvent(), seq: -0 }],
  141. ['unsafe sequence', { ...surfaceEvent(), seq: Number.MAX_SAFE_INTEGER + 1 }],
  142. ['fractional time', { ...surfaceEvent(), time: 0.5 }],
  143. ['invalid ignorable marker', { ...surfaceEvent(), ignorable: false }],
  144. ...['system/message', 'user/message', 'assistant/message', 'tool/result'].map((type): [string, unknown] => [
  145. `missing ${type} surface marker`,
  146. { type, seq: 10, time: 10, data: {} },
  147. ]),
  148. ...['turn/start', 'assistant/attempt', 'request/header', 'request/context', 'session/title', 'extension/event'].flatMap(type => [
  149. [`non-surface ${type} operation`, { ...surfaceEvent(type) }],
  150. [`non-surface ${type} sources`, { type, seq: 10, time: 10, data: {}, sourceEventSeqs: [0] }],
  151. ] as [string, unknown][]),
  152. ...['turn/start', 'assistant/attempt', 'request/context', 'tool/ptc-dispatch', 'session/title'].flatMap(type => [
  153. [`known ignorable ${type} operation`, { type, seq: 10, time: 10, data: {}, ignorable: true, surfaceOp: { opaque: true } }],
  154. [`known ignorable ${type} sources`, { type, seq: 10, time: 10, data: {}, ignorable: true, sourceEventSeqs: { opaque: true } }],
  155. ] as [string, unknown][]),
  156. ['assistant sources', { ...surfaceEvent('assistant/message'), sourceEventSeqs: [0] }],
  157. ...[[], [0, 0], [-1], [-0], [0.5], [10], [11], [Number.MAX_SAFE_INTEGER + 1]].map(
  158. (sourceEventSeqs): [string, unknown] => [`invalid sources ${JSON.stringify(sourceEventSeqs)}`, { ...surfaceEvent(), sourceEventSeqs }],
  159. ),
  160. ...[
  161. null, {}, 'replace',
  162. { op: 'replace', start: 0, end: 1 },
  163. { op: 'replace', startSeq: 0, end: 1 },
  164. { op: 'replace', start: 0, endSeq: 1 },
  165. { op: 'replace', startSeq: 0, endSeq: 1, start: 0, end: 1 },
  166. { op: 'replace', startSeq: 0, endSeq: 1, extra: true },
  167. { op: 'replace', startSeq: 0 },
  168. { op: 'replace', endSeq: 1 },
  169. ...[-1, -0, 0.5, 10, Number.MAX_SAFE_INTEGER + 1].flatMap(seq => [
  170. { op: 'replace', startSeq: seq, endSeq: 1 },
  171. { op: 'replace', startSeq: 0, endSeq: seq },
  172. ]),
  173. ].map((surfaceOp, index): [string, unknown] => [`invalid replacement ${index}`, { ...surfaceEvent(), surfaceOp }]),
  174. ...[{ system: '' }, { system: ' ' }, { system: 'prompt' }, { system: null }, { system: {} }, { tools: [] }, { adapterDefaults: {} }].map((optional): [string, unknown] => [
  175. `empty request header ${JSON.stringify(optional)}`,
  176. { type: 'request/header', seq: 10, time: 10, data: {
  177. reason: 'initial', header: { config: { provider: 'mock', model: 'mock' }, ...optional },
  178. } },
  179. ]),
  180. ...[false, undefined].map((isError): [string, unknown] => [
  181. `contradictory tool error ${String(isError)}`,
  182. { ...surfaceEvent('tool/result'), data: {
  183. message: { content: [{ type: 'tool-result', content: [], ...(isError === undefined ? {} : { isError }) }] },
  184. error: { name: 'Error', code: 'FAILURE' },
  185. } },
  186. ]),
  187. ]
  188. describe.each(['snapshot', 'live', 'page'] as const)('Session %s wire acceptance', (path) => {
  189. it.each(invalidWireEvents)('refuses %s without publishing or retrying', async (_name, event) => {
  190. // The Remote mock is the decoded JSON transport, not a typed same-process producer.
  191. const record = { type: 'event', event } as SessionHistoryRecord
  192. const opening = snapshot(path === 'live' ? 9 : 11, [entry(path === 'live' ? 9 : 11)])
  193. const remote = new ScriptedSessionRemote([{
  194. frames: path === 'snapshot' ? [snapshot(10, [record])]
  195. : path === 'live' ? [opening, record] : [opening],
  196. hold: true,
  197. }], path === 'page' ? [{ ok: true, value: page([record]) }] : [])
  198. const publish = vi.fn()
  199. const failed = vi.fn()
  200. const carrierFailed = vi.fn()
  201. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, { publish, failed, carrierFailed })
  202. try {
  203. if (path === 'snapshot') {
  204. await expect(stream.open({})).rejects.toThrow()
  205. expect(publish).not.toHaveBeenCalled()
  206. } else {
  207. await stream.open({})
  208. if (path === 'page') await expect(stream.prepend({})).rejects.toThrow()
  209. else await vi.waitFor(() => { expect(failed).toHaveBeenCalledOnce() })
  210. expect(publish).toHaveBeenCalledOnce()
  211. }
  212. expect(remote.followRequests).toHaveLength(1)
  213. expect(carrierFailed).not.toHaveBeenCalled()
  214. } finally {
  215. await stream.dispose()
  216. }
  217. })
  218. })
  219. describe('Session Client stream adapters', () => {
  220. it('preserves current envelopes and payloads without normalization across every journal path', async () => {
  221. const events: SessionWireEvent[] = [
  222. surfaceEvent(),
  223. surfaceEvent('system/message'),
  224. { ...surfaceEvent('system/message'), surfaceOp: { op: 'replace', startSeq: 2, endSeq: 2 }, sourceEventSeqs: [2], data: { message: { source: { plugin: 'system', extra: true }, content: [] }, extra: { retained: true } } },
  225. { ...surfaceEvent(), sourceEventSeqs: [0, 2] },
  226. { ...surfaceEvent(), surfaceOp: { op: 'replace', startSeq: 2, endSeq: 0 }, sourceEventSeqs: [2, 0] },
  227. { ...surfaceEvent('assistant/message'), data: { turn: 1, step: 1, message: {}, stream: [] } },
  228. { ...surfaceEvent('tool/result'), data: {
  229. message: { content: [{ type: 'tool-result', content: [], isError: true }] },
  230. error: { name: 'Error', code: 'FAILURE' }, meta: { extension: ['retained'] },
  231. } },
  232. { ...surfaceEvent('tool/result'), sourceEventSeqs: [0], data: {
  233. message: { content: [{ type: 'tool-result', content: [], isError: true }] },
  234. } },
  235. { type: 'request/header', seq: 10, time: 10, data: {
  236. reason: 'initial', header: { config: { provider: 'mock', model: 'mock' } },
  237. } },
  238. { type: 'request/header', seq: 10, time: 10, data: {
  239. reason: 'change', header: {
  240. config: { provider: 'mock', model: 'mock' }, extension: { nested: ['retained'] },
  241. tools: [{ name: 'fixture' }], adapterDefaults: { temperature: 1 },
  242. },
  243. } },
  244. ...['extension/event', 'tool/code-dispatch', 'tool/code-dispatch-start'].map(type => ({
  245. type, seq: 10, time: 10, data: { nested: [null, true] }, ignorable: true,
  246. surfaceOp: { opaque: ['retained'] }, sourceEventSeqs: { opaque: [null] },
  247. }) as unknown as SessionWireEvent),
  248. ]
  249. for (const event of events) {
  250. const before = structuredClone(event)
  251. const record: SessionHistoryRecord = { type: 'event', event }
  252. const remote = new ScriptedSessionRemote([{
  253. frames: [snapshot(10, [record]), { type: 'event', event: { ...event, seq: 11 } }], hold: true,
  254. }], [{ ok: true, value: page([{ type: 'event', event: { ...event, seq: 9 } }]) }])
  255. const changes: SessionJournalChange[] = []
  256. let appended!: () => void
  257. const ready = new Promise<void>((resolve) => { appended = resolve })
  258. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  259. publish: (change) => { changes.push(change); if (change.type === 'append') appended() },
  260. failed: vi.fn(),
  261. })
  262. try {
  263. await stream.open({})
  264. await ready
  265. await stream.prepend({})
  266. expect(changes.map(change => change.type)).toEqual(['replace', 'append', 'prepend'])
  267. expect(changes[0]).toMatchObject({ page: { records: [record] } })
  268. expect(event).toEqual(before)
  269. } finally {
  270. await stream.dispose()
  271. }
  272. }
  273. })
  274. it('opts into assistant notifications and publishes the reconnect baseline plus live frame', async () => {
  275. const attemptId = LlmAttemptId('transport-attempt')
  276. const baseline: SessionAssistantStreamBaseline = {
  277. revision: 2,
  278. activeAttempt: {
  279. attemptId,
  280. startedAfterSeq: -1,
  281. turn: 1,
  282. step: 1,
  283. nextIndex: 1,
  284. stream: [{ type: 'text-chunks', time0: 0, index: 0, dt: [], texts: ['a'] }],
  285. },
  286. }
  287. const frame: SessionAssistantStreamFrame = {
  288. type: 'chunk', attemptId, revision: 3, index: 1,
  289. time: 1, chunk: { type: 'text-delta', index: 0, text: 'b' },
  290. }
  291. const remote = new ScriptedSessionRemote(
  292. [{ frames: [snapshot(0, [entry(0)], false, baseline), assistantFrame(frame)], hold: true }],
  293. [],
  294. )
  295. const changes: SessionJournalChange[] = []
  296. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  297. publish: (change) => { changes.push(change) },
  298. failed: vi.fn(),
  299. })
  300. await stream.open({})
  301. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  302. expect(remote.followRequests).toEqual([{ address: ADDRESS, assistantStream: true }])
  303. expect(changes).toMatchObject([
  304. { type: 'replace', page: { assistantStream: baseline } },
  305. { type: 'assistant-stream', frame },
  306. ])
  307. await stream.dispose()
  308. })
  309. it('rejects an opted-in opening that omits its Assistant baseline', async () => {
  310. const remote = new ScriptedSessionRemote([{
  311. frames: [{
  312. type: 'snapshot',
  313. header: {
  314. version: SESSION_FORMAT_VERSION,
  315. id: ADDRESS.sessionId,
  316. createdAt: 0,
  317. isSeeded: false,
  318. },
  319. cursor: -1,
  320. records: [],
  321. hasMore: false,
  322. projections: { asOfSeq: -1, values: {} },
  323. }],
  324. }], [])
  325. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  326. publish: vi.fn(),
  327. failed: vi.fn(),
  328. })
  329. try {
  330. await expect(stream.open({})).rejects.toMatchObject({
  331. code: 'gateway/internal',
  332. message: 'session assistant stream omitted its opted-in opening baseline',
  333. })
  334. } finally {
  335. await stream.dispose()
  336. }
  337. })
  338. it('rejects an Assistant frame that arrives before the opening baseline', async () => {
  339. const remote = new ScriptedSessionRemote([{
  340. frames: [assistantFrame({
  341. type: 'start', attemptId: LlmAttemptId('pre-opening-attempt'),
  342. revision: 1, startedAfterSeq: -1, turn: 1, step: 1,
  343. })],
  344. }], [])
  345. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  346. publish: vi.fn(),
  347. failed: vi.fn(),
  348. })
  349. try {
  350. await expect(stream.open({})).rejects.toMatchObject({
  351. code: 'gateway/internal',
  352. message: 'session event stream emitted an entry before its opening cursor',
  353. })
  354. expect(remote.followRequests).toEqual([{ address: ADDRESS, assistantStream: true }])
  355. } finally {
  356. await stream.dispose()
  357. }
  358. })
  359. it('rebaselines after a transient assistant revision gap without advancing the durable cursor', async () => {
  360. const attemptId = LlmAttemptId('gapped-attempt')
  361. const start: SessionAssistantStreamFrame = {
  362. type: 'start', attemptId, revision: 1, startedAfterSeq: -1,
  363. turn: 1, step: 1,
  364. }
  365. const gap: SessionAssistantStreamFrame = {
  366. type: 'chunk', attemptId, revision: 3, index: 0,
  367. time: 1, chunk: { type: 'text-delta', index: 0, text: 'lost predecessor' },
  368. }
  369. const replacement: SessionAssistantStreamBaseline = {
  370. revision: 3,
  371. activeAttempt: {
  372. attemptId,
  373. startedAfterSeq: -1,
  374. turn: 1,
  375. step: 1,
  376. nextIndex: 1,
  377. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['lost predecessor'] }],
  378. },
  379. }
  380. const remote = new ScriptedSessionRemote([
  381. {
  382. frames: [snapshot(0, [entry(0)]), assistantFrame(start), assistantFrame(gap)],
  383. },
  384. { frames: [snapshot(0, [entry(0)], false, replacement)], hold: true },
  385. ], [])
  386. const changes: SessionJournalChange[] = []
  387. const carrierFailed = vi.fn()
  388. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  389. publish: (change) => { changes.push(change) },
  390. carrierFailed,
  391. failed: vi.fn(),
  392. })
  393. await stream.open({})
  394. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  395. expect(changes.map(change => change.type)).toEqual([
  396. 'replace', 'assistant-stream', 'replace',
  397. ])
  398. expect(changes.at(-1)).toMatchObject({
  399. type: 'replace', page: { assistantStream: replacement },
  400. })
  401. expect(remote.pageRequests).toEqual([])
  402. expect(carrierFailed).toHaveBeenCalledWith(expect.objectContaining({
  403. message: 'session assistant stream skipped revision 2',
  404. }))
  405. await stream.dispose()
  406. })
  407. it('rebaselines when a replacement Agent lifecycle restarts at revision one', async () => {
  408. const attemptId = LlmAttemptId('replacement-lifecycle-attempt')
  409. const previous: SessionAssistantStreamBaseline = {
  410. revision: 2,
  411. activeAttempt: {
  412. attemptId,
  413. startedAfterSeq: -1,
  414. turn: 1,
  415. step: 1,
  416. nextIndex: 1,
  417. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['old'] }],
  418. },
  419. }
  420. const replacementStart: SessionAssistantStreamFrame = {
  421. type: 'start', attemptId, revision: 1, startedAfterSeq: -1,
  422. turn: 2, step: 1,
  423. }
  424. const replacement: SessionAssistantStreamBaseline = {
  425. revision: 1,
  426. activeAttempt: {
  427. attemptId,
  428. startedAfterSeq: -1,
  429. turn: 2,
  430. step: 1,
  431. nextIndex: 0,
  432. stream: [],
  433. },
  434. }
  435. const remote = new ScriptedSessionRemote([
  436. {
  437. frames: [snapshot(0, [entry(0)], false, previous), assistantFrame(replacementStart)],
  438. },
  439. { frames: [snapshot(0, [entry(0)], false, replacement)], hold: true },
  440. ], [])
  441. const changes: SessionJournalChange[] = []
  442. const carrierFailed = vi.fn()
  443. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  444. publish: (change) => { changes.push(change) },
  445. carrierFailed,
  446. failed: vi.fn(),
  447. })
  448. try {
  449. await stream.open({})
  450. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  451. expect(changes).toMatchObject([
  452. { type: 'replace', page: { assistantStream: previous } },
  453. { type: 'replace', page: { assistantStream: replacement } },
  454. ])
  455. expect(carrierFailed).toHaveBeenCalledWith(expect.objectContaining({
  456. message: 'session assistant stream skipped revision 3',
  457. }))
  458. } finally {
  459. await stream.dispose()
  460. }
  461. })
  462. it('validates one scalar current-event range before publishing Client entries', async () => {
  463. const remote = new ScriptedSessionRemote(
  464. [{ frames: [snapshot(2, [entry(0), entry(1), entry(2)]), entry(3)], hold: true }],
  465. [],
  466. )
  467. const changes: SessionJournalChange[] = []
  468. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  469. publish: (change) => { changes.push(change) },
  470. failed: vi.fn(),
  471. })
  472. await stream.open({})
  473. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  474. expect(changes[0]).toMatchObject({
  475. type: 'replace',
  476. entries: [
  477. entry(0),
  478. entry(1),
  479. entry(2),
  480. ],
  481. })
  482. expect(changes[1]).toEqual({ type: 'append', entry: entry(3) })
  483. await stream.dispose()
  484. })
  485. it('binds an event journal to one address and publishes replace, append, and prepend changes', async () => {
  486. const remote = new ScriptedSessionRemote(
  487. [{
  488. frames: [
  489. snapshot(3, [entry(2), entry(3)], true),
  490. entry(3),
  491. entry(4),
  492. ],
  493. hold: true,
  494. }],
  495. [
  496. { ok: true, value: page([entry(0), entry(1)], false) },
  497. ],
  498. )
  499. const changes: SessionJournalChange[] = []
  500. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  501. publish: (change) => { changes.push(change) },
  502. failed: vi.fn(),
  503. })
  504. await stream.open({ maxMessages: 50 })
  505. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  506. await stream.prepend({ beforeSeq: 2, maxMessages: 50 })
  507. expect(remote.followRequests).toEqual([{
  508. address: ADDRESS, assistantStream: true, maxMessages: 50,
  509. }])
  510. expect(remote.pageRequests).toEqual([
  511. { address: ADDRESS, throughSeq: 4, beforeSeq: 2, maxMessages: 50 },
  512. ])
  513. expect(changes).toMatchObject([
  514. { type: 'replace', entries: [entry(2), entry(3)], hasMore: true },
  515. { type: 'append', entry: entry(4) },
  516. { type: 'prepend', entries: [entry(0), entry(1)], hasMore: false },
  517. ])
  518. await stream.dispose()
  519. expect(remote.signals[0]?.aborted).toBe(true)
  520. })
  521. it('replaces the retained window from each reconnect snapshot', async () => {
  522. const lost = new RemoteStreamCarrierError('lost')
  523. const remote = new ScriptedSessionRemote(
  524. [
  525. {
  526. frames: [snapshot(1, [entry(0), entry(1)]), entry(2)],
  527. terminal: lost,
  528. },
  529. { frames: [snapshot(4, [entry(0), entry(1), entry(2), entry(3), entry(4)])], hold: true },
  530. ],
  531. [],
  532. )
  533. const changes: SessionJournalChange[] = []
  534. const carrierFailed = vi.fn()
  535. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  536. publish: (change) => { changes.push(change) },
  537. carrierFailed,
  538. failed: vi.fn(),
  539. })
  540. await stream.open({ maxMessages: 50 })
  541. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  542. expect(remote.followRequests).toEqual([
  543. { address: ADDRESS, assistantStream: true, maxMessages: 50 },
  544. { address: ADDRESS, assistantStream: true, maxMessages: 50 },
  545. ])
  546. expect(remote.pageRequests).toEqual([])
  547. expect(changes.map(change => change.type)).toEqual(['replace', 'append', 'replace'])
  548. expect(carrierFailed).toHaveBeenCalledWith(lost)
  549. await stream.dispose()
  550. })
  551. it('repairs a resumed event stream without an optional message limit', async () => {
  552. const finish = Promise.withResolvers<undefined>()
  553. const remote = new ScriptedSessionRemote(
  554. [
  555. {
  556. frames: [snapshot(0, [entry(0)])],
  557. waitAfterFrames: finish.promise,
  558. terminal: new RemoteStreamCarrierError('lost'),
  559. },
  560. { frames: [snapshot(1, [entry(0), entry(1)])], hold: true },
  561. ],
  562. [],
  563. )
  564. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  565. publish: vi.fn(),
  566. failed: vi.fn(),
  567. })
  568. await stream.open({})
  569. finish.resolve(undefined)
  570. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  571. expect(remote.followRequests).toEqual([
  572. { address: ADDRESS, assistantStream: true },
  573. { address: ADDRESS, assistantStream: true },
  574. ])
  575. expect(remote.pageRequests).toEqual([])
  576. await stream.dispose()
  577. })
  578. it.each([{}, { maxMessages: 50 }])('repairs a live gap preserving message limit %j', async (request) => {
  579. const remote = new ScriptedSessionRemote(
  580. [{ frames: [snapshot(0, [entry(0)]), entry(2)], hold: true }],
  581. [{ ok: true, value: page([entry(0), entry(1), entry(2)]) }],
  582. )
  583. const changes: SessionJournalChange[] = []
  584. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  585. publish: (change) => { changes.push(change) },
  586. failed: vi.fn(),
  587. })
  588. try {
  589. await stream.open(request)
  590. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  591. expect(remote.pageRequests).toEqual([{ address: ADDRESS, throughSeq: 2, ...request }])
  592. } finally {
  593. await stream.dispose()
  594. }
  595. })
  596. it('turns a pagination failure into a typed stream failure', async () => {
  597. const failure = new RemoteError('session/not-found', 'missing', { sessionId: 'session-1' as never })
  598. const remote = new ScriptedSessionRemote(
  599. [{ frames: [snapshot(-1, [])], hold: true }],
  600. [{ ok: false, error: failure }],
  601. )
  602. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  603. publish: vi.fn(),
  604. failed: vi.fn(),
  605. })
  606. await stream.open({})
  607. await expect(stream.prepend({})).rejects.toMatchObject({ code: 'session/not-found' })
  608. await expect(stream.open({})).rejects.toThrow('already opened')
  609. expect(remote.signals[0]?.aborted).toBe(false)
  610. expect(remote.pageRequests).toEqual([{ address: ADDRESS, throughSeq: -1 }])
  611. await stream.dispose()
  612. expect(remote.signals[0]?.aborted).toBe(true)
  613. })
  614. it('maps the Host-wide control baseline and deltas into one snapshot stream', async () => {
  615. const baseline: SessionControlFrame = {
  616. type: 'baseline',
  617. value: { queues: {}, jobs: {}, projections: {} },
  618. }
  619. const update: SessionControlFrame = {
  620. type: 'queue', sessionId: 'session-1' as never, items: [],
  621. }
  622. const remote = new ScriptedSessionRemote([], [], [baseline, update])
  623. const accept = vi.fn<(frame: SessionControlFrame) => void>()
  624. const stream = createSessionControlStream(sessionClient(remote), {
  625. accept,
  626. failed: vi.fn(),
  627. })
  628. stream.start()
  629. stream.start()
  630. await vi.waitFor(() => { expect(accept).toHaveBeenCalledTimes(2) })
  631. expect(accept.mock.calls.map(([frame]) => frame)).toEqual([baseline, update])
  632. await stream.dispose()
  633. await stream.dispose()
  634. })
  635. it('classifies control streams that end before and after their opening baseline', async () => {
  636. const beforeFailed = vi.fn()
  637. const before = createSessionControlStream(
  638. sessionClient(new ScriptedSessionRemote([], [], [], false)),
  639. { accept: vi.fn(), failed: beforeFailed },
  640. )
  641. before.start()
  642. await vi.waitFor(() => { expect(beforeFailed).toHaveBeenCalledOnce() })
  643. expect(beforeFailed.mock.calls[0]?.[0]).toMatchObject({
  644. message: 'session control stream ended before its opening snapshot',
  645. })
  646. await before.dispose()
  647. const baseline: SessionControlFrame = {
  648. type: 'baseline',
  649. value: { queues: {}, jobs: {}, projections: {} },
  650. }
  651. const carrierFailed = vi.fn()
  652. const failed = vi.fn()
  653. const afterRemote = new ScriptedSessionRemote([], [], [baseline], false)
  654. const after = createSessionControlStream(sessionClient(afterRemote), {
  655. accept: vi.fn(),
  656. carrierFailed: (error) => {
  657. carrierFailed(error)
  658. void after.dispose()
  659. },
  660. failed,
  661. })
  662. after.start()
  663. await vi.waitFor(() => { expect(carrierFailed).toHaveBeenCalledOnce() })
  664. expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
  665. message: 'session control stream ended without a terminal result',
  666. })
  667. expect(failed).not.toHaveBeenCalled()
  668. await after.dispose()
  669. })
  670. })