stream-protocol.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399
  1. /** Wire messages for Gateway-owned Remote streams and event-result RPCs. */
  2. import type { Branded } from '@deepseek-ai/dsh-brand'
  3. /** Exact WebSocket route carrying every Typert Remote stream. */
  4. export const REMOTE_STREAM_MUX_PATH = '/api/remote.mux'
  5. /** Gateway-internal logical stream carrying application-selected Cordis events. */
  6. export const REMOTE_EVENT_STREAM_ENDPOINT = '$events'
  7. /** Gateway-internal unary endpoint returning one Client Remote Event outcome. */
  8. export const REMOTE_EVENT_RESULT_ENDPOINT = '$events/result'
  9. /** Empty standard Remote payload used to open the forwarded-event stream. */
  10. export const REMOTE_EVENT_STREAM_PAYLOAD = { args: {} } as const
  11. /** Discriminator for the first item proving the Host event source is ready. */
  12. export const REMOTE_EVENT_STREAM_READY = { type: 'ready' } as const
  13. /** Opaque identity for one active Client Remote Event generation. */
  14. export type RemoteEventClientId = Branded<'RemoteEventClientId'>
  15. /** Opaque correlation id for one pending Host-to-Client Remote Event. */
  16. export type RemoteEventId = Branded<'RemoteEventId'>
  17. /** Opening item that binds later HTTP results to this active event stream. */
  18. export interface RemoteEventReadyFrame {
  19. readonly type: 'ready'
  20. readonly clientId: RemoteEventClientId
  21. }
  22. /** Opaque Agent identity carried by one scoped Remote Event. */
  23. export type RemoteEventAgentId = Branded<'RemoteEventAgentId'>
  24. /** One Host notification delivered to a Client generation. */
  25. export interface RemoteEventEmitFrame {
  26. readonly type: 'emit'
  27. readonly event: string
  28. readonly args: readonly unknown[]
  29. }
  30. /** One pending Agent-scoped waterfall delivered to a Client generation. */
  31. export interface RemoteEventInvocationFrame {
  32. readonly type: 'waterfall'
  33. readonly event: string
  34. readonly eventId: RemoteEventId
  35. readonly agentId: RemoteEventAgentId
  36. readonly request: Readonly<Record<string, unknown>>
  37. }
  38. /** Cancellation of a pending waterfall previously delivered under the same id. */
  39. export interface RemoteEventCancellationFrame {
  40. readonly type: 'cancel'
  41. readonly eventId: RemoteEventId
  42. }
  43. /** Every item carried by the Gateway-internal forwarded-event stream. */
  44. export type RemoteEventDownlinkFrame =
  45. | RemoteEventReadyFrame
  46. | RemoteEventEmitFrame
  47. | RemoteEventInvocationFrame
  48. | RemoteEventCancellationFrame
  49. /** JSON request fields plus the Host cancellation lifetime removed for transport. */
  50. export interface ProjectedRemoteEventRequest {
  51. readonly request: Readonly<Record<string, unknown>>
  52. readonly signal?: AbortSignal
  53. }
  54. /** Error fields retained when a Client listener rejects a Host waterfall. */
  55. export interface RemoteEventRejection {
  56. readonly name: string
  57. readonly message: string
  58. readonly code?: string
  59. readonly details?: unknown
  60. }
  61. /** Client response to one scoped Remote Event delivery. */
  62. export interface RemoteEventResult {
  63. readonly clientId: RemoteEventClientId
  64. readonly eventId: RemoteEventId
  65. readonly outcome:
  66. | { readonly kind: 'next' }
  67. | { readonly kind: 'result'; readonly value?: unknown }
  68. | { readonly kind: 'rejected'; readonly error: RemoteEventRejection }
  69. }
  70. /**
  71. * Parse one result sent through the Client's `$events/result` HTTP RPC.
  72. * @param value - untrusted result payload.
  73. * @returns validated event correlation and outcome fields.
  74. */
  75. export function parseRemoteEventResult(value: unknown): RemoteEventResult {
  76. if (!isRecord(value)
  77. || !exactKeys(value, ['clientId', 'eventId', 'outcome'])
  78. || !isRemoteEventClientId(value.clientId)
  79. || !isRemoteEventId(value.eventId)
  80. || !isRecord(value.outcome)) {
  81. throw new Error('api gateway: invalid Remote event result')
  82. }
  83. const outcome = value.outcome
  84. if (outcome.kind === 'next' && exactKeys(outcome, ['kind'])) {
  85. return {
  86. clientId: value.clientId,
  87. eventId: value.eventId,
  88. outcome: { kind: 'next' },
  89. }
  90. }
  91. if (outcome.kind === 'result'
  92. && (exactKeys(outcome, ['kind']) || exactKeys(outcome, ['kind', 'value']))
  93. && (!Object.hasOwn(outcome, 'value') || isRemoteJsonValue(outcome.value))) {
  94. return {
  95. clientId: value.clientId,
  96. eventId: value.eventId,
  97. outcome: Object.hasOwn(outcome, 'value')
  98. ? { kind: 'result', value: outcome.value }
  99. : { kind: 'result' },
  100. }
  101. }
  102. if (outcome.kind === 'rejected'
  103. && exactKeys(outcome, ['kind', 'error'])) {
  104. return {
  105. clientId: value.clientId,
  106. eventId: value.eventId,
  107. outcome: { kind: 'rejected', error: parseRemoteEventRejection(outcome.error) },
  108. }
  109. }
  110. throw new Error('api gateway: invalid Remote event result')
  111. }
  112. /**
  113. * Remove the direct Agent and cancellation fields from one waterfall request.
  114. * @param value - request object before the waterfall's `next` callback.
  115. * @param subject - Agent used by the Cordis scope carrier.
  116. * @returns JSON-safe request fields and the optional Host cancellation signal.
  117. */
  118. export function projectRemoteEventRequest(
  119. value: unknown,
  120. subject: object,
  121. ): ProjectedRemoteEventRequest {
  122. if (!isPlainRecord(value) || !Object.hasOwn(value, 'agent') || value.agent !== subject) {
  123. throw new TypeError('api gateway: Remote event request must carry its scoped Agent directly')
  124. }
  125. const signal = value.signal
  126. if (signal !== undefined && !(signal instanceof AbortSignal)) {
  127. throw new TypeError('api gateway: Remote event request signal must be an AbortSignal')
  128. }
  129. const request: Record<string, unknown> = Object.create(null) as Record<string, unknown>
  130. for (const key of Reflect.ownKeys(value)) {
  131. if (key === 'agent' || key === 'signal') continue
  132. const descriptor = typeof key === 'string' ? Object.getOwnPropertyDescriptor(value, key) : undefined
  133. if (typeof key !== 'string' || descriptor?.enumerable !== true) {
  134. throw new TypeError('api gateway: Remote event request has a non-JSON property')
  135. }
  136. request[key] = Reflect.get(value, key)
  137. }
  138. if (!isRemoteJsonValue(request)) {
  139. throw new TypeError('api gateway: Remote event request is not lossless JSON data')
  140. }
  141. return {
  142. request,
  143. ...(signal === undefined ? {} : { signal }),
  144. }
  145. }
  146. /**
  147. * Project an arbitrary rejection to stable, JSON-safe error fields.
  148. * @param reason - value thrown or rejected by a Client listener.
  149. * @returns wire-safe rejection fields.
  150. */
  151. export function projectRemoteEventRejection(reason: unknown): RemoteEventRejection {
  152. const record = typeof reason === 'object' && reason !== null ? reason : undefined
  153. const name = stringProperty(record, 'name') ?? 'Error'
  154. const message = stringProperty(record, 'message') ?? String(reason)
  155. const code = stringProperty(record, 'code')
  156. const details = record === undefined ? undefined : Reflect.get(record, 'details') as unknown
  157. return {
  158. name,
  159. message,
  160. ...(code === undefined ? {} : { code }),
  161. ...(details === undefined || !isRemoteJsonValue(details) ? {} : { details }),
  162. }
  163. }
  164. /**
  165. * Recreate a Client rejection for the Host continuation.
  166. * @param rejection - validated wire-safe error fields.
  167. * @returns an Error preserving the remote name, code, and JSON-safe details.
  168. */
  169. export function restoreRemoteEventRejection(rejection: RemoteEventRejection): Error {
  170. const error = new Error(rejection.message) as Error & { code?: string; details?: unknown }
  171. error.name = rejection.name
  172. if (rejection.code !== undefined) error.code = rejection.code
  173. if (rejection.details !== undefined) error.details = rejection.details
  174. return error
  175. }
  176. /**
  177. * Test whether a value crosses JSON transport without coercion or omission.
  178. * @param value - candidate boundary value.
  179. * @returns whether the value is losslessly JSON-compatible.
  180. */
  181. export function isRemoteJsonValue(value: unknown): boolean {
  182. return visitJsonValue(value, new Set<object>())
  183. }
  184. /**
  185. * Recognize a non-empty Remote Event correlation id at a wire boundary.
  186. * @param value - untrusted wire value.
  187. * @returns whether the value is a valid Remote Event id.
  188. */
  189. export function isRemoteEventId(value: unknown): value is RemoteEventId {
  190. return typeof value === 'string' && value.length > 0
  191. }
  192. /**
  193. * Recognize a non-empty Remote Event Client id at a wire boundary.
  194. * @param value - untrusted wire value.
  195. * @returns whether the value identifies one event-stream generation.
  196. */
  197. export function isRemoteEventClientId(value: unknown): value is RemoteEventClientId {
  198. return typeof value === 'string' && value.length > 0
  199. }
  200. /**
  201. * Recognize the direct Agent identity used by a scoped Remote Event.
  202. * @param value - untrusted wire value.
  203. * @returns whether the value is a non-empty Agent identity.
  204. */
  205. export function isRemoteEventAgentId(value: unknown): value is RemoteEventAgentId {
  206. return typeof value === 'string' && value.length > 0
  207. }
  208. /** One logical stream request sent from the browser. */
  209. export type RemoteStreamClientMessage =
  210. | {
  211. readonly type: 'open'
  212. readonly streamId: string
  213. readonly endpoint: string
  214. readonly payload: unknown
  215. }
  216. | { readonly type: 'cancel'; readonly streamId: string }
  217. /** Carrier-safe failure delivered by the Host. */
  218. export interface RemoteStreamFailure {
  219. readonly code: string
  220. readonly message: string
  221. readonly details: object
  222. }
  223. /** One logical stream frame sent from the Host. */
  224. export type RemoteStreamServerMessage =
  225. | { readonly type: 'item'; readonly streamId: string; readonly value?: unknown }
  226. | { readonly type: 'error'; readonly streamId: string; readonly error: RemoteStreamFailure }
  227. | { readonly type: 'end'; readonly streamId: string }
  228. /**
  229. * Parse and validate one browser-to-Host text message.
  230. * @param text - complete WebSocket text message.
  231. * @returns the validated logical-stream request.
  232. */
  233. export function parseRemoteStreamClientMessage(text: string): RemoteStreamClientMessage {
  234. return parseMessage(text, (value) => {
  235. if (value.type === 'cancel' && exactKeys(value, ['type', 'streamId']) && validId(value.streamId)) {
  236. return value as unknown as RemoteStreamClientMessage
  237. }
  238. if (value.type === 'open'
  239. && exactKeys(value, ['type', 'streamId', 'endpoint', 'payload'])
  240. && validId(value.streamId)
  241. && typeof value.endpoint === 'string'
  242. && value.endpoint.length > 0) {
  243. return value as unknown as RemoteStreamClientMessage
  244. }
  245. throw new Error('api gateway: invalid Remote stream client message')
  246. })
  247. }
  248. /**
  249. * Parse and validate one Host-to-browser text message.
  250. * @param text - complete WebSocket text message.
  251. * @returns the validated logical-stream frame.
  252. */
  253. export function parseRemoteStreamServerMessage(text: string): RemoteStreamServerMessage {
  254. return parseMessage(text, (value) => {
  255. if (value.type === 'item'
  256. && (exactKeys(value, ['type', 'streamId']) || exactKeys(value, ['type', 'streamId', 'value']))
  257. && validId(value.streamId)) {
  258. return value as unknown as RemoteStreamServerMessage
  259. }
  260. if (value.type === 'end' && exactKeys(value, ['type', 'streamId']) && validId(value.streamId)) {
  261. return value as unknown as RemoteStreamServerMessage
  262. }
  263. if (value.type === 'error'
  264. && exactKeys(value, ['type', 'streamId', 'error'])
  265. && validId(value.streamId)
  266. && isRecord(value.error)
  267. && exactKeys(value.error, ['code', 'message', 'details'])
  268. && typeof value.error.code === 'string'
  269. && typeof value.error.message === 'string'
  270. && isRecord(value.error.details)) {
  271. return value as unknown as RemoteStreamServerMessage
  272. }
  273. throw new Error('api gateway: invalid Remote stream server message')
  274. })
  275. }
  276. function parseMessage<T>(text: string, validate: (value: Record<string, unknown>) => T): T {
  277. let decoded: unknown
  278. try {
  279. decoded = JSON.parse(text) as unknown
  280. } catch (cause) {
  281. throw new Error('api gateway: Remote stream message is not JSON', { cause })
  282. }
  283. if (!isRecord(decoded)) throw new Error('api gateway: Remote stream message must be an object')
  284. return validate(decoded)
  285. }
  286. function isRecord(value: unknown): value is Record<string, unknown> {
  287. return typeof value === 'object'
  288. && value !== null
  289. && !Array.isArray(value)
  290. }
  291. function isPlainRecord(value: unknown): value is Record<string, unknown> {
  292. if (!isRecord(value)) return false
  293. const prototype: unknown = Object.getPrototypeOf(value)
  294. return prototype === Object.prototype || prototype === null
  295. }
  296. function exactKeys(value: Record<string, unknown>, expected: readonly string[]): boolean {
  297. const keys = Reflect.ownKeys(value)
  298. return keys.length === expected.length && expected.every(key => Object.hasOwn(value, key))
  299. }
  300. function validId(value: unknown): value is string {
  301. return typeof value === 'string' && value.length > 0
  302. }
  303. function parseRemoteEventRejection(value: unknown): RemoteEventRejection {
  304. if (!isRecord(value)
  305. || !hasOnlyKeys(value, ['name', 'message'], ['code', 'details'])
  306. || typeof value.name !== 'string'
  307. || value.name.length === 0
  308. || typeof value.message !== 'string'
  309. || (Object.hasOwn(value, 'code') && typeof value.code !== 'string')
  310. || (Object.hasOwn(value, 'details') && !isRemoteJsonValue(value.details))) {
  311. throw new Error('api gateway: invalid Remote event rejection')
  312. }
  313. return {
  314. name: value.name,
  315. message: value.message,
  316. ...(typeof value.code === 'string' ? { code: value.code } : {}),
  317. ...(Object.hasOwn(value, 'details') ? { details: value.details } : {}),
  318. }
  319. }
  320. function hasOnlyKeys(
  321. value: Record<string, unknown>,
  322. required: readonly string[],
  323. optional: readonly string[],
  324. ): boolean {
  325. const keys = Reflect.ownKeys(value)
  326. return required.every(key => Object.hasOwn(value, key))
  327. && keys.every(key => typeof key === 'string' && (required.includes(key) || optional.includes(key)))
  328. }
  329. function stringProperty(value: object | undefined, key: string): string | undefined {
  330. if (value === undefined) return undefined
  331. const candidate: unknown = Reflect.get(value, key)
  332. return typeof candidate === 'string' ? candidate : undefined
  333. }
  334. function visitJsonValue(value: unknown, ancestors: Set<object>): boolean {
  335. if (value === null || typeof value === 'string' || typeof value === 'boolean') return true
  336. if (typeof value === 'number') return Number.isFinite(value) && !Object.is(value, -0)
  337. if (typeof value !== 'object') return false
  338. if (ancestors.has(value)) return false
  339. ancestors.add(value)
  340. try {
  341. if (Array.isArray(value)) {
  342. if (Object.getPrototypeOf(value) !== Array.prototype
  343. || Reflect.ownKeys(value).length !== value.length + 1) return false
  344. for (let index = 0; index < value.length; index++) {
  345. if (!Object.hasOwn(value, index) || !visitJsonValue(value[index], ancestors)) return false
  346. }
  347. return true
  348. }
  349. const prototype: unknown = Object.getPrototypeOf(value)
  350. if (prototype !== Object.prototype && prototype !== null) return false
  351. for (const key of Reflect.ownKeys(value)) {
  352. if (typeof key !== 'string') return false
  353. const descriptor = Object.getOwnPropertyDescriptor(value, key)
  354. if (descriptor?.enumerable !== true || !visitJsonValue(Reflect.get(value, key), ancestors)) return false
  355. }
  356. return true
  357. } finally {
  358. ancestors.delete(value)
  359. }
  360. }