diff --git a/.changeset/transparent-payment-durability.md b/.changeset/transparent-payment-durability.md new file mode 100644 index 0000000..4ee9827 --- /dev/null +++ b/.changeset/transparent-payment-durability.md @@ -0,0 +1,17 @@ +--- +'@contextvm/sdk': patch +--- + +Transparent payment robustness (CEP-8): prevent double charges and lost paid results. + +- **No double charge on redelivery:** the pending-payment entry is now retained until TTL once an invoice has been issued, instead of being deleted on every failure. Previously a `verifyPayment` timeout (or any post-invoice failure) disarmed the dedup, so a spec-blessed client retry of the same request event minted a second invoice that the client auto-paid — violating CEP-8's "MUST NOT charge more than once for the same transparent request event". Pre-invoice failures still delete the entry so retries are free. +- **No paid-but-undelivered:** a failed `payment_accepted` publish (relay error, or the idle paying client's session LRU-evicted mid-payment) no longer aborts the forward. The acceptance notification is a SHOULD; the capability result is the point. +- **Pending capacity fails closed:** at `maxPendingPayments` capacity the middleware purges expired entries and then refuses new priced requests (best-effort `payment_rejected`, no invoice minted) instead of silently evicting a live payment's dedup entry. `maxPendingPayments: 0` now refuses all priced requests. +- **Route survival for paid requests:** the correlation route (and session) is snapshotted when an invoice is issued; the response router falls back to that snapshot on route miss, so paid results still reach the client after duplicate-delivery cleanup popped the route or the paying client's session was evicted. The dead "recreate session if active routes" guard in the eviction handler was removed (snapshots make it unnecessary, and re-inserting into the session LRU from inside its own eviction callback would corrupt capacity accounting). +- **No leaked open-stream writers on dropped requests:** a request dropped by middleware (payment gating, policy) or failed by a throwing middleware chain never produces the normal-path response that would reap its writer, so its `OpenStreamWriter` reservation leaked until transport teardown. Drop cleanup now releases the writer alongside the correlation route. +- **Non-positive TTLs fall back to the default:** `paymentTtlMs: 0` (both lifecycles) and invoice `ttl: 0` (explicit gating) previously birth-expired the pending/grant entry while the invoice stayed payable, disarming the redelivery dedup. They now fall back to the default window, matching the guard `getVerificationTimeoutMs` already applied. +- **Double registration refuses:** calling `withServerPayments` twice on the same transport used to silently register a second middleware pair with its own dedup closures, minting two invoices and double-charging every priced request. It now throws. +- **Client handler failures resolve the request:** a throwing in-band payment handler previously surfaced only on `onerror` — the pending MCP request hung until the server TTL, while policy/handler _declines_ resolved it with a synthesized error. Handler crashes now synthesize the same decline error (and still surface on `onerror`). +- **Verify-task crash cleanup:** the detached explicit-gating verification task's outer catch now clears the pending identity (belt-and-suspenders; idempotent), so a store failure inside the inner catch cannot pin the identity as pending for the whole TTL. +- **Gating errors carry the client's request id:** `-32042`/`-32043` responses previously leaked the server's internal routing key (the Nostr event id) into the JSON-RPC `id` field, because the targeted-response exit path skipped the id restore the normal response path performs. The restore now runs on both exit paths, so the wire conforms to JSON-RPC 2.0 and CEP-8's examples. SDK clients are unaffected either way (the client transport restores ids independently). +- **Client: double-wrap guard, retry floor, bounded retry counters, type-aware cache keys** (reported by the rs-sdk maintainers): `withClientPayments` now refuses a second wrap of the same transport (a chained double-wrap double-paid every offer through two independent pipelines with separate dedup sets); `-32043` retries are floored at `minRetryDelayMs` (default 1 s — `retry_after: 0` could otherwise re-send a byte-identical Nostr event, which relays and servers swallow as a duplicate, silently losing the retry); retry counters are bounded by the same LRU capacity as the raw-request cache; and the raw-request cache plus retry counters key by `JSON.stringify(id)` so numeric `5` and string `"5"` no longer collide (a collision retried the wrong request's payload). diff --git a/src/__mocks__/mock-relay-server.ts b/src/__mocks__/mock-relay-server.ts index 5b10ec9..c6286b5 100644 --- a/src/__mocks__/mock-relay-server.ts +++ b/src/__mocks__/mock-relay-server.ts @@ -45,6 +45,7 @@ type ConnectionInstance = { cleanup: () => void; cleanupWithoutClosingSocket: () => void; closeSocket: (code: number, reason: string) => void; + terminateSocket: () => void; handle: (message: string) => void; send: (message: NostrRelayMessage) => void; }; @@ -102,6 +103,19 @@ export function startMockRelay( } } + terminateSocket(): void { + // Simulate the relay process dying: hard reset with no close frame, so + // clients observe an abnormal closure (1006 / wasClean=false) — the + // signal reconnect logic keys on. A frame-based close(1011) reads as + // wasClean=true on bun >= 1.4.0, which stopped clients from + // reconnecting after a simulated relay outage. + try { + this.socket.terminate(); + } catch { + this.closeSocket(1011, 'Relay paused'); + } + } + cleanup(): void { // Used by stop()/pause() to actively take the relay offline. // Close the socket first, then drop all subscriptions. @@ -298,9 +312,13 @@ export function startMockRelay( relayUrl: `ws://127.0.0.1:${port}`, httpUrl: `http://127.0.0.1:${port}`, stop: () => { - // Close any existing WebSocket connections before stopping the server. + // Drop connections *uncleanly* (hard reset) before stopping the server: + // tests use stop()+restart-on-same-port to simulate outages/partitions, + // and clients only reconnect on abnormal closures. A graceful 1001 reads + // as a clean close on bun >= 1.4.0 (no remap), which stopped pools from + // reconnecting after recovery. for (const instance of state.connections.values()) { - instance.closeSocket(1001, 'Relay stopping'); + instance.terminateSocket(); instance.cleanupWithoutClosingSocket(); } state.connections.clear(); @@ -310,8 +328,9 @@ export function startMockRelay( pause: () => { runtime.acceptingWs = false; for (const instance of state.connections.values()) { - // Close *uncleanly* so clients reconnect (applesauce-relay only retries on !wasClean). - instance.closeSocket(1011, 'Relay paused'); + // Drop *uncleanly* so clients reconnect (applesauce-relay only retries + // on !wasClean). + instance.terminateSocket(); instance.cleanupWithoutClosingSocket(); } state.connections.clear(); diff --git a/src/payments/authorization-store.ts b/src/payments/authorization-store.ts index 28c1939..7899276 100644 --- a/src/payments/authorization-store.ts +++ b/src/payments/authorization-store.ts @@ -17,6 +17,11 @@ interface PaidAuthorization { * meaning it is strictly single-process. For multi-process horizontal scaling, * implementers should use a distributed lock (e.g. Redis Redlock) keyed by * the canonical invocation identity to prevent duplicate payments. + * + * NOTE: Atomic reserve-then-dispatch (`claim()` here, `trySetPending()` in the + * gating middlewares) relies on JavaScript run-to-completion: there is no + * interleaving point between the two calls. Porters to await-capable runtimes + * must compose them into a single critical section. */ export class AuthorizationStore { private readonly authorizations: LruCache; diff --git a/src/payments/client-payments.test.ts b/src/payments/client-payments.test.ts index 4d1fb78..fb039c8 100644 --- a/src/payments/client-payments.test.ts +++ b/src/payments/client-payments.test.ts @@ -1,21 +1,13 @@ import { describe, expect, test } from 'bun:test'; -import type { Transport } from '@contextvm/mcp-sdk/shared/transport'; import type { JSONRPCMessage } from '@contextvm/mcp-sdk/types.js'; import { withClientPayments } from './client-payments.js'; import type { PaymentHandlerRequest } from './types.js'; import { NostrClientTransport } from '../transport/nostr-client-transport.js'; +import type { TransportWithContext } from '../transport/nostr-client-transport.js'; import { PrivateKeySigner } from '../signer/private-key-signer.js'; import { EncryptionMode } from '../core/interfaces.js'; import { MockRelayHub } from '../__mocks__/mock-relay-handler.js'; -/** Minimal fake transport that exposes onmessageWithContext for unit tests. */ -type TransportWithContext = Transport & { - onmessageWithContext?: ( - message: JSONRPCMessage, - ctx: { eventId: string; correlatedEventId?: string }, - ) => void; -}; - const createMockNostrTransport = (): NostrClientTransport => { const hub = new MockRelayHub(); return new NostrClientTransport({ @@ -267,6 +259,74 @@ describe('withClientPayments()', () => { expect(handleCalls).toBe(0); }); + test('synthesizes JSON-RPC error when the payment handler throws, instead of hanging the request', async () => { + const transport = createMockNostrTransport(); + + const observed: JSONRPCMessage[] = []; + const errors: Error[] = []; + + const paid = withClientPayments(transport, { + handlers: [ + { + pmi: 'fake', + async handle(): Promise { + throw new Error('wallet exploded'); + }, + }, + ], + }); + + paid.onmessage = (msg) => observed.push(msg); + paid.onerror = (err) => errors.push(err); + + await paid.start(); + + ( + transport as unknown as { + correlationStore: { + registerRequest: (eventId: string, req: unknown) => void; + }; + } + ).correlationStore.registerRequest('req-event-id', { + originalRequestId: 42, + isInitialize: false, + progressToken: undefined, + originalRequestContext: { method: 'tools/call', capability: 'tool:add' }, + }); + + const paymentRequired: JSONRPCMessage = { + jsonrpc: '2.0', + method: 'notifications/payment_required', + params: { amount: 1, pay_req: 'x', pmi: 'fake' }, + }; + + (transport as unknown as TransportWithContext).onmessageWithContext?.( + paymentRequired, + { + eventId: 'evt', + correlatedEventId: 'req-event-id', + }, + ); + + await new Promise((r) => setTimeout(r, 0)); + + const errResp = observed.find( + ( + m, + ): m is { + jsonrpc: '2.0'; + id: number; + error: { code: number; message: string; data?: unknown }; + } => 'id' in m && m.id === 42 && 'error' in m, + ); + expect(errResp?.error?.code).toBe(-32000); + expect(errResp?.error?.message).toBe( + 'Payment handler failed: wallet exploded', + ); + // The crash still surfaces on onerror. + expect(errors[0]?.message).toMatch(/wallet exploded/); + }); + test('synthesizes JSON-RPC error when canHandle declines and correlation exists', async () => { const transport = createMockNostrTransport(); @@ -923,6 +983,7 @@ describe('withClientPayments()', () => { const paid = withClientPayments(transport, { handlers: [{ pmi: 'fake', async handle(): Promise {} }], paymentInteraction: 'explicit_gating', + minRetryDelayMs: 1, }); paid.onmessage = (msg) => observed.push(msg); await paid.start(); @@ -1058,6 +1119,7 @@ describe('withClientPayments()', () => { handlers: [{ pmi: 'fake', async handle(): Promise {} }], paymentInteraction: 'explicit_gating', maxPendingRetries: 2, + minRetryDelayMs: 1, }); paid.onmessage = (msg) => observed.push(msg); await paid.start(); @@ -1107,4 +1169,130 @@ describe('withClientPayments()', () => { await paid.close(); }); + + test('refuses wrapping a transport that already has client payments', () => { + const transport = createMockNostrTransport(); + const paid = withClientPayments(transport, { handlers: [] }); + + // A chained second wrap double-pays every offer; a sibling wrap silently + // kills the first wrapper's trampolines. Both must fail fast. + expect(() => withClientPayments(paid, { handlers: [] })).toThrow( + /already called on this transport/, + ); + expect(() => withClientPayments(transport, { handlers: [] })).toThrow( + /already called on this transport/, + ); + }); + + test('floors -32043 retry delay so retry_after: 0 cannot re-send within the same second', async () => { + const transport = createMockNostrTransport(); + transport + .getInternalStateForTesting() + .correlationStore.registerRequest('req-event-id-floor', { + originalRequestId: 77, + isInitialize: false, + originalRequestContext: { method: 'tools/call' }, + }); + + let sentMessage: JSONRPCMessage | undefined; + transport.send = async (msg) => { + sentMessage = msg; + }; + + const paid = withClientPayments(transport, { + handlers: [{ pmi: 'fake', async handle(): Promise {} }], + minRetryDelayMs: 50, + }); + paid.onmessage = (): void => {}; + await paid.start(); + + await paid.send({ + jsonrpc: '2.0', + id: 77, + method: 'tools/call', + params: { name: 'floored' }, + }); + sentMessage = undefined; // Reset so only the retry is observed + (transport as unknown as TransportWithContext).onmessageWithContext?.( + { + jsonrpc: '2.0', + id: 77, + error: { + code: -32043, + message: 'Payment Pending', + data: { retry_after: 0 }, + }, + }, + { eventId: 'evt-floor', correlatedEventId: 'req-event-id-floor' }, + ); + + // Before the floor elapses: no retry. + await new Promise((r) => setTimeout(r, 10)); + expect(sentMessage).toBeUndefined(); + + // After the floor: the original request is retried. + await new Promise((r) => setTimeout(r, 120)); + expect(sentMessage as unknown).toEqual({ + jsonrpc: '2.0', + id: 77, + method: 'tools/call', + params: { name: 'floored' }, + }); + + await paid.close(); + }); + + test('keeps numeric and string request ids with the same text form distinct', async () => { + const sent: JSONRPCMessage[] = []; + const baseTransport: TransportWithContext = { + onmessage: undefined, + onmessageWithContext: undefined, + onerror: undefined, + onclose: undefined, + async start(): Promise {}, + async send(message: JSONRPCMessage): Promise { + sent.push(message); + }, + async close(): Promise {}, + }; + + const paid = withClientPayments(baseTransport, { minRetryDelayMs: 1 }); + await paid.start(); + + await paid.send({ + jsonrpc: '2.0', + id: 5, + method: 'tools/call', + params: { name: 'numeric' }, + }); + await paid.send({ + jsonrpc: '2.0', + id: '5', + method: 'tools/call', + params: { name: 'string' }, + }); + sent.length = 0; + + // -32043 answering the NUMERIC id 5 must retry the numeric request — + // with String(id) keys the later string write overwrote it. + baseTransport.onmessageWithContext?.( + { + jsonrpc: '2.0', + id: 5, + error: { + code: -32043, + message: 'Payment Pending', + data: { retry_after: 0 }, + }, + }, + { eventId: 'evt', correlatedEventId: 'req' }, + ); + + await new Promise((r) => setTimeout(r, 30)); + + expect(sent).toHaveLength(1); + expect((sent[0] as { params?: { name?: string } }).params?.name).toBe( + 'numeric', + ); + }); }); diff --git a/src/payments/client-payments.ts b/src/payments/client-payments.ts index 49c2949..916a305 100644 --- a/src/payments/client-payments.ts +++ b/src/payments/client-payments.ts @@ -9,6 +9,7 @@ import { type JSONRPCErrorResponse, } from '@contextvm/mcp-sdk/types.js'; import { NostrClientTransport } from '../transport/nostr-client-transport.js'; +import type { TransportWithContext } from '../transport/nostr-client-transport.js'; import type { PaymentHandler, PaymentRejectedNotification, @@ -66,6 +67,17 @@ export interface ClientPaymentsOptions { * @default DEFAULT_PAYMENT_TTL_MS (300_000 ms) */ defaultPaymentTtlMs?: number; + /** + * Minimum delay before a `-32043` Payment Pending retry is re-sent + * (milliseconds), applied after the server-provided `retry_after` backoff. + * + * Guards against `retry_after: 0`: an immediate retry within the same second + * can produce a byte-identical Nostr event (same content, same tags, same + * `created_at` at second resolution), which relays and servers swallow as a + * duplicate — silently losing the retry. + * @default 1000 + */ + minRetryDelayMs?: number; /** * Optional policy hook invoked when a `payment_required` notification is received. @@ -130,13 +142,6 @@ type SyntheticProgressEntry = { wireProgressToken: string | number; }; -type TransportWithContext = Transport & { - onmessageWithContext?: ( - message: JSONRPCMessage, - ctx: { eventId: string; correlatedEventId?: string }, - ) => void; -}; - function supportsOnmessageWithContext( transport: Transport, ): transport is TransportWithContext { @@ -182,6 +187,9 @@ function isPaymentRequiredNotification( return isJSONRPCNotification(msg) && msg.method === PAYMENT_REQUIRED_METHOD; } +/** Transports already wrapped by {@link withClientPayments}. */ +const clientPaymentsWrapped = new WeakSet(); + /** * Wraps a transport to automatically handle CEP-8 payment requests. * @@ -197,6 +205,12 @@ export function withClientPayments( transport: Transport, options: ClientPaymentsOptions, ): Transport { + // A second wrap silently double-pays every offer (each wrapper runs its own + // pipeline with its own dedup on separate closures) or, for sibling wraps, + // silently kills the first wrapper's trampolines. Refuse outright. + if (clientPaymentsWrapped.has(transport)) { + throw new Error('withClientPayments already called on this transport'); + } const logger = createLogger('client-payments'); const syntheticProgressIntervalMs = @@ -206,6 +220,8 @@ export function withClientPayments( const defaultPaymentTtlMs = options.defaultPaymentTtlMs ?? DEFAULT_PAYMENT_TTL_MS; + const minRetryDelayMs = options.minRetryDelayMs ?? 1000; + const syntheticProgress = new Map(); let syntheticProgressScheduler: ReturnType | undefined; @@ -257,7 +273,10 @@ export function withClientPayments( }; const pendingTimers = new Set>(); - const retryCounts = new Map(); + // Bounded together with rawRequestCache: an abandoned request's counter must + // not accumulate forever (terminal responses and close() clear these, but a + // request that never terminates otherwise would leak). + const retryCounts = new LruCache(1000); const rawRequestCache = new LruCache(1000); const MAX_RETRIES = options.maxPendingRetries ?? 10; @@ -379,7 +398,9 @@ export function withClientPayments( const requestId = errorMsg.id; const rawRequest = - requestId != null ? rawRequestCache.get(String(requestId)) : undefined; + requestId != null + ? rawRequestCache.get(JSON.stringify(requestId)) + : undefined; if (!rawRequest) { logger.warn( 'missing raw original request, cannot retry explicit payment', @@ -446,7 +467,9 @@ export function withClientPayments( const requestId = errorMsg.id; const rawRequest = - requestId != null ? rawRequestCache.get(String(requestId)) : undefined; + requestId != null + ? rawRequestCache.get(JSON.stringify(requestId)) + : undefined; if (!rawRequest) { logger.warn( 'missing raw original request, cannot retry explicit payment pending', @@ -456,7 +479,9 @@ export function withClientPayments( return; } - const requestIdKey = errorMsg.id as string | number; + // Type-aware key: JSON.stringify keeps numeric 5 and string "5" + // distinct, and must be used consistently for both maps below. + const requestIdKey = JSON.stringify(errorMsg.id); const retries = retryCounts.get(requestIdKey) ?? 0; if (retries >= MAX_RETRIES) { logger.error('max explicit payment retries exceeded', { @@ -478,7 +503,13 @@ export function withClientPayments( const baseDelayMs = (retryAfterSeconds ?? 1) * 1000; const exponentialMultiplier = Math.pow(1.5, retries); - const delayMs = Math.min(baseDelayMs * exponentialMultiplier, 10000); + // Floor at minRetryDelayMs: a sub-second retry can produce a + // byte-identical Nostr event (created_at has second resolution), which + // relays and servers swallow as a duplicate. + const delayMs = Math.min( + Math.max(baseDelayMs * exponentialMultiplier, minRetryDelayMs), + 10000, + ); const timer = setTimeout(() => { pendingTimers.delete(timer); @@ -618,26 +649,27 @@ export function withClientPayments( } inFlightPayReqs.add(message.params.pay_req); - try { - const req: PaymentHandlerRequest = { - amount: message.params.amount, - pay_req: message.params.pay_req, - pmi: message.params.pmi, - description: message.params.description, - ttl: message.params.ttl, - _meta: message.params._meta, - requestEventId, - }; - const synthesizeClientDeclineError = (params: { - message: string; - }): void => { - if (pending?.progressToken) { - stopSyntheticProgress(pending.progressToken); - } - synthesizePaymentDecline(pending, params.message, req.pmi, req.amount); - }; + const req: PaymentHandlerRequest = { + amount: message.params.amount, + pay_req: message.params.pay_req, + pmi: message.params.pmi, + description: message.params.description, + ttl: message.params.ttl, + _meta: message.params._meta, + requestEventId, + }; + + const synthesizeClientDeclineError = (params: { + message: string; + }): void => { + if (pending?.progressToken) { + stopSyntheticProgress(pending.progressToken); + } + synthesizePaymentDecline(pending, params.message, req.pmi, req.amount); + }; + try { logger.info('processing payment_required', { requestEventId, pmi: message.params.pmi, @@ -692,6 +724,14 @@ export function withClientPayments( pmi: message.params.pmi, error: error instanceof Error ? error.message : String(error), }); + // Same-severity outcomes must behave the same: policy and canHandle + // declines resolve the pending request with a synthesized error. A + // handler crash must too, or the MCP request hangs until the server TTL. + synthesizeClientDeclineError({ + message: `Payment handler failed: ${ + error instanceof Error ? error.message : String(error) + }`, + }); throw error; } finally { inFlightPayReqs.delete(message.params.pay_req); @@ -786,8 +826,8 @@ export function withClientPayments( !isExplicitPaymentPendingError(message) ) { const reqId = message.id as string | number; - rawRequestCache.delete(String(reqId)); - retryCounts.delete(reqId); + rawRequestCache.delete(JSON.stringify(reqId)); + retryCounts.delete(JSON.stringify(reqId)); } } @@ -855,7 +895,10 @@ export function withClientPayments( async send(message: JSONRPCMessage): Promise { if ('method' in message && 'id' in message && message.id != null) { - rawRequestCache.set(String(message.id), message as JSONRPCRequest); + rawRequestCache.set( + JSON.stringify(message.id), + message as JSONRPCRequest, + ); } await transport.send(message); }, @@ -866,5 +909,7 @@ export function withClientPayments( }, }; + clientPaymentsWrapped.add(transport); + clientPaymentsWrapped.add(wrapped); return wrapped; } diff --git a/src/payments/server-explicit-gating.test.ts b/src/payments/server-explicit-gating.test.ts index 02f2401..7303283 100644 --- a/src/payments/server-explicit-gating.test.ts +++ b/src/payments/server-explicit-gating.test.ts @@ -735,4 +735,49 @@ describe('Explicit Gating Middleware', () => { expect(forwarded).toBe(true); expect(sentResponses).toHaveLength(0); }); + + test('falls back to the default grant TTL when the invoice carries ttl 0', async () => { + const store = new AuthorizationStore(); + const sentResponses: JSONRPCErrorResponse[] = []; + + // ttl: 0 must not birth-expire the pending grant while the invoice stays + // payable. Verification hangs so the pending entry keeps the TTL the + // middleware applied via updatePendingTtl. + const ttlZeroProcessor = { + ...processor, + async createPaymentRequired(params: { amount: number }) { + return { + amount: params.amount, + pay_req: 'pay_req', + pmi: 'fake', + ttl: 0, + }; + }, + async verifyPayment(): Promise<{ _meta?: Record }> { + return new Promise(() => {}); + }, + }; + + const mw = createExplicitGatingMiddleware({ + options: { + processors: [ttlZeroProcessor], + pricedCapabilities: [...pricedCapabilities], + }, + authorizationStore: store, + sendResponse: async (_pubkey, response) => { + sentResponses.push(response); + }, + }); + + await mw(message, ctx, async () => {}); + + expect(sentResponses[0]?.error.code).toBe(PAYMENT_REQUIRED_ERROR_CODE); + const identity = computeCanonicalInvocationIdentity( + ctx.clientPubkey, + message.method, + message.params, + ); + // Falls back to the 5-minute default window, not 0. + expect(store.getPendingRemainingMs(identity)).toBeGreaterThan(60_000); + }); }); diff --git a/src/payments/server-explicit-gating.ts b/src/payments/server-explicit-gating.ts index 5603343..8844907 100644 --- a/src/payments/server-explicit-gating.ts +++ b/src/payments/server-explicit-gating.ts @@ -71,7 +71,12 @@ export function createExplicitGatingMiddleware( return; } - const paymentTtlMs = options.paymentTtlMs ?? 300_000; + // Non-positive TTLs would birth-expire the pending grant while the + // invoice stays payable; fall back to the configured/default TTL. + const paymentTtlMs = + options.paymentTtlMs && options.paymentTtlMs > 0 + ? options.paymentTtlMs + : 300_000; // 2. Try to set pending state atomically // We use a safe default TTL here, but will override it below if the payment option has a specific TTL @@ -162,7 +167,9 @@ export function createExplicitGatingMiddleware( // verifyTimeoutMs includes standard bounds, but for grants we want to honor the // payment option's TTL explicitly if it is smaller, or the fallback paymentTtlMs. const grantTtlMs = - paymentRequired.ttl !== undefined + paymentRequired.ttl !== undefined && + Number.isFinite(paymentRequired.ttl) && + paymentRequired.ttl > 0 ? paymentRequired.ttl * 1000 : paymentTtlMs; @@ -240,6 +247,10 @@ export function createExplicitGatingMiddleware( controller.abort(); } })().catch((err) => { + // Belt-and-suspenders: if clearPending/grant itself threw inside the + // inner catch, pending would otherwise stick until TTL answering every + // retry with -32043. Idempotent on the normal path. + authorizationStore.clearPending(identity); logger.error('unhandled exception in async payment verification', { requestEventId, pmi: paymentRequired.pmi, diff --git a/src/payments/server-payments-utils.ts b/src/payments/server-payments-utils.ts index e36be75..8775c8d 100644 --- a/src/payments/server-payments-utils.ts +++ b/src/payments/server-payments-utils.ts @@ -86,7 +86,7 @@ function isResolvePriceWaiver( return 'waive' in quote && quote.waive; } -function resolvePaymentProcessor( +export function resolvePaymentProcessor( clientPmis: readonly string[] | undefined, processorsByPmi: Map, processors: readonly PaymentProcessor[], diff --git a/src/payments/server-payments.test.ts b/src/payments/server-payments.test.ts new file mode 100644 index 0000000..03c60b8 --- /dev/null +++ b/src/payments/server-payments.test.ts @@ -0,0 +1,317 @@ +import { describe, expect, test } from 'bun:test'; +import type { JSONRPCRequest } from '@contextvm/mcp-sdk/types.js'; +import { createServerPaymentsMiddleware } from './server-payments.js'; +import type { ServerPaymentsOptions } from './server-payments.js'; +import type { + CorrelatedNotificationSender, + PaymentProcessor, + PaymentRequired, + PricedCapability, +} from './types.js'; + +const PRICED: PricedCapability = { + method: 'tools/call', + name: 'expensive_tool', + amount: 1000, + currencyUnit: 'millisats', +}; + +interface FakeProcessorOptions { + /** Verification behavior once the invoice has been issued. */ + verify: 'hang' | 'resolve' | 'throw'; + verifyDelayMs?: number; +} + +interface ProcessorSpy { + processor: PaymentProcessor; + /** pay_req values minted per request event id, in order. */ + invoicesByEvent: Map; +} + +function makeProcessor(opts: FakeProcessorOptions): ProcessorSpy { + let n = 0; + const invoicesByEvent = new Map(); + const processor: PaymentProcessor = { + pmi: 'test-pmi', + async createPaymentRequired(params) { + n += 1; + const payReq = `invoice-${n}`; + const list = invoicesByEvent.get(params.requestEventId) ?? []; + list.push(payReq); + invoicesByEvent.set(params.requestEventId, list); + const required: PaymentRequired = { + amount: params.amount, + pay_req: payReq, + pmi: 'test-pmi', + }; + return required; + }, + async verifyPayment(params) { + void params; + await new Promise((r) => setTimeout(r, opts.verifyDelayMs ?? 20)); + if (opts.verify === 'throw') { + throw new Error('payment rail unreachable'); + } + if (opts.verify === 'hang') { + // Never settles: withTimeout turns this into a verification timeout. + await new Promise(() => {}); + } + return {}; + }, + }; + return { processor, invoicesByEvent }; +} + +interface SenderSpy { + sender: CorrelatedNotificationSender; + methods: string[]; + notifications: Array<{ method: string; params: Record }>; +} + +function makeSender(failOnMethod?: string): SenderSpy { + const methods: string[] = []; + const notifications: SenderSpy['notifications'] = []; + const sender: CorrelatedNotificationSender = { + async sendNotification(_pubkey, notification, _eventId) { + methods.push(notification.method); + notifications.push({ + method: notification.method, + params: notification.params as Record, + }); + if (notification.method === failOnMethod) { + throw new Error('publish failed'); + } + }, + }; + return { sender, methods, notifications }; +} + +interface MiddlewareHarness { + run: (requestEventId: string) => Promise; + forwards: () => number; + invoiceCount: (requestEventId: string) => number; +} + +function buildHarness( + spy: ProcessorSpy, + sender: CorrelatedNotificationSender, + extraOptions?: Partial, + onInvoiceIssued?: (params: { + requestEventId: string; + snapshotTtlMs: number; + }) => void, +): MiddlewareHarness { + let forwards = 0; + const middleware = createServerPaymentsMiddleware({ + sender, + options: { + processors: [spy.processor], + pricedCapabilities: [PRICED], + ...extraOptions, + }, + onInvoiceIssued, + }); + const request = (id: string): JSONRPCRequest => ({ + jsonrpc: '2.0', + id, + method: 'tools/call', + params: { name: 'expensive_tool' }, + }); + return { + run: (id) => + middleware(request(id), { clientPubkey: 'client' }, async () => { + forwards += 1; + }), + forwards: () => forwards, + invoiceCount: (id) => (spy.invoicesByEvent.get(id) ?? []).length, + }; +} + +describe('createServerPaymentsMiddleware pending-payment retention', () => { + test('payment-rail failure after invoice issuance keeps the entry, so redelivery does not mint a second invoice', async () => { + const spy = makeProcessor({ verify: 'throw', verifyDelayMs: 20 }); + const { sender } = makeSender(); + const harness = buildHarness(spy, sender, { paymentTtlMs: 5000 }); + + await expect(harness.run('evt1')).rejects.toThrow( + 'payment rail unreachable', + ); + expect(harness.invoiceCount('evt1')).toBe(1); + + // Redelivery of the same request event while the entry is within its TTL: + // the duplicate awaits the (rejected) in-flight lifecycle instead of + // minting a second invoice. (After TTL expiry a retry re-invoices by + // design — idempotency is TTL-bounded.) + await expect(harness.run('evt1')).rejects.toThrow( + 'payment rail unreachable', + ); + expect(harness.invoiceCount('evt1')).toBe(1); + }); + + test('pre-invoice failure deletes the entry, so the retry mints exactly one invoice', async () => { + const spy = makeProcessor({ verify: 'resolve' }); + const { sender } = makeSender(); + let quotes = 0; + const harness = buildHarness(spy, sender, { + paymentTtlMs: 5000, + resolvePrice: async () => { + quotes += 1; + if (quotes === 1) { + throw new Error('pricing backend down'); + } + return { amount: 1000 }; + }, + }); + + await expect(harness.run('evt1')).rejects.toThrow('pricing backend down'); + expect(harness.invoiceCount('evt1')).toBe(0); + + await harness.run('evt1'); + expect(harness.invoiceCount('evt1')).toBe(1); + expect(harness.forwards()).toBe(1); + }); + + test('redelivery after success neither re-invoices nor double-forwards', async () => { + const spy = makeProcessor({ verify: 'resolve' }); + const { sender } = makeSender(); + const harness = buildHarness(spy, sender, { paymentTtlMs: 5000 }); + + await harness.run('evt1'); + await harness.run('evt1'); + + expect(harness.invoiceCount('evt1')).toBe(1); + expect(harness.forwards()).toBe(1); + }); +}); + +describe('createServerPaymentsMiddleware payment_accepted publish failure', () => { + test('forwards the request anyway: the result is the point, not the SHOULD notification', async () => { + const spy = makeProcessor({ verify: 'resolve' }); + const { sender } = makeSender('notifications/payment_accepted'); + const harness = buildHarness(spy, sender, { paymentTtlMs: 5000 }); + + // Resolves (no throw): the publish failure is logged, the forward happens. + await harness.run('evt1'); + expect(harness.forwards()).toBe(1); + + // Entry was retained (post-invoice): redelivery must not re-invoice. + await harness.run('evt1'); + expect(harness.invoiceCount('evt1')).toBe(1); + expect(harness.forwards()).toBe(1); + }); +}); + +describe('createServerPaymentsMiddleware pending capacity', () => { + test('refuses new priced requests at capacity without minting an invoice', async () => { + const spy = makeProcessor({ verify: 'hang' }); + const senderSpy = makeSender(); + const harness = buildHarness(spy, senderSpy.sender, { + paymentTtlMs: 10_000, + maxPendingPayments: 2, + }); + + const inFlight = [ + harness.run('evt1').catch(() => {}), + harness.run('evt2').catch(() => {}), + ]; + await new Promise((r) => setTimeout(r, 10)); + + await harness.run('evt3'); + + expect(harness.invoiceCount('evt3')).toBe(0); + expect(harness.forwards()).toBe(0); + // Refused without charging: the client is told instead of hanging. + expect(senderSpy.methods).toContain('notifications/payment_rejected'); + const rejection = senderSpy.notifications.find( + (n) => n.method === 'notifications/payment_rejected', + ); + expect(rejection?.params.pmi).toBe('test-pmi'); + void inFlight; + }); + + test('purges expired entries at capacity instead of refusing', async () => { + const spy = makeProcessor({ verify: 'hang' }); + const { sender } = makeSender(); + const harness = buildHarness(spy, sender, { + paymentTtlMs: 60, + maxPendingPayments: 2, + }); + + await Promise.allSettled([harness.run('evt1'), harness.run('evt2')]); + // Both entries have now expired. + await new Promise((r) => setTimeout(r, 10)); + + // evt3 is accepted (invoice minted) rather than refused; its own + // verification still times out afterwards, which is fine here. + await harness.run('evt3').catch(() => {}); + expect(harness.invoiceCount('evt3')).toBe(1); + }); + + test('maxPendingPayments 0 refuses all priced requests', async () => { + const spy = makeProcessor({ verify: 'resolve' }); + const { sender } = makeSender(); + const harness = buildHarness(spy, sender, { maxPendingPayments: 0 }); + + await harness.run('evt1'); + + expect(harness.invoiceCount('evt1')).toBe(0); + expect(harness.forwards()).toBe(0); + }); + + test('paymentTtlMs 0 falls back to the default so redelivery dedup survives', async () => { + // A zero TTL would birth-expire the pending entry while the invoice stays + // payable — disarming the duplicate-request dedup (CEP-8). It must fall + // back to the default window instead. + const spy = makeProcessor({ verify: 'throw' }); + const { sender } = makeSender(); + const harness = buildHarness(spy, sender, { paymentTtlMs: 0 }); + + await expect(harness.run('evt-ttl0')).rejects.toThrow( + 'payment rail unreachable', + ); + + // Redelivery while the (paid) invoice is still outstanding: deduped by + // the surviving pending entry, not re-invoiced. + await expect(harness.run('evt-ttl0')).rejects.toThrow( + 'payment rail unreachable', + ); + expect(harness.invoiceCount('evt-ttl0')).toBe(1); + }); +}); + +describe('createServerPaymentsMiddleware onInvoiceIssued', () => { + test('fires when an invoice is issued, with a snapshot TTL covering the payment window', async () => { + const spy = makeProcessor({ verify: 'resolve' }); + const { sender } = makeSender(); + const issued: Array<{ requestEventId: string; snapshotTtlMs: number }> = []; + const harness = buildHarness(spy, sender, { paymentTtlMs: 5000 }, (p) => + issued.push(p), + ); + + await harness.run('evt1'); + + expect(issued).toEqual([ + { requestEventId: 'evt1', snapshotTtlMs: 5000 + 60_000 }, + ]); + }); + + test('does not fire for rejections or waivers', async () => { + const spy = makeProcessor({ verify: 'resolve' }); + const { sender } = makeSender(); + const issued: Array<{ requestEventId: string; snapshotTtlMs: number }> = []; + const harness = buildHarness( + spy, + sender, + { + paymentTtlMs: 5000, + resolvePrice: async () => ({ reject: true, message: 'no funds' }), + }, + (p) => issued.push(p), + ); + + await harness.run('evt1'); + + expect(issued).toEqual([]); + expect(harness.forwards()).toBe(0); + }); +}); diff --git a/src/payments/server-payments.ts b/src/payments/server-payments.ts index cfdfe31..dc9b64a 100644 --- a/src/payments/server-payments.ts +++ b/src/payments/server-payments.ts @@ -23,6 +23,7 @@ import { buildProcessorsByPmi, matchPricedCapability, resolveAndInitiatePayment, + resolvePaymentProcessor, } from './server-payments-utils.js'; export interface ServerPaymentsOptions { @@ -45,7 +46,12 @@ export interface ServerPaymentsOptions { /** * Maximum number of concurrent pending-payment request ids to track. * - * This is a DoS/memory-safety guardrail. + * This is a DoS/memory-safety guardrail. Once reached, new priced requests + * are refused (best-effort `payment_rejected`, no invoice minted) rather than + * evicting a live entry — an evicted live payment loses its dedup and a + * redelivery would mint a second invoice (CEP-8: MUST NOT charge twice). + * `0` refuses all priced requests. + * * @default 1000 */ maxPendingPayments?: number; @@ -81,6 +87,12 @@ type PendingPaymentState = { inFlight: Promise; }; +/** + * Extra lifetime granted to a transport route snapshot beyond the payment TTL, + * so a settled response can still be routed after a slow handler finishes. + */ +const ROUTE_SNAPSHOT_GRACE_MS = 60_000; + function createPaymentRequiredNotification(params: { amount: number; pay_req: string; @@ -128,15 +140,32 @@ export function createServerPaymentsMiddleware(params: { options: ServerPaymentsOptions; /** Pre-built PMI → processor map. Built locally when omitted (standalone use). */ processorsByPmi?: Map; + /** + * Fired once an invoice has been issued for a request, before verification + * begins. Lets the transport snapshot the correlation route/session a paid + * response may need after its route is popped (duplicate-delivery cleanup) + * or its idle session is LRU-evicted mid-payment (CEP-8). + */ + onInvoiceIssued?: (params: { + requestEventId: string; + snapshotTtlMs: number; + }) => void; }): ServerMiddlewareFn { const { sender, options } = params; const logger = createLogger('server-payments'); const processorsByPmi = params.processorsByPmi ?? buildProcessorsByPmi(options.processors, logger); - const paymentTtlMs = options.paymentTtlMs ?? DEFAULT_PAYMENT_TTL_MS; + // Non-positive TTLs would birth-expire every pending entry (disarming the + // redelivery dedup while invoices stay payable); fall back to the default, + // mirroring getVerificationTimeoutMs's guard. + const paymentTtlMs = + options.paymentTtlMs && options.paymentTtlMs > 0 + ? options.paymentTtlMs + : DEFAULT_PAYMENT_TTL_MS; + const maxPendingPayments = options.maxPendingPayments ?? 1000; const pending = new LruCache( - options.maxPendingPayments ?? 1000, + Math.max(1, maxPendingPayments), ); return async (message, ctx, forward) => { @@ -184,7 +213,50 @@ export function createServerPaymentsMiddleware(params: { return; } + // Capacity guard, before any invoice is minted: silently evicting a live + // entry here would disarm its dedup, and a redelivery of the evicted + // request would be charged twice (CEP-8: MUST NOT charge twice). + if (pending.size >= maxPendingPayments) { + purgeExpiredPending({ + pending, + nowMs: now, + maxToCheck: Number.POSITIVE_INFINITY, + }); + } + if (pending.size >= maxPendingPayments) { + logger.warn('pending payment capacity reached, refusing priced request', { + requestEventId, + method: message.method, + maxPendingPayments, + }); + // Refuse without charging: best-effort payment_rejected so the client + // stops waiting instead of hanging until its own TTL. + try { + const processor = resolvePaymentProcessor( + ctx.clientPmis, + processorsByPmi, + options.processors, + ); + await sender.sendNotification( + ctx.clientPubkey, + createPaymentRejectedNotification({ + pmi: processor.pmi, + amount: priced.amount, + message: 'payment capacity reached, retry later', + }), + requestEventId, + ); + } catch (err) { + logger.warn('failed to send payment capacity rejection', { + requestEventId, + error: err instanceof Error ? err.message : String(err), + }); + } + return; + } + // IMPORTANT: set pending state synchronously before any await to make idempotency atomic. + let invoiceIssued = false; const inFlight = (async (): Promise => { const initResult = await resolveAndInitiatePayment({ message, @@ -232,6 +304,15 @@ export function createServerPaymentsMiddleware(params: { const { paymentRequired, mergedMeta, processor, verifyTimeoutMs } = initResult; + // An invoice now exists on the payment rail — from here on, every + // outcome keeps the pending entry until TTL: the client's money may + // already be gone, and a redelivery MUST NOT mint a second invoice. + invoiceIssued = true; + params.onInvoiceIssued?.({ + requestEventId, + snapshotTtlMs: paymentTtlMs + ROUTE_SNAPSHOT_GRACE_MS, + }); + const requiredNotification = createPaymentRequiredNotification({ amount: paymentRequired.amount, pay_req: paymentRequired.pay_req, @@ -287,11 +368,23 @@ export function createServerPaymentsMiddleware(params: { _meta: verified._meta, }); - await sender.sendNotification( - ctx.clientPubkey, - acceptedNotification, - requestEventId, - ); + // payment_accepted is a SHOULD (CEP-8); the capability result is the + // point. A publish failure here (relay error, or the idle paying + // client's session LRU-evicted mid-payment) MUST NOT abort the forward — + // that would be paid-but-undelivered. Do not turn this back into an + // early return. + try { + await sender.sendNotification( + ctx.clientPubkey, + acceptedNotification, + requestEventId, + ); + } catch (err) { + logger.warn('failed to publish payment_accepted, forwarding anyway', { + requestEventId, + error: err instanceof Error ? err.message : String(err), + }); + } logger.debug('forwarding priced request after payment', { requestEventId, @@ -313,8 +406,13 @@ export function createServerPaymentsMiddleware(params: { // This guards against relay redelivery triggering a second charge. // purgeExpiredPending handles eventual cleanup. } catch (err) { - // On failure, remove immediately so the client can retry. - pending.delete(requestEventId); + // Pre-invoice failures never billed anyone — delete so the retry is free. + // Post-invoice failures keep the entry until TTL for the same reason as + // success: the client may already have paid, and a redelivery of the + // same request event MUST NOT mint a second invoice (CEP-8). + if (!invoiceIssued) { + pending.delete(requestEventId); + } throw err; } }; diff --git a/src/payments/server-transport-payments.test.ts b/src/payments/server-transport-payments.test.ts new file mode 100644 index 0000000..531a661 --- /dev/null +++ b/src/payments/server-transport-payments.test.ts @@ -0,0 +1,64 @@ +import { describe, expect, test } from 'bun:test'; +import { withServerPayments } from './server-transport-payments.js'; +import type { NostrServerTransport } from '../transport/nostr-server-transport.js'; + +describe('withServerPayments registration guard', () => { + function makeStubTransport(): { + stub: NostrServerTransport; + inboundCount: () => number; + } { + let inbound = 0; + const stub = { + setAnnouncementExtraTags: () => {}, + setAnnouncementPricingTags: () => {}, + setSupportedPaymentInteraction: () => {}, + addInboundMiddleware: () => { + inbound += 1; + }, + } as unknown as NostrServerTransport; + return { stub, inboundCount: () => inbound }; + } + + const options = { + processors: [ + { + pmi: 'fake', + async createPaymentRequired() { + return { + amount: 1, + pay_req: 'pay_req', + pmi: 'fake', + ttl: 300, + }; + }, + async verifyPayment() { + return { _meta: {} }; + }, + }, + ], + pricedCapabilities: [], + paymentInteraction: 'optional' as const, + }; + + test('registers the payments middleware pair exactly once on first call', () => { + const { stub, inboundCount } = makeStubTransport(); + withServerPayments(stub, options); + // 'optional' policy registers the transparent + explicit-gating pair. + expect(inboundCount()).toBe(2); + }); + + test('refuses double registration instead of double-charging', () => { + const { stub } = makeStubTransport(); + withServerPayments(stub, options); + expect(() => withServerPayments(stub, options)).toThrow( + /already called on this transport/, + ); + }); + + test('different transports register independently', () => { + const a = makeStubTransport(); + const b = makeStubTransport(); + withServerPayments(a.stub, options); + expect(() => withServerPayments(b.stub, options)).not.toThrow(); + }); +}); diff --git a/src/payments/server-transport-payments.ts b/src/payments/server-transport-payments.ts index d9e42dd..3169c33 100644 --- a/src/payments/server-transport-payments.ts +++ b/src/payments/server-transport-payments.ts @@ -19,10 +19,24 @@ import { NOSTR_TAGS } from '../core/constants.js'; * the notification-based flow. Pass `paymentInteraction: 'transparent'` for a * transparent-only server. */ +// Transports that already have the payments lifecycle registered. The dedup +// closures live in per-call middleware instances, so a second registration +// silently double-charges every priced request (two invoices, two forwards). +// ponytail: WeakSet keyed by transport; promote to an instance flag if a +// legitimate multi-registration use case appears. +const paymentsRegistered = new WeakSet(); + export function withServerPayments( transport: NostrServerTransport, options: ServerPaymentsOptions, ): NostrServerTransport { + if (paymentsRegistered.has(transport)) { + throw new Error( + 'withServerPayments already called on this transport; registering the payments lifecycle twice double-charges priced requests', + ); + } + paymentsRegistered.add(transport); + // Build the PMI → processor map once and share it across both middlewares. const processorsByPmi = buildProcessorsByPmi( options.processors, @@ -56,6 +70,10 @@ export function withServerPayments( sender: transport, options, processorsByPmi, + // Snapshot the correlation route when an invoice goes out so the paid + // response survives route-pop (duplicate cleanup) and session eviction. + onInvoiceIssued: ({ requestEventId, snapshotTtlMs }) => + transport.capturePaymentRouteSnapshot(requestEventId, snapshotTtlMs), }), ); diff --git a/src/transport/nostr-client-transport.ts b/src/transport/nostr-client-transport.ts index 840133a..2b00be8 100644 --- a/src/transport/nostr-client-transport.ts +++ b/src/transport/nostr-client-transport.ts @@ -108,6 +108,22 @@ export interface NostrTransportOptions extends Omit< * Implements the Transport interface from the @contextvm/mcp-sdk. * Handles request/response correlation and optional stateless mode emulation. */ +/** + * Correlation context delivered alongside inbound messages by + * `NostrClientTransport.onmessageWithContext`. `eventId` identifies the Nostr + * event that carried the message; `correlatedEventId` references the request + * event for responses and payment notifications. + */ +export type MessageContext = { + eventId: string; + correlatedEventId?: string; +}; + +/** A {@link Transport} extended with the context-aware inbound path. */ +export type TransportWithContext = Transport & { + onmessageWithContext?: (message: JSONRPCMessage, ctx: MessageContext) => void; +}; + export class NostrClientTransport extends BaseNostrTransport implements Transport @@ -129,10 +145,7 @@ export class NostrClientTransport */ public onmessage?: (message: JSONRPCMessage) => void; public onmessageWithContext: - | (( - message: JSONRPCMessage, - ctx: { eventId: string; correlatedEventId?: string }, - ) => void) + | ((message: JSONRPCMessage, ctx: MessageContext) => void) | undefined = undefined; public onclose?: () => void; public onerror?: (error: Error) => void; diff --git a/src/transport/nostr-server-transport.drop-cleanup.test.ts b/src/transport/nostr-server-transport.drop-cleanup.test.ts new file mode 100644 index 0000000..d2ff880 --- /dev/null +++ b/src/transport/nostr-server-transport.drop-cleanup.test.ts @@ -0,0 +1,127 @@ +import { describe, it, expect } from 'bun:test'; +import type { NostrEvent } from 'nostr-tools'; +import { + finalizeEvent, + generateSecretKey, + getPublicKey, +} from 'nostr-tools/pure'; +import type { RelayHandler } from '../core/interfaces.js'; +import { NostrServerTransport } from './nostr-server-transport.js'; +import { PrivateKeySigner } from '../signer/private-key-signer.js'; +import { EncryptionMode } from '../core/interfaces.js'; +import { CTXVM_MESSAGES_KIND } from '../core/index.js'; + +/** + * Regression tests for per-request state release when an inbound request is + * dropped by middleware (gating) or fails the chain: the correlation route AND + * the open-stream writer reservation must both be released. A dropped request + * never produces the normal-path response that would otherwise reap them. + */ +describe.serial('NostrServerTransport dropped-request state release', () => { + function makeTransport(): NostrServerTransport { + const relayHandler = { + async connect() {}, + async disconnect() {}, + async publish() {}, + async subscribe() { + return () => {}; + }, + } as unknown as RelayHandler; + + return new NostrServerTransport({ + signer: new PrivateKeySigner('1'.repeat(64)), + relayHandler, + encryptionMode: EncryptionMode.DISABLED, + openStream: { enabled: true }, + }); + } + + const clientSk = generateSecretKey(); + const serverPubkey = getPublicKey( + Uint8Array.from(Buffer.from('1'.repeat(64), 'hex')), + ); + + let seq = 0; + function pricedToolRequestEvent(): NostrEvent { + seq += 1; + return finalizeEvent( + { + kind: CTXVM_MESSAGES_KIND, + created_at: Math.floor(Date.now() / 1000) + seq, + tags: [['p', serverPubkey]], + content: JSON.stringify({ + jsonrpc: '2.0' as const, + id: seq, + method: 'tools/call', + params: { + name: 'expensive_tool', + arguments: {}, + _meta: { progressToken: `tok-${seq}` }, + }, + }), + }, + clientSk, + ); + } + + async function settle(): Promise { + await new Promise((resolve) => setTimeout(resolve, 50)); + } + + it('releases the route and writer reservation when middleware drops a request', async () => { + const transport = makeTransport(); + transport.onmessage = () => {}; + // Simulate any gating drop: response sent out-of-band, request never forwarded. + transport.addInboundMiddleware(async () => {}); + + const event = pricedToolRequestEvent(); + await transport['processIncomingEvent'](event); + await settle(); + + const state = transport.getInternalStateForTesting(); + expect(state.correlationStore.getEventRoute(event.id)).toBeUndefined(); + expect(transport.getOpenStreams()).toEqual([]); + await transport.close(); + }); + + it('releases the route and writer reservation when the middleware chain throws', async () => { + const transport = makeTransport(); + transport.onmessage = () => {}; + transport.onerror = () => {}; // swallow chain error + transport.addInboundMiddleware(async () => { + throw new Error('middleware exploded'); + }); + + const event = pricedToolRequestEvent(); + await transport['processIncomingEvent'](event); + await settle(); + + const state = transport.getInternalStateForTesting(); + expect(state.correlationStore.getEventRoute(event.id)).toBeUndefined(); + expect(transport.getOpenStreams()).toEqual([]); + await transport.close(); + }); + + it('keeps reaping via the normal response path when the request is forwarded', async () => { + const transport = makeTransport(); + transport.onmessage = () => {}; + transport.addInboundMiddleware(async (msg, _ctx, forward) => { + await forward(msg); + }); + + const event = pricedToolRequestEvent(); + await transport['processIncomingEvent'](event); + await settle(); + // Normal path: writer exists until the response routes through the router. + expect(transport.getOpenStreams().length).toBe(1); + + await transport.send({ + jsonrpc: '2.0', + id: event.id, + result: { ok: true }, + }); + + expect(transport.getOpenStreams()).toEqual([]); + await transport.close(); + }); +}); diff --git a/src/transport/nostr-server-transport.ts b/src/transport/nostr-server-transport.ts index a8b3e7c..08a6d8d 100644 --- a/src/transport/nostr-server-transport.ts +++ b/src/transport/nostr-server-transport.ts @@ -248,26 +248,17 @@ export class NostrServerTransport this.sessionStore = new SessionStore({ maxSessions: options.maxSessions ?? 1000, onSessionEvicted: (clientPubkey, session) => { - // Clean up all correlation data for evicted session + // Eviction is final for the session's routes. In-flight paid requests + // survive via route snapshots captured at invoice issuance + // (capturePaymentRouteSnapshot), so no veto/recreate is needed here. + // (Re-inserting into the session LRU from inside its own eviction + // callback would also corrupt the cache's capacity accounting.) const removedCount = this.correlationStore.removeRoutesForClient(clientPubkey); this.logger.info( `Evicted session for ${clientPubkey} (removed ${removedCount} routes)`, ); - // If there are still active routes (evicted early), recreate the session - // to prevent losing track of in-flight requests - if (this.correlationStore.hasActiveRoutesForClient(clientPubkey)) { - this.logger.debug( - `Recreating session ${clientPubkey} due to active routes`, - ); - this.sessionStore.getOrCreateSession( - clientPubkey, - session.isEncrypted, - ); - return; // Don't call onClientSessionEvicted for vetoed eviction - } - if (this.onClientSessionEvicted) { Promise.resolve( this.onClientSessionEvicted({ clientPubkey, session }), @@ -751,6 +742,24 @@ export class NostrServerTransport await this.outboundNotificationBroadcaster.broadcast(notification); } + /** + * Snapshots the correlation route (and current session) for a priced + * request so the settled response can be delivered even if the route was + * popped by duplicate-delivery cleanup or the idle paying client's session + * was LRU-evicted while payment settled (CEP-8). + */ + public capturePaymentRouteSnapshot( + requestEventId: string, + ttlMs: number, + ): void { + const route = this.correlationStore.getEventRoute(requestEventId); + if (!route) { + return; + } + const session = this.sessionStore.getSession(route.clientPubkey); + this.correlationStore.captureRouteSnapshot(requestEventId, session, ttlMs); + } + /** * Sends a notification to a specific client by their public key. * @param clientPubkey The public key of the target client. diff --git a/src/transport/nostr-server/correlation-store.test.ts b/src/transport/nostr-server/correlation-store.test.ts index b13e964..21fa8e1 100644 --- a/src/transport/nostr-server/correlation-store.test.ts +++ b/src/transport/nostr-server/correlation-store.test.ts @@ -473,4 +473,82 @@ describe('CorrelationStore', () => { expect(store.hasActiveRoutesForClient('c2')).toBe(true); }); }); + + describe('route snapshots', () => { + const session = { + isInitialized: true, + isEncrypted: true, + hasSentCommonTags: true, + supportsEncryption: true, + supportsEphemeralEncryption: false, + supportsOversizedTransfer: false, + supportsOpenStream: false, + }; + + it('captures and consumes a snapshot with a copy of route and session', () => { + const store = new CorrelationStore(); + store.registerEventRoute('e1', 'client1', 'req1'); + + store.captureRouteSnapshot('e1', session, 60_000); + const snapshot = store.takeRouteSnapshot('e1'); + + expect(snapshot).toBeDefined(); + expect(snapshot!.route.clientPubkey).toBe('client1'); + expect(snapshot!.route.originalRequestId).toBe('req1'); + expect(snapshot!.session).toEqual(session); + // Copy, not the same reference: later session mutation cannot rewrite history. + expect(snapshot!.session).not.toBe(session); + + // Take is consuming: a second take finds nothing. + expect(store.takeRouteSnapshot('e1')).toBeUndefined(); + }); + + it('is a no-op when no route exists at capture time', () => { + const store = new CorrelationStore(); + + store.captureRouteSnapshot('unknown', session, 60_000); + + expect(store.takeRouteSnapshot('unknown')).toBeUndefined(); + }); + + it('returns undefined for an expired snapshot', async () => { + const store = new CorrelationStore(); + store.registerEventRoute('e1', 'client1', 'req1'); + + store.captureRouteSnapshot('e1', undefined, 1); + await new Promise((r) => setTimeout(r, 5)); + + expect(store.takeRouteSnapshot('e1')).toBeUndefined(); + }); + + it('snapshot survives removeRoutesForClient (session eviction path)', () => { + const store = new CorrelationStore(); + store.registerEventRoute('e1', 'client1', 'req1'); + store.captureRouteSnapshot('e1', session, 60_000); + + store.removeRoutesForClient('client1'); + + expect(store.getEventRoute('e1')).toBeUndefined(); + expect(store.takeRouteSnapshot('e1')).toBeDefined(); + }); + + it('snapshot survives popEventRoute from cleanup paths; the router drops it via dropRouteSnapshot', () => { + const store = new CorrelationStore(); + store.registerEventRoute('e1', 'client1', 'req1'); + store.captureRouteSnapshot('e1', session, 60_000); + + // cleanupDroppedRequest / removeRoutesForClient go through popEventRoute + // and must leave the fallback snapshot intact. + const popped = store.popEventRoute('e1'); + expect(popped).toBeDefined(); + expect(store.takeRouteSnapshot('e1')).toBeDefined(); + + // After the response router uses the live route it drops the snapshot. + store.registerEventRoute('e1', 'client1', 'req1'); + store.captureRouteSnapshot('e1', session, 60_000); + store.popEventRoute('e1'); + store.dropRouteSnapshot('e1'); + expect(store.takeRouteSnapshot('e1')).toBeUndefined(); + }); + }); }); diff --git a/src/transport/nostr-server/correlation-store.ts b/src/transport/nostr-server/correlation-store.ts index 50692fd..95f7b82 100644 --- a/src/transport/nostr-server/correlation-store.ts +++ b/src/transport/nostr-server/correlation-store.ts @@ -7,6 +7,7 @@ import { DEFAULT_LRU_SIZE } from '../../core/constants.js'; import { LruCache } from '../../core/utils/lru-cache.js'; import type { NostrEvent } from 'nostr-tools'; +import type { ClientSession } from './session-store.js'; /** * Represents a route for an in-flight request. @@ -29,6 +30,22 @@ export interface EventRoute { wrapKind?: number; } +/** + * A bounded snapshot of an event route plus the client session at capture time. + * + * Used for priced requests (CEP-8): a paid response can outlive both its route + * (duplicate-delivery cleanup pops it) and its session (LRU eviction of an + * idle paying client), so the snapshot carries everything needed to deliver + * the eventual result. + */ +export interface RouteSnapshot { + /** Copy of the route at capture time. */ + route: EventRoute; + /** Shallow copy of the session at capture time, if one existed. */ + session?: ClientSession; + expiresAtMs: number; +} + /** * Options for configuring the CorrelationStore. */ @@ -54,6 +71,7 @@ export interface CorrelationStoreOptions { */ export class CorrelationStore { private readonly eventRoutes: LruCache; + private readonly routeSnapshots: LruCache; private readonly progressTokenToEventId: Map; private readonly clientEventIds: Map>; private readonly onEventRouteRemoved?: ( @@ -71,6 +89,7 @@ export class CorrelationStore { this.progressTokenToEventId = new Map(); this.clientEventIds = new Map>(); this.onEventRouteRemoved = onEventRouteRemoved; + this.routeSnapshots = new LruCache(maxEventRoutes); this.eventRoutes = new LruCache( maxEventRoutes, @@ -118,9 +137,61 @@ export class CorrelationStore { } this.eventRoutes.delete(eventId); + // Note: route snapshots deliberately survive route removal here — the + // snapshot IS the fallback for routes removed by duplicate-delivery + // cleanup or session eviction. It is consumed only via takeRouteSnapshot() + // or dropped by the response router once the live route has been used. this.onEventRouteRemoved?.(eventId, route); } + /** + * Snapshots the current route for an event, plus the session at capture + * time, so the response can still be delivered after the route or session + * is gone (priced requests, CEP-8). No-op when no route exists. + */ + captureRouteSnapshot( + eventId: string, + session: ClientSession | undefined, + ttlMs: number, + ): void { + const route = this.eventRoutes.get(eventId); + if (!route) { + return; + } + this.routeSnapshots.set(eventId, { + route: { ...route }, + session: session ? { ...session } : undefined, + expiresAtMs: Date.now() + ttlMs, + }); + } + + /** + * Drops a route snapshot without returning it. Called by the response + * router once a response was routed through the live route, so a later + * duplicate response cannot be delivered from the leftover snapshot. + */ + dropRouteSnapshot(eventId: string): void { + this.routeSnapshots.delete(eventId); + } + + /** + * Atomically consumes an unexpired route snapshot, if any. + * + * @param eventId The Nostr event ID + * @returns The snapshot, or undefined when absent or expired + */ + takeRouteSnapshot(eventId: string): RouteSnapshot | undefined { + const snapshot = this.routeSnapshots.get(eventId); + if (!snapshot) { + return undefined; + } + this.routeSnapshots.delete(eventId); + if (snapshot.expiresAtMs <= Date.now()) { + return undefined; + } + return snapshot; + } + /** * Registers a new event route for an incoming request. * @@ -308,6 +379,7 @@ export class CorrelationStore { */ clear(): void { this.eventRoutes.clear(); + this.routeSnapshots.clear(); this.progressTokenToEventId.clear(); this.clientEventIds.clear(); } diff --git a/src/transport/nostr-server/inbound-coordinator.ts b/src/transport/nostr-server/inbound-coordinator.ts index 2eefe25..debdbbf 100644 --- a/src/transport/nostr-server/inbound-coordinator.ts +++ b/src/transport/nostr-server/inbound-coordinator.ts @@ -18,7 +18,10 @@ import { injectClientPubkey, injectRequestEventId, } from '../../core/utils/utils.js'; -import { learnPeerCapabilities, mirrorRequestWrapKind } from '../capability-negotiator.js'; +import { + learnPeerCapabilities, + mirrorRequestWrapKind, +} from '../capability-negotiator.js'; import { CTXVM_MESSAGES_KIND, INITIALIZE_METHOD, @@ -336,6 +339,9 @@ export class ServerInboundCoordinator { this.deps.onerror?.( err instanceof Error ? err : new Error('inboundMiddleware failed'), ); + // A failed chain never forwards and (normally) never responds: the + // request's reserved state must be released here too, or it leaks. + this.cleanupDroppedRequest(inboundMessage); }); } catch (error) { this.deps.logger.error('Error in authorizeAndProcessEvent', { @@ -395,13 +401,18 @@ export class ServerInboundCoordinator { } /** - * Cleans up request correlation for a request that was dropped by middleware. + * Cleans up per-request state reserved for a request that was dropped by + * middleware or failed the chain. Single owner: releases the correlation + * route and the open-stream writer reservation together, since a dropped + * request never produces the normal-path response that would reap them. */ public cleanupDroppedRequest(message: JSONRPCMessage): void { if (!isJSONRPCRequest(message)) { return; } - this.deps.correlationStore.popEventRoute(String(message.id)); + const eventId = String(message.id); + this.deps.correlationStore.popEventRoute(eventId); + this.deps.openStreamFactory.releaseUnusedWriter(eventId); } /** diff --git a/src/transport/nostr-server/open-stream-factory.test.ts b/src/transport/nostr-server/open-stream-factory.test.ts index 824e8b8..fe24661 100644 --- a/src/transport/nostr-server/open-stream-factory.test.ts +++ b/src/transport/nostr-server/open-stream-factory.test.ts @@ -144,6 +144,40 @@ describe('ServerOpenStreamFactory.deferIfStreamActive', () => { }); }); +describe('ServerOpenStreamFactory.releaseUnusedWriter', () => { + test('releases a never-started writer and its metadata, idempotently', () => { + const { factory } = createFactory(); + + factory.createWriterIfEnabled('evt-rel', 'a'.repeat(64), 'token-rel'); + expect(factory.getWriter('evt-rel')).toBeDefined(); + + factory.releaseUnusedWriter('evt-rel'); + factory.releaseUnusedWriter('evt-rel'); // second call is a no-op + + expect(factory.getWriter('evt-rel')).toBeUndefined(); + expect(factory.getWritersMap().has('evt-rel')).toBe(false); + expect(factory.getOpenStreams().length).toBe(0); + }); + + test('disposes a started writer, clearing keepalive state', async () => { + const { factory } = createFactory(); + + const writer = factory.createWriterIfEnabled( + 'evt-started', + 'a'.repeat(64), + 'token-started', + ); + expect(writer).toBeDefined(); + await writer!.start(); + expect(writer!.hasStarted).toBe(true); + + factory.releaseUnusedWriter('evt-started'); + + expect(writer!.isActive).toBe(false); + expect(factory.getWriter('evt-started')).toBeUndefined(); + }); +}); + describe('ServerOpenStreamFactory.getOpenStreams', () => { test('lists active writers with resolved client context and metadata', async () => { const { factory } = createFactory(); diff --git a/src/transport/nostr-server/open-stream-factory.ts b/src/transport/nostr-server/open-stream-factory.ts index 1d9f92a..085d0b0 100644 --- a/src/transport/nostr-server/open-stream-factory.ts +++ b/src/transport/nostr-server/open-stream-factory.ts @@ -222,13 +222,26 @@ export class ServerOpenStreamFactory { // the tool never streamed to it. Drop the unused writer so the response is // sent normally instead of being deferred indefinitely, and so it does not // leak. See docs/ISSUE-open-stream-progress-token-conflict.md (Part A). - if (existingWriter) { - this.writers.delete(eventId); - this.writerMeta.delete(eventId); - } + this.releaseUnusedWriter(eventId); return false; } + /** + * Releases the writer reserved for a request that will never produce a + * normal-path response (middleware drop or chain failure). This is the only + * reaper besides a routed response and full transport teardown, so every + * drop path must call it or the entry leaks until close(). + */ + public releaseUnusedWriter(eventId: string): void { + const writer = this.writers.get(eventId); + if (!writer) { + return; + } + this.writers.delete(eventId); + this.writerMeta.delete(eventId); + writer.dispose(); + } + /** * Conditionally creates a new OpenStreamWriter if the client supports it. * Returns the created writer (or `undefined` when creation was skipped) so diff --git a/src/transport/nostr-server/outbound-response-router.test.ts b/src/transport/nostr-server/outbound-response-router.test.ts index 1a27922..c2a5f3e 100644 --- a/src/transport/nostr-server/outbound-response-router.test.ts +++ b/src/transport/nostr-server/outbound-response-router.test.ts @@ -23,9 +23,7 @@ const testLogger: Logger = { const CLIENT_PUBKEY = 'a'.repeat(64); -function createSession( - overrides: Partial = {}, -): ClientSession { +function createSession(overrides: Partial = {}): ClientSession { return { isInitialized: true, isEncrypted: true, @@ -44,6 +42,7 @@ interface CapturedDeps { deps: OutboundResponseRouterDeps; chooseCalls: Array<{ fallbackWrapKind?: number }>; sentGiftWrapKinds: Array; + sentMessages: JSONRPCMessage[]; } function createRouterWithCapturedDeps( @@ -53,6 +52,7 @@ function createRouterWithCapturedDeps( ): CapturedDeps { const chooseCalls: CapturedDeps['chooseCalls'] = []; const sentGiftWrapKinds: CapturedDeps['sentGiftWrapKinds'] = []; + const sentMessages: CapturedDeps['sentMessages'] = []; const deps = { correlationStore, @@ -87,6 +87,7 @@ function createRouterWithCapturedDeps( if (options.failSend) { throw new Error('send failed'); } + sentMessages.push(_message); sentGiftWrapKinds.push(giftWrapKind); return 'inner-event-id'; }, @@ -95,7 +96,7 @@ function createRouterWithCapturedDeps( logger: testLogger, } as unknown as OutboundResponseRouterDeps; - return { deps, chooseCalls, sentGiftWrapKinds }; + return { deps, chooseCalls, sentGiftWrapKinds, sentMessages }; } const gatingErrorResponse: JSONRPCErrorResponse = { @@ -114,10 +115,8 @@ describe('OutboundResponseRouter.routeTargeted', () => { undefined, EPHEMERAL_GIFT_WRAP_KIND, ); - const { deps, chooseCalls, sentGiftWrapKinds } = createRouterWithCapturedDeps( - correlationStore, - createSession(), - ); + const { deps, chooseCalls, sentGiftWrapKinds } = + createRouterWithCapturedDeps(correlationStore, createSession()); await new OutboundResponseRouter(deps).routeTargeted( CLIENT_PUBKEY, @@ -135,10 +134,8 @@ describe('OutboundResponseRouter.routeTargeted', () => { }); test('sends without a wrap-kind hint when no route is recorded', async () => { - const { deps, chooseCalls, sentGiftWrapKinds } = createRouterWithCapturedDeps( - new CorrelationStore({}), - createSession(), - ); + const { deps, chooseCalls, sentGiftWrapKinds } = + createRouterWithCapturedDeps(new CorrelationStore({}), createSession()); await new OutboundResponseRouter(deps).routeTargeted( CLIENT_PUBKEY, @@ -150,12 +147,56 @@ describe('OutboundResponseRouter.routeTargeted', () => { expect(sentGiftWrapKinds).toEqual([undefined]); }); - test('does not send when the client has no active session', async () => { - const { deps, chooseCalls, sentGiftWrapKinds } = createRouterWithCapturedDeps( + test('restores the client request id the coordinator rewrote to the event id', async () => { + const correlationStore = new CorrelationStore({}); + correlationStore.registerEventRoute( + 'evt-rewritten', + CLIENT_PUBKEY, + 42, + undefined, + EPHEMERAL_GIFT_WRAP_KIND, + ); + const { deps, sentMessages } = createRouterWithCapturedDeps( + correlationStore, + createSession(), + ); + + // The gating middleware echoes the id it saw: the event id, post-rewrite. + const errorWithEventId: JSONRPCErrorResponse = { + ...gatingErrorResponse, + id: 'evt-rewritten', + }; + await new OutboundResponseRouter(deps).routeTargeted( + CLIENT_PUBKEY, + errorWithEventId, + 'evt-rewritten', + ); + + // On the wire: the caller's original JSON-RPC id, not the routing key. + expect((sentMessages[0] as JSONRPCErrorResponse).id).toBe(42); + // Non-mutating: the caller's response object is untouched. + expect(errorWithEventId.id).toBe('evt-rewritten'); + }); + + test('keeps the response id as-is when no route is recorded', async () => { + const { deps, sentMessages } = createRouterWithCapturedDeps( new CorrelationStore({}), createSession(), ); + await new OutboundResponseRouter(deps).routeTargeted( + CLIENT_PUBKEY, + { ...gatingErrorResponse, id: 'evt-unknown' }, + 'evt-unknown', + ); + + expect((sentMessages[0] as JSONRPCErrorResponse).id).toBe('evt-unknown'); + }); + + test('does not send when the client has no active session', async () => { + const { deps, chooseCalls, sentGiftWrapKinds } = + createRouterWithCapturedDeps(new CorrelationStore({}), createSession()); + await new OutboundResponseRouter(deps).routeTargeted( 'b'.repeat(64), gatingErrorResponse, @@ -206,3 +247,107 @@ describe('OutboundResponseRouter.routeTargeted', () => { expect(correlationStore.getRequestEvent('evt-a3')).toBe(requestEvent); }); }); + +describe('OutboundResponseRouter.route payment route snapshot fallback', () => { + // Fresh object per test: route() restores the original request id by + // mutating the response in place, so a shared literal would leak state. + const paidResult = () => ({ + jsonrpc: '2.0' as const, + id: 'evt-paid', + result: { content: [] }, + }); + + test('delivers from the snapshot when the live route was popped by duplicate-delivery cleanup', async () => { + const correlationStore = new CorrelationStore({}); + correlationStore.registerEventRoute( + 'evt-paid', + CLIENT_PUBKEY, + 'original-request-id', + undefined, + EPHEMERAL_GIFT_WRAP_KIND, + ); + const session = createSession(); + correlationStore.captureRouteSnapshot('evt-paid', session, 60_000); + // Duplicate delivery popped the route while payment settled. + correlationStore.popEventRoute('evt-paid'); + + const { deps, sentGiftWrapKinds } = createRouterWithCapturedDeps( + correlationStore, + session, + ); + const sent: JSONRPCMessage[] = []; + ( + deps as { sendMcpMessage: OutboundResponseRouterDeps['sendMcpMessage'] } + ).sendMcpMessage = async ( + message, + _target, + _kind, + _tags, + _encrypt, + _onCreate, + giftWrapKind, + ) => { + sent.push(message); + sentGiftWrapKinds.push(giftWrapKind); + return 'inner-event-id'; + }; + + await new OutboundResponseRouter(deps).route(paidResult()); + + expect(sent).toHaveLength(1); + // The client sees the result under its original request id. + expect((sent[0] as { id?: string }).id).toBe('original-request-id'); + expect(sentGiftWrapKinds).toEqual([EPHEMERAL_GIFT_WRAP_KIND]); + }); + + test('delivers from the snapshot session when the paying client was evicted', async () => { + const correlationStore = new CorrelationStore({}); + correlationStore.registerEventRoute( + 'evt-paid', + CLIENT_PUBKEY, + 'original-request-id', + ); + correlationStore.captureRouteSnapshot( + 'evt-paid', + createSession({ isEncrypted: false }), + 60_000, + ); + // Session eviction removed the route entirely. + correlationStore.removeRoutesForClient(CLIENT_PUBKEY); + + const { deps } = createRouterWithCapturedDeps( + correlationStore, + createSession(), + ); + // The live session store no longer knows the client. + (deps.sessionStore as { getSession: (p: string) => unknown }).getSession = + () => undefined; + + const sent: JSONRPCMessage[] = []; + ( + deps as { sendMcpMessage: OutboundResponseRouterDeps['sendMcpMessage'] } + ).sendMcpMessage = async (message) => { + sent.push(message); + return 'inner-event-id'; + }; + + await new OutboundResponseRouter(deps).route(paidResult()); + + expect(sent).toHaveLength(1); + expect((sent[0] as { id?: string }).id).toBe('original-request-id'); + }); + + test('still errors when neither route nor snapshot exists', async () => { + const { deps } = createRouterWithCapturedDeps( + new CorrelationStore({}), + createSession(), + ); + const errors: Error[] = []; + deps.onerror = (e) => errors.push(e); + + await new OutboundResponseRouter(deps).route(paidResult()); + + expect(errors).toHaveLength(1); + expect(errors[0].message).toContain('No pending request found'); + }); +}); diff --git a/src/transport/nostr-server/outbound-response-router.ts b/src/transport/nostr-server/outbound-response-router.ts index 130b8e5..fb79ade 100644 --- a/src/transport/nostr-server/outbound-response-router.ts +++ b/src/transport/nostr-server/outbound-response-router.ts @@ -112,7 +112,21 @@ export class OutboundResponseRouter { return; } - const route = this.deps.correlationStore.popEventRoute(nostrEventId); + const poppedRoute = this.deps.correlationStore.popEventRoute(nostrEventId); + if (poppedRoute) { + // The live route was used: drop any snapshot so a later duplicate + // response for this id cannot be delivered from it. + this.deps.correlationStore.dropRouteSnapshot(nostrEventId); + } + // Route miss fallback: a paid request's route may have been popped by + // duplicate-delivery cleanup or removed with its evicted session while + // payment settled. The snapshot captured at invoice issuance (CEP-8) + // carries the routing fields — and the session — needed to deliver the + // settled result anyway. + const snapshot = poppedRoute + ? undefined + : this.deps.correlationStore.takeRouteSnapshot(nostrEventId); + const route = poppedRoute ?? snapshot?.route; if (!route) { this.deps.onerror?.( @@ -121,10 +135,18 @@ export class OutboundResponseRouter { return; } + if (snapshot) { + this.deps.logger.info('Delivering response from payment route snapshot', { + eventId: nostrEventId, + clientPubkey: route.clientPubkey, + }); + } + const pendingEviction = this.deps.openStreamFactory.takePendingEviction(nostrEventId); const session = this.deps.sessionStore.getSession(route.clientPubkey) ?? + snapshot?.session ?? pendingEviction?.session; if (!session) { @@ -323,16 +345,26 @@ export class OutboundResponseRouter { this.maybeAppendPaymentInteractionDisclosure(tags, session); + const route = this.deps.correlationStore.getEventRoute(requestEventId); + const giftWrapKind = this.deps.chooseGiftWrapKind({ session, // Non-destructive read: the route must stay registered for the normal // response/cleanup lifecycle that runs after this early rejection. - fallbackWrapKind: - this.deps.correlationStore.getEventRoute(requestEventId)?.wrapKind, + fallbackWrapKind: route?.wrapKind, }); + // Restore the client's original request id before publishing: the inbound + // coordinator rewrites request ids to the event id for routing, and this + // early-rejection exit path must not leak that rewrite onto the wire + // (mirrors route()'s restore; JSON-RPC responses MUST echo the caller's + // id, and both SDK clients ignore the wire id anyway). + const responseToSend = route + ? { ...response, id: route.originalRequestId } + : response; + await this.deps.sendMcpMessage( - response, + responseToSend, clientPubkey, CTXVM_MESSAGES_KIND, tags,