invariant.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344
  1. import { describe, expect, it } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { createScope, scopeTarget } from '@deepseek-ai/dsh-scope'
  4. import { CallId } from '@deepseek-ai/dsh-llm'
  5. import SessionStore, { SessionId, TOOL_NOT_STARTED } from '@deepseek-ai/dsh-session'
  6. import * as SessionInvariant from '@deepseek-ai/dsh-session/invariant'
  7. import InvariantService, { InvariantError } from '@deepseek-ai/dsh-invariants'
  8. async function setup(): Promise<{ ctx: Context; fiber: Awaited<ReturnType<Context['plugin']>> }> {
  9. const ctx = new Context()
  10. await ctx.plugin(SessionStore)
  11. await ctx.plugin(InvariantService)
  12. const fiber = await ctx.plugin(SessionInvariant)
  13. return { ctx, fiber }
  14. }
  15. describe('session-log invariants', () => {
  16. it('keeps registration global when the companion is mounted under a scope', async () => {
  17. const ctx = new Context()
  18. await ctx.plugin(SessionStore)
  19. await ctx.plugin(InvariantService)
  20. let scopedCtx!: Context
  21. await ctx.plugin(Object.assign((inner: Context) => {
  22. scopedCtx = createScope(inner, {}).ctx
  23. }, { inject: ['sessions', 'invariants'] }))
  24. await scopedCtx.plugin(SessionInvariant)
  25. const session = ctx.sessions.create(SessionId('global-under-scoped-invariants'))
  26. expect(() => {
  27. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  28. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  29. }).not.toThrow()
  30. })
  31. it('accepts a well-formed turn, step, and tool sequence', async () => {
  32. const { ctx } = await setup()
  33. const session = ctx.sessions.create()
  34. expect(() => {
  35. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  36. session.append('user/message', { content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }, { surfaceOp: 'append' })
  37. session.append('step/start', { turn: 1, step: 1 })
  38. session.append('assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'h' } })
  39. session.append('assistant/message', {
  40. provenance: { provider: 'mock', model: 'mock' },
  41. turn: 1,
  42. step: 1,
  43. content: [{ type: 'tool-call', id: CallId('c1'), name: 'echo', arguments: '{}' }],
  44. }, { surfaceOp: 'append' })
  45. session.append('tool/call', { turn: 1, step: 1, callId: CallId('c1'), name: 'echo', arguments: '{}' })
  46. session.append('tool/result', { turn: 1, step: 1, callId: CallId('c1'), content: [], isError: false }, { surfaceOp: 'append' })
  47. session.append('step/end', { turn: 1, step: 1 })
  48. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  49. }).not.toThrow()
  50. })
  51. it('does not advance committed trace state when a later dispatch listener vetoes', async () => {
  52. const { ctx } = await setup()
  53. const session = ctx.sessions.create(SessionId('dispatch-veto-rollback'))
  54. let veto = true
  55. ctx.on('internal/dispatch', (_mode, name) => {
  56. if (name !== 'session/event' || !veto) return
  57. veto = false
  58. throw new Error('later dispatch veto')
  59. })
  60. expect(() => session.append('turn/start', {
  61. turn: 1,
  62. trigger: { kind: 'message', source: { kind: 'user' } },
  63. })).toThrow('later dispatch veto')
  64. expect(session.events).toEqual([])
  65. expect(() => {
  66. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  67. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  68. }).not.toThrow()
  69. })
  70. it('applies the committed transition after another postcommit observer throws', async () => {
  71. const { ctx } = await setup()
  72. const warnings: string[] = []
  73. ctx.logger.warn = ((message: unknown) => { warnings.push(String(message)) }) as typeof ctx.logger.warn
  74. const session = ctx.sessions.create(SessionId('postcommit-peer'))
  75. ctx.on('session/event', () => { throw new Error('hostile observer') }, { prepend: true })
  76. expect(() => {
  77. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  78. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  79. }).not.toThrow()
  80. expect(warnings).toHaveLength(2)
  81. })
  82. it('rejects non-monotonic event sequence numbers', async () => {
  83. const { ctx } = await setup()
  84. const session = ctx.sessions.create()
  85. ctx.emit(scopeTarget(session, undefined), 'session/event', session, {
  86. type: 'turn/start',
  87. seq: 0,
  88. time: 1,
  89. data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
  90. } as never)
  91. expect(() => { ctx.emit(scopeTarget(session, undefined), 'session/event', session, {
  92. type: 'turn/end',
  93. seq: 0,
  94. time: 2,
  95. data: { turn: 1, reason: { kind: 'completed' } },
  96. } as never) }).toThrow(/seq must strictly increase/)
  97. })
  98. it('enforces turn numbering and encloses events other than idle context', async () => {
  99. const first = await setup()
  100. const open = first.ctx.sessions.create()
  101. open.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  102. expect(() => open.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } }))
  103. .toThrow(/turn 1 is still open/)
  104. expect(() => open.append('turn/end', { turn: 2, reason: { kind: 'completed' } }))
  105. .toThrow(/does not match open turn 1/)
  106. const second = (await setup()).ctx.sessions.create()
  107. second.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  108. second.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  109. expect(() => second.append('turn/start', { turn: 3, trigger: { kind: 'message', source: { kind: 'user' } } }))
  110. .toThrow(/expected turn 2, got 3/)
  111. const outside = (await setup()).ctx.sessions.create()
  112. expect(() => outside.append('user/message', {
  113. content: [{ type: 'text', text: 'idle context' }],
  114. source: { kind: 'plugin', plugin: 'test' },
  115. }, { surfaceOp: 'append' })).not.toThrow()
  116. expect(() => outside.append('steering/message', {
  117. turn: 1,
  118. content: [{ type: 'text', text: 'go' }],
  119. source: { kind: 'user' },
  120. }, { surfaceOp: 'append' })).toThrow(/outside any open turn/)
  121. // Merge-extensible session events use the same default enclosure branch.
  122. const appendUnknown = outside.append.bind(outside) as (type: string, data: unknown) => unknown
  123. expect(() => { appendUnknown('plugin/marker', {}) }).toThrow(/outside any open turn/)
  124. })
  125. it('enforces open-step identity and numbering', async () => {
  126. const wrongTurn = (await setup()).ctx.sessions.create()
  127. wrongTurn.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  128. expect(() => wrongTurn.append('step/start', { turn: 2, step: 1 })).toThrow(/open turn is 1/)
  129. const nested = (await setup()).ctx.sessions.create()
  130. nested.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  131. nested.append('step/start', { turn: 1, step: 1 })
  132. expect(() => nested.append('step/start', { turn: 1, step: 2 })).toThrow(/while step 1 is still open/)
  133. expect(() => nested.append('turn/end', { turn: 1, reason: { kind: 'completed' } }))
  134. .toThrow(/while step 1 is still open/)
  135. expect(() => nested.append('step/end', { turn: 1, step: 2 })).toThrow(/open is turn 1\/step 1/)
  136. expect(() => nested.append('assistant/message', {
  137. provenance: { provider: 'mock', model: 'mock' },
  138. turn: 1,
  139. step: 2,
  140. content: [],
  141. }, { surfaceOp: 'append' })).toThrow(/open is turn 1\/step 1/)
  142. const skipped = (await setup()).ctx.sessions.create()
  143. skipped.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  144. skipped.append('step/start', { turn: 1, step: 1 })
  145. skipped.append('step/end', { turn: 1, step: 1 })
  146. expect(() => skipped.append('step/start', { turn: 1, step: 3 }))
  147. .toThrow(/expected step 2 in turn 1, got 3/)
  148. })
  149. it('requires step-scoped stream and tool events to name the open step', async () => {
  150. const chunk = (await setup()).ctx.sessions.create()
  151. chunk.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  152. expect(() => chunk.append('assistant/chunk', {
  153. turn: 1,
  154. step: 1,
  155. chunk: { type: 'text-delta', index: 0, text: 'x' },
  156. })).toThrow(/open is turn 1\/step null/)
  157. const tool = (await setup()).ctx.sessions.create()
  158. tool.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  159. tool.append('step/start', { turn: 1, step: 1 })
  160. expect(() => tool.append('tool/result', {
  161. turn: 1,
  162. step: 1,
  163. callId: CallId('ghost'),
  164. content: [],
  165. isError: false,
  166. }, { surfaceOp: 'append' })).toThrow(/no prior tool\/call/)
  167. })
  168. it('keeps fresh tool-result appends open-step checked', async () => {
  169. const { ctx } = await setup()
  170. const session = ctx.sessions.create()
  171. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  172. expect(() => session.append('tool/result', {
  173. turn: 1,
  174. step: 1,
  175. callId: CallId('closed'),
  176. content: [],
  177. isError: false,
  178. }, { surfaceOp: 'append' })).toThrow(/open is turn 1\/step null/)
  179. })
  180. it('treats a validated tool-result replacement as a turn-enclosed rewrite', async () => {
  181. const { ctx } = await setup()
  182. const session = ctx.sessions.create()
  183. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  184. session.append('step/start', { turn: 1, step: 1 })
  185. session.append('tool/call', {
  186. turn: 1,
  187. step: 1,
  188. callId: CallId('rewrite'),
  189. name: 'echo',
  190. arguments: '{}',
  191. })
  192. const original = session.append('tool/result', {
  193. turn: 1,
  194. step: 1,
  195. callId: CallId('rewrite'),
  196. content: [{ type: 'text', text: 'original' }],
  197. isError: false,
  198. }, { surfaceOp: 'append' })
  199. session.append('step/end', { turn: 1, step: 1 })
  200. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  201. session.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } })
  202. expect(() => session.append('tool/result', {
  203. ...original.data,
  204. content: [{ type: 'text', text: 'pruned' }],
  205. }, {
  206. surfaceOp: { op: 'replace', start: original.seq, end: original.seq },
  207. sourceEventSeqs: [original.seq],
  208. })).not.toThrow()
  209. })
  210. it('rejects a tool-result replacement outside a turn', async () => {
  211. const { ctx } = await setup()
  212. const session = ctx.sessions.create()
  213. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  214. session.append('step/start', { turn: 1, step: 1 })
  215. session.append('tool/call', {
  216. turn: 1,
  217. step: 1,
  218. callId: CallId('rewrite'),
  219. name: 'echo',
  220. arguments: '{}',
  221. })
  222. const original = session.append('tool/result', {
  223. turn: 1,
  224. step: 1,
  225. callId: CallId('rewrite'),
  226. content: [{ type: 'text', text: 'original' }],
  227. isError: false,
  228. }, { surfaceOp: 'append' })
  229. session.append('step/end', { turn: 1, step: 1 })
  230. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  231. expect(() => session.append('tool/result', {
  232. ...original.data,
  233. content: [{ type: 'text', text: 'pruned' }],
  234. }, {
  235. surfaceOp: { op: 'replace', start: original.seq, end: original.seq },
  236. sourceEventSeqs: [original.seq],
  237. })).toThrow(/outside any open turn/)
  238. })
  239. it('allows not-started repair results and unresolved calls at step end', async () => {
  240. const repaired = (await setup()).ctx.sessions.create()
  241. expect(() => {
  242. repaired.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  243. repaired.append('step/start', { turn: 1, step: 1 })
  244. repaired.append('tool/result', {
  245. turn: 1,
  246. step: 1,
  247. callId: CallId('crashed'),
  248. content: [],
  249. isError: true,
  250. error: { name: 'ToolNotStartedError', code: TOOL_NOT_STARTED },
  251. }, { surfaceOp: 'append' })
  252. repaired.append('step/end', { turn: 1, step: 1 })
  253. repaired.append('turn/end', { turn: 1, reason: { kind: 'interrupted' } })
  254. }).not.toThrow()
  255. const unresolved = (await setup()).ctx.sessions.create()
  256. expect(() => {
  257. unresolved.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  258. unresolved.append('step/start', { turn: 1, step: 1 })
  259. unresolved.append('tool/call', { turn: 1, step: 1, callId: CallId('c1'), name: 'echo', arguments: '{}' })
  260. unresolved.append('step/end', { turn: 1, step: 1 })
  261. unresolved.append('turn/end', { turn: 1, reason: { kind: 'error', step: 1, message: 'boom' } })
  262. }).not.toThrow()
  263. })
  264. it('does not let a result in a later step satisfy an earlier call', async () => {
  265. const { ctx } = await setup()
  266. const session = ctx.sessions.create()
  267. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  268. session.append('step/start', { turn: 1, step: 1 })
  269. session.append('tool/call', { turn: 1, step: 1, callId: CallId('c1'), name: 'echo', arguments: '{}' })
  270. session.append('step/end', { turn: 1, step: 1 })
  271. session.append('step/start', { turn: 1, step: 2 })
  272. expect(() => session.append('tool/result', {
  273. turn: 1,
  274. step: 2,
  275. callId: CallId('c1'),
  276. content: [],
  277. isError: false,
  278. }, { surfaceOp: 'append' })).toThrow(/no prior tool\/call in this step/)
  279. })
  280. it('replays seeded sessions and tracks each session independently', async () => {
  281. const { ctx } = await setup()
  282. const badSeed = [
  283. { type: 'turn/start' as const, seq: 0, time: 0, data: { turn: 1, trigger: { kind: 'message' as const, source: { kind: 'user' as const } } } },
  284. { type: 'turn/start' as const, seq: 1, time: 0, data: { turn: 2, trigger: { kind: 'message' as const, source: { kind: 'user' as const } } } },
  285. ]
  286. expect(() => ctx.sessions.create(undefined, { seed: badSeed })).toThrow(InvariantError)
  287. const a = ctx.sessions.create(SessionId('a'))
  288. const b = ctx.sessions.create(SessionId('b'))
  289. a.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  290. expect(() => b.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }))
  291. .not.toThrow()
  292. })
  293. it('rebuilds trace state for sessions that exist when the companion reloads', async () => {
  294. const { ctx, fiber } = await setup()
  295. const session = ctx.sessions.create()
  296. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  297. session.append('step/start', { turn: 1, step: 1 })
  298. await fiber.dispose()
  299. await ctx.plugin(SessionInvariant)
  300. expect(() => session.append('assistant/chunk', {
  301. turn: 1,
  302. step: 1,
  303. chunk: { type: 'text-delta', index: 0, text: 'h' },
  304. })).not.toThrow()
  305. expect(() => session.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } }))
  306. .toThrow(/turn 1 is still open/)
  307. })
  308. it('removes all listeners when the companion is disposed', async () => {
  309. const { ctx, fiber } = await setup()
  310. const session = ctx.sessions.create()
  311. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  312. await fiber.dispose()
  313. expect(() => session.append('turn/start', {
  314. turn: 2,
  315. trigger: { kind: 'message', source: { kind: 'user' } },
  316. })).not.toThrow()
  317. })
  318. })