client-apply.spec.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398
  1. /**
  2. * Connection plugin browser-half apply: ctx.connection handle mounting, mode
  3. * selection off the page URL, and the single-consumer stream-loop ownership.
  4. */
  5. import { Context } from '@deepseek-ai/cordis'
  6. import { afterEach, describe, expect, it, vi } from 'vitest'
  7. import { apply, type ConnectionHandle } from '../src/client/index.ts'
  8. import type { RpcMessage } from '../src/client/api.ts'
  9. import { RpcId } from '../src/client/api.ts'
  10. import { FixtureApiClient } from '../src/client/fixture.ts'
  11. import { WebApiClient } from '../src/client/web-api-client.ts'
  12. type Win = { location?: { hostname: string; search: string; origin?: string } }
  13. type WebSocketGlobal = { WebSocket?: typeof WebSocket }
  14. const originalWebSocket = globalThis.WebSocket
  15. const sockets: FakeWebSocket[] = []
  16. class FakeWebSocket extends EventTarget {
  17. static readonly CONNECTING = 0
  18. static readonly OPEN = 1
  19. static readonly CLOSING = 2
  20. static readonly CLOSED = 3
  21. readonly url: string
  22. readyState = FakeWebSocket.CONNECTING
  23. constructor(url: string | URL) {
  24. super()
  25. this.url = String(url)
  26. sockets.push(this)
  27. queueMicrotask(() => {
  28. if (this.readyState !== FakeWebSocket.CONNECTING) return
  29. this.readyState = FakeWebSocket.OPEN
  30. this.dispatchEvent(new Event('open'))
  31. })
  32. }
  33. close(): void {
  34. if (this.readyState === FakeWebSocket.CLOSED) return
  35. this.readyState = FakeWebSocket.CLOSED
  36. this.dispatchEvent(new Event('close'))
  37. }
  38. receive(data: unknown): void {
  39. this.dispatchEvent(new MessageEvent('message', { data }))
  40. }
  41. }
  42. afterEach(() => {
  43. delete (globalThis as Win).location
  44. sockets.length = 0
  45. if (originalWebSocket === undefined) delete (globalThis as WebSocketGlobal).WebSocket
  46. else globalThis.WebSocket = originalWebSocket
  47. })
  48. async function mount(): Promise<ConnectionHandle> {
  49. const ctx = new Context()
  50. await ctx.plugin({ apply, inject: [] })
  51. const handle = ctx.get('connection') as ConnectionHandle | undefined
  52. if (handle === undefined) throw new Error('ctx.connection not provided')
  53. return handle
  54. }
  55. describe('connection client apply', () => {
  56. it('mounts ctx.connection with the real client when no ?fixture switch is present', async () => {
  57. ;(globalThis as Win).location = { hostname: 'localhost', search: '' }
  58. const handle = await mount()
  59. expect(handle.api).toBeInstanceOf(WebApiClient)
  60. expect(handle.isLoopback).toBe(true)
  61. })
  62. it('selects the fixture client under ?fixture (and with no location at all stays real)', async () => {
  63. ;(globalThis as Win).location = { hostname: '127.0.0.1', search: '?fixture' }
  64. expect((await mount()).api).toBeInstanceOf(FixtureApiClient)
  65. delete (globalThis as Win).location
  66. const handle = await mount()
  67. expect(handle.api).toBeInstanceOf(WebApiClient)
  68. expect(handle.isLoopback).toBe(true)
  69. })
  70. it('reports non-loopback page authority through the connection handle', async () => {
  71. ;(globalThis as Win).location = { hostname: '192.0.2.20', search: '' }
  72. expect((await mount()).isLoopback).toBe(false)
  73. })
  74. it('start() hands out one loop, rejects a second consumer, and stop() aborts the streams', async () => {
  75. ;(globalThis as Win).location = { hostname: 'localhost', search: '?fixture' }
  76. const handle = await mount()
  77. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  78. const descriptions: Array<boolean | undefined> = []
  79. const stopThrowing = handle.hostDescription.subscribe(() => { throw new Error('subscriber bug') })
  80. const stopDescription = handle.hostDescription.subscribe(() => {
  81. descriptions.push(handle.hostDescription.getSnapshot()?.canOpenPath)
  82. })
  83. expect(handle.hostDescription.getSnapshot()).toBeUndefined()
  84. // config omitted: the `config ?? {}` default arm is part of the surface.
  85. let connected = 0
  86. const loop = handle.start({ onConnected: () => { connected++ } })
  87. expect(() => handle.start({})).toThrow(/already owned by another consumer/)
  88. await vi.waitFor(() => {
  89. expect(handle.hostDescription.getSnapshot()?.canOpenPath).toBe(true)
  90. })
  91. loop.stop() // teardown must not throw; the fixture streams abort quietly
  92. expect(handle.hostDescription.getSnapshot()).toBeUndefined()
  93. expect(descriptions).toEqual([true, undefined])
  94. expect(connected).toBe(1)
  95. expect(errorSpy).toHaveBeenCalledTimes(2)
  96. stopThrowing()
  97. stopDescription()
  98. errorSpy.mockRestore()
  99. })
  100. it('does not announce a generation synchronously stopped by a description subscriber', async () => {
  101. ;(globalThis as Win).location = { hostname: 'localhost', search: '?fixture' }
  102. const handle = await mount()
  103. const owner: { loop?: ReturnType<ConnectionHandle['start']> } = {}
  104. let sawDescription = false
  105. const stopDescription = handle.hostDescription.subscribe(() => {
  106. if (handle.hostDescription.getSnapshot() === undefined) return
  107. sawDescription = true
  108. owner.loop?.stop()
  109. })
  110. const connected = vi.fn()
  111. const loop = handle.start({ onConnected: connected })
  112. owner.loop = loop
  113. try {
  114. await vi.waitFor(() => { expect(sawDescription).toBe(true) })
  115. expect(handle.hostDescription.getSnapshot()).toBeUndefined()
  116. expect(connected).not.toHaveBeenCalled()
  117. } finally {
  118. stopDescription()
  119. loop.stop()
  120. }
  121. })
  122. it('retracts the host description while reconnecting and republishes the next generation', async () => {
  123. ;(globalThis as Win).location = { hostname: 'localhost', search: '?fixture' }
  124. const handle = await mount()
  125. const descriptions: Array<boolean | undefined> = []
  126. const reconnectSnapshots: Array<boolean | undefined> = []
  127. const stopDescription = handle.hostDescription.subscribe(() => {
  128. descriptions.push(handle.hostDescription.getSnapshot()?.canOpenPath)
  129. })
  130. const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => undefined)
  131. const loop = handle.start({
  132. onStateChange: (state) => {
  133. if (state === 'reconnecting') {
  134. reconnectSnapshots.push(handle.hostDescription.getSnapshot()?.canOpenPath)
  135. }
  136. },
  137. }, { backoffBaseMs: 10, backoffFactor: 1, backoffMaxMs: 10, streamOpenTimeoutMs: 500 })
  138. try {
  139. await vi.waitFor(() => {
  140. expect(handle.hostDescription.getSnapshot()?.canOpenPath).toBe(true)
  141. })
  142. const timing = (globalThis as Record<string, unknown>).__fxTiming as
  143. | { breakStreams(): void }
  144. | undefined
  145. if (timing === undefined) throw new Error('fixture timing hooks missing')
  146. timing.breakStreams()
  147. await vi.waitFor(() => { expect(reconnectSnapshots).toEqual([undefined]) })
  148. await vi.waitFor(() => { expect(descriptions).toEqual([true, undefined, true]) })
  149. expect(handle.hostDescription.getSnapshot()?.canOpenPath).toBe(true)
  150. } finally {
  151. stopDescription()
  152. loop.stop()
  153. warnSpy.mockRestore()
  154. }
  155. })
  156. it('WebApiClient keeps unary calls and respond on globalThis.fetch', async () => {
  157. ;(globalThis as Win).location = { hostname: 'localhost', search: '' }
  158. const handle = await mount()
  159. const original = globalThis.fetch
  160. const seen: string[] = []
  161. globalThis.fetch = (input: URL | RequestInfo) => {
  162. seen.push(typeof input === 'string' ? input : input instanceof URL ? input.href : input.url)
  163. return Promise.resolve(new Response('{}', { status: 200 }))
  164. }
  165. try {
  166. // Schema rejection is fine — the transport hop is the assertion.
  167. await (handle.api as WebApiClient).host.describe({}).catch(() => undefined)
  168. await handle.api.respond({
  169. type: 'client-response',
  170. rpcId: RpcId('response-over-http'),
  171. result: { ok: true, value: {} },
  172. }).catch(() => undefined)
  173. } finally {
  174. globalThis.fetch = original
  175. }
  176. expect(seen.some(u => u.includes('/api/host.describe'))).toBe(true)
  177. expect(seen.some(u => u.includes('/api/respond'))).toBe(true)
  178. })
  179. it('opens one WebSocket per downlink, parses frames, and aborts both without using fetch', async () => {
  180. ;(globalThis as Win).location = {
  181. hostname: 'localhost', search: '', origin: 'http://localhost:3080',
  182. }
  183. ;(globalThis as WebSocketGlobal).WebSocket = FakeWebSocket as unknown as typeof WebSocket
  184. const fetch = vi.spyOn(globalThis, 'fetch')
  185. const client = (await mount()).api as WebApiClient
  186. const envelopes: RpcMessage[][] = []
  187. client.subscribeEnvelopes((batch) => { envelopes.push([...batch]) })
  188. const opened: string[] = []
  189. const muxAbort = new AbortController()
  190. const hostAbort = new AbortController()
  191. const mux = client.events.mux({}, muxAbort.signal, () => { opened.push('mux') })[Symbol.asyncIterator]()
  192. const host = client.events.host({}, hostAbort.signal, () => { opened.push('host') })[Symbol.asyncIterator]()
  193. const muxFrame = mux.next()
  194. const hostFrame = host.next()
  195. await vi.waitFor(() => { expect(sockets).toHaveLength(2) })
  196. expect(sockets.map(socket => socket.url)).toEqual([
  197. 'ws://localhost:3080/api/events.mux',
  198. 'ws://localhost:3080/api/events.host',
  199. ])
  200. await vi.waitFor(() => { expect(opened).toEqual(['mux', 'host']) })
  201. const errors = vi.spyOn(console, 'error').mockImplementation(() => {})
  202. sockets[0]!.receive(new Uint8Array([1, 2, 3]))
  203. sockets[1]!.receive(JSON.stringify({ type: 'server-request', rpcId: 'bad', method: 'host/session-status', payload: {} }))
  204. sockets[0]!.receive(JSON.stringify({
  205. type: 'server-request',
  206. rpcId: 'mux-browser',
  207. method: 'session/subscribed',
  208. payload: { type: 'session/subscribed', sessionId: 'session-browser', lastSeq: 8 },
  209. }))
  210. sockets[1]!.receive(JSON.stringify({
  211. type: 'server-request',
  212. rpcId: 'host-browser',
  213. method: 'host/remote-event',
  214. payload: { type: 'host/remote-event', event: 'commands/change', args: [] },
  215. }))
  216. expect(await muxFrame).toMatchObject({
  217. value: { rpcId: 'mux-browser', payload: { type: 'session/subscribed', lastSeq: 8 } },
  218. })
  219. expect(await hostFrame).toMatchObject({
  220. value: { rpcId: 'host-browser', payload: { type: 'host/remote-event', event: 'commands/change' } },
  221. })
  222. expect(errors).toHaveBeenCalledTimes(2)
  223. await vi.waitFor(() => { expect(envelopes.flat()).toHaveLength(2) })
  224. expect(fetch).not.toHaveBeenCalled()
  225. const muxEnd = mux.next()
  226. const hostEnd = host.next()
  227. muxAbort.abort()
  228. hostAbort.abort()
  229. await expect(muxEnd).resolves.toMatchObject({ done: true })
  230. await expect(hostEnd).resolves.toMatchObject({ done: true })
  231. expect(sockets.every(socket => socket.readyState === FakeWebSocket.CLOSED)).toBe(true)
  232. errors.mockRestore()
  233. fetch.mockRestore()
  234. })
  235. it('maps an HTTPS page origin to a secure WebSocket URL', async () => {
  236. ;(globalThis as Win).location = {
  237. hostname: 'harness.example', search: '', origin: 'https://harness.example',
  238. }
  239. ;(globalThis as WebSocketGlobal).WebSocket = FakeWebSocket as unknown as typeof WebSocket
  240. const client = (await mount()).api
  241. const abort = new AbortController()
  242. const iterator = client.events.mux({}, abort.signal)[Symbol.asyncIterator]()
  243. const pending = iterator.next()
  244. await vi.waitFor(() => { expect(sockets[0]?.url).toBe('wss://harness.example/api/events.mux') })
  245. abort.abort()
  246. await expect(pending).resolves.toMatchObject({ done: true })
  247. })
  248. it('closes a WebSocket immediately when its signal was already aborted', async () => {
  249. ;(globalThis as Win).location = {
  250. hostname: 'localhost', search: '', origin: 'http://localhost:3080',
  251. }
  252. ;(globalThis as WebSocketGlobal).WebSocket = FakeWebSocket as unknown as typeof WebSocket
  253. const client = (await mount()).api
  254. const abort = new AbortController()
  255. abort.abort()
  256. const iterator = client.events.mux({}, abort.signal)[Symbol.asyncIterator]()
  257. await expect(iterator.next()).resolves.toMatchObject({ done: true })
  258. expect(sockets).toHaveLength(1)
  259. expect(sockets[0]?.readyState).toBe(FakeWebSocket.CLOSED)
  260. })
  261. it('carries RPC calls without requiring secure-context randomUUID', async () => {
  262. ;(globalThis as Win).location = { hostname: 'localhost', search: '' }
  263. vi.stubGlobal('crypto', {
  264. getRandomValues(bytes: Uint8Array) {
  265. return bytes.fill(0)
  266. },
  267. })
  268. const handle = await mount()
  269. const original = globalThis.fetch
  270. const seen: { url: string; body: unknown }[] = []
  271. globalThis.fetch = async (input: URL | RequestInfo, init?: RequestInit) => {
  272. const url = typeof input === 'string' ? input : input instanceof URL ? input.href : input.url
  273. if (typeof init?.body !== 'string') throw new TypeError('expected a JSON string request body')
  274. const body = JSON.parse(init.body) as { rpcId: string }
  275. seen.push({ url, body })
  276. return Response.json({
  277. type: 'server-response',
  278. rpcId: body.rpcId,
  279. result: { ok: true, value: { ref: 'goal-1' } },
  280. })
  281. }
  282. try {
  283. await expect(handle.rpc.call('/api', 'goals/create', { args: { agentId: 'agent-1' } }))
  284. .resolves.toEqual({ ok: true, value: { ref: 'goal-1' } })
  285. } finally {
  286. globalThis.fetch = original
  287. vi.unstubAllGlobals()
  288. }
  289. expect(seen).toHaveLength(1)
  290. expect(seen[0]?.url).toBe('http://dsh.internal/api/goals/create')
  291. expect(seen[0]?.body).toMatchObject({
  292. type: 'client-request',
  293. rpcId: '00000000-0000-4000-8000-000000000000',
  294. method: 'goals/create',
  295. payload: { args: { agentId: 'agent-1' } },
  296. })
  297. })
  298. it('validates generic RPC transport failures, correlation, and targets', async () => {
  299. ;(globalThis as Win).location = {
  300. hostname: 'harness.example', search: '', origin: 'https://harness.example',
  301. }
  302. const handle = await mount()
  303. const original = globalThis.fetch
  304. const abort = new AbortController()
  305. globalThis.fetch = vi.fn().mockResolvedValue(new Response('unavailable', { status: 503 }))
  306. try {
  307. await expect(handle.rpc.call('/api', 'goals/create', {}, abort.signal))
  308. .rejects.toThrow('HTTP 503')
  309. expect(globalThis.fetch).toHaveBeenCalledWith(
  310. new URL('https://harness.example/api/goals/create'),
  311. expect.objectContaining({ signal: abort.signal }),
  312. )
  313. ;(globalThis as Win).location = { hostname: 'localhost', search: '', origin: 'null' }
  314. globalThis.fetch = vi.fn().mockResolvedValue(Response.json({
  315. type: 'server-response',
  316. rpcId: 'different-rpc',
  317. result: { ok: true, value: null },
  318. }))
  319. await expect(handle.rpc.call('/api', 'goals/create', {})).rejects.toThrow('rpcId mismatch')
  320. const fetch = vi.mocked(globalThis.fetch)
  321. expect(fetch.mock.calls[0]?.[0]).toEqual(new URL('http://dsh.internal/api/goals/create'))
  322. expect(fetch.mock.calls[0]?.[1]).not.toHaveProperty('signal')
  323. } finally {
  324. globalThis.fetch = original
  325. }
  326. for (const [channel, endpoint] of [
  327. ['api2', 'goals/create'],
  328. ['/api/path', 'goals/create'],
  329. ['/api', ''],
  330. ['/api', '.'],
  331. ['/api', '..'],
  332. ['/api', 'goals//create'],
  333. ['/api', 'goals/create?unsafe'],
  334. ] as const) {
  335. await expect(handle.rpc.call(channel, endpoint, {})).rejects.toThrow('invalid RPC target')
  336. }
  337. })
  338. it('carries Goal Remotes over the same state as the client-only fixture API', async () => {
  339. ;(globalThis as Win).location = { hostname: 'localhost', search: '?fixture' }
  340. const handle = await mount()
  341. const created = await handle.rpc.call('/api', 'goals/create', {
  342. args: { agentId: 'fx-alpha', request: { objective: 'fixture remote' } },
  343. })
  344. expect(created).toMatchObject({ ok: true, value: { ref: { revision: 1 } } })
  345. if (!created.ok) throw new Error('fixture Goal create failed')
  346. const ref = (created.value as { ref: { id: string; revision: number } }).ref
  347. const edited = await handle.rpc.call('/api', 'goals/edit', {
  348. args: { agentId: 'fx-alpha', ref, request: { objective: 'edited fixture remote' } },
  349. })
  350. expect(edited).toMatchObject({ ok: true, value: { objective: 'edited fixture remote', revision: 2 } })
  351. const editedRef = { id: ref.id, revision: 2 }
  352. const paused = await handle.rpc.call('/api', 'goals/pause', {
  353. args: { agentId: 'fx-alpha', ref: editedRef },
  354. })
  355. expect(paused).toMatchObject({ ok: true, value: { phase: 'paused', activation: 'disarmed', revision: 3 } })
  356. const resumed = await handle.rpc.call('/api', 'goals/resume', {
  357. args: { agentId: 'fx-alpha', ref: { id: ref.id, revision: 3 } },
  358. })
  359. expect(resumed).toMatchObject({ ok: true, value: { phase: 'active', activation: 'armed', revision: 4 } })
  360. const completed = await handle.rpc.call('/api', 'goals/complete', {
  361. args: { agentId: 'fx-alpha', ref: { id: ref.id, revision: 4 } },
  362. })
  363. expect(completed).toMatchObject({ ok: true, value: { phase: 'complete', activation: 'disarmed', revision: 5 } })
  364. await expect(handle.rpc.call('/api', 'goals/clear', {
  365. args: { agentId: 'fx-alpha', ref: { id: ref.id, revision: 5 } },
  366. })).resolves.toEqual({ ok: true, value: { id: ref.id, revision: 6 } })
  367. await expect(handle.rpc.call('/other', 'goals/create', {})).rejects.toThrow(/channel.*unavailable/)
  368. await expect(handle.rpc.call('/api', 'unknown/read', { args: { agentId: 'fx-alpha' } }))
  369. .rejects.toThrow(/endpoint.*unavailable/)
  370. })
  371. })