Przeglądaj źródła

test(api-gateway): cover clientless Remote event replay

imccyu 1 miesiąc temu
rodzic
commit
7d21611391

+ 2 - 0
packages/api/gateway/src/index.ts

@@ -513,6 +513,8 @@ export class TypertGatewayService extends Service implements TypertGateway {
     result: ReturnType<typeof parseRemoteEventResult>,
   ): void {
     const pending = this.pendingRemoteEvents.get(result.eventId)
+    // Settlement and Client replacement may race the result request. Results
+    // from a completed event or a superseded delivery are idempotent no-ops.
     if (pending === undefined || !pending.deliveries.has(client)) return
     this.removeRemoteEventDelivery(pending, client)
     if (result.outcome.kind === 'result') {

+ 33 - 0
packages/api/gateway/tests/gateway-stream.host.spec.ts

@@ -714,6 +714,39 @@ describe('Typert Remote streams', () => {
     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()