inbox.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278
  1. import { Context } from '@deepseek-ai/cordis'
  2. import { agentEvents, type Agent } from '@deepseek-ai/dsh-agent'
  3. import { createUserMessage, freezeMessage } from '@deepseek-ai/dsh-llm'
  4. import SessionStore, { Session, SessionId } from '@deepseek-ai/dsh-session'
  5. import type { UserMessage } from '@deepseek-ai/dsh-session'
  6. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  7. import { describe, expect, it, onTestFinished } from 'vitest'
  8. import { inboxProjectionDefinition, ReactLoopInbox } from '../src/inbox.ts'
  9. function unsupportedInbox(): Agent['inbox'] {
  10. const rejectMutation = (): never => {
  11. throw new Error('this test Agent does not support Inbox mutations')
  12. }
  13. return {
  14. nextTurn: [], nextStep: [], clear: rejectMutation, append: rejectMutation,
  15. prepend: rejectMutation, replace: rejectMutation, remove: rejectMutation, splice: rejectMutation,
  16. }
  17. }
  18. function stubAgent(rawId: string, overrides: Partial<Agent> = {}): Agent {
  19. const id = SessionId(rawId)
  20. const session = overrides.session ?? Session.create(id)
  21. const ctx = overrides.ctx ?? new Context()
  22. return {
  23. id,
  24. options: {},
  25. session,
  26. inbox: unsupportedInbox(),
  27. status: 'idle',
  28. ctx,
  29. send: () => {},
  30. followup: () => {},
  31. steer: () => {},
  32. inject: () => {},
  33. cancel() {},
  34. runMaintenance: task => task(new AbortController().signal),
  35. whenIdle: () => Promise.resolve(),
  36. ...overrides,
  37. }
  38. }
  39. async function inboxAgent(rawId: string): Promise<{
  40. ctx: Context
  41. session: Session
  42. agent: Agent
  43. inbox: ReactLoopInbox
  44. }> {
  45. const ctx = new Context()
  46. onTestFinished(() => ctx.fiber.dispose())
  47. await ctx.plugin(SessionStore)
  48. await ctx.plugin(SessionProjectionRegistry)
  49. ctx.sessionProjections.register(inboxProjectionDefinition)
  50. const session = ctx.sessions.create(SessionId(rawId))
  51. const agent = stubAgent(rawId, { ctx, session })
  52. const inbox = new ReactLoopInbox(ctx.sessionProjections, session, agentEvents(ctx, agent))
  53. Object.assign(agent, { inbox })
  54. return { ctx, session, agent, inbox }
  55. }
  56. async function reconstructPersistedInbox(
  57. rawId: string,
  58. populate: (session: Session) => void,
  59. ): Promise<Error> {
  60. const ctx = new Context()
  61. onTestFinished(() => ctx.fiber.dispose())
  62. await ctx.plugin(SessionStore)
  63. const session = ctx.sessions.create(SessionId(rawId))
  64. populate(session)
  65. await ctx.plugin(SessionProjectionRegistry)
  66. ctx.sessionProjections.register(inboxProjectionDefinition)
  67. const agent = stubAgent(rawId, { ctx, session })
  68. const inbox = new ReactLoopInbox(ctx.sessionProjections, session, agentEvents(ctx, agent))
  69. try {
  70. void inbox.nextTurn
  71. } catch (error: unknown) {
  72. if (error instanceof Error) return error
  73. throw error
  74. }
  75. throw new Error('persisted inbox reconstruction unexpectedly succeeded')
  76. }
  77. describe('ReactLoopInbox', () => {
  78. it('reads the shared projection without owning its registration', async () => {
  79. const ctx = new Context()
  80. onTestFinished(() => ctx.fiber.dispose())
  81. await ctx.plugin(SessionStore)
  82. await ctx.plugin(SessionProjectionRegistry)
  83. const unregister = ctx.sessionProjections.register(inboxProjectionDefinition)
  84. const session = ctx.sessions.create(SessionId('inbox-projection'))
  85. const pending = createUserMessage({
  86. content: [{ type: 'text', text: 'pending' }],
  87. source: { kind: 'user' },
  88. })
  89. session.append('agent/inbox/spliced', {
  90. target: 'next-turn', start: 0, inserted: [pending],
  91. })
  92. const agent = stubAgent('inbox-projection', { ctx, session })
  93. const dispatch = agentEvents(ctx, agent)
  94. const first = new ReactLoopInbox(ctx.sessionProjections, session, dispatch)
  95. const second = new ReactLoopInbox(ctx.sessionProjections, session, dispatch)
  96. expect(first.nextTurn).toEqual([pending])
  97. expect(second.nextTurn).toEqual([pending])
  98. expect(ctx.sessionProjections.snapshot(session).values.inbox).toEqual({
  99. 'next-turn': [pending],
  100. 'next-step': [],
  101. })
  102. unregister()
  103. expect(ctx.sessionProjections.stateOf(session, 'inbox')).toBeUndefined()
  104. expect(() => first.nextTurn).toThrow('its projection registration is not active')
  105. })
  106. it('rejects invalid durable coordinates and duplicate identities during reconstruction', async () => {
  107. const outOfRange = await reconstructPersistedInbox('invalid-inbox-range', (session) => {
  108. session.append('agent/inbox/spliced', {
  109. target: 'next-turn', start: 0, removedCount: 1, inserted: [],
  110. })
  111. })
  112. expect(outOfRange.message).toBe('invalid persisted inbox splice at session seq 0')
  113. expect((outOfRange.cause as Error).message).toBe('invalid inbox splice')
  114. const pending = createUserMessage({
  115. content: [{ type: 'text', text: 'duplicate' }],
  116. source: { kind: 'user' },
  117. })
  118. const duplicate = await reconstructPersistedInbox('invalid-inbox-duplicate', (session) => {
  119. session.append('agent/inbox/spliced', {
  120. target: 'next-turn', start: 0, inserted: [pending],
  121. })
  122. session.append('agent/inbox/spliced', {
  123. target: 'next-step', start: 0, inserted: [pending],
  124. })
  125. })
  126. expect(duplicate.message).toBe('invalid persisted inbox splice at session seq 1')
  127. expect((duplicate.cause as Error).message).toBe(`message "${pending.id}" is already pending`)
  128. })
  129. it('projects inherited inbox events in a forked session', async () => {
  130. const ctx = new Context()
  131. onTestFinished(() => ctx.fiber.dispose())
  132. await ctx.plugin(SessionStore)
  133. await ctx.plugin(SessionProjectionRegistry)
  134. ctx.sessionProjections.register(inboxProjectionDefinition)
  135. const parent = ctx.sessions.create(SessionId('inbox-fork-parent'))
  136. const parentAgent = stubAgent('inbox-fork-parent', { ctx, session: parent })
  137. const parentInbox = new ReactLoopInbox(ctx.sessionProjections, parent, agentEvents(ctx, parentAgent))
  138. const inherited = createUserMessage({
  139. content: [{ type: 'text', text: 'parent pending' }],
  140. source: { kind: 'user' },
  141. })
  142. parentInbox.append('next-turn', inherited)
  143. const child = ctx.sessions.fork(parent, undefined, SessionId('inbox-fork-child'))
  144. const childAgent = stubAgent('inbox-fork-child', { ctx, session: child })
  145. const childInbox = new ReactLoopInbox(ctx.sessionProjections, child, agentEvents(ctx, childAgent))
  146. expect(child.inheritedEventCount).toBe(parent.snapshotEvents().length)
  147. expect(childInbox.nextTurn).toEqual([inherited])
  148. const own = createUserMessage({
  149. content: [{ type: 'text', text: 'child pending' }],
  150. source: { kind: 'user' },
  151. })
  152. childInbox.append('next-turn', own)
  153. expect(childInbox.nextTurn).toEqual([inherited, own])
  154. })
  155. it('updates the projection cell before session observers run', async () => {
  156. const { ctx, session, inbox } = await inboxAgent('inbox-live-projection')
  157. const pending = createUserMessage({
  158. content: [{ type: 'text', text: 'direct' }],
  159. source: { kind: 'user' },
  160. })
  161. let observed: readonly UserMessage[] | undefined
  162. ctx.on('session/event', (subject, event) => {
  163. if (subject === session && event.type === 'agent/inbox/spliced') {
  164. observed = ctx.sessionProjections.stateOf(session, 'inbox')?.['next-turn']
  165. }
  166. })
  167. inbox.append('next-turn', pending)
  168. expect(observed).toEqual([pending])
  169. expect(ctx.sessionProjections.snapshot(session).values.inbox).toEqual({
  170. 'next-turn': [pending], 'next-step': [],
  171. })
  172. })
  173. it('replaces a pending message by identity across both lists', async () => {
  174. const { ctx, agent } = await inboxAgent('replace-inbox')
  175. const inserted: UserMessage[] = []
  176. const discarded: UserMessage[] = []
  177. ctx.on('agent/inbox/inserted', ({ message }) => void inserted.push(message))
  178. ctx.on('agent/inbox/discarded', ({ message }) => void discarded.push(message))
  179. const original = createUserMessage({
  180. content: [{ type: 'text', text: 'original' }],
  181. source: { kind: 'user' },
  182. })
  183. const nextStep = createUserMessage({
  184. content: [{ type: 'text', text: 'step' }],
  185. source: { kind: 'user' },
  186. })
  187. const replacement = createUserMessage({
  188. content: [{ type: 'text', text: 'replacement' }],
  189. source: { kind: 'user' },
  190. })
  191. const editedStep = freezeMessage({
  192. ...nextStep,
  193. content: [{ type: 'text', text: 'edited step' }],
  194. })
  195. agent.inbox.append('next-turn', original)
  196. agent.inbox.append('next-step', nextStep)
  197. expect(agent.inbox.replace(createUserMessage({
  198. content: [{ type: 'text', text: 'missing' }],
  199. source: { kind: 'user' },
  200. }).id, replacement)).toBe(false)
  201. expect(agent.inbox.replace(original.id, replacement)).toBe(true)
  202. expect(agent.inbox.replace(nextStep.id, editedStep)).toBe(true)
  203. expect(agent.inbox.nextTurn).toEqual([replacement])
  204. expect(agent.inbox.nextStep).toEqual([editedStep])
  205. expect(discarded).toEqual([original, nextStep])
  206. expect(inserted).toEqual([original, nextStep, replacement, editedStep])
  207. expect(() => { agent.inbox.replace(editedStep.id, replacement) })
  208. .toThrow(`message "${replacement.id}" is already pending`)
  209. })
  210. it('normalizes splice coordinates, rejects duplicate identities, and reports missing removals', async () => {
  211. const { agent } = await inboxAgent('splice-inbox')
  212. const first = createUserMessage({
  213. content: [{ type: 'text', text: 'first' }],
  214. source: { kind: 'user' },
  215. })
  216. const second = createUserMessage({
  217. content: [{ type: 'text', text: 'second' }],
  218. source: { kind: 'user' },
  219. })
  220. const prefixed = createUserMessage({
  221. content: [{ type: 'text', text: 'prefixed' }],
  222. source: { kind: 'user' },
  223. })
  224. agent.inbox.splice('next-turn', Number.NaN, Number.NaN, [first, second])
  225. expect(agent.inbox.nextTurn).toEqual([first, second])
  226. expect(agent.inbox.splice('next-turn', -1, 1, [])).toEqual([second])
  227. agent.inbox.prepend('next-turn', prefixed)
  228. expect(agent.inbox.nextTurn).toEqual([prefixed, first])
  229. expect(agent.inbox.remove(second.id)).toBe(false)
  230. expect(() => { agent.inbox.append('next-step', first) }).toThrow(`message "${first.id}" is already pending`)
  231. })
  232. it('clears both pending lists as durable cancellations', async () => {
  233. const { ctx, session, agent } = await inboxAgent('clear-inbox')
  234. const discarded: UserMessage[] = []
  235. ctx.on('agent/inbox/discarded', ({ message }) => void discarded.push(message))
  236. const nextTurn = createUserMessage({ content: [{ type: 'text', text: 'turn' }], source: { kind: 'user' } })
  237. const nextStep = createUserMessage({ content: [{ type: 'text', text: 'step' }], source: { kind: 'user' } })
  238. agent.inbox.append('next-turn', nextTurn)
  239. agent.inbox.append('next-step', nextStep)
  240. const beforeClear = session.snapshotEvents().length
  241. agent.inbox.clear()
  242. expect(agent.inbox.nextTurn).toEqual([])
  243. expect(agent.inbox.nextStep).toEqual([])
  244. expect(discarded).toEqual([nextStep, nextTurn])
  245. expect(session.snapshotEvents().slice(beforeClear).map(event => event.type === 'agent/inbox/spliced'
  246. ? event.data
  247. : event.type)).toEqual([
  248. { target: 'next-step', start: 0, removedCount: 1, inserted: [], outcome: 'canceled' },
  249. { target: 'next-turn', start: 0, removedCount: 1, inserted: [], outcome: 'canceled' },
  250. ])
  251. agent.inbox.clear()
  252. expect(session.snapshotEvents()).toHaveLength(beforeClear + 2)
  253. })
  254. })