import { describe, expect, expectTypeOf, it } from 'vitest' import { Context, Service, symbols } from 'cordis' import type { Events } from 'cordis' import { Session, SessionId } from '@deepseek-ai/dsh-session' import AgentRegistry, { AgentId, agentEvents, observeAgentStart } from '@deepseek-ai/dsh-agent' import type { Agent, AgentFactory, ContinuationStop, CreateAgentOptions, ResumeAgentOptions } from '@deepseek-ai/dsh-agent' function stubAgent(rawId: string): Agent { const id = AgentId(rawId) return { id, options: {}, session: new Session(SessionId(`${id}-session`)), status: 'idle', ctx: new Context(), send() {}, steer() {}, inject() {}, cancel() {}, whenIdle() { return Promise.resolve() }, } } describe('AgentRegistry', () => { it('keeps terminal stop decisions synchronous', () => { type TurnStopListener = Events['agent/turn-stop'] type AsyncTurnStopListener = () => Promise expectTypeOf().not.toExtend() expectTypeOf>().toEqualTypeOf() }) it('registers exact entries, emits lifecycle events, and unregisters on owner disposal', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const lifecycle: string[] = [] ctx.on('agent/created', agent => void lifecycle.push(`created:${agent.id}`)) ctx.on('agent/disposed', agent => void lifecycle.push(`disposed:${agent.id}`)) const agent = stubAgent('a1') const dispose = ctx.agents.register(agent) expect(ctx.agents.get(agent.id)).toBe(agent) expect(ctx.agents.list()).toEqual([agent]) expect(() => ctx.agents.register(stubAgent('a1'))).toThrow(/already registered/) dispose() expect(ctx.agents.get(agent.id)).toBeUndefined() expect(lifecycle).toEqual(['created:a1', 'disposed:a1']) }) it('rolls an entry back and pairs a partially delivered creation when a listener throws', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const lifecycle: string[] = [] ctx.on('agent/created', agent => void lifecycle.push(`created:${agent.id}`)) ctx.on('agent/created', () => { throw new Error('creation veto') }) ctx.on('agent/disposed', agent => void lifecycle.push(`disposed:${agent.id}`)) expect(() => ctx.agents.register(stubAgent('vetoed'))).toThrow('creation veto') expect(ctx.agents.get(AgentId('vetoed'))).toBeUndefined() expect(lifecycle).toEqual(['created:vetoed', 'disposed:vetoed']) }) it('contains asynchronous creation rejection and every disposal-listener failure', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const warnings: string[] = [] const heard: string[] = [] ctx.logger.warn = ((message: unknown) => { warnings.push(String(message)) }) as typeof ctx.logger.warn ctx.on('agent/created', () => Promise.reject(new Error('created async')) as never) ctx.on('agent/disposed', () => { throw new Error('disposed sync') }) ctx.on('agent/disposed', () => Promise.reject(new Error('disposed async')) as never) ctx.on('agent/disposed', agent => void heard.push(agent.id)) const dispose = ctx.agents.register(stubAgent('contained')) await Promise.resolve() dispose() await Promise.resolve() expect(heard).toEqual(['contained']) expect(warnings).toEqual([ 'agent "contained": agent/created listener rejected: Error: created async', 'agent "contained": agent/disposed listener threw: Error: disposed sync', 'agent "contained": agent/disposed listener rejected: Error: disposed async', ]) }) it('retains declarative startup failures, contains listeners, and clears stale records', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const warnings: string[] = [] const heard: string[] = [] ctx.logger.warn = ((message: unknown) => { warnings.push(String(message)) }) as typeof ctx.logger.warn ctx.on('agent/start-failed', () => { throw new Error('start sync') }) ctx.on('agent/start-failed', () => Promise.reject(new Error('start async')) as never) ctx.on('agent/start-failed', id => void heard.push(id)) const firstError = new Error('first failure') const disposeFirst = ctx.agents.reportStartFailure(AgentId('main'), firstError) expect(ctx.agents.getStartFailure(AgentId('main'))).toBe(firstError) await Promise.resolve() expect(heard).toEqual(['main']) expect(warnings).toEqual([ 'agent "main": agent/start-failed listener threw: Error: start sync', 'agent "main": agent/start-failed listener rejected: Error: start async', ]) disposeFirst() expect(ctx.agents.getStartFailure(AgentId('main'))).toBeUndefined() const disposeSecond = ctx.agents.reportStartFailure(AgentId('main'), new Error('second failure')) const thirdError = new Error('third failure') const disposeThird = ctx.agents.reportStartFailure(AgentId('main'), thirdError) expect(ctx.agents.getStartFailure(AgentId('main'))).toBe(thirdError) disposeSecond() expect(ctx.agents.getStartFailure(AgentId('main'))).toBe(thirdError) disposeThird() const disposeOccupiedAgent = ctx.agents.register(stubAgent('occupied')) const occupiedError = new Error('occupied failure') const disposeOccupiedFailure = ctx.agents.reportStartFailure(AgentId('occupied'), occupiedError) expect(() => { ctx.agents.enter(stubAgent('occupied')) }).toThrow(/already registered/) expect(ctx.agents.getStartFailure(AgentId('occupied'))).toBe(occupiedError) disposeOccupiedFailure() disposeOccupiedAgent() const disposeCleared = ctx.agents.reportStartFailure(AgentId('main'), new Error('cleared failure')) ctx.agents.register(stubAgent('main'))() expect(ctx.agents.getStartFailure(AgentId('main'))).toBeUndefined() disposeCleared() await ctx.fiber.dispose() }) it('observes immediate, retained, live, unrelated, and cancelled startup outcomes', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const outcomes: string[] = [] const existing = stubAgent('existing') const disposeExisting = ctx.agents.register(existing) observeAgentStart(ctx, existing.id, { onStarted: agent => void outcomes.push(`started:${agent.id}`), onFailed: error => void outcomes.push(`failed:${error.message}`), })() const retainedError = new Error('retained') const disposeRetained = ctx.agents.reportStartFailure(AgentId('retained'), retainedError) observeAgentStart(ctx, AgentId('retained'), { onStarted: agent => void outcomes.push(`started:${agent.id}`), onFailed: error => void outcomes.push(`failed:${error.message}`), })() const stopLiveStart = observeAgentStart(ctx, AgentId('live-start'), { onStarted: agent => void outcomes.push(`started:${agent.id}`), onFailed: error => void outcomes.push(`failed:${error.message}`), }) const disposeUnrelatedAgent = ctx.agents.register(stubAgent('unrelated')) const disposeUnrelatedFailure = ctx.agents.reportStartFailure(AgentId('unrelated'), new Error('unrelated')) const disposeLiveAgent = ctx.agents.register(stubAgent('live-start')) stopLiveStart() const stopLiveFailure = observeAgentStart(ctx, AgentId('live-failure'), { onStarted: agent => void outcomes.push(`started:${agent.id}`), onFailed: error => void outcomes.push(`failed:${error.message}`), }) const disposeLiveFailure = ctx.agents.reportStartFailure(AgentId('live-failure'), new Error('live')) stopLiveFailure() const stopCancelled = observeAgentStart(ctx, AgentId('cancelled'), { onStarted: agent => void outcomes.push(`started:${agent.id}`), onFailed: error => void outcomes.push(`failed:${error.message}`), }) stopCancelled() const disposeCancelledFailure = ctx.agents.reportStartFailure(AgentId('cancelled'), new Error('cancelled')) expect(outcomes).toEqual([ 'started:existing', 'failed:retained', 'started:live-start', 'failed:live', ]) disposeCancelledFailure() disposeLiveFailure() disposeLiveAgent() disposeUnrelatedFailure() disposeUnrelatedAgent() disposeRetained() disposeExisting() await ctx.fiber.dispose() }) it('separates entry from announcement and stale/idempotent detach cannot remove a replacement', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const lifecycle: string[] = [] ctx.on('agent/created', agent => void lifecycle.push(`created:${agent.id}`)) ctx.on('agent/disposed', agent => void lifecycle.push(`disposed:${agent.id}`)) const first = stubAgent('split') const detachFirst = ctx.agents.enter(first) expect(lifecycle).toEqual([]) ctx.agents.announce(first) expect(() => { ctx.agents.announce(first) }).toThrow(/already announced/) detachFirst() detachFirst() const replacement = stubAgent('split') const detachReplacement = ctx.agents.enter(replacement) detachFirst() expect(ctx.agents.get(replacement.id)).toBe(replacement) expect(() => { ctx.agents.announce(first) }).toThrow(/not live/) detachReplacement() expect(lifecycle).toEqual(['created:split', 'disposed:split']) }) it('defers detach requested by a creation listener until that dispatch unwinds', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const order: string[] = [] const agent = stubAgent('reentrant') ctx.on('agent/created', () => { order.push(`first:${ctx.agents.get(agent.id) === agent}`) detach() order.push(`after-detach:${ctx.agents.get(agent.id) === agent}`) }) ctx.on('agent/created', () => void order.push(`second:${ctx.agents.get(agent.id) === agent}`)) ctx.on('agent/disposed', () => void order.push('disposed')) const detach = ctx.agents.enter(agent) ctx.agents.announce(agent) expect(order).toEqual(['first:true', 'after-detach:true', 'second:true', 'disposed']) expect(ctx.agents.get(agent.id)).toBeUndefined() }) }) describe('agentEvents()', () => { it('contains each synchronous throw and returned-promise rejection', async () => { const ctx = new Context() const warnings: string[] = [] const heard: string[] = [] ctx.logger.warn = ((message: unknown) => { warnings.push(String(message)) }) as typeof ctx.logger.warn const agent = stubAgent('event') ctx.on('agent/status', () => { throw new Error('sync listener') }) ctx.on('agent/status', () => Promise.reject(new Error('async listener')) as never) ctx.on('agent/status', (_agent, status) => void heard.push(status)) agentEvents(ctx, agent).emit('agent/status', 'running') await Promise.resolve() expect(heard).toEqual(['running']) expect(warnings).toEqual([ 'agent event "agent/status" listener threw: Error: sync listener', 'agent event "agent/status" listener rejected: Error: async listener', ]) }) }) describe('AgentRegistry factory seam', () => { function stubFactory() { const calls: { create: Array<{ ownerCtx: Context; options: CreateAgentOptions }> resume: Array<{ ownerCtx: Context; options: ResumeAgentOptions }> } = { create: [], resume: [] } const factory: AgentFactory = { async createAgent(ownerCtx, options) { calls.create.push({ ownerCtx, options }) return { agent: stubAgent(options.agentId), dispose: () => Promise.resolve() } }, async resume(ownerCtx, options) { calls.resume.push({ ownerCtx, options }) return { agent: stubAgent(options.agentId), dispose: () => Promise.resolve() } }, } return { factory, calls } } it('requires a factory and delegates through the calling context', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) await expect(ctx.agents.create({ agentId: AgentId('a'), sessionId: SessionId('s') })).rejects.toThrow(/no agent factory/) const { factory, calls } = stubFactory() ctx.agents.setFactory(factory) let callerFiber: Context['fiber'] | undefined await ctx.plugin(Object.assign(async (inner: Context) => { callerFiber = inner.fiber await inner.agents.create({ agentId: AgentId('create'), sessionId: SessionId('create-s') }) await inner.agents.resume({ agentId: AgentId('resume'), resumeSessionId: SessionId('resume-s') }) }, { inject: ['agents'] })) expect(calls.create[0]?.ownerCtx.fiber).toBe(callerFiber) expect(calls.resume[0]?.ownerCtx.fiber).toBe(callerFiber) }) it('rejects a second factory and clears the slot with its owner (HMR)', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const owner = await ctx.plugin(Object.assign((inner: Context) => { inner.agents.setFactory(stubFactory().factory) expect(() => inner.agents.setFactory(stubFactory().factory)).toThrow(/already registered/) }, { inject: ['agents'] })) await expect(ctx.agents.create({ agentId: AgentId('before'), sessionId: SessionId('before-s') })).resolves.toBeDefined() await owner.dispose() await expect(ctx.agents.create({ agentId: AgentId('after'), sessionId: SessionId('after-s') })).rejects.toThrow(/no agent factory/) }) it('canonicalizes an already traced Service before tracing it for the caller', async () => { const ctx = new Context() await ctx.plugin(AgentRegistry) const states = new WeakMap() class TracedFactory extends Service implements AgentFactory { constructor(inner: Context) { super(inner, 'tracedFactory') states.set(this, []) } private calls(): string[] { const original = (this as unknown as { [symbols.original]?: TracedFactory })[symbols.original] ?? this const calls = states.get(original) if (calls === undefined) throw new Error('factory receiver was not canonicalized') return calls } async createAgent(_ownerCtx: Context, options: CreateAgentOptions) { this.calls().push('create') return { agent: stubAgent(options.agentId), dispose: () => Promise.resolve() } } async resume(_ownerCtx: Context, options: ResumeAgentOptions) { this.calls().push('resume') return { agent: stubAgent(options.agentId), dispose: () => Promise.resolve() } } } await ctx.plugin(TracedFactory) const traced = (ctx as Context & { tracedFactory: TracedFactory }).tracedFactory ctx.agents.setFactory(traced) await ctx.agents.create({ agentId: AgentId('create'), sessionId: SessionId('create-s') }) await ctx.agents.resume({ agentId: AgentId('resume'), resumeSessionId: SessionId('resume-s') }) const raw = (traced as unknown as { [symbols.original]?: TracedFactory })[symbols.original] expect(states.get(raw!)).toEqual(['create', 'resume']) }) })