| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248 |
- /** Host registry and HTTP adapter for generic Connection RPC channels. */
- import { Context, Service } from '@deepseek-ai/cordis'
- import type { WebRoute } from '@deepseek-ai/dsh-host-webserver'
- import {
- clientRequestSchema,
- RpcId,
- type ClientRequest,
- type RpcError,
- type RpcErrorDetailsMap,
- type RpcId as RpcIdType,
- } from '@deepseek-ai/dsh-host-apiproxy/api'
- import { bridge, type FetchHandler } from './http-bridge.ts'
- import { isTrustedApiRequest } from './api-request-trust.ts'
- import { API_PATH } from './api-path.ts'
- import type { BrowserAuth } from './browser-auth.ts'
- import type {
- ConnectionIndexRequest,
- ConnectionIndexResponse,
- ConnectionRpcEndpointMatcher,
- ConnectionRpcHandler,
- ConnectionRpcResult,
- ConnectionRequestRejection,
- ConnectionTrustRequest,
- HostConnectionHandle,
- HostConnectionRpc,
- } from './rpc.ts'
- const INVALID_REQUEST_RPC_ID = RpcId('invalid-request')
- const CHANNEL_PATTERN = /^\/[A-Za-z0-9._~-]+$/
- const ENDPOINT_SEGMENT_PATTERN = /^[A-Za-z0-9_$.-]+$/
- interface ConnectionRpcInterceptor {
- readonly matches: ConnectionRpcEndpointMatcher
- readonly fetchHandler: FetchHandler
- }
- interface ConnectionServerResponse {
- readonly type: 'server-response'
- readonly rpcId: RpcIdType
- readonly result: ConnectionRpcResult<unknown>
- }
- declare module '@deepseek-ai/cordis' {
- interface Context {
- /** Host Connection transport and RPC registrations. */
- connection: HostConnectionHandle
- }
- }
- /** Host Connection service whose channel registrations belong to the caller fiber. */
- export class HostConnectionService extends Service implements HostConnectionHandle {
- private readonly interceptors = new Map<string, ConnectionRpcInterceptor>()
- /**
- * Provide the Host half over the active HTTP server.
- * @param ctx - owning Connection plugin context.
- * @param trustedHosts - deployment authorities accepted by the Host/Origin fence.
- * @param browserAuth - process token and persistent browser-session owner.
- */
- constructor(
- ctx: Context,
- private readonly trustedHosts: readonly string[],
- private readonly browserAuth: BrowserAuth,
- ) {
- super(ctx, 'connection')
- }
- /** Generic channel registry scoped to the Context reading this service. */
- get rpc(): HostConnectionRpc {
- const owner = this.ctx
- return {
- handle: (channel, handler) => this.register(owner, channel, handler),
- intercept: (channel, matches, handler) =>
- this.registerInterceptor(owner, channel, matches, handler),
- }
- }
- /** Apply the configured Host/Origin fence, then browser authentication. */
- requestRejection(request: ConnectionTrustRequest): ConnectionRequestRejection {
- if (!isTrustedApiRequest(request, this.trustedHosts)) return 403
- return this.browserAuth.isAuthenticated(request) ? undefined : 401
- }
- /** Authenticate an index request through the process-token exchange or cookie. */
- authorizeIndex(request: ConnectionIndexRequest, response: ConnectionIndexResponse): boolean {
- return this.browserAuth.authorizeIndex(request, response)
- }
- /** Add this process's launch token to the clean application URL. */
- authenticatedUrl(baseUrl: string): string {
- return this.browserAuth.authenticatedUrl(baseUrl)
- }
- /**
- * Compose one shared-channel Fetch handler from its interceptor and fallback.
- * @param channel - shared channel mounted by Connection.
- * @param fallback - handler for endpoints not claimed by the interceptor.
- * @returns Fetch handler that selects exactly one target for each request.
- */
- createSharedFetchHandler(
- channel: '/api',
- fallback: FetchHandler,
- ): FetchHandler {
- return {
- fetch: (request) => {
- const endpoint = endpointFromPath(channel, new URL(request.url).pathname)
- const interceptor = this.interceptors.get(channel)
- if (endpoint === undefined || interceptor === undefined || !interceptor.matches(endpoint)) {
- return fallback.fetch(request)
- }
- return interceptor.fetchHandler.fetch(request)
- },
- }
- }
- private register(
- owner: Context,
- channel: string,
- handler: ConnectionRpcHandler,
- ): () => Promise<void> {
- assertChannel(channel)
- const fetchHandler = rpcFetchHandler(channel, handler)
- const route: WebRoute = {
- kind: 'prefix',
- path: channel,
- handler: async (req, res) => {
- const rejection = this.requestRejection(req)
- if (rejection !== undefined) {
- res.writeHead(rejection)
- res.end(rejection === 401 ? 'unauthorized' : 'forbidden')
- return
- }
- await bridge(req, res, fetchHandler)
- },
- }
- return owner.effect(
- () => owner.webServer.register(route),
- `client-connection: ${channel} rpc channel`,
- )
- }
- private registerInterceptor(
- owner: Context,
- channel: string,
- matches: ConnectionRpcEndpointMatcher,
- handler: ConnectionRpcHandler,
- ): () => Promise<void> {
- if (channel !== API_PATH) {
- throw new Error(`connection: invalid shared RPC channel ${JSON.stringify(channel)}`)
- }
- const interceptor: ConnectionRpcInterceptor = {
- matches,
- fetchHandler: rpcFetchHandler(channel, handler),
- }
- return owner.effect(() => {
- if (this.interceptors.has(channel)) {
- throw new Error(`connection: shared RPC channel ${JSON.stringify(channel)} already has an interceptor`)
- }
- this.interceptors.set(channel, interceptor)
- return () => {
- this.interceptors.delete(channel)
- }
- }, `client-connection: ${channel} rpc interceptor`)
- }
- }
- function rpcFetchHandler(
- channel: string,
- handler: ConnectionRpcHandler,
- ): FetchHandler {
- return {
- async fetch(request: Request): Promise<Response> {
- const endpoint = endpointFromPath(channel, new URL(request.url).pathname)
- if (request.method !== 'POST' || endpoint === undefined) {
- return new Response('not found', { status: 404 })
- }
- const mediaType = request.headers.get('content-type')?.split(';', 1)[0]?.trim().toLowerCase()
- if (mediaType !== 'application/json') {
- return new Response('content type must be application/json', { status: 415 })
- }
- let body: unknown
- try {
- body = await request.json()
- } catch {
- return new Response('body is not JSON', { status: 400 })
- }
- const envelope = clientRequestSchema.safeParse(body)
- if (!envelope.success) {
- return invalidEnvelopeResponse(body, envelope.error.issues)
- }
- const message: ClientRequest = envelope.data
- if (message.method !== endpoint) {
- return errorResponse(message.rpcId, {
- code: 'bad-request',
- message: `method ${JSON.stringify(message.method)} does not match endpoint ${JSON.stringify(endpoint)}`,
- details: { issues: [] },
- })
- }
- try {
- const result = await handler(endpoint, message.payload, request.signal)
- return fullResponse(message.rpcId, result)
- } catch (error) {
- return new Response(`handler failure: ${String(error)}`, { status: 500 })
- }
- },
- }
- }
- function invalidEnvelopeResponse(body: unknown, issues: RpcErrorDetailsMap['bad-request']['issues']): Response {
- const rawId = (body as { rpcId?: unknown } | null)?.rpcId
- const rpcId = typeof rawId === 'string' ? RpcId(rawId) : INVALID_REQUEST_RPC_ID
- return errorResponse(rpcId, {
- code: 'bad-request',
- message: 'invalid client-request message',
- details: { issues },
- })
- }
- function endpointFromPath(channel: string, pathname: string): string | undefined {
- if (!pathname.startsWith(`${channel}/`)) return undefined
- const endpoint = pathname.slice(channel.length + 1)
- const segments = endpoint.split('/')
- if (segments.some(segment =>
- segment === '' || segment === '.' || segment === '..' || !ENDPOINT_SEGMENT_PATTERN.test(segment))) {
- return undefined
- }
- return endpoint
- }
- function errorResponse(rpcId: RpcIdType, error: RpcError): Response {
- return fullResponse(rpcId, { ok: false, error })
- }
- function fullResponse(rpcId: RpcIdType, result: ConnectionRpcResult<unknown>): Response {
- const body: ConnectionServerResponse = { type: 'server-response', rpcId, result }
- return Response.json(body)
- }
- function assertChannel(channel: string): void {
- if (!CHANNEL_PATTERN.test(channel) || channel === '/api') {
- throw new Error(`connection: invalid or reserved RPC channel ${JSON.stringify(channel)}`)
- }
- }
|