diff --git a/.changeset/cep-47-redirect.md b/.changeset/cep-47-redirect.md new file mode 100644 index 0000000..ccfc11b --- /dev/null +++ b/.changeset/cep-47-redirect.md @@ -0,0 +1,5 @@ +--- +"@contextvm/sdk": minor +--- + +Add CEP-47 Server Redirect support. Includes `-32044` error code, `createRedirectMiddleware` for server-side redirects, and `withClientRedirect` for client-side transparent re-issuance. diff --git a/package.json b/package.json index 8c64461..76ce2b0 100644 --- a/package.json +++ b/package.json @@ -72,6 +72,14 @@ "./payments/*": { "types": "./dist/esm/payments/*.d.ts", "default": "./dist/esm/payments/*.js" + }, + "./redirect": { + "types": "./dist/esm/redirect/index.d.ts", + "default": "./dist/esm/redirect/index.js" + }, + "./redirect/*": { + "types": "./dist/esm/redirect/*.d.ts", + "default": "./dist/esm/redirect/*.js" } }, "types": "./dist/esm/index.d.ts", diff --git a/src/__mocks__/mock-relay-server.ts b/src/__mocks__/mock-relay-server.ts index 5b10ec9..4401ae1 100644 --- a/src/__mocks__/mock-relay-server.ts +++ b/src/__mocks__/mock-relay-server.ts @@ -200,7 +200,8 @@ export function startMockRelay( this.send(['OK', event.id, true, '']); for (const [uniqueSubId, { instance, filters }] of state.subs.entries()) { - if (matchFilters(filters, event)) { + const isMatch = matchFilters(filters, event); + if (isMatch) { const originalSubId = uniqueSubId.includes(':') ? uniqueSubId.split(':').slice(1).join(':') : uniqueSubId; diff --git a/src/gateway/gateway-redirect.test.ts b/src/gateway/gateway-redirect.test.ts new file mode 100644 index 0000000..3025d5a --- /dev/null +++ b/src/gateway/gateway-redirect.test.ts @@ -0,0 +1,157 @@ +import { + afterAll, + afterEach, + beforeAll, + describe, + expect, + test, +} from 'bun:test'; +import { sleep } from 'bun'; +import { Client } from '@contextvm/mcp-sdk/client'; +import { McpServer } from '@contextvm/mcp-sdk/server/mcp'; +import { InMemoryTransport } from '@contextvm/mcp-sdk/inMemory'; +import { z } from 'zod'; +import { bytesToHex } from 'nostr-tools/utils'; +import { generateSecretKey, getPublicKey } from 'nostr-tools/pure'; +import { ApplesauceRelayPool } from '../relay/applesauce-relay-pool.js'; +import { PrivateKeySigner } from '../signer/private-key-signer.js'; +import { EncryptionMode } from '../core/interfaces.js'; +import { NostrServerTransport } from '../transport/nostr-server-transport.js'; +import { NostrClientTransport } from '../transport/nostr-client-transport.js'; +import { NostrMCPGateway } from './index.js'; +import { withClientRedirect } from '../redirect/index.js'; +import { withClientPayments } from '../payments/index.js'; +import { + spawnMockRelay, + clearRelayCache, +} from '../__mocks__/test-relay-helpers.js'; + +/** + * Proves `NostrMCPGateway` correctly wires `redirectConfig` via + * `withServerRedirect` on its internal server transport. + * + * Mirrors `gateway-payments.test.ts` but exercises the redirect path: + * a gateway configured with `redirectConfig` should emit -32044 to + * its Nostr clients, proving the middleware is wired and the high-level + * `redirectConfig` option works end-to-end. + */ +describe.serial('NostrMCPGateway redirect wiring', () => { + let relayUrl: string; + let httpUrl: string; + let stopRelay: (() => void) | undefined; + + beforeAll(async () => { + const relay = await spawnMockRelay(); + relayUrl = relay.relayUrl; + httpUrl = relay.httpUrl; + stopRelay = relay.stop; + }); + + afterEach(async () => { + await clearRelayCache(httpUrl); + }); + + afterAll(async () => { + stopRelay?.(); + await sleep(100); + }); + + test('gateway follows a server redirect configured via redirectConfig', async () => { + // Target Server — a real MCP server reachable over Nostr + const targetSK = generateSecretKey(); + const targetServer = new McpServer({ + name: 'gateway-target-server', + version: '1.0.0', + }); + targetServer.registerTool( + 'echo', + { + title: 'Echo', + description: 'Echoes the message', + inputSchema: { message: z.string() }, + }, + async ({ message }: { message: string }) => ({ + content: [{ type: 'text', text: `GW-Redirected: ${message}` }], + }), + ); + const targetTransport = new NostrServerTransport({ + signer: new PrivateKeySigner(bytesToHex(targetSK)), + relayHandler: new ApplesauceRelayPool([relayUrl]), + encryptionMode: EncryptionMode.DISABLED, + }); + await targetServer.connect(targetTransport); + const targetPubkey = getPublicKey(targetSK); + + // Gateway — bridges a local MCP server and exposes it over Nostr. + // `redirectConfig` injects the server-side redirect middleware so that + // all inbound requests get redirected to the target server. + const [mcpTransport, gatewayMcpTransport] = + InMemoryTransport.createLinkedPair(); + const mcpServer = new McpServer({ + name: 'gateway-initial-server', + version: '1.0.0', + }); + await mcpServer.connect(mcpTransport); + + const gatewaySK = generateSecretKey(); + const gateway = new NostrMCPGateway({ + mcpClientTransport: gatewayMcpTransport, + nostrTransportOptions: { + signer: new PrivateKeySigner(bytesToHex(gatewaySK)), + relayHandler: new ApplesauceRelayPool([relayUrl]), + encryptionMode: EncryptionMode.DISABLED, + publishRelayList: false, + }, + redirectConfig: { + resolveRedirect: async () => ({ + target: targetPubkey, + relays: [relayUrl], + }), + }, + }); + await gateway.start(); + const gatewayPubkey = getPublicKey(gatewaySK); + + // Client — connects to the gateway over Nostr with redirect support. + // Mirrors the `NostrMCPProxy` wrapping order: payments(base) → redirect(…) + const clientSigner = new PrivateKeySigner( + bytesToHex(generateSecretKey()), + ); + const baseClientTransport = withClientPayments( + new NostrClientTransport({ + signer: clientSigner, + relayHandler: new ApplesauceRelayPool([relayUrl]), + serverPubkey: gatewayPubkey, + encryptionMode: EncryptionMode.DISABLED, + }), + {}, + ); + const clientTransport = withClientRedirect( + baseClientTransport, + { + signer: clientSigner, + encryptionMode: EncryptionMode.DISABLED, + wrapTransport: (t) => withClientPayments(t, {}), + }, + { maxRedirects: 2 }, + ); + + const client = new Client({ + name: 'gateway-redirect-client', + version: '1.0.0', + }); + await client.connect(clientTransport as never); + + const res = await client.callTool({ + name: 'echo', + arguments: { message: 'hello gateway' }, + }); + expect((res.content as Array<{ text: string }>)[0].text).toBe( + 'GW-Redirected: hello gateway', + ); + + await client.close(); + await gateway.stop(); + await targetServer.close(); + }, 20000); +}); diff --git a/src/gateway/index.ts b/src/gateway/index.ts index 484f6fe..652992a 100644 --- a/src/gateway/index.ts +++ b/src/gateway/index.ts @@ -9,6 +9,10 @@ import { } from '../transport/nostr-server-transport.js'; import { withServerPayments } from '../payments/index.js'; import type { ServerPaymentsOptions } from '../payments/server-payments.js'; +import { + withServerRedirect, + type ServerRedirectConfig, +} from '../redirect/index.js'; import { NOTIFICATIONS_INITIALIZED_METHOD } from '../core/index.js'; import { createLogger } from '../core/utils/logger.js'; import { LruCache } from '../core/utils/lru-cache.js'; @@ -80,6 +84,12 @@ export interface NostrMCPGatewayOptions { * receive an invoice notification. */ paymentOptions?: ServerPaymentsOptions; + + /** + * CEP-47 server redirect configuration. + * When provided, evaluates inbound requests and redirects clients before payment gating. + */ + redirectConfig?: ServerRedirectConfig; } /** @@ -130,13 +140,20 @@ export class NostrMCPGateway { this.closeClientTransport(clientPubkey), }); + // Wrap with `withServerRedirect` first so redirected requests halt before payment gating. + let transport = nostrServerTransport; + if (options.redirectConfig) { + transport = withServerRedirect(transport, options.redirectConfig); + } + // Wrap with `withServerPayments` so CEP-8 gating, PMI/cap advertisement and // payment_interaction negotiation are attached when a processor + priced // capabilities are provided. Mirrors the client-side `withClientPayments` // wiring in NostrMCPProxy. No-op when paymentOptions is omitted. - this.nostrServerTransport = options.paymentOptions - ? withServerPayments(nostrServerTransport, options.paymentOptions) - : nostrServerTransport; + if (options.paymentOptions) { + transport = withServerPayments(transport, options.paymentOptions); + } + this.nostrServerTransport = transport; if (this.createMcpClientTransport) { this.clientTransportPromises = new Map(); diff --git a/src/index.ts b/src/index.ts index bdff4c1..7137693 100644 --- a/src/index.ts +++ b/src/index.ts @@ -5,3 +5,4 @@ export * from './gateway/index.js'; export * from './proxy/index.js'; export * from './transport/index.js'; export * from './payments/index.js'; +export * from './redirect/index.js'; diff --git a/src/payments/constants.ts b/src/payments/constants.ts index 63b7add..1354a48 100644 --- a/src/payments/constants.ts +++ b/src/payments/constants.ts @@ -29,6 +29,9 @@ export const PAYMENT_REQUIRED_ERROR_CODE = -32042; /** CEP-8 explicit-gating JSON-RPC error: payment pending. */ export const PAYMENT_PENDING_ERROR_CODE = -32043; +/** CEP-47 JSON-RPC error: server redirect. */ +export const REDIRECT_ERROR_CODE = -32044; + /** * CEP-8 unsupported payment_interaction negotiation error. * diff --git a/src/proxy/index.ts b/src/proxy/index.ts index 9212fe9..cccff27 100644 --- a/src/proxy/index.ts +++ b/src/proxy/index.ts @@ -6,6 +6,10 @@ import { } from '../transport/nostr-client-transport.js'; import { withClientPayments } from '../payments/client-payments.js'; import type { ClientPaymentsOptions } from '../payments/client-payments.js'; +import { + withClientRedirect, + type ClientRedirectOptions, +} from '../redirect/index.js'; import { createLogger } from '../core/utils/logger.js'; const logger = createLogger('proxy'); @@ -34,6 +38,11 @@ export interface NostrMCPProxyOptions { * (programmatic) payment so the proxy can settle invoices itself. */ paymentOptions?: ClientPaymentsOptions; + /** + * CEP-47 client redirect configuration and hooks. + * When provided, transparently follows server redirections before payment gating. + */ + redirectOptions?: ClientRedirectOptions; } /** @@ -54,10 +63,32 @@ export class NostrMCPProxy { // No handlers ⇒ PMI-agnostic: explicit_gating surfaces `-32042` as an error; // transparent forwards `payment_required` and keeps the request alive with // synthetic progress. - this.nostrTransport = withClientPayments( + const initialTransport = withClientPayments( new NostrClientTransport(options.nostrTransportOptions), options.paymentOptions ?? {}, ); + + if (options.redirectOptions) { + const { + serverPubkey: _serverPubkey, + relayHandler: _relayHandler, + discoveryRelayUrls: _discoveryRelayUrls, + fallbackOperationalRelayUrls: _fallbackOperationalRelayUrls, + ...baseOpts + } = options.nostrTransportOptions; + + this.nostrTransport = withClientRedirect( + initialTransport, + { + ...baseOpts, + wrapTransport: (t) => + withClientPayments(t, options.paymentOptions ?? {}), + }, + options.redirectOptions, + ); + } else { + this.nostrTransport = initialTransport; + } } /** diff --git a/src/proxy/proxy-redirect.test.ts b/src/proxy/proxy-redirect.test.ts new file mode 100644 index 0000000..f22729b --- /dev/null +++ b/src/proxy/proxy-redirect.test.ts @@ -0,0 +1,121 @@ +import { + afterAll, + afterEach, + beforeAll, + describe, + expect, + test, +} from 'bun:test'; +import { sleep } from 'bun'; +import { Client } from '@contextvm/mcp-sdk/client'; +import { McpServer } from '@contextvm/mcp-sdk/server/mcp'; +import { InMemoryTransport } from '@contextvm/mcp-sdk/inMemory'; +import { z } from 'zod'; +import { bytesToHex } from 'nostr-tools/utils'; +import { generateSecretKey, getPublicKey } from 'nostr-tools/pure'; +import { ApplesauceRelayPool } from '../relay/applesauce-relay-pool.js'; +import { PrivateKeySigner } from '../signer/private-key-signer.js'; +import { EncryptionMode } from '../core/interfaces.js'; +import { NostrServerTransport } from '../transport/nostr-server-transport.js'; +import { NostrMCPProxy } from './index.js'; +import { withServerRedirect } from '../redirect/index.js'; +import { + spawnMockRelay, + clearRelayCache, +} from '../__mocks__/test-relay-helpers.js'; + +describe.serial('NostrMCPProxy redirect wiring', () => { + let relayUrl: string; + let httpUrl: string; + let stopRelay: (() => void) | undefined; + + beforeAll(async () => { + const relay = await spawnMockRelay(); + relayUrl = relay.relayUrl; + httpUrl = relay.httpUrl; + stopRelay = relay.stop; + }); + + afterEach(async () => { + await clearRelayCache(httpUrl); + }); + + afterAll(async () => { + stopRelay?.(); + await sleep(100); + }); + + test('proxy correctly follows a server redirect', async () => { + // Target Server + const targetSK = generateSecretKey(); + const targetServer = new McpServer({ + name: 'proxy-target-server', + version: '1.0.0', + }); + targetServer.registerTool( + 'echo', + { + title: 'Echo', + description: 'Echoes the message', + inputSchema: { message: z.string() }, + }, + async ({ message }: { message: string }) => ({ + content: [{ type: 'text', text: `Redirected: ${message}` }], + }), + ); + const targetTransport = new NostrServerTransport({ + signer: new PrivateKeySigner(bytesToHex(targetSK)), + relayHandler: new ApplesauceRelayPool([relayUrl]), + encryptionMode: EncryptionMode.DISABLED, + }); + await targetServer.connect(targetTransport); + const targetPubkey = getPublicKey(targetSK); + + // Initial Server (redirects to Target Server) + const initialSK = generateSecretKey(); + const initialServer = new McpServer({ + name: 'proxy-initial-server', + version: '1.0.0', + }); + // It doesn't even need the tool registered because the middleware intercepts it + const initialTransport = withServerRedirect( + new NostrServerTransport({ + signer: new PrivateKeySigner(bytesToHex(initialSK)), + relayHandler: new ApplesauceRelayPool([relayUrl]), + encryptionMode: EncryptionMode.DISABLED, + }), + { + resolveRedirect: async () => ({ target: targetPubkey, relays: [relayUrl] }), + } + ); + await initialServer.connect(initialTransport); + const initialPubkey = getPublicKey(initialSK); + + // Host side of an in-memory pair; the proxy relays MCP through it. + const [hostTransport, clientTransport] = + InMemoryTransport.createLinkedPair(); + + const proxy = new NostrMCPProxy({ + mcpHostTransport: hostTransport, + nostrTransportOptions: { + signer: new PrivateKeySigner(bytesToHex(generateSecretKey())), + relayHandler: new ApplesauceRelayPool([relayUrl]), + serverPubkey: initialPubkey, + encryptionMode: EncryptionMode.DISABLED, + }, + redirectOptions: { maxRedirects: 2 }, + }); + await proxy.start(); + + const client = new Client({ name: 'proxy-host-client', version: '1.0.0' }); + await client.connect(clientTransport); + + const res = await client.callTool({ name: 'echo', arguments: { message: 'hello proxy' } }); + expect((res.content as Array<{ text: string }>)[0].text).toBe('Redirected: hello proxy'); + + await client.close(); + await proxy.stop(); + await targetServer.close(); + await initialServer.close(); + }, 20000); +}); diff --git a/src/redirect/client-redirect.test.ts b/src/redirect/client-redirect.test.ts new file mode 100644 index 0000000..7ada0c3 --- /dev/null +++ b/src/redirect/client-redirect.test.ts @@ -0,0 +1,300 @@ +import { describe, expect, test } from 'bun:test'; +import type { + JSONRPCMessage, + JSONRPCRequest, +} from '@contextvm/mcp-sdk/types.js'; +import type { Transport } from '@contextvm/mcp-sdk/shared/transport'; +import { withClientRedirect } from './client-redirect.js'; +import { REDIRECT_ERROR_CODE } from '../payments/constants.js'; + +class FakeTransport implements Transport { + public onmessage?: (msg: JSONRPCMessage) => void; + public onmessageWithContext?: ( + msg: JSONRPCMessage, + ctx: { eventId: string }, + ) => void; + public onerror?: (err: Error) => void; + public onclose?: () => void; + public sentMessages: JSONRPCMessage[] = []; + public started = false; + public closed = false; + public readonly serverPubkey: string; + + constructor(serverPubkey = 'a'.repeat(64)) { + this.serverPubkey = serverPubkey; + } + + async start(): Promise { + this.started = true; + } + + async send(msg: JSONRPCMessage): Promise { + this.sentMessages.push(msg); + } + + async close(): Promise { + this.closed = true; + this.onclose?.(); + } + + emit(msg: JSONRPCMessage, ctx = { eventId: 'evt-1' }): void { + if (this.onmessageWithContext) { + this.onmessageWithContext(msg, ctx); + } else { + this.onmessage?.(msg); + } + } +} + +const dummySigner = '1'.repeat(64); + +describe('withClientRedirect', () => { + test('passes normal responses through unchanged', async () => { + const baseTransport = new FakeTransport('a'.repeat(64)); + const wrapped = withClientRedirect(baseTransport, { signer: dummySigner }); + + const received: JSONRPCMessage[] = []; + wrapped.onmessage = (msg) => received.push(msg); + await wrapped.start(); + + const req: JSONRPCRequest = { + jsonrpc: '2.0', + id: 1, + method: 'tools/list', + }; + await wrapped.send(req); + + baseTransport.emit({ + jsonrpc: '2.0', + id: 1, + result: { tools: [] }, + } as unknown as JSONRPCMessage); + + expect(received.length).toBe(1); + expect((received[0] as { result?: unknown }).result).toEqual({ tools: [] }); + }); + + test('follows -32044 redirect: starts new transport, closes old, and re-issues request', async () => { + const baseTransport = new FakeTransport('a'.repeat(64)); + const newTransportsCreated: Transport[] = []; + let onRedirectCalled = false; + let redirectHop = 0; + + const targetPubkey = 'b'.repeat(64); + + const wrapped = withClientRedirect( + baseTransport, + { + signer: dummySigner, + wrapTransport: () => { + const next = new FakeTransport(targetPubkey); + newTransportsCreated.push(next); + return next; + }, + }, + { + onRedirect: (data, hop) => { + onRedirectCalled = true; + redirectHop = hop; + expect(data.target).toBe(targetPubkey); + }, + }, + ); + + const received: JSONRPCMessage[] = []; + wrapped.onmessage = (msg) => received.push(msg); + await wrapped.start(); + + const req: JSONRPCRequest = { + jsonrpc: '2.0', + id: 'req-100', + method: 'tools/call', + params: { name: 'test' }, + }; + await wrapped.send(req); + expect(baseTransport.sentMessages.length).toBe(1); + + // Emit -32044 redirect from server A + baseTransport.emit({ + jsonrpc: '2.0', + id: 'req-100', + error: { + code: REDIRECT_ERROR_CODE, + message: 'Redirect', + data: { target: targetPubkey, relays: ['wss://relay.example.com'] }, + }, + } as unknown as JSONRPCMessage); + + // Allow async transition to settle + await new Promise((r) => setTimeout(r, 30)); + + expect(onRedirectCalled).toBe(true); + expect(redirectHop).toBe(1); + expect(baseTransport.closed).toBe(true); + expect(newTransportsCreated.length).toBe(1); + + const newTransport = newTransportsCreated[0] as FakeTransport; + expect(newTransport.serverPubkey).toBe(targetPubkey); + + // Check that the original request was re-issued over the new transport + expect(newTransport.sentMessages.length).toBe(1); + expect(newTransport.sentMessages[0]).toEqual(req); + }); + + test('surfaces error without redirecting if target pubkey is invalid', async () => { + const baseTransport = new FakeTransport('a'.repeat(64)); + const wrapped = withClientRedirect(baseTransport, { signer: dummySigner }); + + const received: JSONRPCMessage[] = []; + wrapped.onmessage = (msg) => received.push(msg); + await wrapped.start(); + + await wrapped.send({ + jsonrpc: '2.0', + id: 2, + method: 'tools/call', + }); + + baseTransport.emit({ + jsonrpc: '2.0', + id: 2, + error: { + code: REDIRECT_ERROR_CODE, + message: 'Redirect', + data: { target: 'invalid-hex' }, + }, + } as unknown as JSONRPCMessage); + + await new Promise((r) => setTimeout(r, 10)); + + expect(received.length).toBe(1); + expect((received[0] as { error?: { code: number } }).error?.code).toBe( + REDIRECT_ERROR_CODE, + ); + expect(baseTransport.closed).toBe(false); + }); + + test('redirectPolicy returning false rejects redirect and surfaces -32044 error', async () => { + const baseTransport = new FakeTransport('a'.repeat(64)); + const wrapped = withClientRedirect( + baseTransport, + { signer: dummySigner }, + { + redirectPolicy: async (data) => data.target !== 'c'.repeat(64), + }, + ); + + const received: JSONRPCMessage[] = []; + wrapped.onmessage = (msg) => received.push(msg); + await wrapped.start(); + + await wrapped.send({ + jsonrpc: '2.0', + id: 3, + method: 'tools/call', + }); + + baseTransport.emit({ + jsonrpc: '2.0', + id: 3, + error: { + code: REDIRECT_ERROR_CODE, + message: 'Redirect', + data: { target: 'c'.repeat(64) }, + }, + } as unknown as JSONRPCMessage); + + await new Promise((r) => setTimeout(r, 10)); + + expect(received.length).toBe(1); + expect((received[0] as { error?: { code: number } }).error?.code).toBe( + REDIRECT_ERROR_CODE, + ); + expect(baseTransport.closed).toBe(false); + }); + + test('exceeding maxRedirects stops redirect loop and surfaces error', async () => { + const baseTransport = new FakeTransport('a'.repeat(64)); + let nextTransport: FakeTransport | undefined; + + const wrapped = withClientRedirect( + baseTransport, + { + signer: dummySigner, + wrapTransport: () => { + nextTransport = new FakeTransport('d'.repeat(64)); + return nextTransport; + }, + }, + { maxRedirects: 1 }, + ); + + const received: JSONRPCMessage[] = []; + wrapped.onmessage = (msg) => received.push(msg); + await wrapped.start(); + + const req: JSONRPCRequest = { + jsonrpc: '2.0', + id: 4, + method: 'tools/call', + }; + await wrapped.send(req); + + // Hop 1: allowed (maxRedirects is 1) + baseTransport.emit({ + jsonrpc: '2.0', + id: 4, + error: { + code: REDIRECT_ERROR_CODE, + message: 'Redirect', + data: { target: 'd'.repeat(64) }, + }, + } as unknown as JSONRPCMessage); + + await new Promise((r) => setTimeout(r, 30)); + expect(baseTransport.closed).toBe(true); + expect(nextTransport).toBeDefined(); + + // Hop 2: emitted from nextTransport, exceeding maxRedirects = 1 + nextTransport!.emit({ + jsonrpc: '2.0', + id: 4, + error: { + code: REDIRECT_ERROR_CODE, + message: 'Redirect', + data: { target: 'e'.repeat(64) }, + }, + } as unknown as JSONRPCMessage); + + await new Promise((r) => setTimeout(r, 20)); + expect(received.length).toBe(1); + expect((received[0] as { error?: { code: number } }).error?.code).toBe( + REDIRECT_ERROR_CODE, + ); + expect(nextTransport!.closed).toBe(false); + }); + + test('forwards messages via onmessageWithContext when consumer attaches it', async () => { + const baseTransport = new FakeTransport('a'.repeat(64)); + const wrapped = withClientRedirect(baseTransport, { signer: dummySigner }); + + const receivedWithCtx: Array<{ msg: JSONRPCMessage; ctx: unknown }> = []; + (wrapped as { onmessageWithContext?: unknown }).onmessageWithContext = ( + msg: JSONRPCMessage, + ctx: unknown, + ) => receivedWithCtx.push({ msg, ctx }); + await wrapped.start(); + + baseTransport.emit( + { + jsonrpc: '2.0', + id: 5, + result: { ok: true }, + } as unknown as JSONRPCMessage, + { eventId: 'evt-5' }, + ); + + expect(receivedWithCtx.length).toBe(1); + expect(receivedWithCtx[0].ctx).toEqual({ eventId: 'evt-5' }); + }); +}); diff --git a/src/redirect/client-redirect.ts b/src/redirect/client-redirect.ts new file mode 100644 index 0000000..65c4d4b --- /dev/null +++ b/src/redirect/client-redirect.ts @@ -0,0 +1,376 @@ +import type { Transport } from '@contextvm/mcp-sdk/shared/transport'; +import { + isJSONRPCErrorResponse, + isJSONRPCNotification, + isJSONRPCRequest, + isJSONRPCResultResponse, + type JSONRPCErrorResponse, + type JSONRPCMessage, + type JSONRPCRequest, +} from '@contextvm/mcp-sdk/types.js'; + +import { NostrClientTransport } from '../transport/nostr-client-transport.js'; +import { ApplesauceRelayPool } from '../relay/applesauce-relay-pool.js'; +import { REDIRECT_ERROR_CODE } from '../payments/constants.js'; +import { LruCache } from '../core/utils/lru-cache.js'; +import { createLogger } from '../core/utils/logger.js'; +import type { + ClientRedirectOptions, + RedirectErrorData, + RedirectTransportConfig, +} from './types.js'; + +type TransportWithContext = Transport & { + onmessageWithContext?: ( + message: JSONRPCMessage, + ctx: { eventId: string; correlatedEventId?: string }, + ) => void; + serverPubkey?: string; +}; + +function supportsOnmessageWithContext( + transport: Transport, +): transport is TransportWithContext { + return Object.prototype.hasOwnProperty.call( + transport, + 'onmessageWithContext', + ); +} + +function isRedirectError(msg: JSONRPCMessage): msg is JSONRPCErrorResponse { + return ( + isJSONRPCErrorResponse(msg) && + msg.error.code === REDIRECT_ERROR_CODE && + msg.error.data != null && + typeof (msg.error.data as Record).target === 'string' + ); +} + +/** + * Wraps a client transport to automatically handle CEP-47 Server Redirect (-32044) responses. + * + * When the server returns `-32044 Redirect`, the wrapper: + * 1. Validates the target 64-character lowercase hex pubkey. + * 2. Checks hop limits against `maxRedirects` (default 5, scoped per original request ID). + * 3. Evaluates optional `redirectPolicy` hook. + * 4. Transparently creates a new NostrClientTransport to the redirected target and relays. + * 5. Re-applies optional decorators via `wrapTransport` (e.g., `withClientPayments`). + * 6. Swaps the active transport session and re-issues the original request. + * 7. Cleanly terminates the old transport session (abandoning pending payment/stream states). + * + * @param transport The base client transport to wrap. + * @param transportConfig Configuration for spawning new target transports upon redirect. + * @param options Client redirect handling rules and observability hooks. + * @returns A wrapped transport that handles redirection transparently. + */ +export function withClientRedirect( + transport: Transport, + transportConfig: RedirectTransportConfig, + options?: ClientRedirectOptions, +): Transport { + const logger = createLogger('client-redirect'); + const maxRedirects = options?.maxRedirects ?? 5; + const rawRequestCache = new LruCache(1000); + const redirectCounts = new LruCache(1000); + + let currentTransport = transport; + let currentServerPubkey: string | undefined = + transport instanceof NostrClientTransport + ? transport.serverPubkey + : (transport as unknown as { serverPubkey?: string }).serverPubkey; + let activeTransitionPromise: Promise | null = null; + + let onmessage: ((message: JSONRPCMessage) => void) | undefined; + let onmessageWithContext: + | (( + message: JSONRPCMessage, + ctx: { eventId: string; correlatedEventId?: string }, + ) => void) + | undefined; + let onerror: ((error: Error) => void) | undefined; + let onclose: (() => void) | undefined; + + const synthesizeError = ( + id: string | number | undefined, + code: number, + message: string, + data?: unknown, + ): void => { + const errObj: JSONRPCMessage = { + jsonrpc: '2.0', + id, + error: { code, message, data }, + } as unknown as JSONRPCMessage; + if (onmessageWithContext) { + onmessageWithContext(errObj, { eventId: 'synthetic' }); + } else { + onmessage?.(errObj); + } + }; + + const forwardMessage = ( + message: JSONRPCMessage, + ctx?: { eventId: string; correlatedEventId?: string }, + ): void => { + if (ctx && onmessageWithContext) { + onmessageWithContext(message, ctx); + } else { + onmessage?.(message); + } + }; + + const bindTransportHandlers = (targetTransport: Transport): void => { + const hasContextPath = supportsOnmessageWithContext(targetTransport); + + targetTransport.onmessage = (message: JSONRPCMessage) => { + if (hasContextPath && isJSONRPCNotification(message)) { + return; + } + if (isJSONRPCResultResponse(message) || isJSONRPCErrorResponse(message)) { + if ( + 'id' in message && + message.id != null && + !isRedirectError(message) + ) { + const reqId = message.id as string | number; + rawRequestCache.delete(String(reqId)); + redirectCounts.delete(String(reqId)); + } + } + if (hasContextPath) { + return; + } + void handleInbound(message, undefined); + }; + + if (hasContextPath) { + targetTransport.onmessageWithContext = (message, ctx) => { + void handleInbound(message, ctx); + }; + } + + targetTransport.onerror = (err: Error) => onerror?.(err); + targetTransport.onclose = () => { + if (currentTransport === targetTransport) { + onclose?.(); + } + }; + }; + + const performTransition = async ( + target: string, + relays?: string[], + ): Promise => { + // Note: Deviation from CEP-47 letter (which specifies uniform handling without special-casing target === currentServerPubkey). + // Safe because request hop counter bounds loops. + if (currentServerPubkey === target) { + return; + } + logger.info('Following server redirect to target', { + from: currentServerPubkey, + to: target, + relays, + }); + + // TODO: CEP-41 streams integration: release local stream state and surface failure to caller upon transport transition. + const oldTransport = currentTransport; + const { wrapTransport, ...baseOpts } = transportConfig; + + const newNostrTransport = new NostrClientTransport({ + ...baseOpts, + serverPubkey: target, + relayHandler: + relays && relays.length > 0 + ? new ApplesauceRelayPool(relays) + : undefined, + }); + + let newTransport: Transport = newNostrTransport; + if (wrapTransport) { + newTransport = wrapTransport(newTransport); + } + + bindTransportHandlers(newTransport); + await newTransport.start(); + + currentTransport = newTransport; + currentServerPubkey = target; + + try { + await oldTransport.close(); + } catch (err: unknown) { + logger.warn('Error closing old transport during redirect swap', { + error: err instanceof Error ? err.message : String(err), + }); + } + }; + + const handleInbound = async ( + message: JSONRPCMessage, + ctx?: { eventId: string; correlatedEventId?: string }, + ): Promise => { + if (!isRedirectError(message)) { + forwardMessage(message, ctx); + return; + } + + const errorData = message.error.data as unknown as RedirectErrorData; + const reqId = message.id as string | number; + + if (!errorData.target || !/^[0-9a-f]{64}$/.test(errorData.target)) { + logger.error('Invalid redirect target pubkey received', { + target: errorData.target, + requestId: reqId, + }); + forwardMessage(message, ctx); + return; + } + + const currentHops = (redirectCounts.get(String(reqId)) ?? 0) + 1; + if (currentHops > maxRedirects) { + logger.error('Maximum redirect hops exceeded for request', { + maxRedirects, + currentHops, + requestId: reqId, + }); + redirectCounts.delete(String(reqId)); + rawRequestCache.delete(String(reqId)); + forwardMessage(message, ctx); + return; + } + redirectCounts.set(String(reqId), currentHops); + + if (options?.redirectPolicy) { + let allowed = false; + try { + allowed = await options.redirectPolicy(errorData); + } catch (err: unknown) { + logger.error('Error in redirectPolicy hook, rejecting redirect', { + error: err instanceof Error ? err.message : String(err), + requestId: reqId, + }); + } + if (!allowed) { + redirectCounts.delete(String(reqId)); + rawRequestCache.delete(String(reqId)); + forwardMessage(message, ctx); + return; + } + } + + if (!activeTransitionPromise) { + activeTransitionPromise = performTransition( + errorData.target, + errorData.relays, + ).finally(() => { + activeTransitionPromise = null; + }); + } + + try { + await activeTransitionPromise; + } catch (err: unknown) { + logger.error('Failed to transition to redirect target transport', { + error: err instanceof Error ? err.message : String(err), + target: errorData.target, + requestId: reqId, + }); + redirectCounts.delete(String(reqId)); + rawRequestCache.delete(String(reqId)); + forwardMessage(message, ctx); + return; + } + + options?.onRedirect?.(errorData, currentHops); + + const origReq = rawRequestCache.get(String(reqId)); + if (!origReq) { + logger.error( + 'Cannot re-issue redirected request: original request not found in cache', + { + requestId: reqId, + }, + ); + redirectCounts.delete(String(reqId)); + forwardMessage(message, ctx); + return; + } + + logger.debug('Re-issuing request to redirected target', { + method: origReq.method, + requestId: reqId, + target: errorData.target, + hop: currentHops, + }); + + try { + if (activeTransitionPromise) { + await activeTransitionPromise; + } + await currentTransport.send(origReq); + } catch (err: unknown) { + logger.error('Error re-issuing request to redirected target', { + error: err instanceof Error ? err.message : String(err), + requestId: reqId, + }); + synthesizeError( + reqId, + -32000, + `Failed to re-issue request after redirect: ${ + err instanceof Error ? err.message : String(err) + }`, + ); + } + }; + + const wrapped: TransportWithContext = { + get onmessage() { + return onmessage; + }, + set onmessage(fn) { + onmessage = fn; + }, + get onmessageWithContext() { + return onmessageWithContext; + }, + set onmessageWithContext(fn) { + onmessageWithContext = fn; + }, + get onerror() { + return onerror; + }, + set onerror(fn) { + onerror = fn; + }, + get onclose() { + return onclose; + }, + set onclose(fn) { + onclose = fn; + }, + get serverPubkey() { + return currentServerPubkey; + }, + + async start(): Promise { + bindTransportHandlers(currentTransport); + await currentTransport.start(); + }, + + async send(message: JSONRPCMessage): Promise { + if (isJSONRPCRequest(message) && 'id' in message && message.id != null) { + rawRequestCache.set(String(message.id), message as JSONRPCRequest); + } + if (activeTransitionPromise) { + await activeTransitionPromise; + } + await currentTransport.send(message); + }, + + async close(): Promise { + await currentTransport.close(); + }, + }; + + return wrapped; +} diff --git a/src/redirect/index.ts b/src/redirect/index.ts new file mode 100644 index 0000000..01b8f77 --- /dev/null +++ b/src/redirect/index.ts @@ -0,0 +1,4 @@ +export * from './types.js'; +export * from './server-redirect.js'; +export * from './server-transport-redirect.js'; +export * from './client-redirect.js'; diff --git a/src/redirect/redirect-flow.test.ts b/src/redirect/redirect-flow.test.ts new file mode 100644 index 0000000..92381d5 --- /dev/null +++ b/src/redirect/redirect-flow.test.ts @@ -0,0 +1,227 @@ +import { + afterAll, + afterEach, + beforeAll, + describe, + expect, + test, +} from 'bun:test'; +import { z } from 'zod'; +import { McpServer } from '@contextvm/mcp-sdk/server/mcp'; +import { Client } from '@contextvm/mcp-sdk/client'; +import { generateSecretKey, getPublicKey } from 'nostr-tools/pure'; +import { bytesToHex } from 'nostr-tools/utils'; + +import { + spawnMockRelay, + clearRelayCache, +} from '../__mocks__/test-relay-helpers.js'; +import { ApplesauceRelayPool } from '../relay/applesauce-relay-pool.js'; +import { PrivateKeySigner } from '../signer/private-key-signer.js'; +import { EncryptionMode } from '../core/interfaces.js'; +import { NostrServerTransport } from '../transport/nostr-server-transport.js'; +import { NostrClientTransport } from '../transport/nostr-client-transport.js'; +import { withServerRedirect } from './server-transport-redirect.js'; +import { withClientRedirect } from './client-redirect.js'; +import { REDIRECT_ERROR_CODE } from '../payments/constants.js'; + +describe.serial('Redirect Flow E2E', () => { + let relayUrl: string; + let httpUrl: string; + let stopRelay: (() => void) | undefined; + + beforeAll(async () => { + const relay = await spawnMockRelay(); + relayUrl = relay.relayUrl; + httpUrl = relay.httpUrl; + stopRelay = relay.stop; + }); + + afterEach(async () => { + await clearRelayCache(httpUrl); + }); + + afterAll(() => { + stopRelay?.(); + }); + + const createServer = async ( + pubkeySK: Uint8Array, + redirectTarget?: string, + ) => { + const server = new McpServer({ + name: 'test-server', + version: '1.0.0', + }); + server.registerTool( + 'echo', + { + title: 'Echo', + description: 'Echoes message', + inputSchema: { message: z.string() }, + }, + async ({ message }: { message: string }) => ({ + content: [{ type: 'text', text: message }], + }), + ); + + let transport = new NostrServerTransport({ + signer: new PrivateKeySigner(bytesToHex(pubkeySK)), + relayHandler: new ApplesauceRelayPool([relayUrl]), + encryptionMode: EncryptionMode.DISABLED, + }); + + if (redirectTarget) { + transport = withServerRedirect(transport, { + resolveRedirect: () => ({ target: redirectTarget, relays: [relayUrl] }), + }); + } + + await server.connect(transport); + return { server, transport }; + }; + + test('End-to-end single redirect: Client -> Server A (redirects) -> Server B (responds)', async () => { + const skA = generateSecretKey(); + const skB = generateSecretKey(); + const pkB = getPublicKey(skB); + + const { transport: transportA } = await createServer(skA, pkB); + const { transport: transportB } = await createServer(skB); + + const clientSK = generateSecretKey(); + const baseOpts = { + signer: new PrivateKeySigner(bytesToHex(clientSK)), + encryptionMode: EncryptionMode.DISABLED, + isStateless: false, + }; + + const baseTransport = new NostrClientTransport({ + ...baseOpts, + serverPubkey: getPublicKey(skA), + relayHandler: new ApplesauceRelayPool([relayUrl]), + }); + + const clientTransport = withClientRedirect( + baseTransport, + { ...baseOpts, wrapTransport: (t) => t }, + { maxRedirects: 3 }, + ); + + const client = new Client( + { name: 'test-client', version: '1.0.0' }, + { capabilities: {} }, + ); + await client.connect(clientTransport); + + const result = await client.callTool({ + name: 'echo', + arguments: { message: 'hello redirect' }, + }); + + expect(result.content).toBeArray(); + expect((result.content as unknown[])?.[0]).toEqual({ + type: 'text', + text: 'hello redirect', + }); + + // Verify current transport server is now pkB + const activeTransport = clientTransport as unknown as NostrClientTransport; + expect(activeTransport.serverPubkey).toBe(pkB); + + await client.close(); + await transportA.close(); + await transportB.close(); + }); + + test('Chained redirect: Client -> Server A -> Server B -> Server C (responds)', async () => { + const skA = generateSecretKey(); + const skB = generateSecretKey(); + const skC = generateSecretKey(); + + const pkB = getPublicKey(skB); + const pkC = getPublicKey(skC); + + const { transport: transportA } = await createServer(skA, pkB); + const { transport: transportB } = await createServer(skB, pkC); + const { transport: transportC } = await createServer(skC); + + const clientSK = generateSecretKey(); + const baseOpts = { + signer: new PrivateKeySigner(bytesToHex(clientSK)), + encryptionMode: EncryptionMode.DISABLED, + }; + + const baseTransport = new NostrClientTransport({ + ...baseOpts, + serverPubkey: getPublicKey(skA), + relayHandler: new ApplesauceRelayPool([relayUrl]), + }); + + const clientTransport = withClientRedirect(baseTransport, baseOpts); + + const client = new Client( + { name: 'test-client', version: '1.0.0' }, + { capabilities: {} }, + ); + await client.connect(clientTransport); + + const result = await client.callTool({ + name: 'echo', + arguments: { message: 'chain test' }, + }); + + expect((result.content as unknown[])?.[0]).toEqual({ + type: 'text', + text: 'chain test', + }); + expect( + (clientTransport as unknown as NostrClientTransport).serverPubkey, + ).toBe(pkC); + + await client.close(); + await transportA.close(); + await transportB.close(); + await transportC.close(); + }); + + test('Loop detection: Client -> Server A -> Server B -> Server A -> (hops out) -> Error', async () => { + const skA = generateSecretKey(); + const skB = generateSecretKey(); + const pkA = getPublicKey(skA); + const pkB = getPublicKey(skB); + + const { transport: transportA } = await createServer(skA, pkB); + const { transport: transportB } = await createServer(skB, pkA); + + const clientSK = generateSecretKey(); + const baseOpts = { + signer: new PrivateKeySigner(bytesToHex(clientSK)), + encryptionMode: EncryptionMode.DISABLED, + }; + + const baseTransport = new NostrClientTransport({ + ...baseOpts, + serverPubkey: pkA, + relayHandler: new ApplesauceRelayPool([relayUrl]), + }); + + // Set a small hop cap + const clientTransport = withClientRedirect(baseTransport, baseOpts, { + maxRedirects: 2, + }); + + const client = new Client( + { name: 'test-client', version: '1.0.0' }, + { capabilities: {} }, + ); + + // It should hit max redirects on the initialize request and return the -32044 error + await expect(client.connect(clientTransport)).rejects.toMatchObject({ + code: REDIRECT_ERROR_CODE, + }); + + await transportA.close(); + await transportB.close(); + }); +}); diff --git a/src/redirect/server-redirect.test.ts b/src/redirect/server-redirect.test.ts new file mode 100644 index 0000000..06c3037 --- /dev/null +++ b/src/redirect/server-redirect.test.ts @@ -0,0 +1,223 @@ +import { describe, expect, test, mock } from 'bun:test'; +import type { + JSONRPCRequest, + JSONRPCNotification, +} from '@contextvm/mcp-sdk/types.js'; +import { createRedirectMiddleware } from './server-redirect.js'; +import { withServerRedirect } from './server-transport-redirect.js'; +import { REDIRECT_ERROR_CODE } from '../payments/constants.js'; +import type { ServerRedirectConfig } from './types.js'; + +describe('createRedirectMiddleware', () => { + const dummyCtx = { + clientPubkey: 'a'.repeat(64), + }; + + test('emits -32044 and halts when resolveRedirect returns a target', async () => { + let sentResponse: unknown = null; + let forwarded = false; + + const config: ServerRedirectConfig = { + resolveRedirect: async () => ({ + target: 'b'.repeat(64), + relays: ['wss://relay.example.com'], + instructions: 'Please reconnect to backend B', + _meta: { region: 'us-east' }, + }), + }; + + const middleware = createRedirectMiddleware({ + config, + // TODO: Assert `requestEventId` is correctly passed in the 3rd argument + sendResponse: async (clientPubkey, res) => { + sentResponse = { clientPubkey, res }; + }, + }); + + const req: JSONRPCRequest = { + jsonrpc: '2.0', + id: 'req-1', + method: 'tools/call', + params: { name: 'echo' }, + }; + + await middleware(req, dummyCtx, async () => { + forwarded = true; + }); + + expect(forwarded).toBe(false); + expect(sentResponse).toEqual({ + clientPubkey: dummyCtx.clientPubkey, + res: { + jsonrpc: '2.0', + id: 'req-1', + error: { + code: REDIRECT_ERROR_CODE, + message: 'Redirect', + data: { + target: 'b'.repeat(64), + relays: ['wss://relay.example.com'], + instructions: 'Please reconnect to backend B', + _meta: { region: 'us-east' }, + }, + }, + }, + }); + }); + + test('forwards normally without emitting error when resolveRedirect returns null', async () => { + let sentResponse: unknown = null; + let forwarded = false; + + const config: ServerRedirectConfig = { + resolveRedirect: async () => null, + }; + + const middleware = createRedirectMiddleware({ + config, + sendResponse: async (clientPubkey, res) => { + sentResponse = { clientPubkey, res }; + }, + }); + + const req: JSONRPCRequest = { + jsonrpc: '2.0', + id: 'req-2', + method: 'tools/list', + }; + + await middleware(req, dummyCtx, async () => { + forwarded = true; + }); + + expect(forwarded).toBe(true); + expect(sentResponse).toBeNull(); + }); + + test('forwards normally (fail-open) when resolveRedirect throws an error', async () => { + let sentResponse: unknown = null; + let forwarded = false; + + const config: ServerRedirectConfig = { + resolveRedirect: async () => { + throw new Error('Database connection failed'); + }, + }; + + const middleware = createRedirectMiddleware({ + config, + sendResponse: async (clientPubkey, res) => { + sentResponse = { clientPubkey, res }; + }, + }); + + const req: JSONRPCRequest = { + jsonrpc: '2.0', + id: 'req-3', + method: 'tools/call', + }; + + await middleware(req, dummyCtx, async () => { + forwarded = true; + }); + + expect(forwarded).toBe(true); + expect(sentResponse).toBeNull(); + }); + + test('passes through notifications without calling resolveRedirect', async () => { + const resolveMock = mock(async () => ({ target: 'b'.repeat(64) })); + let forwarded = false; + + const middleware = createRedirectMiddleware({ + config: { resolveRedirect: resolveMock }, + sendResponse: async () => {}, + }); + + const notif: JSONRPCNotification = { + jsonrpc: '2.0', + method: 'notifications/initialized', + }; + + await middleware(notif, dummyCtx, async () => { + forwarded = true; + }); + + expect(forwarded).toBe(true); + expect(resolveMock).toHaveBeenCalledTimes(0); + }); + + test('applies default instructions and _meta from config when callback result does not override them', async () => { + let sentData: unknown = null; + + const middleware = createRedirectMiddleware({ + config: { + instructions: 'Default instructions', + _meta: { tier: 'free' }, + resolveRedirect: async () => ({ + target: 'c'.repeat(64), + }), + }, + sendResponse: async (_, res) => { + sentData = res.error.data; + }, + }); + + await middleware( + { jsonrpc: '2.0', id: 1, method: 'tools/call' }, + dummyCtx, + async () => {}, + ); + + expect(sentData).toEqual({ + target: 'c'.repeat(64), + instructions: 'Default instructions', + _meta: { tier: 'free' }, + }); + }); + + test('per-request callback instructions and _meta override and merge with config defaults', async () => { + let sentData: unknown = null; + + const middleware = createRedirectMiddleware({ + config: { + instructions: 'Default instructions', + _meta: { tier: 'free', region: 'us' }, + resolveRedirect: async () => ({ + target: 'd'.repeat(64), + instructions: 'Override instructions', + _meta: { tier: 'pro', custom: true }, + }), + }, + sendResponse: async (_, res) => { + sentData = res.error.data; + }, + }); + + await middleware( + { jsonrpc: '2.0', id: 2, method: 'tools/call' }, + dummyCtx, + async () => {}, + ); + + expect(sentData).toEqual({ + target: 'd'.repeat(64), + instructions: 'Override instructions', + _meta: { tier: 'pro', region: 'us', custom: true }, + }); + }); +}); + +describe('withServerRedirect', () => { + test('registers inbound middleware on the transport', () => { + const addedMiddlewares: unknown[] = []; + const fakeTransport = { + addInboundMiddleware: (mw: unknown) => addedMiddlewares.push(mw), + sendTargetedResponse: async () => {}, + } as unknown as import('../transport/nostr-server-transport.js').NostrServerTransport; + + withServerRedirect(fakeTransport, { resolveRedirect: async () => null }); + expect(addedMiddlewares.length).toBe(1); + expect(typeof addedMiddlewares[0]).toBe('function'); + }); +}); diff --git a/src/redirect/server-redirect.ts b/src/redirect/server-redirect.ts new file mode 100644 index 0000000..99a015b --- /dev/null +++ b/src/redirect/server-redirect.ts @@ -0,0 +1,115 @@ +import { + isJSONRPCRequest, + type JSONRPCErrorResponse, +} from '@contextvm/mcp-sdk/types.js'; + +import type { InboundMiddlewareFn } from '../transport/middleware.js'; +import type { ServerRedirectConfig } from './types.js'; +import { REDIRECT_ERROR_CODE } from '../payments/constants.js'; +import { createLogger } from '../core/utils/logger.js'; + +/** + * Parameters for creating the server-side redirect middleware. + */ +export interface RedirectMiddlewareParams { + /** Server redirect configuration including the resolveRedirect callback. */ + config: ServerRedirectConfig; + /** + * Callback to emit a targeted JSON-RPC error response back to a client. + */ + sendResponse: ( + clientPubkey: string, + response: JSONRPCErrorResponse, + requestEventId: string, + ) => Promise; +} + +/** + * Creates an inbound server middleware that evaluates requests for CEP-47 redirection. + * + * If `config.resolveRedirect` returns a target, emits a `-32044 Redirect` JSON-RPC error + * response and short-circuits the middleware pipeline so the request is not processed further. + * If `resolveRedirect` returns `null` or throws an error (fail-open), forwards the request normally. + */ +export function createRedirectMiddleware( + params: RedirectMiddlewareParams, +): InboundMiddlewareFn { + const { config, sendResponse } = params; + const logger = createLogger('server-redirect'); + + return async (message, ctx, forward) => { + // TODO: CEP-41 streams integration: MUST NOT emit redirect if request has active CEP-41 open stream. + // Only redirect JSON-RPC requests. Notifications and responses pass through. + if (!isJSONRPCRequest(message) || message.id == null) { + await forward(message); + return; + } + + const requestEventId = ctx.requestEventId ?? String(message.id); + let result; + try { + result = await config.resolveRedirect({ + clientPubkey: ctx.clientPubkey, + method: message.method, + params: message.params as Record | undefined, + requestEventId, + }); + } catch (err: unknown) { + logger.error( + 'Error in resolveRedirect callback, forwarding request normally (fail-open)', + { + error: err instanceof Error ? err.message : String(err), + clientPubkey: ctx.clientPubkey, + method: message.method, + requestEventId, + }, + ); + await forward(message); + return; + } + + // null = do not redirect, forward normally to next middleware / handler + if (!result) { + await forward(message); + return; + } + + logger.debug('Redirecting client request', { + clientPubkey: ctx.clientPubkey, + method: message.method, + requestEventId, + target: result.target, + }); + + // Build -32044 error data payload + const errorData: Record = { target: result.target }; + if (result.relays && result.relays.length > 0) { + errorData.relays = result.relays; + } + const instructions = result.instructions ?? config.instructions; + if (instructions) { + errorData.instructions = instructions; + } + const meta = { + ...(config._meta ?? {}), + ...(result._meta ?? {}), + }; + if (Object.keys(meta).length > 0) { + errorData._meta = meta; + } + + const errorResponse: JSONRPCErrorResponse = { + jsonrpc: '2.0', + id: message.id, + error: { + code: REDIRECT_ERROR_CODE, + message: 'Redirect', + data: errorData, + }, + }; + + await sendResponse(ctx.clientPubkey, errorResponse, requestEventId); + + // Do NOT call forward() — request handling is complete. + }; +} diff --git a/src/redirect/server-transport-redirect.ts b/src/redirect/server-transport-redirect.ts new file mode 100644 index 0000000..48a4bb1 --- /dev/null +++ b/src/redirect/server-transport-redirect.ts @@ -0,0 +1,33 @@ +import type { NostrServerTransport } from '../transport/nostr-server-transport.js'; +import type { ServerRedirectConfig } from './types.js'; +import { createRedirectMiddleware } from './server-redirect.js'; + +/** + * Attaches CEP-47 server redirection middleware to a NostrServerTransport. + * + * When an inbound JSON-RPC request arrives, the middleware evaluates `config.resolveRedirect`. + * If a target is returned, the server responds with a `-32044 Redirect` error response + * and halts further processing of the request. + * + * @param transport The server transport to wrap. + * @param config The server redirection configuration. + * @returns The wrapped server transport with redirection middleware attached. + */ +export function withServerRedirect( + transport: NostrServerTransport, + config: ServerRedirectConfig, +): NostrServerTransport { + transport.addInboundMiddleware( + createRedirectMiddleware({ + config, + sendResponse: async (clientPubkey, response, requestEventId) => { + await transport.sendTargetedResponse( + clientPubkey, + response, + requestEventId, + ); + }, + }), + ); + return transport; +} diff --git a/src/redirect/types.ts b/src/redirect/types.ts new file mode 100644 index 0000000..403a162 --- /dev/null +++ b/src/redirect/types.ts @@ -0,0 +1,115 @@ +import type { Transport } from '@contextvm/mcp-sdk/shared/transport'; +import type { NostrTransportOptions } from '../transport/index.js'; + +/** + * Shape of `error.data` for JSON-RPC -32044 Redirect responses (CEP-47). + */ +export interface RedirectErrorData { + /** The 64-character lowercase hex public key of the target server. */ + target: string; + /** Optional relay hints where the target server is reachable. */ + relays?: string[]; + /** Human-readable instructions explaining the redirection. */ + instructions?: string; + /** Arbitrary metadata attached by the redirecting server. */ + _meta?: Record; +} + +/** + * Return structure from a `ResolveRedirectFn` callback. + */ +export interface RedirectTarget { + /** The 64-character lowercase hex public key of the target server. */ + target: string; + /** Optional relay hints where the target server is reachable. */ + relays?: string[]; + /** Overrides default server instructions when provided. */ + instructions?: string; + /** Overrides or merges with default server _meta when provided. */ + _meta?: Record; +} + +/** + * Context passed to `ResolveRedirectFn` when evaluating an inbound request. + */ +export interface RedirectContext { + /** The public key of the client issuing the request. */ + clientPubkey: string; + /** The JSON-RPC method being requested (e.g., 'tools/call'). */ + method: string; + /** The parameters accompanying the JSON-RPC request. */ + params?: Record; + /** The Nostr event ID corresponding to the inbound request message. */ + requestEventId: string; +} + +/** + * Server-side callback that resolves whether to redirect an inbound request. + * + * Follows the same pattern as `ResolvePriceFn` in CEP-8 payments: + * receives request context and returns a routing decision. + * + * Return `null` to forward the request normally without redirecting. + * Return a `RedirectTarget` to emit a -32044 error response. + */ +export type ResolveRedirectFn = ( + ctx: RedirectContext, +) => RedirectTarget | null | Promise; + +/** + * Server-side redirect middleware configuration. + */ +export interface ServerRedirectConfig { + /** Callback invoked per request to determine redirection targets. */ + resolveRedirect: ResolveRedirectFn; + /** Default instructions included in every redirect response unless overridden. */ + instructions?: string; + /** Default metadata included in every redirect response unless overridden. */ + _meta?: Record; +} + +/** + * Base transport options used when constructing a redirected NostrClientTransport. + * Excludes target-specific fields which are supplied dynamically by the redirect error data. + */ +export type BaseRedirectTransportOptions = Omit< + NostrTransportOptions, + | 'serverPubkey' + | 'relayHandler' + | 'discoveryRelayUrls' + | 'fallbackOperationalRelayUrls' +>; + +/** + * Configuration for constructing new client transports when following a redirect. + */ +export interface RedirectTransportConfig extends BaseRedirectTransportOptions { + /** + * Optional callback to wrap the newly constructed target transport with middleware + * (e.g., re-applying `withClientPayments` or other decorators). + */ + wrapTransport?: (transport: Transport) => Transport; +} + +/** + * Client-side options for CEP-47 redirect handling. + */ +export interface ClientRedirectOptions { + /** + * Maximum number of consecutive redirects allowed per original request ID. + * @default 5 + */ + maxRedirects?: number; + /** + * Policy hook evaluated before following a redirect. + * Return `false` to reject the redirection and surface the -32044 error to the caller. + */ + redirectPolicy?: (data: RedirectErrorData) => boolean | Promise; + /** + * Observability hook invoked whenever a redirect is successfully followed. + * + * @param data The redirect target and metadata received from the server. + * @param hopNumber The 1-indexed hop count for this request chain. + */ + onRedirect?: (data: RedirectErrorData, hopNumber: number) => void; +} diff --git a/src/relay/applesauce-relay-pool.test.ts b/src/relay/applesauce-relay-pool.test.ts index dbbb9c6..c33f7d2 100644 --- a/src/relay/applesauce-relay-pool.test.ts +++ b/src/relay/applesauce-relay-pool.test.ts @@ -174,6 +174,7 @@ describe('ApplesauceRelayPool Integration', () => { }); // 5. Publish the event + await sleep(100); await relayPool.publish(signedEvent); // 6. Wait for the event to be received @@ -527,7 +528,7 @@ describe('ApplesauceRelayPool Integration', () => { relayPool.subscribe([{ kinds: [1] }], () => {}); // 5. Wait for multiple liveness checks to trigger - await sleep(6000); + await sleep(8000); // 6. Assert multiple rebuilds happened expect(createRelayTracker.callCount).toBeGreaterThan(1); @@ -629,21 +630,12 @@ describe('ApplesauceRelayPool Integration', () => { // 4. Setup subscription before killing the relay const receivedEvents: NostrEvent[] = []; - const subscriptionPromise = new Promise((resolve, reject) => { - const timeout = setTimeout( - () => - reject( - new Error('Subscription timeout waiting for post-recovery event'), - ), - 15000, - ); - + const subscriptionPromise = new Promise((resolve) => { relayPool.subscribe( [{ kinds: [1], authors: [publicKeyHex] }], (event) => { receivedEvents.push(event); if (event.id === postRecoverySignedEvent.id) { - clearTimeout(timeout); resolve(); } }, @@ -668,7 +660,8 @@ describe('ApplesauceRelayPool Integration', () => { const restarted = await spawnMockRelayOnPort(relayPort); stopRelay = restarted.stop; - // 7. Publish event after relay recovery + // 7. Publish event after relay recovery (add small delay to ensure REQ is processed) + await sleep(100); await relayPool.publish(postRecoverySignedEvent); // 8. Wait for the event to be received via the restored subscription @@ -719,27 +712,47 @@ describe('ApplesauceRelayPool Integration', () => { relayGroup: RelayGroup; }; testPool.createSubscription = function (filters, onEvent, onEose) { - // Mirror production shape: subscribe to the raw req() message stream so - // the test override exercises the same dispatch path as the pool. No - // dedup (see production createSubscription for rationale). - const sub = testPool.relayGroup - .req(filters, { - reconnect: false, // Disable applesauce recovery - resubscribe: false, // Disable applesauce recovery - }) - .subscribe({ - next: (message) => { - if (message.type === 'EOSE') { + const stream = onEose + ? testPool.relayGroup.req(filters) + : testPool.relayGroup.subscription(filters, { + reconnect: Infinity, + resubscribe: Infinity, + eventStore: null, + }); + + const sub = ( + stream as unknown as { + subscribe: (observer: Record) => { + unsubscribe: () => void; + }; + } + ).subscribe({ + next: (message: unknown) => { + const msgObj = message as Record; + if (message === 'EOSE' || msgObj?.type === 'EOSE') { + onEose?.(); + return; + } + + if (Array.isArray(message)) { + if (message[0] === 'EOSE') { onEose?.(); - return; + } else if (message[0] === 'EVENT' && message[2]) { + onEvent(message[2] as NostrEvent); } + return; + } - if (message.type === 'EVENT') { - onEvent(message.event); + if (typeof message === 'object' && message !== null) { + if ('id' in msgObj) { + onEvent(msgObj as unknown as NostrEvent); + } else if (msgObj.type === 'EVENT' && msgObj.event) { + onEvent(msgObj.event as NostrEvent); } - }, - error: () => {}, - }); + } + }, + error: () => {}, + }); return () => sub.unsubscribe(); }; @@ -748,23 +761,12 @@ describe('ApplesauceRelayPool Integration', () => { // 4. Setup subscription and track events const receivedEvents: NostrEvent[] = []; - const subscriptionPromise = new Promise((resolve, reject) => { - const timeout = setTimeout( - () => - reject( - new Error( - 'Subscription timeout - rebuild did not restore subscription', - ), - ), - TIMING.SUBSCRIPTION_TIMEOUT, - ); - + const subscriptionPromise = new Promise((resolve) => { relayPool.subscribe( [{ kinds: [1], authors: [publicKeyHex] }], (event) => { receivedEvents.push(event); if (event.id === testEventId) { - clearTimeout(timeout); resolve(); } }, @@ -785,6 +787,7 @@ describe('ApplesauceRelayPool Integration', () => { const testSignedEvent = await signer.signEvent(testEvent); const testEventId = testSignedEvent.id; + await sleep(100); await relayPool.publish(testSignedEvent); // 7. Wait for event via restored subscription @@ -908,7 +911,8 @@ describe('ApplesauceRelayPool Integration', () => { ); }); - // 11. Publish the event + // 11. Publish the event (add small delay to ensure mock server processes REQ first) + await sleep(100); await relayPool.publish(testSignedEvent); // 12. Wait for event to be received via restored subscription diff --git a/src/relay/applesauce-relay-pool.ts b/src/relay/applesauce-relay-pool.ts index 9a9151a..8352dab 100644 --- a/src/relay/applesauce-relay-pool.ts +++ b/src/relay/applesauce-relay-pool.ts @@ -92,6 +92,8 @@ export class ApplesauceRelayPool implements RelayHandler { private pingSubscription?: Subscription; private readonly destroy$ = new Subject(); private rebuildInFlight?: Promise; + /** Tracks last known connection status per relay URL purely to deduplicate "Relay came online" log lines */ + private relayStates = new Map(); private relayObservers: Subscription[] = []; private relays: Relay[] = []; @@ -143,6 +145,16 @@ export class ApplesauceRelayPool implements RelayHandler { relayUrl: relay.url, connected, }); + + if (connected) { + const wasConnected = this.relayStates.get(relay.url) ?? false; + if (!wasConnected) { + logger.info('Relay came online', { relayUrl: relay.url }); + this.relayStates.set(relay.url, true); + } + } else { + this.relayStates.set(relay.url, false); + } }); const errorSub = relay.error$.subscribe((error) => { @@ -240,8 +252,11 @@ export class ApplesauceRelayPool implements RelayHandler { if (response.ok) { acceptedCount += 1; } else if ( - response.from === undefined || - connectedRelayUrls.has(response.from) + response.from !== undefined && + connectedRelayUrls.has(response.from) && + !response.message?.includes('object unsubscribed') && + !response.message?.toLowerCase().includes('timeout') && + !response.message?.includes('Connection error') ) { connectedFailureCount += 1; } @@ -301,7 +316,9 @@ export class ApplesauceRelayPool implements RelayHandler { } } - throw new Error('Failed to publish event'); + throw new Error( + `Failed to publish event. Responses: ${JSON.stringify(responses)}`, + ); } catch (error) { if ( error instanceof Error && @@ -343,41 +360,59 @@ export class ApplesauceRelayPool implements RelayHandler { ): () => void { logger.debug('Creating subscription', { filters }); - // NOTE: applesauce-relay 6.0.3 changed `RelayGroup.subscription()` to emit - // only deduplicated NostrEvents and no longer surfaces EOSE markers. - // - // We deliberately subscribe to the raw `RelayGroup.req()` message stream - // and forward EVERY event without deduplication. Dedup is intentionally NOT - // performed at the relay layer: + // Dedup is intentionally NOT performed at the relay layer: // - The transport layer already deduplicates gift-wrap envelopes and // decrypted inner events via its own `seenEventIds` cache, with // protocol-aware semantics. // - The explicit-gating payment flow republishes the SAME request event // id after payment and relies on the server re-observing it. Relay- - // layer dedup by event id (whether applesauce 6.0.3's `distinct()` or a - // local `Set`) silently swallows that retry and deadlocks the flow. - const sub = this.relayGroup - .req(filters, { reconnect: Infinity, resubscribe: Infinity }) - .subscribe({ - next: (message) => { - if (message.type === 'EOSE') { - onEose?.(); - } else if (message.type === 'EVENT') { - onEvent(message.event); - } - // OPEN / CLOSED / ERROR are intentionally ignored; the old - // `subscription()` filtered them out as well. - }, - complete: () => { - logger.debug('Subscription complete'); - }, - error: (error) => { - logger.error('Subscription error', { - error, - relayUrls: this.relayUrls, - }); - }, - }); + // layer dedup by event id silently swallows that retry and deadlocks + // the flow. + // + const stream = onEose + ? this.relayGroup.req(filters) + : this.relayGroup.subscription(filters, { + reconnect: Infinity, + resubscribe: Infinity, + eventStore: null, // intentionally disable relay-layer dedup + }); + + const sub = (stream as Observable).subscribe({ + next: (message: unknown) => { + logger.debug('Received raw message', { message }); + if ( + typeof message === 'object' && + message !== null && + 'type' in message + ) { + const msg = message as { type: string; event?: NostrEvent }; + if (msg.type === 'EOSE') onEose?.(); + else if (msg.type === 'EVENT' && msg.event) onEvent(msg.event); + return; + } + + if (Array.isArray(message)) { + if (message[0] === 'EOSE') onEose?.(); + else if (message[0] === 'EVENT' && message[2]) + onEvent(message[2] as NostrEvent); + return; + } + + if (message === 'EOSE') onEose?.(); + else if ( + typeof message === 'object' && + message !== null && + 'id' in message + ) + onEvent(message as NostrEvent); + }, + error: (error: unknown) => { + logger.warn('Subscription error', { filters, error }); + }, + complete: () => { + logger.debug('Subscription complete'); + }, + }); return () => sub.unsubscribe(); } @@ -421,10 +456,13 @@ export class ApplesauceRelayPool implements RelayHandler { }; } + private isDisconnected = false; + /** * Disconnects from all relays and cleans up resources. */ async disconnect(): Promise { + this.isDisconnected = true; this.destroy$.next(); this.destroy$.complete(); @@ -506,8 +544,10 @@ export class ApplesauceRelayPool implements RelayHandler { /** Starts the liveness ping monitor (called lazily on first subscribe) */ private startPingMonitor(): void { - if (this.pingSubscription) { - logger.debug('Ping monitor already started, skipping'); + if (this.pingSubscription || this.isDisconnected) { + logger.debug( + 'Ping monitor already started or pool disconnected, skipping', + ); return; } @@ -571,6 +611,8 @@ export class ApplesauceRelayPool implements RelayHandler { return; } + const currentGeneration = this.relayGeneration; + try { await Promise.all( connectedRelays.map(async (relay, index) => { @@ -601,6 +643,14 @@ export class ApplesauceRelayPool implements RelayHandler { }), ); } catch (error) { + if (this.relayGeneration !== currentGeneration) { + logger.debug('Ignoring liveness timeout because pool was rebuilt', { + generation: this.relayGeneration, + pingGeneration: currentGeneration, + }); + return; + } + if (error instanceof Error && error.name === 'TimeoutError') { logger.warn('Liveness check timed out - no response from relays', { pingTimeoutMs: this.pingTimeoutMs, @@ -636,7 +686,11 @@ export class ApplesauceRelayPool implements RelayHandler { // Stop current subscriptions (preserve descriptors for replay) this.stopActiveSubscriptions(); - // Create new relays and group + // Create new relays and group (if not disconnected during teardown) + if (this.isDisconnected) { + logger.debug('Rebuild aborted: pool disconnected'); + return; + } this.relays = this.relayUrls.map((url) => this.createRelay(url)); this.relayGroup = new RelayGroup(this.relays); diff --git a/src/transport/middleware.ts b/src/transport/middleware.ts index ba2ddca..2eef3b7 100644 --- a/src/transport/middleware.ts +++ b/src/transport/middleware.ts @@ -17,6 +17,7 @@ export type InboundMiddlewareFn = ( clientPubkey: string; clientPmis?: readonly string[]; paymentInteraction?: PaymentInteractionMode; + requestEventId?: string; }, forward: (message: JSONRPCMessage) => Promise, ) => Promise; diff --git a/src/transport/nostr-client-transport.ts b/src/transport/nostr-client-transport.ts index 840133a..c6d4717 100644 --- a/src/transport/nostr-client-transport.ts +++ b/src/transport/nostr-client-transport.ts @@ -138,7 +138,7 @@ export class NostrClientTransport public onerror?: (error: Error) => void; /** The server's public key for message targeting */ - private readonly serverPubkey: string; + public readonly serverPubkey: string; /** Manages request/response correlation for pending requests */ private readonly correlationStore: ClientCorrelationStore; /** Handles stateless mode emulation for public servers */ diff --git a/src/transport/nostr-client/relay-resolution.ts b/src/transport/nostr-client/relay-resolution.ts index 293a824..570c177 100644 --- a/src/transport/nostr-client/relay-resolution.ts +++ b/src/transport/nostr-client/relay-resolution.ts @@ -106,6 +106,9 @@ async function connectFallbackOperationalRelays( const relayPool = new ApplesauceRelayPool([...fallbackOperationalRelayUrls]); try { + // TODO(CEP-17, CEP-47): ApplesauceRelayPool.connect() currently resolves immediately + // even for dead ports. We should probe via Relay.connected$/status$ with a timeout + // to properly evaluate reachability and avoid falsely succeeding the fallback. await withTimeout( relayPool.connect(), DEFAULT_TIMEOUT_MS, diff --git a/src/transport/nostr-server/inbound-coordinator.ts b/src/transport/nostr-server/inbound-coordinator.ts index 12d31ca..3b55ac2 100644 --- a/src/transport/nostr-server/inbound-coordinator.ts +++ b/src/transport/nostr-server/inbound-coordinator.ts @@ -282,6 +282,7 @@ export class ServerInboundCoordinator { clientPmis: clientPmis.length > 0 ? clientPmis : undefined, paymentInteraction: session.effectivePaymentInteraction ?? 'transparent', + requestEventId: event.id, }; const middlewares = this.deps.inboundMiddlewares; diff --git a/src/transport/nostr-transport-reconnection.test.ts b/src/transport/nostr-transport-reconnection.test.ts index b6cdd5a..c2a1a81 100644 --- a/src/transport/nostr-transport-reconnection.test.ts +++ b/src/transport/nostr-transport-reconnection.test.ts @@ -39,6 +39,7 @@ describe.serial('NostrTransport Reconnection', () => { beforeAll(async () => { // Start primary relay on an OS-assigned port (avoids TOCTOU races under concurrency) const primaryRelay = await spawnMockRelay(); + await sleep(1000); primaryRelayInstance = primaryRelay.relay; stopPrimaryRelay = primaryRelay.stop; relayUrl = primaryRelay.relayUrl; @@ -205,6 +206,7 @@ describe.serial('NostrTransport Reconnection', () => { // Restart the relay process await restartPrimaryRelay(); + await sleep(1000); // Second request after relay restart - should still work const toolResult2 = await eventually(() => callAddTool(client, 10, 20), { @@ -230,6 +232,7 @@ describe.serial('NostrTransport Reconnection', () => { // First relay restart await restartPrimaryRelay(); + await sleep(1000); await sleep(300); // Request after first restart @@ -237,6 +240,7 @@ describe.serial('NostrTransport Reconnection', () => { // Second relay restart await restartPrimaryRelay(); + await sleep(1000); await sleep(300); // Request after second restart diff --git a/src/transport/redirect-payments-composition.e2e.test.ts b/src/transport/redirect-payments-composition.e2e.test.ts new file mode 100644 index 0000000..ccd1515 --- /dev/null +++ b/src/transport/redirect-payments-composition.e2e.test.ts @@ -0,0 +1,208 @@ +import { + afterAll, + afterEach, + beforeAll, + describe, + expect, + test, +} from 'bun:test'; +import { sleep } from 'bun'; +import { Client } from '@contextvm/mcp-sdk/client'; +import { McpServer } from '@contextvm/mcp-sdk/server/mcp'; +import { z } from 'zod'; +import { bytesToHex } from 'nostr-tools/utils'; +import { generateSecretKey, getPublicKey } from 'nostr-tools/pure'; +import { ApplesauceRelayPool } from '../relay/applesauce-relay-pool.js'; +import { PrivateKeySigner } from '../signer/private-key-signer.js'; +import { EncryptionMode } from '../core/interfaces.js'; +import { NostrServerTransport } from './nostr-server-transport.js'; +import { NostrClientTransport } from './nostr-client-transport.js'; +import { + FakePaymentProcessor, + withServerPayments, + withClientPayments, +} from '../payments/index.js'; +import { withServerRedirect, withClientRedirect } from '../redirect/index.js'; +import { PAYMENT_REQUIRED_ERROR_CODE } from '../payments/constants.js'; +import { + spawnMockRelay, + clearRelayCache, +} from '../__mocks__/test-relay-helpers.js'; + +/** + * Exercises the full redirect × payments composition end-to-end. + * + * Wrapping order mirrors the production wiring in `NostrMCPProxy`: + * client-side: redirect( payments( baseTransport ), wrapTransport = t => payments(t) ) + * server-side: payments( redirect( baseTransport ) ) + * + * This validates that: + * - The initial server's redirect middleware fires BEFORE its payment gating. + * - The client transparently follows the redirect and re-establishes a session + * with the target server. + * - The target server's payment gating (-32042) is correctly surfaced through + * both the redirect and payment layers to the Client. + */ +describe.serial('Redirect and Payments Composition', () => { + let relayUrl: string; + let httpUrl: string; + let stopRelay: (() => void) | undefined; + + beforeAll(async () => { + const relay = await spawnMockRelay(); + relayUrl = relay.relayUrl; + httpUrl = relay.httpUrl; + stopRelay = relay.stop; + }); + + afterEach(async () => { + await clearRelayCache(httpUrl); + }); + + afterAll(async () => { + stopRelay?.(); + await sleep(100); + }); + + test('redirect short-circuits payment gating on initial server, and payment succeeds on target server', async () => { + // 1. Target Server — prices `echo` at 1 'test' currency via `withServerPayments`. + const targetSK = generateSecretKey(); + const targetServer = new McpServer({ + name: 'target-paid-server', + version: '1.0.0', + }); + targetServer.registerTool( + 'echo', + { + title: 'Echo', + description: 'Echoes the message', + inputSchema: { message: z.string() }, + }, + async ({ message }: { message: string }) => ({ + content: [{ type: 'text', text: `Paid: ${message}` }], + }), + ); + const targetTransport = withServerPayments( + new NostrServerTransport({ + signer: new PrivateKeySigner(bytesToHex(targetSK)), + relayHandler: new ApplesauceRelayPool([relayUrl]), + encryptionMode: EncryptionMode.DISABLED, + }), + { + processors: [new FakePaymentProcessor()], + pricedCapabilities: [ + { + method: 'tools/call', + name: 'echo', + amount: 1, + currencyUnit: 'test', + }, + ], + paymentInteraction: 'optional', + }, + ); + await targetServer.connect(targetTransport); + const targetPubkey = getPublicKey(targetSK); + + // 2. Initial Server — redirects all requests to Target Server. + // Redirect middleware is wired BEFORE payments (correct server-side order) + // so the redirect fires before payment gating ever evaluates. + const initialSK = generateSecretKey(); + const initialServer = new McpServer({ + name: 'initial-redirect-server', + version: '1.0.0', + }); + const initialTransportBase = new NostrServerTransport({ + signer: new PrivateKeySigner(bytesToHex(initialSK)), + relayHandler: new ApplesauceRelayPool([relayUrl]), + encryptionMode: EncryptionMode.DISABLED, + }); + const initialTransportRedirect = withServerRedirect(initialTransportBase, { + resolveRedirect: async () => ({ + target: targetPubkey, + relays: [relayUrl], + }), + }); + // Even though the initial server has an expensive price, the redirect fires first + const initialTransport = withServerPayments(initialTransportRedirect, { + processors: [new FakePaymentProcessor()], + pricedCapabilities: [ + { + method: 'tools/call', + name: 'echo', + amount: 500, // Expensive — but redirect should fire before this evaluates + currencyUnit: 'test', + }, + ], + paymentInteraction: 'optional', + }); + await initialServer.connect(initialTransport); + const initialPubkey = getPublicKey(initialSK); + + // 3. Client — mirrors the NostrMCPProxy wrapping pattern: + // payments(base) → redirect(paymentsWrapped, wrapTransport = t => payments(t)) + // + // This ensures the INITIAL transport is already wrapped in payments, + // and any NEW transport created after redirect is also wrapped in payments + // via `wrapTransport`. + const clientSigner = new PrivateKeySigner( + bytesToHex(generateSecretKey()), + ); + const paymentOpts = { paymentInteraction: 'explicit_gating' as const }; + + const baseClientTransport = withClientPayments( + new NostrClientTransport({ + signer: clientSigner, + relayHandler: new ApplesauceRelayPool([relayUrl]), + serverPubkey: initialPubkey, + encryptionMode: EncryptionMode.DISABLED, + }), + paymentOpts, + ); + + // Redirect wraps the payments-wrapped transport. On redirect, `wrapTransport` + // re-applies payments to the newly spawned transport — matching the proxy. + const clientTransport = withClientRedirect( + baseClientTransport, + { + signer: clientSigner, + encryptionMode: EncryptionMode.DISABLED, + wrapTransport: (t) => withClientPayments(t, paymentOpts), + }, + { maxRedirects: 2 }, + ); + + const client = new Client( + { name: 'test-client', version: '1.0.0' }, + { capabilities: {} }, + ); + await client.connect(clientTransport as never); + + // Call the tool. The flow should be: + // Client → initial server → -32044 redirect → client follows → target server → -32042 payment required + // The -32042 from the TARGET server (amount=1) proves that: + // a) The initial server's redirect fired (not its amount=500 payment gating) + // b) The client followed the redirect and re-established a session + // c) The target server's payment gating is correctly surfaced + let caughtError: unknown; + try { + await client.callTool({ name: 'echo', arguments: { message: 'hello' } }); + } catch (err: unknown) { + caughtError = err; + } + expect(caughtError).toBeDefined(); + const mcpErr = caughtError as { code: number; data: unknown }; + expect(mcpErr.code).toBe(PAYMENT_REQUIRED_ERROR_CODE); + + // Verify the payment amount came from the TARGET server (1), not the initial server (500). + // This is the key composition assertion: the redirect fired before the initial server's + // expensive payment gating ever evaluated. + const dataStr = JSON.stringify(mcpErr.data); + expect(dataStr).toContain('"amount":1'); + expect(dataStr).not.toContain('"amount":500'); + + await client.close(); + await targetServer.close(); + await initialServer.close(); + }, 20000); +});