| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259 |
- import { Context } from '@deepseek-ai/cordis'
- import type { Fiber } from '@deepseek-ai/cordis'
- import type {
- RemoteEventHostInfo,
- TypertRemoteEventInvocation,
- TypertRemoteEventSource,
- } from '@deepseek-ai/dsh-api-gateway'
- import { scopeTarget } from '@deepseek-ai/dsh-scope'
- import { describe, expect, it } from 'vitest'
- import { apply, inject } from '../src/index.ts'
- interface GatewayProbe {
- source: TypertRemoteEventSource | undefined
- host: RemoteEventHostInfo | undefined
- removals: number
- registerRemoteEvents(
- source: TypertRemoteEventSource,
- host: RemoteEventHostInfo,
- ): () => Promise<void>
- }
- async function setup(): Promise<{
- readonly ctx: Context
- readonly gateway: GatewayProbe
- readonly fiber: Fiber
- }> {
- const ctx = new Context()
- const gateway: GatewayProbe = {
- source: undefined,
- host: undefined,
- removals: 0,
- registerRemoteEvents(source, host) {
- gateway.source = source
- gateway.host = host
- return async () => {
- if (gateway.source !== source) return
- gateway.source = undefined
- gateway.host = undefined
- gateway.removals += 1
- }
- },
- }
- ctx.reflect.provide('typertGateway', gateway)
- const fiber = ctx.plugin({ inject: [...inject], apply })
- await fiber
- return { ctx, gateway, fiber }
- }
- function sourceOf(gateway: GatewayProbe): TypertRemoteEventSource {
- if (gateway.source === undefined) throw new Error('fixture Gateway has no Remote event source')
- return gateway.source
- }
- function emitRaw(ctx: Context, event: string, args: readonly unknown[]): void {
- const emit = ctx.emit.bind(ctx) as unknown as (name: string, ...values: readonly unknown[]) => void
- emit(event, ...args)
- }
- function waterfallRaw(
- ctx: Context,
- target: object,
- event: string,
- args: readonly unknown[],
- next: () => Promise<unknown>,
- ): Promise<unknown> {
- const waterfall = ctx.waterfall.bind(ctx) as unknown as (
- receiver: object,
- name: string,
- ...values: readonly unknown[]
- ) => Promise<unknown>
- return waterfall(target, event, ...args, next)
- }
- function invocationOf(value: unknown): TypertRemoteEventInvocation {
- if (typeof value !== 'object' || value === null || !Object.hasOwn(value, 'context')) {
- throw new Error('fixture did not receive a scoped Remote Event invocation')
- }
- return value as TypertRemoteEventInvocation
- }
- describe('Remote event Host source', () => {
- it('registers the Host home used by Client connection generations', async () => {
- const { gateway, fiber } = await setup()
- expect(gateway.host?.home).toBeTypeOf('string')
- expect(gateway.host?.home.length).toBeGreaterThan(0)
- await fiber.dispose()
- expect(gateway.host).toBeUndefined()
- })
- it('gives each Client stream an independent allowlisted event queue', async () => {
- const { ctx, gateway, fiber } = await setup()
- const firstAbort = new AbortController()
- const secondAbort = new AbortController()
- const first = sourceOf(gateway)(firstAbort.signal)[Symbol.asyncIterator]()
- const second = sourceOf(gateway)(secondAbort.signal)[Symbol.asyncIterator]()
- emitRaw(ctx, 'settings/document-updated', ['ui-theme', 1])
- await expect(first.next()).resolves.toEqual({
- done: false,
- value: { event: 'settings/document-updated', args: ['ui-theme', 1] },
- })
- await expect(second.next()).resolves.toEqual({
- done: false,
- value: { event: 'settings/document-updated', args: ['ui-theme', 1] },
- })
- emitRaw(ctx, 'goal/activation-changed', [{
- sessionId: 'session-1',
- goal: { id: 'goal-1', revision: 1, activation: 'disarmed' },
- }])
- await expect(first.next()).resolves.toEqual({
- done: false,
- value: {
- event: 'goal/activation-changed',
- args: [{ sessionId: 'session-1', goal: { id: 'goal-1', revision: 1, activation: 'disarmed' } }],
- },
- })
- await expect(second.next()).resolves.toEqual({
- done: false,
- value: {
- event: 'goal/activation-changed',
- args: [{ sessionId: 'session-1', goal: { id: 'goal-1', revision: 1, activation: 'disarmed' } }],
- },
- })
- const firstDone = first.next()
- firstAbort.abort(new Error('first Client disconnected'))
- emitRaw(ctx, 'commands/change', [])
- await expect(firstDone).resolves.toEqual({ done: true, value: undefined })
- await expect(second.next()).resolves.toEqual({
- done: false,
- value: { event: 'commands/change', args: [] },
- })
- const secondDone = second.next()
- secondAbort.abort(new Error('second Client disconnected'))
- await expect(secondDone).resolves.toEqual({ done: true, value: undefined })
- await fiber.dispose()
- expect(gateway.source).toBeUndefined()
- expect(gateway.removals).toBe(1)
- await ctx.fiber.dispose()
- })
- it('rejects a non-JSON argument without poisoning the stream', async () => {
- const { ctx, gateway } = await setup()
- const abort = new AbortController()
- const iterator = sourceOf(gateway)(abort.signal)[Symbol.asyncIterator]()
- const pending = iterator.next()
- expect(() => {
- emitRaw(ctx, 'settings/document-updated', ['ui-theme', 1n])
- }).toThrow('argument 1 is not lossless JSON data')
- emitRaw(ctx, 'settings/document-updated', ['ui-theme', 2])
- await expect(pending).resolves.toEqual({
- done: false,
- value: { event: 'settings/document-updated', args: ['ui-theme', 2] },
- })
- const done = iterator.next()
- abort.abort()
- await expect(done).resolves.toEqual({ done: true, value: undefined })
- const alreadyAborted = new AbortController()
- alreadyAborted.abort()
- await expect(sourceOf(gateway)(alreadyAborted.signal)[Symbol.asyncIterator]().next())
- .resolves.toEqual({ done: true, value: undefined })
- await ctx.fiber.dispose()
- })
- it('bridges scoped waterfall result, next delegation, and rejection', async () => {
- const { ctx, gateway } = await setup()
- const abort = new AbortController()
- const iterator = sourceOf(gateway)(abort.signal)[Symbol.asyncIterator]()
- const agentCtx = ctx.extend()
- const agent = { id: 'agent-1', ctx: agentCtx }
- const target = scopeTarget(ctx, agent)
- const request = { questions: [], agent }
- await expect(async () => waterfallRaw(
- ctx,
- target,
- 'user-questions/request',
- [{ questions: [], agent: { id: 'agent-2', ctx: ctx.extend() } }],
- () => Promise.resolve('host fallback'),
- )).rejects.toThrow('must carry its Agent directly')
- const claimed = waterfallRaw(
- ctx,
- target,
- 'user-questions/request',
- [request],
- () => Promise.resolve('host fallback'),
- )
- const claimedDispatch = invocationOf((await iterator.next()).value)
- expect(claimedDispatch).toMatchObject({
- event: 'user-questions/request',
- request,
- context: { value: agentCtx, subject: agent, agentId: 'agent-1' },
- })
- claimedDispatch.resolve({ kind: 'result', value: 'client answer' })
- await expect(claimed).resolves.toBe('client answer')
- const delegated = waterfallRaw(
- ctx,
- target,
- 'user-questions/request',
- [request],
- () => Promise.resolve('host fallback'),
- )
- const delegatedDispatch = invocationOf((await iterator.next()).value)
- delegatedDispatch.resolve({ kind: 'next' })
- await expect(delegated).resolves.toBe('host fallback')
- const rejection = Object.assign(new Error('the user cancelled ask_user_question'), {
- code: 'ASK_CANCELLED',
- })
- const rejected = waterfallRaw(
- ctx,
- target,
- 'user-questions/request',
- [request],
- () => Promise.resolve('host fallback'),
- )
- const rejectedAssertion = expect(rejected).rejects.toBe(rejection)
- const rejectedDispatch = invocationOf((await iterator.next()).value)
- rejectedDispatch.reject(rejection)
- await rejectedAssertion
- const done = iterator.next()
- abort.abort()
- await expect(done).resolves.toEqual({ done: true, value: undefined })
- await ctx.fiber.dispose()
- })
- it('rejects a queued scoped waterfall when its source is withdrawn', async () => {
- const { ctx, gateway, fiber } = await setup()
- const abort = new AbortController()
- const iterator = sourceOf(gateway)(abort.signal)[Symbol.asyncIterator]()
- const delivery = iterator.next()
- const agent = { id: 'agent-1', ctx: ctx.extend() }
- const reason = new Error('forwarded event source removed')
- const pending = waterfallRaw(
- ctx,
- scopeTarget(ctx, agent),
- 'user-questions/request',
- [{ questions: [], agent }],
- () => Promise.resolve('host fallback'),
- )
- const rejected = expect(pending).rejects.toBe(reason)
- abort.abort(reason)
- await rejected
- await expect(delivery).resolves.toEqual({ done: true, value: undefined })
- await fiber.dispose()
- await ctx.fiber.dispose()
- })
- })
|