diff --git a/mobile/src/transport/rpc-incompatible-reply-error.ts b/mobile/src/transport/rpc-incompatible-reply-error.ts new file mode 100644 index 00000000000..026732ca3b5 --- /dev/null +++ b/mobile/src/transport/rpc-incompatible-reply-error.ts @@ -0,0 +1,27 @@ +import type { RpcDecodeIssue } from './rpc-operation-contract' + +const INCOMPATIBLE_REPLY_MESSAGE_PREFIX = 'incompatible_reply: ' + +// Why: a reply the operation's reader cannot read says nothing about what the host did. +// On a mutation it is NOT evidence the mutation failed and authorizes no retry — only a +// host-negotiated idempotency capability inside its dedupe window does (see +// tasks/worktree-create-retry.ts). So this error is deliberately neither marked +// delivery-unknown nor shaped like the cutover error the retry loops replay on. +export class RpcIncompatibleReplyError extends Error { + constructor( + readonly operationName: string, + readonly method: string, + readonly issues: readonly RpcDecodeIssue[] + ) { + super(`${INCOMPATIBLE_REPLY_MESSAGE_PREFIX}${operationName} (${method})`) + } +} + +// Why: instanceof can miss across bundle copies, so also match by message, mirroring +// isLogicalClientCutoverError. +export function isRpcIncompatibleReplyError(error: unknown): boolean { + return ( + error instanceof RpcIncompatibleReplyError || + (error instanceof Error && error.message.startsWith(INCOMPATIBLE_REPLY_MESSAGE_PREFIX)) + ) +} diff --git a/mobile/src/transport/rpc-operation-barrier.test.ts b/mobile/src/transport/rpc-operation-barrier.test.ts new file mode 100644 index 00000000000..43c74198b44 --- /dev/null +++ b/mobile/src/transport/rpc-operation-barrier.test.ts @@ -0,0 +1,138 @@ +import { describe, expect, it } from 'vitest' +import type { RpcResponse } from './types' +import { FakeSession } from './mobile-endpoint-supervisor-test-fakes' +import { isRpcDeliveryUnknown, markRpcDeliveryUnknown } from './rpc-delivery-ambiguity' +import { interpretAtRpcBarrier, startRpcOperation } from './rpc-operation' +import { + rpcRefusal, + rpcSuccess, + terminalListAtBarrier, + workspaceListAtBarrier, + worktreePsProbeAtBarrier +} from './rpc-operation-test-families' + +const rows = { worktrees: [{ id: 'w1' }] } + +function deferred(): { promise: Promise; resolve: (value: T) => void } { + let resolve!: (value: T) => void + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise + }) + return { promise, resolve } +} + +function replying(response: RpcResponse): FakeSession { + const session = new FakeSession('connected') + session.sendRequest.mockResolvedValue(response) + return session +} + +function settleAfter(milliseconds: number): Promise<'still waiting'> { + return new Promise((resolve) => setTimeout(() => resolve('still waiting'), milliseconds)) +} + +describe('the post-barrier combinator', () => { + it('starts every request before anything is awaited', () => { + const first = replying(rpcSuccess(rows)) + const second = replying(rpcSuccess({ terminals: [] })) + + startRpcOperation(first, workspaceListAtBarrier, {}) + startRpcOperation(second, terminalListAtBarrier, {}) + + expect(first.sendRequest).toHaveBeenCalledTimes(1) + expect(second.sendRequest).toHaveBeenCalledTimes(1) + }) + + it('yields one verdict per operation, in declared order', async () => { + const verdicts = await interpretAtRpcBarrier([ + startRpcOperation(replying(rpcSuccess(rows)), workspaceListAtBarrier, {}), + startRpcOperation(replying(rpcSuccess({ terminals: [] })), terminalListAtBarrier, {}), + startRpcOperation(replying(rpcSuccess(rows)), worktreePsProbeAtBarrier, {}) + ]) + + expect(verdicts).toEqual([rows, { terminals: [] }, false]) + }) + + // The bug class this exists to remove: whichever peer lost the race used to decide which + // error the user saw. Here the second request fails first in time and the first one refuses + // afterwards, and the declaration still decides. + it('interprets in declared order rather than completion order', async () => { + const lateRefusal = deferred() + const refusing = new FakeSession('connected') + refusing.sendRequest.mockReturnValue(lateRefusal.promise) + const dropped = new FakeSession('connected') + dropped.sendRequest.mockRejectedValue(new Error('socket closed first')) + + const barrier = interpretAtRpcBarrier([ + startRpcOperation(refusing, workspaceListAtBarrier, {}), + startRpcOperation(dropped, terminalListAtBarrier, {}) + ]) + lateRefusal.resolve(rpcRefusal('method_not_found', 'no such method')) + + await expect(barrier).rejects.toThrow('method_not_found: no such method') + }) + + it('keeps the middle operation error when a later one also fails', async () => { + const barrier = interpretAtRpcBarrier([ + startRpcOperation(replying(rpcSuccess(rows)), workspaceListAtBarrier, {}), + startRpcOperation( + replying(rpcRefusal('conflict', 'middle refused')), + workspaceListAtBarrier, + {} + ), + startRpcOperation( + replying(rpcRefusal('runtime_error', 'later refused')), + workspaceListAtBarrier, + {} + ) + ]) + + await expect(barrier).rejects.toThrow('conflict: middle refused') + }) + + it('does not interpret until every raw request has settled', async () => { + const refusing = replying(rpcRefusal('runtime_error', 'boom')) + const pending = deferred() + const stalled = new FakeSession('connected') + stalled.sendRequest.mockReturnValue(pending.promise) + + const barrier = interpretAtRpcBarrier([ + startRpcOperation(refusing, workspaceListAtBarrier, {}), + startRpcOperation(stalled, terminalListAtBarrier, {}) + ]) + const raced = await Promise.race([ + barrier.then( + () => 'resolved' as const, + () => 'rejected' as const + ), + settleAfter(50) + ]) + expect(raced).toBe('still waiting') + + pending.resolve(rpcSuccess({ terminals: [] })) + await expect(barrier).rejects.toThrow('runtime_error: boom') + }) + + it('rethrows a captured transport rejection as the original error object', async () => { + const error = markRpcDeliveryUnknown(new Error('socket closed before response')) + const dropped = new FakeSession('connected') + dropped.sendRequest.mockRejectedValue(error) + + const caught = await interpretAtRpcBarrier([ + startRpcOperation(replying(rpcSuccess(rows)), workspaceListAtBarrier, {}), + startRpcOperation(dropped, terminalListAtBarrier, {}) + ]).catch((thrown: unknown) => thrown) + + expect(caught).toBe(error) + expect(isRpcDeliveryUnknown(caught)).toBe(true) + }) + + it('applies the policy each family declared, at the barrier', async () => { + const verdicts = await interpretAtRpcBarrier([ + startRpcOperation(replying(rpcRefusal('runtime_error')), terminalListAtBarrier, {}), + startRpcOperation(replying(rpcRefusal('method_not_found')), worktreePsProbeAtBarrier, {}) + ]) + + expect(verdicts).toEqual([null, true]) + }) +}) diff --git a/mobile/src/transport/rpc-operation-compile-fence.ts b/mobile/src/transport/rpc-operation-compile-fence.ts new file mode 100644 index 00000000000..2e4449da73c --- /dev/null +++ b/mobile/src/transport/rpc-operation-compile-fence.ts @@ -0,0 +1,131 @@ +import type { RpcClient } from './rpc-client' +import type { RpcMethodName } from './rpc-params-contract' +import { defineRpcOperation, runRpcOperation, startRpcOperation } from './rpc-operation' +import { + workspaceListAtBarrier, + workspaceListOrNull, + workspaceRowsReader, + worktreePsProbe, + type WorkspaceRows +} from './rpc-operation-test-families' +import type { + CapabilityProbeRpcDefinition, + ObjectResultRpcDefinition, + RequireResultRpcDefinition, + RpcAcceptanceName, + RpcCompatibleReader +} from './rpc-operation-contract' + +// Why this file exists: the descriptor's whole point is that a call site cannot pick the +// acceptance policy, the interpretation barrier, or the send-side params for itself. Every +// expect-error directive below is that claim as an assertion — tsc fails on a directive that +// stops catching an error, so `pnpm --dir mobile typecheck` is the gate. Nothing here runs and +// no app code imports it. The FENCE markers are pinned by rpc-operation.test.ts. + +declare const client: RpcClient + +// FENCE: probe-cannot-take-a-reader +export const fenceProbeWithReader: CapabilityProbeRpcDefinition<'worktree.ps', 'on-settle'> = { + name: 'fence.probeWithReader', + method: 'worktree.ps', + acceptance: 'method-not-found-refusal', + barrier: 'on-settle', + consumes: [], + schedules: [], + // @ts-expect-error a refusal-code probe reads no payload, so it cannot carry a reader + read: workspaceRowsReader +} + +// FENCE: decoding-policy-needs-a-reader +// @ts-expect-error 'require-result-or-throw' has no value to return without a reader +export const fenceDecodingWithoutReader: RequireResultRpcDefinition< + 'worktree.ps', + 'rows', + WorkspaceRows, + 'on-settle' +> = { + name: 'fence.decodingWithoutReader', + method: 'worktree.ps', + acceptance: 'require-result-or-throw', + barrier: 'on-settle', + consumes: [], + schedules: [] +} + +// FENCE: acceptance-must-be-a-named-policy +// @ts-expect-error the four policies in rpc-acceptance-policies.ts are the whole vocabulary +export const fenceInventedPolicy: RpcAcceptanceName = 'no-error-means-fine' + +// A reader for a payload no acceptance policy here admits, i.e. one belonging to some other +// family's shape. +const fenceTextReader: RpcCompatibleReader = (raw) => ({ + compatible: true, + variant: 'text', + value: raw, + salvage: { droppedPaths: [], droppedCount: 0 } +}) + +// FENCE: object-policy-reader-sees-an-object +export const fenceObjectPolicyWrongReader: ObjectResultRpcDefinition< + 'worktree.ps', + 'text', + string, + 'on-settle' +> = { + name: 'fence.objectPolicyWrongReader', + method: 'worktree.ps', + acceptance: 'object-result-or-null', + barrier: 'on-settle', + consumes: [], + schedules: [], + // @ts-expect-error the policy admits a non-null object, not the string this reader expects + read: fenceTextReader +} + +// FENCE: define-rejects-a-mismatched-definition +export const fenceDefineRejectsMismatch = defineRpcOperation({ + name: 'fence.defineRejectsMismatch', + method: 'worktree.ps', + // @ts-expect-error no overload of defineRpcOperation pairs a probe with a payload reader + acceptance: 'method-not-found-refusal', + barrier: 'on-settle', + consumes: [], + schedules: [], + // @ts-expect-error ... and the reader it would need is exactly what the probe overload bans + read: workspaceRowsReader +}) + +// FENCE: method-must-exist-in-the-catalog +// @ts-expect-error only generated catalog method names are addressable +export const fenceUnknownMethod: RpcMethodName = 'worktree.nope' + +export async function fenceBarrierAndParams(): Promise { + // FENCE: declared-barrier-cannot-be-moved-earlier + await runRpcOperation( + client, + // @ts-expect-error this family interprets after all requests, so it has no on-settle run + workspaceListAtBarrier, + {} + ) + // FENCE: on-settle-operation-cannot-defer-to-a-barrier + startRpcOperation( + client, + // @ts-expect-error an on-settle family must not be parked behind someone else's barrier + worktreePsProbe, + {} + ) + // FENCE: params-are-fixed-by-the-method + await runRpcOperation( + client, + workspaceListOrNull, + // @ts-expect-error worktree.ps takes a numeric limit + { limit: 'ten' } + ) +} + +export async function fenceVerdictTypes(): Promise { + // FENCE: probe-verdict-is-not-a-decoded-value + // @ts-expect-error the probe's policy yields a boolean, not the other family's rows + const rows: WorkspaceRows = await runRpcOperation(client, worktreePsProbe, {}) + void rows +} diff --git a/mobile/src/transport/rpc-operation-contract.ts b/mobile/src/transport/rpc-operation-contract.ts new file mode 100644 index 00000000000..578b17f755b --- /dev/null +++ b/mobile/src/transport/rpc-operation-contract.ts @@ -0,0 +1,160 @@ +import type { RpcMethodName } from './rpc-params-contract' +import type { RpcFailure, RpcResponse, RpcSuccess } from './types' + +// An operation descriptor fixes the method, the acceptance policy and the interpretation +// barrier at definition time. Per-call freedom over those three is what produced acceptance +// drift and settlement-order drift across mobile's RPC call sites, so none of them is a +// parameter of any send helper. + +/** One of the named policies in rpc-acceptance-policies.ts, chosen per operation family. */ +export type RpcAcceptanceName = + | 'require-result-or-throw' + | 'object-result-or-null' + | 'method-not-found-refusal' + | 'streaming-opener' + +/** Where a settled reply may become a value or a throw. */ +export type RpcInterpretationBarrier = 'on-settle' | 'after-all-requests' + +export type RpcDecodeIssue = { readonly path: string; readonly message: string } + +/** Bounded salvage diagnostics for a reply that decoded with parts dropped. */ +export type RpcSalvageReport = { + readonly droppedPaths: readonly string[] + readonly droppedCount: number +} + +export type RpcReadResult = + | { + readonly compatible: true + readonly variant: Variant + readonly value: Value + readonly salvage: RpcSalvageReport + } + | { readonly compatible: false; readonly issues: readonly RpcDecodeIssue[] } + +/** Reads the payload its acceptance policy admits into one declared semantic variant. */ +export type RpcCompatibleReader = ( + raw: Raw +) => RpcReadResult + +export type RpcStreamOpenerReply = RpcSuccess & { streaming: true } + +// Only a fulfilled outer envelope is classified. Transport rejection stays on the promise +// channel, so an operation in a Promise.all still fails the group immediately instead of +// waiting for a peer and letting a later policy surface a different error. +export type RpcRequestOutcome = + | { + readonly kind: 'outer-refused' + readonly error: RpcFailure['error'] + readonly raw: RpcResponse + } + | { + readonly kind: 'decoded' + readonly variant: Variant + readonly value: Value + readonly raw: RpcResponse + readonly salvage: RpcSalvageReport + } + | { + readonly kind: 'incompatible' + readonly raw: RpcResponse + readonly issues: readonly RpcDecodeIssue[] + } + +export type RpcOperation< + Method extends RpcMethodName, + Acceptance extends RpcAcceptanceName, + Variant extends string, + Value, + Barrier extends RpcInterpretationBarrier +> = { + /** Family name, not the method: two families may share a method with different acceptance. */ + readonly name: string + readonly method: Method + readonly acceptance: Acceptance + readonly barrier: Barrier + /** Reply fields this family reads — the host contract a wire change has to respect. */ + readonly consumes: readonly string[] + /** Refresh/poll schedules this family owns; empty when it only runs on user intent. */ + readonly schedules: readonly string[] + /** Absent for policies that read no result payload. */ + readonly read?: RpcCompatibleReader +} + +// Covariant top type: an operation's reader may narrow the variant and the value. +export type AnyRpcOperation = RpcOperation< + RpcMethodName, + RpcAcceptanceName, + string, + unknown, + RpcInterpretationBarrier +> + +/** The verdict the declared policy yields. Not a per-call choice. */ +export type RpcVerdict< + Acceptance extends RpcAcceptanceName, + Value +> = Acceptance extends 'require-result-or-throw' + ? Value + : Acceptance extends 'object-result-or-null' + ? Value | null + : Acceptance extends 'method-not-found-refusal' + ? boolean + : Acceptance extends 'streaming-opener' + ? RpcStreamOpenerReply | null + : never + +export type RpcOperationSettlement = + | { readonly status: 'fulfilled'; readonly outcome: RpcRequestOutcome } + | { readonly status: 'rejected'; readonly error: unknown } + +type RpcOperationDefinition< + Method extends RpcMethodName, + Barrier extends RpcInterpretationBarrier +> = { + name: string + method: Method + barrier: Barrier + consumes: readonly string[] + schedules: readonly string[] +} + +export type RequireResultRpcDefinition< + Method extends RpcMethodName, + Variant extends string, + Value, + Barrier extends RpcInterpretationBarrier +> = RpcOperationDefinition & { + acceptance: 'require-result-or-throw' + read: RpcCompatibleReader +} + +export type ObjectResultRpcDefinition< + Method extends RpcMethodName, + Variant extends string, + Value, + Barrier extends RpcInterpretationBarrier +> = RpcOperationDefinition & { + acceptance: 'object-result-or-null' + // Raw is the non-null object rpcObjectResultOrNull admits; anything else is incompatible. + read: RpcCompatibleReader, Variant, Value> +} + +export type CapabilityProbeRpcDefinition< + Method extends RpcMethodName, + Barrier extends RpcInterpretationBarrier +> = RpcOperationDefinition & { + acceptance: 'method-not-found-refusal' + /** A probe answers from the refusal code alone, so a reader would have nothing to read. */ + read?: never +} + +export type StreamOpenerRpcDefinition< + Method extends RpcMethodName, + Barrier extends RpcInterpretationBarrier +> = RpcOperationDefinition & { + acceptance: 'streaming-opener' + /** The opener's value is the reply itself; frames arrive on the subscription, not here. */ + read?: never +} diff --git a/mobile/src/transport/rpc-operation-result-reader.test.ts b/mobile/src/transport/rpc-operation-result-reader.test.ts new file mode 100644 index 00000000000..fb51a3d6d34 --- /dev/null +++ b/mobile/src/transport/rpc-operation-result-reader.test.ts @@ -0,0 +1,147 @@ +import { describe, expect, it } from 'vitest' +import { z } from 'zod' +import { salvagingArray } from '../../../src/shared/zod-salvage' +import { FakeSession } from './mobile-endpoint-supervisor-test-fakes' +import { captureRpcOperationSettlement, defineRpcOperation } from './rpc-operation' +import { rpcResultVariant, rpcResultVariants } from './rpc-operation-result-reader' +import { + WORKSPACE_ROWS_SCHEMA, + rpcSuccess, + workspaceRowsReader +} from './rpc-operation-test-families' + +const SALVAGING_ROWS_SCHEMA = z.object({ + worktrees: salvagingArray(z.object({ id: z.string() })) +}) + +const salvagingReader = rpcResultVariant('rows', SALVAGING_ROWS_SCHEMA) + +describe('a single-variant reader', () => { + it('decodes a matching payload and reports nothing dropped', () => { + expect(workspaceRowsReader({ worktrees: [{ id: 'w1' }] })).toEqual({ + compatible: true, + variant: 'rows', + value: { worktrees: [{ id: 'w1' }] }, + salvage: { droppedPaths: [], droppedCount: 0 } + }) + }) + + it('reports dotted issue paths for a payload it cannot read', () => { + const result = workspaceRowsReader({ worktrees: [{ id: 1 }] }) + + expect(result.compatible).toBe(false) + if (result.compatible) { + throw new Error('expected an incompatible read') + } + expect(result.issues).toEqual([{ path: 'worktrees.0.id', message: expect.any(String) }]) + }) + + it('carries the salvage report when the schema drops an element', () => { + const result = salvagingReader({ worktrees: [{ id: 'w1' }, { id: 7 }] }) + + expect(result).toEqual({ + compatible: true, + variant: 'rows', + value: { worktrees: [{ id: 'w1' }] }, + // zod-salvage reports the path relative to the salvaging container, not the envelope. + salvage: { droppedPaths: ['1'], droppedCount: 1 } + }) + }) + + // zod-salvage keeps its collector at module level, so a leak here would blame the next + // reply for the previous one's drops. + it('does not leak drop diagnostics into the next read', () => { + salvagingReader({ worktrees: [{ id: 7 }] }) + + expect(salvagingReader({ worktrees: [{ id: 'w1' }] })).toMatchObject({ + salvage: { droppedPaths: [], droppedCount: 0 } + }) + }) + + it('bounds the issues it reports and says how many it dropped', () => { + const wide = { worktrees: Array.from({ length: 25 }, () => ({ id: 1 })) } + const result = workspaceRowsReader(wide) + + if (result.compatible) { + throw new Error('expected an incompatible read') + } + expect(result.issues).toHaveLength(21) + expect(result.issues[20]).toEqual({ path: '', message: '5 further issues omitted' }) + }) + + // The caveat that comes with zod-salvage: it wraps a synchronous parse only. An async + // schema must read as incompatible rather than leaking a promise into the outcome. + it('reads an async schema as incompatible instead of leaking a promise', () => { + const asyncReader = rpcResultVariant( + 'rows', + z.object({ id: z.string() }).refine(async () => true) + ) + + const result = asyncReader({ id: 'w1' }) + + expect(result.compatible).toBe(false) + if (result.compatible) { + throw new Error('expected an incompatible read') + } + expect(result.issues[0].message).toContain('synchronous parse') + }) +}) + +describe('a multi-variant reader', () => { + const reader = rpcResultVariants<'rows' | 'legacy-array', unknown>([ + workspaceRowsReader, + rpcResultVariant('legacy-array', z.array(z.object({ id: z.string() }))) + ]) + + it('takes the first declared variant that reads', () => { + expect(reader({ worktrees: [{ id: 'w1' }] })).toMatchObject({ variant: 'rows' }) + }) + + it('falls through to a later variant', () => { + expect(reader([{ id: 'w1' }])).toMatchObject({ + variant: 'legacy-array', + value: [{ id: 'w1' }] + }) + }) + + it('tags every variant it tried when none of them reads', () => { + const result = reader('neither shape') + + if (result.compatible) { + throw new Error('expected an incompatible read') + } + expect(result.issues.map((issue) => issue.path)).toEqual(['rows', 'legacy-array']) + }) +}) + +describe('salvage through a descriptor', () => { + const salvagingList = defineRpcOperation({ + name: 'test.salvagingWorkspaceList', + method: 'worktree.ps', + acceptance: 'require-result-or-throw', + barrier: 'on-settle', + consumes: ['worktrees.id'], + schedules: [], + read: salvagingReader + }) + + it('reports the dropped paths on the decoded outcome', async () => { + const session = new FakeSession('connected') + session.sendRequest.mockResolvedValue(rpcSuccess({ worktrees: [{ id: 'w1' }, { id: 7 }] })) + + const settlement = await captureRpcOperationSettlement(session, salvagingList, {}) + + expect(settlement.status === 'fulfilled' && settlement.outcome).toMatchObject({ + kind: 'decoded', + value: { worktrees: [{ id: 'w1' }] }, + salvage: { droppedPaths: ['1'], droppedCount: 1 } + }) + }) +}) + +describe('the shared workspace schema', () => { + it('is the strict shape the salvaging variant relaxes', () => { + expect(WORKSPACE_ROWS_SCHEMA.safeParse({ worktrees: [{ id: 7 }] }).success).toBe(false) + expect(SALVAGING_ROWS_SCHEMA.safeParse({ worktrees: [{ id: 7 }] }).success).toBe(true) + }) +}) diff --git a/mobile/src/transport/rpc-operation-result-reader.ts b/mobile/src/transport/rpc-operation-result-reader.ts new file mode 100644 index 00000000000..0b7e3bef376 --- /dev/null +++ b/mobile/src/transport/rpc-operation-result-reader.ts @@ -0,0 +1,86 @@ +import { z } from 'zod' +import { collectSalvageDrops } from '../../../src/shared/zod-salvage' +import type { RpcCompatibleReader, RpcDecodeIssue } from './rpc-operation-contract' + +const MAX_REPORTED_DECODE_ISSUES = 20 + +/** A reader that names the semantic variant it decodes, so a combinator can tag its issues. */ +export type NamedRpcResultReader = RpcCompatibleReader< + unknown, + Variant, + Value +> & { readonly variant: Variant } + +/** Builds a compatible reader for one semantic variant of a reply payload. */ +export function rpcResultVariant( + variant: Variant, + schema: Schema +): NamedRpcResultReader> { + const read: RpcCompatibleReader> = (raw) => { + try { + // Why: zod-salvage holds module-level collector state and wraps a *synchronous* + // parse only; safeParse throws on an async schema, which reads as incompatible. + const parsed = collectSalvageDrops(() => schema.safeParse(raw)) + if (!parsed.value.success) { + return { compatible: false, issues: decodeIssues(parsed.value.error) } + } + return { + compatible: true, + variant, + value: parsed.value.data as z.output, + salvage: { droppedPaths: parsed.droppedPaths, droppedCount: parsed.droppedCount } + } + } catch (error) { + return { compatible: false, issues: [{ path: '', message: describeThrow(error) }] } + } + } + return Object.assign(read, { variant }) +} + +/** Tries each variant in declared order and takes the first that reads. */ +export function rpcResultVariants( + readers: readonly NamedRpcResultReader[] +): RpcCompatibleReader { + return (raw) => { + const issues: RpcDecodeIssue[] = [] + for (const reader of readers) { + const result = reader(raw) + if (result.compatible) { + return result + } + for (const issue of result.issues) { + issues.push({ path: joinPath(reader.variant, issue.path), message: issue.message }) + } + } + return { compatible: false, issues: boundIssues(issues) } + } +} + +function decodeIssues(error: z.ZodError): RpcDecodeIssue[] { + return boundIssues( + error.issues.map((issue) => ({ + path: issue.path.map((segment) => String(segment)).join('.'), + message: issue.message + })) + ) +} + +// Why: a hostile or very foreign reply can issue per element; report a bounded sample and +// say how many were dropped rather than letting the diagnostic grow with the payload. +function boundIssues(issues: readonly RpcDecodeIssue[]): RpcDecodeIssue[] { + if (issues.length <= MAX_REPORTED_DECODE_ISSUES) { + return [...issues] + } + return [ + ...issues.slice(0, MAX_REPORTED_DECODE_ISSUES), + { path: '', message: `${issues.length - MAX_REPORTED_DECODE_ISSUES} further issues omitted` } + ] +} + +function joinPath(variant: string, path: string): string { + return path ? `${variant}.${path}` : variant +} + +function describeThrow(error: unknown): string { + return error instanceof Error ? error.message : String(error) +} diff --git a/mobile/src/transport/rpc-operation-test-families.ts b/mobile/src/transport/rpc-operation-test-families.ts new file mode 100644 index 00000000000..8adde3565eb --- /dev/null +++ b/mobile/src/transport/rpc-operation-test-families.ts @@ -0,0 +1,113 @@ +import { z } from 'zod' +import type { RpcResponse } from './types' +import type { RpcCompatibleReader } from './rpc-operation-contract' +import { rpcResultVariant, rpcResultVariants } from './rpc-operation-result-reader' +import { defineRpcOperation } from './rpc-operation' + +// Operation families used by the rpc-operation suites and by the compile fence. Kept in one +// place so the tests and the fence assert against the same descriptors, and so tsc sees them +// (the app tsconfig excludes *.test.ts). No app code imports this module. + +export const WORKSPACE_ROWS_SCHEMA = z.object({ + worktrees: z.array(z.object({ id: z.string() })) +}) + +const LEGACY_WORKSPACE_ROWS_SCHEMA = z.array(z.object({ id: z.string() })) + +export type WorkspaceRows = z.output +export type LegacyWorkspaceRows = z.output + +export const workspaceRowsReader = rpcResultVariant('rows', WORKSPACE_ROWS_SCHEMA) + +/** Two semantic variants: the modern envelope, then a host that answered a bare array. */ +export const workspaceRowsOrLegacyReader: RpcCompatibleReader< + unknown, + 'rows' | 'legacy-array', + WorkspaceRows | LegacyWorkspaceRows +> = rpcResultVariants<'rows' | 'legacy-array', WorkspaceRows | LegacyWorkspaceRows>([ + workspaceRowsReader, + rpcResultVariant('legacy-array', LEGACY_WORKSPACE_ROWS_SCHEMA) +]) + +export const workspaceListOrThrow = defineRpcOperation({ + name: 'test.workspaceListOrThrow', + method: 'worktree.ps', + acceptance: 'require-result-or-throw', + barrier: 'on-settle', + consumes: ['worktrees.id'], + schedules: [], + read: workspaceRowsOrLegacyReader +}) + +// Same method, a different family: main's callers disagreed about acceptance, so both rules +// stay named rather than being unified behind one descriptor. +export const workspaceListOrNull = defineRpcOperation({ + name: 'test.workspaceListOrNull', + method: 'worktree.ps', + acceptance: 'object-result-or-null', + barrier: 'on-settle', + consumes: ['worktrees.id'], + schedules: [], + read: workspaceRowsReader +}) + +export const worktreePsProbe = defineRpcOperation({ + name: 'test.worktreePsProbe', + method: 'worktree.ps', + acceptance: 'method-not-found-refusal', + barrier: 'on-settle', + consumes: [], + schedules: [] +}) + +export const terminalStreamOpener = defineRpcOperation({ + name: 'test.terminalStreamOpener', + method: 'terminal.subscribe', + acceptance: 'streaming-opener', + barrier: 'on-settle', + consumes: [], + schedules: [] +}) + +export const workspaceListAtBarrier = defineRpcOperation({ + name: 'test.workspaceListAtBarrier', + method: 'worktree.ps', + acceptance: 'require-result-or-throw', + barrier: 'after-all-requests', + consumes: ['worktrees.id'], + schedules: ['test-poll'], + read: workspaceRowsReader +}) + +export const terminalListAtBarrier = defineRpcOperation({ + name: 'test.terminalListAtBarrier', + method: 'terminal.list', + acceptance: 'object-result-or-null', + barrier: 'after-all-requests', + consumes: ['terminals'], + schedules: [], + read: rpcResultVariant('terminals', z.object({ terminals: z.array(z.unknown()) })) +}) + +export const worktreePsProbeAtBarrier = defineRpcOperation({ + name: 'test.worktreePsProbeAtBarrier', + method: 'worktree.ps', + acceptance: 'method-not-found-refusal', + barrier: 'after-all-requests', + consumes: [], + schedules: [] +}) + +export function rpcSuccess(result: unknown, streaming?: true): RpcResponse { + return { + id: 'rpc-1', + ok: true, + result, + _meta: { runtimeId: 'runtime-1' }, + ...(streaming ? { streaming } : {}) + } +} + +export function rpcRefusal(code: string, message = 'Nope'): RpcResponse { + return { id: 'rpc-1', ok: false, error: { code, message }, _meta: { runtimeId: 'runtime-1' } } +} diff --git a/mobile/src/transport/rpc-operation.test.ts b/mobile/src/transport/rpc-operation.test.ts new file mode 100644 index 00000000000..4b938a83598 --- /dev/null +++ b/mobile/src/transport/rpc-operation.test.ts @@ -0,0 +1,371 @@ +import { readFileSync } from 'node:fs' +import { fileURLToPath } from 'node:url' +import { describe, expect, it } from 'vitest' +import type { RpcResponse } from './types' +import { FakeSession } from './mobile-endpoint-supervisor-test-fakes' +import { createStableLogicalRpcClient } from './stable-logical-rpc-client' +import { isLogicalClientCutoverError } from './stable-logical-rpc-client' +import { isRpcDeliveryUnknown, markRpcDeliveryUnknown } from './rpc-delivery-ambiguity' +import { + RpcIncompatibleReplyError, + isRpcIncompatibleReplyError +} from './rpc-incompatible-reply-error' +import { captureRpcOperationSettlement, runRpcOperation } from './rpc-operation' +import { + rpcRefusal, + rpcSuccess, + terminalListAtBarrier, + terminalStreamOpener, + workspaceListOrNull, + workspaceListOrThrow, + worktreePsProbe +} from './rpc-operation-test-families' + +function connectedSession(response?: RpcResponse): FakeSession { + const session = new FakeSession('connected') + if (response) { + session.sendRequest.mockResolvedValue(response) + } + return session +} + +const rows = { worktrees: [{ id: 'w1' }] } + +describe('request classification', () => { + it('decodes a compatible reply, naming the variant and keeping the raw envelope', async () => { + const response = rpcSuccess(rows) + const settlement = await captureRpcOperationSettlement( + connectedSession(response), + workspaceListOrThrow, + {} + ) + + expect(settlement).toEqual({ + status: 'fulfilled', + outcome: { + kind: 'decoded', + variant: 'rows', + value: rows, + raw: response, + salvage: { droppedPaths: [], droppedCount: 0 } + } + }) + }) + + it('names the legacy variant when the host answered the older shape', async () => { + const settlement = await captureRpcOperationSettlement( + connectedSession(rpcSuccess([{ id: 'w1' }])), + workspaceListOrThrow, + {} + ) + + expect(settlement.status === 'fulfilled' && settlement.outcome.kind).toBe('decoded') + expect(settlement.status === 'fulfilled' && settlement.outcome).toMatchObject({ + variant: 'legacy-array', + value: [{ id: 'w1' }] + }) + }) + + it('classifies a refusal as outer-refused rather than throwing', async () => { + const response = rpcRefusal('runtime_error', 'boom') + const settlement = await captureRpcOperationSettlement( + connectedSession(response), + workspaceListOrThrow, + {} + ) + + expect(settlement).toEqual({ + status: 'fulfilled', + outcome: { + kind: 'outer-refused', + error: { code: 'runtime_error', message: 'boom' }, + raw: response + } + }) + }) + + it('classifies a reply the reader cannot read as incompatible, with bounded issues', async () => { + const response = rpcSuccess({ worktrees: [{ id: 7 }] }) + const settlement = await captureRpcOperationSettlement( + connectedSession(response), + workspaceListOrThrow, + {} + ) + + expect(settlement.status).toBe('fulfilled') + if (settlement.status !== 'fulfilled' || settlement.outcome.kind !== 'incompatible') { + throw new Error('expected an incompatible outcome') + } + expect(settlement.outcome.raw).toBe(response) + expect(settlement.outcome.issues.map((issue) => issue.path)).toContain('rows.worktrees.0.id') + }) + + it('treats a reply that is not an object as incompatible for the nullable family', async () => { + const settlement = await captureRpcOperationSettlement( + connectedSession(rpcSuccess('not an object')), + workspaceListOrNull, + {} + ) + + expect(settlement.status === 'fulfilled' && settlement.outcome.kind).toBe('incompatible') + }) + + it('treats a throwing reader as incompatible, never as a transport failure', async () => { + const exploding = { + ...workspaceListOrThrow, + read: () => { + throw new Error('reader exploded') + } + } + const settlement = await captureRpcOperationSettlement( + connectedSession(rpcSuccess(rows)), + exploding, + {} + ) + + expect(settlement.status === 'fulfilled' && settlement.outcome.kind).toBe('incompatible') + }) +}) + +describe('transport rejection stays on the promise channel', () => { + it('rejects with the original error object', async () => { + const error = new Error('socket closed') + const session = new FakeSession('connected') + session.sendRequest.mockRejectedValue(error) + + await expect(runRpcOperation(session, workspaceListOrThrow, {})).rejects.toBe(error) + }) + + it('keeps a delivery-unknown mark readable through the descriptor', async () => { + const error = markRpcDeliveryUnknown(new Error('socket closed before response')) + const session = new FakeSession('connected') + session.sendRequest.mockRejectedValue(error) + + const caught = await runRpcOperation(session, workspaceListOrThrow, {}).catch( + (thrown: unknown) => thrown + ) + + expect(caught).toBe(error) + expect(isRpcDeliveryUnknown(caught)).toBe(true) + }) + + // The cutover predicate matches class or exact message because instanceof misses across + // bundle copies; a clone from another copy must still read as a cutover through the descriptor. + it('keeps a cutover error from another bundle copy recognisable', async () => { + class ForeignBundleCutoverError extends Error {} + const error = new ForeignBundleCutoverError('RPC interrupted by connection migration') + const session = new FakeSession('connected') + session.sendRequest.mockRejectedValue(error) + + const caught = await runRpcOperation(session, workspaceListOrThrow, {}).catch( + (thrown: unknown) => thrown + ) + + expect(caught).toBe(error) + expect(isLogicalClientCutoverError(caught)).toBe(true) + }) + + it('fails a Promise.all group immediately instead of waiting for a stalled peer', async () => { + const error = new Error('socket closed') + const failing = new FakeSession('connected') + failing.sendRequest.mockRejectedValue(error) + const stalled = new FakeSession('connected') + stalled.sendRequest.mockReturnValue(new Promise(() => {})) + + const group = Promise.all([ + runRpcOperation(failing, workspaceListOrThrow, {}), + runRpcOperation(stalled, workspaceListOrThrow, {}) + ]) + const raced = await Promise.race([ + group.then( + () => 'resolved' as const, + (caught: unknown) => caught + ), + new Promise((resolve) => setTimeout(() => resolve('still waiting'), 50)) + ]) + + expect(raced).toBe(error) + }) + + it('only captures a rejection when a caller names the all-settled helper', async () => { + const error = new Error('socket closed') + const session = new FakeSession('connected') + session.sendRequest.mockRejectedValue(error) + + await expect(captureRpcOperationSettlement(session, workspaceListOrThrow, {})).resolves.toEqual( + { status: 'rejected', error } + ) + }) +}) + +describe('the send path', () => { + it('carries the worktree.ps capability stamp, because it goes through the logical client', async () => { + const session = connectedSession(rpcSuccess(rows)) + const logical = createStableLogicalRpcClient(session, 'lan') + + await runRpcOperation(logical, workspaceListOrThrow, { limit: 500 }) + + expect(session.sendRequest).toHaveBeenCalledWith( + 'worktree.ps', + { limit: 500, supportsWorktreeVisibilitySourceDefaults: true }, + undefined + ) + }) + + it('sends the caller params object untouched for a method with no projection', async () => { + const session = connectedSession(rpcSuccess({ terminals: [] })) + const logical = createStableLogicalRpcClient(session, 'lan') + const params = { worktree: 'w1' } + + await captureRpcOperationSettlement(logical, terminalListAtBarrier, params, { + timeoutMs: 1234 + }) + + expect(session.sendRequest.mock.calls[0][0]).toBe('terminal.list') + expect(session.sendRequest.mock.calls[0][1]).toBe(params) + expect(session.sendRequest.mock.calls[0][2]).toEqual({ timeoutMs: 1234 }) + }) +}) + +describe('acceptance is a property of the family', () => { + const refusal = rpcRefusal('method_not_found', 'no such method') + + it('surfaces the refusal as a coded error for the throwing family', async () => { + await expect( + runRpcOperation(connectedSession(refusal), workspaceListOrThrow, {}) + ).rejects.toThrow('method_not_found: no such method') + }) + + it('answers null to the same refusal for the nullable family', async () => { + await expect( + runRpcOperation(connectedSession(refusal), workspaceListOrNull, {}) + ).resolves.toBeNull() + }) + + it('answers true to the same refusal for the capability probe', async () => { + await expect(runRpcOperation(connectedSession(refusal), worktreePsProbe, {})).resolves.toBe( + true + ) + }) + + it('keeps the probe false for another refusal code and for a success', async () => { + await expect( + runRpcOperation(connectedSession(rpcRefusal('runtime_error')), worktreePsProbe, {}) + ).resolves.toBe(false) + await expect( + runRpcOperation(connectedSession(rpcSuccess(rows)), worktreePsProbe, {}) + ).resolves.toBe(false) + }) + + it('returns the decoded value for the throwing family', async () => { + await expect( + runRpcOperation(connectedSession(rpcSuccess(rows)), workspaceListOrThrow, {}) + ).resolves.toEqual(rows) + }) + + it('returns the reply itself only when it opened a stream', async () => { + const opener = rpcSuccess({ subscriptionId: 's1' }, true) + await expect( + runRpcOperation(connectedSession(opener), terminalStreamOpener, { terminal: 't1' }) + ).resolves.toBe(opener) + await expect( + runRpcOperation( + connectedSession(rpcSuccess({ subscriptionId: 's1' })), + terminalStreamOpener, + { + terminal: 't1' + } + ) + ).resolves.toBeNull() + await expect( + runRpcOperation(connectedSession(rpcRefusal('runtime_error')), terminalStreamOpener, { + terminal: 't1' + }) + ).resolves.toBeNull() + }) +}) + +describe('an incompatible reply', () => { + const incompatible = rpcSuccess({ worktrees: [{ id: 7 }] }) + + it('throws a named incompatible-reply error for the throwing family', async () => { + const caught = await runRpcOperation( + connectedSession(incompatible), + workspaceListOrThrow, + {} + ).catch((thrown: unknown) => thrown) + + expect(caught).toBeInstanceOf(RpcIncompatibleReplyError) + expect(isRpcIncompatibleReplyError(caught)).toBe(true) + expect((caught as RpcIncompatibleReplyError).method).toBe('worktree.ps') + expect((caught as RpcIncompatibleReplyError).operationName).toBe('test.workspaceListOrThrow') + expect((caught as RpcIncompatibleReplyError).issues.length).toBeGreaterThan(0) + }) + + // A reply nobody can read says nothing about what the host did, so it must not look like + // either of the two errors the mutation retry loops replay on. + it('authorizes no retry', async () => { + const caught = await runRpcOperation( + connectedSession(incompatible), + workspaceListOrThrow, + {} + ).catch((thrown: unknown) => thrown) + + expect(isRpcDeliveryUnknown(caught)).toBe(false) + expect(isLogicalClientCutoverError(caught)).toBe(false) + }) + + it('answers null for the nullable family', async () => { + await expect( + runRpcOperation(connectedSession(incompatible), workspaceListOrNull, {}) + ).resolves.toBeNull() + }) +}) + +describe('a descriptor', () => { + it('cannot have its policy or barrier swapped at runtime', () => { + expect(Object.isFrozen(workspaceListOrThrow)).toBe(true) + expect(() => { + ;(workspaceListOrThrow as { acceptance: string }).acceptance = 'object-result-or-null' + }).toThrow(TypeError) + expect(() => { + ;(workspaceListOrThrow as { barrier: string }).barrier = 'after-all-requests' + }).toThrow(TypeError) + }) + + it('publishes the fields it consumes and the schedules it owns', () => { + expect(workspaceListOrThrow.consumes).toEqual(['worktrees.id']) + expect(workspaceListOrThrow.schedules).toEqual([]) + expect(Object.isFrozen(workspaceListOrThrow.consumes)).toBe(true) + }) +}) + +// The fence itself is checked by `tsc --noEmit` (an expect-error directive that stops catching +// an error fails the typecheck). This pins its coverage so the file cannot be quietly gutted. +describe('the compile fence', () => { + const fenceSource = readFileSync( + fileURLToPath(new URL('./rpc-operation-compile-fence.ts', import.meta.url)), + 'utf8' + ) + const expectErrorDirective = `@ts-${'expect-error'}` + + it('still asserts every rejection it is meant to', () => { + const markers = [...fenceSource.matchAll(/\/\/ FENCE: (?[a-z-]+)/g)].map( + (match) => match.groups?.name + ) + + expect(markers).toEqual([ + 'probe-cannot-take-a-reader', + 'decoding-policy-needs-a-reader', + 'acceptance-must-be-a-named-policy', + 'object-policy-reader-sees-an-object', + 'define-rejects-a-mismatched-definition', + 'method-must-exist-in-the-catalog', + 'declared-barrier-cannot-be-moved-earlier', + 'on-settle-operation-cannot-defer-to-a-barrier', + 'params-are-fixed-by-the-method', + 'probe-verdict-is-not-a-decoded-value' + ]) + expect(fenceSource.split(expectErrorDirective).length - 1).toBe(11) + }) +}) diff --git a/mobile/src/transport/rpc-operation.ts b/mobile/src/transport/rpc-operation.ts new file mode 100644 index 00000000000..5d2660b7b72 --- /dev/null +++ b/mobile/src/transport/rpc-operation.ts @@ -0,0 +1,284 @@ +import type { RpcClient, SendRequestOptions } from './rpc-client' +import type { RpcMethodName, RpcParams } from './rpc-params-contract' +import type { RpcResponse } from './types' +import { + isMethodNotFoundRefusal, + isStreamingOpenerReply, + requireRpcResultOrThrowCodedError, + rpcObjectResultOrNull +} from './rpc-acceptance-policies' +import { RpcIncompatibleReplyError } from './rpc-incompatible-reply-error' +import type { + AnyRpcOperation, + CapabilityProbeRpcDefinition, + ObjectResultRpcDefinition, + RpcAcceptanceName, + RpcCompatibleReader, + RpcDecodeIssue, + RpcInterpretationBarrier, + RpcOperation, + RpcOperationSettlement, + RpcRequestOutcome, + RpcSalvageReport, + RequireResultRpcDefinition, + StreamOpenerRpcDefinition, + RpcVerdict +} from './rpc-operation-contract' + +const NOTHING_SALVAGED: RpcSalvageReport = { droppedPaths: [], droppedCount: 0 } + +type RpcOperationDefinitionInput = + | RequireResultRpcDefinition + | ObjectResultRpcDefinition + | CapabilityProbeRpcDefinition + | StreamOpenerRpcDefinition + +export function defineRpcOperation< + Method extends RpcMethodName, + Variant extends string, + Value, + Barrier extends RpcInterpretationBarrier +>( + definition: RequireResultRpcDefinition +): RpcOperation +export function defineRpcOperation< + Method extends RpcMethodName, + Variant extends string, + Value, + Barrier extends RpcInterpretationBarrier +>( + definition: ObjectResultRpcDefinition +): RpcOperation +export function defineRpcOperation< + Method extends RpcMethodName, + Barrier extends RpcInterpretationBarrier +>( + definition: CapabilityProbeRpcDefinition +): RpcOperation +export function defineRpcOperation< + Method extends RpcMethodName, + Barrier extends RpcInterpretationBarrier +>( + definition: StreamOpenerRpcDefinition +): RpcOperation +export function defineRpcOperation(definition: RpcOperationDefinitionInput): AnyRpcOperation { + // Why: frozen so no call site can swap the policy or the barrier on a shared descriptor. + return Object.freeze({ + name: definition.name, + method: definition.method, + acceptance: definition.acceptance, + barrier: definition.barrier, + consumes: Object.freeze([...definition.consumes]), + schedules: Object.freeze([...definition.schedules]), + // Why: classifyReply only ever hands a reader the payload its own policy admitted, so + // the object policy's narrower parameter is sound to store as unknown. + read: definition.read as RpcCompatibleReader | undefined + }) +} + +/** Sends the operation without interpreting it; transport rejection stays on the promise. */ +async function request( + client: RpcClient, + operation: AnyRpcOperation, + params: unknown, + options?: SendRequestOptions +): Promise> { + // Why: no try/catch here. A transport failure must reach the caller as the original error + // object — isLogicalClientCutoverError and isRpcDeliveryUnknown both die on a wrapper — + // and an always-settled send would make Promise.all wait for a peer where today the group + // fails immediately, letting a later policy surface a different error. + const response = await client.sendRequest(operation.method, params, options) + return classifyReply(operation, response) +} + +type AdmittedPayload = + | { readonly admitted: true; readonly value: unknown } + | { readonly admitted: false; readonly issues: readonly RpcDecodeIssue[] } + +// The payload the operation's own acceptance policy admits from a fulfilled success. +function admitPayload(operation: AnyRpcOperation, response: RpcResponse): AdmittedPayload { + switch (operation.acceptance) { + case 'object-result-or-null': { + const object = rpcObjectResultOrNull(response) + return object === null + ? { admitted: false, issues: [{ path: 'result', message: 'not a non-null object' }] } + : { admitted: true, value: object } + } + case 'streaming-opener': + return isStreamingOpenerReply(response) + ? { admitted: true, value: response } + : { admitted: false, issues: [{ path: 'streaming', message: 'reply opened no stream' }] } + default: + // Reuses the policy rather than reading `.result` again; a success never throws here. + return { admitted: true, value: requireRpcResultOrThrowCodedError(response) } + } +} + +const READERLESS_VARIANTS: Record = { + 'method-not-found-refusal': 'accepted', + 'streaming-opener': 'stream-opened' +} + +function classifyReply( + operation: AnyRpcOperation, + response: RpcResponse +): RpcRequestOutcome { + if (!response.ok) { + return { kind: 'outer-refused', error: response.error, raw: response } + } + const payload = admitPayload(operation, response) + if (!payload.admitted) { + return { kind: 'incompatible', raw: response, issues: payload.issues } + } + const read = operation.read + if (!read) { + return { + kind: 'decoded', + variant: READERLESS_VARIANTS[operation.acceptance] ?? 'accepted', + value: payload.value, + raw: response, + salvage: NOTHING_SALVAGED + } + } + let result: ReturnType + try { + result = read(payload.value) + } catch (error) { + // A reader that throws is an incompatible reply, never a transport failure. + return { + kind: 'incompatible', + raw: response, + issues: [{ path: '', message: error instanceof Error ? error.message : String(error) }] + } + } + if (!result.compatible) { + return { kind: 'incompatible', raw: response, issues: result.issues } + } + return { + kind: 'decoded', + variant: result.variant, + value: result.value, + raw: response, + salvage: result.salvage + } +} + +// Applies the operation's declared acceptance policy. Private on purpose: there is no +// free-standing callOrThrow, so no call site can pick a different rule for the same reply. +function interpret( + operation: AnyRpcOperation, + settled: RpcRequestOutcome +): unknown { + const acceptance: RpcAcceptanceName = operation.acceptance + switch (acceptance) { + case 'require-result-or-throw': + if (settled.kind === 'outer-refused') { + // Reuses the policy so the thrown `code: message` text cannot drift from main's. + return requireRpcResultOrThrowCodedError(settled.raw) + } + if (settled.kind === 'incompatible') { + throw new RpcIncompatibleReplyError(operation.name, operation.method, settled.issues) + } + return settled.value + case 'object-result-or-null': + return settled.kind === 'decoded' ? settled.value : null + case 'method-not-found-refusal': + return settled.kind === 'outer-refused' ? isMethodNotFoundRefusal(settled.raw) : false + case 'streaming-opener': + return settled.kind === 'decoded' && isStreamingOpenerReply(settled.raw) ? settled.raw : null + } +} + +function interpretSettlement( + operation: AnyRpcOperation, + settlement: RpcOperationSettlement +): unknown { + if (settlement.status === 'rejected') { + // Why: rethrow the original object — isRpcDeliveryUnknown is a WeakSet on identity and + // isLogicalClientCutoverError matches class or exact message; a wrapper loses both. + throw settlement.error + } + return interpret(operation, settlement.outcome) +} + +/** Sends and interprets at the operation's own barrier. Only for barrier 'on-settle'. */ +export async function runRpcOperation< + Method extends RpcMethodName, + Acceptance extends RpcAcceptanceName, + Variant extends string, + Value +>( + client: RpcClient, + operation: RpcOperation, + params: RpcParams, + options?: SendRequestOptions +): Promise> { + const outcome = await request(client, operation, params, options) + return interpret(operation, outcome) as RpcVerdict +} + +/** The named opt-in to all-settled semantics. Yields an outcome, never a verdict: the + * verdict still comes only from the declared policy, at the declared barrier. */ +export async function captureRpcOperationSettlement< + Method extends RpcMethodName, + Acceptance extends RpcAcceptanceName, + Variant extends string, + Value, + Barrier extends RpcInterpretationBarrier +>( + client: RpcClient, + operation: RpcOperation, + params: RpcParams, + options?: SendRequestOptions +): Promise> { + try { + const outcome = await request(client, operation, params, options) + return { status: 'fulfilled', outcome: outcome as RpcRequestOutcome } + } catch (error) { + return { status: 'rejected', error } + } +} + +export type PendingRpcOperation = { + readonly operation: Op + readonly settlement: Promise> +} + +/** Starts a request whose interpretation is deferred to the barrier it declared. */ +export function startRpcOperation< + Method extends RpcMethodName, + Acceptance extends RpcAcceptanceName, + Variant extends string, + Value +>( + client: RpcClient, + operation: RpcOperation, + params: RpcParams, + options?: SendRequestOptions +): PendingRpcOperation> { + return { + operation, + settlement: captureRpcOperationSettlement(client, operation, params, options) + } +} + +type RpcBarrierVerdicts[]> = { + [Index in keyof Pending]: Pending[Index] extends PendingRpcOperation< + RpcOperation + > + ? RpcVerdict + : never +} + +/** Awaits every raw request, then interprets in declared order. */ +export async function interpretAtRpcBarrier< + Pending extends readonly PendingRpcOperation[] +>(pending: Pending): Promise> { + // Why: interpreting as each request lands would let whichever peer failed first decide the + // error the user sees and how long the screen spins. Declared order makes that a property + // of the definition instead of a race. + const settlements = await Promise.all(pending.map((entry) => entry.settlement)) + return pending.map((entry, index) => + interpretSettlement(entry.operation, settlements[index]) + ) as RpcBarrierVerdicts +}