feat(remote): deliver allowlisted Host events through ctx.remote.$on
api/remotes owns the allowlist and its type projection; type-meta owns the shape predicate, the selection seat, and the internal remote/host-event carrier signal; api/gateway's Client half turns that signal into $on callbacks through a private dispatch. apiproxy forwards each allowlisted emission verbatim in one host/remote-event frame, registered ahead of the derived invalidation frames so frame order is unchanged, and drops the three per-event variants it replaces. Owner packages move their Events declarations into client-safe ./types exports, so a consumer's listener signature is the Host's own declaration.
This commit is contained in:
@@ -5,7 +5,7 @@
|
||||
*/
|
||||
|
||||
import { Service } from '@deepseek-ai/cordis'
|
||||
import type { Context } from '@deepseek-ai/cordis'
|
||||
import type { Context, Events } from '@deepseek-ai/cordis'
|
||||
import type { ConnectionHandle, RpcError } from '@deepseek-ai/dsh-client-connection/client'
|
||||
import type {
|
||||
InvocationDescriptor,
|
||||
@@ -13,6 +13,7 @@ import type {
|
||||
TypeRTCodec,
|
||||
TypeRTDisposer,
|
||||
TypeRTRemoteContribution,
|
||||
TypeRTRemoteEvent,
|
||||
} from '@deepseek-ai/dsh-type-meta'
|
||||
|
||||
interface MountToken {
|
||||
@@ -71,14 +72,20 @@ export function apply(ctx: Context): void {
|
||||
new ClientRemoteService(ctx)
|
||||
}
|
||||
|
||||
/** One subscribed listener after `$on` erased its per-event argument list. */
|
||||
type RemoteEventListener = (...args: never[]) => void
|
||||
|
||||
class ClientRemoteService extends Service implements TypeRTClientRemote {
|
||||
private readonly ownerCtx: Context
|
||||
private readonly namespaces = new Map<string, RemoteNamespaceHandle>()
|
||||
private readonly subscriptions = new Map<string, Set<RemoteEventListener>>()
|
||||
private mutations = Promise.resolve()
|
||||
|
||||
constructor(ctx: Context) {
|
||||
super(ctx, 'remote')
|
||||
this.ownerCtx = ctx
|
||||
ctx.on('remote/host-event', (event, args) => { this.dispatch(event, args) })
|
||||
ctx.effect(() => () => { this.subscriptions.clear() }, 'api-gateway.client.subscriptions')
|
||||
}
|
||||
|
||||
async $mount(contribution: TypeRTRemoteContribution): ReturnType<TypeRTClientRemote['$mount']> {
|
||||
@@ -91,6 +98,49 @@ class ClientRemoteService extends Service implements TypeRTClientRemote {
|
||||
return async () => { await owned() }
|
||||
}
|
||||
|
||||
$on<Event extends TypeRTRemoteEvent>(
|
||||
event: Event,
|
||||
listener: Events[Event],
|
||||
): ReturnType<TypeRTClientRemote['$on']> {
|
||||
// The table is keyed by the runtime event name, so the argument list this
|
||||
// signature pins per event cannot survive in it; `$deliver` restores it
|
||||
// from the frame the Host emitted for that same name.
|
||||
const erased: RemoteEventListener = listener
|
||||
const owned = this.ctx.effect(() => {
|
||||
const listeners = this.listeners(event)
|
||||
listeners.add(erased)
|
||||
return () => { listeners.delete(erased) }
|
||||
}, `api-gateway.client.$on(${JSON.stringify(event)})`)
|
||||
return () => { void owned() }
|
||||
}
|
||||
|
||||
/**
|
||||
* Deliver one forwarded event in registration order, isolating a throwing
|
||||
* listener; an event name nobody subscribes to is dropped, since the wire
|
||||
* carries whatever the Host forwarding allowlist selected.
|
||||
*/
|
||||
private dispatch(event: string, args: readonly unknown[]): void {
|
||||
const listeners = this.subscriptions.get(event)
|
||||
if (listeners === undefined) return
|
||||
for (const listener of listeners) {
|
||||
try {
|
||||
listener(...args as never[])
|
||||
} catch (error) {
|
||||
console.error(`client api: Remote event ${JSON.stringify(event)} listener threw:`, error)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Subscription set for one event name; empty sets are retained, bounded by the Host's selection. */
|
||||
private listeners(event: string): Set<RemoteEventListener> {
|
||||
let listeners = this.subscriptions.get(event)
|
||||
if (listeners === undefined) {
|
||||
listeners = new Set()
|
||||
this.subscriptions.set(event, listeners)
|
||||
}
|
||||
return listeners
|
||||
}
|
||||
|
||||
private enqueue<T>(operation: () => T | Promise<T>): Promise<T> {
|
||||
const result = this.mutations.then(operation, operation)
|
||||
this.mutations = result.then(() => undefined, () => undefined)
|
||||
|
||||
@@ -16,7 +16,8 @@ export const inject = ['invariants']
|
||||
|
||||
/**
|
||||
* No runtime invariant: Host calls re-read authoritative Cordis and TypeRT
|
||||
* state, while Client methods and descriptors mutate in one owned effect.
|
||||
* state, while Client methods, descriptors, and `$on` subscriptions mutate in
|
||||
* one owned effect.
|
||||
*/
|
||||
const install: InvariantInstaller = () => {}
|
||||
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { Context, Service } from '@deepseek-ai/cordis'
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { Fiber } from '@deepseek-ai/cordis'
|
||||
import { describe, expect, expectTypeOf, it, vi } from 'vitest'
|
||||
import { z } from 'zod'
|
||||
import type { ConnectionHandle } from '@deepseek-ai/dsh-client-connection/client'
|
||||
import type {
|
||||
@@ -10,9 +11,32 @@ import type {
|
||||
TypeRTRemoteNamespace,
|
||||
} from '@deepseek-ai/dsh-type-meta'
|
||||
import TypertRegistry from '@deepseek-ai/dsh-typert-registry'
|
||||
import type { ClientRemote } from '../src/client/index.ts'
|
||||
import { apply, inject } from '../src/client/index.ts'
|
||||
|
||||
declare module '@deepseek-ai/cordis' {
|
||||
interface Events {
|
||||
/**
|
||||
* Test-only forwarded Host event.
|
||||
* @param namespace - marker payload recorded by listeners.
|
||||
*/
|
||||
'fixture/changed'(namespace: string): void
|
||||
/**
|
||||
* Test-only forwarded Host event nobody subscribes to.
|
||||
* @param count - marker payload never observed.
|
||||
*/
|
||||
'fixture/idle'(count: number): void
|
||||
/**
|
||||
* Test-only event the Host assembly does not forward.
|
||||
* @param flag - marker payload never delivered.
|
||||
*/
|
||||
'fixture/unselected'(flag: boolean): void
|
||||
}
|
||||
}
|
||||
|
||||
declare module '@deepseek-ai/dsh-type-meta' {
|
||||
interface TypeRTRemoteEventSelection extends Record<'fixture/changed' | 'fixture/idle', true> {}
|
||||
|
||||
interface TypeRTContextMap {
|
||||
fixture: TypeRTContext<string>
|
||||
}
|
||||
@@ -43,6 +67,19 @@ type FixtureContext = Omit<Context, 'remote'> & {
|
||||
readonly remote: TypeRTClientRemote & TypeRTRemoteScopeApi<'fixture'>
|
||||
}
|
||||
|
||||
// Compile-time contract of `$on`: the key face is the forwarding selection and
|
||||
// the listener signature is the owning package's own Cordis declaration.
|
||||
function remoteEventContracts(remote: ClientRemote): void {
|
||||
remote.$on('fixture/changed', (namespace) => { void namespace })
|
||||
// @ts-expect-error -- declared in Events but outside the forwarding selection.
|
||||
remote.$on('fixture/unselected', () => {})
|
||||
// @ts-expect-error -- not declared in Events at all.
|
||||
remote.$on('fixture/absent', () => {})
|
||||
// @ts-expect-error -- the listener signature comes from the event declaration.
|
||||
remote.$on('fixture/changed', (count: number) => { void count })
|
||||
}
|
||||
void remoteEventContracts
|
||||
|
||||
const idSchema = z.string().min(1)
|
||||
const requestSchema = z.object({ objective: z.string().min(1) })
|
||||
const createResultSchema = z.object({ ref: z.string().min(1) })
|
||||
@@ -96,11 +133,19 @@ function contextDescriptor(): InvocationDescriptor {
|
||||
}
|
||||
|
||||
async function bench(call: ConnectionHandle['rpc']['call']): Promise<Context> {
|
||||
const { ctx } = await benchFiber(call)
|
||||
return ctx
|
||||
}
|
||||
|
||||
async function benchFiber(
|
||||
call: ConnectionHandle['rpc']['call'],
|
||||
): Promise<{ readonly ctx: Context; readonly client: Fiber }> {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(TypertRegistry)
|
||||
ctx.provide('connection', { rpc: { call } } as unknown as ConnectionHandle)
|
||||
await ctx.plugin({ inject, apply })
|
||||
return ctx
|
||||
const client = ctx.plugin({ inject, apply })
|
||||
await client
|
||||
return { ctx, client }
|
||||
}
|
||||
|
||||
describe('Client TypeRT API', () => {
|
||||
@@ -570,4 +615,64 @@ describe('Client TypeRT API', () => {
|
||||
expect(failure.message).toContain('internal: host failed')
|
||||
expect(failure.cause).toBe(rpcError)
|
||||
})
|
||||
|
||||
it('owns each $on subscription in the calling fiber', async () => {
|
||||
const { ctx, client } = await benchFiber(vi.fn<ConnectionHandle['rpc']['call']>())
|
||||
const seen: string[] = []
|
||||
const subscriber = ctx.plugin(Object.assign(
|
||||
(scope: Context) => { scope.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) }) },
|
||||
{ inject: ['remote'] },
|
||||
))
|
||||
await subscriber
|
||||
|
||||
ctx.emit('remote/host-event', 'fixture/changed', ['settings'])
|
||||
expect(seen).toEqual(['settings'])
|
||||
|
||||
await subscriber.dispose()
|
||||
ctx.emit('remote/host-event', 'fixture/changed', ['after fiber disposal'])
|
||||
expect(seen).toEqual(['settings'])
|
||||
|
||||
await client.dispose()
|
||||
expect(ctx.get('remote')).toBeUndefined()
|
||||
})
|
||||
|
||||
it('isolates a throwing listener from the rest of the same event', async () => {
|
||||
const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
|
||||
const consoleError = vi.spyOn(console, 'error').mockImplementation(() => undefined)
|
||||
const seen: string[] = []
|
||||
const disposeFirst = ctx.remote.$on('fixture/changed', () => {
|
||||
throw new Error('fixture listener failure')
|
||||
})
|
||||
ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
|
||||
try {
|
||||
ctx.emit('remote/host-event', 'fixture/changed', ['credentials'])
|
||||
|
||||
expect(seen).toEqual(['credentials'])
|
||||
expect(consoleError).toHaveBeenCalledWith(
|
||||
'client api: Remote event "fixture/changed" listener threw:',
|
||||
expect.any(Error),
|
||||
)
|
||||
disposeFirst()
|
||||
ctx.emit('remote/host-event', 'fixture/changed', ['commands'])
|
||||
expect(seen).toEqual(['credentials', 'commands'])
|
||||
expect(consoleError).toHaveBeenCalledTimes(1)
|
||||
} finally {
|
||||
consoleError.mockRestore()
|
||||
}
|
||||
})
|
||||
|
||||
it('exposes subscription as the only forwarded-event verb', () => {
|
||||
expectTypeOf<ClientRemote>().toHaveProperty('$on')
|
||||
expectTypeOf<ClientRemote>().not.toHaveProperty('$dispatch')
|
||||
})
|
||||
|
||||
it('drops a forwarded event nobody subscribes to', async () => {
|
||||
const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
|
||||
const seen: string[] = []
|
||||
ctx.remote.$on('fixture/changed', (namespace) => { seen.push(namespace) })
|
||||
|
||||
ctx.emit('remote/host-event', 'fixture/idle', [1])
|
||||
|
||||
expect(seen).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user