import { StorageError, ValidationError } from '../errors'; import { AuthContext, KeyValueStore, Logger, ResourceDefinition, StoredSubscription } from '../types'; import { generateSubscriptionId, generateWebhookUrl } from '../utils'; export class SubscriptionManager { constructor( private store: KeyValueStore, private resources: ResourceDefinition[], private publicUrl: string, private logger?: Logger ) {} /** * Create a new subscription */ async createSubscription(params: { uri: string; clientCallbackUrl: string; clientCallbackSecret?: string; source?: string; context: AuthContext; }): Promise<{ subscriptionId: string; status: string }> { const { uri, clientCallbackUrl, clientCallbackSecret, source, context } = params; // Find resource definition const resource = this.findResourceForUri(uri); if (!resource) { throw new ValidationError(`No resource found for URI: ${uri}`); } if (!resource.subscription) { throw new ValidationError(`Resource ${resource.name} does not support subscriptions`); } // Check for existing subscription for this URI, user, AND source — only clean up same-source const existingSubscription = await this.findSubscriptionByUri(uri, context.userId, source); if (existingSubscription) { this.logger?.info('Found existing subscription for URI (same source), cleaning up first', { existingSubscriptionId: existingSubscription.subscriptionId, uri, source, }); try { await this.deleteSubscription(existingSubscription.subscriptionId, context); } catch (error) { this.logger?.warn('Failed to delete existing subscription, continuing anyway', { existingSubscriptionId: existingSubscription.subscriptionId, error: error instanceof Error ? error.message : String(error), }); } } // Generate subscription ID const subscriptionId = generateSubscriptionId(); const thirdPartyWebhookUrl = generateWebhookUrl(this.publicUrl, subscriptionId); this.logger?.info('Creating subscription', { subscriptionId, uri, clientCallbackUrl, source, }); try { // Call resource's onSubscribe handler const metadata = await resource.subscription.onSubscribe( uri, subscriptionId, thirdPartyWebhookUrl, context ); // Store subscription data const subscriptionData: StoredSubscription = { uri, resourceType: resource.name, clientCallbackUrl, clientCallbackSecret, userId: context.userId, thirdPartyWebhookId: metadata.thirdPartyWebhookId, metadata: metadata.metadata, source, createdAt: Date.now(), }; await this.storeSubscription(subscriptionId, subscriptionData, context.userId); this.logger?.info('Subscription created', { subscriptionId }); return { subscriptionId, status: 'active', }; } catch (error) { this.logger?.error('Failed to create subscription', { subscriptionId, error: error instanceof Error ? error.message : String(error), }); throw new StorageError('Failed to create subscription', { cause: error }); } } /** * Delete a subscription */ async deleteSubscription(subscriptionId: string, context: AuthContext): Promise { this.logger?.info('Deleting subscription', { subscriptionId }); // Load subscription const subscription = await this.loadSubscription(subscriptionId); if (!subscription) { throw new ValidationError(`Subscription ${subscriptionId} not found`); } // Verify ownership if (subscription.userId !== context.userId) { throw new ValidationError('Unauthorized to delete this subscription'); } // Find resource const resource = this.findResourceForUri(subscription.uri); if (!resource?.subscription) { throw new StorageError('Resource subscription configuration not found'); } try { // Call resource's onUnsubscribe handler await resource.subscription.onUnsubscribe( subscription.uri, subscriptionId, { thirdPartyWebhookId: subscription.thirdPartyWebhookId, metadata: subscription.metadata, }, context ); // Delete from store await this.removeSubscription(subscriptionId, context.userId); this.logger?.info('Subscription deleted', { subscriptionId }); } catch (error) { this.logger?.error('Failed to delete subscription', { subscriptionId, error: error instanceof Error ? error.message : String(error), }); throw new StorageError('Failed to delete subscription', { cause: error }); } } /** * Get subscription by ID */ async getSubscription(subscriptionId: string): Promise { return this.loadSubscription(subscriptionId); } /** * List subscriptions for a user */ async listSubscriptions(userId: string): Promise { const indexKey = `user:${userId}:subscriptions`; const indexData = await this.store.get(indexKey); if (!indexData) { return []; } const subscriptionIds: string[] = JSON.parse(indexData); const subscriptions = await Promise.all( subscriptionIds.map((id) => this.loadSubscription(id)) ); return subscriptions.filter((s): s is StoredSubscription => s !== null); } /** * Find existing subscription for a URI, user, and optionally source. * When source is provided, only matches subscriptions with the same source. * When source is undefined, only matches subscriptions without a source (backward compat). */ async findSubscriptionByUri(uri: string, userId: string, source?: string): Promise<(StoredSubscription & { subscriptionId: string }) | null> { const indexKey = `user:${userId}:subscriptions`; const indexData = await this.store.get(indexKey); if (!indexData) { return null; } const subscriptionIds: string[] = JSON.parse(indexData); for (const subscriptionId of subscriptionIds) { const subscription = await this.loadSubscription(subscriptionId); if (subscription && subscription.uri === uri) { // Match source: both must be same (including both undefined) if ((subscription.source || undefined) === source) { return { ...subscription, subscriptionId }; } } } return null; } /** * Find ALL subscriptions for a URI across all users and sources. * Used to notify all subscribers when a webhook arrives. */ async findAllSubscriptionsByUri(uri: string, excludeSubscriptionId?: string): Promise<(StoredSubscription & { subscriptionId: string })[]> { const results: (StoredSubscription & { subscriptionId: string })[] = []; // Scan all subscription keys to find matching URIs const allKeys = await this.store.scan?.('subscription:*') || []; for (const key of allKeys) { const subscriptionId = key.replace('subscription:', ''); if (subscriptionId === excludeSubscriptionId) continue; const subscription = await this.loadSubscription(subscriptionId); if (subscription && subscription.uri === uri) { results.push({ ...subscription, subscriptionId }); } } return results; } /** * Store subscription data */ private async storeSubscription( subscriptionId: string, data: StoredSubscription, userId: string ): Promise { const key = `subscription:${subscriptionId}`; await this.store.set(key, JSON.stringify(data)); // Update user index const indexKey = `user:${userId}:subscriptions`; const existing = await this.store.get(indexKey); const subscriptionIds: string[] = existing ? JSON.parse(existing) : []; if (!subscriptionIds.includes(subscriptionId)) { subscriptionIds.push(subscriptionId); await this.store.set(indexKey, JSON.stringify(subscriptionIds)); } } /** * Load subscription data */ private async loadSubscription(subscriptionId: string): Promise { const key = `subscription:${subscriptionId}`; const data = await this.store.get(key); if (!data) { return null; } return JSON.parse(data); } /** * Remove subscription data */ private async removeSubscription(subscriptionId: string, userId: string): Promise { const key = `subscription:${subscriptionId}`; await this.store.delete(key); // Update user index const indexKey = `user:${userId}:subscriptions`; const existing = await this.store.get(indexKey); if (existing) { const subscriptionIds: string[] = JSON.parse(existing); const filtered = subscriptionIds.filter((id) => id !== subscriptionId); await this.store.set(indexKey, JSON.stringify(filtered)); } } /** * Find resource definition for URI */ private findResourceForUri(uri: string): ResourceDefinition | undefined { return this.resources.find((resource) => { // Convert URI template to regex pattern // Step 1: Replace template variables with a placeholder const withPlaceholders = resource.uri.replace(/\{[^}]+\}/g, '__PLACEHOLDER__'); // Step 2: Escape special regex characters (but not our placeholders) const escapedPattern = withPlaceholders.replace(/[-\\^$*+?.()|[\]{}]/g, '\\$&'); // Step 3: Replace placeholders with capture groups const patternString = escapedPattern.replace(/__PLACEHOLDER__/g, '([^/]+)'); // Step 4: Create regex with anchors const pattern = new RegExp('^' + patternString + '$'); console.log('Pattern for resource', resource.name, ':', pattern); console.log('Testing URI:', uri); return pattern.test(uri); }); } }