1
0

rpc-host.ts 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248
  1. /** Host registry and HTTP adapter for generic Connection RPC channels. */
  2. import { Context, Service } from '@deepseek-ai/cordis'
  3. import type { WebRoute } from '@deepseek-ai/dsh-host-webserver'
  4. import {
  5. clientRequestSchema,
  6. RpcId,
  7. type ClientRequest,
  8. type RpcError,
  9. type RpcErrorDetailsMap,
  10. type RpcId as RpcIdType,
  11. } from '@deepseek-ai/dsh-host-apiproxy/api'
  12. import { bridge, type FetchHandler } from './http-bridge.ts'
  13. import { isTrustedApiRequest } from './api-request-trust.ts'
  14. import { API_PATH } from './api-path.ts'
  15. import type { BrowserAuth } from './browser-auth.ts'
  16. import type {
  17. ConnectionIndexRequest,
  18. ConnectionIndexResponse,
  19. ConnectionRpcEndpointMatcher,
  20. ConnectionRpcHandler,
  21. ConnectionRpcResult,
  22. ConnectionRequestRejection,
  23. ConnectionTrustRequest,
  24. HostConnectionHandle,
  25. HostConnectionRpc,
  26. } from './rpc.ts'
  27. const INVALID_REQUEST_RPC_ID = RpcId('invalid-request')
  28. const CHANNEL_PATTERN = /^\/[A-Za-z0-9._~-]+$/
  29. const ENDPOINT_SEGMENT_PATTERN = /^[A-Za-z0-9_$.-]+$/
  30. interface ConnectionRpcInterceptor {
  31. readonly matches: ConnectionRpcEndpointMatcher
  32. readonly fetchHandler: FetchHandler
  33. }
  34. interface ConnectionServerResponse {
  35. readonly type: 'server-response'
  36. readonly rpcId: RpcIdType
  37. readonly result: ConnectionRpcResult<unknown>
  38. }
  39. declare module '@deepseek-ai/cordis' {
  40. interface Context {
  41. /** Host Connection transport and RPC registrations. */
  42. connection: HostConnectionHandle
  43. }
  44. }
  45. /** Host Connection service whose channel registrations belong to the caller fiber. */
  46. export class HostConnectionService extends Service implements HostConnectionHandle {
  47. private readonly interceptors = new Map<string, ConnectionRpcInterceptor>()
  48. /**
  49. * Provide the Host half over the active HTTP server.
  50. * @param ctx - owning Connection plugin context.
  51. * @param trustedHosts - deployment authorities accepted by the Host/Origin fence.
  52. * @param browserAuth - process token and persistent browser-session owner.
  53. */
  54. constructor(
  55. ctx: Context,
  56. private readonly trustedHosts: readonly string[],
  57. private readonly browserAuth: BrowserAuth,
  58. ) {
  59. super(ctx, 'connection')
  60. }
  61. /** Generic channel registry scoped to the Context reading this service. */
  62. get rpc(): HostConnectionRpc {
  63. const owner = this.ctx
  64. return {
  65. handle: (channel, handler) => this.register(owner, channel, handler),
  66. intercept: (channel, matches, handler) =>
  67. this.registerInterceptor(owner, channel, matches, handler),
  68. }
  69. }
  70. /** Apply the configured Host/Origin fence, then browser authentication. */
  71. requestRejection(request: ConnectionTrustRequest): ConnectionRequestRejection {
  72. if (!isTrustedApiRequest(request, this.trustedHosts)) return 403
  73. return this.browserAuth.isAuthenticated(request) ? undefined : 401
  74. }
  75. /** Authenticate an index request through the process-token exchange or cookie. */
  76. authorizeIndex(request: ConnectionIndexRequest, response: ConnectionIndexResponse): boolean {
  77. return this.browserAuth.authorizeIndex(request, response)
  78. }
  79. /** Add this process's launch token to the clean application URL. */
  80. authenticatedUrl(baseUrl: string): string {
  81. return this.browserAuth.authenticatedUrl(baseUrl)
  82. }
  83. /**
  84. * Compose one shared-channel Fetch handler from its interceptor and fallback.
  85. * @param channel - shared channel mounted by Connection.
  86. * @param fallback - handler for endpoints not claimed by the interceptor.
  87. * @returns Fetch handler that selects exactly one target for each request.
  88. */
  89. createSharedFetchHandler(
  90. channel: '/api',
  91. fallback: FetchHandler,
  92. ): FetchHandler {
  93. return {
  94. fetch: (request) => {
  95. const endpoint = endpointFromPath(channel, new URL(request.url).pathname)
  96. const interceptor = this.interceptors.get(channel)
  97. if (endpoint === undefined || interceptor === undefined || !interceptor.matches(endpoint)) {
  98. return fallback.fetch(request)
  99. }
  100. return interceptor.fetchHandler.fetch(request)
  101. },
  102. }
  103. }
  104. private register(
  105. owner: Context,
  106. channel: string,
  107. handler: ConnectionRpcHandler,
  108. ): () => Promise<void> {
  109. assertChannel(channel)
  110. const fetchHandler = rpcFetchHandler(channel, handler)
  111. const route: WebRoute = {
  112. kind: 'prefix',
  113. path: channel,
  114. handler: async (req, res) => {
  115. const rejection = this.requestRejection(req)
  116. if (rejection !== undefined) {
  117. res.writeHead(rejection)
  118. res.end(rejection === 401 ? 'unauthorized' : 'forbidden')
  119. return
  120. }
  121. await bridge(req, res, fetchHandler)
  122. },
  123. }
  124. return owner.effect(
  125. () => owner.webServer.register(route),
  126. `client-connection: ${channel} rpc channel`,
  127. )
  128. }
  129. private registerInterceptor(
  130. owner: Context,
  131. channel: string,
  132. matches: ConnectionRpcEndpointMatcher,
  133. handler: ConnectionRpcHandler,
  134. ): () => Promise<void> {
  135. if (channel !== API_PATH) {
  136. throw new Error(`connection: invalid shared RPC channel ${JSON.stringify(channel)}`)
  137. }
  138. const interceptor: ConnectionRpcInterceptor = {
  139. matches,
  140. fetchHandler: rpcFetchHandler(channel, handler),
  141. }
  142. return owner.effect(() => {
  143. if (this.interceptors.has(channel)) {
  144. throw new Error(`connection: shared RPC channel ${JSON.stringify(channel)} already has an interceptor`)
  145. }
  146. this.interceptors.set(channel, interceptor)
  147. return () => {
  148. this.interceptors.delete(channel)
  149. }
  150. }, `client-connection: ${channel} rpc interceptor`)
  151. }
  152. }
  153. function rpcFetchHandler(
  154. channel: string,
  155. handler: ConnectionRpcHandler,
  156. ): FetchHandler {
  157. return {
  158. async fetch(request: Request): Promise<Response> {
  159. const endpoint = endpointFromPath(channel, new URL(request.url).pathname)
  160. if (request.method !== 'POST' || endpoint === undefined) {
  161. return new Response('not found', { status: 404 })
  162. }
  163. const mediaType = request.headers.get('content-type')?.split(';', 1)[0]?.trim().toLowerCase()
  164. if (mediaType !== 'application/json') {
  165. return new Response('content type must be application/json', { status: 415 })
  166. }
  167. let body: unknown
  168. try {
  169. body = await request.json()
  170. } catch {
  171. return new Response('body is not JSON', { status: 400 })
  172. }
  173. const envelope = clientRequestSchema.safeParse(body)
  174. if (!envelope.success) {
  175. return invalidEnvelopeResponse(body, envelope.error.issues)
  176. }
  177. const message: ClientRequest = envelope.data
  178. if (message.method !== endpoint) {
  179. return errorResponse(message.rpcId, {
  180. code: 'bad-request',
  181. message: `method ${JSON.stringify(message.method)} does not match endpoint ${JSON.stringify(endpoint)}`,
  182. details: { issues: [] },
  183. })
  184. }
  185. try {
  186. const result = await handler(endpoint, message.payload, request.signal)
  187. return fullResponse(message.rpcId, result)
  188. } catch (error) {
  189. return new Response(`handler failure: ${String(error)}`, { status: 500 })
  190. }
  191. },
  192. }
  193. }
  194. function invalidEnvelopeResponse(body: unknown, issues: RpcErrorDetailsMap['bad-request']['issues']): Response {
  195. const rawId = (body as { rpcId?: unknown } | null)?.rpcId
  196. const rpcId = typeof rawId === 'string' ? RpcId(rawId) : INVALID_REQUEST_RPC_ID
  197. return errorResponse(rpcId, {
  198. code: 'bad-request',
  199. message: 'invalid client-request message',
  200. details: { issues },
  201. })
  202. }
  203. function endpointFromPath(channel: string, pathname: string): string | undefined {
  204. if (!pathname.startsWith(`${channel}/`)) return undefined
  205. const endpoint = pathname.slice(channel.length + 1)
  206. const segments = endpoint.split('/')
  207. if (segments.some(segment =>
  208. segment === '' || segment === '.' || segment === '..' || !ENDPOINT_SEGMENT_PATTERN.test(segment))) {
  209. return undefined
  210. }
  211. return endpoint
  212. }
  213. function errorResponse(rpcId: RpcIdType, error: RpcError): Response {
  214. return fullResponse(rpcId, { ok: false, error })
  215. }
  216. function fullResponse(rpcId: RpcIdType, result: ConnectionRpcResult<unknown>): Response {
  217. const body: ConnectionServerResponse = { type: 'server-response', rpcId, result }
  218. return Response.json(body)
  219. }
  220. function assertChannel(channel: string): void {
  221. if (!CHANNEL_PATTERN.test(channel) || channel === '/api') {
  222. throw new Error(`connection: invalid or reserved RPC channel ${JSON.stringify(channel)}`)
  223. }
  224. }