feat(connection): dispatch TypeRT remotes through shared API
This commit is contained in:
@@ -2,5 +2,5 @@
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write packages/client/connection/README.md
|
||||
README.md: 1393e79aacecbbf7b186f19e4c42269595854b0e
|
||||
README.zh.md: 70380ceba1b16b2970e947fb6cd9b2af9085ae51
|
||||
README.md: 161e34c4b6018625fb690e178eb9a9f8ac0ef21b
|
||||
README.zh.md: d17012cc89c02a1b11f16d126b7c0cafe67fb2a0
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
English | [中文](README.zh.md)
|
||||
|
||||
Wire consumer layer: the client plugin's apply mounts `ctx.connection` (shared api client + current-page loopback state + single-consumer stream-loop starter); the export face carries the wire contract types, the `AbstractApiClient` seam, and the loop's sink/config types. The browser carrier uses HTTP POST for unary and respond operations and opens one downlink-only WebSocket each for `events.mux` and `events.host`; the in-process carrier satisfies the same two-stream abstraction. Loopback hostname classification stays package-internal: the `/api` Host fence and WebSocket upgrades use it directly, while other client plugins consume the derived `ctx.connection.isLoopback` state. The node half's `/api` route pins the privileged method set (`host.pickDirectory`, `host.openPath`, and the whole configuration plane — `settings.describe`/`openDocument`/`update`/`replace`/`mutate` and `credentials.describe`/`set`/`unset`; reads and native actions included, since describing returns the exposed configuration, opening acts on the Host desktop, and probing an arbitrary reference reports where a credential comes from) to loopback by passing the trust fence with an empty trust list — a declared `trustedHosts` authority reaches every other method, while these stay loopback-local until a real authentication layer exists. The platform carriers and ConnectionController loop are package-internal; apply selects and drives them. The downlink boundary is documented in the [WebSocket downlink carrier Agent Note](../../../.agents/notes/implemented/architecture/2026-08-04-websocket-downlink-carrier.md); the protocol contract is api-contracts v3 §3.
|
||||
Wire consumer layer: the client plugin's apply mounts `ctx.connection` (shared api client + current-page loopback state + single-consumer stream-loop starter); the export face carries the wire contract types, the `AbstractApiClient` seam, and the loop's sink/config types. The browser carrier uses HTTP POST for unary and respond operations and opens one downlink-only WebSocket each for `events.mux` and `events.host`; the in-process carrier satisfies the same two-stream abstraction. The Host half owns the single `/api` route and its Fetch bridge; a registered TypeRT interceptor claims its Remote endpoints before the API Proxy fallback. Loopback hostname classification stays package-internal: the `/api` Host fence and WebSocket upgrades use it directly, while other client plugins consume the derived `ctx.connection.isLoopback` state. The node half's `/api` route pins the privileged method set (`host.pickDirectory`, `host.openPath`, and the whole configuration plane — `settings.describe`/`openDocument`/`update`/`replace`/`mutate` and `credentials.describe`/`set`/`unset`; reads and native actions included, since describing returns the exposed configuration, opening acts on the Host desktop, and probing an arbitrary reference reports where a credential comes from) to loopback by passing the trust fence with an empty trust list — a declared `trustedHosts` authority reaches every other method, while these stay loopback-local until a real authentication layer exists. The platform carriers and ConnectionController loop are package-internal; apply selects and drives them. The downlink boundary is documented in the [WebSocket downlink carrier Agent Note](../../../.agents/notes/implemented/architecture/2026-08-04-websocket-downlink-carrier.md); the protocol contract is api-contracts v3 §3.
|
||||
|
||||
## /api browser-trust fence
|
||||
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
[English](README.md) | 中文
|
||||
|
||||
协议消费层:客户端插件的 apply 会挂载 `ctx.connection`(共享 API 客户端 + 当前页面的 loopback 状态 + 单消费方流循环启动器);导出表层携带协议契约类型、`AbstractApiClient` seam,以及循环的 sink/配置类型。浏览器载体以 HTTP POST 发送 unary/respond,并为 `events.mux` 与 `events.host` 各开一条只下行的 WebSocket;进程内载体满足同一双流抽象。Loopback hostname 判定逻辑留在包内部:`/api` Host fence 与 WebSocket upgrade 会直接使用它,其他客户端插件则消费派生的 `ctx.connection.isLoopback` 状态。node 半侧的 `/api` 路由让特权方法集(`host.pickDirectory`、`host.openPath`,以及整个配置面——`settings.describe`/`openDocument`/`update`/`replace`/`mutate` 与 `credentials.describe`/`set`/`unset`;读取与原生操作也在内,因为 describe 会返回已暴露的配置、打开操作会作用于 Host 桌面,而探测任意引用会报出某条凭据来自何处)以空信任表过信任 fence,从而钉在回环——已声明的 `trustedHosts` 授权可达其余全部方法,而这些方法在真正的认证层出现之前仍只限回环本机。平台载体与 ConnectionController 循环属于包内部;apply 负责选择并驱动它们。下行边界见 [WebSocket 下行载体 Agent Note](../../../.agents/notes/implemented/architecture/2026-08-04-websocket-downlink-carrier.md);协议契约见 api-contracts v3 §3。
|
||||
协议消费层:客户端插件的 apply 会挂载 `ctx.connection`(共享 API 客户端 + 当前页面的 loopback 状态 + 单消费方流循环启动器);导出表层携带协议契约类型、`AbstractApiClient` seam,以及循环的 sink/配置类型。浏览器载体以 HTTP POST 发送 unary/respond,并为 `events.mux` 与 `events.host` 各开一条只下行的 WebSocket;进程内载体满足同一双流抽象。Host half 持有唯一 `/api` route 及其 Fetch bridge;已注册的 TypeRT interceptor 会先认领自己的 Remote endpoint,未认领请求再回退 API Proxy。Loopback hostname 判定逻辑留在包内部:`/api` Host fence 与 WebSocket upgrade 会直接使用它,其他客户端插件则消费派生的 `ctx.connection.isLoopback` 状态。node 半侧的 `/api` 路由让特权方法集(`host.pickDirectory`、`host.openPath`,以及整个配置面——`settings.describe`/`openDocument`/`update`/`replace`/`mutate` 与 `credentials.describe`/`set`/`unset`;读取与原生操作也在内,因为 describe 会返回已暴露的配置、打开操作会作用于 Host 桌面,而探测任意引用会报出某条凭据来自何处)以空信任表过信任 fence,从而钉在回环——已声明的 `trustedHosts` 授权可达其余全部方法,而这些方法在真正的认证层出现之前仍只限回环本机。平台载体与 ConnectionController 循环属于包内部;apply 负责选择并驱动它们。下行边界见 [WebSocket 下行载体 Agent Note](../../../.agents/notes/implemented/architecture/2026-08-04-websocket-downlink-carrier.md);协议契约见 api-contracts v3 §3。
|
||||
|
||||
## /api 浏览器信任栅栏
|
||||
|
||||
|
||||
@@ -16,12 +16,13 @@
|
||||
import type { IncomingHttpHeaders } from 'node:http'
|
||||
import { isLoopbackHostname } from './loopback-hostname.ts'
|
||||
|
||||
/** The request facts the fence reads (structural subset of IncomingMessage). */
|
||||
/** The request facts the fence reads from either HTTP representation. */
|
||||
interface ApiTrustRequest {
|
||||
headers: IncomingHttpHeaders
|
||||
headers: IncomingHttpHeaders | Headers
|
||||
}
|
||||
|
||||
function header(headers: IncomingHttpHeaders, name: string): string | undefined {
|
||||
function header(headers: IncomingHttpHeaders | Headers, name: string): string | undefined {
|
||||
if (headers instanceof Headers) return headers.get(name) ?? undefined
|
||||
const value = headers[name]
|
||||
return typeof value === 'string' ? value : undefined
|
||||
}
|
||||
@@ -88,7 +89,7 @@ function isTrustedAuthority(hostUrl: URL, trustedHosts: readonly string[]): bool
|
||||
|
||||
/**
|
||||
* Decide whether one /api request may reach the RPC bridge.
|
||||
* @param request - node HTTP request facts (headers).
|
||||
* @param request - Node HTTP or Fetch request facts (headers).
|
||||
* @param trustedHosts - non-loopback authorities this deployment serves: exact `host:port`, or port-less `host` matching any port.
|
||||
* @returns true when the Host is ours (loopback or trusted) and any attached browser markers are same-origin.
|
||||
*/
|
||||
|
||||
@@ -5,7 +5,13 @@
|
||||
|
||||
import type { IncomingMessage, ServerResponse } from 'node:http'
|
||||
|
||||
interface FetchHandler {
|
||||
/** Transport-independent request handler consumed by the Host HTTP bridge. */
|
||||
export interface FetchHandler {
|
||||
/**
|
||||
* Handle one standard Fetch request.
|
||||
* @param request - request produced by the active transport bridge.
|
||||
* @returns complete or streaming Fetch response.
|
||||
*/
|
||||
fetch(request: Request): Promise<Response>
|
||||
}
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ import { rejectWebSocketUpgrade, WebSocketDownlinks } from './websocket-downlink
|
||||
|
||||
export type {
|
||||
ConnectionRpcAuthority,
|
||||
ConnectionRpcEndpointMatcher,
|
||||
ConnectionRpcHandler,
|
||||
ConnectionRpcHandlerOptions,
|
||||
HostConnectionHandle,
|
||||
@@ -24,7 +25,7 @@ export { API_PATH, HOST_EVENTS_PATH, MUX_EVENTS_PATH } from './api-path.ts'
|
||||
/** Stable Cordis plugin name. */
|
||||
export const name = 'client-connection'
|
||||
|
||||
/** Services required before providing Connection; legacy `/api` attaches when apiProxy is present. */
|
||||
/** Services required before providing Connection; API Proxy is an optional `/api` fallback. */
|
||||
export const inject = ['httpServer']
|
||||
|
||||
/** Plugin config: the deployment's non-loopback serving authorities. */
|
||||
@@ -93,35 +94,44 @@ export function apply(ctx: Context, config?: ConnectionConfig): void {
|
||||
// Config boundary: a malformed entry fails the load loudly here rather than
|
||||
// silently authorizing its hostname prefix at request time.
|
||||
for (const entry of trustedHosts) assertTrustedAuthority(entry)
|
||||
new HostConnectionService(ctx, trustedHosts)
|
||||
const connection = new HostConnectionService(ctx, trustedHosts)
|
||||
const fetchHandler = connection.createSharedFetchHandler(API_PATH, {
|
||||
async fetch(request) {
|
||||
const pathname = new URL(request.url).pathname
|
||||
const method = pathname.startsWith(`${API_PATH}/`)
|
||||
? pathname.slice(API_PATH.length + 1)
|
||||
: undefined
|
||||
if (method !== undefined
|
||||
&& PRIVILEGED_METHODS.has(method)
|
||||
&& !isTrustedApiRequest(request, [])) {
|
||||
return new Response('forbidden', { status: 403 })
|
||||
}
|
||||
if (request.method === 'GET' && (pathname === MUX_EVENTS_PATH || pathname === HOST_EVENTS_PATH)) {
|
||||
return new Response('upgrade required', {
|
||||
status: 426,
|
||||
headers: { connection: 'Upgrade', upgrade: 'websocket' },
|
||||
})
|
||||
}
|
||||
const apiProxy = ctx.get('apiProxy')
|
||||
if (apiProxy === undefined) return new Response('not found', { status: 404 })
|
||||
return toFetchHandler(apiProxy).fetch(request)
|
||||
},
|
||||
})
|
||||
const route: WebRoute = {
|
||||
kind: 'prefix',
|
||||
path: API_PATH,
|
||||
handler: async (req, res) => {
|
||||
if (!isTrustedApiRequest(req, trustedHosts)) {
|
||||
res.writeHead(403)
|
||||
res.end('forbidden')
|
||||
return
|
||||
}
|
||||
await bridge(req, res, fetchHandler)
|
||||
},
|
||||
}
|
||||
ctx.effect(() => ctx.httpServer.register(route), 'client-connection: /api route')
|
||||
ctx.inject(['apiProxy'], (apiCtx) => {
|
||||
const apiHandler = toFetchHandler(apiCtx.apiProxy)
|
||||
const downlinks = new WebSocketDownlinks(apiCtx.apiProxy)
|
||||
const route: WebRoute = {
|
||||
kind: 'prefix',
|
||||
path: API_PATH,
|
||||
handler: async (req, res) => {
|
||||
const pathname = new URL(req.url ?? '/', 'http://dsh.internal').pathname
|
||||
const method = pathname.startsWith(`${API_PATH}/`)
|
||||
? pathname.slice(API_PATH.length + 1)
|
||||
: undefined
|
||||
const allowed = method !== undefined && PRIVILEGED_METHODS.has(method)
|
||||
? isTrustedApiRequest(req, [])
|
||||
: isTrustedApiRequest(req, trustedHosts)
|
||||
if (!allowed) {
|
||||
res.writeHead(403)
|
||||
res.end('forbidden')
|
||||
return
|
||||
}
|
||||
if (req.method === 'GET' && (pathname === MUX_EVENTS_PATH || pathname === HOST_EVENTS_PATH)) {
|
||||
res.writeHead(426, { connection: 'Upgrade', upgrade: 'websocket' })
|
||||
res.end('upgrade required')
|
||||
return
|
||||
}
|
||||
await bridge(req, res, apiHandler)
|
||||
},
|
||||
}
|
||||
apiCtx.effect(() => apiCtx.httpServer.register(route), 'client-connection: /api route')
|
||||
const registerDownlink = (
|
||||
path: string,
|
||||
handle: WebUpgradeRoute['handler'],
|
||||
|
||||
@@ -11,9 +11,11 @@ import {
|
||||
type RpcId as RpcIdType,
|
||||
type ServerResponse as RpcServerResponse,
|
||||
} from '@deepseek-ai/dsh-host-apiproxy/api'
|
||||
import { bridge } from './http-bridge.ts'
|
||||
import { bridge, type FetchHandler } from './http-bridge.ts'
|
||||
import { isTrustedApiRequest } from './api-request-trust.ts'
|
||||
import { API_PATH } from './api-path.ts'
|
||||
import type {
|
||||
ConnectionRpcEndpointMatcher,
|
||||
ConnectionRpcHandler,
|
||||
ConnectionRpcHandlerOptions,
|
||||
HostConnectionHandle,
|
||||
@@ -24,8 +26,23 @@ 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
|
||||
readonly options: ConnectionRpcHandlerOptions
|
||||
}
|
||||
|
||||
declare module '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.
|
||||
@@ -40,6 +57,33 @@ export class HostConnectionService extends Service implements HostConnectionHand
|
||||
const owner = this.ctx
|
||||
return {
|
||||
handle: (channel, handler, options) => this.register(owner, channel, handler, options),
|
||||
intercept: (channel, matches, handler, options) =>
|
||||
this.registerInterceptor(owner, channel, matches, handler, options),
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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)
|
||||
}
|
||||
if (interceptor.options.authority === 'loopback' && !isTrustedApiRequest(request, [])) {
|
||||
return Promise.resolve(new Response('forbidden', { status: 403 }))
|
||||
}
|
||||
return interceptor.fetchHandler.fetch(request)
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -69,12 +113,38 @@ export class HostConnectionService extends Service implements HostConnectionHand
|
||||
`client-connection: ${channel} rpc channel`,
|
||||
)
|
||||
}
|
||||
|
||||
private registerInterceptor(
|
||||
owner: Context,
|
||||
channel: string,
|
||||
matches: ConnectionRpcEndpointMatcher,
|
||||
handler: ConnectionRpcHandler,
|
||||
options: ConnectionRpcHandlerOptions,
|
||||
): () => 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),
|
||||
options,
|
||||
}
|
||||
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,
|
||||
): { fetch(request: Request): Promise<Response> } {
|
||||
): FetchHandler {
|
||||
return {
|
||||
async fetch(request: Request): Promise<Response> {
|
||||
const endpoint = endpointFromPath(channel, new URL(request.url).pathname)
|
||||
|
||||
@@ -18,11 +18,14 @@ export type ConnectionRpcHandler = (
|
||||
signal: AbortSignal,
|
||||
) => Promise<RpcResult<unknown>>
|
||||
|
||||
/** Synchronous ownership test for one endpoint on a shared RPC channel. */
|
||||
export type ConnectionRpcEndpointMatcher = (endpoint: string) => boolean
|
||||
|
||||
/** Host registry for logical RPC channels carried by the current transport. */
|
||||
export interface HostConnectionRpc {
|
||||
/**
|
||||
* Register one absolute channel prefix and its trust policy.
|
||||
* @param channel - absolute logical channel such as `/api2`.
|
||||
* @param channel - absolute logical channel such as `/rpc`.
|
||||
* @param handler - decoded endpoint handler returning the existing RPC result shape.
|
||||
* @param options - channel trust policy.
|
||||
* @returns asynchronous disposer removing the channel and its physical route.
|
||||
@@ -32,6 +35,21 @@ export interface HostConnectionRpc {
|
||||
handler: ConnectionRpcHandler,
|
||||
options: ConnectionRpcHandlerOptions,
|
||||
): () => Promise<void>
|
||||
|
||||
/**
|
||||
* Intercept owned endpoints on the shared `/api` channel before its fallback.
|
||||
* @param channel - reserved shared channel; currently `/api`.
|
||||
* @param matches - synchronous endpoint ownership test.
|
||||
* @param handler - decoded endpoint handler returning the existing RPC result shape.
|
||||
* @param options - trust policy for every endpoint claimed by this interceptor.
|
||||
* @returns asynchronous disposer removing the interceptor.
|
||||
*/
|
||||
intercept(
|
||||
channel: '/api',
|
||||
matches: ConnectionRpcEndpointMatcher,
|
||||
handler: ConnectionRpcHandler,
|
||||
options: ConnectionRpcHandlerOptions,
|
||||
): () => Promise<void>
|
||||
}
|
||||
|
||||
/** Host `ctx.connection` shape consumed by transport-independent adapters. */
|
||||
@@ -44,7 +62,7 @@ export interface HostConnectionHandle {
|
||||
export interface ClientConnectionRpc {
|
||||
/**
|
||||
* Call one endpoint through an already registered logical channel.
|
||||
* @param channel - absolute logical channel such as `/api2`.
|
||||
* @param channel - absolute logical channel such as `/api`.
|
||||
* @param endpoint - channel-relative endpoint such as `goals/create`.
|
||||
* @param payload - channel-owned request payload.
|
||||
* @param signal - optional caller cancellation.
|
||||
|
||||
@@ -204,7 +204,7 @@ describe('connection client apply', () => {
|
||||
expect(sockets[0]?.readyState).toBe(FakeWebSocket.CLOSED)
|
||||
})
|
||||
|
||||
it('carries generic RPC calls over the isolated channel with rpcId echo validation', async () => {
|
||||
it('carries RPC calls over the shared API channel with rpcId echo validation', async () => {
|
||||
;(globalThis as Win).location = { hostname: 'localhost', search: '' }
|
||||
const handle = await mount()
|
||||
const original = globalThis.fetch
|
||||
@@ -221,13 +221,13 @@ describe('connection client apply', () => {
|
||||
})
|
||||
}
|
||||
try {
|
||||
await expect(handle.rpc.call('/api2', 'goals/create', { args: { agentId: 'agent-1' } }))
|
||||
await expect(handle.rpc.call('/api', 'goals/create', { args: { agentId: 'agent-1' } }))
|
||||
.resolves.toEqual({ ok: true, value: { ref: 'goal-1' } })
|
||||
} finally {
|
||||
globalThis.fetch = original
|
||||
}
|
||||
expect(seen).toHaveLength(1)
|
||||
expect(seen[0]?.url).toBe('http://dsh.internal/api2/goals/create')
|
||||
expect(seen[0]?.url).toBe('http://dsh.internal/api/goals/create')
|
||||
expect(seen[0]?.body).toMatchObject({
|
||||
type: 'client-request',
|
||||
method: 'goals/create',
|
||||
@@ -244,10 +244,10 @@ describe('connection client apply', () => {
|
||||
const abort = new AbortController()
|
||||
globalThis.fetch = vi.fn().mockResolvedValue(new Response('unavailable', { status: 503 }))
|
||||
try {
|
||||
await expect(handle.rpc.call('/api2', 'goals/create', {}, abort.signal))
|
||||
await expect(handle.rpc.call('/api', 'goals/create', {}, abort.signal))
|
||||
.rejects.toThrow('HTTP 503')
|
||||
expect(globalThis.fetch).toHaveBeenCalledWith(
|
||||
new URL('https://harness.example/api2/goals/create'),
|
||||
new URL('https://harness.example/api/goals/create'),
|
||||
expect.objectContaining({ signal: abort.signal }),
|
||||
)
|
||||
|
||||
@@ -257,9 +257,9 @@ describe('connection client apply', () => {
|
||||
rpcId: 'different-rpc',
|
||||
result: { ok: true, value: null },
|
||||
}))
|
||||
await expect(handle.rpc.call('/api2', 'goals/create', {})).rejects.toThrow('rpcId mismatch')
|
||||
await expect(handle.rpc.call('/api', 'goals/create', {})).rejects.toThrow('rpcId mismatch')
|
||||
const fetch = vi.mocked(globalThis.fetch)
|
||||
expect(fetch.mock.calls[0]?.[0]).toEqual(new URL('http://dsh.internal/api2/goals/create'))
|
||||
expect(fetch.mock.calls[0]?.[0]).toEqual(new URL('http://dsh.internal/api/goals/create'))
|
||||
expect(fetch.mock.calls[0]?.[1]).not.toHaveProperty('signal')
|
||||
} finally {
|
||||
globalThis.fetch = original
|
||||
@@ -267,12 +267,12 @@ describe('connection client apply', () => {
|
||||
|
||||
for (const [channel, endpoint] of [
|
||||
['api2', 'goals/create'],
|
||||
['/api2/path', 'goals/create'],
|
||||
['/api2', ''],
|
||||
['/api2', '.'],
|
||||
['/api2', '..'],
|
||||
['/api2', 'goals//create'],
|
||||
['/api2', 'goals/create?unsafe'],
|
||||
['/api/path', 'goals/create'],
|
||||
['/api', ''],
|
||||
['/api', '.'],
|
||||
['/api', '..'],
|
||||
['/api', 'goals//create'],
|
||||
['/api', 'goals/create?unsafe'],
|
||||
] as const) {
|
||||
await expect(handle.rpc.call(channel, endpoint, {})).rejects.toThrow('invalid RPC target')
|
||||
}
|
||||
@@ -281,6 +281,6 @@ describe('connection client apply', () => {
|
||||
it('keeps generic Remote calls unavailable in the client-only fixture', async () => {
|
||||
;(globalThis as Win).location = { hostname: 'localhost', search: '?fixture' }
|
||||
const handle = await mount()
|
||||
await expect(handle.rpc.call('/api2', 'goals/create', {})).rejects.toThrow(/unavailable in fixture mode/)
|
||||
await expect(handle.rpc.call('/api', 'goals/create', {})).rejects.toThrow(/unavailable in fixture mode/)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -195,35 +195,36 @@ describe('connection node half', () => {
|
||||
await dispose()
|
||||
})
|
||||
|
||||
it('provides a disposable generic RPC channel without requiring apiProxy', async () => {
|
||||
it('provides a disposable dedicated RPC channel without requiring apiProxy', async () => {
|
||||
const ctx = new Context()
|
||||
const routes: WebRoute[] = []
|
||||
ctx.provide('httpServer', fakeHttpServer(routes, []) as HttpServerService)
|
||||
const fiber = ctx.plugin({ inject: [...inject], apply })
|
||||
await fiber.await()
|
||||
expect(routes).toHaveLength(0)
|
||||
expect(routes).toHaveLength(1)
|
||||
expect(routes[0]).toMatchObject({ kind: 'prefix', path: API_PATH })
|
||||
|
||||
const connection = ctx.get('connection') as HostConnectionHandle
|
||||
const calls: unknown[] = []
|
||||
const remove = connection.rpc.handle('/api2', async (endpoint, payload) => {
|
||||
const remove = connection.rpc.handle('/rpc', async (endpoint, payload) => {
|
||||
calls.push({ endpoint, payload })
|
||||
return { ok: true, value: { accepted: true } }
|
||||
}, { authority: 'trusted-host' })
|
||||
const route = routes.find(candidate => candidate.path === '/api2')
|
||||
const route = routes.find(candidate => candidate.path === '/rpc')
|
||||
expect(route).toBeDefined()
|
||||
|
||||
const request: ClientRequest = {
|
||||
type: 'client-request',
|
||||
rpcId: RpcId('rpc-api2'),
|
||||
rpcId: RpcId('rpc-dedicated'),
|
||||
method: 'goals/create',
|
||||
payload: { args: { agentId: 'agent-1' } },
|
||||
}
|
||||
const result = fakeResponse()
|
||||
await route!.handler(fakePost({ host: '127.0.0.1:3080' }, '/api2/goals/create', request), result.response)
|
||||
await route!.handler(fakePost({ host: '127.0.0.1:3080' }, '/rpc/goals/create', request), result.response)
|
||||
expect(result.state.status).toBe(200)
|
||||
expect(JSON.parse(String(result.state.body))).toEqual({
|
||||
type: 'server-response',
|
||||
rpcId: 'rpc-api2',
|
||||
rpcId: 'rpc-dedicated',
|
||||
result: { ok: true, value: { accepted: true } },
|
||||
})
|
||||
expect(calls).toEqual([{
|
||||
@@ -231,11 +232,90 @@ describe('connection node half', () => {
|
||||
payload: { args: { agentId: 'agent-1' } },
|
||||
}])
|
||||
|
||||
expect(() => connection.rpc.handle('/api2', async () => ({ ok: true, value: null }), {
|
||||
expect(() => connection.rpc.handle('/rpc', async () => ({ ok: true, value: null }), {
|
||||
authority: 'trusted-host',
|
||||
})).toThrow(/duplicate route/)
|
||||
await remove()
|
||||
expect(routes.map(candidate => candidate.path)).toEqual([API_PATH])
|
||||
await fiber.dispose()
|
||||
expect(routes).toHaveLength(0)
|
||||
})
|
||||
|
||||
it('dispatches claimed /api endpoints before the API Proxy fallback and withdraws the claim', async () => {
|
||||
const ctx = new Context()
|
||||
const routes: WebRoute[] = []
|
||||
ctx.provide('httpServer', fakeHttpServer(routes, []) as HttpServerService)
|
||||
ctx.provide('apiProxy', {} as unknown as ApiProxy)
|
||||
const fiber = ctx.plugin({ inject: [...inject], apply }, { trustedHosts: ['harness.example'] })
|
||||
await fiber.await()
|
||||
const connection = ctx.get('connection') as HostConnectionHandle
|
||||
const calls: unknown[] = []
|
||||
const remove = connection.rpc.intercept(
|
||||
'/api',
|
||||
endpoint => endpoint === 'goals/create',
|
||||
async (endpoint, payload) => {
|
||||
calls.push({ endpoint, payload })
|
||||
return { ok: true, value: { accepted: true } }
|
||||
},
|
||||
{ authority: 'trusted-host' },
|
||||
)
|
||||
expect(() => connection.rpc.intercept(
|
||||
'/api',
|
||||
() => true,
|
||||
async () => ({ ok: true, value: null }),
|
||||
{ authority: 'trusted-host' },
|
||||
)).toThrow('already has an interceptor')
|
||||
expect(() => connection.rpc.intercept(
|
||||
'/rpc' as '/api',
|
||||
() => true,
|
||||
async () => ({ ok: true, value: null }),
|
||||
{ authority: 'trusted-host' },
|
||||
)).toThrow('invalid shared RPC channel')
|
||||
const route = routes.find(candidate => candidate.path === API_PATH)!
|
||||
const request: ClientRequest = {
|
||||
type: 'client-request',
|
||||
rpcId: RpcId('rpc-shared'),
|
||||
method: 'goals/create',
|
||||
payload: { args: { agentId: 'agent-1' } },
|
||||
}
|
||||
|
||||
const claimed = fakeResponse()
|
||||
await route.handler(fakePost({ host: '127.0.0.1:3080' }, '/api/goals/create', request), claimed.response)
|
||||
expect(JSON.parse(String(claimed.state.body))).toEqual({
|
||||
type: 'server-response',
|
||||
rpcId: 'rpc-shared',
|
||||
result: { ok: true, value: { accepted: true } },
|
||||
})
|
||||
expect(calls).toEqual([{
|
||||
endpoint: 'goals/create',
|
||||
payload: { args: { agentId: 'agent-1' } },
|
||||
}])
|
||||
|
||||
const denied = fakeResponse()
|
||||
await route.handler(fakePost({ host: 'other.example' }, '/api/goals/create', request), denied.response)
|
||||
expect(denied.state).toMatchObject({ status: 403, body: 'forbidden' })
|
||||
expect(calls).toHaveLength(1)
|
||||
|
||||
const unclaimed = fakeResponse()
|
||||
await route.handler(fakeRequest({ host: '127.0.0.1:3080' }, '/api/session.list'), unclaimed.response)
|
||||
expect(unclaimed.state.status).toBe(404)
|
||||
|
||||
await remove()
|
||||
const withdrawn = fakeResponse()
|
||||
await route.handler(fakePost({ host: '127.0.0.1:3080' }, '/api/goals/create', request), withdrawn.response)
|
||||
expect(withdrawn.state.status).toBe(404)
|
||||
expect(calls).toHaveLength(1)
|
||||
|
||||
const removeLoopback = connection.rpc.intercept(
|
||||
'/api',
|
||||
endpoint => endpoint === 'goals/create',
|
||||
async () => ({ ok: true, value: null }),
|
||||
{ authority: 'loopback' },
|
||||
)
|
||||
const loopbackOnly = fakeResponse()
|
||||
await route.handler(fakePost({ host: 'harness.example' }, '/api/goals/create', request), loopbackOnly.response)
|
||||
expect(loopbackOnly.state.status).toBe(403)
|
||||
await removeLoopback()
|
||||
await fiber.dispose()
|
||||
})
|
||||
|
||||
@@ -246,20 +326,20 @@ describe('connection node half', () => {
|
||||
const fiber = ctx.plugin({ inject: [...inject], apply }, { trustedHosts: ['harness.example'] })
|
||||
await fiber.await()
|
||||
const connection = ctx.get('connection') as HostConnectionHandle
|
||||
const remove = connection.rpc.handle('/api2', async (endpoint) => {
|
||||
const remove = connection.rpc.handle('/rpc', async (endpoint) => {
|
||||
if (endpoint === 'fail') throw new Error('handler broke')
|
||||
return { ok: true, value: null }
|
||||
}, {
|
||||
authority: 'trusted-host',
|
||||
})
|
||||
const route = routes[0]!
|
||||
const route = routes.find(candidate => candidate.path === '/rpc')!
|
||||
|
||||
const denied = fakeResponse()
|
||||
await route.handler(fakePost({ host: 'other.example' }, '/api2/goals/create', {}), denied.response)
|
||||
await route.handler(fakePost({ host: 'other.example' }, '/rpc/goals/create', {}), denied.response)
|
||||
expect(denied.state).toMatchObject({ status: 403, body: 'forbidden' })
|
||||
|
||||
const methodMismatch = fakeResponse()
|
||||
await route.handler(fakePost({ host: 'harness.example' }, '/api2/goals/create', {
|
||||
await route.handler(fakePost({ host: 'harness.example' }, '/rpc/goals/create', {
|
||||
type: 'client-request', rpcId: 'rpc-bad', method: 'other', payload: {},
|
||||
}), methodMismatch.response)
|
||||
expect(JSON.parse(String(methodMismatch.state.body))).toMatchObject({
|
||||
@@ -268,12 +348,12 @@ describe('connection node half', () => {
|
||||
})
|
||||
|
||||
for (const [request, status] of [
|
||||
[fakeRequest({ host: 'harness.example' }, '/api2/goals/create'), 404],
|
||||
[fakeRequest({ host: 'harness.example' }, '/rpc/goals/create'), 404],
|
||||
[fakePost({ host: 'harness.example' }, '/outside/goals/create', {}), 404],
|
||||
[fakePost({ host: 'harness.example' }, '/api2/goals//create', {}), 404],
|
||||
[fakeRawPost({ host: 'harness.example' }, '/api2/goals/create', '{}'), 415],
|
||||
[fakeRawPost({ host: 'harness.example', 'content-type': 'text/plain' }, '/api2/goals/create', '{}'), 415],
|
||||
[fakeRawPost({ host: 'harness.example', 'content-type': 'application/json; charset=utf-8' }, '/api2/goals/create', '{'), 400],
|
||||
[fakePost({ host: 'harness.example' }, '/rpc/goals//create', {}), 404],
|
||||
[fakeRawPost({ host: 'harness.example' }, '/rpc/goals/create', '{}'), 415],
|
||||
[fakeRawPost({ host: 'harness.example', 'content-type': 'text/plain' }, '/rpc/goals/create', '{}'), 415],
|
||||
[fakeRawPost({ host: 'harness.example', 'content-type': 'application/json; charset=utf-8' }, '/rpc/goals/create', '{'), 400],
|
||||
] as const) {
|
||||
const response = fakeResponse()
|
||||
await route.handler(request, response.response)
|
||||
@@ -286,7 +366,7 @@ describe('connection node half', () => {
|
||||
[null, 'invalid-request'],
|
||||
] as const) {
|
||||
const response = fakeResponse()
|
||||
await route.handler(fakePost({ host: 'harness.example' }, '/api2/goals/create', body), response.response)
|
||||
await route.handler(fakePost({ host: 'harness.example' }, '/rpc/goals/create', body), response.response)
|
||||
expect(JSON.parse(String(response.state.body))).toMatchObject({
|
||||
rpcId,
|
||||
result: { ok: false, error: { code: 'bad-request' } },
|
||||
@@ -294,7 +374,7 @@ describe('connection node half', () => {
|
||||
}
|
||||
|
||||
const failed = fakeResponse()
|
||||
await route.handler(fakePost({ host: 'harness.example' }, '/api2/fail', {
|
||||
await route.handler(fakePost({ host: 'harness.example' }, '/rpc/fail', {
|
||||
type: 'client-request', rpcId: 'rpc-fail', method: 'fail', payload: {},
|
||||
}), failed.response)
|
||||
expect(failed.state).toMatchObject({ status: 500, body: 'handler failure: Error: handler broke' })
|
||||
|
||||
@@ -2,5 +2,5 @@
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write packages/host/api-gateway/README.md
|
||||
README.md: 3ef926ace2ee4d6008b1d6c18b1e070fa39bc176
|
||||
README.zh.md: 77b8b8a87d5f511000aac5cf9f75ebca5fcdfbca
|
||||
README.md: cc80bb19fec15414aa0857154a8a36fb4f642672
|
||||
README.zh.md: 6febb1cfe4fc7fa4c5a17e1e4f6a21e2ee03e295
|
||||
|
||||
@@ -10,13 +10,13 @@ Two-sided Remote control for Host and Client Cordis environments. The Host entry
|
||||
|
||||
Strict mode reads generated invocation descriptors from `ctx.typert.local`. Lookup parameters use registered `ctx.typert.lookups` providers, while `@RemoteContext` resolves its receiver through a registered Host Context provider. SRC mode is a development fallback for endpoints that have never had a strict definition; it parses simple parameter names and accepts only JSON-safe values for non-lookup parameters. Withdrawing an observed strict definition fails instead of weakening validation.
|
||||
|
||||
The Host entry registers the trusted-host `/api2` unary RPC channel when Connection is available. Direct `invoke()` calls preserve business errors; `TypertGatewayError` distinguishes failures owned by dispatch, binding, providers, lookup, Context, arguments, and codecs.
|
||||
The Host entry registers a trusted-host interceptor on Connection's shared `/api` FetchHandler. Connection passes this composite handler through its HTTP bridge; the handler dispatches claimed endpoints to Gateway and unclaimed endpoints to API Proxy. Direct `invoke()` calls preserve business errors; `TypertGatewayError` distinguishes failures owned by dispatch, binding, providers, lookup, Context, arguments, and codecs.
|
||||
|
||||
## Client service: `ClientApi` (ctx key: `api`)
|
||||
|
||||
`ctx.api.mount()` validates and registers a generated Host-for-Client contribution, then installs concrete direct and scoped methods for the calling Cordis fiber. Duplicate endpoints, namespace collisions, and descriptors without strict generated codecs fail before methods become callable.
|
||||
|
||||
Each call validates positional inputs, constructs the descriptor's exact named `args`, and sends it through `ctx.connection.rpc.call('/api2', endpoint, ...)`. The returned value is validated before reaching application code. Withdrawing a contribution removes its descriptors and methods together, aborts in-flight calls, and makes retained method handles reject.
|
||||
Each call validates positional inputs, constructs the descriptor's exact named `args`, and sends it through `ctx.connection.rpc.call('/api', endpoint, ...)`. The returned value is validated before reaching application code. Withdrawing a contribution removes its descriptors and methods together, aborts in-flight calls, and makes retained method handles reject.
|
||||
|
||||
Generated declaration merges provide the TypeScript API. The Client entry contains no Host Service or Host Cordis interface merge, and method lookup and invocation use ordinary objects and functions rather than a JavaScript Proxy.
|
||||
|
||||
|
||||
@@ -10,13 +10,13 @@
|
||||
|
||||
严格模式从 `ctx.typert.local` 读取生成的调用描述符。查找参数使用已向 `ctx.typert.lookups` 注册的提供方,`@RemoteContext` 则通过已注册的 Host Context 提供方解析其接收者。SRC 模式是开发阶段的回退路径,适用于从未具备严格定义的端点;它解析简单参数名,并且只允许非查找参数使用可安全表示为 JSON 的值。已观测到的严格定义一旦撤回,系统会直接报错,而不会降低校验强度。
|
||||
|
||||
Connection 可用时,Host 入口会注册 trusted-host 的 `/api2` 一元 RPC 通道。直接调用 `invoke()` 会保留业务错误;`TypertGatewayError` 可区分分发、绑定、提供方、查找、Context、参数和编解码器各自负责的故障。
|
||||
Connection 可用时,Host 入口会在 Connection 共享的 `/api` FetchHandler 上注册 trusted-host interceptor。Connection 把这个复合 handler 交给 HTTP bridge;handler 将已认领 endpoint 分发给 Gateway,未认领 endpoint 则交给 API Proxy。直接调用 `invoke()` 会保留业务错误;`TypertGatewayError` 可区分分发、绑定、提供方、查找、Context、参数和编解码器各自负责的故障。
|
||||
|
||||
## Client 服务:`ClientApi`(ctx key:`api`)
|
||||
|
||||
`ctx.api.mount()` 会校验并注册生成的 Host-for-Client 贡献项,然后为发起调用的 Cordis fiber 安装具体的直接方法和作用域方法。重复端点、命名空间冲突,以及缺少生成的严格编解码器的描述符,都会在方法可调用前报错。
|
||||
|
||||
每次调用都会校验位置参数,构造与描述符完全匹配的具名 `args`,再通过 `ctx.connection.rpc.call('/api2', endpoint, ...)` 发送。返回值经过校验后才会交给应用代码。撤回贡献项会同时移除其描述符和方法、中止正在进行的调用,并使外部仍持有的方法句柄在调用时返回拒绝。
|
||||
每次调用都会校验位置参数,构造与描述符完全匹配的具名 `args`,再通过 `ctx.connection.rpc.call('/api', endpoint, ...)` 发送。返回值经过校验后才会交给应用代码。撤回贡献项会同时移除其描述符和方法、中止正在进行的调用,并使外部仍持有的方法句柄在调用时返回拒绝。
|
||||
|
||||
生成的声明合并提供 TypeScript API。Client 入口不包含 Host 服务或 Host Cordis 接口合并;方法查找和调用使用普通对象与函数,而不使用 JavaScript Proxy。
|
||||
|
||||
|
||||
@@ -248,7 +248,7 @@ class ClientApiService extends Service implements ClientApi {
|
||||
})
|
||||
const connection = this.ownerCtx.get('connection') as ConnectionHandle | undefined
|
||||
if (connection === undefined) throw new Error(`client api: ${endpoint} has no active Connection`)
|
||||
const result = await connection.rpc.call('/api2', endpoint, { args }, token.abort.signal)
|
||||
const result = await connection.rpc.call('/api', endpoint, { args }, token.abort.signal)
|
||||
if (!mountActive(token)) throw new Error(`client api: Remote method ${endpoint} was withdrawn during invocation`)
|
||||
if (!result.ok) throw remoteFailure(endpoint, result.error)
|
||||
return parse(descriptor.result, result.value, endpoint, 'result')
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
*/
|
||||
|
||||
import { Context, Service, symbols } from 'cordis'
|
||||
import type { ConnectionRpcHandler } from '@deepseek-ai/dsh-client-connection'
|
||||
import {
|
||||
remoteMethods,
|
||||
type InvocationDescriptor,
|
||||
@@ -35,26 +36,7 @@ interface ResolvedBinding {
|
||||
readonly original: object
|
||||
}
|
||||
|
||||
type ConnectionRpcResult =
|
||||
| { readonly ok: true; readonly value: unknown }
|
||||
| {
|
||||
readonly ok: false
|
||||
readonly error: {
|
||||
readonly code: 'internal'
|
||||
readonly message: string
|
||||
readonly details: Record<never, never>
|
||||
}
|
||||
}
|
||||
|
||||
interface HostConnectionLike {
|
||||
readonly rpc: {
|
||||
handle(
|
||||
channel: string,
|
||||
handler: (endpoint: string, payload: unknown, signal: AbortSignal) => Promise<ConnectionRpcResult>,
|
||||
options: { readonly authority: 'trusted-host' | 'loopback' },
|
||||
): () => Promise<void>
|
||||
}
|
||||
}
|
||||
type ConnectionRpcResult = Awaited<ReturnType<ConnectionRpcHandler>>
|
||||
|
||||
/** Dispatch failure produced outside the invoked business method. */
|
||||
export class TypertGatewayError extends Error {
|
||||
@@ -101,15 +83,32 @@ export class TypertGatewayService extends Service implements TypertGateway {
|
||||
constructor(ctx: Context) {
|
||||
super(ctx, 'typertGateway')
|
||||
ctx.inject(['connection'], (connectionCtx) => {
|
||||
const connection = connectionCtx.get('connection') as unknown as HostConnectionLike
|
||||
connection.rpc.handle(
|
||||
'/api2',
|
||||
connectionCtx.connection.rpc.intercept(
|
||||
'/api',
|
||||
endpoint => this.claimsEndpoint(endpoint),
|
||||
(endpoint, payload, signal) => this.dispatchRpc(endpoint, payload, signal),
|
||||
{ authority: 'trusted-host' },
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
private claimsEndpoint(endpoint: string): boolean {
|
||||
const segments = endpoint.split('/')
|
||||
if (segments.length !== 2 || segments[0] === '' || segments[1] === '') return false
|
||||
const [namespace, method] = segments as [string, string]
|
||||
if (this.ctx.typert.local.get(endpoint) !== undefined || this.ctx.typert.local.hasSeen(endpoint)) return true
|
||||
for (const [serviceKey, definition] of Object.entries(this.ctx.reflect.props)) {
|
||||
if (definition.type !== 'service') continue
|
||||
const receiver = this.ctx.get(serviceKey) as unknown
|
||||
if (!isObject(receiver)) continue
|
||||
const original = originalOf(receiver)
|
||||
const binding = Reflect.get(original, 'typertGateway') as unknown
|
||||
if (!isObject(binding) || Reflect.get(binding, 'namespace') !== namespace) continue
|
||||
if (remoteMethods(original).some(candidate => (candidate.exportName ?? candidate.method) === method)) return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
/**
|
||||
* Invoke one live Remote method through strict generated reflection or SRC markers.
|
||||
* @param request - decoded endpoint and exact named wire arguments.
|
||||
|
||||
@@ -109,7 +109,7 @@ describe('Client TypeRT API', () => {
|
||||
|
||||
await expect(ctx.api.goals.create('agent-1', { objective: 'ship' })).resolves.toEqual({ ref: 'goal-1' })
|
||||
expect(call).toHaveBeenCalledWith(
|
||||
'/api2',
|
||||
'/api',
|
||||
'goals/create',
|
||||
{ args: { agentId: 'agent-1', request: { objective: 'ship' } } },
|
||||
expect.any(AbortSignal),
|
||||
@@ -144,7 +144,7 @@ describe('Client TypeRT API', () => {
|
||||
|
||||
await expect(agentCtx.goals.create({ objective: 'ship scoped' })).resolves.toEqual({ ref: 'goal-2' })
|
||||
expect(call).toHaveBeenCalledWith(
|
||||
'/api2',
|
||||
'/api',
|
||||
'goals/create',
|
||||
{ args: { agentId: 'agent-2', request: { objective: 'ship scoped' } } },
|
||||
expect.any(AbortSignal),
|
||||
@@ -175,7 +175,7 @@ describe('Client TypeRT API', () => {
|
||||
|
||||
await expect(agentCtx.goals.rename({ objective: 'land' })).resolves.toEqual({ renamed: true })
|
||||
expect(call).toHaveBeenCalledWith(
|
||||
'/api2',
|
||||
'/api',
|
||||
'goals/rename',
|
||||
{ args: { agentId: 'agent-2', request: { objective: 'land' } } },
|
||||
expect.any(AbortSignal),
|
||||
|
||||
@@ -96,6 +96,7 @@ type FakeRpcHandler = (endpoint: string, payload: unknown, signal: AbortSignal)
|
||||
class FakeConnectionService extends Service {
|
||||
channel: string | undefined
|
||||
authority: string | undefined
|
||||
matches: ((endpoint: string) => boolean) | undefined
|
||||
handler: FakeRpcHandler | undefined
|
||||
|
||||
constructor(ctx: Context) {
|
||||
@@ -105,14 +106,21 @@ class FakeConnectionService extends Service {
|
||||
get rpc() {
|
||||
const owner = this.ctx
|
||||
return {
|
||||
handle: (channel: string, handler: FakeRpcHandler, options: { readonly authority: string }) =>
|
||||
intercept: (
|
||||
channel: string,
|
||||
matches: (endpoint: string) => boolean,
|
||||
handler: FakeRpcHandler,
|
||||
options: { readonly authority: string },
|
||||
) =>
|
||||
owner.effect(() => {
|
||||
this.channel = channel
|
||||
this.authority = options.authority
|
||||
this.matches = matches
|
||||
this.handler = handler
|
||||
return () => {
|
||||
this.channel = undefined
|
||||
this.authority = undefined
|
||||
this.matches = undefined
|
||||
this.handler = undefined
|
||||
}
|
||||
}),
|
||||
@@ -820,7 +828,7 @@ describe('TypertGatewayService', () => {
|
||||
}), 'invocation-unavailable')
|
||||
})
|
||||
|
||||
it('mounts /api2 through an optional Connection and returns existing RPC results', async () => {
|
||||
it('mounts a shared /api interceptor through an optional Connection and returns existing RPC results', async () => {
|
||||
const ctx = new Context().extend({ fixtureScope: 'rpc-caller' })
|
||||
await ctx.plugin(TypertRegistry)
|
||||
await ctx.plugin(FakeConnectionService)
|
||||
@@ -828,13 +836,18 @@ describe('TypertGatewayService', () => {
|
||||
await gatewayFiber
|
||||
await ctx.plugin(GoalService)
|
||||
const connection = rawConnection(ctx)
|
||||
expect(connection).toMatchObject({ channel: '/api2', authority: 'trusted-host' })
|
||||
expect(connection).toMatchObject({ channel: '/api', authority: 'trusted-host' })
|
||||
|
||||
registerAgentLookup(ctx, { id: 'agent-1' })
|
||||
registerStrict(ctx, [createDescriptor()])
|
||||
expect(connection.matches?.('goals/create')).toBe(true)
|
||||
expect(connection.matches?.('goals/passthrough')).toBe(true)
|
||||
expect(connection.matches?.('goals')).toBe(false)
|
||||
expect(connection.matches?.('goals/missing')).toBe(false)
|
||||
expect(connection.matches?.('legacy/list')).toBe(false)
|
||||
const signal = new AbortController().signal
|
||||
const handler = connection.handler
|
||||
if (handler === undefined) throw new Error('fixture Connection did not retain the /api2 handler')
|
||||
if (handler === undefined) throw new Error('fixture Connection did not retain the /api interceptor')
|
||||
await expect(handler('goals/create', {
|
||||
args: { agentId: 'agent-1', request: { title: 'ship' } },
|
||||
}, signal)).resolves.toEqual({
|
||||
@@ -873,7 +886,7 @@ describe('TypertGatewayService', () => {
|
||||
expect(connection.handler).toBeUndefined()
|
||||
})
|
||||
|
||||
it('dispatches a generated invocation through the real /api2 HTTP carrier', async () => {
|
||||
it('dispatches claimed invocations through /api and leaves unclaimed endpoints to its fallback', async () => {
|
||||
const ctx = new Context().extend({ fixtureScope: 'http-caller' })
|
||||
const routes: WebRoute[] = []
|
||||
ctx.provide('httpServer', fakeHttpServer(routes) as HttpServerService)
|
||||
@@ -886,11 +899,12 @@ describe('TypertGatewayService', () => {
|
||||
await goalFiber
|
||||
const removeLookup = registerAgentLookup(ctx, { id: 'agent-1' })
|
||||
const removeStrict = registerStrict(ctx, [createDescriptor()])
|
||||
let strictActive = true
|
||||
expect(routes).toHaveLength(1)
|
||||
const server = await serveRoute(routes[0]!)
|
||||
|
||||
try {
|
||||
const response = await fetch(`${server.origin}/api2/goals/create`, {
|
||||
const response = await fetch(`${server.origin}/api/goals/create`, {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
@@ -909,9 +923,54 @@ describe('TypertGatewayService', () => {
|
||||
value: { agentId: 'agent-1', title: 'ship', scope: 'http-caller' },
|
||||
},
|
||||
})
|
||||
|
||||
const invalid = await fetch(`${server.origin}/api/goals/create`, {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
type: 'client-request',
|
||||
rpcId: 'rpc-invalid',
|
||||
method: 'goals/create',
|
||||
payload: { invalid: true },
|
||||
}),
|
||||
})
|
||||
expect(invalid.status).toBe(200)
|
||||
await expect(invalid.json()).resolves.toMatchObject({
|
||||
type: 'server-response',
|
||||
rpcId: 'rpc-invalid',
|
||||
result: {
|
||||
ok: false,
|
||||
error: { code: 'internal', message: expect.stringContaining('plain-object args field') },
|
||||
},
|
||||
})
|
||||
|
||||
await removeStrict()
|
||||
strictActive = false
|
||||
const withdrawn = await fetch(`${server.origin}/api/goals/create`, {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
type: 'client-request',
|
||||
rpcId: 'rpc-withdrawn',
|
||||
method: 'goals/create',
|
||||
payload: { args: { agentId: 'agent-1', request: { title: 'ship' } } },
|
||||
}),
|
||||
})
|
||||
expect(withdrawn.status).toBe(200)
|
||||
await expect(withdrawn.json()).resolves.toMatchObject({
|
||||
type: 'server-response',
|
||||
rpcId: 'rpc-withdrawn',
|
||||
result: {
|
||||
ok: false,
|
||||
error: { code: 'internal', message: expect.stringContaining('strict definition was withdrawn') },
|
||||
},
|
||||
})
|
||||
|
||||
const unclaimed = await fetch(`${server.origin}/api/legacy/list`, { method: 'POST' })
|
||||
expect(unclaimed.status).toBe(404)
|
||||
} finally {
|
||||
await server.close()
|
||||
await removeStrict()
|
||||
if (strictActive) await removeStrict()
|
||||
await removeLookup()
|
||||
await goalFiber.dispose()
|
||||
await gatewayFiber.dispose()
|
||||
|
||||
Reference in New Issue
Block a user