返回 CodeWhale
rpc.ts
1 /**
2 * JSON-RPC 2.0 peer over the `CWX1` channel.
3 *
4 * Every inbound request gets an `AbortController`; `$/cancel {id}` from the
5 * core aborts it. The core resolves a cancelled call on its own side after a
6 * 500 ms grace, so a late answer here is harmless (the core drops it).
7 */
8 import { ErrorCode, MAX_INFLIGHT, validateMessage, type HostTier, type Message, type RpcErrorWire } from './protocol.ts'
9
10 export class RpcError extends Error {
11 constructor(
12 readonly code: number,
13 message: string,
14 readonly data?: unknown,
15 ) {
16 super(message)
17 this.name = 'RpcError'
18 }
19
20 toWire(): RpcErrorWire {
21 return this.data === undefined
22 ? { code: this.code, message: this.message }
23 : { code: this.code, message: this.message, data: this.data as never }
24 }
25 }
26
27 export interface RequestContext {
28 id: number
29 signal: AbortSignal
30 }
31
32 type RequestHandler = (params: any, cx: RequestContext) => Promise<unknown> | unknown
33 type NotificationHandler = (params: any) => void
34
35 export class RpcPeer {
36 private nextId = 1
37 private readonly pending = new Map<number, { resolve: (v: any) => void; reject: (e: Error) => void }>()
38 private readonly inbound = new Map<number, AbortController>()
39 private readonly requestHandlers = new Map<string, RequestHandler>()
40 private readonly notificationHandlers = new Map<string, NotificationHandler>()
41 private closed = false
42
43 constructor(
44 private readonly send: (message: Message) => void,
45 /** The trust tier this host serves: methods reserved for the other tier are neither sent nor accepted. */
46 private readonly tier: HostTier,
47 ) {}
48
49 onRequest(method: string, handler: RequestHandler) {
50 this.requestHandlers.set(method, handler)
51 }
52
53 onNotification(method: string, handler: NotificationHandler) {
54 this.notificationHandlers.set(method, handler)
55 }
56
57 /**
58 * Send a host→core request. Outbound messages are validated strictly first.
59 * When `signal` aborts before the answer, the core is sent `$/cancel` for it
60 * and the promise rejects as cancelled at once; an answer that still arrives
61 * is dropped (the core drops its own after a cancel as well).
62 */
63 request<T = any>(method: string, params: unknown, signal?: AbortSignal): Promise<T> {
64 if (this.closed) return Promise.reject(new RpcError(ErrorCode.NotAvailable, 'channel closed'))
65 if (signal?.aborted) return Promise.reject(new RpcError(ErrorCode.Cancelled, 'cancelled'))
66 // The core's per-direction limit.
67 if (this.pending.size >= MAX_INFLIGHT) {
68 return Promise.reject(new RpcError(ErrorCode.Internal, `more than ${MAX_INFLIGHT} requests in flight`))
69 }
70 const id = this.nextId++
71 const message = { jsonrpc: '2.0' as const, id, method, params }
72 validateMessage(message, 'host_to_core', this.tier)
73 return new Promise<T>((resolve, reject) => {
74 const onAbort = () => {
75 if (!this.pending.delete(id)) return
76 this.notify('$/cancel', { id })
77 reject(new RpcError(ErrorCode.Cancelled, 'cancelled'))
78 }
79 this.pending.set(id, {
80 resolve: (value) => {
81 signal?.removeEventListener('abort', onAbort)
82 resolve(value)
83 },
84 reject: (error) => {
85 signal?.removeEventListener('abort', onAbort)
86 reject(error)
87 },
88 })
89 signal?.addEventListener('abort', onAbort, { once: true })
90 this.send(message)
91 })
92 }
93
94 notify(method: string, params: unknown) {
95 if (this.closed) return
96 const message = { jsonrpc: '2.0' as const, method, params }
97 validateMessage(message, 'host_to_core', this.tier)
98 this.send(message)
99 }
100
101 /** Dispatch one decoded core→host message. Throws `ProtocolError` for malformed input. */
102 handle(raw: unknown) {
103 const message = validateMessage(raw, 'core_to_host', this.tier) as any
104 if ('method' in message) {
105 if (message.method === '$/cancel') {
106 this.inbound.get(message.params.id)?.abort()
107 return
108 }
109 if ('id' in message) {
110 void this.dispatchRequest(message.id, message.method, message.params ?? {})
111 return
112 }
113 this.notificationHandlers.get(message.method)?.(message.params ?? {})
114 return
115 }
116 const waiter = this.pending.get(message.id)
117 if (!waiter) return
118 this.pending.delete(message.id)
119 if ('error' in message) {
120 waiter.reject(new RpcError(message.error.code, message.error.message, message.error.data))
121 } else {
122 waiter.resolve(message.result)
123 }
124 }
125
126 /** Fail every outbound request; used when the channel closes. */
127 close(reason: string) {
128 this.closed = true
129 for (const [, waiter] of this.pending) waiter.reject(new RpcError(ErrorCode.NotAvailable, reason))
130 this.pending.clear()
131 for (const [, controller] of this.inbound) controller.abort()
132 }
133
134 private async dispatchRequest(id: number, method: string, params: unknown) {
135 const handler = this.requestHandlers.get(method)
136 if (!handler) {
137 this.reply(id, undefined, new RpcError(ErrorCode.MethodNotFound, `no handler for \`${method}\``))
138 return
139 }
140 const controller = new AbortController()
141 this.inbound.set(id, controller)
142 try {
143 const result = await handler(params, { id, signal: controller.signal })
144 this.reply(id, result ?? {})
145 } catch (error) {
146 this.reply(id, undefined, toRpcError(error, controller.signal))
147 } finally {
148 this.inbound.delete(id)
149 }
150 }
151
152 private reply(id: number, result: unknown, error?: RpcError) {
153 if (this.closed) return
154 this.send(error ? { jsonrpc: '2.0', id, error: error.toWire() } : { jsonrpc: '2.0', id, result })
155 }
156 }
157
158 export function toRpcError(error: unknown, signal?: AbortSignal): RpcError {
159 if (error instanceof RpcError) return error
160 if (signal?.aborted) return new RpcError(ErrorCode.Cancelled, 'cancelled')
161 const message = error instanceof Error ? `${error.name}: ${error.message}` : String(error)
162 return new RpcError(ErrorCode.ExecutionFailed, message.slice(0, 4096))
163 }
164
164 lines TYPESCRIPT