gateway.client.spec.ts 92 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381
  1. import { Context, Service } from '@deepseek-ai/cordis'
  2. import type { Fiber } from '@deepseek-ai/cordis'
  3. import { describe, expect, expectTypeOf, it, vi } from 'vitest'
  4. import { z } from 'zod'
  5. import {
  6. apply as applyConnection,
  7. type ConnectionGenerationSource,
  8. type ConnectionHandle,
  9. } from '@deepseek-ai/dsh-client-connection/client'
  10. import type {
  11. InvocationDescriptor,
  12. RemoteResult,
  13. TypertContextMap,
  14. TypertContextWire,
  15. TypertContext,
  16. TypertLookup,
  17. TypertRemoteScopeApi,
  18. TypertRemoteNamespace,
  19. } from '@deepseek-ai/dsh-typert-protocol'
  20. import TypertRegistry from '@deepseek-ai/dsh-typert-registry'
  21. import type { ClientRemote } from '../src/client/index.ts'
  22. import { apply, inject, RemoteStream } from '../src/client/index.ts'
  23. import {
  24. RemoteStreamCarrierError,
  25. RemoteStreamError,
  26. RemoteStreamMuxClient,
  27. } from '../src/client/stream-client.ts'
  28. type FixtureApprovalOutcome = 'allowed' | 'unavailable'
  29. const fixtureContextTag = Symbol('fixture-context-tag')
  30. type AgentWireId = TypertContextWire<TypertContextMap['agent']>
  31. const agentId = (value: string): AgentWireId => value as AgentWireId
  32. interface FixtureAgent {
  33. readonly agentId: string
  34. }
  35. declare module '@deepseek-ai/cordis' {
  36. interface Events {
  37. /**
  38. * Test-only forwarded Host event.
  39. * @param namespace - marker payload recorded by listeners.
  40. */
  41. 'fixture/changed'(namespace: string): void
  42. /**
  43. * Test-only forwarded Host event nobody subscribes to.
  44. * @param count - marker payload never observed.
  45. */
  46. 'fixture/idle'(count: number): void
  47. /**
  48. * Test-only scoped waterfall forwarded through the existing Remote Event stream.
  49. * @param request - JSON-safe request payload.
  50. * @param next - delegates to the next Client listener or Host waterfall.
  51. * @returns the claimed or delegated outcome.
  52. */
  53. 'fixture/approval'(
  54. this: Context,
  55. request: {
  56. readonly prompt: string
  57. readonly agent: FixtureAgent
  58. readonly signal?: AbortSignal
  59. },
  60. next: () => Promise<FixtureApprovalOutcome>,
  61. ): Promise<FixtureApprovalOutcome>
  62. /**
  63. * Test-only event the Host assembly does not forward.
  64. * @param flag - marker payload never delivered.
  65. */
  66. 'fixture/unselected'(flag: boolean): void
  67. }
  68. }
  69. declare module '@deepseek-ai/dsh-typert-protocol' {
  70. interface TypertRemoteEventSelection extends
  71. Record<'fixture/changed' | 'fixture/idle' | 'fixture/approval', true> {}
  72. interface TypertContextMap {
  73. fixture: TypertContext<string>
  74. }
  75. interface TypertLookupMap {
  76. fixture: TypertLookup<FixtureAgent, string>
  77. }
  78. interface TypertRemoteMap {
  79. 'probe/create': (
  80. agentId: string,
  81. request: { readonly objective: string },
  82. signal?: AbortSignal,
  83. ) => Promise<RemoteResult<{ readonly ref: string }>>
  84. 'probe/maybe': (value: string | null | undefined) => Promise<RemoteResult<string | null | undefined>>
  85. 'probe/watch': (topic: string, signal?: AbortSignal) => AsyncIterable<string>
  86. }
  87. interface TypertRemoteScopeMap {
  88. 'fixture:probe/create': (
  89. request: { readonly objective: string },
  90. signal?: AbortSignal,
  91. ) => Promise<RemoteResult<{ readonly ref: string }>>
  92. 'fixture:probe/rename': (
  93. request: { readonly objective: string },
  94. ) => Promise<RemoteResult<{ readonly renamed: boolean }>>
  95. }
  96. interface TypertRemoteNamespaceMap {
  97. probe: TypertRemoteNamespace<'probe'>
  98. }
  99. }
  100. type FixtureContext = Omit<Context, 'remote'> & {
  101. readonly remote: ClientRemote & TypertRemoteScopeApi<'fixture'>
  102. }
  103. // Compile-time contract of `$on`: the key face is the forwarding selection and
  104. // the listener signature is the owning package's own Cordis declaration.
  105. function remoteEventContracts(remote: ClientRemote): void {
  106. remote.$on('fixture/changed', (namespace) => { void namespace })
  107. remote.$on('fixture/approval', async function (request, next) {
  108. expectTypeOf(this).toEqualTypeOf<Context>()
  109. expectTypeOf(request.agent).toEqualTypeOf<Context>()
  110. expectTypeOf(request.signal).toEqualTypeOf<AbortSignal | undefined>()
  111. return request.prompt === '' ? next() : 'allowed'
  112. })
  113. // @ts-expect-error -- declared in Events but outside the forwarding selection.
  114. remote.$on('fixture/unselected', () => {})
  115. // @ts-expect-error -- not declared in Events at all.
  116. remote.$on('fixture/absent', () => {})
  117. // @ts-expect-error -- the listener signature comes from the event declaration.
  118. remote.$on('fixture/changed', (count: number) => { void count })
  119. }
  120. void remoteEventContracts
  121. const idSchema = z.string().min(1)
  122. const requestSchema = z.object({ objective: z.string().min(1) })
  123. const createResultSchema = z.object({ ref: z.string().min(1) })
  124. const renameResultSchema = z.object({ renamed: z.boolean() })
  125. function directDescriptor(): InvocationDescriptor {
  126. return {
  127. id: '@fixture/probe#probe/create',
  128. service: 'probe',
  129. namespace: 'probe',
  130. method: 'create',
  131. invocation: { kind: 'direct' },
  132. scope: { context: 'fixture', wire: 'agentId' },
  133. parameters: [{
  134. name: 'agent',
  135. wire: 'agentId',
  136. source: 'lookup',
  137. lookup: 'fixture',
  138. codec: { mode: 'strict', typeSymbol: '@fixture#AgentId', schema: idSchema },
  139. }, {
  140. name: 'request',
  141. wire: 'request',
  142. source: 'json',
  143. codec: { mode: 'strict', typeSymbol: '@fixture#CreateRequest', schema: requestSchema },
  144. }],
  145. cancellation: { parameter: 'signal' },
  146. result: { mode: 'strict', typeSymbol: '@fixture#CreateResult', schema: createResultSchema },
  147. }
  148. }
  149. function contextDescriptor(): InvocationDescriptor {
  150. return {
  151. id: '@fixture/probe#probe/rename',
  152. service: 'probe',
  153. namespace: 'probe',
  154. method: 'rename',
  155. invocation: {
  156. kind: 'context',
  157. context: 'fixture',
  158. wire: 'agentId',
  159. codec: { mode: 'strict', typeSymbol: '@fixture#AgentId', schema: idSchema },
  160. },
  161. parameters: [{
  162. name: 'request',
  163. wire: 'request',
  164. source: 'json',
  165. codec: { mode: 'strict', typeSymbol: '@fixture#RenameRequest', schema: requestSchema },
  166. }],
  167. result: { mode: 'strict', typeSymbol: '@fixture#RenameResult', schema: renameResultSchema },
  168. }
  169. }
  170. function maybeDescriptor(): InvocationDescriptor {
  171. const schema = z.union([z.string(), z.null(), z.undefined()])
  172. return {
  173. id: '@fixture/probe#probe/maybe',
  174. service: 'probe',
  175. namespace: 'probe',
  176. method: 'maybe',
  177. invocation: { kind: 'direct' },
  178. parameters: [{
  179. name: 'value',
  180. wire: 'value',
  181. source: 'json',
  182. acceptsUndefined: true,
  183. codec: { mode: 'strict', typeSymbol: '@fixture#MaybeValue', schema },
  184. }],
  185. result: { mode: 'strict', typeSymbol: '@fixture#MaybeValue', schema },
  186. }
  187. }
  188. function streamDescriptor(): InvocationDescriptor {
  189. return {
  190. id: '@fixture/probe#probe/watch',
  191. service: 'probe',
  192. namespace: 'probe',
  193. method: 'watch',
  194. mode: 'stream',
  195. invocation: { kind: 'direct' },
  196. parameters: [{
  197. name: 'topic',
  198. wire: 'topic',
  199. source: 'json',
  200. codec: { mode: 'strict', typeSymbol: '@fixture#Topic', schema: z.string().min(1) },
  201. }],
  202. cancellation: { parameter: 'signal' },
  203. result: { mode: 'strict', typeSymbol: '@fixture#WatchItem', schema: z.string().min(1) },
  204. }
  205. }
  206. type WebSocketGlobal = { WebSocket?: typeof WebSocket }
  207. class FakeWebSocket extends EventTarget {
  208. static readonly CONNECTING = 0
  209. static readonly OPEN = 1
  210. static readonly CLOSING = 2
  211. static readonly CLOSED = 3
  212. static readonly sockets: FakeWebSocket[] = []
  213. static autoOpen = true
  214. static dispatchClose = true
  215. readonly url: string
  216. readonly sent: string[] = []
  217. readonly closedWith: { readonly code?: number; readonly reason?: string }[] = []
  218. readyState = FakeWebSocket.CONNECTING
  219. constructor(url: string | URL) {
  220. super()
  221. this.url = String(url)
  222. FakeWebSocket.sockets.push(this)
  223. queueMicrotask(() => {
  224. if (FakeWebSocket.autoOpen) this.open()
  225. })
  226. }
  227. open(): void {
  228. if (this.readyState !== FakeWebSocket.CONNECTING) return
  229. this.readyState = FakeWebSocket.OPEN
  230. this.dispatchEvent(new Event('open'))
  231. }
  232. fail(): void {
  233. this.dispatchEvent(new Event('error'))
  234. }
  235. send(data: string): void {
  236. if (this.readyState !== FakeWebSocket.OPEN) throw new Error('fixture socket is not open')
  237. this.sent.push(data)
  238. }
  239. close(code?: number, reason?: string): void {
  240. this.closedWith.push({
  241. ...(code === undefined ? {} : { code }),
  242. ...(reason === undefined ? {} : { reason }),
  243. })
  244. if (this.readyState === FakeWebSocket.CLOSED) return
  245. if (!FakeWebSocket.dispatchClose) {
  246. this.readyState = FakeWebSocket.CLOSING
  247. return
  248. }
  249. this.drop()
  250. }
  251. drop(): void {
  252. if (this.readyState === FakeWebSocket.CLOSED) return
  253. this.readyState = FakeWebSocket.CLOSED
  254. this.dispatchEvent(new Event('close'))
  255. }
  256. receive(value: unknown): void {
  257. this.receiveRaw(typeof value === 'string' ? value : JSON.stringify(value))
  258. }
  259. receiveRaw(data: unknown): void {
  260. this.dispatchEvent(new MessageEvent('message', {
  261. data,
  262. }))
  263. }
  264. }
  265. async function bench(
  266. call: ConnectionHandle['rpc']['call'],
  267. carrier: 'in-process' | 'web' = 'in-process',
  268. ): Promise<Context> {
  269. const { ctx } = await benchFiber(call, carrier)
  270. return ctx
  271. }
  272. async function benchFiber(
  273. call: ConnectionHandle['rpc']['call'],
  274. carrier: 'in-process' | 'web' = 'in-process',
  275. open: NonNullable<ConnectionHandle['rpc']['open']> = () => unexpectedInProcessStream(),
  276. ): Promise<{
  277. readonly ctx: Context
  278. readonly client: Fiber
  279. readonly generation: GenerationHarness
  280. }> {
  281. const ctx = new Context()
  282. await ctx.plugin(TypertRegistry)
  283. const rpc = carrier === 'web'
  284. ? { call }
  285. : { call, open }
  286. const generation = new GenerationHarness()
  287. ctx.provide('connection', {
  288. rpc,
  289. registerGenerationSource: generation.register,
  290. start: () => ({ stop: () => {} }),
  291. } as unknown as ConnectionHandle)
  292. const client = ctx.plugin({ inject, apply })
  293. await client
  294. return { ctx, client, generation }
  295. }
  296. async function *unexpectedInProcessStream(): AsyncGenerator<never> {
  297. throw new Error('fixture did not install an in-process stream')
  298. }
  299. interface GenerationRun {
  300. readonly signal: AbortSignal
  301. readonly ready: Promise<void>
  302. readonly done: Promise<void>
  303. abort(reason?: unknown): void
  304. }
  305. class GenerationHarness {
  306. private source: ConnectionGenerationSource | undefined
  307. private active: AbortController | undefined
  308. readonly register = (source: ConnectionGenerationSource): (() => void) => {
  309. if (this.source !== undefined) throw new Error('fixture generation source already registered')
  310. this.source = source
  311. return () => {
  312. if (this.source !== source) return
  313. this.source = undefined
  314. this.active?.abort(new Error('fixture generation source removed'))
  315. this.active = undefined
  316. }
  317. }
  318. start(): GenerationRun {
  319. if (this.source === undefined) throw new Error('fixture generation source is not registered')
  320. if (this.active !== undefined) throw new Error('fixture generation is already active')
  321. const source = this.source
  322. const controller = new AbortController()
  323. this.active = controller
  324. let reportReady!: () => void
  325. const ready = new Promise<void>((resolve) => { reportReady = resolve })
  326. const done = Promise.resolve()
  327. .then(() => source(controller.signal, reportReady))
  328. .finally(() => {
  329. if (this.active === controller) this.active = undefined
  330. })
  331. void done.catch(() => undefined)
  332. return {
  333. signal: controller.signal,
  334. ready,
  335. done,
  336. abort: (reason) => { controller.abort(reason) },
  337. }
  338. }
  339. startOverlapping(): GenerationRun {
  340. if (this.source === undefined) throw new Error('fixture generation source is not registered')
  341. const controller = new AbortController()
  342. let reportReady!: () => void
  343. const ready = new Promise<void>((resolve) => { reportReady = resolve })
  344. const done = Promise.resolve().then(() => this.source?.(controller.signal, reportReady))
  345. .then(() => undefined)
  346. void done.catch(() => undefined)
  347. return {
  348. signal: controller.signal,
  349. ready,
  350. done,
  351. abort: (reason) => { controller.abort(reason) },
  352. }
  353. }
  354. }
  355. function deferredReadiness(): {
  356. readonly promise: Promise<void>
  357. readonly resolve: () => void
  358. readonly reject: (error: unknown) => void
  359. } {
  360. let resolve!: () => void
  361. let reject!: (error: unknown) => void
  362. const promise = new Promise<void>((accept, decline) => {
  363. resolve = accept
  364. reject = decline
  365. })
  366. return { promise, resolve, reject }
  367. }
  368. async function loaderReadinessBench(readiness: Promise<unknown>): Promise<{
  369. readonly client: Fiber
  370. readonly start: ReturnType<typeof vi.fn<ConnectionHandle['start']>>
  371. readonly stop: ReturnType<typeof vi.fn<() => void>>
  372. }> {
  373. const ctx = new Context()
  374. await ctx.plugin(TypertRegistry)
  375. const generation = new GenerationHarness()
  376. const stop = vi.fn<() => void>()
  377. const start = vi.fn<ConnectionHandle['start']>(() => ({ stop }))
  378. ctx.provide('connection', {
  379. rpc: {
  380. call: vi.fn<ConnectionHandle['rpc']['call']>(),
  381. open: () => unexpectedInProcessStream(),
  382. },
  383. registerGenerationSource: generation.register,
  384. start,
  385. } as unknown as ConnectionHandle)
  386. ctx.provide('loader', { await: () => readiness })
  387. const client = ctx.plugin({ inject, apply })
  388. await client
  389. return { client, start, stop }
  390. }
  391. type EventStreamItem =
  392. | { readonly kind: 'frame'; readonly value: unknown }
  393. | { readonly kind: 'end' }
  394. | { readonly kind: 'fail'; readonly error: unknown }
  395. interface EventStreamConnection {
  396. readonly items: EventStreamItem[]
  397. wake: (() => void) | undefined
  398. }
  399. class RemoteEventCarrier {
  400. readonly calls: {
  401. readonly channel: string
  402. readonly endpoint: string
  403. readonly payload: unknown
  404. readonly signal: AbortSignal
  405. }[] = []
  406. private readonly connections = new Set<EventStreamConnection>()
  407. private nextClient = 1
  408. get activeConnections(): number {
  409. return this.connections.size
  410. }
  411. readonly open: NonNullable<ConnectionHandle['rpc']['open']> = (channel, endpoint, payload, signal) => {
  412. this.calls.push({ channel, endpoint, payload, signal })
  413. return this.iterate(signal)
  414. }
  415. emit(value: unknown): void {
  416. this.feed({ kind: 'frame', value })
  417. }
  418. end(): void {
  419. this.feed({ kind: 'end' })
  420. }
  421. fail(error: unknown): void {
  422. this.feed({ kind: 'fail', error })
  423. }
  424. private feed(item: EventStreamItem): void {
  425. for (const connection of this.connections) {
  426. connection.items.push(item)
  427. connection.wake?.()
  428. }
  429. }
  430. private async *iterate(signal: AbortSignal): AsyncGenerator {
  431. signal.throwIfAborted()
  432. const clientId = `event-client-${String(this.nextClient++)}`
  433. const connection: EventStreamConnection = { items: [], wake: undefined }
  434. this.connections.add(connection)
  435. const abort = (): void => { connection.wake?.() }
  436. signal.addEventListener('abort', abort, { once: true })
  437. try {
  438. yield { type: 'ready', clientId }
  439. while (!signal.aborted) {
  440. while (connection.items.length > 0) {
  441. const item = connection.items.shift() as EventStreamItem
  442. if (item.kind === 'end') return
  443. if (item.kind === 'fail') throw item.error
  444. yield item.value
  445. }
  446. if (signal.aborted) return
  447. await new Promise<void>((resolve) => { connection.wake = resolve })
  448. connection.wake = undefined
  449. }
  450. } finally {
  451. signal.removeEventListener('abort', abort)
  452. this.connections.delete(connection)
  453. }
  454. }
  455. }
  456. async function eventBench(
  457. call: ConnectionHandle['rpc']['call'] = vi.fn<ConnectionHandle['rpc']['call']>()
  458. .mockResolvedValue({ ok: true, value: undefined }),
  459. ): Promise<{
  460. readonly ctx: Context
  461. readonly client: Fiber
  462. readonly carrier: RemoteEventCarrier
  463. readonly generation: GenerationHarness
  464. readonly run: GenerationRun
  465. readonly call: ConnectionHandle['rpc']['call']
  466. }> {
  467. const carrier = new RemoteEventCarrier()
  468. const { ctx, client, generation } = await benchFiber(
  469. call,
  470. 'in-process',
  471. carrier.open,
  472. )
  473. const run = generation.start()
  474. await run.ready
  475. return { ctx, client, carrier, generation, run, call }
  476. }
  477. function approvalFrame(eventId: string, agentId: string, prompt: string): object {
  478. return {
  479. type: 'waterfall',
  480. event: 'fixture/approval',
  481. eventId,
  482. agentId,
  483. request: { prompt },
  484. }
  485. }
  486. describe('Client Remote transport readiness', () => {
  487. it('creates logical stream supervisors against the installed Connection', async () => {
  488. const { ctx, client } = await benchFiber(vi.fn<ConnectionHandle['rpc']['call']>())
  489. const stream = ctx.remote.$stream({
  490. name: 'fixture stream',
  491. open: () => unexpectedInProcessStream(),
  492. ended: () => new Error('fixture stream ended'),
  493. })
  494. expect(stream).toBeInstanceOf(RemoteStream)
  495. await stream.dispose()
  496. await client.dispose()
  497. })
  498. it('starts after Loader settlement and stops the owned loop on disposal', async () => {
  499. const readiness = deferredReadiness()
  500. const { client, start, stop } = await loaderReadinessBench(readiness.promise)
  501. expect(start).not.toHaveBeenCalled()
  502. readiness.resolve()
  503. await vi.waitFor(() => { expect(start).toHaveBeenCalledTimes(1) })
  504. await client.dispose()
  505. expect(stop).toHaveBeenCalledTimes(1)
  506. })
  507. it('does not start when disposal wins the Loader-settlement race', async () => {
  508. const readiness = deferredReadiness()
  509. const { client, start, stop } = await loaderReadinessBench(readiness.promise)
  510. await client.dispose()
  511. readiness.resolve()
  512. await Promise.resolve()
  513. expect(start).not.toHaveBeenCalled()
  514. expect(stop).not.toHaveBeenCalled()
  515. })
  516. it('leaves the transport stopped when Loader settlement rejects', async () => {
  517. const readiness = deferredReadiness()
  518. const { client, start, stop } = await loaderReadinessBench(readiness.promise)
  519. readiness.reject(new Error('fixture Loader failed'))
  520. await Promise.resolve()
  521. await Promise.resolve()
  522. expect(start).not.toHaveBeenCalled()
  523. await client.dispose()
  524. expect(stop).not.toHaveBeenCalled()
  525. })
  526. })
  527. describe('Client Typert API', () => {
  528. it('mounts concrete direct methods, validates inputs, and withdraws retained handles', async () => {
  529. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  530. .mockResolvedValue({ ok: true, value: { ref: 'goal-1' } })
  531. const ctx = await bench(call)
  532. const businessProbe = { owner: 'host business service' }
  533. const disposeBusinessProbe = ctx.provide('probe', businessProbe)
  534. const assembly = ctx.plugin(Object.assign(
  535. (scope: Context) => scope.remote.$mount({ package: '@fixture/probe', descriptors: [directDescriptor()] }),
  536. { inject: ['remote'] },
  537. ))
  538. await assembly
  539. const retained = ctx.remote.probe.create
  540. await expect(ctx.remote.probe.create('agent-1', { objective: 'ship' }))
  541. .resolves.toEqual({ ok: true, value: { ref: 'goal-1' } })
  542. expect(call).toHaveBeenCalledWith(
  543. '/api',
  544. 'probe/create',
  545. { args: { agentId: 'agent-1', request: { objective: 'ship' } } },
  546. expect.any(AbortSignal),
  547. )
  548. const callerAbort = new AbortController()
  549. await expect(ctx.remote.probe.create(
  550. 'agent-1',
  551. { objective: 'cancel me' },
  552. callerAbort.signal,
  553. )).resolves.toEqual({ ok: true, value: { ref: 'goal-1' } })
  554. const combinedSignal = call.mock.calls.at(-1)?.[3]
  555. expect(combinedSignal).toBeInstanceOf(AbortSignal)
  556. expect(combinedSignal).not.toBe(callerAbort.signal)
  557. const cancellation = new Error('caller cancelled')
  558. callerAbort.abort(cancellation)
  559. expect(combinedSignal?.aborted).toBe(true)
  560. expect(combinedSignal?.reason).toBe(cancellation)
  561. await expect(ctx.remote.probe.create('', { objective: 'ship' })).rejects.toThrow('rejected "agentId"')
  562. call.mockResolvedValueOnce({ ok: true, value: { ref: 1 } })
  563. await expect(ctx.remote.probe.create('agent-1', { objective: 'ship' })).resolves.toEqual({
  564. ok: true,
  565. value: { ref: 1 },
  566. })
  567. await assembly.dispose()
  568. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  569. expect(ctx.get('remote.probe')).toBeUndefined()
  570. expect(ctx.get('probe')).toBe(businessProbe)
  571. expect(ctx.typert.remotes.list()).toEqual([])
  572. await expect(retained?.('agent-1', { objective: 'ship' })).resolves.toEqual({
  573. ok: false,
  574. error: {
  575. code: 'internal',
  576. message: 'client api: Remote method probe/create is no longer mounted',
  577. details: {},
  578. },
  579. })
  580. disposeBusinessProbe()
  581. })
  582. it('encodes declared undefined as an omitted argument and distinguishes it from null results', async () => {
  583. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  584. .mockResolvedValueOnce({ ok: true, value: undefined })
  585. .mockResolvedValueOnce({ ok: true, value: null })
  586. const ctx = await bench(call)
  587. const dispose = await ctx.remote.$mount({
  588. package: '@fixture/maybe',
  589. descriptors: [maybeDescriptor()],
  590. })
  591. await expect(ctx.remote.probe.maybe(undefined)).resolves.toStrictEqual({ ok: true, value: undefined })
  592. expect(call).toHaveBeenNthCalledWith(
  593. 1,
  594. '/api',
  595. 'probe/maybe',
  596. { args: {} },
  597. expect.any(AbortSignal),
  598. )
  599. await expect(ctx.remote.probe.maybe(null)).resolves.toStrictEqual({ ok: true, value: null })
  600. expect(call).toHaveBeenNthCalledWith(
  601. 2,
  602. '/api',
  603. 'probe/maybe',
  604. { args: { value: null } },
  605. expect.any(AbortSignal),
  606. )
  607. await dispose()
  608. })
  609. it('projects one direct lookup descriptor onto an Agent-scoped alias', async () => {
  610. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  611. .mockResolvedValue({ ok: true, value: { ref: 'goal-2' } })
  612. const ctx = await bench(call)
  613. const agentCtx = ctx.extend({ fixtureId: 'agent-2' }) as FixtureContext
  614. ctx.typert.contexts.registerClient('fixture', {
  615. identity: candidate => (candidate as Context & { fixtureId?: string }).fixtureId,
  616. resolve: id => id === 'agent-2' ? agentCtx : undefined,
  617. })
  618. const assembly = ctx.plugin(Object.assign(
  619. (scope: Context) => scope.remote.$mount({ package: '@fixture/probe', descriptors: [directDescriptor()] }),
  620. { inject: ['remote'] },
  621. ))
  622. await assembly
  623. await expect(agentCtx.remote.probe.create({ objective: 'ship scoped' }))
  624. .resolves.toEqual({ ok: true, value: { ref: 'goal-2' } })
  625. expect(call).toHaveBeenCalledWith(
  626. '/api',
  627. 'probe/create',
  628. { args: { agentId: 'agent-2', request: { objective: 'ship scoped' } } },
  629. expect.any(AbortSignal),
  630. )
  631. await expect((ctx as FixtureContext).remote.probe.create({ objective: 'wrong scope' }))
  632. .rejects.toThrow('expected 2 business argument(s)')
  633. await assembly.dispose()
  634. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  635. expect(ctx.get('remote.probe')).toBeUndefined()
  636. })
  637. it('uses the caller Context identity for scoped namespace methods', async () => {
  638. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  639. .mockResolvedValue({ ok: true, value: { renamed: true } })
  640. const ctx = await bench(call)
  641. const agentCtx = ctx.extend({ fixtureId: 'agent-2' }) as FixtureContext
  642. ctx.typert.contexts.registerClient('fixture', {
  643. identity: candidate => (candidate as Context & { fixtureId?: string }).fixtureId,
  644. resolve: id => id === 'agent-2' ? agentCtx : undefined,
  645. })
  646. const assembly = ctx.plugin(Object.assign(
  647. (scope: Context) => scope.remote.$mount({ package: '@fixture/probe', descriptors: [contextDescriptor()] }),
  648. { inject: ['remote'] },
  649. ))
  650. await assembly
  651. await expect(agentCtx.remote.probe.rename({ objective: 'land' }))
  652. .resolves.toEqual({ ok: true, value: { renamed: true } })
  653. expect(call).toHaveBeenCalledWith(
  654. '/api',
  655. 'probe/rename',
  656. { args: { agentId: 'agent-2', request: { objective: 'land' } } },
  657. expect.any(AbortSignal),
  658. )
  659. await expect((ctx as FixtureContext).remote.probe.rename({ objective: 'land' }))
  660. .rejects.toThrow('requires a "fixture" Context')
  661. await assembly.dispose()
  662. expect(ctx.get('remote.probe')).toBeUndefined()
  663. })
  664. it('accepts weak result codecs and rejects namespace collisions before registration', async () => {
  665. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  666. const weak: InvocationDescriptor = {
  667. ...directDescriptor(),
  668. result: { mode: 'src-json' },
  669. }
  670. const disposeWeak = await ctx.remote.$mount({ package: '@fixture/weak', descriptors: [weak] })
  671. await disposeWeak()
  672. await expect(ctx.remote.$mount({
  673. package: '@fixture/conflict',
  674. descriptors: [{ ...directDescriptor(), namespace: '$mount' }],
  675. })).rejects.toThrow('conflicts with the Remote service')
  676. expect(ctx.typert.remotes.list()).toEqual([])
  677. })
  678. it('rejects duplicate, live, scoped-service, and Context namespace collisions', async () => {
  679. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  680. .mockResolvedValue({ ok: true, value: { renamed: true } })
  681. const ctx = await bench(call)
  682. const agentCtx = ctx.extend({ fixtureId: 'agent-remounted' }) as FixtureContext
  683. ctx.typert.contexts.registerClient('fixture', {
  684. identity: candidate => (candidate as Context & { fixtureId?: string }).fixtureId,
  685. resolve: id => id === 'agent-remounted' ? agentCtx : undefined,
  686. })
  687. const direct = directDescriptor()
  688. const context = contextDescriptor()
  689. await expect(ctx.remote.$mount({
  690. package: '@fixture/direct-duplicates',
  691. descriptors: [direct, { ...direct, id: '@fixture/probe#probe/create-again' }],
  692. })).rejects.toThrow('repeats direct method')
  693. await expect(ctx.remote.$mount({
  694. package: '@fixture/scoped-duplicates',
  695. descriptors: [context, { ...context, id: '@fixture/probe#probe/rename-again' }],
  696. })).rejects.toThrow('repeats scoped method')
  697. const disposeDirect = await ctx.remote.$mount({ package: '@fixture/direct-live', descriptors: [direct] })
  698. await expect(ctx.remote.$mount({
  699. package: '@fixture/direct-conflict', descriptors: [{ ...direct, id: '@fixture/other#probe/create' }],
  700. })).rejects.toThrow('direct method probe/create is already mounted')
  701. await disposeDirect()
  702. const disposeScoped = await ctx.remote.$mount({ package: '@fixture/scoped-live', descriptors: [context] })
  703. await expect(ctx.remote.$mount({
  704. package: '@fixture/scoped-conflict', descriptors: [{ ...context, id: '@fixture/other#probe/rename' }],
  705. })).rejects.toThrow('scoped method probe/rename is already mounted')
  706. await expect(ctx.remote.$mount({
  707. package: '@fixture/service-method-conflict',
  708. descriptors: [{ ...context, id: '@fixture/probe#probe/remove', method: 'remove' }],
  709. })).rejects.toThrow('conflicts with its namespace service')
  710. const scopedService = ctx.get('remote.probe') as unknown as object
  711. Object.defineProperty(scopedService, 'custom', { configurable: true, value: () => undefined })
  712. await expect(ctx.remote.$mount({
  713. package: '@fixture/service-own-property-conflict',
  714. descriptors: [{ ...direct, id: '@fixture/probe#probe/custom', method: 'custom' }],
  715. })).rejects.toThrow('conflicts with its namespace service')
  716. Reflect.deleteProperty(scopedService, 'custom')
  717. await disposeScoped()
  718. const disposeRemoteTypert = ctx.reflect.provide('remote.typert', { owner: 'fixture' })
  719. await expect(ctx.remote.$mount({
  720. package: '@fixture/context-property-conflict',
  721. descriptors: [{ ...context, namespace: 'typert' }],
  722. })).rejects.toThrow('conflicts with an existing Remote namespace')
  723. await disposeRemoteTypert()
  724. const disposeMultipleScoped = await ctx.remote.$mount({
  725. package: '@fixture/multiple-scoped',
  726. descriptors: [directDescriptor(), contextDescriptor()],
  727. })
  728. await expect(agentCtx.remote.probe.rename({ objective: 'remounted' }))
  729. .resolves.toEqual({ ok: true, value: { renamed: true } })
  730. expect(call).toHaveBeenLastCalledWith(
  731. '/api',
  732. 'probe/rename',
  733. { args: { agentId: 'agent-remounted', request: { objective: 'remounted' } } },
  734. expect.any(AbortSignal),
  735. )
  736. await disposeMultipleScoped()
  737. })
  738. it('rolls back earlier descriptors when a later descriptor fails to install', async () => {
  739. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  740. const { scope: _scope, ...first } = directDescriptor()
  741. const second: InvocationDescriptor = {
  742. ...first,
  743. id: '@fixture/probe#probe/archive',
  744. method: 'archive',
  745. }
  746. const defineProperty = Object.defineProperty
  747. const spy = vi.spyOn(Object, 'defineProperty').mockImplementation((target, key, attributes) => {
  748. if (key === 'archive') throw new Error('fixture later-descriptor failure')
  749. return defineProperty(target, key, attributes)
  750. })
  751. try {
  752. await expect(ctx.remote.$mount({ package: '@fixture/failing-batch', descriptors: [first, second] }))
  753. .rejects.toThrow('fixture later-descriptor failure')
  754. } finally {
  755. spy.mockRestore()
  756. }
  757. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  758. await vi.waitFor(() => { expect(ctx.typert.remotes.list()).toEqual([]) })
  759. const retry = await ctx.remote.$mount({ package: '@fixture/retry-batch', descriptors: [first, second] })
  760. expect(ctx.remote.probe.create).toBeTypeOf('function')
  761. expect((ctx.remote.probe as unknown as Record<string, unknown>).archive).toBeTypeOf('function')
  762. await retry()
  763. })
  764. it('rolls back earlier namespaces when a later namespace fails to install', async () => {
  765. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  766. const { scope: _scope, ...first } = directDescriptor()
  767. const second: InvocationDescriptor = {
  768. ...first,
  769. id: '@fixture/archive#archive/store',
  770. namespace: 'archive',
  771. method: 'store',
  772. }
  773. const defineProperty = Object.defineProperty
  774. const spy = vi.spyOn(Object, 'defineProperty').mockImplementation((target, key, attributes) => {
  775. if (key === 'store') throw new Error('fixture later-namespace failure')
  776. return defineProperty(target, key, attributes)
  777. })
  778. try {
  779. await expect(ctx.remote.$mount({
  780. package: '@fixture/failing-namespaces',
  781. descriptors: [first, second],
  782. })).rejects.toThrow('fixture later-namespace failure')
  783. } finally {
  784. spy.mockRestore()
  785. }
  786. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  787. expect((ctx.remote as unknown as Record<string, unknown>).archive).toBeUndefined()
  788. await vi.waitFor(() => { expect(ctx.typert.remotes.list()).toEqual([]) })
  789. const retry = await ctx.remote.$mount({
  790. package: '@fixture/retry-namespaces',
  791. descriptors: [first, second],
  792. })
  793. expect(ctx.remote.probe.create).toBeTypeOf('function')
  794. expect((ctx.remote as unknown as Record<string, Record<string, unknown>>).archive?.store).toBeTypeOf('function')
  795. await retry()
  796. })
  797. it('rolls back a direct projection when its scoped projection fails to install', async () => {
  798. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  799. const disposeContext = await ctx.remote.$mount({
  800. package: '@fixture/context-anchor',
  801. descriptors: [contextDescriptor()],
  802. })
  803. const namespace = ctx.get('remote.probe') as unknown as {
  804. installScoped: (...args: unknown[]) => void
  805. readonly create?: unknown
  806. }
  807. const installScoped = vi.spyOn(namespace, 'installScoped').mockImplementation(() => {
  808. throw new Error('fixture scoped projection failure')
  809. })
  810. try {
  811. await expect(ctx.remote.$mount({
  812. package: '@fixture/direct-projection-failure',
  813. descriptors: [directDescriptor()],
  814. })).rejects.toThrow('fixture scoped projection failure')
  815. } finally {
  816. installScoped.mockRestore()
  817. }
  818. expect(namespace.create).toBeUndefined()
  819. await disposeContext()
  820. })
  821. it('unwinds an already-installed namespace when a later namespace fails to install', async () => {
  822. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  823. const { scope: _scope, ...probe } = directDescriptor()
  824. const vault: InvocationDescriptor = {
  825. ...probe,
  826. id: '@fixture/vault#vault/seal',
  827. service: 'vault',
  828. namespace: 'vault',
  829. method: 'seal',
  830. }
  831. const defineProperty = Object.defineProperty
  832. const spy = vi.spyOn(Object, 'defineProperty').mockImplementation((target, key, attributes) => {
  833. if (key === 'seal') throw new Error('fixture later-namespace failure')
  834. return defineProperty(target, key, attributes)
  835. })
  836. try {
  837. await expect(ctx.remote.$mount({ package: '@fixture/two-namespaces', descriptors: [probe, vault] }))
  838. .rejects.toThrow('fixture later-namespace failure')
  839. } finally {
  840. spy.mockRestore()
  841. }
  842. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  843. expect((ctx.remote as unknown as Record<string, unknown>).vault).toBeUndefined()
  844. await vi.waitFor(() => { expect(ctx.typert.remotes.list()).toEqual([]) })
  845. })
  846. it('rolls back an earlier scoped projection when a later descriptor fails to install', async () => {
  847. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  848. const { scope: _scope, ...direct } = directDescriptor()
  849. const failing: InvocationDescriptor = {
  850. ...direct,
  851. id: '@fixture/probe#probe/archive',
  852. method: 'archive',
  853. }
  854. const defineProperty = Object.defineProperty
  855. const spy = vi.spyOn(Object, 'defineProperty').mockImplementation((target, key, attributes) => {
  856. if (key === 'archive') throw new Error('fixture trailing failure')
  857. return defineProperty(target, key, attributes)
  858. })
  859. try {
  860. await expect(ctx.remote.$mount({
  861. package: '@fixture/scoped-then-failing',
  862. descriptors: [contextDescriptor(), failing],
  863. })).rejects.toThrow('fixture trailing failure')
  864. } finally {
  865. spy.mockRestore()
  866. }
  867. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  868. await vi.waitFor(() => { expect(ctx.typert.remotes.list()).toEqual([]) })
  869. })
  870. it('keeps a namespace another contribution still populates when a group leaves', async () => {
  871. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  872. .mockResolvedValue({ ok: true, value: { renamed: true } })
  873. const ctx = await bench(call)
  874. const { scope: _scope, ...direct } = directDescriptor()
  875. const disposeDirect = await ctx.remote.$mount({ package: '@fixture/direct-owner', descriptors: [direct] })
  876. const disposeScoped = await ctx.remote.$mount({ package: '@fixture/scoped-owner', descriptors: [contextDescriptor()] })
  877. await disposeDirect()
  878. // The namespace survives its first group: the second contribution still owns methods on it.
  879. const surviving = ctx.get('remote.probe') as unknown as Record<string, unknown> | undefined
  880. expect(surviving).toBeDefined()
  881. expect(surviving?.create).toBeUndefined()
  882. await disposeScoped()
  883. expect(ctx.get('remote.probe')).toBeUndefined()
  884. })
  885. it('unparks a namespace dependent only after its contribution methods exist', async () => {
  886. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  887. const { scope: _scope, ...direct } = directDescriptor()
  888. let observed: string | undefined
  889. // Parked before the mount: the unpark moment is the observation — the
  890. // atomic-visibility guarantee says the service never appears methodless.
  891. const parked = ctx.inject(['remote.probe'], (probeCtx) => {
  892. observed = typeof (probeCtx.get('remote.probe') as { create?: unknown } | undefined)?.create
  893. })
  894. const dispose = await ctx.remote.$mount({ package: '@fixture/atomic-visibility', descriptors: [direct] })
  895. await parked
  896. expect(observed).toBe('function')
  897. await dispose()
  898. })
  899. it('rejects weak parameter and Context codecs plus malformed scope projections', async () => {
  900. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  901. const direct = directDescriptor()
  902. const context = contextDescriptor()
  903. await expect(ctx.remote.$mount({
  904. package: '@fixture/weak-parameter',
  905. descriptors: [{
  906. ...direct,
  907. parameters: direct.parameters.map((parameter, index) => index === 0
  908. ? { ...parameter, codec: { mode: 'src-json' } }
  909. : parameter),
  910. }],
  911. })).rejects.toThrow('has no strict codec')
  912. await expect(ctx.remote.$mount({
  913. package: '@fixture/weak-context',
  914. descriptors: [{
  915. ...context,
  916. invocation: { ...context.invocation, codec: { mode: 'src-json' } },
  917. } as InvocationDescriptor],
  918. })).rejects.toThrow('has no strict codec')
  919. await expect(ctx.remote.$mount({
  920. package: '@fixture/malformed-scope',
  921. descriptors: [{ ...direct, scope: { context: 'fixture', wire: 'missingId' } }],
  922. })).rejects.toThrow('scope must select its only lookup parameter')
  923. await expect(ctx.remote.$mount({
  924. package: '@fixture/ambiguous-scope',
  925. descriptors: [{
  926. ...direct,
  927. parameters: [...direct.parameters, {
  928. name: 'other', wire: 'otherId', source: 'lookup', lookup: 'fixture',
  929. codec: { mode: 'strict', typeSymbol: '@fixture#AgentId', schema: idSchema },
  930. }],
  931. }],
  932. })).rejects.toThrow('scope must select its only lookup parameter')
  933. })
  934. it('validates invocation arity, required adapters, live Connection, and mutable descriptor codecs', async () => {
  935. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  936. .mockResolvedValue({ ok: true, value: { ref: 'goal-1' } })
  937. const ctx = await bench(call)
  938. const descriptor = directDescriptor()
  939. const dispose = await ctx.remote.$mount({
  940. package: '@fixture/probe',
  941. descriptors: [descriptor, contextDescriptor()],
  942. })
  943. const create = ctx.remote.probe.create as unknown as (...args: unknown[]) => Promise<unknown>
  944. const probe = (ctx as FixtureContext).remote.probe
  945. const rename = probe.rename as unknown as (...args: unknown[]) => Promise<unknown>
  946. await expect(create('agent-1')).rejects.toThrow('expected 2 business argument(s) plus an optional AbortSignal, got 1')
  947. await expect(create('agent-1', { objective: 'ship' }, undefined, 'extra'))
  948. .rejects.toThrow('got 4')
  949. await expect(rename.call(probe)).rejects.toThrow('expected 1 argument(s), got 0')
  950. await expect((ctx as FixtureContext).remote.probe.create({ objective: 'ship' }))
  951. .rejects.toThrow('expected 2 business argument(s)')
  952. await expect((ctx as FixtureContext).remote.probe.rename({ objective: 'ship' }))
  953. .rejects.toThrow('no Client Context adapter')
  954. ;(descriptor.parameters[0] as { codec: { mode: string } }).codec.mode = 'src-json'
  955. await expect(ctx.remote.probe.create('agent-1', { objective: 'ship' })).rejects.toThrow('has no strict codec')
  956. ;(descriptor.parameters[0] as { codec: { mode: string } }).codec.mode = 'strict'
  957. ctx.set('connection', undefined)
  958. await expect(ctx.remote.probe.create('agent-1', { objective: 'ship' })).rejects.toThrow('no active Connection')
  959. await dispose()
  960. })
  961. it('withdraws a pending invocation and preserves a direct namespace until its last method leaves', async () => {
  962. let resolveCall!: (result: Awaited<ReturnType<ConnectionHandle['rpc']['call']>>) => void
  963. const pending = new Promise<Awaited<ReturnType<ConnectionHandle['rpc']['call']>>>((resolve) => {
  964. resolveCall = resolve
  965. })
  966. const call = vi.fn<ConnectionHandle['rpc']['call']>().mockReturnValue(pending)
  967. const ctx = await bench(call)
  968. const { scope: _scope, ...first } = directDescriptor()
  969. const second: InvocationDescriptor = {
  970. ...first,
  971. id: '@fixture/probe#probe/archive',
  972. method: 'archive',
  973. }
  974. const dispose = await ctx.remote.$mount({ package: '@fixture/probe', descriptors: [first, second] })
  975. const invocation = ctx.remote.probe.create('agent-1', { objective: 'ship' })
  976. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(1) })
  977. await dispose()
  978. resolveCall({ ok: true, value: { ref: 'goal-1' } })
  979. await expect(invocation).resolves.toEqual({
  980. ok: false,
  981. error: {
  982. code: 'internal',
  983. message: 'client api: Remote method probe/create is no longer mounted',
  984. details: {},
  985. },
  986. })
  987. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  988. })
  989. it('keeps a namespace while another contribution still owns a method', async () => {
  990. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  991. const disposeCreate = await ctx.remote.$mount({
  992. package: '@fixture/create-contribution',
  993. descriptors: [directDescriptor()],
  994. })
  995. const disposeMaybe = await ctx.remote.$mount({
  996. package: '@fixture/maybe-contribution',
  997. descriptors: [maybeDescriptor()],
  998. })
  999. const namespace = ctx.get('remote.probe') as unknown as Record<string, unknown>
  1000. await disposeCreate()
  1001. expect(ctx.get('remote.probe') !== undefined).toBe(true)
  1002. expect(namespace.create).toBeUndefined()
  1003. expect(namespace.maybe).toBeTypeOf('function')
  1004. await disposeMaybe()
  1005. expect(ctx.get('remote.probe')).toBeUndefined()
  1006. })
  1007. it('fails a method obtained from a withdrawn namespace getter', async () => {
  1008. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  1009. const dispose = await ctx.remote.$mount({ package: '@fixture/probe', descriptors: [directDescriptor()] })
  1010. const namespace = ctx.get('remote.probe') as unknown as object
  1011. const getWithdrawn = Object.getOwnPropertyDescriptor(namespace, 'create')?.get?.bind(namespace)
  1012. await dispose()
  1013. expect(getWithdrawn).toBeTypeOf('function')
  1014. const withdrawn = getWithdrawn?.() as (...args: unknown[]) => Promise<unknown>
  1015. expect(() => withdrawn('agent-1', { objective: 'ship' }))
  1016. .toThrow('Remote method is no longer mounted')
  1017. })
  1018. it('preserves a __proto__ wire parameter as an own named argument', async () => {
  1019. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  1020. .mockResolvedValue({ ok: true, value: { ref: 'goal-1' } })
  1021. const ctx = await bench(call)
  1022. const { scope: _scope, ...base } = directDescriptor()
  1023. const descriptor: InvocationDescriptor = {
  1024. ...base,
  1025. id: '@fixture/probe#probe/prototype',
  1026. method: 'prototype',
  1027. parameters: [{
  1028. name: 'value',
  1029. wire: '__proto__',
  1030. source: 'json',
  1031. codec: { mode: 'strict', typeSymbol: '@fixture#PrototypeValue', schema: z.string() },
  1032. }],
  1033. }
  1034. const dispose = await ctx.remote.$mount({ package: '@fixture/prototype', descriptors: [descriptor] })
  1035. const method = (ctx.remote.probe as unknown as Record<string, (...args: unknown[]) => Promise<unknown>>).prototype
  1036. await expect(method?.('wire-value')).resolves.toEqual({ ok: true, value: { ref: 'goal-1' } })
  1037. const payload = call.mock.calls[0]?.[2] as { readonly args: Record<string, unknown> }
  1038. expect(Object.getPrototypeOf(payload.args)).toBeNull()
  1039. expect(Object.hasOwn(payload.args, '__proto__')).toBe(true)
  1040. expect(payload.args.__proto__).toBe('wire-value')
  1041. await dispose()
  1042. })
  1043. it('rolls back Remote registration when namespace Service startup fails', async () => {
  1044. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  1045. const defineProperty = Object.defineProperty
  1046. const spy = vi.spyOn(Object, 'defineProperty').mockImplementation((target, key, attributes) => {
  1047. if (key === Service.tracker) throw new Error('fixture namespace startup failure')
  1048. return defineProperty(target, key, attributes)
  1049. })
  1050. try {
  1051. await expect(ctx.remote.$mount({ package: '@fixture/probe', descriptors: [directDescriptor()] }))
  1052. .rejects.toThrow('fixture namespace startup failure')
  1053. await vi.waitFor(() => { expect(ctx.typert.remotes.list()).toEqual([]) })
  1054. } finally {
  1055. spy.mockRestore()
  1056. }
  1057. const retry = await ctx.remote.$mount({ package: '@fixture/probe-retry', descriptors: [directDescriptor()] })
  1058. expect(ctx.remote.probe.create).toBeTypeOf('function')
  1059. await retry()
  1060. })
  1061. it('withdraws a fresh direct namespace when its first method fails to install', async () => {
  1062. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  1063. const defineProperty = Object.defineProperty
  1064. const spy = vi.spyOn(Object, 'defineProperty').mockImplementation((target, key, attributes) => {
  1065. if (key === 'create') throw new Error('fixture direct method installation failure')
  1066. return defineProperty(target, key, attributes)
  1067. })
  1068. try {
  1069. await expect(ctx.remote.$mount({
  1070. package: '@fixture/direct-method-failure',
  1071. descriptors: [directDescriptor()],
  1072. })).rejects.toThrow('fixture direct method installation failure')
  1073. } finally {
  1074. spy.mockRestore()
  1075. }
  1076. expect((ctx.remote as unknown as Record<string, unknown>).probe).toBeUndefined()
  1077. await vi.waitFor(() => { expect(ctx.typert.remotes.list()).toEqual([]) })
  1078. const retry = await ctx.remote.$mount({
  1079. package: '@fixture/direct-method-retry',
  1080. descriptors: [directDescriptor()],
  1081. })
  1082. expect(ctx.remote.probe.create).toBeTypeOf('function')
  1083. await retry()
  1084. })
  1085. it('withdraws a fresh scoped Service when its first method fails to install', async () => {
  1086. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  1087. const defineProperty = Object.defineProperty
  1088. const spy = vi.spyOn(Object, 'defineProperty').mockImplementation((target, key, attributes) => {
  1089. if (key === 'rename') throw new Error('fixture scoped installation failure')
  1090. return defineProperty(target, key, attributes)
  1091. })
  1092. try {
  1093. await expect(ctx.remote.$mount({ package: '@fixture/scoped-failure', descriptors: [contextDescriptor()] }))
  1094. .rejects.toThrow('fixture scoped installation failure')
  1095. } finally {
  1096. spy.mockRestore()
  1097. }
  1098. expect(ctx.get('remote.probe')).toBeUndefined()
  1099. await vi.waitFor(() => { expect(ctx.typert.remotes.list()).toEqual([]) })
  1100. const retry = await ctx.remote.$mount({ package: '@fixture/scoped-retry', descriptors: [contextDescriptor()] })
  1101. expect((ctx.get('remote.probe') as unknown as Record<string, unknown>).rename).toBeTypeOf('function')
  1102. await retry()
  1103. })
  1104. it('unregisters an empty scoped namespace so another provider can claim its name', async () => {
  1105. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  1106. const dispose = await ctx.remote.$mount({ package: '@fixture/scoped', descriptors: [contextDescriptor()] })
  1107. expect(ctx.get('remote.probe')).toBeDefined()
  1108. await dispose()
  1109. expect(ctx.get('remote.probe')).toBeUndefined()
  1110. const replacement = { owner: 'replacement' }
  1111. const disposeReplacement = ctx.reflect.provide('remote.probe', replacement)
  1112. expect(ctx.get('remote.probe')).toBe(replacement)
  1113. await disposeReplacement()
  1114. })
  1115. it('delivers an RPC failure in the error branch with the Host error verbatim', async () => {
  1116. const rpcError = { code: 'internal' as const, message: 'host failed', details: {} }
  1117. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>().mockResolvedValue({ ok: false, error: rpcError }))
  1118. await ctx.remote.$mount({ package: '@fixture/probe', descriptors: [directDescriptor()] })
  1119. const outcome = await ctx.remote.probe.create('agent-1', { objective: 'ship' })
  1120. expect(outcome.ok).toBe(false)
  1121. if (outcome.ok) throw new Error('expected the Client API invocation to report a failure')
  1122. expect(outcome.error).toBe(rpcError)
  1123. })
  1124. it('folds a transport throw into the error branch', async () => {
  1125. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>()
  1126. .mockRejectedValue(new Error('carrier offline')))
  1127. await ctx.remote.$mount({ package: '@fixture/probe', descriptors: [directDescriptor()] })
  1128. await expect(ctx.remote.probe.create('agent-1', { objective: 'ship' })).resolves.toEqual({
  1129. ok: false,
  1130. error: {
  1131. code: 'internal',
  1132. message: 'client api: probe/create failed: carrier offline',
  1133. details: {},
  1134. },
  1135. })
  1136. })
  1137. it('folds a carrier throw that is not an Error into the error branch', async () => {
  1138. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>()
  1139. .mockRejectedValue('carrier exploded'))
  1140. await ctx.remote.$mount({ package: '@fixture/probe', descriptors: [directDescriptor()] })
  1141. await expect(ctx.remote.probe.create('agent-1', { objective: 'ship' })).resolves.toEqual({
  1142. ok: false,
  1143. error: {
  1144. code: 'internal',
  1145. message: 'client api: probe/create failed: carrier exploded',
  1146. details: {},
  1147. },
  1148. })
  1149. })
  1150. it('owns each $on subscription in the calling fiber', async () => {
  1151. const { ctx, client, carrier } = await eventBench()
  1152. const seen: string[] = []
  1153. const subscriber = ctx.plugin(Object.assign(
  1154. (scope: Context) => { scope.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) }) },
  1155. { inject: ['remote'] },
  1156. ))
  1157. await subscriber
  1158. expect(carrier.calls).toEqual([expect.objectContaining({
  1159. channel: '/api', endpoint: '$events', payload: { args: {} },
  1160. })])
  1161. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['settings'] })
  1162. await vi.waitFor(() => { expect(seen).toEqual(['settings']) })
  1163. await subscriber.dispose()
  1164. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['after fiber disposal'] })
  1165. await Promise.resolve()
  1166. expect(seen).toEqual(['settings'])
  1167. await client.dispose()
  1168. expect(ctx.get('remote')).toBeUndefined()
  1169. })
  1170. it('isolates throwing and rejected notification listeners', async () => {
  1171. const { ctx, client, carrier } = await eventBench()
  1172. const consoleError = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  1173. const seen: string[] = []
  1174. const failingListener = (namespace: string): unknown => {
  1175. if (namespace === 'sync') throw new Error('fixture listener failure')
  1176. return Promise.reject(new Error('fixture async failure'))
  1177. }
  1178. const disposeThrowing = ctx.remote.$on('fixture/changed', failingListener)
  1179. ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
  1180. try {
  1181. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['sync'] })
  1182. await vi.waitFor(() => { expect(seen).toEqual(['sync']) })
  1183. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['async'] })
  1184. await vi.waitFor(() => { expect(seen).toEqual(['sync', 'async']) })
  1185. expect(consoleError).toHaveBeenCalledTimes(2)
  1186. disposeThrowing()
  1187. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['survivor'] })
  1188. await vi.waitFor(() => { expect(seen).toEqual(['sync', 'async', 'survivor']) })
  1189. } finally {
  1190. consoleError.mockRestore()
  1191. await client.dispose()
  1192. }
  1193. })
  1194. it('retires only its own registration when one listener subscribes twice', async () => {
  1195. const { ctx, client, carrier } = await eventBench()
  1196. const seen: string[] = []
  1197. const listener = (namespace: string): void => { seen.push(namespace) }
  1198. const disposeFirst = ctx.remote.$on('fixture/changed', listener)
  1199. ctx.remote.$on('fixture/changed', listener)
  1200. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['both'] })
  1201. await vi.waitFor(() => { expect(seen).toEqual(['both', 'both']) })
  1202. disposeFirst()
  1203. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['survivor'] })
  1204. await vi.waitFor(() => { expect(seen).toEqual(['both', 'both', 'survivor']) })
  1205. disposeFirst()
  1206. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['still here'] })
  1207. await vi.waitFor(() => {
  1208. expect(seen).toEqual(['both', 'both', 'survivor', 'still here'])
  1209. })
  1210. await client.dispose()
  1211. })
  1212. it('keeps the carrier handoff private', () => {
  1213. expectTypeOf<ClientRemote>().toHaveProperty('$on')
  1214. expectTypeOf<'$dispatch' extends keyof ClientRemote ? true : false>().toEqualTypeOf<false>()
  1215. })
  1216. it('drops an unobserved notification and accepts a null-prototype frame', async () => {
  1217. const { ctx, client, carrier } = await eventBench()
  1218. const seen: string[] = []
  1219. ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
  1220. carrier.emit({ type: 'emit', event: 'fixture/idle', args: [1] })
  1221. carrier.emit(Object.assign(Object.create(null) as Record<string, unknown>, {
  1222. type: 'emit',
  1223. event: 'fixture/changed',
  1224. args: ['null prototype'],
  1225. }))
  1226. await vi.waitFor(() => { expect(seen).toEqual(['null prototype']) })
  1227. await client.dispose()
  1228. })
  1229. it('delegates immediately when the Agent adapter or Context is unavailable', async () => {
  1230. const { ctx, client, carrier, call } = await eventBench()
  1231. carrier.emit(approvalFrame('event-no-adapter', 'agent-late', 'no adapter'))
  1232. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(1) })
  1233. const target = ctx.extend()
  1234. const resolve = vi.fn((id: unknown) => id === 'agent-found' ? target : undefined)
  1235. ctx.typert.contexts.registerClient('agent', {
  1236. identity: candidate => candidate === target ? agentId('agent-found') : undefined,
  1237. resolve,
  1238. })
  1239. carrier.emit(approvalFrame('event-missing-context', 'agent-missing', 'missing'))
  1240. carrier.emit(approvalFrame('event-no-listener', 'agent-found', 'delegate'))
  1241. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(3) })
  1242. expect(resolve).toHaveBeenCalledTimes(2)
  1243. for (const eventId of ['event-no-adapter', 'event-missing-context', 'event-no-listener']) {
  1244. expect(call).toHaveBeenCalledWith(
  1245. '/api',
  1246. '$events/result',
  1247. { args: { clientId: 'event-client-1', eventId, outcome: { kind: 'next' } } },
  1248. expect.any(AbortSignal),
  1249. )
  1250. }
  1251. await client.dispose()
  1252. })
  1253. it('reports Agent Context resolution failures and delegates', async () => {
  1254. const { ctx, client, carrier, call } = await eventBench()
  1255. ctx.typert.contexts.registerClient('agent', {
  1256. identity: () => undefined,
  1257. resolve: () => { throw new Error('fixture Context lookup failed') },
  1258. })
  1259. const consoleError = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  1260. try {
  1261. carrier.emit(approvalFrame('event-resolve-error', 'agent-error', 'resolve'))
  1262. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(1) })
  1263. expect(consoleError).toHaveBeenCalledWith(
  1264. 'client api: Remote event "fixture/approval" listener threw:',
  1265. expect.objectContaining({ message: 'fixture Context lookup failed' }),
  1266. )
  1267. expect(call).toHaveBeenCalledWith(
  1268. '/api',
  1269. '$events/result',
  1270. {
  1271. args: {
  1272. clientId: 'event-client-1',
  1273. eventId: 'event-resolve-error',
  1274. outcome: { kind: 'next' },
  1275. },
  1276. },
  1277. expect.any(AbortSignal),
  1278. )
  1279. } finally {
  1280. consoleError.mockRestore()
  1281. await client.dispose()
  1282. }
  1283. })
  1284. it('normalizes undefined waterfall results and rejects non-JSON results', async () => {
  1285. const { ctx, client, carrier, call } = await eventBench()
  1286. const target = ctx.extend()
  1287. ctx.typert.contexts.registerClient('agent', {
  1288. identity: candidate => candidate === target ? agentId('agent-results') : undefined,
  1289. resolve: id => id === 'agent-results' ? target : undefined,
  1290. })
  1291. target.remote.$on('fixture/approval', async request => request.prompt === 'undefined'
  1292. ? undefined as unknown as FixtureApprovalOutcome
  1293. : Symbol('not JSON') as unknown as FixtureApprovalOutcome)
  1294. carrier.emit(approvalFrame('event-undefined', 'agent-results', 'undefined'))
  1295. carrier.emit(approvalFrame('event-invalid', 'agent-results', 'invalid'))
  1296. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(2) })
  1297. expect(call).toHaveBeenCalledWith(
  1298. '/api',
  1299. '$events/result',
  1300. { args: { clientId: 'event-client-1', eventId: 'event-undefined', outcome: { kind: 'result' } } },
  1301. expect.any(AbortSignal),
  1302. )
  1303. expect(call).toHaveBeenCalledWith(
  1304. '/api',
  1305. '$events/result',
  1306. {
  1307. args: {
  1308. clientId: 'event-client-1',
  1309. eventId: 'event-invalid',
  1310. outcome: {
  1311. kind: 'rejected',
  1312. error: {
  1313. name: 'TypeError',
  1314. message: 'Remote event listener result is not lossless JSON data',
  1315. },
  1316. },
  1317. },
  1318. },
  1319. expect.any(AbortSignal),
  1320. )
  1321. await client.dispose()
  1322. })
  1323. it('fails the Connection generation when a result RPC is rejected', async () => {
  1324. const call = vi.fn<ConnectionHandle['rpc']['call']>().mockResolvedValue({
  1325. ok: false,
  1326. error: { code: 'internal', message: 'fixture result rejected', details: {} },
  1327. })
  1328. const { client, carrier, run } = await eventBench(call)
  1329. carrier.emit(approvalFrame('event-result-rejected', 'agent-missing', 'respond'))
  1330. await expect(run.done).rejects.toThrow('fixture result rejected')
  1331. await client.dispose()
  1332. })
  1333. it('filters scoped waterfall listeners and returns the first claimed result', async () => {
  1334. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  1335. .mockResolvedValue({ ok: true, value: undefined })
  1336. const { ctx, client, carrier } = await eventBench(call)
  1337. const target = ctx.extend({
  1338. [Context.filter](candidate: Context): boolean {
  1339. const tag = (candidate as Context & { [fixtureContextTag]?: string })[fixtureContextTag]
  1340. return tag === undefined || tag === 'agent-1'
  1341. },
  1342. })
  1343. ctx.typert.contexts.registerClient('agent', {
  1344. identity: candidate => candidate === target ? agentId('agent-1') : undefined,
  1345. resolve: id => id === 'agent-1' ? target : undefined,
  1346. })
  1347. const matching = ctx.extend({ [fixtureContextTag]: 'agent-1' })
  1348. const excluded = ctx.extend({ [fixtureContextTag]: 'agent-2' })
  1349. const seen: string[] = []
  1350. ctx.remote.$on('fixture/approval', async function (request, next) {
  1351. expect(this).toBe(target)
  1352. expect(request.agent).toBe(target)
  1353. expect(request.signal).toBeInstanceOf(AbortSignal)
  1354. seen.push('root')
  1355. return next()
  1356. })
  1357. matching.remote.$on('fixture/approval', async (_request, next) => {
  1358. seen.push('matching-next')
  1359. return next()
  1360. })
  1361. excluded.remote.$on('fixture/approval', async () => {
  1362. seen.push('excluded')
  1363. return 'unavailable'
  1364. })
  1365. matching.remote.$on('fixture/approval', async () => {
  1366. seen.push('matching-result')
  1367. return 'allowed'
  1368. })
  1369. carrier.emit(approvalFrame('event-1', 'agent-1', 'ship'))
  1370. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(1) })
  1371. expect(seen).toEqual(['root', 'matching-next', 'matching-result'])
  1372. expect(call).toHaveBeenCalledWith(
  1373. '/api',
  1374. '$events/result',
  1375. {
  1376. args: {
  1377. clientId: 'event-client-1',
  1378. eventId: 'event-1',
  1379. outcome: { kind: 'result', value: 'allowed' },
  1380. },
  1381. },
  1382. expect.any(AbortSignal),
  1383. )
  1384. await client.dispose()
  1385. })
  1386. it('returns a scoped listener rejection to the Host', async () => {
  1387. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  1388. .mockResolvedValue({ ok: true, value: undefined })
  1389. const { ctx, client, carrier } = await eventBench(call)
  1390. const target = ctx.extend()
  1391. ctx.typert.contexts.registerClient('agent', {
  1392. identity: candidate => candidate === target ? agentId('agent-rejected') : undefined,
  1393. resolve: id => id === 'agent-rejected' ? target : undefined,
  1394. })
  1395. const rejection = Object.assign(new Error('the user cancelled ask_user_question'), {
  1396. name: 'UserQuestionError',
  1397. code: 'ASK_CANCELLED',
  1398. details: { questionId: 'question-1' },
  1399. })
  1400. target.remote.$on('fixture/approval', () => Promise.reject(rejection))
  1401. carrier.emit(approvalFrame('event-rejected', 'agent-rejected', 'cancelled'))
  1402. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(1) })
  1403. expect(call).toHaveBeenCalledWith(
  1404. '/api',
  1405. '$events/result',
  1406. {
  1407. args: {
  1408. clientId: 'event-client-1',
  1409. eventId: 'event-rejected',
  1410. outcome: {
  1411. kind: 'rejected',
  1412. error: {
  1413. name: 'UserQuestionError',
  1414. message: 'the user cancelled ask_user_question',
  1415. code: 'ASK_CANCELLED',
  1416. details: { questionId: 'question-1' },
  1417. },
  1418. },
  1419. },
  1420. },
  1421. expect.any(AbortSignal),
  1422. )
  1423. await client.dispose()
  1424. })
  1425. it('returns Context-filter failures as rejections', async () => {
  1426. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  1427. .mockResolvedValue({ ok: true, value: undefined })
  1428. const { ctx, client, carrier } = await eventBench(call)
  1429. const target = ctx.extend({
  1430. [Context.filter](): boolean {
  1431. throw new Error('fixture Context filter failed')
  1432. },
  1433. })
  1434. ctx.typert.contexts.registerClient('agent', {
  1435. identity: candidate => candidate === target ? agentId('agent-filter-failure') : undefined,
  1436. resolve: id => id === 'agent-filter-failure' ? target : undefined,
  1437. })
  1438. ctx.remote.$on('fixture/approval', async (_request, next) => next())
  1439. carrier.emit(approvalFrame('event-filter-failure', 'agent-filter-failure', 'filter'))
  1440. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(1) })
  1441. expect(call).toHaveBeenCalledWith(
  1442. '/api',
  1443. '$events/result',
  1444. {
  1445. args: {
  1446. clientId: 'event-client-1',
  1447. eventId: 'event-filter-failure',
  1448. outcome: {
  1449. kind: 'rejected',
  1450. error: {
  1451. name: 'Error',
  1452. message: 'fixture Context filter failed',
  1453. },
  1454. },
  1455. },
  1456. },
  1457. expect.any(AbortSignal),
  1458. )
  1459. await client.dispose()
  1460. })
  1461. it('cancels a pending Client listener without returning a late result', async () => {
  1462. const { ctx, client, carrier, call } = await eventBench()
  1463. const target = ctx.extend()
  1464. ctx.typert.contexts.registerClient('agent', {
  1465. identity: candidate => candidate === target ? agentId('agent-cancel') : undefined,
  1466. resolve: id => id === 'agent-cancel' ? target : undefined,
  1467. })
  1468. const entered = Promise.withResolvers<AbortSignal>()
  1469. target.remote.$on('fixture/approval', async (request) => {
  1470. const signal = request.signal as AbortSignal
  1471. entered.resolve(signal)
  1472. await new Promise<void>((resolve) => {
  1473. if (signal.aborted) resolve()
  1474. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  1475. })
  1476. return 'allowed'
  1477. })
  1478. carrier.emit(approvalFrame('event-cancel', 'agent-cancel', 'wait'))
  1479. const deliverySignal = await entered.promise
  1480. carrier.emit({ type: 'cancel', eventId: 'event-cancel' })
  1481. await vi.waitFor(() => { expect(deliverySignal.aborted).toBe(true) })
  1482. await Promise.resolve()
  1483. expect(call).not.toHaveBeenCalled()
  1484. await client.dispose()
  1485. })
  1486. it('drops a settled listener result when cancellation wins before reply', async () => {
  1487. const { ctx, client, carrier, call } = await eventBench()
  1488. const target = ctx.extend()
  1489. ctx.typert.contexts.registerClient('agent', {
  1490. identity: candidate => candidate === target ? agentId('agent-cancel-race') : undefined,
  1491. resolve: id => id === 'agent-cancel-race' ? target : undefined,
  1492. })
  1493. const entered = Promise.withResolvers<AbortSignal>()
  1494. const release = Promise.withResolvers<undefined>()
  1495. target.remote.$on('fixture/approval', async (request) => {
  1496. entered.resolve(request.signal as AbortSignal)
  1497. await release.promise
  1498. return 'allowed'
  1499. })
  1500. carrier.emit(approvalFrame('event-cancel-race', 'agent-cancel-race', 'wait'))
  1501. const deliverySignal = await entered.promise
  1502. release.resolve(undefined)
  1503. carrier.emit({ type: 'cancel', eventId: 'event-cancel-race' })
  1504. await vi.waitFor(() => { expect(deliverySignal.aborted).toBe(true) })
  1505. await Promise.resolve()
  1506. expect(call).not.toHaveBeenCalled()
  1507. await client.dispose()
  1508. })
  1509. it('cancels pending listener work when the generation ends', async () => {
  1510. const { ctx, client, carrier, run, call } = await eventBench()
  1511. const target = ctx.extend()
  1512. ctx.typert.contexts.registerClient('agent', {
  1513. identity: candidate => candidate === target ? agentId('agent-generation') : undefined,
  1514. resolve: id => id === 'agent-generation' ? target : undefined,
  1515. })
  1516. const entered = Promise.withResolvers<AbortSignal>()
  1517. target.remote.$on('fixture/approval', async (request) => {
  1518. const signal = request.signal as AbortSignal
  1519. entered.resolve(signal)
  1520. await new Promise<void>((resolve) => {
  1521. if (signal.aborted) resolve()
  1522. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  1523. })
  1524. return 'allowed'
  1525. })
  1526. carrier.emit(approvalFrame('event-generation', 'agent-generation', 'wait'))
  1527. const deliverySignal = await entered.promise
  1528. run.abort(new Error('fixture generation ended'))
  1529. await expect(run.done).resolves.toBeUndefined()
  1530. expect(deliverySignal.aborted).toBe(true)
  1531. expect(call).not.toHaveBeenCalled()
  1532. await client.dispose()
  1533. })
  1534. it('contains a result transport failure after the generation is cancelled', async () => {
  1535. const response = Promise.withResolvers<never>()
  1536. const call = vi.fn<ConnectionHandle['rpc']['call']>(() => response.promise)
  1537. const { client, carrier, run } = await eventBench(call)
  1538. carrier.emit(approvalFrame('event-late-result', 'agent-missing', 'respond'))
  1539. await vi.waitFor(() => { expect(call).toHaveBeenCalledOnce() })
  1540. run.abort(new Error('fixture generation cancelled'))
  1541. response.reject(new Error('fixture late result failure'))
  1542. await expect(run.done).resolves.toBeUndefined()
  1543. await client.dispose()
  1544. })
  1545. it('normalizes a non-Error result transport failure', async () => {
  1546. const call = vi.fn<ConnectionHandle['rpc']['call']>().mockRejectedValue('fixture transport failure')
  1547. const { client, carrier, run } = await eventBench(call)
  1548. carrier.emit(approvalFrame('event-result-throw', 'agent-missing', 'respond'))
  1549. await expect(run.done).rejects.toMatchObject({
  1550. message: 'client api: Remote event result delivery failed',
  1551. cause: 'fixture transport failure',
  1552. })
  1553. await client.dispose()
  1554. })
  1555. it('keeps the newer generation tracked when an overlapping generation settles', async () => {
  1556. const { client, generation, run } = await eventBench()
  1557. const overlapping = generation.startOverlapping()
  1558. await overlapping.ready
  1559. run.abort(new Error('fixture older generation ended'))
  1560. await expect(run.done).resolves.toBeUndefined()
  1561. overlapping.abort(new Error('fixture newer generation ended'))
  1562. await expect(overlapping.done).resolves.toBeUndefined()
  1563. await client.dispose()
  1564. })
  1565. it('opens the forwarded-event stream on the browser Remote mux', async () => {
  1566. await withFakeWebSocket('https://harness.example', async () => {
  1567. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  1568. .mockResolvedValue({ ok: true, value: undefined })
  1569. const { ctx, client, generation } = await benchFiber(call, 'web')
  1570. const seen: string[] = []
  1571. const target = ctx.extend()
  1572. ctx.typert.contexts.registerClient('agent', {
  1573. identity: candidate => candidate === target ? agentId('agent-browser') : undefined,
  1574. resolve: id => id === 'agent-browser' ? target : undefined,
  1575. })
  1576. ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
  1577. target.remote.$on('fixture/approval', async function (request) {
  1578. expect(this).toBe(target)
  1579. expect(request.agent).toBe(this)
  1580. expect(request.signal).toBeInstanceOf(AbortSignal)
  1581. return 'allowed'
  1582. })
  1583. const run = generation.start()
  1584. await vi.waitFor(() => { expect(FakeWebSocket.sockets[0]?.sent).toHaveLength(1) })
  1585. const socket = FakeWebSocket.sockets[0]!
  1586. const opened = JSON.parse(socket.sent[0]!) as { streamId: string }
  1587. expect(opened).toMatchObject({
  1588. type: 'open', endpoint: '$events', payload: { args: {} },
  1589. })
  1590. socket.receive({
  1591. type: 'item',
  1592. streamId: opened.streamId,
  1593. value: { type: 'ready', clientId: 'browser-client' },
  1594. })
  1595. await run.ready
  1596. socket.receive({
  1597. type: 'item',
  1598. streamId: opened.streamId,
  1599. value: { type: 'emit', event: 'fixture/changed', args: ['browser'] },
  1600. })
  1601. await vi.waitFor(() => { expect(seen).toEqual(['browser']) })
  1602. socket.receive({
  1603. type: 'item',
  1604. streamId: opened.streamId,
  1605. value: {
  1606. type: 'waterfall',
  1607. event: 'fixture/approval',
  1608. eventId: 'event-browser',
  1609. agentId: 'agent-browser',
  1610. request: { prompt: 'browser approval' },
  1611. },
  1612. })
  1613. await vi.waitFor(() => { expect(call).toHaveBeenCalledTimes(1) })
  1614. expect(socket.sent).toHaveLength(1)
  1615. expect(call).toHaveBeenCalledWith(
  1616. '/api',
  1617. '$events/result',
  1618. {
  1619. args: {
  1620. clientId: 'browser-client',
  1621. eventId: 'event-browser',
  1622. outcome: { kind: 'result', value: 'allowed' },
  1623. },
  1624. },
  1625. expect.any(AbortSignal),
  1626. )
  1627. await client.dispose()
  1628. })
  1629. })
  1630. it('publishes the Fixture Host description after Remote events report ready', async () => {
  1631. const locationDescriptor = Object.getOwnPropertyDescriptor(globalThis, 'location')
  1632. Object.defineProperty(globalThis, 'location', {
  1633. configurable: true,
  1634. value: { hostname: '127.0.0.1', search: '?fixture' },
  1635. })
  1636. const ctx = new Context()
  1637. try {
  1638. await ctx.plugin(TypertRegistry)
  1639. await ctx.plugin({ inject: [], apply: applyConnection })
  1640. await ctx.plugin({ inject, apply })
  1641. const connection = ctx.get('connection') as ConnectionHandle | undefined
  1642. if (connection === undefined) throw new Error('fixture Connection service is unavailable')
  1643. await vi.waitFor(() => {
  1644. expect(connection.hostDescription.getSnapshot()?.home).toBe('/home/fixture')
  1645. })
  1646. } finally {
  1647. await ctx.fiber.dispose()
  1648. if (locationDescriptor === undefined) Reflect.deleteProperty(globalThis, 'location')
  1649. else Object.defineProperty(globalThis, 'location', locationDescriptor)
  1650. }
  1651. })
  1652. it.each([
  1653. null,
  1654. [],
  1655. {},
  1656. { type: 'pending' },
  1657. { type: 'ready' },
  1658. { type: 'ready', clientId: '' },
  1659. { type: 'ready', clientId: 'client', extra: true },
  1660. { type: 'emit', event: 'fixture/changed', args: ['too early'] },
  1661. ])('rejects malformed forwarded-event readiness item %#', async (opening) => {
  1662. const open: NonNullable<ConnectionHandle['rpc']['open']> = () => (async function *() {
  1663. yield opening
  1664. })()
  1665. const { client, generation } = await benchFiber(
  1666. vi.fn<ConnectionHandle['rpc']['call']>(),
  1667. 'in-process',
  1668. open,
  1669. )
  1670. const run = generation.start()
  1671. try {
  1672. await expect(run.done).rejects.toThrow('forwarded Remote event stream did not begin with ready')
  1673. } finally {
  1674. await client.dispose()
  1675. }
  1676. })
  1677. it('propagates physical carrier failure and opens events for the replacement generation', async () => {
  1678. const { ctx, client, carrier, generation, run } = await eventBench()
  1679. const seen: string[] = []
  1680. ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
  1681. expect(carrier.calls).toHaveLength(1)
  1682. carrier.fail(new RemoteStreamCarrierError('fixture generation lost'))
  1683. await expect(run.done).rejects.toThrow('fixture generation lost')
  1684. const replacement = generation.start()
  1685. await replacement.ready
  1686. await vi.waitFor(() => { expect(carrier.calls).toHaveLength(2) })
  1687. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['replacement'] })
  1688. await vi.waitFor(() => { expect(seen).toEqual(['replacement']) })
  1689. await client.dispose()
  1690. })
  1691. it.each([
  1692. {
  1693. name: 'Host failure',
  1694. stop: (carrier: RemoteEventCarrier) => {
  1695. carrier.fail(new RemoteStreamError('internal', 'fixture Host failed', {}))
  1696. },
  1697. message: 'fixture Host failed',
  1698. },
  1699. {
  1700. name: 'normal end',
  1701. stop: (carrier: RemoteEventCarrier) => { carrier.end() },
  1702. message: 'forwarded Remote event stream ended unexpectedly',
  1703. },
  1704. ])('fails the active generation after $name', async ({ stop, message }) => {
  1705. const { ctx, client, carrier, run } = await eventBench()
  1706. const seen: string[] = []
  1707. ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
  1708. stop(carrier)
  1709. await expect(run.done).rejects.toThrow(message)
  1710. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['too late'] })
  1711. await Promise.resolve()
  1712. expect(carrier.calls).toHaveLength(1)
  1713. expect(seen).toEqual([])
  1714. await client.dispose()
  1715. })
  1716. it.each([
  1717. 'not an object',
  1718. null,
  1719. [],
  1720. {},
  1721. { type: 'unknown' },
  1722. { type: 'emit', event: 'fixture/changed' },
  1723. { type: 'emit', event: 'fixture/changed', args: [], extra: true },
  1724. { type: 'emit', event: 1, args: [] },
  1725. { type: 'emit', event: '', args: [] },
  1726. { type: 'emit', event: 'fixture/changed', args: {} },
  1727. { type: 'emit', event: 'fixture/changed', args: [1n] },
  1728. { type: 'waterfall', event: 'fixture/approval', eventId: '', agentId: 'agent-1', request: {} },
  1729. { type: 'waterfall', event: 'fixture/approval', eventId: 'event-1', agentId: '', request: {} },
  1730. {
  1731. type: 'waterfall', event: 'fixture/approval', eventId: 'event-1', agentId: 'agent-1', request: { agent: null },
  1732. },
  1733. {
  1734. type: 'waterfall', event: 'fixture/approval', eventId: 'event-1', agentId: 'agent-1', request: { signal: null },
  1735. },
  1736. { type: 'cancel', eventId: '' },
  1737. { type: 'cancel', eventId: 'event-1', extra: true },
  1738. ])('rejects malformed forwarded-event frame %# and stops that stream', async (frame) => {
  1739. const { ctx, client, carrier, run } = await eventBench()
  1740. const seen: string[] = []
  1741. ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
  1742. carrier.emit(frame)
  1743. await expect(run.done).rejects.toThrow('client api: invalid forwarded Remote event frame')
  1744. carrier.emit({ type: 'emit', event: 'fixture/changed', args: ['too late'] })
  1745. await Promise.resolve()
  1746. expect(carrier.calls).toHaveLength(1)
  1747. expect(seen).toEqual([])
  1748. await client.dispose()
  1749. })
  1750. it('aborts and awaits forwarded-event delivery during disposal', async () => {
  1751. const { ctx, client, carrier, run } = await eventBench()
  1752. ctx.remote.$on('fixture/changed', () => {})
  1753. expect(carrier.activeConnections).toBe(1)
  1754. const signal = carrier.calls[0]?.signal
  1755. await client.dispose()
  1756. await expect(run.done).resolves.toBeUndefined()
  1757. expect(signal?.aborted).toBe(true)
  1758. expect(carrier.activeConnections).toBe(0)
  1759. expect(ctx.get('remote')).toBeUndefined()
  1760. })
  1761. it('rejects a generation when its Connection has been withdrawn', async () => {
  1762. const carrier = new RemoteEventCarrier()
  1763. const { ctx, client, generation } = await benchFiber(
  1764. vi.fn<ConnectionHandle['rpc']['call']>(),
  1765. 'in-process',
  1766. carrier.open,
  1767. )
  1768. ctx.set('connection', undefined)
  1769. const run = generation.start()
  1770. await expect(run.done).rejects.toThrow('$events has no active Connection')
  1771. expect(carrier.calls).toEqual([])
  1772. await client.dispose()
  1773. })
  1774. it('guards stream iteration across mount and Connection withdrawal', async () => {
  1775. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  1776. const ctx = await bench(call)
  1777. const firstDispose = await ctx.remote.$mount({
  1778. package: '@fixture/stream-first', descriptors: [streamDescriptor()],
  1779. })
  1780. const withdrawn = ctx.remote.probe.watch('withdrawn')[Symbol.asyncIterator]()
  1781. await firstDispose()
  1782. await expect(withdrawn.next()).rejects.toThrow('Remote method probe/watch is no longer mounted')
  1783. const secondDispose = await ctx.remote.$mount({
  1784. package: '@fixture/stream-second', descriptors: [streamDescriptor()],
  1785. })
  1786. ctx.set('connection', undefined)
  1787. await expect(ctx.remote.probe.watch('offline')[Symbol.asyncIterator]().next())
  1788. .rejects.toThrow('probe/watch has no active Connection')
  1789. let release!: () => void
  1790. const released = new Promise<void>((resolve) => { release = resolve })
  1791. let markStarted!: () => void
  1792. const started = new Promise<void>((resolve) => { markStarted = resolve })
  1793. const source = async function *(): AsyncIterable<string> {
  1794. markStarted()
  1795. await released
  1796. yield 'late item'
  1797. }
  1798. ctx.set('connection', {
  1799. rpc: { call, open: () => source() },
  1800. } as unknown as ConnectionHandle)
  1801. const active = ctx.remote.probe.watch('active')[Symbol.asyncIterator]()
  1802. const pending = active.next()
  1803. await started
  1804. await secondDispose()
  1805. release()
  1806. await expect(pending).rejects.toThrow('Remote method probe/watch is no longer mounted')
  1807. })
  1808. it('publishes a namespace only after every contributed method is installed', async () => {
  1809. const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
  1810. let visible: string[] | undefined
  1811. const consumer = ctx.plugin({
  1812. inject: ['remote.probe'],
  1813. apply(scope) {
  1814. const namespace = scope.get('remote.probe') as unknown as Record<string, unknown>
  1815. visible = [typeof namespace.watch, typeof namespace.archive]
  1816. },
  1817. })
  1818. const archive: InvocationDescriptor = {
  1819. ...streamDescriptor(),
  1820. id: '@fixture/probe#probe/archive',
  1821. method: 'archive',
  1822. }
  1823. const dispose = await ctx.remote.$mount({
  1824. package: '@fixture/atomic-namespace',
  1825. descriptors: [streamDescriptor(), archive],
  1826. })
  1827. await consumer.await()
  1828. expect(visible).toEqual(['function', 'function'])
  1829. await dispose()
  1830. })
  1831. it('normalizes worker-local structural stream failures without sharing class identity', async () => {
  1832. const cases = [{
  1833. failure: Object.assign(new Error('fixture Host rejected the stream'), {
  1834. dshRemoteStreamFailure: {
  1835. kind: 'remote' as const,
  1836. code: 'fixture-rejected',
  1837. details: { retry: false },
  1838. },
  1839. }),
  1840. assert: (error: unknown) => {
  1841. expect(error).toBeInstanceOf(RemoteStreamError)
  1842. expect(error).toMatchObject({
  1843. code: 'fixture-rejected',
  1844. message: 'fixture Host rejected the stream',
  1845. details: { retry: false },
  1846. })
  1847. },
  1848. }, {
  1849. failure: Object.assign(new Error('worker carrier stopped'), {
  1850. dshRemoteStreamFailure: { kind: 'carrier' as const },
  1851. }),
  1852. assert: (error: unknown) => {
  1853. expect(error).toBeInstanceOf(RemoteStreamCarrierError)
  1854. expect(error).toMatchObject({ message: 'worker carrier stopped' })
  1855. },
  1856. }, {
  1857. failure: 'caller abort sentinel',
  1858. assert: (error: unknown) => { expect(error).toBe('caller abort sentinel') },
  1859. }]
  1860. for (const testCase of cases) {
  1861. const open: NonNullable<ConnectionHandle['rpc']['open']> = () => (async function *(): AsyncGenerator {
  1862. throw testCase.failure
  1863. })()
  1864. const { ctx, client } = await benchFiber(
  1865. vi.fn<ConnectionHandle['rpc']['call']>(),
  1866. 'in-process',
  1867. open,
  1868. )
  1869. const dispose = await ctx.remote.$mount({ package: '@fixture/worker-stream', descriptors: [streamDescriptor()] })
  1870. try {
  1871. const error = await ctx.remote.probe.watch('failure')[Symbol.asyncIterator]().next()
  1872. .then(() => undefined, (reason: unknown) => reason)
  1873. testCase.assert(error)
  1874. } finally {
  1875. await dispose()
  1876. await client.dispose()
  1877. }
  1878. }
  1879. })
  1880. it('multiplexes Remote streams without using the Connection RPC caller', async () => {
  1881. const originalWebSocket = globalThis.WebSocket
  1882. const locationDescriptor = Object.getOwnPropertyDescriptor(globalThis, 'location')
  1883. ;(globalThis as WebSocketGlobal).WebSocket = FakeWebSocket as unknown as typeof WebSocket
  1884. Object.defineProperty(globalThis, 'location', {
  1885. configurable: true,
  1886. value: { origin: 'https://harness.example' },
  1887. })
  1888. FakeWebSocket.sockets.length = 0
  1889. const call = vi.fn<ConnectionHandle['rpc']['call']>()
  1890. const ctx = await bench(call, 'web')
  1891. expect(FakeWebSocket.sockets).toHaveLength(1)
  1892. const dispose = await ctx.remote.$mount({ package: '@fixture/stream', descriptors: [streamDescriptor()] })
  1893. try {
  1894. const first = ctx.remote.probe.watch('alpha')[Symbol.asyncIterator]()
  1895. const firstItem = first.next()
  1896. await vi.waitFor(() => { expect(FakeWebSocket.sockets[0]?.sent).toHaveLength(1) })
  1897. const socket = FakeWebSocket.sockets[0]!
  1898. expect(socket.url).toBe('wss://harness.example/api/remote.mux')
  1899. const opened = JSON.parse(socket.sent[0]!) as { streamId: string }
  1900. expect(opened).toMatchObject({
  1901. type: 'open',
  1902. endpoint: 'probe/watch',
  1903. payload: { args: { topic: 'alpha' } },
  1904. })
  1905. socket.receive({ type: 'item', streamId: opened.streamId, value: 'alpha:one' })
  1906. await expect(firstItem).resolves.toEqual({ done: false, value: 'alpha:one' })
  1907. const firstEnd = first.next()
  1908. socket.receive({ type: 'end', streamId: opened.streamId })
  1909. await expect(firstEnd).resolves.toEqual({ done: true, value: undefined })
  1910. const failed = ctx.remote.probe.watch('failure')[Symbol.asyncIterator]()
  1911. const failedItem = failed.next()
  1912. await vi.waitFor(() => { expect(socket.sent).toHaveLength(2) })
  1913. const failedOpen = JSON.parse(socket.sent[1]!) as { streamId: string }
  1914. socket.receive({
  1915. type: 'error',
  1916. streamId: failedOpen.streamId,
  1917. error: {
  1918. code: 'lookup-unavailable',
  1919. message: 'fixture stream failed',
  1920. details: { lookup: 'missing' },
  1921. },
  1922. })
  1923. await expect(failedItem).rejects.toMatchObject({
  1924. name: 'RemoteStreamError',
  1925. code: 'lookup-unavailable',
  1926. message: 'fixture stream failed',
  1927. details: { lookup: 'missing' },
  1928. })
  1929. const abort = new AbortController()
  1930. const cancelled = ctx.remote.probe.watch('cancel', abort.signal)[Symbol.asyncIterator]()
  1931. const cancelledItem = cancelled.next()
  1932. await vi.waitFor(() => { expect(socket.sent).toHaveLength(3) })
  1933. const cancelledOpen = JSON.parse(socket.sent[2]!) as { streamId: string }
  1934. const cancellation = new Error('caller cancelled')
  1935. socket.receive({ type: 'item', streamId: cancelledOpen.streamId, value: 'already queued' })
  1936. abort.abort(cancellation)
  1937. socket.receive({ type: 'item', streamId: cancelledOpen.streamId, value: 'after cancellation' })
  1938. await expect(cancelledItem).rejects.toBe(cancellation)
  1939. await vi.waitFor(() => {
  1940. expect(socket.sent.map(text => JSON.parse(text) as unknown)).toContainEqual({
  1941. type: 'cancel', streamId: cancelledOpen.streamId,
  1942. })
  1943. })
  1944. expect(call).not.toHaveBeenCalled()
  1945. } finally {
  1946. await dispose()
  1947. await ctx.fiber.dispose()
  1948. FakeWebSocket.sockets.length = 0
  1949. FakeWebSocket.autoOpen = true
  1950. FakeWebSocket.dispatchClose = true
  1951. if (originalWebSocket === undefined) delete (globalThis as WebSocketGlobal).WebSocket
  1952. else globalThis.WebSocket = originalWebSocket
  1953. if (locationDescriptor === undefined) Reflect.deleteProperty(globalThis, 'location')
  1954. else Object.defineProperty(globalThis, 'location', locationDescriptor)
  1955. }
  1956. })
  1957. })
  1958. describe('Remote stream client carrier lifecycle', () => {
  1959. it('connects without a logical stream, reconnects after failures, and stops permanently', async () => {
  1960. await withFakeWebSocket('https://harness.example', async () => {
  1961. FakeWebSocket.autoOpen = false
  1962. vi.useFakeTimers()
  1963. const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
  1964. try {
  1965. const client = new RemoteStreamMuxClient()
  1966. client.start()
  1967. client.start()
  1968. expect(FakeWebSocket.sockets).toHaveLength(1)
  1969. const failed = FakeWebSocket.sockets[0]!
  1970. failed.fail()
  1971. await vi.advanceTimersByTimeAsync(500)
  1972. expect(FakeWebSocket.sockets).toHaveLength(2)
  1973. const connected = FakeWebSocket.sockets[1]!
  1974. connected.open()
  1975. await vi.advanceTimersByTimeAsync(0)
  1976. expect(connected.sent).toEqual([])
  1977. connected.fail()
  1978. await vi.advanceTimersByTimeAsync(500)
  1979. expect(FakeWebSocket.sockets).toHaveLength(3)
  1980. const replacement = FakeWebSocket.sockets[2]!
  1981. replacement.open()
  1982. replacement.drop()
  1983. await vi.advanceTimersByTimeAsync(500)
  1984. expect(FakeWebSocket.sockets).toHaveLength(4)
  1985. const final = FakeWebSocket.sockets[3]!
  1986. final.open()
  1987. await vi.advanceTimersByTimeAsync(0)
  1988. await client.close()
  1989. await client.close()
  1990. client.start()
  1991. await expect(client.open('feed/follow', {}, new AbortController().signal)
  1992. [Symbol.asyncIterator]().next()).rejects.toThrow('Remote stream client disposed')
  1993. await vi.advanceTimersByTimeAsync(20_000)
  1994. expect(FakeWebSocket.sockets).toHaveLength(4)
  1995. expect(final.closedWith).toContainEqual({ code: 1000, reason: 'disposed' })
  1996. expect(warn).toHaveBeenCalledTimes(3)
  1997. const stopping = new RemoteStreamMuxClient()
  1998. stopping.start()
  1999. const racing = FakeWebSocket.sockets[4]!
  2000. racing.open()
  2001. racing.drop()
  2002. await stopping.close()
  2003. await vi.advanceTimersByTimeAsync(20_000)
  2004. expect(FakeWebSocket.sockets).toHaveLength(5)
  2005. } finally {
  2006. warn.mockRestore()
  2007. vi.useRealTimers()
  2008. }
  2009. })
  2010. })
  2011. it('shares an in-flight connection and uses the internal ws URL without a browser origin', async () => {
  2012. await withFakeWebSocket(undefined, async () => {
  2013. FakeWebSocket.autoOpen = false
  2014. const client = new RemoteStreamMuxClient()
  2015. const first = client.open('feed/follow', { label: 'first' }, new AbortController().signal)
  2016. [Symbol.asyncIterator]()
  2017. const second = client.open('feed/follow', { label: 'second' }, new AbortController().signal)
  2018. [Symbol.asyncIterator]()
  2019. const firstPending = first.next()
  2020. const secondPending = second.next()
  2021. expect(FakeWebSocket.sockets).toHaveLength(1)
  2022. const socket = FakeWebSocket.sockets[0]!
  2023. expect(socket.url).toBe('ws://dsh.internal/api/remote.mux')
  2024. socket.open()
  2025. await vi.waitFor(() => { expect(socket.sent).toHaveLength(2) })
  2026. const streamIds = socket.sent.map(text => (JSON.parse(text) as { streamId: string }).streamId)
  2027. socket.receive({ type: 'end', streamId: streamIds[0] })
  2028. socket.receive({ type: 'end', streamId: streamIds[1] })
  2029. await expect(firstPending).resolves.toEqual({ done: true, value: undefined })
  2030. await expect(secondPending).resolves.toEqual({ done: true, value: undefined })
  2031. await client.close()
  2032. })
  2033. })
  2034. it('keeps waiters across failed attempts and contains waiter cancellation', async () => {
  2035. await withFakeWebSocket('null', async () => {
  2036. FakeWebSocket.autoOpen = false
  2037. vi.useFakeTimers()
  2038. const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
  2039. try {
  2040. const closedClient = new RemoteStreamMuxClient()
  2041. const closed = closedClient.open('feed/follow', {}, new AbortController().signal)
  2042. [Symbol.asyncIterator]().next()
  2043. FakeWebSocket.sockets[0]!.drop()
  2044. await vi.advanceTimersByTimeAsync(500)
  2045. const replacement = FakeWebSocket.sockets[1]!
  2046. replacement.open()
  2047. await vi.advanceTimersByTimeAsync(0)
  2048. const { streamId } = JSON.parse(replacement.sent[0]!) as { streamId: string }
  2049. replacement.receive({ type: 'end', streamId })
  2050. await expect(closed).resolves.toEqual({ done: true, value: undefined })
  2051. await closedClient.close()
  2052. const disposedClient = new RemoteStreamMuxClient()
  2053. const disposed = disposedClient.open('feed/follow', {}, new AbortController().signal)
  2054. [Symbol.asyncIterator]().next()
  2055. FakeWebSocket.sockets[2]!.fail()
  2056. await disposedClient.close()
  2057. await expect(disposed).rejects.toThrow('Remote stream client disposed')
  2058. const abortedClient = new RemoteStreamMuxClient()
  2059. const abort = new AbortController()
  2060. const aborted = abortedClient.open('feed/follow', {}, abort.signal)[Symbol.asyncIterator]().next()
  2061. abort.abort('cancelled while connecting')
  2062. await expect(aborted).rejects.toBe('cancelled while connecting')
  2063. await abortedClient.close()
  2064. expect(FakeWebSocket.sockets[3]?.url).toBe('ws://dsh.internal/api/remote.mux')
  2065. } finally {
  2066. warn.mockRestore()
  2067. vi.useRealTimers()
  2068. }
  2069. })
  2070. })
  2071. it('fails active streams on an invalid frame and ignores later frames', async () => {
  2072. await withFakeWebSocket('https://harness.example', async () => {
  2073. const client = new RemoteStreamMuxClient()
  2074. const stream = client.open('feed/follow', {}, new AbortController().signal)[Symbol.asyncIterator]()
  2075. const pending = stream.next()
  2076. await vi.waitFor(() => { expect(FakeWebSocket.sockets[0]?.sent).toHaveLength(1) })
  2077. const socket = FakeWebSocket.sockets[0]!
  2078. const { streamId } = JSON.parse(socket.sent[0]!) as { streamId: string }
  2079. FakeWebSocket.dispatchClose = false
  2080. socket.receiveRaw(new Uint8Array([1, 2, 3]))
  2081. socket.receive({ type: 'item', streamId, value: 'too late' })
  2082. socket.drop()
  2083. await expect(pending).rejects.toMatchObject({
  2084. name: 'RemoteStreamCarrierError', message: 'api gateway: invalid Remote stream frame',
  2085. })
  2086. expect(socket.closedWith).toContainEqual({ code: 4002, reason: 'invalid Remote stream frame' })
  2087. await client.close()
  2088. })
  2089. })
  2090. it('completes a stream and drops a frame racing with cancellation', async () => {
  2091. await withFakeWebSocket('https://harness.example', async () => {
  2092. const client = new RemoteStreamMuxClient()
  2093. const completed = client.open('feed/follow', {}, new AbortController().signal)
  2094. [Symbol.asyncIterator]()
  2095. const completedPending = completed.next()
  2096. await vi.waitFor(() => { expect(FakeWebSocket.sockets[0]?.sent).toHaveLength(1) })
  2097. const socket = FakeWebSocket.sockets[0]!
  2098. const completedOpen = JSON.parse(socket.sent[0]!) as { streamId: string }
  2099. socket.receive({ type: 'end', streamId: completedOpen.streamId })
  2100. await expect(completedPending).resolves.toEqual({ done: true, value: undefined })
  2101. const abort = new AbortController()
  2102. const cancelled = client.open('feed/follow', {}, abort.signal)[Symbol.asyncIterator]().next()
  2103. await vi.waitFor(() => { expect(socket.sent).toHaveLength(2) })
  2104. const cancelledOpen = JSON.parse(socket.sent[1]!) as { streamId: string }
  2105. const reason = new Error('fixture cancellation race')
  2106. abort.abort(reason)
  2107. socket.receive({ type: 'item', streamId: cancelledOpen.streamId, value: 'too late' })
  2108. await expect(cancelled).rejects.toBe(reason)
  2109. await client.close()
  2110. })
  2111. })
  2112. it('contains non-Error cancellation reasons and late socket close events', async () => {
  2113. await withFakeWebSocket('http://harness.example', async () => {
  2114. const cancelledClient = new RemoteStreamMuxClient()
  2115. const abort = new AbortController()
  2116. const cancelled = cancelledClient.open('feed/follow', {}, abort.signal)[Symbol.asyncIterator]().next()
  2117. await vi.waitFor(() => { expect(FakeWebSocket.sockets[0]?.sent).toHaveLength(1) })
  2118. abort.abort('caller cancelled')
  2119. await expect(cancelled).rejects.toThrow('caller cancelled')
  2120. await cancelledClient.close()
  2121. FakeWebSocket.dispatchClose = false
  2122. const disposedClient = new RemoteStreamMuxClient()
  2123. const disposed = disposedClient.open('feed/follow', {}, new AbortController().signal)
  2124. [Symbol.asyncIterator]().next()
  2125. await vi.waitFor(() => { expect(FakeWebSocket.sockets[1]?.sent).toHaveLength(1) })
  2126. const disposedSocket = FakeWebSocket.sockets[1]!
  2127. await disposedClient.close()
  2128. disposedSocket.receive({ type: 'end', streamId: 'stale' })
  2129. disposedSocket.drop()
  2130. await expect(disposed).rejects.toThrow('Remote stream client disposed')
  2131. })
  2132. })
  2133. })
  2134. async function withFakeWebSocket(
  2135. origin: string | undefined,
  2136. run: () => Promise<void>,
  2137. ): Promise<void> {
  2138. const originalWebSocket = globalThis.WebSocket
  2139. const locationDescriptor = Object.getOwnPropertyDescriptor(globalThis, 'location')
  2140. ;(globalThis as WebSocketGlobal).WebSocket = FakeWebSocket as unknown as typeof WebSocket
  2141. if (origin === undefined) Reflect.deleteProperty(globalThis, 'location')
  2142. else Object.defineProperty(globalThis, 'location', { configurable: true, value: { origin } })
  2143. FakeWebSocket.sockets.length = 0
  2144. FakeWebSocket.autoOpen = true
  2145. FakeWebSocket.dispatchClose = true
  2146. try {
  2147. await run()
  2148. } finally {
  2149. FakeWebSocket.sockets.length = 0
  2150. FakeWebSocket.autoOpen = true
  2151. FakeWebSocket.dispatchClose = true
  2152. if (originalWebSocket === undefined) delete (globalThis as WebSocketGlobal).WebSocket
  2153. else globalThis.WebSocket = originalWebSocket
  2154. if (locationDescriptor === undefined) Reflect.deleteProperty(globalThis, 'location')
  2155. else Object.defineProperty(globalThis, 'location', locationDescriptor)
  2156. }
  2157. }