| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083 |
- import { randomUUID } from 'node:crypto'
- import { once } from 'node:events'
- import { afterEach, describe, expect, it, vi } from 'vitest'
- import WebSocket, { type RawData } from 'ws'
- import { Context, Service, symbols } from '@deepseek-ai/cordis'
- import { apply as applyConnection, inject as connectionInject } from '@deepseek-ai/dsh-client-connection'
- import WebServer from '@deepseek-ai/dsh-host-webserver'
- import {
- bindTypertRemote,
- Remote,
- type InvocationDescriptor,
- type TypertContextMap,
- type TypertContextWire,
- TypertRemoteFailure,
- } from '@deepseek-ai/dsh-typert-protocol'
- import TypertRegistry from '@deepseek-ai/dsh-typert-registry'
- import TypertGatewayService, {
- TypertGatewayError,
- type TypertRemoteEventDispatch,
- type TypertRemoteEventInvocation,
- type TypertRemoteEventOutcome,
- } from '@deepseek-ai/dsh-api-gateway'
- import { z } from 'zod'
- import type {
- RemoteEventClientId,
- RemoteEventInvocationFrame,
- } from '../src/stream-protocol.ts'
- vi.mock('node:crypto', async (importOriginal) => {
- const actual = await importOriginal<typeof import('node:crypto')>()
- return { ...actual, randomUUID: vi.fn(actual.randomUUID) }
- })
- const randomUuid = vi.mocked(randomUUID)
- type AgentWireId = TypertContextWire<TypertContextMap['agent']>
- const agentId = (value: string): AgentWireId => value as AgentWireId
- class FeedService extends Service {
- readonly typertRemote = bindTypertRemote(this, 'feed')
- readonly signals: AbortSignal[] = []
- returns = 0
- constructor(ctx: Context) {
- super(ctx, 'feed')
- }
- @Remote({ mode: 'stream' })
- async *follow(label: string, signal: AbortSignal): AsyncIterable<string> {
- this.signals.push(signal)
- try {
- yield `${label}:ready`
- await new Promise<void>((resolve) => {
- if (signal.aborted) resolve()
- else signal.addEventListener('abort', () => { resolve() }, { once: true })
- })
- } finally {
- this.returns += 1
- }
- }
- @Remote({ mode: 'stream' })
- *sync(label: string): Iterable<string> {
- yield `${label}:one`
- yield `${label}:two`
- }
- @Remote({ mode: 'stream' })
- *invalid(): Iterable<string> {
- yield 42 as unknown as string
- }
- @Remote({ mode: 'stream' })
- *nonJson(): Iterable<unknown> {
- yield 1n
- }
- @Remote({ mode: 'stream' })
- missing(): Iterable<string> {
- return null as unknown as Iterable<string>
- }
- @Remote({ mode: 'stream' })
- *src(label: string): Iterable<string> {
- yield `${label}:src`
- }
- @Remote({ mode: 'stream' })
- abortBeforeOpen(signal: AbortSignal): Iterable<string> {
- if (signal.aborted) throw new Error('fixture observed pre-open cancellation')
- return []
- }
- @Remote({ mode: 'stream' })
- reject(): Iterable<string> {
- throw new TypertRemoteFailure({
- code: 'fixture-rejected', message: 'fixture rejected the stream', details: { retryable: false },
- })
- }
- @Remote({ mode: 'stream' })
- rejectWithNonJsonDetails(): Iterable<string> {
- throw new TypertRemoteFailure({
- code: 'fixture-broken', message: 'fixture emitted invalid details', details: { count: 1n },
- })
- }
- unary(label: string): string {
- return label
- }
- }
- const roots: Context[] = []
- class RemoteEventSourceProbe {
- readonly source = (signal: AbortSignal): AsyncIterable<TypertRemoteEventDispatch> => {
- this.signal = signal
- return this.iterate(signal)
- }
- signal: AbortSignal | undefined
- private readonly dispatches: TypertRemoteEventDispatch[] = []
- private wake: (() => void) | undefined
- push(dispatch: TypertRemoteEventDispatch): void {
- this.dispatches.push(dispatch)
- this.wake?.()
- this.wake = undefined
- }
- private async *iterate(signal: AbortSignal): AsyncGenerator<TypertRemoteEventDispatch> {
- const aborted = (): void => {
- this.wake?.()
- this.wake = undefined
- }
- signal.addEventListener('abort', aborted, { once: true })
- try {
- while (!signal.aborted) {
- while (this.dispatches.length > 0) {
- yield this.dispatches.shift() as TypertRemoteEventDispatch
- }
- if (signal.aborted) return
- await new Promise<void>((resolve) => { this.wake = resolve })
- this.wake = undefined
- }
- } finally {
- signal.removeEventListener('abort', aborted)
- }
- }
- }
- interface PendingInvocationProbe {
- readonly dispatch: TypertRemoteEventInvocation
- readonly outcome: Promise<TypertRemoteEventOutcome>
- readonly resolve: (outcome: TypertRemoteEventOutcome) => void
- readonly reject: (reason: unknown) => void
- }
- function pendingInvocation(
- context: Context,
- signal?: AbortSignal,
- prompt = 'ship',
- ): PendingInvocationProbe {
- const subject = { ctx: context }
- const settled = Promise.withResolvers<TypertRemoteEventOutcome>()
- const resolve = vi.fn((outcome: TypertRemoteEventOutcome) => {
- settled.resolve(outcome)
- })
- const reject = vi.fn((reason: unknown) => {
- settled.reject(reason)
- })
- return {
- dispatch: {
- event: 'fixture/approval',
- request: { prompt, agent: subject, ...(signal === undefined ? {} : { signal }) },
- context: { value: context, subject },
- resolve,
- reject,
- },
- outcome: settled.promise,
- resolve,
- reject,
- }
- }
- afterEach(async () => {
- randomUuid.mockClear()
- await Promise.all(roots.splice(0).map(ctx => ctx.fiber.dispose()))
- })
- describe('Typert Remote streams', () => {
- it('opens decoded carrier payloads through the in-process wire adapter', async () => {
- const { ctx } = await setup(false)
- const source = await ctx.typertGateway.wireStream.open(
- 'feed/sync',
- { args: { label: 'wire' } },
- new AbortController().signal,
- )
- await expect(collect(source)).resolves.toEqual(['wire:one', 'wire:two'])
- })
- it('passes Iterable and AsyncIterable items through and returns the iterator on cancellation', async () => {
- const { ctx, service } = await setup(false)
- const abort = new AbortController()
- const source = await ctx.typertGateway.stream({
- namespace: 'feed',
- method: 'follow',
- args: { label: 'a' },
- signal: abort.signal,
- })
- const iterator = source[Symbol.asyncIterator]()
- await expect(iterator.next()).resolves.toEqual({ done: false, value: 'a:ready' })
- const pending = iterator.next()
- abort.abort(new Error('fixture cancellation'))
- await expect(pending).rejects.toThrow('Remote invocation "feed/follow" was aborted')
- expect(service.signals).toEqual([abort.signal])
- expect(service.returns).toBe(1)
- await expect(collect(await ctx.typertGateway.stream({
- namespace: 'feed', method: 'sync', args: { label: 'b' },
- }))).resolves.toEqual(['b:one', 'b:two'])
- await expect(collect(await ctx.typertGateway.stream({
- namespace: 'feed', method: 'invalid', args: {},
- }))).resolves.toEqual([42])
- await expect(collect(await ctx.typertGateway.stream({
- namespace: 'feed', method: 'nonJson', args: {},
- }))).resolves.toEqual([1n])
- await expect(ctx.typertGateway.stream({
- namespace: 'feed', method: 'missing', args: {},
- })).rejects.toMatchObject({ code: 'result-invalid' })
- await expect(collect(await ctx.typertGateway.stream({
- namespace: 'feed', method: 'src', args: { label: 'c' },
- }))).resolves.toEqual(['c:src'])
- const abortedBeforeOpen = new AbortController()
- abortedBeforeOpen.abort(new Error('cancelled before open'))
- await expect(ctx.typertGateway.stream({
- namespace: 'feed', method: 'abortBeforeOpen', args: {}, signal: abortedBeforeOpen.signal,
- })).rejects.toThrow('Remote invocation "feed/abortBeforeOpen" was aborted')
- const abortedBeforeIteration = new AbortController()
- abortedBeforeIteration.abort(new Error('cancelled before iteration'))
- const preCancelled = await ctx.typertGateway.stream({
- namespace: 'feed', method: 'sync', args: { label: 'ignored' }, signal: abortedBeforeIteration.signal,
- })
- await expect(collect(preCancelled)).rejects.toThrow('Remote invocation "feed/sync" was aborted')
- })
- it('keeps unary and stream invocation modes distinct', async () => {
- const { ctx } = await setup(false)
- await expect(ctx.typertGateway.invoke({
- namespace: 'feed', method: 'sync', args: { label: 'a' },
- })).rejects.toMatchObject({ code: 'signature-invalid' } satisfies Partial<TypertGatewayError>)
- await expect(ctx.typertGateway.stream({
- namespace: 'feed', method: 'unary', args: { label: 'a' },
- })).rejects.toMatchObject({ code: 'signature-invalid' } satisfies Partial<TypertGatewayError>)
- })
- it('multiplexes independent streams over one WebSocket and propagates cancellation', async () => {
- const { ctx, service } = await setup(true)
- const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`)
- await once(socket, 'open')
- const frames: Record<string, unknown>[] = []
- socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
- sendOpen(socket, 'a', 'feed/follow', { label: 'a' })
- sendOpen(socket, 'b', 'feed/follow', { label: 'b' })
- await vi.waitFor(() => {
- expect(frames).toEqual(expect.arrayContaining([
- { type: 'item', streamId: 'a', value: 'a:ready' },
- { type: 'item', streamId: 'b', value: 'b:ready' },
- ]))
- })
- expect(service.signals.map(signal => signal.aborted)).toEqual([false, false])
- expect(service.returns).toBe(0)
- socket.send(JSON.stringify({ type: 'cancel', streamId: 'a' }))
- await vi.waitFor(() => { expect(service.returns).toBe(1) })
- expect(service.signals[0]?.aborted).toBe(true)
- expect(service.signals[1]?.aborted).toBe(false)
- sendOpen(socket, 'sync', 'feed/sync', { label: 's' })
- sendOpen(socket, 'invalid', 'feed/invalid', {})
- sendOpen(socket, 'non-json', 'feed/nonJson', {})
- sendOpen(socket, 'rejected', 'feed/reject', {})
- await vi.waitFor(() => {
- expect(frames.filter(frame => frame.streamId === 'sync')).toEqual([
- { type: 'item', streamId: 'sync', value: 's:one' },
- { type: 'item', streamId: 'sync', value: 's:two' },
- { type: 'end', streamId: 'sync' },
- ])
- expect(frames.filter(frame => frame.streamId === 'invalid')).toEqual([
- { type: 'item', streamId: 'invalid', value: 42 },
- { type: 'end', streamId: 'invalid' },
- ])
- expect(frames.find(frame => frame.streamId === 'non-json')).toMatchObject({
- type: 'error', error: { code: 'internal' },
- })
- expect(frames.find(frame => frame.streamId === 'rejected')).toEqual({
- type: 'error',
- streamId: 'rejected',
- error: {
- code: 'fixture-rejected',
- message: 'fixture rejected the stream',
- details: { retryable: false },
- },
- })
- })
- const closed = once(socket, 'close')
- sendOpen(socket, 'broken-error', 'feed/rejectWithNonJsonDetails', {})
- const closeEvent = await closed
- expect(closeEvent[0]).toBe(1011)
- expect(String(closeEvent[1])).toBe('Remote stream failure could not be delivered')
- await vi.waitFor(() => { expect(service.returns).toBe(2) })
- expect(service.signals[1]?.aborted).toBe(true)
- })
- it('carries the registered Remote event source and withdraws its active stream', async () => {
- const { ctx } = await setup(true)
- let sourceSignal: AbortSignal | undefined
- const sourceClosed = vi.fn()
- const publish = Promise.withResolvers<undefined>()
- const source = (signal: AbortSignal): AsyncIterable<{ event: string; args: readonly unknown[] }> => {
- sourceSignal = signal
- return (async function *() {
- try {
- await publish.promise
- yield { event: 'fixture/changed', args: ['settings'] }
- await new Promise<void>((resolve) => {
- if (signal.aborted) resolve()
- else signal.addEventListener('abort', () => { resolve() }, { once: true })
- })
- } finally {
- sourceClosed()
- }
- })()
- }
- const unregister = ctx.typertGateway.registerRemoteEvents(source)
- expect(() => { ctx.typertGateway.registerRemoteEvents(source) })
- .toThrow('forwarded Remote event source is already registered')
- const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`)
- await once(socket, 'open')
- const frames: Record<string, unknown>[] = []
- socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
- sendOpen(socket, 'events', '$events', {})
- await vi.waitFor(() => {
- const eventFrames = frames.filter(frame => frame.streamId === 'events')
- expect(eventFrames).toHaveLength(1)
- expect(eventFrames[0]).toMatchObject({
- type: 'item', streamId: 'events', value: { type: 'ready' },
- })
- expect(typeof Reflect.get(eventFrames[0]!.value as object, 'clientId')).toBe('string')
- })
- publish.resolve(undefined)
- await vi.waitFor(() => {
- const eventFrames = frames.filter(frame => frame.streamId === 'events').slice(0, 2)
- expect(eventFrames).toHaveLength(2)
- expect(eventFrames[0]).toMatchObject({
- type: 'item', streamId: 'events', value: { type: 'ready' },
- })
- expect(typeof Reflect.get(eventFrames[0]!.value as object, 'clientId')).toBe('string')
- expect(eventFrames[1]).toEqual({
- type: 'item', streamId: 'events', value: {
- type: 'emit', event: 'fixture/changed', args: ['settings'],
- },
- })
- })
- expect(sourceSignal?.aborted).toBe(false)
- await unregister()
- expect(sourceClosed).toHaveBeenCalledOnce()
- await vi.waitFor(() => {
- expect(sourceSignal?.aborted).toBe(true)
- expect(frames).toContainEqual({ type: 'end', streamId: 'events' })
- })
- const unregisterReplacement = ctx.typertGateway.registerRemoteEvents(source)
- await unregister()
- expect(() => { ctx.typertGateway.registerRemoteEvents(source) })
- .toThrow('forwarded Remote event source is already registered')
- await unregisterReplacement()
- socket.close()
- })
- it('rejects a scoped dispatch yielded after its Remote event source is withdrawn', async () => {
- const { ctx } = await setup(false)
- const publish = Promise.withResolvers<undefined>()
- const agent = ctx.extend()
- const pending = pendingInvocation(agent)
- const source = (): AsyncIterable<TypertRemoteEventDispatch> => (async function* () {
- await publish.promise
- yield pending.dispatch
- })()
- const unregister = ctx.typertGateway.registerRemoteEvents(source)
- const rejected = expect(pending.outcome).rejects.toThrow(
- 'forwarded Remote event source was removed',
- )
- publish.resolve(undefined)
- await unregister()
- await rejected
- expect(pending.reject).toHaveBeenCalledTimes(1)
- expect(pending.resolve).not.toHaveBeenCalled()
- })
- it('cancels a pending waterfall when its source rejects during removal', async () => {
- const { ctx } = await setup(true)
- const agent = ctx.extend()
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === agent ? agentId('agent-removal') : undefined,
- resolve: id => id === 'agent-removal' ? agent : undefined,
- })
- const pending = pendingInvocation(agent)
- const rejected = expect(pending.outcome).rejects.toThrow(
- 'forwarded Remote event source was removed',
- )
- const unregister = ctx.typertGateway.registerRemoteEvents(signal => (async function* () {
- yield pending.dispatch
- await new Promise<void>((resolve) => {
- if (signal.aborted) resolve()
- else signal.addEventListener('abort', () => { resolve() }, { once: true })
- })
- throw new Error('fixture source rejected during removal')
- })())
- const client = await openEventClient(ctx, 'events-removal')
- await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
- await unregister()
- await rejected
- expect(pending.reject).toHaveBeenCalledTimes(1)
- expect(pending.resolve).not.toHaveBeenCalled()
- await vi.waitFor(() => {
- expect(client.frames).toContainEqual({ type: 'end', streamId: client.streamId })
- })
- client.socket.close()
- })
- it('delegates unavailable Contexts and rejects malformed scoped invocations', async () => {
- const { ctx } = await setup(false)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- for (const event of [42, ''] as const) {
- const invalidName = pendingInvocation(ctx)
- const rejected = expect(invalidName.outcome).rejects.toThrow(
- 'Remote event name must be a nonempty string',
- )
- source.push({
- ...invalidName.dispatch,
- event: event as unknown as string,
- })
- await rejected
- }
- const unavailable = pendingInvocation(ctx)
- source.push(unavailable.dispatch)
- await expect(unavailable.outcome).resolves.toEqual({ kind: 'next' })
- expect(unavailable.reject).not.toHaveBeenCalled()
- let selected = ctx.extend()
- let identity: unknown = 1n
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === selected ? identity as AgentWireId : undefined,
- resolve: () => selected,
- })
- const nonJsonIdentity = pendingInvocation(selected)
- const nonJsonRejected = expect(nonJsonIdentity.outcome).rejects.toThrow(
- 'require a non-empty Agent identity',
- )
- source.push(nonJsonIdentity.dispatch)
- await nonJsonRejected
- identity = 'agent-invalid-request'
- const invalidRequest = pendingInvocation(selected)
- const invalidRequestRejected = expect(invalidRequest.outcome).rejects.toThrow(
- 'must carry its scoped Agent directly',
- )
- source.push({
- ...invalidRequest.dispatch,
- request: {},
- })
- await invalidRequestRejected
- const staleFiber = ctx.plugin(() => {})
- await staleFiber
- selected = staleFiber.ctx
- identity = 'agent-stale'
- await staleFiber.dispose()
- const stale = pendingInvocation(selected)
- source.push(stale.dispatch)
- await expect(stale.outcome).resolves.toEqual({ kind: 'next' })
- expect(stale.reject).not.toHaveBeenCalled()
- selected = ctx.extend()
- identity = 'agent-cancelled'
- const abort = new AbortController()
- abort.abort('fixture non-error cancellation')
- const cancelled = pendingInvocation(selected, abort.signal)
- const cancelledOutcome = expect(cancelled.outcome).rejects.toMatchObject({
- message: 'typert gateway: Remote event was cancelled',
- cause: 'fixture non-error cancellation',
- })
- source.push(cancelled.dispatch)
- await cancelledOutcome
- await unregister()
- })
- it('rejects notification arguments that are not lossless JSON arrays', async () => {
- const { ctx } = await setup(false)
- const frames = [
- { event: 'fixture/changed', args: {} },
- { event: 'fixture/changed', args: [1n] },
- ]
- for (const frame of frames) {
- let sourceSignal: AbortSignal | undefined
- const unregister = ctx.typertGateway.registerRemoteEvents((signal) => {
- sourceSignal = signal
- return (async function* () {
- yield frame as unknown as TypertRemoteEventDispatch
- })()
- })
- await vi.waitFor(() => { expect(sourceSignal?.aborted).toBe(true) })
- const reason: unknown = sourceSignal?.reason
- if (!(reason instanceof Error)) throw new Error('Remote event source did not fail with an Error')
- expect(reason.message).toContain('arguments are not lossless JSON data')
- await unregister()
- }
- })
- it('retries a colliding Remote event id before publishing the second waterfall', async () => {
- const { ctx } = await setup(false)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const agent = ctx.extend()
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === agent ? agentId('agent-collision') : undefined,
- resolve: id => id === 'agent-collision' ? agent : undefined,
- })
- const firstId = '00000000-0000-4000-8000-000000000001' as ReturnType<typeof randomUUID>
- const secondId = '00000000-0000-4000-8000-000000000002' as ReturnType<typeof randomUUID>
- randomUuid.mockReturnValueOnce(firstId).mockReturnValueOnce(firstId).mockReturnValueOnce(secondId)
- const firstAbort = new AbortController()
- const secondAbort = new AbortController()
- const first = pendingInvocation(agent, firstAbort.signal, 'first')
- const second = pendingInvocation(agent, secondAbort.signal, 'second')
- source.push(first.dispatch)
- await vi.waitFor(() => { expect(randomUuid).toHaveBeenCalledTimes(1) })
- source.push(second.dispatch)
- await vi.waitFor(() => { expect(randomUuid).toHaveBeenCalledTimes(3) })
- const firstReason = new Error('cancel first collision fixture')
- const secondReason = new Error('cancel second collision fixture')
- const firstRejected = expect(first.outcome).rejects.toBe(firstReason)
- const secondRejected = expect(second.outcome).rejects.toBe(secondReason)
- firstAbort.abort(firstReason)
- secondAbort.abort(secondReason)
- await firstRejected
- await secondRejected
- await unregister()
- })
- it('retries a colliding Remote event Client id before opening the second generation', async () => {
- const { ctx } = await setup(true)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const firstId = '00000000-0000-4000-8000-000000000011' as ReturnType<typeof randomUUID>
- const secondId = '00000000-0000-4000-8000-000000000012' as ReturnType<typeof randomUUID>
- randomUuid.mockReturnValueOnce(firstId).mockReturnValueOnce(firstId).mockReturnValueOnce(secondId)
- const first = await openEventClient(ctx, 'events-client-id-a')
- const second = await openEventClient(ctx, 'events-client-id-b')
- expect(first.clientId).toBe(firstId)
- expect(second.clientId).toBe(secondId)
- expect(randomUuid).toHaveBeenCalledTimes(3)
- first.socket.close()
- second.socket.close()
- await unregister()
- })
- it('fans one scoped waterfall out and accepts the first Client result', async () => {
- const { ctx } = await setup(true)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const agent = ctx.extend()
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === agent ? agentId('agent-1') : undefined,
- resolve: id => id === 'agent-1' ? agent : undefined,
- })
- const first = await openEventClient(ctx, 'events-a')
- const second = await openEventClient(ctx, 'events-b')
- const pending = pendingInvocation(agent)
- source.push(pending.dispatch)
- await vi.waitFor(() => {
- expect(deliveredInvocation(first)).toBeDefined()
- expect(deliveredInvocation(second)).toBeDefined()
- })
- const firstFrame = deliveredInvocation(first)!
- const secondFrame = deliveredInvocation(second)!
- expect(firstFrame.eventId).toBe(secondFrame.eventId)
- expect(firstFrame).toMatchObject({
- type: 'waterfall',
- event: 'fixture/approval',
- agentId: 'agent-1',
- request: { prompt: 'ship' },
- })
- expect(firstFrame).not.toHaveProperty('deliveryId')
- expect(secondFrame).not.toHaveProperty('deliveryId')
- await sendEventResult(second, secondFrame, {
- kind: 'result', value: 'allowed',
- })
- await expect(pending.outcome).resolves.toEqual({ kind: 'result', value: 'allowed' })
- await vi.waitFor(() => {
- expect(first.frames).toContainEqual({
- type: 'item',
- streamId: first.streamId,
- value: { type: 'cancel', eventId: firstFrame.eventId },
- })
- })
- await sendEventResult(first, firstFrame, {
- kind: 'result', value: 'rejected',
- })
- expect(pending.resolve).toHaveBeenCalledTimes(1)
- expect(pending.reject).not.toHaveBeenCalled()
- first.socket.close()
- second.socket.close()
- await unregister()
- })
- it('rejects the Host waterfall with the first Client listener rejection', async () => {
- const { ctx } = await setup(true)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const agent = ctx.extend()
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === agent ? agentId('agent-rejected') : undefined,
- resolve: id => id === 'agent-rejected' ? agent : undefined,
- })
- const client = await openEventClient(ctx, 'events-rejected')
- const pending = pendingInvocation(agent)
- source.push(pending.dispatch)
- await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
- const frame = deliveredInvocation(client)!
- const rejected = expect(pending.outcome).rejects.toMatchObject({
- name: 'UserQuestionError',
- message: 'the user cancelled ask_user_question',
- code: 'ASK_CANCELLED',
- details: { questionId: 'question-1' },
- })
- await sendEventResult(client, frame, {
- kind: 'rejected',
- error: {
- name: 'UserQuestionError',
- message: 'the user cancelled ask_user_question',
- code: 'ASK_CANCELLED',
- details: { questionId: 'question-1' },
- },
- })
- await rejected
- expect(pending.reject).toHaveBeenCalledTimes(1)
- expect(pending.resolve).not.toHaveBeenCalled()
- client.socket.close()
- await unregister()
- })
- it('delegates to the Host only after every active Client returns next', async () => {
- const { ctx } = await setup(true)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const agent = ctx.extend()
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === agent ? agentId('agent-1') : undefined,
- resolve: id => id === 'agent-1' ? agent : undefined,
- })
- const first = await openEventClient(ctx, 'events-next-a')
- const second = await openEventClient(ctx, 'events-next-b')
- const pending = pendingInvocation(agent)
- source.push(pending.dispatch)
- await vi.waitFor(() => {
- expect(deliveredInvocation(first)).toBeDefined()
- expect(deliveredInvocation(second)).toBeDefined()
- })
- const firstFrame = deliveredInvocation(first)!
- const secondFrame = deliveredInvocation(second)!
- await sendEventResult(first, firstFrame, { kind: 'next' })
- expect(pending.resolve).not.toHaveBeenCalled()
- await sendEventResult(second, secondFrame, { kind: 'next' })
- await expect(pending.outcome).resolves.toEqual({ kind: 'next' })
- expect(pending.resolve).toHaveBeenCalledTimes(1)
- expect(pending.reject).not.toHaveBeenCalled()
- first.socket.close()
- second.socket.close()
- await unregister()
- })
- it('delivers a pending waterfall to the first Client that connects', async () => {
- const { ctx } = await setup(true)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const agent = ctx.extend()
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === agent ? agentId('agent-late-client') : undefined,
- resolve: id => id === 'agent-late-client' ? agent : undefined,
- })
- const pending = pendingInvocation(agent, undefined, 'before-connect')
- source.push(pending.dispatch)
- await vi.waitFor(() => { expect(randomUuid).toHaveBeenCalledTimes(1) })
- const client = await openEventClient(ctx, 'events-first-client')
- await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
- const frame = deliveredInvocation(client)!
- expect(frame).toMatchObject({
- type: 'waterfall',
- event: 'fixture/approval',
- agentId: 'agent-late-client',
- request: { prompt: 'before-connect' },
- })
- await sendEventResult(client, frame, { kind: 'result', value: 'allowed' })
- await expect(pending.outcome).resolves.toEqual({ kind: 'result', value: 'allowed' })
- client.socket.close()
- await unregister()
- })
- it('replays a pending event id to a replacement Client generation', async () => {
- const { ctx } = await setup(true)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const agent = ctx.extend()
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: candidate => candidate === agent ? agentId('agent-1') : undefined,
- resolve: id => id === 'agent-1' ? agent : undefined,
- })
- const original = await openEventClient(ctx, 'events-original')
- const pending = pendingInvocation(agent)
- source.push(pending.dispatch)
- await vi.waitFor(() => { expect(deliveredInvocation(original)).toBeDefined() })
- const originalFrame = deliveredInvocation(original)!
- const closed = once(original.socket, 'close')
- original.socket.close()
- await closed
- const replacement = await openEventClient(ctx, 'events-replacement')
- await vi.waitFor(() => { expect(deliveredInvocation(replacement)).toBeDefined() })
- const replayed = deliveredInvocation(replacement)!
- expect(replayed.eventId).toBe(originalFrame.eventId)
- expect(replayed).not.toHaveProperty('deliveryId')
- await sendEventResult(replacement, replayed, {
- kind: 'result', value: 'allowed',
- })
- await expect(pending.outcome).resolves.toEqual({ kind: 'result', value: 'allowed' })
- replacement.socket.close()
- await unregister()
- })
- it('cancels pending deliveries when the Host signal or Context ends', async () => {
- const { ctx } = await setup(true)
- const source = new RemoteEventSourceProbe()
- const unregister = ctx.typertGateway.registerRemoteEvents(source.source)
- const signalAgent = ctx.extend()
- const contextFiber = ctx.plugin(() => {})
- await contextFiber
- const contextAgent = contextFiber.ctx
- ctx.typert.contexts.registerHost('agent', {
- wire: 'agentId',
- wireTypeSymbol: '@fixture#AgentId',
- identity: (candidate) => {
- if (candidate === signalAgent) return agentId('agent-signal')
- if (candidate === contextAgent) return agentId('agent-context')
- return undefined
- },
- resolve: (id) => {
- if (id === 'agent-signal') return signalAgent
- if (id === 'agent-context') return contextAgent
- return undefined
- },
- })
- const client = await openEventClient(ctx, 'events-cancel')
- const abort = new AbortController()
- const signalPending = pendingInvocation(signalAgent, abort.signal, 'signal')
- source.push(signalPending.dispatch)
- await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
- const signalFrame = deliveredInvocation(client)!
- expect(signalFrame).toMatchObject({
- type: 'waterfall',
- agentId: 'agent-signal',
- request: { prompt: 'signal' },
- })
- const signalReason = new Error('Host caller cancelled')
- const signalOutcome = expect(signalPending.outcome).rejects.toBe(signalReason)
- abort.abort(signalReason)
- await signalOutcome
- await vi.waitFor(() => {
- expect(client.frames).toContainEqual({
- type: 'item',
- streamId: client.streamId,
- value: { type: 'cancel', eventId: signalFrame.eventId },
- })
- })
- const contextPending = pendingInvocation(contextAgent, undefined, 'context')
- source.push(contextPending.dispatch)
- let contextFrame: RemoteEventInvocationFrame | undefined
- await vi.waitFor(() => {
- contextFrame = client.frames
- .filter(frame => frame.type === 'item' && frame.streamId === client.streamId)
- .map(frame => frame.value)
- .find(value => typeof value === 'object'
- && value !== null
- && Reflect.get(value, 'event') === 'fixture/approval'
- && Reflect.get(value, 'eventId') !== signalFrame.eventId) as RemoteEventInvocationFrame | undefined
- expect(contextFrame).toBeDefined()
- })
- const contextOutcome = expect(contextPending.outcome).rejects.toThrow('Context "agent" was released')
- await contextFiber.dispose()
- await contextOutcome
- await vi.waitFor(() => {
- expect(client.frames).toContainEqual({
- type: 'item',
- streamId: client.streamId,
- value: { type: 'cancel', eventId: contextFrame!.eventId },
- })
- })
- client.socket.close()
- await unregister()
- })
- it('validates the internal Remote event request and reports an absent source', async () => {
- const { ctx } = await setup(true)
- const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`)
- await once(socket, 'open')
- const frames: Record<string, unknown>[] = []
- socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
- sendOpen(socket, 'missing', '$events', {})
- await vi.waitFor(() => {
- expect(frames.find(frame => frame.streamId === 'missing')?.type).toBe('error')
- expect(streamErrorMessage(frames, 'missing')).toContain('source is unavailable')
- })
- let sourceCalls = 0
- const unregister = ctx.typertGateway.registerRemoteEvents(() => {
- sourceCalls += 1
- return (async function *(): AsyncIterable<never> {})()
- })
- const invalidPayloads: readonly unknown[] = [
- null,
- [],
- {},
- { other: {} },
- { args: null },
- { args: [] },
- { args: { extra: true } },
- ]
- invalidPayloads.forEach((payload, index) => {
- socket.send(JSON.stringify({
- type: 'open', streamId: `invalid-${String(index)}`, endpoint: '$events', payload,
- }))
- })
- await vi.waitFor(() => {
- expect(frames.filter(frame => String(frame.streamId).startsWith('invalid-'))).toHaveLength(invalidPayloads.length)
- })
- for (const [index] of invalidPayloads.entries()) {
- const streamId = `invalid-${String(index)}`
- expect(frames.find(frame => frame.streamId === streamId)?.type).toBe('error')
- expect(streamErrorMessage(frames, streamId)).toContain('requires an empty args object')
- }
- expect(sourceCalls).toBe(1)
- await unregister()
- socket.close()
- })
- it('applies Connection trusted-host policy before accepting the Gateway socket', async () => {
- const { ctx } = await setup(true)
- const socket = new WebSocket(
- `ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`,
- { headers: { host: 'untrusted.example' } },
- )
- socket.on('error', () => {})
- const responseEvent: unknown[] = await once(socket, 'unexpected-response')
- const request = responseEvent[0]
- const response = responseEvent[1]
- const rejected = response as { statusCode?: number; resume(): void }
- expect(rejected.statusCode).toBe(403)
- rejected.resume()
- ;(request as { abort(): void }).abort()
- })
- })
- async function setup(transport: boolean): Promise<{ readonly ctx: Context; readonly service: FeedService }> {
- const ctx = new Context()
- roots.push(ctx)
- if (transport) {
- await ctx.plugin(WebServer, { host: '127.0.0.1', port: 0 })
- }
- await ctx.plugin(TypertRegistry)
- await ctx.plugin(TypertGatewayService)
- if (transport) {
- await ctx.plugin({ inject: [...connectionInject], apply: applyConnection })
- }
- await ctx.plugin(FeedService)
- ctx.typert.register({
- package: '@fixture/feed',
- face: 'host',
- schemas: [],
- model: { services: [], events: [], objects: [] },
- invocations: descriptors(),
- })
- const receiver = ctx.get('feed') as unknown as FeedService & { [symbols.original]?: FeedService }
- return { ctx, service: receiver[symbols.original] ?? receiver }
- }
- function descriptors(): InvocationDescriptor[] {
- const label = {
- name: 'label',
- wire: 'label',
- source: 'json' as const,
- codec: { mode: 'strict' as const, typeSymbol: '@fixture/feed#Label', schema: z.string() },
- }
- const stream = (method: string, parameters: InvocationDescriptor['parameters'], schema: z.ZodType): InvocationDescriptor => ({
- id: `@fixture/feed#feed/${method}`,
- service: 'feed',
- namespace: 'feed',
- method,
- mode: 'stream',
- invocation: { kind: 'direct' },
- parameters,
- result: { mode: 'strict', typeSymbol: '@fixture/feed#Item', schema },
- })
- return [
- { ...stream('follow', [label], z.string()), cancellation: { parameter: 'signal' } },
- stream('sync', [label], z.string()),
- stream('invalid', [], z.string()),
- stream('nonJson', [], z.unknown()),
- stream('missing', [], z.string()),
- { ...stream('abortBeforeOpen', [], z.string()), cancellation: { parameter: 'signal' } },
- stream('reject', [], z.string()),
- stream('rejectWithNonJsonDetails', [], z.string()),
- {
- id: '@fixture/feed#feed/unary',
- service: 'feed',
- namespace: 'feed',
- method: 'unary',
- invocation: { kind: 'direct' },
- parameters: [label],
- result: { mode: 'strict', typeSymbol: '@fixture/feed#Item', schema: z.string() },
- },
- ]
- }
- interface RemoteEventTestClient {
- readonly socket: WebSocket
- readonly frames: Record<string, unknown>[]
- readonly streamId: string
- readonly clientId: RemoteEventClientId
- readonly origin: string
- }
- async function openEventClient(ctx: Context, streamId: string): Promise<RemoteEventTestClient> {
- const origin = `http://127.0.0.1:${String(ctx.webServer.port)}`
- const socket = new WebSocket(`${origin.replace('http:', 'ws:')}/api/remote.mux`)
- await once(socket, 'open')
- const frames: Record<string, unknown>[] = []
- socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
- sendOpen(socket, streamId, '$events', {})
- let clientId: RemoteEventClientId | undefined
- await vi.waitFor(() => {
- const ready = frames.find(frame => frame.type === 'item'
- && frame.streamId === streamId
- && typeof frame.value === 'object'
- && frame.value !== null
- && Reflect.get(frame.value, 'type') === 'ready')
- const candidate: unknown = ready === undefined ? undefined : Reflect.get(ready.value as object, 'clientId')
- expect(typeof candidate).toBe('string')
- if (typeof candidate === 'string') clientId = candidate as RemoteEventClientId
- })
- if (clientId === undefined) throw new Error('Remote event stream omitted its Client id')
- return { socket, frames, streamId, clientId, origin }
- }
- function deliveredInvocation(client: RemoteEventTestClient): RemoteEventInvocationFrame | undefined {
- for (const frame of client.frames) {
- if (frame.type !== 'item' || frame.streamId !== client.streamId) continue
- const value = frame.value
- if (typeof value !== 'object' || value === null || !Object.hasOwn(value, 'eventId')) continue
- return value as RemoteEventInvocationFrame
- }
- return undefined
- }
- async function sendEventResult(
- client: RemoteEventTestClient,
- frame: RemoteEventInvocationFrame,
- outcome:
- | { readonly kind: 'next' }
- | { readonly kind: 'result'; readonly value?: unknown }
- | {
- readonly kind: 'rejected'
- readonly error: {
- readonly name: string
- readonly message: string
- readonly code?: string
- readonly details?: unknown
- }
- },
- ): Promise<void> {
- const rpcId = `remote-event-result-${client.streamId}`
- const response = await fetch(`${client.origin}/api/$events/result`, {
- method: 'POST',
- headers: { 'content-type': 'application/json' },
- body: JSON.stringify({
- type: 'client-request',
- rpcId,
- method: '$events/result',
- payload: {
- args: { clientId: client.clientId, eventId: frame.eventId, outcome },
- },
- }),
- })
- expect(response.status).toBe(200)
- const body = await response.json() as { readonly result?: { readonly ok?: boolean; readonly error?: { message?: string } } }
- if (body.result?.ok !== true) {
- throw new Error(body.result?.error?.message ?? 'Remote event result failed')
- }
- }
- function sendOpen(socket: WebSocket, streamId: string, endpoint: string, args: object): void {
- socket.send(JSON.stringify({ type: 'open', streamId, endpoint, payload: { args } }))
- }
- function rawText(data: RawData): string {
- if (Array.isArray(data)) return Buffer.concat(data).toString('utf8')
- if (data instanceof ArrayBuffer) return Buffer.from(data).toString('utf8')
- return Buffer.from(data).toString('utf8')
- }
- function streamErrorMessage(frames: readonly Record<string, unknown>[], streamId: string): string | undefined {
- const error = frames.find(frame => frame.streamId === streamId)?.error
- if (typeof error !== 'object' || error === null) return undefined
- const message = Reflect.get(error, 'message') as unknown
- return typeof message === 'string' ? message : undefined
- }
- async function collect(source: AsyncIterable<unknown>): Promise<unknown[]> {
- const values: unknown[] = []
- for await (const value of source) values.push(value)
- return values
- }
|