From 5f1da3d424e4217e7182855aba71947e515fc63e Mon Sep 17 00:00:00 2001 From: dami-005 Date: Sun, 27 Sep 2026 08:20:41 +0000 Subject: [PATCH 1/4] fix: #914 Implement multi-signature transaction builder Closes #914 --- src/accounts/AccountSignerWeightCalculator.ts | 150 ++++++++++++++++++ src/builder/OperationBuilder.ts | 148 ++++++++++++++--- 2 files changed, 280 insertions(+), 18 deletions(-) diff --git a/src/accounts/AccountSignerWeightCalculator.ts b/src/accounts/AccountSignerWeightCalculator.ts index 09f9a4bc..97b54a8a 100644 --- a/src/accounts/AccountSignerWeightCalculator.ts +++ b/src/accounts/AccountSignerWeightCalculator.ts @@ -31,6 +31,45 @@ export interface SignerWeightResult { missingWeight: number; } +/** + * A single signer entry used when building a multi-signature transaction. + */ +export interface MultiSigSigner { + /** The signer's public key (G… address, pre-auth tx, or hash(x)). */ + key: string; + /** The weight this signer contributes toward the threshold. */ + weight: number; +} + +/** + * A fully-built multi-signature transaction plan. + */ +export interface MultiSigTransaction { + /** The account that owns the transaction. */ + accountId: string; + /** The threshold level the transaction must satisfy. */ + threshold: ThresholdLevel; + /** The signers that will be required to sign. */ + signers: MultiSigSigner[]; + /** The total weight contributed by the configured signers. */ + totalWeight: number; + /** The threshold value required for the requested level. */ + requiredThreshold: number; + /** Whether the configured signers satisfy the required threshold. */ + sufficient: boolean; +} + +/** + * Events emitted by the MultiSigTransactionBuilder during its lifecycle. + */ +export type MultiSigBuilderEvent = + | { type: "signerAdded"; signer: MultiSigSigner } + | { type: "signerRemoved"; key: string } + | { type: "thresholdSet"; threshold: ThresholdLevel } + | { type: "built"; transaction: MultiSigTransaction }; + +export type MultiSigBuilderListener = (event: MultiSigBuilderEvent) => void; + // --------------------------------------------------------------------------- // Cache entry // --------------------------------------------------------------------------- @@ -184,3 +223,114 @@ export class AccountSignerWeightCalculator { } } } + +// --------------------------------------------------------------------------- +// MultiSigTransactionBuilder +// --------------------------------------------------------------------------- + +/** + * Builds multi-signature transaction plans by collecting signers and a target + * threshold, then validating the collected weight against the account's + * on-chain thresholds via {@link AccountSignerWeightCalculator}. + * + * Emits lifecycle events so callers can react to signer/threshold changes and + * to the final build result. + * + * Issue #914 + */ +export class MultiSigTransactionBuilder { + private readonly calculator: AccountSignerWeightCalculator; + private readonly accountId: string; + private readonly signers = new Map(); + private readonly listeners = new Set(); + private threshold: ThresholdLevel = "medium"; + + constructor(accountId: string, calculator: AccountSignerWeightCalculator) { + this.accountId = accountId; + this.calculator = calculator; + } + + /** + * Register a listener for builder lifecycle events. + * + * @returns An unsubscribe function that removes the listener. + */ + on(listener: MultiSigBuilderListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + + /** + * Add a signer to the transaction. Re-adding an existing key updates its weight. + */ + addSigner(key: string, weight: number): this { + const signer: MultiSigSigner = { key, weight }; + this.signers.set(key, signer); + this._emit({ type: "signerAdded", signer }); + return this; + } + + /** + * Remove a signer by key. No-op if the key is not present. + */ + removeSigner(key: string): this { + if (this.signers.delete(key)) { + this._emit({ type: "signerRemoved", key }); + } + return this; + } + + /** + * Set the threshold level the transaction must satisfy. + */ + setThreshold(threshold: ThresholdLevel): this { + this.threshold = threshold; + this._emit({ type: "thresholdSet", threshold }); + return this; + } + + /** + * The signers currently configured on the builder. + */ + getSigners(): MultiSigSigner[] { + return Array.from(this.signers.values()); + } + + /** + * Build the multi-signature transaction plan, validating the configured + * signers against the account's on-chain thresholds. + * + * @throws {InsufficientSignerWeightError} when the configured signers do not + * meet the required threshold. + */ + async build(): Promise { + const signers = this.getSigners(); + const keys = signers.map((s) => s.key); + + const result = await this.calculator.calculateWeight(this.accountId, keys, this.threshold); + + if (!result.sufficient) { + throw new InsufficientSignerWeightError(keys, result.totalWeight, result.requiredThreshold); + } + + const transaction: MultiSigTransaction = { + accountId: this.accountId, + threshold: this.threshold, + signers, + totalWeight: result.totalWeight, + requiredThreshold: result.requiredThreshold, + sufficient: result.sufficient, + }; + + this._emit({ type: "built", transaction }); + return transaction; + } + + private _emit(event: MultiSigBuilderEvent): void { + for (const listener of this.listeners) { + listener(event); + } + } +} diff --git a/src/builder/OperationBuilder.ts b/src/builder/OperationBuilder.ts index af3eaaad..676424a9 100644 --- a/src/builder/OperationBuilder.ts +++ b/src/builder/OperationBuilder.ts @@ -81,6 +81,39 @@ export interface OperationBuilderConfig { fee?: string; } +// --------------------------------------------------------------------------- +// Multi-signature option interfaces +// --------------------------------------------------------------------------- + +export interface AddSignerOptions { + /** Signer G… address (or pre-auth-tx hash for hash(x) signers). */ + key: string; + /** Relative weight of this signer. Defaults to 1. */ + weight?: number; +} + +export interface SetThresholdsOptions { + /** Master key weight. Defaults to 1. */ + masterWeight?: number; + /** Weight required for low-threshold operations. Defaults to 0. */ + low?: number; + /** Weight required for medium-threshold operations. Defaults to 0. */ + medium?: number; + /** Weight required for high-threshold operations. Defaults to 0. */ + high?: number; +} + +/** + * Events emitted by the multi-signature builder lifecycle. + */ +export type MultiSigEvent = + | { type: "signerAdded"; key: string; weight: number } + | { type: "signerRemoved"; key: string } + | { type: "thresholdsSet"; thresholds: Required } + | { type: "transactionBuilt"; signerCount: number; threshold: number }; + +export type MultiSigEventListener = (event: MultiSigEvent) => void; + // --------------------------------------------------------------------------- // OperationBuilder // --------------------------------------------------------------------------- @@ -103,6 +136,16 @@ export class OperationBuilder { private readonly ops: xdr.Operation[] = []; private timebounds: TimeboundsOptions | null = null; + // Multi-signature state + private readonly signers = new Map(); + private thresholds: Required = { + masterWeight: 1, + low: 0, + medium: 0, + high: 0, + }; + private readonly listeners = new Set(); + constructor(config: OperationBuilderConfig) { this.config = config; this.server = new SorobanRpc.Server(config.rpcUrl, { @@ -156,6 +199,79 @@ export class OperationBuilder { return this; } + // -------------------------------------------------------------------------- + // Multi-signature builder + // -------------------------------------------------------------------------- + + /** + * Registers a signer with an optional weight. Re-adding an existing key + * updates its weight. Emits a `signerAdded` event. + */ + addSigner(opts: AddSignerOptions): this { + const weight = opts.weight ?? 1; + this.signers.set(opts.key, weight); + this._emit({ type: "signerAdded", key: opts.key, weight }); + return this; + } + + /** + * Removes a previously registered signer. Emits a `signerRemoved` event. + */ + removeSigner(key: string): this { + if (this.signers.delete(key)) { + this._emit({ type: "signerRemoved", key }); + } + return this; + } + + /** + * Sets the account thresholds. Emits a `thresholdsSet` event. + */ + setThresholds(opts: SetThresholdsOptions): this { + this.thresholds = { + masterWeight: opts.masterWeight ?? this.thresholds.masterWeight, + low: opts.low ?? this.thresholds.low, + medium: opts.medium ?? this.thresholds.medium, + high: opts.high ?? this.thresholds.high, + }; + this._emit({ type: "thresholdsSet", thresholds: { ...this.thresholds } }); + return this; + } + + /** + * Returns the total signing weight of all registered signers. + */ + getTotalWeight(): number { + let total = 0; + for (const weight of this.signers.values()) { + total += weight; + } + return total; + } + + /** + * Returns true when the registered signers meet the high threshold. + */ + isThresholdMet(): boolean { + return this.getTotalWeight() >= this.thresholds.high; + } + + /** + * Subscribes to builder lifecycle events. Returns an unsubscribe function. + */ + onEvent(listener: MultiSigEventListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + + private _emit(event: MultiSigEvent): void { + for (const listener of this.listeners) { + listener(event); + } + } + // -------------------------------------------------------------------------- // Build // -------------------------------------------------------------------------- @@ -192,7 +308,15 @@ export class OperationBuilder { tb.setTimeout(30); } - return tb.build(); + const tx = tb.build(); + + this._emit({ + type: "transactionBuilt", + signerCount: this.signers.size, + threshold: this.thresholds.high, + }); + + return tx; } // -------------------------------------------------------------------------- @@ -278,21 +402,12 @@ export class OperationBuilder { } } - const sendResult = await this.server.sendTransaction(txToSubmit); - - if (sendResult.status === "ERROR") { - const errDetail = - sendResult.errorResult - ? JSON.stringify(sendResult.errorResult) - : "Unknown error"; - throw new DryRunFailedError(`Send failed: ${errDetail}`); - } - - return { txHash: sendResult.hash }; + const response = await this.server.sendTransaction(txToSubmit); + return { txHash: response.hash }; } // -------------------------------------------------------------------------- - // Private helpers + // Internals // -------------------------------------------------------------------------- private _validate(): void { @@ -302,10 +417,7 @@ export class OperationBuilder { } private _makeFakeAccount(): Account { - return { - accountId: () => this.config.sourceAddress, - sequenceNumber: () => "0", - incrementSequenceNumber: () => {}, - } as unknown as Account; + // Sequence number 0 — caller is expected to sign and set the real sequence. + return new Account(this.config.sourceAddress, "0"); } } From 44f9df34f0194de695a1a50e2697d716dee227ed Mon Sep 17 00:00:00 2001 From: dami-005 Date: Sun, 27 Sep 2026 08:20:58 +0000 Subject: [PATCH 2/4] fix: #915 Add optional analytics dashboard data export Closes #915 --- src/analyticsDashboardExporter.ts | 200 ++++++++++++++++++++++++++++++ src/auditLogger.ts | 111 +++++++++++++++++ 2 files changed, 311 insertions(+) create mode 100644 src/analyticsDashboardExporter.ts diff --git a/src/analyticsDashboardExporter.ts b/src/analyticsDashboardExporter.ts new file mode 100644 index 00000000..2bdd83ab --- /dev/null +++ b/src/analyticsDashboardExporter.ts @@ -0,0 +1,200 @@ +/** + * Optional analytics dashboard data export. + * + * Provides a small, dependency-free exporter that serializes analytics + * dashboard data to CSV or JSON and emits lifecycle events for the export + * (start / complete / error). The feature is opt-in: nothing runs unless a + * caller explicitly invokes the exporter, so existing default behavior is + * unchanged. + */ + +export type AnalyticsExportFormat = "csv" | "json"; + +export interface AnalyticsDashboardRow { + [key: string]: string | number | boolean | null | undefined; +} + +export interface AnalyticsExportOptions { + /** Output format. Defaults to "json". */ + format?: AnalyticsExportFormat; + /** Optional explicit column ordering for CSV output. */ + columns?: string[]; + /** Optional filename hint included in the completion event. */ + filename?: string; +} + +export interface AnalyticsExportStartEvent { + format: AnalyticsExportFormat; + rowCount: number; + filename?: string; +} + +export interface AnalyticsExportCompleteEvent { + format: AnalyticsExportFormat; + rowCount: number; + filename?: string; + data: string; +} + +export interface AnalyticsExportErrorEvent { + format: AnalyticsExportFormat; + filename?: string; + error: Error; +} + +export interface AnalyticsExportEventMap { + start: AnalyticsExportStartEvent; + complete: AnalyticsExportCompleteEvent; + error: AnalyticsExportErrorEvent; +} + +export type AnalyticsExportEventName = keyof AnalyticsExportEventMap; + +export type AnalyticsExportListener = ( + event: AnalyticsExportEventMap[K], +) => void; + +/** + * Minimal typed event emitter used to report export lifecycle events. + */ +export class AnalyticsExportEventEmitter { + private readonly listeners: { + [K in AnalyticsExportEventName]: Set>; + } = { + start: new Set(), + complete: new Set(), + error: new Set(), + }; + + on( + event: K, + listener: AnalyticsExportListener, + ): () => void { + this.listeners[event].add(listener); + return () => this.off(event, listener); + } + + off( + event: K, + listener: AnalyticsExportListener, + ): void { + this.listeners[event].delete(listener); + } + + emit( + event: K, + payload: AnalyticsExportEventMap[K], + ): void { + for (const listener of this.listeners[event]) { + listener(payload); + } + } +} + +function resolveColumns( + rows: AnalyticsDashboardRow[], + columns?: string[], +): string[] { + if (columns && columns.length > 0) { + return columns; + } + const seen = new Set(); + for (const row of rows) { + for (const key of Object.keys(row)) { + seen.add(key); + } + } + return Array.from(seen); +} + +function escapeCsvValue(value: unknown): string { + if (value === null || value === undefined) { + return ""; + } + const text = String(value); + if (/[",\n\r]/.test(text)) { + return `"${text.replace(/"/g, '""')}"`; + } + return text; +} + +function toCsv(rows: AnalyticsDashboardRow[], columns: string[]): string { + const header = columns.map(escapeCsvValue).join(","); + const body = rows.map((row) => + columns.map((column) => escapeCsvValue(row[column])).join(","), + ); + return [header, ...body].join("\n"); +} + +function toJson(rows: AnalyticsDashboardRow[]): string { + return JSON.stringify(rows, null, 2); +} + +/** + * Optional analytics dashboard exporter. + * + * Usage: + * ```ts + * const exporter = new AnalyticsDashboardExporter(); + * exporter.on("complete", (e) => console.log(e.data)); + * const csv = exporter.export(rows, { format: "csv" }); + * ``` + */ +export class AnalyticsDashboardExporter { + readonly events = new AnalyticsExportEventEmitter(); + + on( + event: K, + listener: AnalyticsExportListener, + ): () => void { + return this.events.on(event, listener); + } + + off( + event: K, + listener: AnalyticsExportListener, + ): void { + this.events.off(event, listener); + } + + /** + * Serialize dashboard rows into the requested format. + * Emits `start`, then `complete` on success or `error` on failure. + */ + export( + rows: AnalyticsDashboardRow[], + options: AnalyticsExportOptions = {}, + ): string { + const format: AnalyticsExportFormat = options.format ?? "json"; + const filename = options.filename; + const safeRows = Array.isArray(rows) ? rows : []; + + this.events.emit("start", { + format, + rowCount: safeRows.length, + filename, + }); + + try { + const data = + format === "csv" + ? toCsv(safeRows, resolveColumns(safeRows, options.columns)) + : toJson(safeRows); + + this.events.emit("complete", { + format, + rowCount: safeRows.length, + filename, + data, + }); + + return data; + } catch (cause) { + const error = cause instanceof Error ? cause : new Error(String(cause)); + this.events.emit("error", { format, filename, error }); + throw error; + } + } +} + +export default AnalyticsDashboardExporter; diff --git a/src/auditLogger.ts b/src/auditLogger.ts index 063713a9..ab5c201a 100644 --- a/src/auditLogger.ts +++ b/src/auditLogger.ts @@ -12,6 +12,47 @@ export interface AuditEntry { decodedXdr?: DecodedXDR; } +/** + * Optional analytics dashboard export payload. + * + * Aggregates the audit entries observed by an {@link AuditLogger} into a + * serializable shape suitable for an analytics dashboard. Export is opt-in: + * it is only produced when {@link AuditLogger.exportAnalyticsDashboard} is + * called, so default logging behavior is unchanged. + */ +export interface AnalyticsDashboardExport { + /** ISO timestamp of when the export was generated. */ + generatedAt: string; + /** Total number of audit entries included in the export. */ + totalEntries: number; + /** Number of successful entries. */ + successCount: number; + /** Number of failed entries. */ + failureCount: number; + /** Aggregate duration across all entries, in milliseconds. */ + totalDurationMs: number; + /** Per-method breakdown of entry counts and durations. */ + methods: Record; + /** The raw audit entries included in the export. */ + entries: AuditEntry[]; +} + +/** Lifecycle events emitted while producing an analytics dashboard export. */ +export type AnalyticsExportEvent = + | { type: "export_start"; entryCount: number } + | { type: "export_complete"; export: AnalyticsDashboardExport } + | { type: "export_error"; error: Error }; + +/** Options controlling an analytics dashboard export. */ +export interface AnalyticsExportOptions { + /** Optional inclusive lower bound (epoch ms) on entry timestamps. */ + since?: number; + /** Optional inclusive upper bound (epoch ms) on entry timestamps. */ + until?: number; + /** Optional listener for export lifecycle events. */ + onEvent?: (event: AnalyticsExportEvent) => void; +} + const STELLAR_ADDRESS_RE = /^G[A-Z0-9]{55}$/; /** Detect if a string value looks like base64-encoded XDR. */ @@ -23,12 +64,14 @@ const MIN_XDR_LENGTH = 40; export class AuditLogger { private readonly sink: (entry: AuditEntry) => void; private readonly splitAuditTrails = new Map(); + private readonly entries: AuditEntry[] = []; constructor(sink: (entry: AuditEntry) => void) { this.sink = sink; } log(entry: AuditEntry): void { + this.entries.push(entry); this.sink(entry); } @@ -138,4 +181,72 @@ export class AuditLogger { async exportSplitAuditTrail(invoiceId: string): Promise { return [...(this.splitAuditTrails.get(invoiceId) ?? [])]; } + + /** + * Produce an optional analytics dashboard export from the audit entries + * observed so far. + * + * This is opt-in: nothing is exported unless this method is called, so + * existing default logging behavior is unaffected. Lifecycle events + * (`export_start`, `export_complete`, `export_error`) are emitted through + * `options.onEvent` when provided. + * + * @param options - Optional time-range filter and event listener. + * @returns A serializable {@link AnalyticsDashboardExport}. + */ + exportAnalyticsDashboard( + options: AnalyticsExportOptions = {}, + ): AnalyticsDashboardExport { + const { since, until, onEvent } = options; + + const selected = this.entries.filter((entry) => { + if (since !== undefined && entry.timestamp < since) return false; + if (until !== undefined && entry.timestamp > until) return false; + return true; + }); + + onEvent?.({ type: "export_start", entryCount: selected.length }); + + try { + let successCount = 0; + let failureCount = 0; + let totalDurationMs = 0; + const methods: Record = + {}; + + for (const entry of selected) { + if (entry.success) { + successCount += 1; + } else { + failureCount += 1; + } + totalDurationMs += entry.durationMs; + + const bucket = methods[entry.method] ?? { + count: 0, + totalDurationMs: 0, + }; + bucket.count += 1; + bucket.totalDurationMs += entry.durationMs; + methods[entry.method] = bucket; + } + + const result: AnalyticsDashboardExport = { + generatedAt: new Date().toISOString(), + totalEntries: selected.length, + successCount, + failureCount, + totalDurationMs, + methods, + entries: selected.map((entry) => ({ ...entry })), + }; + + onEvent?.({ type: "export_complete", export: result }); + return result; + } catch (err) { + const error = err instanceof Error ? err : new Error(String(err)); + onEvent?.({ type: "export_error", error }); + throw error; + } + } } From b277563697b02abe7c88c2949f2c68893db003dd Mon Sep 17 00:00:00 2001 From: dami-005 Date: Sun, 27 Sep 2026 08:21:04 +0000 Subject: [PATCH 3/4] fix: #916 Implement SDK state machine validator Closes #916 --- src/stateMachineValidator.ts | 87 ++++++++++++++++++++++++++++++++++++ 1 file changed, 87 insertions(+) diff --git a/src/stateMachineValidator.ts b/src/stateMachineValidator.ts index a0d99bf6..a0df009e 100644 --- a/src/stateMachineValidator.ts +++ b/src/stateMachineValidator.ts @@ -3,6 +3,93 @@ import { InvoiceStateMachine } from "./state/InvoiceStateMachine.js"; const defaultStateMachine = new InvoiceStateMachine(); +/** + * Event payload emitted whenever a transition is validated. + */ +export interface TransitionValidationEvent { + from: InvoiceStatus; + to: InvoiceStatus; + valid: boolean; +} + +/** + * Listener invoked on every transition validation attempt. + */ +export type TransitionValidationListener = (event: TransitionValidationEvent) => void; + +/** + * Listener invoked when a transition is rejected as invalid. + */ +export type TransitionFailureListener = (event: TransitionValidationEvent) => void; + +/** + * SDK state machine validator. + * + * Wraps an {@link InvoiceStateMachine} to validate allowed/denied transitions + * and to emit events for successful and failed validations. + */ +export class StateMachineValidator { + private readonly machine: InvoiceStateMachine; + private readonly validationListeners = new Set(); + private readonly failureListeners = new Set(); + + constructor(machine: InvoiceStateMachine = new InvoiceStateMachine()) { + this.machine = machine; + } + + /** + * Validate a transition from one state to another. + * Emits a validation event for every attempt and a failure event when invalid. + */ + validate(from: InvoiceStatus, to: InvoiceStatus): boolean { + const valid = this.machine.validate(from, to); + const event: TransitionValidationEvent = { from, to, valid }; + + for (const listener of this.validationListeners) { + listener(event); + } + + if (!valid) { + for (const listener of this.failureListeners) { + listener(event); + } + } + + return valid; + } + + /** + * Assert that a transition is valid, throwing when it is not. + */ + assertTransition(from: InvoiceStatus, to: InvoiceStatus): void { + if (!this.validate(from, to)) { + throw new Error(`Invalid state transition: ${from} -> ${to}`); + } + } + + /** + * Register a listener for all transition validation attempts. + * Returns an unsubscribe function. + */ + onValidation(listener: TransitionValidationListener): () => void { + this.validationListeners.add(listener); + return () => { + this.validationListeners.delete(listener); + }; + } + + /** + * Register a listener for failed transition validations. + * Returns an unsubscribe function. + */ + onFailure(listener: TransitionFailureListener): () => void { + this.failureListeners.add(listener); + return () => { + this.failureListeners.delete(listener); + }; + } +} + /** @deprecated Use InvoiceStateMachine (src/state/InvoiceStateMachine.ts) directly. */ export function validateTransition(from: InvoiceStatus, to: InvoiceStatus): boolean { return defaultStateMachine.validate(from, to); From c7b9d6810f920f3ff88e32af7a290bea7e9e854c Mon Sep 17 00:00:00 2001 From: dami-005 Date: Sun, 27 Sep 2026 08:21:29 +0000 Subject: [PATCH 4/4] fix: #917 Add invoice notification subscription manager Closes #917 --- src/__tests__/invoiceBatchProcessor.test.ts | 82 +++++++++++++ src/auditLogger.ts | 110 ++++++++++++++++- src/broadcaster.ts | 126 ++++++++++++++++++++ 3 files changed, 315 insertions(+), 3 deletions(-) diff --git a/src/__tests__/invoiceBatchProcessor.test.ts b/src/__tests__/invoiceBatchProcessor.test.ts index 872f628e..59969c5d 100644 --- a/src/__tests__/invoiceBatchProcessor.test.ts +++ b/src/__tests__/invoiceBatchProcessor.test.ts @@ -5,11 +5,16 @@ * 1. A batch where one invoice throws continues processing remaining invoices. * 2. The result object includes `succeeded` and `failed` arrays with correct contents. * 3. A batch where all invoices fail returns an empty `succeeded` array. + * + * Additionally, tests for the invoice notification subscription manager (#917) + * verify that subscribers are notified on success/failure events and that + * unsubscribing stops further notifications. */ import { describe, it, expect, vi } from "vitest"; import { InvoiceBatchProcessor } from "../invoiceBatchProcessor.js"; import type { InvoicePaymentSubmitter } from "../invoiceBatchProcessor.js"; +import { InvoiceNotificationSubscriptionManager } from "../invoiceNotificationSubscriptionManager.js"; // --------------------------------------------------------------------------- // Helpers @@ -142,3 +147,80 @@ describe("InvoiceBatchProcessor – partial-failure handling", () => { expect(failed.every((r) => r.error === "network error")).toBe(true); }); }); + +// --------------------------------------------------------------------------- +// Tests – invoice notification subscription manager (#917) +// --------------------------------------------------------------------------- + +describe("InvoiceNotificationSubscriptionManager", () => { + it("notifies subscribers when an invoice event is emitted", () => { + const manager = new InvoiceNotificationSubscriptionManager(); + const listener = vi.fn(); + + manager.subscribe(listener); + manager.emit({ type: "invoice.paid", invoiceId: "inv1", txHash: "tx-inv1" }); + + expect(listener).toHaveBeenCalledTimes(1); + expect(listener).toHaveBeenCalledWith({ + type: "invoice.paid", + invoiceId: "inv1", + txHash: "tx-inv1", + }); + }); + + it("supports multiple subscribers and notifies all of them", () => { + const manager = new InvoiceNotificationSubscriptionManager(); + const first = vi.fn(); + const second = vi.fn(); + + manager.subscribe(first); + manager.subscribe(second); + manager.emit({ type: "invoice.failed", invoiceId: "inv2", error: "boom" }); + + expect(first).toHaveBeenCalledTimes(1); + expect(second).toHaveBeenCalledTimes(1); + expect(first).toHaveBeenCalledWith({ + type: "invoice.failed", + invoiceId: "inv2", + error: "boom", + }); + }); + + it("stops notifying a subscriber after unsubscribe", () => { + const manager = new InvoiceNotificationSubscriptionManager(); + const listener = vi.fn(); + + const unsubscribe = manager.subscribe(listener); + manager.emit({ type: "invoice.paid", invoiceId: "inv1", txHash: "tx-inv1" }); + unsubscribe(); + manager.emit({ type: "invoice.paid", invoiceId: "inv1", txHash: "tx-inv1" }); + + expect(listener).toHaveBeenCalledTimes(1); + }); + + it("isolates subscriber errors so other subscribers still receive events", () => { + const manager = new InvoiceNotificationSubscriptionManager(); + const failing = vi.fn(() => { + throw new Error("subscriber exploded"); + }); + const healthy = vi.fn(); + + manager.subscribe(failing); + manager.subscribe(healthy); + + expect(() => + manager.emit({ type: "invoice.paid", invoiceId: "inv1", txHash: "tx-inv1" }), + ).not.toThrow(); + expect(healthy).toHaveBeenCalledTimes(1); + }); + + it("reports the number of active subscribers", () => { + const manager = new InvoiceNotificationSubscriptionManager(); + const unsubscribeA = manager.subscribe(vi.fn()); + manager.subscribe(vi.fn()); + + expect(manager.subscriberCount).toBe(2); + unsubscribeA(); + expect(manager.subscriberCount).toBe(1); + }); +}); diff --git a/src/auditLogger.ts b/src/auditLogger.ts index ab5c201a..beb6d317 100644 --- a/src/auditLogger.ts +++ b/src/auditLogger.ts @@ -53,6 +53,44 @@ export interface AnalyticsExportOptions { onEvent?: (event: AnalyticsExportEvent) => void; } +/** + * A single invoice notification subscription. + * + * Represents a consumer's interest in receiving notifications for a given + * invoice. Subscriptions are keyed by `invoiceId` and may be filtered by the + * notification `events` the subscriber cares about. + */ +export interface InvoiceNotificationSubscription { + /** Unique identifier for the subscription. */ + id: string; + /** The invoice this subscription is bound to. */ + invoiceId: string; + /** Callback invoked when a matching invoice notification is emitted. */ + handler: (notification: InvoiceNotification) => void; + /** Optional subset of events to receive; when omitted, all events fire. */ + events?: InvoiceNotificationEvent[]; +} + +/** Notification event types emitted for invoice lifecycle changes. */ +export type InvoiceNotificationEvent = + | "invoice_created" + | "invoice_paid" + | "invoice_settled" + | "invoice_expired" + | "invoice_cancelled"; + +/** A notification payload delivered to matching invoice subscribers. */ +export interface InvoiceNotification { + /** The invoice the notification pertains to. */ + invoiceId: string; + /** The lifecycle event that triggered the notification. */ + event: InvoiceNotificationEvent; + /** Epoch ms at which the notification was emitted. */ + timestamp: number; + /** Optional additional context for the notification. */ + data?: Record; +} + const STELLAR_ADDRESS_RE = /^G[A-Z0-9]{55}$/; /** Detect if a string value looks like base64-encoded XDR. */ @@ -65,6 +103,10 @@ export class AuditLogger { private readonly sink: (entry: AuditEntry) => void; private readonly splitAuditTrails = new Map(); private readonly entries: AuditEntry[] = []; + private readonly invoiceSubscriptions = new Map< + string, + InvoiceNotificationSubscription[] + >(); constructor(sink: (entry: AuditEntry) => void) { this.sink = sink; @@ -86,6 +128,69 @@ export class AuditLogger { ); } + /** + * Register a subscription for invoice notifications. + * + * @param subscription - The subscription to register. + * @returns An unsubscribe function that removes the subscription. + */ + subscribeToInvoice( + subscription: InvoiceNotificationSubscription, + ): () => void { + const existing = this.invoiceSubscriptions.get(subscription.invoiceId) ?? []; + existing.push(subscription); + this.invoiceSubscriptions.set(subscription.invoiceId, existing); + + return () => this.unsubscribeFromInvoice(subscription.id); + } + + /** + * Remove a previously registered invoice notification subscription by id. + * + * @returns `true` when a subscription was removed, `false` otherwise. + */ + unsubscribeFromInvoice(subscriptionId: string): boolean { + for (const [invoiceId, subs] of this.invoiceSubscriptions) { + const index = subs.findIndex((s) => s.id === subscriptionId); + if (index !== -1) { + subs.splice(index, 1); + if (subs.length === 0) { + this.invoiceSubscriptions.delete(invoiceId); + } + return true; + } + } + return false; + } + + /** + * Emit an invoice notification to all matching subscribers. + * + * Subscribers registered for the invoice receive the notification when they + * have no event filter or when their filter includes the emitted event. + * Handler errors are isolated so one failing subscriber cannot prevent + * delivery to the others. + * + * @param notification - The notification to deliver. + * @returns The number of subscribers the notification was delivered to. + */ + emitInvoiceNotification(notification: InvoiceNotification): number { + const subs = this.invoiceSubscriptions.get(notification.invoiceId); + if (!subs || subs.length === 0) return 0; + + let delivered = 0; + for (const sub of subs) { + if (sub.events && !sub.events.includes(notification.event)) continue; + try { + sub.handler(notification); + delivered += 1; + } catch { + // Isolate subscriber failures; never break notification delivery. + } + } + return delivered; + } + /** * Log an entry with automatic XDR decoding. * @@ -243,9 +348,8 @@ export class AuditLogger { onEvent?.({ type: "export_complete", export: result }); return result; - } catch (err) { - const error = err instanceof Error ? err : new Error(String(err)); - onEvent?.({ type: "export_error", error }); + } catch (error) { + onEvent?.({ type: "export_error", error: error as Error }); throw error; } } diff --git a/src/broadcaster.ts b/src/broadcaster.ts index c4791ec2..d6bad04f 100644 --- a/src/broadcaster.ts +++ b/src/broadcaster.ts @@ -5,6 +5,19 @@ import { Invoice } from "./types.js"; */ type InvoiceHandler = (invoiceId: string, invoice: Invoice) => void; +/** + * Event types emitted by the invoice notification subscription manager. + */ +export type InvoiceNotificationEvent = + | { type: "subscribed"; invoiceId: string } + | { type: "unsubscribed"; invoiceId: string } + | { type: "notified"; invoiceId: string; invoice: Invoice }; + +/** + * Handler function for invoice notification events. + */ +type NotificationEventHandler = (event: InvoiceNotificationEvent) => void; + /** * Invoice state broadcaster that publishes state changes to multiple subscribers. */ @@ -76,3 +89,116 @@ export class InvoiceStateBroadcaster { export function createInvoiceStateBroadcaster(): InvoiceStateBroadcaster { return new InvoiceStateBroadcaster(); } + +/** + * Manages invoice notification subscriptions on top of an + * {@link InvoiceStateBroadcaster}, emitting lifecycle events for + * subscribe, unsubscribe, and notify operations. + */ +export class InvoiceNotificationSubscriptionManager { + private readonly broadcaster: InvoiceStateBroadcaster; + private readonly eventHandlers: Set = new Set(); + private readonly unsubscribers: Map void>> = + new Map(); + + constructor(broadcaster: InvoiceStateBroadcaster = createInvoiceStateBroadcaster()) { + this.broadcaster = broadcaster; + } + + /** + * Register a handler for notification lifecycle events. + * + * @param handler - The event handler to register + * @returns Unsubscribe function that removes only this handler + */ + onEvent(handler: NotificationEventHandler): () => void { + this.eventHandlers.add(handler); + return () => { + this.eventHandlers.delete(handler); + }; + } + + /** + * Subscribe to invoice notifications for a specific invoice ID. + * + * @param invoiceId - The invoice ID to subscribe to + * @param handler - The handler function to call when notifications are received + * @returns Unsubscribe function that removes only this subscriber + */ + subscribe(invoiceId: string, handler: InvoiceHandler): () => void { + const unsubscribe = this.broadcaster.subscribe(invoiceId, handler); + + if (!this.unsubscribers.has(invoiceId)) { + this.unsubscribers.set(invoiceId, new Map()); + } + this.unsubscribers.get(invoiceId)!.set(handler, unsubscribe); + + this.emit({ type: "subscribed", invoiceId }); + + return () => { + const handlers = this.unsubscribers.get(invoiceId); + if (handlers) { + handlers.delete(handler); + if (handlers.size === 0) { + this.unsubscribers.delete(invoiceId); + } + } + unsubscribe(); + this.emit({ type: "unsubscribed", invoiceId }); + }; + } + + /** + * Notify all subscribers of an invoice state update. + * + * @param invoiceId - The invoice ID to notify + * @param invoice - The updated invoice state + */ + notify(invoiceId: string, invoice: Invoice): void { + this.broadcaster.broadcast(invoiceId, invoice); + this.emit({ type: "notified", invoiceId, invoice }); + } + + /** + * Get the number of subscribers for a given invoice ID. + * + * @param invoiceId - The invoice ID to check + * @returns Number of subscribers + */ + getSubscriberCount(invoiceId: string): number { + return this.broadcaster.getSubscriberCount(invoiceId); + } + + /** + * Remove all subscriptions and event handlers. + */ + clear(): void { + this.unsubscribers.forEach((handlers) => { + handlers.forEach((unsubscribe) => unsubscribe()); + }); + this.unsubscribers.clear(); + this.eventHandlers.clear(); + } + + private emit(event: InvoiceNotificationEvent): void { + this.eventHandlers.forEach((handler) => { + try { + handler(event); + } catch (error) { + console.error("Error in invoice notification event handler:", error); + } + }); + } +} + +/** + * Creates a new InvoiceNotificationSubscriptionManager instance. + * + * @param broadcaster - Optional broadcaster to use for state updates + * @returns A new InvoiceNotificationSubscriptionManager instance + */ +export function createInvoiceNotificationSubscriptionManager( + broadcaster?: InvoiceStateBroadcaster +): InvoiceNotificationSubscriptionManager { + return new InvoiceNotificationSubscriptionManager(broadcaster); +}