diff --git a/.changeset/resource-subscription-lifecycle.md b/.changeset/resource-subscription-lifecycle.md new file mode 100644 index 0000000..90a9e60 --- /dev/null +++ b/.changeset/resource-subscription-lifecycle.md @@ -0,0 +1,5 @@ +--- +'@contextvm/sdk': patch +--- + +Tighten the resource-subscription lifecycle on `NostrServerTransport`. `resources/subscribe`/`resources/unsubscribe` are now recorded only after the request is actually forwarded to the server, so middleware-dropped requests can no longer leave phantom subscription state behind (and a dropped unsubscribe can never re-add a subscription that never existed). A client's subscriptions are also cleared as soon as the authorization policy rejects it, so a revoked client stops receiving resource updates immediately instead of lingering until session eviction. diff --git a/.changeset/resource-update-subscriptions.md b/.changeset/resource-update-subscriptions.md new file mode 100644 index 0000000..52a9751 --- /dev/null +++ b/.changeset/resource-update-subscriptions.md @@ -0,0 +1,5 @@ +--- +'@contextvm/sdk': patch +--- + +Route `notifications/resources/updated` only to clients subscribed to the matching resource URI. Servers can configure matching for parent/sub-resource relationships. diff --git a/src/transport/nostr-server-transport.resource-updates.test.ts b/src/transport/nostr-server-transport.resource-updates.test.ts new file mode 100644 index 0000000..6ba823c --- /dev/null +++ b/src/transport/nostr-server-transport.resource-updates.test.ts @@ -0,0 +1,126 @@ +import { expect, test } from 'bun:test'; +import { Client } from '@contextvm/mcp-sdk/client'; +import { McpServer } from '@contextvm/mcp-sdk/server/mcp'; +import { + ResourceUpdatedNotificationSchema, + SubscribeRequestSchema, + UnsubscribeRequestSchema, +} from '@contextvm/mcp-sdk/types.js'; +import { generateSecretKey, getPublicKey } from 'nostr-tools/pure'; +import { bytesToHex, hexToBytes } from 'nostr-tools/utils'; +import { sleep } from '../core/utils/utils.js'; +import { PrivateKeySigner } from '../signer/private-key-signer.js'; +import { ApplesauceRelayPool } from '../relay/applesauce-relay-pool.js'; +import { spawnMockRelay } from '../__mocks__/test-relay-helpers.js'; +import { NostrClientTransport } from './nostr-client-transport.js'; +import { NostrServerTransport } from './nostr-server-transport.js'; + +test('routes resources/updated only to clients subscribed to that resource', async () => { + const relay = await spawnMockRelay(); + const serverPrivateKey = bytesToHex(generateSecretKey()); + const serverPublicKey = getPublicKey(hexToBytes(serverPrivateKey)); + const alphaPrivateKey = bytesToHex(generateSecretKey()); + const betaPrivateKey = bytesToHex(generateSecretKey()); + + const server = new McpServer( + { name: 'Resource update server', version: '1.0.0' }, + { capabilities: { resources: { subscribe: true } } }, + ); + server.server.setRequestHandler(SubscribeRequestSchema, async () => ({})); + server.server.setRequestHandler(UnsubscribeRequestSchema, async () => ({})); + + const serverTransport = new NostrServerTransport({ + signer: new PrivateKeySigner(serverPrivateKey), + relayHandler: new ApplesauceRelayPool([relay.relayUrl]), + matchesSubResource: (subscribedUri, updatedUri) => + updatedUri.startsWith(`${subscribedUri}/`), + }); + const alphaClient = new Client({ name: 'Alpha client', version: '1.0.0' }); + const betaClient = new Client({ name: 'Beta client', version: '1.0.0' }); + const alphaUpdates: string[] = []; + const betaUpdates: string[] = []; + + alphaClient.setNotificationHandler( + ResourceUpdatedNotificationSchema, + (notification) => { + alphaUpdates.push(notification.params.uri); + }, + ); + betaClient.setNotificationHandler( + ResourceUpdatedNotificationSchema, + (notification) => { + betaUpdates.push(notification.params.uri); + }, + ); + + try { + await server.connect(serverTransport); + await alphaClient.connect( + new NostrClientTransport({ + signer: new PrivateKeySigner(alphaPrivateKey), + relayHandler: new ApplesauceRelayPool([relay.relayUrl]), + serverPubkey: serverPublicKey, + }), + ); + await betaClient.connect( + new NostrClientTransport({ + signer: new PrivateKeySigner(betaPrivateKey), + relayHandler: new ApplesauceRelayPool([relay.relayUrl]), + serverPubkey: serverPublicKey, + }), + ); + + await alphaClient.subscribeResource({ uri: 'resource://alpha' }); + await betaClient.subscribeResource({ uri: 'resource://beta' }); + + await server.server.sendResourceUpdated({ uri: 'resource://alpha' }); + await sleep(150); + + expect(alphaUpdates).toEqual(['resource://alpha']); + expect(betaUpdates).toEqual([]); + + await server.server.sendResourceUpdated({ uri: 'resource://alpha/child' }); + await sleep(150); + + expect(alphaUpdates).toEqual([ + 'resource://alpha', + 'resource://alpha/child', + ]); + expect(betaUpdates).toEqual([]); + + await server.server.sendResourceUpdated({ + uri: 'resource://alphabet/child', + }); + await sleep(150); + + expect(alphaUpdates).toEqual([ + 'resource://alpha', + 'resource://alpha/child', + ]); + expect(betaUpdates).toEqual([]); + + await alphaClient.unsubscribeResource({ uri: 'resource://alpha' }); + await server.server.sendResourceUpdated({ uri: 'resource://alpha' }); + await sleep(150); + + expect(alphaUpdates).toEqual([ + 'resource://alpha', + 'resource://alpha/child', + ]); + expect(betaUpdates).toEqual([]); + + await server.server.sendResourceUpdated({ uri: 'resource://alpha/child' }); + await sleep(150); + + expect(alphaUpdates).toEqual([ + 'resource://alpha', + 'resource://alpha/child', + ]); + expect(betaUpdates).toEqual([]); + } finally { + await alphaClient.close().catch(() => undefined); + await betaClient.close().catch(() => undefined); + await server.close().catch(() => undefined); + relay.stop(); + } +}); diff --git a/src/transport/nostr-server-transport.ts b/src/transport/nostr-server-transport.ts index abb6542..4e53950 100644 --- a/src/transport/nostr-server-transport.ts +++ b/src/transport/nostr-server-transport.ts @@ -22,6 +22,10 @@ import { NostrEvent } from 'nostr-tools'; import { LogLevel } from '../core/utils/logger.js'; import { withTimeout } from '../core/utils/utils.js'; import { CorrelationStore } from './nostr-server/correlation-store.js'; +import { + SubscriptionStore, + type ResourceSubscriptionMatcher, +} from './nostr-server/subscription-store.js'; import { ClientSession, SessionStore } from './nostr-server/session-store.js'; import { LruCache } from '../core/utils/lru-cache.js'; import { ApplesauceRelayPool } from '../relay/applesauce-relay-pool.js'; @@ -58,6 +62,7 @@ import type { InboundMiddlewareFn } from './middleware.js'; import type { PaymentInteractionPolicy } from '../payments/types.js'; export type { InboundMiddlewareFn } from './middleware.js'; +export type { ResourceSubscriptionMatcher } from './nostr-server/subscription-store.js'; /** * Options for configuring the NostrServerTransport. */ @@ -95,6 +100,12 @@ export interface NostrServerTransportOptions extends BaseNostrTransportOptions { logLevel?: LogLevel; /** Maximum number of client sessions to keep in memory. @default 1000 */ maxSessions?: number; + /** + * Resolve server-defined parent/sub-resource relationships for resource updates. + * Exact URI subscriptions always match, even when this callback is omitted. + * The callback is only called for different subscribed and updated URIs. + */ + matchesSubResource?: ResourceSubscriptionMatcher; /** * Whether to inject the client's public key into the _meta field of incoming messages. * @default false @@ -185,6 +196,7 @@ export class NostrServerTransport private readonly sessionStore: SessionStore; private readonly correlationStore: CorrelationStore; + private readonly subscriptionStore: SubscriptionStore; private readonly authorizationPolicy: AuthorizationPolicy; private readonly announcementManager: AnnouncementManager; private readonly injectClientPubkey: boolean; @@ -246,6 +258,8 @@ export class NostrServerTransport isAnnouncedServer: options.isAnnouncedServer ?? options.isPublicServer, }); + this.subscriptionStore = new SubscriptionStore(); + // Initialize session store with eviction callback for correlation cleanup this.sessionStore = new SessionStore({ maxSessions: options.maxSessions ?? 1000, @@ -257,6 +271,7 @@ export class NostrServerTransport // callback would also corrupt the cache's capacity accounting.) const removedCount = this.correlationStore.removeRoutesForClient(clientPubkey); + this.subscriptionStore.removeForClient(clientPubkey); this.logger.info( `Evicted session for ${clientPubkey} (removed ${removedCount} routes)`, ); @@ -350,6 +365,7 @@ export class NostrServerTransport await this.outboundResponseRouter.route(response); }, sessionStore: this.sessionStore, + subscriptionStore: this.subscriptionStore, onClientSessionEvicted: this.onClientSessionEvicted, correlationStore: this.correlationStore, policy: options.openStream?.policy, @@ -359,6 +375,7 @@ export class NostrServerTransport this.inboundCoordinator = new ServerInboundCoordinator({ sessionStore: this.sessionStore, correlationStore: this.correlationStore, + subscriptionStore: this.subscriptionStore, authorizationPolicy: this.authorizationPolicy, openStreamFactory: this.openStreamFactory, inboundMiddlewares: this.inboundMiddlewares, @@ -442,6 +459,8 @@ export class NostrServerTransport this.outboundNotificationBroadcaster = new OutboundNotificationBroadcaster({ correlationStore: this.correlationStore, sessionStore: this.sessionStore, + subscriptionStore: this.subscriptionStore, + matchesSubResource: options.matchesSubResource, sendNotification: this.sendNotification.bind(this), enqueueTask: this.taskQueue.add.bind(this.taskQueue), logger: this.logger, @@ -586,6 +605,7 @@ export class NostrServerTransport await this.disconnect(); this.sessionStore.clear(); this.correlationStore.clear(); + this.subscriptionStore.clear(); this.seenEventIds.clear(); this.oversizedReceiver.clear(); this.openStreamFactory.getReceiver().clear(); diff --git a/src/transport/nostr-server/inbound-coordinator.test.ts b/src/transport/nostr-server/inbound-coordinator.test.ts new file mode 100644 index 0000000..fd814a8 --- /dev/null +++ b/src/transport/nostr-server/inbound-coordinator.test.ts @@ -0,0 +1,187 @@ +import { describe, expect, test } from 'bun:test'; +import type { JSONRPCRequest } from '@contextvm/mcp-sdk/types.js'; +import { isJSONRPCRequest } from '@contextvm/mcp-sdk/types.js'; +import type { NostrEvent } from 'nostr-tools'; +import type { Logger } from '../../core/utils/logger.js'; +import { GiftWrapMode } from '../../core/interfaces.js'; +import { sleep } from '../../core/utils/utils.js'; +import type { InboundMiddlewareFn } from '../middleware.js'; +import { AuthorizationPolicy } from './authorization-policy.js'; +import { CorrelationStore } from './correlation-store.js'; +import { + ServerInboundCoordinator, + type ServerInboundCoordinatorDeps, +} from './inbound-coordinator.js'; +import type { ServerOpenStreamFactory } from './open-stream-factory.js'; +import { SessionStore } from './session-store.js'; +import { SubscriptionStore } from './subscription-store.js'; + +const testLogger: Logger = { + debug: () => undefined, + info: () => undefined, + warn: () => undefined, + error: () => undefined, + withModule: () => testLogger, +}; + +const clientPubkey = 'client-pk'; +const resourceUri = 'resource://alpha'; + +function requestEvent(id: string, pubkey = clientPubkey): NostrEvent { + return { + id, + pubkey, + created_at: Math.floor(Date.now() / 1000), + kind: 1, + tags: [], + content: '', + sig: 'c'.repeat(128), + } as NostrEvent; +} + +function subscriptionRequest( + id: string, + method: string, + uri: string, +): JSONRPCRequest { + return { jsonrpc: '2.0', id, method, params: { uri } }; +} + +function createCoordinator(options?: { + authorizationPolicy?: AuthorizationPolicy; + inboundMiddlewares?: InboundMiddlewareFn[]; +}): { + coordinator: ServerInboundCoordinator; + subscriptions: SubscriptionStore; + correlationStore: CorrelationStore; +} { + const sessionStore = new SessionStore(); + const correlationStore = new CorrelationStore(); + const subscriptions = new SubscriptionStore(); + const deps: ServerInboundCoordinatorDeps = { + sessionStore, + correlationStore, + subscriptionStore: subscriptions, + authorizationPolicy: + options?.authorizationPolicy ?? new AuthorizationPolicy(), + openStreamFactory: { + createWriterIfEnabled: () => undefined, + inputStreamIfEnabled: () => undefined, + releaseUnusedWriter: () => undefined, + } as unknown as ServerOpenStreamFactory, + inboundMiddlewares: options?.inboundMiddlewares ?? [], + injectClientPubkey: false, + shouldInjectRequestEventId: false, + oversizedEnabled: false, + openStreamEnabled: false, + giftWrapMode: GiftWrapMode.OPTIONAL, + sendMcpMessage: async () => 'event-id', + createResponseTags: () => [], + getOrCreateClientSession: (pk, isEncrypted) => + sessionStore.getOrCreateSession(pk, isEncrypted)[0], + forwardMessage: async () => true, + logger: testLogger, + }; + return { + coordinator: new ServerInboundCoordinator(deps), + subscriptions, + correlationStore, + }; +} + +describe('ServerInboundCoordinator resource subscriptions', () => { + test('records a subscription when subscribe is forwarded and removes it when unsubscribe is forwarded', async () => { + const { coordinator, subscriptions } = createCoordinator(); + + await coordinator.authorizeAndProcessEvent( + requestEvent('a'.repeat(64)), + false, + subscriptionRequest('req-1', 'resources/subscribe', resourceUri), + ); + await sleep(0); + + expect(subscriptions.getSubscribers(resourceUri)).toEqual( + new Set([clientPubkey]), + ); + + await coordinator.authorizeAndProcessEvent( + requestEvent('b'.repeat(64)), + false, + subscriptionRequest('req-2', 'resources/unsubscribe', resourceUri), + ); + await sleep(0); + + expect(subscriptions.getSubscribers(resourceUri)).toEqual(new Set()); + }); + + test('does not record a subscription when a middleware drops the request', async () => { + const dropAll: InboundMiddlewareFn = async () => undefined; + const { coordinator, subscriptions, correlationStore } = createCoordinator({ + inboundMiddlewares: [dropAll], + }); + + await coordinator.authorizeAndProcessEvent( + requestEvent('a'.repeat(64)), + false, + subscriptionRequest('req-1', 'resources/subscribe', resourceUri), + ); + await sleep(0); + + expect(subscriptions.getSubscribers(resourceUri)).toEqual(new Set()); + expect(correlationStore.getEventRoute('a'.repeat(64))).toBeUndefined(); + }); + + test('keeps an existing subscription when a dropped unsubscribe never reaches the server', async () => { + const dropUnsubscribe: InboundMiddlewareFn = async ( + message, + _ctx, + forward, + ) => { + if ( + isJSONRPCRequest(message) && + message.method === 'resources/unsubscribe' + ) { + return; + } + await forward(message); + }; + const { coordinator, subscriptions } = createCoordinator({ + inboundMiddlewares: [dropUnsubscribe], + }); + + await coordinator.authorizeAndProcessEvent( + requestEvent('a'.repeat(64)), + false, + subscriptionRequest('req-1', 'resources/subscribe', resourceUri), + ); + await sleep(0); + await coordinator.authorizeAndProcessEvent( + requestEvent('b'.repeat(64)), + false, + subscriptionRequest('req-2', 'resources/unsubscribe', resourceUri), + ); + await sleep(0); + + expect(subscriptions.getSubscribers(resourceUri)).toEqual( + new Set([clientPubkey]), + ); + }); + + test("clears a client's subscriptions when authorization rejects it", async () => { + const { coordinator, subscriptions } = createCoordinator({ + authorizationPolicy: new AuthorizationPolicy({ + allowedPublicKeys: new Set(['other-pk']), + }), + }); + subscriptions.subscribe(clientPubkey, resourceUri); + + await coordinator.authorizeAndProcessEvent( + requestEvent('a'.repeat(64)), + false, + subscriptionRequest('req-1', 'resources/read', resourceUri), + ); + await sleep(0); + + expect(subscriptions.getSubscribers(resourceUri)).toEqual(new Set()); + }); +}); diff --git a/src/transport/nostr-server/inbound-coordinator.ts b/src/transport/nostr-server/inbound-coordinator.ts index 183aa60..2b56141 100644 --- a/src/transport/nostr-server/inbound-coordinator.ts +++ b/src/transport/nostr-server/inbound-coordinator.ts @@ -10,6 +10,7 @@ import { type Logger } from '../../core/utils/logger.js'; import { type SessionStore, type ClientSession } from './session-store.js'; import { NOSTR_TAGS } from '../../core/constants.js'; import { type CorrelationStore } from './correlation-store.js'; +import { type SubscriptionStore } from './subscription-store.js'; import { type AuthorizationPolicy } from './authorization-policy.js'; import { type ServerOpenStreamFactory } from './open-stream-factory.js'; import { type InboundNotificationDispatcher } from './inbound-notification-dispatcher.js'; @@ -38,6 +39,7 @@ import type { export interface ServerInboundCoordinatorDeps { sessionStore: SessionStore; correlationStore: CorrelationStore; + subscriptionStore: SubscriptionStore; authorizationPolicy: AuthorizationPolicy; openStreamFactory: ServerOpenStreamFactory; inboundMiddlewares: InboundMiddlewareFn[]; @@ -109,6 +111,9 @@ export class ServerInboundCoordinator { ); if (!authDecision.allowed) { + // A denied pubkey is a revocation: drop its subscriptions now instead + // of letting it keep receiving resource updates until session eviction. + this.deps.subscriptionStore.removeForClient(event.pubkey); this.deps.logger.error( `Unauthorized message from ${event.pubkey}, message: ${JSON.stringify(mcpMessage)}. Ignoring.`, ); @@ -288,7 +293,11 @@ export class ServerInboundCoordinator { ): Promise => { const mw = middlewares[index]; if (!mw) { - return await this.deps.forwardMessage(msg, event.pubkey); + const forwarded = await this.deps.forwardMessage(msg, event.pubkey); + if (forwarded) { + this.trackSubscription(msg, event.pubkey); + } + return forwarded; } let forwarded = false; await mw(msg, ctx, async (nextMsg) => { @@ -409,13 +418,38 @@ export class ServerInboundCoordinator { (meta as { stream?: OpenStreamWriter }).stream = openStreamWriter; } if (inputStream) { - (meta as { - inputStream?: AsyncIterable<{ value: string; chunkIndex: number }>; - }).inputStream = inputStream; + ( + meta as { + inputStream?: AsyncIterable<{ value: string; chunkIndex: number }>; + } + ).inputStream = inputStream; } } } + /** + * Records a resource subscription for a request that actually reached the + * server. Requests swallowed by middleware never get here, so a dropped + * subscribe/unsubscribe cannot leave phantom subscription state behind. + */ + private trackSubscription( + message: JSONRPCMessage, + clientPubkey: string, + ): void { + if (!isJSONRPCRequest(message)) { + return; + } + const resourceUri = message.params?.uri; + if (typeof resourceUri !== 'string') { + return; + } + if (message.method === 'resources/subscribe') { + this.deps.subscriptionStore.subscribe(clientPubkey, resourceUri); + } else if (message.method === 'resources/unsubscribe') { + this.deps.subscriptionStore.unsubscribe(clientPubkey, resourceUri); + } + } + /** * Cleans up per-request state reserved for a request that was dropped by * middleware or failed the chain. Single owner: releases the correlation diff --git a/src/transport/nostr-server/open-stream-factory.test.ts b/src/transport/nostr-server/open-stream-factory.test.ts index 7092e71..f21ac54 100644 --- a/src/transport/nostr-server/open-stream-factory.test.ts +++ b/src/transport/nostr-server/open-stream-factory.test.ts @@ -4,8 +4,11 @@ import type { JSONRPCResponse, } from '@contextvm/mcp-sdk/types.js'; import type { Logger } from '../../core/utils/logger.js'; -import type { CorrelationStore } from './correlation-store.js'; -import type { ClientSession, SessionStore } from './session-store.js'; +import { CorrelationStore as CorrelationStoreImpl } from './correlation-store.js'; +import { type CorrelationStore } from './correlation-store.js'; +import { type ClientSession, SessionStore } from './session-store.js'; +import { SubscriptionStore } from './subscription-store.js'; +import { OutboundNotificationBroadcaster } from './outbound-notification-broadcaster.js'; import { ServerOpenStreamFactory } from './open-stream-factory.js'; /** Polls `condition` until it returns true or `timeoutMs` elapses. */ @@ -40,6 +43,7 @@ const sessionStore = { getSession: () => undefined, removeSession: () => false, } as unknown as SessionStore; +const subscriptionStore = new SubscriptionStore(); function createFactory(options?: { openStreamEnabled?: boolean }): { factory: ServerOpenStreamFactory; @@ -62,6 +66,7 @@ function createFactory(options?: { openStreamEnabled?: boolean }): { handleResponse: async (response) => { routedResponses.push(response); }, + subscriptionStore, logger: testLogger, }); @@ -176,6 +181,7 @@ describe('ServerOpenStreamFactory.releaseUnusedWriter', () => { notifications.push({ clientPubkey, notification }); }, handleResponse: async () => undefined, + subscriptionStore: new SubscriptionStore(), policy: { idleTimeoutMs: 10, probeTimeoutMs: 10 }, logger: testLogger, }); @@ -237,33 +243,24 @@ describe('ServerOpenStreamFactory.getOpenStreams', () => { expect(factory.getOpenStreams()).toHaveLength(0); }); - test('keepalive probe timeout evicts the session and removes the writer without a manual abort', async () => { - const session: ClientSession = { - isInitialized: true, - isEncrypted: false, - hasSentCommonTags: true, - supportsEncryption: false, - supportsEphemeralEncryption: false, - supportsOversizedTransfer: false, - supportsOpenStream: true, - }; - const sessions = new Map([['pk-1', session]]); + test('probe timeout clears subscriptions until the client subscribes again', async () => { + const clientPubkey = 'pk-1'; + const resourceUri = 'resource://alpha'; const evicted: string[] = []; - const localSessionStore = { - getSession: (pk: string): ClientSession | undefined => sessions.get(pk), - removeSession: (pk: string): boolean => { - const had = sessions.has(pk); - sessions.delete(pk); - return had; - }, - } as unknown as SessionStore; + const localSessionStore = new SessionStore(); + localSessionStore.getOrCreateSession(clientPubkey, false); + localSessionStore.markInitialized(clientPubkey); + const localSubscriptionStore = new SubscriptionStore(); + const localCorrelationStore = new CorrelationStoreImpl(); + localSubscriptionStore.subscribe(clientPubkey, resourceUri); const factory = new ServerOpenStreamFactory({ openStreamEnabled: true, sessionStore: localSessionStore, - correlationStore, + correlationStore: localCorrelationStore, sendNotification: async () => undefined, handleResponse: async () => undefined, + subscriptionStore: localSubscriptionStore, onClientSessionEvicted: async ({ clientPubkey }): Promise => { evicted.push(clientPubkey); }, @@ -271,7 +268,11 @@ describe('ServerOpenStreamFactory.getOpenStreams', () => { logger: testLogger, }); - const writer = factory.createWriterIfEnabled('evt-pt', 'pk-1', 'token-pt'); + const writer = factory.createWriterIfEnabled( + 'evt-pt', + clientPubkey, + 'token-pt', + ); await writer!.start(); // No manual abort, no ackProbe: the writer's own keepalive must drive the @@ -279,10 +280,47 @@ describe('ServerOpenStreamFactory.getOpenStreams', () => { await waitFor(() => !writer!.isActive); expect(writer!.isActive).toBe(false); - expect(sessions.has('pk-1')).toBe(false); - expect(evicted).toEqual(['pk-1']); + expect(localSessionStore.hasSession(clientPubkey)).toBe(false); + expect(localSubscriptionStore.getSubscribers(resourceUri)).toEqual( + new Set(), + ); + expect(evicted).toEqual([clientPubkey]); expect(factory.getWriter('evt-pt')).toBeUndefined(); expect(factory.getOpenStreams()).toHaveLength(0); + + localSessionStore.getOrCreateSession(clientPubkey, false); + localSessionStore.markInitialized(clientPubkey); + const notifications: string[] = []; + const tasks: Promise[] = []; + const broadcaster = new OutboundNotificationBroadcaster({ + correlationStore: localCorrelationStore, + sessionStore: localSessionStore, + subscriptionStore: localSubscriptionStore, + sendNotification: async (recipient) => { + notifications.push(recipient); + }, + enqueueTask: (task) => { + tasks.push(task()); + }, + logger: testLogger, + }); + + await broadcaster.broadcast({ + jsonrpc: '2.0', + method: 'notifications/resources/updated', + params: { uri: resourceUri }, + }); + await Promise.all(tasks); + expect(notifications).toEqual([]); + + localSubscriptionStore.subscribe(clientPubkey, resourceUri); + await broadcaster.broadcast({ + jsonrpc: '2.0', + method: 'notifications/resources/updated', + params: { uri: resourceUri }, + }); + await Promise.all(tasks); + expect(notifications).toEqual([clientPubkey]); }); test('createWriterIfEnabled reuses the existing writer for the same event id', () => { diff --git a/src/transport/nostr-server/open-stream-factory.ts b/src/transport/nostr-server/open-stream-factory.ts index bb0b47c..6d51592 100644 --- a/src/transport/nostr-server/open-stream-factory.ts +++ b/src/transport/nostr-server/open-stream-factory.ts @@ -17,6 +17,7 @@ import { type OpenStreamRegistryOptions } from '../open-stream/registry.js'; import { type Logger } from '../../core/utils/logger.js'; import { type CorrelationStore } from './correlation-store.js'; import { type ClientSession, type SessionStore } from './session-store.js'; +import { type SubscriptionStore } from './subscription-store.js'; import { type JSONRPCMessage, type JSONRPCResponse, @@ -33,6 +34,7 @@ export interface ServerOpenStreamFactoryDeps { ) => Promise; handleResponse: (response: JSONRPCResponse) => Promise; sessionStore: SessionStore; + subscriptionStore: SubscriptionStore; onClientSessionEvicted?: (ctx: { clientPubkey: string; session: ClientSession; @@ -596,6 +598,9 @@ export class ServerOpenStreamFactory { if (!removed && !session) { return; } + if (removed) { + this.deps.subscriptionStore.removeForClient(clientPubkey); + } this.deps.logger.info('Removed session after open-stream probe timeout', { clientPubkey, diff --git a/src/transport/nostr-server/outbound-notification-broadcaster.test.ts b/src/transport/nostr-server/outbound-notification-broadcaster.test.ts new file mode 100644 index 0000000..d7765c9 --- /dev/null +++ b/src/transport/nostr-server/outbound-notification-broadcaster.test.ts @@ -0,0 +1,135 @@ +import { describe, expect, test } from 'bun:test'; +import type { JSONRPCMessage } from '@contextvm/mcp-sdk/types.js'; +import type { Logger } from '../../core/utils/logger.js'; +import { CorrelationStore } from './correlation-store.js'; +import { OutboundNotificationBroadcaster } from './outbound-notification-broadcaster.js'; +import { SessionStore } from './session-store.js'; +import { + SubscriptionStore, + type ResourceSubscriptionMatcher, +} from './subscription-store.js'; + +const testLogger: Logger = { + debug: () => undefined, + info: () => undefined, + warn: () => undefined, + error: () => undefined, + withModule: () => testLogger, +}; + +const resourceUpdated: JSONRPCMessage = { + jsonrpc: '2.0', + method: 'notifications/resources/updated', + params: { uri: 'resource://alpha' }, +}; + +function createBroadcaster(matchesSubResource?: ResourceSubscriptionMatcher): { + broadcaster: OutboundNotificationBroadcaster; + notifications: string[]; + tasks: Promise[]; + subscriptions: SubscriptionStore; +} { + const sessionStore = new SessionStore(); + sessionStore.getOrCreateSession('client-a', false); + sessionStore.getOrCreateSession('client-b', false); + sessionStore.markInitialized('client-a'); + sessionStore.markInitialized('client-b'); + + const notifications: string[] = []; + const tasks: Promise[] = []; + const subscriptions = new SubscriptionStore(); + const broadcaster = new OutboundNotificationBroadcaster({ + correlationStore: new CorrelationStore(), + sessionStore, + subscriptionStore: subscriptions, + matchesSubResource, + sendNotification: async (clientPubkey) => { + notifications.push(clientPubkey); + }, + enqueueTask: (task) => { + tasks.push(task()); + }, + logger: testLogger, + }); + + return { broadcaster, notifications, tasks, subscriptions }; +} + +describe('OutboundNotificationBroadcaster', () => { + test('routes resource updates only to subscribers', async () => { + const { broadcaster, notifications, tasks, subscriptions } = + createBroadcaster(); + subscriptions.subscribe('client-a', 'resource://alpha'); + + await broadcaster.broadcast(resourceUpdated); + await Promise.all(tasks); + + expect(notifications).toEqual(['client-a']); + }); + + test('does not broadcast resource updates for an unknown URI', async () => { + const { broadcaster, notifications, tasks } = createBroadcaster(); + + await broadcaster.broadcast(resourceUpdated); + await Promise.all(tasks); + + expect(notifications).toEqual([]); + }); + + test('routes a resource update to every subscriber, and only those subscribers', async () => { + const { broadcaster, notifications, tasks, subscriptions } = + createBroadcaster(); + subscriptions.subscribe('client-a', 'resource://alpha'); + subscriptions.subscribe('client-b', 'resource://alpha'); + + await broadcaster.broadcast(resourceUpdated); + await Promise.all(tasks); + + expect(notifications).toEqual(['client-a', 'client-b']); + }); + + test('sends one update when a client matches both parent and exact subscriptions', async () => { + const { broadcaster, notifications, tasks, subscriptions } = + createBroadcaster((subscribedUri, updatedUri) => + updatedUri.startsWith(`${subscribedUri}/`), + ); + subscriptions.subscribe('client-a', 'resource://repo'); + subscriptions.subscribe('client-a', 'resource://repo/child'); + subscriptions.subscribe('client-b', 'resource://other'); + + await broadcaster.broadcast({ + jsonrpc: '2.0', + method: 'notifications/resources/updated', + params: { uri: 'resource://repo/child' }, + }); + await Promise.all(tasks); + + expect(notifications).toEqual(['client-a']); + }); + + test('does not fall back to generic broadcast when a resource URI is missing', async () => { + const { broadcaster, notifications, tasks, subscriptions } = + createBroadcaster(); + subscriptions.subscribe('client-a', 'resource://alpha'); + + await broadcaster.broadcast({ + jsonrpc: '2.0', + method: 'notifications/resources/updated', + }); + await Promise.all(tasks); + + expect(notifications).toEqual([]); + }); + + test('continues to broadcast unrelated notifications to initialized sessions', async () => { + const { broadcaster, notifications, tasks } = createBroadcaster(); + + await broadcaster.broadcast({ + jsonrpc: '2.0', + method: 'notifications/resources/list_changed', + }); + await Promise.all(tasks); + + expect(notifications).toEqual(['client-a', 'client-b']); + }); +}); diff --git a/src/transport/nostr-server/outbound-notification-broadcaster.ts b/src/transport/nostr-server/outbound-notification-broadcaster.ts index fa582f6..7ff9ffc 100644 --- a/src/transport/nostr-server/outbound-notification-broadcaster.ts +++ b/src/transport/nostr-server/outbound-notification-broadcaster.ts @@ -5,10 +5,16 @@ import { import { type Logger } from '../../core/utils/logger.js'; import { type CorrelationStore } from './correlation-store.js'; import { type SessionStore } from './session-store.js'; +import { + type ResourceSubscriptionMatcher, + type SubscriptionStore, +} from './subscription-store.js'; export interface OutboundNotificationBroadcasterDeps { correlationStore: CorrelationStore; sessionStore: SessionStore; + subscriptionStore: SubscriptionStore; + matchesSubResource?: ResourceSubscriptionMatcher; sendNotification: ( clientPubkey: string, notification: JSONRPCMessage, @@ -30,8 +36,44 @@ export class OutboundNotificationBroadcaster { */ public async broadcast(notification: JSONRPCMessage): Promise { try { + if ( + isJSONRPCNotification(notification) && + notification.method === 'notifications/resources/updated' + ) { + const resourceUri = notification.params?.uri; + if (typeof resourceUri !== 'string') { + this.deps.logger.warn('Resource update missing resource URI'); + return; + } + + const subscribers = this.deps.subscriptionStore.getSubscribersForUpdate( + resourceUri, + this.deps.matchesSubResource, + ); + if (subscribers.size === 0) { + this.deps.logger.warn('No clients subscribed to resource update', { + uri: resourceUri, + }); + return; + } + + for (const clientPubkey of subscribers) { + this.deps.enqueueTask(async () => { + try { + await this.deps.sendNotification(clientPubkey, notification); + } catch (error) { + this.deps.logger.error('Error sending resource update', { + error: error instanceof Error ? error.message : String(error), + clientPubkey, + uri: resourceUri, + }); + } + }); + } + return; + } + // Special handling for progress notifications - // TODO: Add handling for `notifications/resources/updated`, as they need to be associated with an id if ( isJSONRPCNotification(notification) && notification.method === 'notifications/progress' && diff --git a/src/transport/nostr-server/subscription-store.test.ts b/src/transport/nostr-server/subscription-store.test.ts new file mode 100644 index 0000000..35cc6ab --- /dev/null +++ b/src/transport/nostr-server/subscription-store.test.ts @@ -0,0 +1,121 @@ +import { describe, expect, it } from 'bun:test'; +import { SubscriptionStore } from './subscription-store.js'; + +describe('SubscriptionStore', () => { + it('tracks subscriptions idempotently', () => { + const store = new SubscriptionStore(); + + store.subscribe('client-a', 'resource://alpha'); + store.subscribe('client-a', 'resource://alpha'); + + expect(store.getSubscribers('resource://alpha')).toEqual( + new Set(['client-a']), + ); + }); + + it('removes subscriptions idempotently', () => { + const store = new SubscriptionStore(); + store.subscribe('client-a', 'resource://alpha'); + + store.unsubscribe('client-a', 'resource://alpha'); + store.unsubscribe('client-a', 'resource://alpha'); + + expect(store.getSubscribers('resource://alpha')).toEqual(new Set()); + }); + + it('keeps unrelated subscriptions when one relationship is removed', () => { + const store = new SubscriptionStore(); + store.subscribe('client-a', 'resource://alpha'); + store.subscribe('client-a', 'resource://beta'); + store.subscribe('client-b', 'resource://alpha'); + + store.unsubscribe('client-a', 'resource://alpha'); + + expect(store.getSubscribers('resource://alpha')).toEqual( + new Set(['client-b']), + ); + expect(store.getSubscribers('resource://beta')).toEqual( + new Set(['client-a']), + ); + }); + + it('removes every subscription for an evicted client', () => { + const store = new SubscriptionStore(); + store.subscribe('client-a', 'resource://alpha'); + store.subscribe('client-a', 'resource://beta'); + store.subscribe('client-b', 'resource://alpha'); + + store.removeForClient('client-a'); + + expect(store.getSubscribers('resource://alpha')).toEqual( + new Set(['client-b']), + ); + expect(store.getSubscribers('resource://beta')).toEqual(new Set()); + }); + + it('clears every subscription', () => { + const store = new SubscriptionStore(); + store.subscribe('client-a', 'resource://alpha'); + store.subscribe('client-b', 'resource://beta'); + + store.clear(); + + expect(store.getSubscribers('resource://alpha')).toEqual(new Set()); + expect(store.getSubscribers('resource://beta')).toEqual(new Set()); + }); + + it('does not expose its mutable subscriber index', () => { + const store = new SubscriptionStore(); + store.subscribe('client-a', 'resource://alpha'); + + const subscribers = store.getSubscribers('resource://alpha'); + (subscribers as Set).add('client-b'); + + expect(store.getSubscribers('resource://alpha')).toEqual( + new Set(['client-a']), + ); + }); + + it('matches only exact URIs when no sub-resource matcher is provided', () => { + const store = new SubscriptionStore(); + store.subscribe('client-a', 'resource://repo'); + + expect(store.getSubscribersForUpdate('resource://repo')).toEqual( + new Set(['client-a']), + ); + expect(store.getSubscribersForUpdate('resource://repo/child')).toEqual( + new Set(), + ); + }); + + it('uses server-defined sub-resource matching and sends once per client', () => { + const store = new SubscriptionStore(); + store.subscribe('client-a', 'resource://repo'); + store.subscribe('client-a', 'resource://repo/child'); + store.subscribe('client-b', 'resource://other'); + const matchesSubResource = (subscribedUri: string, updatedUri: string) => + updatedUri.startsWith(`${subscribedUri}/`); + + expect( + store.getSubscribersForUpdate( + 'resource://repo/child', + matchesSubResource, + ), + ).toEqual(new Set(['client-a'])); + expect( + store.getSubscribersForUpdate( + 'resource://repository', + matchesSubResource, + ), + ).toEqual(new Set()); + + store.unsubscribe('client-a', 'resource://repo'); + store.unsubscribe('client-a', 'resource://repo/child'); + expect( + store.getSubscribersForUpdate( + 'resource://repo/child', + matchesSubResource, + ), + ).toEqual(new Set()); + }); +}); diff --git a/src/transport/nostr-server/subscription-store.ts b/src/transport/nostr-server/subscription-store.ts new file mode 100644 index 0000000..6d96aee --- /dev/null +++ b/src/transport/nostr-server/subscription-store.ts @@ -0,0 +1,91 @@ +/** Decides whether an updated resource is a sub-resource of a subscribed URI. */ +export type ResourceSubscriptionMatcher = ( + subscribedUri: string, + updatedUri: string, +) => boolean; + +/** + * Tracks resource subscriptions independently from transient request routes. + */ +export class SubscriptionStore { + private readonly uriToClients = new Map>(); + private readonly clientToUris = new Map>(); + + public subscribe(clientPubkey: string, uri: string): void { + let clients = this.uriToClients.get(uri); + if (!clients) { + clients = new Set(); + this.uriToClients.set(uri, clients); + } + clients.add(clientPubkey); + + let uris = this.clientToUris.get(clientPubkey); + if (!uris) { + uris = new Set(); + this.clientToUris.set(clientPubkey, uris); + } + uris.add(uri); + } + + public unsubscribe(clientPubkey: string, uri: string): void { + const clients = this.uriToClients.get(uri); + clients?.delete(clientPubkey); + if (clients?.size === 0) { + this.uriToClients.delete(uri); + } + + const uris = this.clientToUris.get(clientPubkey); + uris?.delete(uri); + if (uris?.size === 0) { + this.clientToUris.delete(clientPubkey); + } + } + + public removeForClient(clientPubkey: string): void { + const uris = this.clientToUris.get(clientPubkey); + if (!uris) { + return; + } + for (const uri of uris) { + const clients = this.uriToClients.get(uri); + clients?.delete(clientPubkey); + if (clients?.size === 0) { + this.uriToClients.delete(uri); + } + } + this.clientToUris.delete(clientPubkey); + } + + public getSubscribers(uri: string): ReadonlySet { + // Do not expose the set held by the index: callers iterate this while + // sessions can be evicted, and a cast at a call site must not be able to + // mutate subscription state. + return new Set(this.uriToClients.get(uri)); + } + + /** Finds subscribers to an update, including server-defined sub-resources. */ + public getSubscribersForUpdate( + updatedUri: string, + matchesSubResource?: ResourceSubscriptionMatcher, + ): ReadonlySet { + const subscribers = new Set(this.uriToClients.get(updatedUri)); + if (matchesSubResource) { + for (const [subscribedUri, clients] of this.uriToClients) { + if ( + subscribedUri !== updatedUri && + matchesSubResource(subscribedUri, updatedUri) + ) { + for (const clientPubkey of clients) { + subscribers.add(clientPubkey); + } + } + } + } + return subscribers; + } + + public clear(): void { + this.uriToClients.clear(); + this.clientToUris.clear(); + } +}