gateway-stream.host.spec.ts 42 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091
  1. import { randomUUID } from 'node:crypto'
  2. import { once } from 'node:events'
  3. import { afterEach, describe, expect, it, vi } from 'vitest'
  4. import WebSocket, { type RawData } from 'ws'
  5. import { Context, Service, symbols } from '@deepseek-ai/cordis'
  6. import { apply as applyConnection, inject as connectionInject } from '@deepseek-ai/dsh-client-connection'
  7. import WebServer from '@deepseek-ai/dsh-host-webserver'
  8. import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
  9. import {
  10. bindTypertRemote,
  11. Remote,
  12. type InvocationDescriptor,
  13. type TypertContextMap,
  14. type TypertContextWire,
  15. RemoteError,
  16. } from '@deepseek-ai/dsh-typert-protocol'
  17. import TypertRegistry from '@deepseek-ai/dsh-typert-registry'
  18. declare module '@deepseek-ai/dsh-typert-protocol' {
  19. interface RemoteErrorDetailsMap {
  20. 'fixture/rejected': { readonly retryable: boolean }
  21. 'fixture/broken': { readonly count: bigint }
  22. }
  23. }
  24. import { provideBrowserCredentials } from './browser-credentials.ts'
  25. import TypertGatewayService, {
  26. TypertGatewayError,
  27. type Config as GatewayConfig,
  28. type TypertRemoteEventDispatch,
  29. type TypertRemoteEventInvocation,
  30. type TypertRemoteEventOutcome,
  31. } from '@deepseek-ai/dsh-api-gateway'
  32. import { z } from 'zod'
  33. import type {
  34. RemoteEventClientId,
  35. RemoteEventInvocationFrame,
  36. } from '../src/stream-protocol.ts'
  37. vi.mock('node:crypto', async (importOriginal) => {
  38. const actual = await importOriginal<typeof import('node:crypto')>()
  39. return { ...actual, randomUUID: vi.fn(actual.randomUUID) }
  40. })
  41. const randomUuid = vi.mocked(randomUUID)
  42. const browserCookies = new WeakMap<Context, string>()
  43. const REMOTE_HOST = { home: '/home/fixture' } as const
  44. type AgentWireId = TypertContextWire<TypertContextMap['agent']>
  45. const agentId = (value: string): AgentWireId => value as AgentWireId
  46. /** Exchange this test Host's process token for its WebSocket/HTTP Cookie header. */
  47. function browserCookie(ctx: Context): string {
  48. const existing = browserCookies.get(ctx)
  49. if (existing !== undefined) return existing
  50. const origin = `http://127.0.0.1:${String(ctx.webServer.port)}`
  51. const target = new URL(ctx.connection.authenticatedUrl(origin))
  52. let setCookie: string | undefined
  53. ctx.connection.authorizeIndex({
  54. method: 'GET',
  55. url: `${target.pathname}${target.search}`,
  56. headers: { host: target.host },
  57. }, {
  58. writeHead(_status, headers) { setCookie = headers?.['set-cookie'] },
  59. end() {},
  60. })
  61. if (setCookie === undefined) throw new Error('gateway stream fixture did not receive a browser cookie')
  62. const cookie = setCookie.split(';', 1)[0]!
  63. browserCookies.set(ctx, cookie)
  64. return cookie
  65. }
  66. class FeedService extends Service {
  67. readonly typertRemote = bindTypertRemote(this, 'feed')
  68. readonly signals: AbortSignal[] = []
  69. returns = 0
  70. constructor(ctx: Context) {
  71. super(ctx, 'feed')
  72. }
  73. @Remote({ mode: 'stream' })
  74. async *follow(label: string, signal: AbortSignal): AsyncIterable<string> {
  75. this.signals.push(signal)
  76. try {
  77. yield `${label}:ready`
  78. await new Promise<void>((resolve) => {
  79. if (signal.aborted) resolve()
  80. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  81. })
  82. } finally {
  83. this.returns += 1
  84. }
  85. }
  86. @Remote({ mode: 'stream' })
  87. *sync(label: string): Iterable<string> {
  88. yield `${label}:one`
  89. yield `${label}:two`
  90. }
  91. @Remote({ mode: 'stream' })
  92. *invalid(): Iterable<string> {
  93. yield 42 as unknown as string
  94. }
  95. @Remote({ mode: 'stream' })
  96. *nonJson(): Iterable<unknown> {
  97. yield 1n
  98. }
  99. @Remote({ mode: 'stream' })
  100. missing(): Iterable<string> {
  101. return null as unknown as Iterable<string>
  102. }
  103. @Remote({ mode: 'stream' })
  104. *src(label: string): Iterable<string> {
  105. yield `${label}:src`
  106. }
  107. @Remote({ mode: 'stream' })
  108. abortBeforeOpen(signal: AbortSignal): Iterable<string> {
  109. if (signal.aborted) throw new Error('fixture observed pre-open cancellation')
  110. return []
  111. }
  112. @Remote({ mode: 'stream' })
  113. reject(): Iterable<string> {
  114. throw new RemoteError('fixture/rejected', 'fixture rejected the stream', { retryable: false })
  115. }
  116. @Remote({ mode: 'stream' })
  117. rejectWithNonJsonDetails(): Iterable<string> {
  118. throw new RemoteError('fixture/broken', 'fixture emitted invalid details', { count: 1n })
  119. }
  120. unary(label: string): string {
  121. return label
  122. }
  123. }
  124. const roots: Context[] = []
  125. class RemoteEventSourceProbe {
  126. readonly source = (signal: AbortSignal): AsyncIterable<TypertRemoteEventDispatch> => {
  127. this.signal = signal
  128. return this.iterate(signal)
  129. }
  130. signal: AbortSignal | undefined
  131. private readonly dispatches: TypertRemoteEventDispatch[] = []
  132. private wake: (() => void) | undefined
  133. push(dispatch: TypertRemoteEventDispatch): void {
  134. this.dispatches.push(dispatch)
  135. this.wake?.()
  136. this.wake = undefined
  137. }
  138. private async *iterate(signal: AbortSignal): AsyncGenerator<TypertRemoteEventDispatch> {
  139. const aborted = (): void => {
  140. this.wake?.()
  141. this.wake = undefined
  142. }
  143. signal.addEventListener('abort', aborted, { once: true })
  144. try {
  145. while (!signal.aborted) {
  146. while (this.dispatches.length > 0) {
  147. yield this.dispatches.shift() as TypertRemoteEventDispatch
  148. }
  149. if (signal.aborted) return
  150. await new Promise<void>((resolve) => { this.wake = resolve })
  151. this.wake = undefined
  152. }
  153. } finally {
  154. signal.removeEventListener('abort', aborted)
  155. }
  156. }
  157. }
  158. interface PendingInvocationProbe {
  159. readonly dispatch: TypertRemoteEventInvocation
  160. readonly outcome: Promise<TypertRemoteEventOutcome>
  161. readonly resolve: (outcome: TypertRemoteEventOutcome) => void
  162. readonly reject: (reason: unknown) => void
  163. }
  164. function pendingInvocation(
  165. context: Context,
  166. signal?: AbortSignal,
  167. prompt = 'ship',
  168. identity: unknown = agentId('agent-1'),
  169. ): PendingInvocationProbe {
  170. const subject = { ctx: context }
  171. const settled = Promise.withResolvers<TypertRemoteEventOutcome>()
  172. const resolve = vi.fn((outcome: TypertRemoteEventOutcome) => {
  173. settled.resolve(outcome)
  174. })
  175. const reject = vi.fn((reason: unknown) => {
  176. settled.reject(reason)
  177. })
  178. return {
  179. dispatch: {
  180. event: 'fixture/approval',
  181. request: { prompt, agent: subject, ...(signal === undefined ? {} : { signal }) },
  182. context: { value: context, subject, agentId: identity as string },
  183. resolve,
  184. reject,
  185. },
  186. outcome: settled.promise,
  187. resolve,
  188. reject,
  189. }
  190. }
  191. afterEach(async () => {
  192. randomUuid.mockClear()
  193. await Promise.all(roots.splice(0).map(ctx => ctx.fiber.dispose()))
  194. })
  195. describe('Typert Remote streams', () => {
  196. it('validates the WebSocket heartbeat timer range', () => {
  197. expect(TypertGatewayService.Config({})).toEqual({ websocketHeartbeatIntervalMs: 2_000 })
  198. expect(TypertGatewayService.Config({ websocketHeartbeatIntervalMs: MAX_TIMER_DELAY_MS }))
  199. .toEqual({ websocketHeartbeatIntervalMs: MAX_TIMER_DELAY_MS })
  200. for (const websocketHeartbeatIntervalMs of [0, 1.5, MAX_TIMER_DELAY_MS + 1]) {
  201. expect(() => TypertGatewayService.Config({ websocketHeartbeatIntervalMs })).toThrow()
  202. }
  203. })
  204. it('opens decoded carrier payloads through the in-process wire adapter', async () => {
  205. const { ctx } = await setup(false)
  206. const source = await ctx.typertGateway.wireStream.open(
  207. 'feed/sync',
  208. { args: { label: 'wire' } },
  209. new AbortController().signal,
  210. )
  211. await expect(collect(source)).resolves.toEqual(['wire:one', 'wire:two'])
  212. })
  213. it('passes Iterable and AsyncIterable items through and returns the iterator on cancellation', async () => {
  214. const { ctx, service } = await setup(false)
  215. const abort = new AbortController()
  216. const source = await ctx.typertGateway.stream({
  217. namespace: 'feed',
  218. method: 'follow',
  219. args: { label: 'a' },
  220. signal: abort.signal,
  221. })
  222. const iterator = source[Symbol.asyncIterator]()
  223. await expect(iterator.next()).resolves.toEqual({ done: false, value: 'a:ready' })
  224. const pending = iterator.next()
  225. abort.abort(new Error('fixture cancellation'))
  226. await expect(pending).rejects.toThrow('Remote invocation "feed/follow" was aborted')
  227. expect(service.signals).toEqual([abort.signal])
  228. expect(service.returns).toBe(1)
  229. await expect(collect(await ctx.typertGateway.stream({
  230. namespace: 'feed', method: 'sync', args: { label: 'b' },
  231. }))).resolves.toEqual(['b:one', 'b:two'])
  232. await expect(collect(await ctx.typertGateway.stream({
  233. namespace: 'feed', method: 'invalid', args: {},
  234. }))).resolves.toEqual([42])
  235. await expect(collect(await ctx.typertGateway.stream({
  236. namespace: 'feed', method: 'nonJson', args: {},
  237. }))).resolves.toEqual([1n])
  238. await expect(ctx.typertGateway.stream({
  239. namespace: 'feed', method: 'missing', args: {},
  240. })).rejects.toMatchObject({ code: 'gateway/result-invalid' })
  241. await expect(collect(await ctx.typertGateway.stream({
  242. namespace: 'feed', method: 'src', args: { label: 'c' },
  243. }))).resolves.toEqual(['c:src'])
  244. const abortedBeforeOpen = new AbortController()
  245. abortedBeforeOpen.abort(new Error('cancelled before open'))
  246. await expect(ctx.typertGateway.stream({
  247. namespace: 'feed', method: 'abortBeforeOpen', args: {}, signal: abortedBeforeOpen.signal,
  248. })).rejects.toThrow('Remote invocation "feed/abortBeforeOpen" was aborted')
  249. const abortedBeforeIteration = new AbortController()
  250. abortedBeforeIteration.abort(new Error('cancelled before iteration'))
  251. const preCancelled = await ctx.typertGateway.stream({
  252. namespace: 'feed', method: 'sync', args: { label: 'ignored' }, signal: abortedBeforeIteration.signal,
  253. })
  254. await expect(collect(preCancelled)).rejects.toThrow('Remote invocation "feed/sync" was aborted')
  255. })
  256. it('keeps unary and stream invocation modes distinct', async () => {
  257. const { ctx } = await setup(false)
  258. await expect(ctx.typertGateway.invoke({
  259. namespace: 'feed', method: 'sync', args: { label: 'a' },
  260. })).rejects.toMatchObject({ code: 'gateway/signature-invalid' } satisfies Partial<TypertGatewayError>)
  261. await expect(ctx.typertGateway.stream({
  262. namespace: 'feed', method: 'unary', args: { label: 'a' },
  263. })).rejects.toMatchObject({ code: 'gateway/signature-invalid' } satisfies Partial<TypertGatewayError>)
  264. })
  265. it('uses the configured WebSocket heartbeat interval', { timeout: 1_000 }, async () => {
  266. const { ctx } = await setup(true, { websocketHeartbeatIntervalMs: 20 })
  267. const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`, {
  268. headers: { cookie: browserCookie(ctx) },
  269. })
  270. const ping = once(socket, 'ping')
  271. await once(socket, 'open')
  272. expect((await ping)[0]).toEqual(Buffer.alloc(0))
  273. socket.close()
  274. await once(socket, 'close')
  275. })
  276. it('multiplexes independent streams over one WebSocket and propagates cancellation', async () => {
  277. const { ctx, service } = await setup(true)
  278. const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`, {
  279. headers: { cookie: browserCookie(ctx) },
  280. })
  281. await once(socket, 'open')
  282. const frames: Record<string, unknown>[] = []
  283. socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
  284. sendOpen(socket, 'a', 'feed/follow', { label: 'a' })
  285. sendOpen(socket, 'b', 'feed/follow', { label: 'b' })
  286. await vi.waitFor(() => {
  287. expect(frames).toEqual(expect.arrayContaining([
  288. { type: 'item', streamId: 'a', value: 'a:ready' },
  289. { type: 'item', streamId: 'b', value: 'b:ready' },
  290. ]))
  291. })
  292. expect(service.signals.map(signal => signal.aborted)).toEqual([false, false])
  293. expect(service.returns).toBe(0)
  294. socket.send(JSON.stringify({ type: 'cancel', streamId: 'a' }))
  295. await vi.waitFor(() => { expect(service.returns).toBe(1) })
  296. expect(service.signals[0]?.aborted).toBe(true)
  297. expect(service.signals[1]?.aborted).toBe(false)
  298. sendOpen(socket, 'sync', 'feed/sync', { label: 's' })
  299. sendOpen(socket, 'invalid', 'feed/invalid', {})
  300. sendOpen(socket, 'non-json', 'feed/nonJson', {})
  301. sendOpen(socket, 'rejected', 'feed/reject', {})
  302. await vi.waitFor(() => {
  303. expect(frames.filter(frame => frame.streamId === 'sync')).toEqual([
  304. { type: 'item', streamId: 'sync', value: 's:one' },
  305. { type: 'item', streamId: 'sync', value: 's:two' },
  306. { type: 'end', streamId: 'sync' },
  307. ])
  308. expect(frames.filter(frame => frame.streamId === 'invalid')).toEqual([
  309. { type: 'item', streamId: 'invalid', value: 42 },
  310. { type: 'end', streamId: 'invalid' },
  311. ])
  312. expect(frames.find(frame => frame.streamId === 'non-json')).toMatchObject({
  313. type: 'error', error: { code: 'gateway/internal' },
  314. })
  315. expect(frames.find(frame => frame.streamId === 'rejected')).toEqual({
  316. type: 'error',
  317. streamId: 'rejected',
  318. error: {
  319. code: 'fixture/rejected',
  320. message: 'fixture rejected the stream',
  321. details: { retryable: false },
  322. },
  323. })
  324. })
  325. const closed = once(socket, 'close')
  326. sendOpen(socket, 'broken-error', 'feed/rejectWithNonJsonDetails', {})
  327. const closeEvent = await closed
  328. expect(closeEvent[0]).toBe(1011)
  329. expect(String(closeEvent[1])).toBe('Remote stream failure could not be delivered')
  330. await vi.waitFor(() => { expect(service.returns).toBe(2) })
  331. expect(service.signals[1]?.aborted).toBe(true)
  332. })
  333. it('carries the registered Remote event source and withdraws its active stream', async () => {
  334. const { ctx } = await setup(true)
  335. let sourceSignal: AbortSignal | undefined
  336. const sourceClosed = vi.fn()
  337. const publish = Promise.withResolvers<undefined>()
  338. const source = (signal: AbortSignal): AsyncIterable<{ event: string; args: readonly unknown[] }> => {
  339. sourceSignal = signal
  340. return (async function *() {
  341. try {
  342. await publish.promise
  343. yield { event: 'fixture/changed', args: ['settings'] }
  344. await new Promise<void>((resolve) => {
  345. if (signal.aborted) resolve()
  346. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  347. })
  348. } finally {
  349. sourceClosed()
  350. }
  351. })()
  352. }
  353. const unregister = ctx.typertGateway.registerRemoteEvents(source, REMOTE_HOST)
  354. expect(() => { ctx.typertGateway.registerRemoteEvents(source, REMOTE_HOST) })
  355. .toThrow('forwarded Remote event source is already registered')
  356. const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`, {
  357. headers: { cookie: browserCookie(ctx) },
  358. })
  359. await once(socket, 'open')
  360. const frames: Record<string, unknown>[] = []
  361. socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
  362. sendOpen(socket, 'events', '$events', {})
  363. await vi.waitFor(() => {
  364. const eventFrames = frames.filter(frame => frame.streamId === 'events')
  365. expect(eventFrames).toHaveLength(1)
  366. expect(eventFrames[0]).toMatchObject({
  367. type: 'item', streamId: 'events', value: { type: 'ready', host: REMOTE_HOST },
  368. })
  369. expect(typeof Reflect.get(eventFrames[0]!.value as object, 'clientId')).toBe('string')
  370. })
  371. publish.resolve(undefined)
  372. await vi.waitFor(() => {
  373. const eventFrames = frames.filter(frame => frame.streamId === 'events').slice(0, 2)
  374. expect(eventFrames).toHaveLength(2)
  375. expect(eventFrames[0]).toMatchObject({
  376. type: 'item', streamId: 'events', value: { type: 'ready', host: REMOTE_HOST },
  377. })
  378. expect(typeof Reflect.get(eventFrames[0]!.value as object, 'clientId')).toBe('string')
  379. expect(eventFrames[1]).toEqual({
  380. type: 'item', streamId: 'events', value: {
  381. type: 'emit', event: 'fixture/changed', args: ['settings'],
  382. },
  383. })
  384. })
  385. expect(sourceSignal?.aborted).toBe(false)
  386. await unregister()
  387. expect(sourceClosed).toHaveBeenCalledOnce()
  388. await vi.waitFor(() => {
  389. expect(sourceSignal?.aborted).toBe(true)
  390. expect(frames).toContainEqual({ type: 'end', streamId: 'events' })
  391. })
  392. const unregisterReplacement = ctx.typertGateway.registerRemoteEvents(source, REMOTE_HOST)
  393. await unregister()
  394. expect(() => { ctx.typertGateway.registerRemoteEvents(source, REMOTE_HOST) })
  395. .toThrow('forwarded Remote event source is already registered')
  396. await unregisterReplacement()
  397. socket.close()
  398. })
  399. it('rejects a scoped dispatch yielded after its Remote event source is withdrawn', async () => {
  400. const { ctx } = await setup(false)
  401. const publish = Promise.withResolvers<undefined>()
  402. const agent = ctx.extend()
  403. const pending = pendingInvocation(agent)
  404. const source = (): AsyncIterable<TypertRemoteEventDispatch> => (async function* () {
  405. await publish.promise
  406. yield pending.dispatch
  407. })()
  408. const unregister = ctx.typertGateway.registerRemoteEvents(source, REMOTE_HOST)
  409. const rejected = expect(pending.outcome).rejects.toThrow(
  410. 'forwarded Remote event source was removed',
  411. )
  412. publish.resolve(undefined)
  413. await unregister()
  414. await rejected
  415. expect(pending.reject).toHaveBeenCalledTimes(1)
  416. expect(pending.resolve).not.toHaveBeenCalled()
  417. })
  418. it('cancels a pending waterfall when its source rejects during removal', async () => {
  419. const { ctx } = await setup(true)
  420. const agent = ctx.extend()
  421. const pending = pendingInvocation(agent, undefined, 'ship', agentId('agent-removal'))
  422. const rejected = expect(pending.outcome).rejects.toThrow(
  423. 'forwarded Remote event source was removed',
  424. )
  425. const unregister = ctx.typertGateway.registerRemoteEvents(signal => (async function* () {
  426. yield pending.dispatch
  427. await new Promise<void>((resolve) => {
  428. if (signal.aborted) resolve()
  429. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  430. })
  431. throw new Error('fixture source rejected during removal')
  432. })(), REMOTE_HOST)
  433. const client = await openEventClient(ctx, 'events-removal')
  434. await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
  435. await unregister()
  436. await rejected
  437. expect(pending.reject).toHaveBeenCalledTimes(1)
  438. expect(pending.resolve).not.toHaveBeenCalled()
  439. await vi.waitFor(() => {
  440. expect(client.frames).toContainEqual({ type: 'end', streamId: client.streamId })
  441. })
  442. client.socket.close()
  443. })
  444. it('rejects malformed scoped invocations and delegates a released Context', async () => {
  445. const { ctx } = await setup(false)
  446. const source = new RemoteEventSourceProbe()
  447. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  448. for (const event of [42, ''] as const) {
  449. const invalidName = pendingInvocation(ctx)
  450. const rejected = expect(invalidName.outcome).rejects.toThrow(
  451. 'Remote event name must be a nonempty string',
  452. )
  453. source.push({
  454. ...invalidName.dispatch,
  455. event: event as unknown as string,
  456. })
  457. await rejected
  458. }
  459. let selected = ctx.extend()
  460. const nonJsonIdentity = pendingInvocation(selected, undefined, 'ship', 1n)
  461. const nonJsonRejected = expect(nonJsonIdentity.outcome).rejects.toThrow(
  462. 'require a non-empty Agent identity',
  463. )
  464. source.push(nonJsonIdentity.dispatch)
  465. await nonJsonRejected
  466. const invalidRequest = pendingInvocation(selected, undefined, 'ship', agentId('agent-invalid-request'))
  467. const invalidRequestRejected = expect(invalidRequest.outcome).rejects.toThrow(
  468. 'must carry its scoped Agent directly',
  469. )
  470. source.push({
  471. ...invalidRequest.dispatch,
  472. request: {},
  473. })
  474. await invalidRequestRejected
  475. const staleFiber = ctx.plugin(() => {})
  476. await staleFiber
  477. selected = staleFiber.ctx
  478. await staleFiber.dispose()
  479. const stale = pendingInvocation(selected, undefined, 'ship', agentId('agent-stale'))
  480. source.push(stale.dispatch)
  481. await expect(stale.outcome).resolves.toEqual({ kind: 'next' })
  482. expect(stale.reject).not.toHaveBeenCalled()
  483. selected = ctx.extend()
  484. const abort = new AbortController()
  485. abort.abort('fixture non-error cancellation')
  486. const cancelled = pendingInvocation(selected, abort.signal, 'ship', agentId('agent-cancelled'))
  487. const cancelledOutcome = expect(cancelled.outcome).rejects.toMatchObject({
  488. message: 'typert gateway: Remote event was cancelled',
  489. cause: 'fixture non-error cancellation',
  490. })
  491. source.push(cancelled.dispatch)
  492. await cancelledOutcome
  493. await unregister()
  494. })
  495. it('rejects notification arguments that are not lossless JSON arrays', async () => {
  496. const { ctx } = await setup(false)
  497. const frames = [
  498. { event: 'fixture/changed', args: {} },
  499. { event: 'fixture/changed', args: [1n] },
  500. ]
  501. for (const frame of frames) {
  502. let sourceSignal: AbortSignal | undefined
  503. const unregister = ctx.typertGateway.registerRemoteEvents((signal) => {
  504. sourceSignal = signal
  505. return (async function* () {
  506. yield frame as unknown as TypertRemoteEventDispatch
  507. })()
  508. }, REMOTE_HOST)
  509. await vi.waitFor(() => { expect(sourceSignal?.aborted).toBe(true) })
  510. const reason: unknown = sourceSignal?.reason
  511. if (!(reason instanceof Error)) throw new Error('Remote event source did not fail with an Error')
  512. expect(reason.message).toContain('arguments are not lossless JSON data')
  513. await unregister()
  514. }
  515. })
  516. it('retries a colliding Remote event id before publishing the second waterfall', async () => {
  517. const { ctx } = await setup(false)
  518. const source = new RemoteEventSourceProbe()
  519. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  520. const agent = ctx.extend()
  521. const firstId = '00000000-0000-4000-8000-000000000001' as ReturnType<typeof randomUUID>
  522. const secondId = '00000000-0000-4000-8000-000000000002' as ReturnType<typeof randomUUID>
  523. randomUuid.mockReturnValueOnce(firstId).mockReturnValueOnce(firstId).mockReturnValueOnce(secondId)
  524. const firstAbort = new AbortController()
  525. const secondAbort = new AbortController()
  526. const first = pendingInvocation(agent, firstAbort.signal, 'first', agentId('agent-collision'))
  527. const second = pendingInvocation(agent, secondAbort.signal, 'second', agentId('agent-collision'))
  528. source.push(first.dispatch)
  529. await vi.waitFor(() => { expect(randomUuid).toHaveBeenCalledTimes(1) })
  530. source.push(second.dispatch)
  531. await vi.waitFor(() => { expect(randomUuid).toHaveBeenCalledTimes(3) })
  532. const firstReason = new Error('cancel first collision fixture')
  533. const secondReason = new Error('cancel second collision fixture')
  534. const firstRejected = expect(first.outcome).rejects.toBe(firstReason)
  535. const secondRejected = expect(second.outcome).rejects.toBe(secondReason)
  536. firstAbort.abort(firstReason)
  537. secondAbort.abort(secondReason)
  538. await firstRejected
  539. await secondRejected
  540. await unregister()
  541. })
  542. it('retries a colliding Remote event Client id before opening the second generation', async () => {
  543. const { ctx } = await setup(true)
  544. const source = new RemoteEventSourceProbe()
  545. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  546. const firstId = '00000000-0000-4000-8000-000000000011' as ReturnType<typeof randomUUID>
  547. const secondId = '00000000-0000-4000-8000-000000000012' as ReturnType<typeof randomUUID>
  548. randomUuid.mockReturnValueOnce(firstId).mockReturnValueOnce(firstId).mockReturnValueOnce(secondId)
  549. const first = await openEventClient(ctx, 'events-client-id-a')
  550. const second = await openEventClient(ctx, 'events-client-id-b')
  551. expect(first.clientId).toBe(firstId)
  552. expect(second.clientId).toBe(secondId)
  553. expect(randomUuid).toHaveBeenCalledTimes(3)
  554. first.socket.close()
  555. second.socket.close()
  556. await unregister()
  557. })
  558. it('fans one scoped waterfall out and accepts the first Client result', async () => {
  559. const { ctx } = await setup(true)
  560. const source = new RemoteEventSourceProbe()
  561. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  562. const agent = ctx.extend()
  563. const first = await openEventClient(ctx, 'events-a')
  564. const second = await openEventClient(ctx, 'events-b')
  565. const pending = pendingInvocation(agent)
  566. source.push(pending.dispatch)
  567. await vi.waitFor(() => {
  568. expect(deliveredInvocation(first)).toBeDefined()
  569. expect(deliveredInvocation(second)).toBeDefined()
  570. })
  571. const firstFrame = deliveredInvocation(first)!
  572. const secondFrame = deliveredInvocation(second)!
  573. expect(firstFrame.eventId).toBe(secondFrame.eventId)
  574. expect(firstFrame).toMatchObject({
  575. type: 'waterfall',
  576. event: 'fixture/approval',
  577. agentId: 'agent-1',
  578. request: { prompt: 'ship' },
  579. })
  580. expect(firstFrame).not.toHaveProperty('deliveryId')
  581. expect(secondFrame).not.toHaveProperty('deliveryId')
  582. await sendEventResult(second, secondFrame, {
  583. kind: 'result', value: 'allowed',
  584. })
  585. await expect(pending.outcome).resolves.toEqual({ kind: 'result', value: 'allowed' })
  586. await vi.waitFor(() => {
  587. expect(first.frames).toContainEqual({
  588. type: 'item',
  589. streamId: first.streamId,
  590. value: { type: 'cancel', eventId: firstFrame.eventId },
  591. })
  592. })
  593. await sendEventResult(first, firstFrame, {
  594. kind: 'result', value: 'rejected',
  595. })
  596. expect(pending.resolve).toHaveBeenCalledTimes(1)
  597. expect(pending.reject).not.toHaveBeenCalled()
  598. first.socket.close()
  599. second.socket.close()
  600. await unregister()
  601. })
  602. it('rejects the Host waterfall with the first Client listener rejection', async () => {
  603. const { ctx } = await setup(true)
  604. const source = new RemoteEventSourceProbe()
  605. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  606. const agent = ctx.extend()
  607. const client = await openEventClient(ctx, 'events-rejected')
  608. const pending = pendingInvocation(agent, undefined, 'ship', agentId('agent-rejected'))
  609. source.push(pending.dispatch)
  610. await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
  611. const frame = deliveredInvocation(client)!
  612. const rejected = expect(pending.outcome).rejects.toMatchObject({
  613. name: 'UserQuestionError',
  614. message: 'the user cancelled ask_user_question',
  615. code: 'ASK_CANCELLED',
  616. details: { questionId: 'question-1' },
  617. })
  618. await sendEventResult(client, frame, {
  619. kind: 'rejected',
  620. error: {
  621. name: 'UserQuestionError',
  622. message: 'the user cancelled ask_user_question',
  623. code: 'ASK_CANCELLED',
  624. details: { questionId: 'question-1' },
  625. },
  626. })
  627. await rejected
  628. expect(pending.reject).toHaveBeenCalledTimes(1)
  629. expect(pending.resolve).not.toHaveBeenCalled()
  630. client.socket.close()
  631. await unregister()
  632. })
  633. it('delegates to the Host only after every active Client returns next', async () => {
  634. const { ctx } = await setup(true)
  635. const source = new RemoteEventSourceProbe()
  636. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  637. const agent = ctx.extend()
  638. const first = await openEventClient(ctx, 'events-next-a')
  639. const second = await openEventClient(ctx, 'events-next-b')
  640. const pending = pendingInvocation(agent)
  641. source.push(pending.dispatch)
  642. await vi.waitFor(() => {
  643. expect(deliveredInvocation(first)).toBeDefined()
  644. expect(deliveredInvocation(second)).toBeDefined()
  645. })
  646. const firstFrame = deliveredInvocation(first)!
  647. const secondFrame = deliveredInvocation(second)!
  648. await sendEventResult(first, firstFrame, { kind: 'next' })
  649. expect(pending.resolve).not.toHaveBeenCalled()
  650. await sendEventResult(second, secondFrame, { kind: 'next' })
  651. await expect(pending.outcome).resolves.toEqual({ kind: 'next' })
  652. expect(pending.resolve).toHaveBeenCalledTimes(1)
  653. expect(pending.reject).not.toHaveBeenCalled()
  654. first.socket.close()
  655. second.socket.close()
  656. await unregister()
  657. })
  658. it('delivers a pending waterfall to the first Client that connects', async () => {
  659. const { ctx } = await setup(true)
  660. const source = new RemoteEventSourceProbe()
  661. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  662. const agent = ctx.extend()
  663. const pending = pendingInvocation(agent, undefined, 'before-connect', agentId('agent-late-client'))
  664. source.push(pending.dispatch)
  665. await vi.waitFor(() => { expect(randomUuid).toHaveBeenCalledTimes(1) })
  666. const client = await openEventClient(ctx, 'events-first-client')
  667. await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
  668. const frame = deliveredInvocation(client)!
  669. expect(frame).toMatchObject({
  670. type: 'waterfall',
  671. event: 'fixture/approval',
  672. agentId: 'agent-late-client',
  673. request: { prompt: 'before-connect' },
  674. })
  675. await sendEventResult(client, frame, { kind: 'result', value: 'allowed' })
  676. await expect(pending.outcome).resolves.toEqual({ kind: 'result', value: 'allowed' })
  677. client.socket.close()
  678. await unregister()
  679. })
  680. it('replays a pending event id to a replacement Client generation', async () => {
  681. const { ctx } = await setup(true)
  682. const source = new RemoteEventSourceProbe()
  683. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  684. const agent = ctx.extend()
  685. const original = await openEventClient(ctx, 'events-original')
  686. const pending = pendingInvocation(agent)
  687. source.push(pending.dispatch)
  688. await vi.waitFor(() => { expect(deliveredInvocation(original)).toBeDefined() })
  689. const originalFrame = deliveredInvocation(original)!
  690. const closed = once(original.socket, 'close')
  691. original.socket.close()
  692. await closed
  693. const replacement = await openEventClient(ctx, 'events-replacement')
  694. await vi.waitFor(() => { expect(deliveredInvocation(replacement)).toBeDefined() })
  695. const replayed = deliveredInvocation(replacement)!
  696. expect(replayed.eventId).toBe(originalFrame.eventId)
  697. expect(replayed).not.toHaveProperty('deliveryId')
  698. await sendEventResult(replacement, replayed, {
  699. kind: 'result', value: 'allowed',
  700. })
  701. await expect(pending.outcome).resolves.toEqual({ kind: 'result', value: 'allowed' })
  702. replacement.socket.close()
  703. await unregister()
  704. })
  705. it('cancels pending deliveries when the Host signal or Context ends', async () => {
  706. const { ctx } = await setup(true)
  707. const source = new RemoteEventSourceProbe()
  708. const unregister = ctx.typertGateway.registerRemoteEvents(source.source, REMOTE_HOST)
  709. const signalAgent = ctx.extend()
  710. const contextFiber = ctx.plugin(() => {})
  711. await contextFiber
  712. const contextAgent = contextFiber.ctx
  713. const client = await openEventClient(ctx, 'events-cancel')
  714. const abort = new AbortController()
  715. const signalPending = pendingInvocation(signalAgent, abort.signal, 'signal', agentId('agent-signal'))
  716. source.push(signalPending.dispatch)
  717. await vi.waitFor(() => { expect(deliveredInvocation(client)).toBeDefined() })
  718. const signalFrame = deliveredInvocation(client)!
  719. expect(signalFrame).toMatchObject({
  720. type: 'waterfall',
  721. agentId: 'agent-signal',
  722. request: { prompt: 'signal' },
  723. })
  724. const signalReason = new Error('Host caller cancelled')
  725. const signalOutcome = expect(signalPending.outcome).rejects.toBe(signalReason)
  726. abort.abort(signalReason)
  727. await signalOutcome
  728. await vi.waitFor(() => {
  729. expect(client.frames).toContainEqual({
  730. type: 'item',
  731. streamId: client.streamId,
  732. value: { type: 'cancel', eventId: signalFrame.eventId },
  733. })
  734. })
  735. const contextPending = pendingInvocation(contextAgent, undefined, 'context', agentId('agent-context'))
  736. source.push(contextPending.dispatch)
  737. let contextFrame: RemoteEventInvocationFrame | undefined
  738. await vi.waitFor(() => {
  739. contextFrame = client.frames
  740. .filter(frame => frame.type === 'item' && frame.streamId === client.streamId)
  741. .map(frame => frame.value)
  742. .find(value => typeof value === 'object'
  743. && value !== null
  744. && Reflect.get(value, 'event') === 'fixture/approval'
  745. && Reflect.get(value, 'eventId') !== signalFrame.eventId) as RemoteEventInvocationFrame | undefined
  746. expect(contextFrame).toBeDefined()
  747. })
  748. const contextOutcome = expect(contextPending.outcome).rejects.toThrow('Agent Context was released')
  749. await contextFiber.dispose()
  750. await contextOutcome
  751. await vi.waitFor(() => {
  752. expect(client.frames).toContainEqual({
  753. type: 'item',
  754. streamId: client.streamId,
  755. value: { type: 'cancel', eventId: contextFrame!.eventId },
  756. })
  757. })
  758. client.socket.close()
  759. await unregister()
  760. })
  761. it('validates the internal Remote event request and reports an absent source', async () => {
  762. const { ctx } = await setup(true)
  763. const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`, {
  764. headers: { cookie: browserCookie(ctx) },
  765. })
  766. await once(socket, 'open')
  767. const frames: Record<string, unknown>[] = []
  768. socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
  769. sendOpen(socket, 'missing', '$events', {})
  770. await vi.waitFor(() => {
  771. expect(frames.find(frame => frame.streamId === 'missing')?.type).toBe('error')
  772. expect(streamErrorMessage(frames, 'missing')).toContain('source is unavailable')
  773. })
  774. let sourceCalls = 0
  775. const unregister = ctx.typertGateway.registerRemoteEvents(() => {
  776. sourceCalls += 1
  777. return (async function *(): AsyncIterable<never> {})()
  778. }, REMOTE_HOST)
  779. const invalidPayloads: readonly unknown[] = [
  780. null,
  781. [],
  782. {},
  783. { other: {} },
  784. { args: null },
  785. { args: [] },
  786. { args: { extra: true } },
  787. ]
  788. invalidPayloads.forEach((payload, index) => {
  789. socket.send(JSON.stringify({
  790. type: 'open', streamId: `invalid-${String(index)}`, endpoint: '$events', payload,
  791. }))
  792. })
  793. await vi.waitFor(() => {
  794. expect(frames.filter(frame => String(frame.streamId).startsWith('invalid-'))).toHaveLength(invalidPayloads.length)
  795. })
  796. for (const [index] of invalidPayloads.entries()) {
  797. const streamId = `invalid-${String(index)}`
  798. expect(frames.find(frame => frame.streamId === streamId)?.type).toBe('error')
  799. expect(streamErrorMessage(frames, streamId)).toContain('requires an empty args object')
  800. }
  801. expect(sourceCalls).toBe(1)
  802. await unregister()
  803. socket.close()
  804. })
  805. it('applies Connection trusted-host policy before accepting the Gateway socket', async () => {
  806. const { ctx } = await setup(true)
  807. const socket = new WebSocket(
  808. `ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`,
  809. { headers: { host: 'untrusted.example' } },
  810. )
  811. socket.on('error', () => {})
  812. const responseEvent: unknown[] = await once(socket, 'unexpected-response')
  813. const request = responseEvent[0]
  814. const response = responseEvent[1]
  815. const rejected = response as { statusCode?: number; resume(): void }
  816. expect(rejected.statusCode).toBe(403)
  817. rejected.resume()
  818. ;(request as { abort(): void }).abort()
  819. })
  820. it('answers an unauthenticated trusted Host with 401 before opening a stream', async () => {
  821. const { ctx } = await setup(true)
  822. const socket = new WebSocket(`ws://127.0.0.1:${String(ctx.webServer.port)}/api/remote.mux`)
  823. socket.on('error', () => {})
  824. const responseEvent: unknown[] = await once(socket, 'unexpected-response')
  825. const request = responseEvent[0]
  826. const response = responseEvent[1]
  827. const rejected = response as { statusCode?: number; resume(): void }
  828. expect(rejected.statusCode).toBe(401)
  829. rejected.resume()
  830. ;(request as { abort(): void }).abort()
  831. })
  832. })
  833. async function setup(
  834. transport: boolean,
  835. gatewayConfig: GatewayConfig = {},
  836. ): Promise<{ readonly ctx: Context; readonly service: FeedService }> {
  837. const ctx = new Context()
  838. roots.push(ctx)
  839. if (transport) {
  840. await ctx.plugin(WebServer, { host: '127.0.0.1', port: 0 })
  841. provideBrowserCredentials(ctx)
  842. }
  843. await ctx.plugin(TypertRegistry)
  844. await ctx.plugin(TypertGatewayService, gatewayConfig)
  845. if (transport) {
  846. await ctx.plugin({ inject: [...connectionInject], apply: applyConnection })
  847. }
  848. await ctx.plugin(FeedService)
  849. ctx.typert.register({
  850. package: '@fixture/feed',
  851. face: 'host',
  852. schemas: [],
  853. model: { services: [], events: [], objects: [] },
  854. invocations: descriptors(),
  855. })
  856. const receiver = ctx.get('feed') as unknown as FeedService & { [symbols.original]?: FeedService }
  857. return { ctx, service: receiver[symbols.original] ?? receiver }
  858. }
  859. function descriptors(): InvocationDescriptor[] {
  860. const label = {
  861. name: 'label',
  862. wire: 'label',
  863. source: 'json' as const,
  864. codec: { mode: 'strict' as const, typeSymbol: '@fixture/feed#Label', create: () => z.string() },
  865. }
  866. const stream = (method: string, parameters: InvocationDescriptor['parameters'], schema: z.ZodType): InvocationDescriptor => ({
  867. id: `@fixture/feed#feed/${method}`,
  868. service: 'feed',
  869. namespace: 'feed',
  870. method,
  871. mode: 'stream',
  872. invocation: { kind: 'direct' },
  873. parameters,
  874. result: { mode: 'strict', typeSymbol: '@fixture/feed#Item', create: () => schema },
  875. })
  876. return [
  877. { ...stream('follow', [label], z.string()), cancellation: { parameter: 'signal' } },
  878. stream('sync', [label], z.string()),
  879. stream('invalid', [], z.string()),
  880. stream('nonJson', [], z.unknown()),
  881. stream('missing', [], z.string()),
  882. { ...stream('abortBeforeOpen', [], z.string()), cancellation: { parameter: 'signal' } },
  883. stream('reject', [], z.string()),
  884. stream('rejectWithNonJsonDetails', [], z.string()),
  885. {
  886. id: '@fixture/feed#feed/unary',
  887. service: 'feed',
  888. namespace: 'feed',
  889. method: 'unary',
  890. invocation: { kind: 'direct' },
  891. parameters: [label],
  892. result: { mode: 'strict', typeSymbol: '@fixture/feed#Item', create: () => z.string() },
  893. },
  894. ]
  895. }
  896. interface RemoteEventTestClient {
  897. readonly socket: WebSocket
  898. readonly frames: Record<string, unknown>[]
  899. readonly streamId: string
  900. readonly clientId: RemoteEventClientId
  901. readonly origin: string
  902. readonly cookie: string
  903. }
  904. async function openEventClient(ctx: Context, streamId: string): Promise<RemoteEventTestClient> {
  905. const origin = `http://127.0.0.1:${String(ctx.webServer.port)}`
  906. const cookie = browserCookie(ctx)
  907. const socket = new WebSocket(`${origin.replace('http:', 'ws:')}/api/remote.mux`, {
  908. headers: { cookie },
  909. })
  910. await once(socket, 'open')
  911. const frames: Record<string, unknown>[] = []
  912. socket.on('message', (data) => { frames.push(JSON.parse(rawText(data)) as Record<string, unknown>) })
  913. sendOpen(socket, streamId, '$events', {})
  914. let clientId: RemoteEventClientId | undefined
  915. await vi.waitFor(() => {
  916. const ready = frames.find(frame => frame.type === 'item'
  917. && frame.streamId === streamId
  918. && typeof frame.value === 'object'
  919. && frame.value !== null
  920. && Reflect.get(frame.value, 'type') === 'ready')
  921. const candidate: unknown = ready === undefined ? undefined : Reflect.get(ready.value as object, 'clientId')
  922. expect(typeof candidate).toBe('string')
  923. if (typeof candidate === 'string') clientId = candidate as RemoteEventClientId
  924. })
  925. if (clientId === undefined) throw new Error('Remote event stream omitted its Client id')
  926. return { socket, frames, streamId, clientId, origin, cookie }
  927. }
  928. function deliveredInvocation(client: RemoteEventTestClient): RemoteEventInvocationFrame | undefined {
  929. for (const frame of client.frames) {
  930. if (frame.type !== 'item' || frame.streamId !== client.streamId) continue
  931. const value = frame.value
  932. if (typeof value !== 'object' || value === null || !Object.hasOwn(value, 'eventId')) continue
  933. return value as RemoteEventInvocationFrame
  934. }
  935. return undefined
  936. }
  937. async function sendEventResult(
  938. client: RemoteEventTestClient,
  939. frame: RemoteEventInvocationFrame,
  940. outcome:
  941. | { readonly kind: 'next' }
  942. | { readonly kind: 'result'; readonly value?: unknown }
  943. | {
  944. readonly kind: 'rejected'
  945. readonly error: {
  946. readonly name: string
  947. readonly message: string
  948. readonly code?: string
  949. readonly details?: unknown
  950. }
  951. },
  952. ): Promise<void> {
  953. const rpcId = `remote-event-result-${client.streamId}`
  954. const response = await fetch(`${client.origin}/api/$events/result`, {
  955. method: 'POST',
  956. headers: { 'content-type': 'application/json', cookie: client.cookie },
  957. body: JSON.stringify({
  958. type: 'client-request',
  959. rpcId,
  960. method: '$events/result',
  961. payload: {
  962. args: { clientId: client.clientId, eventId: frame.eventId, outcome },
  963. },
  964. }),
  965. })
  966. expect(response.status).toBe(200)
  967. const body = await response.json() as { readonly result?: { readonly ok?: boolean; readonly error?: { message?: string } } }
  968. if (body.result?.ok !== true) {
  969. throw new Error(body.result?.error?.message ?? 'Remote event result failed')
  970. }
  971. }
  972. function sendOpen(socket: WebSocket, streamId: string, endpoint: string, args: object): void {
  973. socket.send(JSON.stringify({ type: 'open', streamId, endpoint, payload: { args } }))
  974. }
  975. function rawText(data: RawData): string {
  976. if (Array.isArray(data)) return Buffer.concat(data).toString('utf8')
  977. if (data instanceof ArrayBuffer) return Buffer.from(data).toString('utf8')
  978. return Buffer.from(data).toString('utf8')
  979. }
  980. function streamErrorMessage(frames: readonly Record<string, unknown>[], streamId: string): string | undefined {
  981. const error = frames.find(frame => frame.streamId === streamId)?.error
  982. if (typeof error !== 'object' || error === null) return undefined
  983. const message = Reflect.get(error, 'message') as unknown
  984. return typeof message === 'string' ? message : undefined
  985. }
  986. async function collect(source: AsyncIterable<unknown>): Promise<unknown[]> {
  987. const values: unknown[] = []
  988. for await (const value of source) values.push(value)
  989. return values
  990. }