perf(api-gateway): skip Remote output decoding
This commit is contained in:
parent
4326dd4bca
commit
2d974b187e
6 changed files with 46 additions and 61 deletions
|
|
@ -411,10 +411,10 @@ class ClientRemoteService extends Service implements ClientRemote {
|
|||
const result = await connection.rpc.call('/api', endpoint, { args: prepared.args }, prepared.signal)
|
||||
if (!mountActive(token)) return withdrawn(endpoint)
|
||||
if (!result.ok) return { ok: false, error: result.error }
|
||||
return { ok: true, value: parse(descriptor.result, result.value, endpoint, 'result') }
|
||||
return { ok: true, value: result.value }
|
||||
} catch (error) {
|
||||
// Carrier throws (offline, abort, a rejected result payload) are outcomes
|
||||
// of the call, not assembly faults, so they join the same error branch.
|
||||
// Carrier throws (offline or abort) are outcomes of the call, not assembly
|
||||
// faults, so they join the same error branch.
|
||||
return carrierFailure(endpoint, error)
|
||||
}
|
||||
}
|
||||
|
|
@ -433,7 +433,7 @@ class ClientRemoteService extends Service implements ClientRemote {
|
|||
const stream = this.openRemoteStream(endpoint, { args: prepared.args }, prepared.signal)
|
||||
for await (const value of stream) {
|
||||
if (!mountActive(token)) throw new Error(withdrawn(endpoint).error.message)
|
||||
yield parse(descriptor.result, value, endpoint, 'result')
|
||||
yield value
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -470,12 +470,12 @@ class ClientRemoteService extends Service implements ClientRemote {
|
|||
if (identity === undefined) {
|
||||
throw new Error(`client api: ${endpoint} requires a ${JSON.stringify(projection.context)} Context`)
|
||||
}
|
||||
args[projection.wire] = parse(projection.codec, identity, endpoint, projection.wire)
|
||||
args[projection.wire] = parseInput(projection.codec, identity, endpoint, projection.wire)
|
||||
}
|
||||
let valueIndex = 0
|
||||
descriptor.parameters.forEach((parameter, parameterIndex) => {
|
||||
if (parameterIndex === projection?.parameterIndex) return
|
||||
const value = parse(parameter.codec, values[valueIndex], endpoint, parameter.wire)
|
||||
const value = parseInput(parameter.codec, values[valueIndex], endpoint, parameter.wire)
|
||||
if (value !== undefined) args[parameter.wire] = value
|
||||
valueIndex += 1
|
||||
})
|
||||
|
|
@ -663,7 +663,6 @@ function scopedProjection(descriptor: InvocationDescriptor): ScopedProjection |
|
|||
|
||||
function requireStrictDescriptor(descriptor: InvocationDescriptor): void {
|
||||
const endpoint = endpointOf(descriptor)
|
||||
requireStrictCodec(descriptor.result, endpoint, 'result')
|
||||
for (const parameter of descriptor.parameters) {
|
||||
requireStrictCodec(parameter.codec, endpoint, parameter.wire)
|
||||
}
|
||||
|
|
@ -678,7 +677,7 @@ function requireStrictCodec(codec: TypertCodec, endpoint: string, field: string)
|
|||
}
|
||||
}
|
||||
|
||||
function parse(codec: TypertCodec, value: unknown, endpoint: string, field: string): unknown {
|
||||
function parseInput(codec: TypertCodec, value: unknown, endpoint: string, field: string): unknown {
|
||||
if (codec.mode !== 'strict') {
|
||||
throw new Error(`client api: generated Remote ${endpoint} field ${JSON.stringify(field)} has no strict codec`)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -269,7 +269,7 @@ export class TypertGatewayService extends Service implements TypertGateway {
|
|||
/**
|
||||
* Invoke one live Remote method through strict generated reflection or SRC markers.
|
||||
* @param request - decoded endpoint and exact named wire arguments.
|
||||
* @returns the validated business result.
|
||||
* @returns the business result without output decoding.
|
||||
* @throws {@link TypertGatewayError} for dispatch, provider, or boundary failures; lookup-policy and business errors retain identity.
|
||||
*/
|
||||
async invoke(request: InvokeRemoteRequest): Promise<unknown> {
|
||||
|
|
@ -282,24 +282,18 @@ export class TypertGatewayService extends Service implements TypertGateway {
|
|||
)
|
||||
}
|
||||
|
||||
let result: unknown
|
||||
try {
|
||||
result = await Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown
|
||||
return await Reflect.apply(prepared.method, prepared.receiver, prepared.args) as unknown
|
||||
} catch (error) {
|
||||
if (request.signal?.aborted === true) throw new RemoteInvocationCancelled(prepared.endpoint, error)
|
||||
throw error
|
||||
}
|
||||
// A weak descriptor declares no return type, so nothing returned is a void
|
||||
// result and rides the wire as an absent value field. A strict descriptor
|
||||
// keeps its schema: there, undefined has to be a declared result.
|
||||
if (result === undefined && prepared.descriptor.result.mode !== 'strict') return result
|
||||
return decode(prepared.descriptor.result, result, 'result-invalid', prepared.endpoint, 'result')
|
||||
}
|
||||
|
||||
/**
|
||||
* Open one live stream Remote method without assuming a physical carrier.
|
||||
* @param request - decoded endpoint and named wire arguments.
|
||||
* @returns an iterable whose items have passed the generated result codec.
|
||||
* @returns a cancellation-aware iterable over the business results.
|
||||
*/
|
||||
async stream(request: InvokeRemoteRequest): Promise<AsyncIterable<unknown>> {
|
||||
const prepared = await this.prepareInvocation(request)
|
||||
|
|
@ -325,9 +319,8 @@ export class TypertGatewayService extends Service implements TypertGateway {
|
|||
{ field: 'result' },
|
||||
)
|
||||
}
|
||||
return validatedStream(
|
||||
return cancellableStream(
|
||||
source,
|
||||
prepared.descriptor.result,
|
||||
prepared.endpoint,
|
||||
request.signal ?? NEVER_ABORTED_SIGNAL,
|
||||
)
|
||||
|
|
@ -772,7 +765,7 @@ export class TypertGatewayService extends Service implements TypertGateway {
|
|||
{ field: invocation.wire },
|
||||
)
|
||||
}
|
||||
const identity = decode(invocation.codec, args[invocation.wire], 'input-invalid', endpoint, invocation.wire)
|
||||
const identity = decode(invocation.codec, args[invocation.wire], endpoint, invocation.wire)
|
||||
let context: Context | undefined
|
||||
try {
|
||||
context = await provider.resolve(identity)
|
||||
|
|
@ -806,7 +799,7 @@ export class TypertGatewayService extends Service implements TypertGateway {
|
|||
// still fails decode. Lookup ids are never omissible, so absence here only
|
||||
// ever belongs to a json parameter.
|
||||
if (!Object.hasOwn(args, parameter.wire)) return undefined
|
||||
const value = decode(parameter.codec, args[parameter.wire], 'input-invalid', endpoint, parameter.wire)
|
||||
const value = decode(parameter.codec, args[parameter.wire], endpoint, parameter.wire)
|
||||
if (parameter.source === 'json') return value
|
||||
const key = parameter.lookup
|
||||
/* v8 ignore next -- registry validation rejects strict descriptors without a key, and SRC derivation always supplies one. */
|
||||
|
|
@ -945,9 +938,8 @@ function isIterable(value: unknown): value is Iterable<unknown> | AsyncIterable<
|
|||
|| typeof Reflect.get(value, Symbol.asyncIterator) === 'function')
|
||||
}
|
||||
|
||||
async function *validatedStream(
|
||||
async function *cancellableStream(
|
||||
source: Iterable<unknown> | AsyncIterable<unknown>,
|
||||
codec: TypertCodec,
|
||||
endpoint: string,
|
||||
signal: AbortSignal,
|
||||
): AsyncGenerator {
|
||||
|
|
@ -967,7 +959,7 @@ async function *validatedStream(
|
|||
while (true) {
|
||||
const next = await Promise.race([Promise.resolve(iterator.next()), aborted])
|
||||
if (next.done === true) return
|
||||
yield decode(codec, next.value, 'result-invalid', endpoint, 'result')
|
||||
yield next.value
|
||||
}
|
||||
} finally {
|
||||
signal.removeEventListener('abort', onAbort)
|
||||
|
|
@ -1128,7 +1120,6 @@ function assertExactArguments(
|
|||
function decode(
|
||||
codec: TypertCodec,
|
||||
value: unknown,
|
||||
code: 'input-invalid' | 'result-invalid',
|
||||
endpoint: string,
|
||||
field: string,
|
||||
): unknown {
|
||||
|
|
@ -1141,11 +1132,9 @@ function decode(
|
|||
return value
|
||||
} catch (cause) {
|
||||
throw new TypertGatewayError(
|
||||
code,
|
||||
'input-invalid',
|
||||
endpoint,
|
||||
code === 'input-invalid'
|
||||
? `wire field ${JSON.stringify(field)} failed boundary validation`
|
||||
: 'business result failed boundary validation',
|
||||
`wire field ${JSON.stringify(field)} failed boundary validation`,
|
||||
{ cause, field },
|
||||
)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -131,7 +131,7 @@ export interface TypertGateway {
|
|||
/**
|
||||
* Invoke one live Remote method without assuming a carrier or response envelope.
|
||||
* @param request - decoded endpoint and named wire arguments.
|
||||
* @returns the validated business result.
|
||||
* @returns the business result without output decoding.
|
||||
* @throws {@link TypertGatewayError} for dispatch, provider, or boundary failures; lookup-policy and business errors retain identity.
|
||||
*/
|
||||
invoke(request: InvokeRemoteRequest): Promise<unknown>
|
||||
|
|
@ -139,7 +139,7 @@ export interface TypertGateway {
|
|||
/**
|
||||
* Open one live stream Remote method without assuming a physical carrier.
|
||||
* @param request - decoded endpoint and named wire arguments.
|
||||
* @returns an iterable whose items have passed the generated result codec.
|
||||
* @returns a cancellation-aware iterable over the business results.
|
||||
*/
|
||||
stream(request: InvokeRemoteRequest): Promise<AsyncIterable<unknown>>
|
||||
}
|
||||
|
|
|
|||
|
|
@ -199,7 +199,7 @@ describe('Typert Remote streams', () => {
|
|||
await expect(collect(source)).resolves.toEqual(['wire:one', 'wire:two'])
|
||||
})
|
||||
|
||||
it('validates Iterable and AsyncIterable items and returns the iterator on cancellation', async () => {
|
||||
it('passes Iterable and AsyncIterable items through and returns the iterator on cancellation', async () => {
|
||||
const { ctx, service } = await setup(false)
|
||||
const abort = new AbortController()
|
||||
const source = await ctx.typertGateway.stream({
|
||||
|
|
@ -221,7 +221,10 @@ describe('Typert Remote streams', () => {
|
|||
}))).resolves.toEqual(['b:one', 'b:two'])
|
||||
await expect(collect(await ctx.typertGateway.stream({
|
||||
namespace: 'feed', method: 'invalid', args: {},
|
||||
}))).rejects.toMatchObject({ code: 'result-invalid' })
|
||||
}))).resolves.toEqual([42])
|
||||
await expect(collect(await ctx.typertGateway.stream({
|
||||
namespace: 'feed', method: 'nonJson', args: {},
|
||||
}))).resolves.toEqual([1n])
|
||||
await expect(ctx.typertGateway.stream({
|
||||
namespace: 'feed', method: 'missing', args: {},
|
||||
})).rejects.toMatchObject({ code: 'result-invalid' })
|
||||
|
|
@ -287,9 +290,10 @@ describe('Typert Remote streams', () => {
|
|||
{ type: 'item', streamId: 'sync', value: 's:two' },
|
||||
{ type: 'end', streamId: 'sync' },
|
||||
])
|
||||
expect(frames.find(frame => frame.streamId === 'invalid')).toMatchObject({
|
||||
type: 'error', error: { code: 'internal' },
|
||||
})
|
||||
expect(frames.filter(frame => frame.streamId === 'invalid')).toEqual([
|
||||
{ type: 'item', streamId: 'invalid', value: 42 },
|
||||
{ type: 'end', streamId: 'invalid' },
|
||||
])
|
||||
expect(frames.find(frame => frame.streamId === 'non-json')).toMatchObject({
|
||||
type: 'error', error: { code: 'internal' },
|
||||
})
|
||||
|
|
|
|||
|
|
@ -587,7 +587,7 @@ describe('Client Remote transport readiness', () => {
|
|||
})
|
||||
|
||||
describe('Client Typert API', () => {
|
||||
it('mounts concrete direct methods, validates both boundaries, and withdraws retained handles', async () => {
|
||||
it('mounts concrete direct methods, validates inputs, and withdraws retained handles', async () => {
|
||||
const call = vi.fn<ConnectionHandle['rpc']['call']>()
|
||||
.mockResolvedValue({ ok: true, value: { ref: 'goal-1' } })
|
||||
const ctx = await bench(call)
|
||||
|
|
@ -625,12 +625,8 @@ describe('Client Typert API', () => {
|
|||
|
||||
call.mockResolvedValueOnce({ ok: true, value: { ref: 1 } })
|
||||
await expect(ctx.remote.probe.create('agent-1', { objective: 'ship' })).resolves.toEqual({
|
||||
ok: false,
|
||||
error: {
|
||||
code: 'internal',
|
||||
message: 'client api: probe/create failed: client api: probe/create rejected "result"',
|
||||
details: {},
|
||||
},
|
||||
ok: true,
|
||||
value: { ref: 1 },
|
||||
})
|
||||
|
||||
await assembly.dispose()
|
||||
|
|
@ -740,15 +736,15 @@ describe('Client Typert API', () => {
|
|||
expect(ctx.get('remote.probe')).toBeUndefined()
|
||||
})
|
||||
|
||||
it('rejects weak descriptors and namespace collisions before registration', async () => {
|
||||
it('accepts weak result codecs and rejects namespace collisions before registration', async () => {
|
||||
const ctx = await bench(vi.fn<ConnectionHandle['rpc']['call']>())
|
||||
const weak: InvocationDescriptor = {
|
||||
...directDescriptor(),
|
||||
result: { mode: 'src-json' },
|
||||
}
|
||||
|
||||
await expect(ctx.remote.$mount({ package: '@fixture/weak', descriptors: [weak] }))
|
||||
.rejects.toThrow('has no strict codec')
|
||||
const disposeWeak = await ctx.remote.$mount({ package: '@fixture/weak', descriptors: [weak] })
|
||||
await disposeWeak()
|
||||
await expect(ctx.remote.$mount({
|
||||
package: '@fixture/conflict',
|
||||
descriptors: [{ ...directDescriptor(), namespace: '$mount' }],
|
||||
|
|
|
|||
|
|
@ -735,7 +735,7 @@ describe('TypertGatewayService', () => {
|
|||
expect(service.calls).toEqual([])
|
||||
})
|
||||
|
||||
it('distinguishes strict input and result validation failures', async () => {
|
||||
it('validates strict input without decoding the business result', async () => {
|
||||
const { ctx, service } = await setup()
|
||||
registerStrict(ctx, [strictOnlyDescriptor()])
|
||||
|
||||
|
|
@ -746,27 +746,23 @@ describe('TypertGatewayService', () => {
|
|||
}), 'input-invalid')
|
||||
|
||||
service.nextResult = { title: 1 }
|
||||
await expectCode(ctx.typertGateway.invoke({
|
||||
await expect(ctx.typertGateway.invoke({
|
||||
namespace: 'goals',
|
||||
method: 'strictOnly',
|
||||
args: { request: { title: 'ship' } },
|
||||
}), 'result-invalid')
|
||||
})).resolves.toEqual({ title: 1 })
|
||||
})
|
||||
|
||||
it('rejects non-JSON values after strict codec validation', async () => {
|
||||
it('does not inspect non-JSON business results', async () => {
|
||||
const { ctx, service } = await setup()
|
||||
const descriptor = strictOnlyDescriptor()
|
||||
registerStrict(ctx, [{
|
||||
...descriptor,
|
||||
result: strictCodec('@fixture/gateway#UnknownResult', z.unknown()),
|
||||
}])
|
||||
registerStrict(ctx, [strictOnlyDescriptor()])
|
||||
service.nextResult = 1n
|
||||
|
||||
await expectCode(ctx.typertGateway.invoke({
|
||||
await expect(ctx.typertGateway.invoke({
|
||||
namespace: 'goals',
|
||||
method: 'strictOnly',
|
||||
args: { request: { title: 'ship' } },
|
||||
}), 'result-invalid')
|
||||
})).resolves.toBe(1n)
|
||||
})
|
||||
|
||||
it.each([
|
||||
|
|
@ -801,7 +797,7 @@ describe('TypertGatewayService', () => {
|
|||
expect(service.calls).toContain('passthrough')
|
||||
})
|
||||
|
||||
it('rejects cyclic SRC input and non-JSON SRC results', async () => {
|
||||
it('rejects cyclic SRC input without inspecting SRC results', async () => {
|
||||
const { ctx, service } = await setup()
|
||||
const cyclic: { self?: unknown } = {}
|
||||
cyclic.self = cyclic
|
||||
|
|
@ -811,12 +807,13 @@ describe('TypertGatewayService', () => {
|
|||
args: { value: cyclic },
|
||||
}), 'input-invalid')
|
||||
|
||||
service.nextResult = new Date(0)
|
||||
await expectCode(ctx.typertGateway.invoke({
|
||||
const result = new Date(0)
|
||||
service.nextResult = result
|
||||
await expect(ctx.typertGateway.invoke({
|
||||
namespace: 'goals',
|
||||
method: 'passthrough',
|
||||
args: { value: null },
|
||||
}), 'result-invalid')
|
||||
})).resolves.toBe(result)
|
||||
})
|
||||
|
||||
it('accepts dense JSON and rejects decorated arrays and object properties', async () => {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue