diff --git a/packages/api/gateway/src/client/index.ts b/packages/api/gateway/src/client/index.ts index ee87914633..4fe46e6d2d 100644 --- a/packages/api/gateway/src/client/index.ts +++ b/packages/api/gateway/src/client/index.ts @@ -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`) } diff --git a/packages/api/gateway/src/index.ts b/packages/api/gateway/src/index.ts index 8271445fa8..27306d9af0 100644 --- a/packages/api/gateway/src/index.ts +++ b/packages/api/gateway/src/index.ts @@ -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 { @@ -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> { 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 | AsyncIterable< || typeof Reflect.get(value, Symbol.asyncIterator) === 'function') } -async function *validatedStream( +async function *cancellableStream( source: Iterable | AsyncIterable, - 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 }, ) } diff --git a/packages/api/gateway/src/types.ts b/packages/api/gateway/src/types.ts index 615eaf90ea..b41f35e905 100644 --- a/packages/api/gateway/src/types.ts +++ b/packages/api/gateway/src/types.ts @@ -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 @@ -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> } diff --git a/packages/api/gateway/tests/gateway-stream.host.spec.ts b/packages/api/gateway/tests/gateway-stream.host.spec.ts index fa810f6103..cf95c84f27 100644 --- a/packages/api/gateway/tests/gateway-stream.host.spec.ts +++ b/packages/api/gateway/tests/gateway-stream.host.spec.ts @@ -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' }, }) diff --git a/packages/api/gateway/tests/gateway.client.spec.ts b/packages/api/gateway/tests/gateway.client.spec.ts index 287cff2392..92eeecb5e0 100644 --- a/packages/api/gateway/tests/gateway.client.spec.ts +++ b/packages/api/gateway/tests/gateway.client.spec.ts @@ -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() .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()) 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' }], diff --git a/packages/api/gateway/tests/gateway.host.spec.ts b/packages/api/gateway/tests/gateway.host.spec.ts index 1dd83bdc47..93651f12e8 100644 --- a/packages/api/gateway/tests/gateway.host.spec.ts +++ b/packages/api/gateway/tests/gateway.host.spec.ts @@ -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 () => {