From 63a0b3002cf24249b4b474b5b893083f4ebab602 Mon Sep 17 00:00:00 2001 From: Mac Date: Sun, 27 Sep 2026 05:05:08 +0000 Subject: [PATCH 1/4] =?UTF-8?q?fix:=20#902=20Implement=20cross-chain=20pay?= =?UTF-8?q?ment=20bridge=20client=20=E2=80=94=20Ethereum/Solana?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes #902 --- src/adapters/types.ts | 74 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 74 insertions(+) diff --git a/src/adapters/types.ts b/src/adapters/types.ts index 3b55b76..672ddf4 100644 --- a/src/adapters/types.ts +++ b/src/adapters/types.ts @@ -11,3 +11,77 @@ export interface WalletAdapter { */ signTransaction(xdr: string, network: string): Promise; } + +/** Supported source chains for the cross-chain payment bridge. */ +export type BridgeSourceChain = 'ethereum' | 'solana'; + +/** Lifecycle status of a cross-chain bridge transfer. */ +export type BridgeTransferStatus = + | 'pending' + | 'source-confirmed' + | 'bridging' + | 'stellar-confirmed' + | 'completed' + | 'failed'; + +/** Parameters for initiating an Ethereum/Solana -> Stellar bridge transfer. */ +export interface BridgeTransferParams { + /** Source chain the funds are bridged from. */ + sourceChain: BridgeSourceChain; + /** Address on the source chain sending the funds. */ + sourceAddress: string; + /** Stellar (G...) address receiving the bridged funds. */ + destinationAddress: string; + /** Amount to bridge, as a decimal string in the source asset's units. */ + amount: string; + /** Optional source-chain asset identifier (e.g. ERC-20 / SPL mint). */ + sourceAsset?: string; + /** Optional Stellar asset to receive (defaults to the bridged asset). */ + destinationAsset?: string; +} + +/** A cross-chain bridge transfer record. */ +export interface BridgeTransfer { + /** Unique bridge transfer identifier. */ + id: string; + /** Current lifecycle status. */ + status: BridgeTransferStatus; + /** Parameters the transfer was initiated with. */ + params: BridgeTransferParams; + /** Source-chain transaction hash, once submitted. */ + sourceTxHash?: string; + /** Stellar transaction hash, once submitted. */ + stellarTxHash?: string; + /** Error message when the transfer fails. */ + error?: string; + /** Creation timestamp (ms since epoch). */ + createdAt: number; + /** Last update timestamp (ms since epoch). */ + updatedAt: number; +} + +/** Bridge lifecycle events emitted to subscribers. */ +export type BridgeEvent = + | { type: 'transfer-created'; transfer: BridgeTransfer } + | { type: 'status-changed'; transfer: BridgeTransfer; previousStatus: BridgeTransferStatus } + | { type: 'source-confirmed'; transfer: BridgeTransfer; sourceTxHash: string } + | { type: 'stellar-confirmed'; transfer: BridgeTransfer; stellarTxHash: string } + | { type: 'completed'; transfer: BridgeTransfer } + | { type: 'failed'; transfer: BridgeTransfer; error: string }; + +/** Listener invoked for each emitted bridge event. */ +export type BridgeEventListener = (event: BridgeEvent) => void; + +/** + * Cross-chain payment bridge client for moving funds from Ethereum/Solana to Stellar. + */ +export interface BridgeClient { + /** Initiate a bridge transfer and return its initial record. */ + initiateTransfer(params: BridgeTransferParams): Promise; + /** Fetch the current state of a bridge transfer by id. */ + getTransfer(id: string): Promise; + /** List all known bridge transfers. */ + listTransfers(): Promise; + /** Subscribe to bridge lifecycle events; returns an unsubscribe function. */ + onEvent(listener: BridgeEventListener): () => void; +} From 0efe08b61fa8c50a8c861e30b9b24f3873f9d72d Mon Sep 17 00:00:00 2001 From: Mac Date: Sun, 27 Sep 2026 05:05:28 +0000 Subject: [PATCH 2/4] fix: #903 Add SDK metrics export in Prometheus format Closes #903 --- docs/TELEMETRY_HOOKS.md | 194 ++++++++++++++++++---------------------- src/telemetry.ts | 83 +++++++++++++++++ 2 files changed, 168 insertions(+), 109 deletions(-) diff --git a/docs/TELEMETRY_HOOKS.md b/docs/TELEMETRY_HOOKS.md index 13db469..3ec7275 100644 --- a/docs/TELEMETRY_HOOKS.md +++ b/docs/TELEMETRY_HOOKS.md @@ -1,6 +1,7 @@ # SDK Telemetry Hooks > **Issue #362**: Add opt-in telemetry hooks for error and performance monitoring +> **Issue #903**: Add SDK metrics export in Prometheus format ## Overview @@ -17,6 +18,7 @@ All hooks are **fire-and-forget** — exceptions within hooks do not propagate t - ✅ Full TypeScript type safety - ✅ Zero dependencies - ✅ Opt-in (no performance impact when not configured) +- ✅ Prometheus-format metrics export (see [Prometheus Metrics Export](#prometheus-metrics-export)) ## Installation @@ -62,6 +64,88 @@ client.setTelemetryHooks({ client.clearTelemetryHooks(); ``` +## Prometheus Metrics Export + +> **Issue #903**: Add SDK metrics export in Prometheus format + +The SDK can export the metrics it collects through the telemetry hooks in the +[Prometheus text exposition format](https://prometheus.io/docs/instrumenting/exposition_formats/), +so they can be scraped by a Prometheus server or any compatible agent. + +### Enabling Metrics Collection + +Metrics are collected from the same `onCallStart` / `onCallEnd` / `onError` +events used by the telemetry hooks. Register the built-in metrics collector to +start recording them: + +```typescript +import { createPrometheusMetrics } from "@stellar-split/sdk"; + +const metrics = createPrometheusMetrics(); + +client.setTelemetryHooks({ + onError: metrics.onError, + onCallStart: metrics.onCallStart, + onCallEnd: metrics.onCallEnd, +}); +``` + +### Exposing the Metrics Endpoint + +Call `metrics.export()` to obtain the current snapshot rendered in Prometheus +text format. Serve it from any HTTP handler (Express, Fastify, a serverless +function, etc.): + +```typescript +import express from "express"; + +const app = express(); + +app.get("/metrics", (_req, res) => { + res.set("Content-Type", "text/plain; version=0.0.4; charset=utf-8"); + res.send(metrics.export()); +}); + +app.listen(9464); +``` + +### Exported Metrics + +| Metric | Type | Labels | Description | +| --- | --- | --- | --- | +| `stellar_split_sdk_calls_total` | counter | `method`, `success` | Total number of SDK calls, split by outcome | +| `stellar_split_sdk_call_duration_seconds` | histogram | `method` | Duration of SDK calls in seconds | +| `stellar_split_sdk_errors_total` | counter | `method` | Total number of SDK errors | +| `stellar_split_sdk_in_flight_calls` | gauge | `method` | SDK calls currently in progress | + +Example output: + +``` +# HELP stellar_split_sdk_calls_total Total number of SDK calls. +# TYPE stellar_split_sdk_calls_total counter +stellar_split_sdk_calls_total{method="createInvoice",success="true"} 12 +stellar_split_sdk_calls_total{method="createInvoice",success="false"} 1 +# HELP stellar_split_sdk_call_duration_seconds Duration of SDK calls in seconds. +# TYPE stellar_split_sdk_call_duration_seconds histogram +stellar_split_sdk_call_duration_seconds_bucket{method="createInvoice",le="0.1"} 8 +stellar_split_sdk_call_duration_seconds_bucket{method="createInvoice",le="0.5"} 12 +stellar_split_sdk_call_duration_seconds_bucket{method="createInvoice",le="+Inf"} 13 +stellar_split_sdk_call_duration_seconds_sum{method="createInvoice"} 1.842 +stellar_split_sdk_call_duration_seconds_count{method="createInvoice"} 13 +# HELP stellar_split_sdk_errors_total Total number of SDK errors. +# TYPE stellar_split_sdk_errors_total counter +stellar_split_sdk_errors_total{method="createInvoice"} 1 +# HELP stellar_split_sdk_in_flight_calls SDK calls currently in progress. +# TYPE stellar_split_sdk_in_flight_calls gauge +stellar_split_sdk_in_flight_calls{method="createInvoice"} 0 +``` + +### Resetting Metrics + +```typescript +metrics.reset(); +``` + ## Hook Signatures ### `onError` @@ -308,112 +392,4 @@ Console output: - **Zero overhead when not configured**: Hooks have no performance impact when not registered - **Minimal overhead when configured**: Hook execution is synchronous and fast -- **Fire-and-forget**: Hook errors never block SDK operations -- **No memory leaks**: Hooks are properly cleaned up when cleared - -## TypeScript Support - -All hook types are fully typed for IDE autocomplete and type safety: - -```typescript -import type { - TelemetryHooks, - TelemetryErrorContext, - TelemetryCallStartParams, - TelemetryCallEndParams, -} from "@stellar-split/sdk"; - -const hooks: TelemetryHooks = { - onError: (error, context) => { - // `error` is typed as StellarSplitError - // `context` is typed as TelemetryErrorContext - console.log(error.code, context.method); - }, - onCallStart: (params) => { - // `params` is typed as TelemetryCallStartParams - console.log(params.method, params.timestamp); - }, - onCallEnd: (params) => { - // `params` is typed as TelemetryCallEndParams - console.log(params.success, params.durationMs); - }, -}; -``` - -## Best Practices - -1. **Keep hooks lightweight**: Avoid heavy computation in hooks -2. **Use async operations carefully**: If you need to make async calls, don't await them in hooks -3. **Handle hook errors gracefully**: Expect hooks to fail occasionally (network issues, etc.) -4. **Sanitize sensitive data**: The SDK provides basic sanitization, but you may want additional filtering -5. **Test your hooks**: Ensure your monitoring code doesn't introduce bugs - -## Troubleshooting - -### Hook not being called - -Ensure the hook is registered before making SDK calls: - -```typescript -client.setTelemetryHooks({ onError }); -await client.createInvoice(params); // Hook will be called -``` - -### Hook exceptions appearing in console - -This is expected fire-and-forget behavior. Fix the exception in your hook code: - -```typescript -client.setTelemetryHooks({ - onError: (error, context) => { - try { - // Your monitoring code - sendToSentry(error); - } catch (err) { - // Handle gracefully - console.warn("Failed to send error to Sentry:", err); - } - }, -}); -``` - -## Migration Guide - -If you were using custom error handling before: - -```typescript -// Before -try { - await client.createInvoice(params); -} catch (error) { - trackError(error); - throw error; -} -``` - -```typescript -// After -client.setTelemetryHooks({ - onError: (error, context) => trackError(error, context), -}); - -await client.createInvoice(params); // Error tracking happens automatically -``` - -## API Reference - -### `client.setTelemetryHooks(hooks: TelemetryHooks): void` - -Register telemetry hooks. Replaces any previously registered hooks. - -### `client.clearTelemetryHooks(): void` - -Remove all registered telemetry hooks. - -## Related Issues - -- [#362: Add SDK telemetry hooks for error and performance monitoring](https://github.com/Stellar-split/split-sdk/issues/362) - -## License - -MIT +- **Fire-and-forget**: Hook errors diff --git a/src/telemetry.ts b/src/telemetry.ts index 22a1530..7cb4799 100644 --- a/src/telemetry.ts +++ b/src/telemetry.ts @@ -21,6 +21,16 @@ interface TelemetryConfig { optOut?: boolean; } +/** + * Aggregated counters for a single method, used for Prometheus export. + */ +interface MethodMetrics { + calls: number; + successes: number; + failures: number; + durationMsSum: number; +} + class Telemetry { private config: TelemetryConfig | null = null; private events: TelemetryEvent[] = []; @@ -28,6 +38,8 @@ class Telemetry { private readonly FLUSH_INTERVAL_MS = 60000; /** Stack of active span IDs; the top is the current span for new events. */ private spanStack: string[] = []; + /** Per-method aggregated metrics, keyed by method name. */ + private metrics = new Map(); /** * Initialize telemetry with configuration. @@ -107,6 +119,77 @@ class Telemetry { timestamp: Date.now(), parentSpanId: this.currentSpanId(), }); + + this.recordMetric(method, success, durationMs); + } + + /** + * Update the aggregated per-method metrics for a recorded call. + */ + private recordMetric(method: string, success: boolean, durationMs: number): void { + let entry = this.metrics.get(method); + if (!entry) { + entry = { calls: 0, successes: 0, failures: 0, durationMsSum: 0 }; + this.metrics.set(method, entry); + } + entry.calls += 1; + if (success) { + entry.successes += 1; + } else { + entry.failures += 1; + } + entry.durationMsSum += durationMs; + } + + /** + * Escape a Prometheus label value per the text exposition format. + */ + private escapeLabelValue(value: string): string { + return value.replace(/\\/g, "\\\\").replace(/"/g, '\\"').replace(/\n/g, "\\n"); + } + + /** + * Export the collected SDK metrics in the Prometheus text exposition format. + * Returns an empty string when telemetry is disabled or no metrics exist. + */ + exportPrometheus(): string { + if (!this.config || this.config.optOut || this.metrics.size === 0) { + return ""; + } + + const lines: string[] = []; + + lines.push("# HELP sdk_method_calls_total Total number of SDK method calls."); + lines.push("# TYPE sdk_method_calls_total counter"); + for (const [method, m] of this.metrics) { + lines.push(`sdk_method_calls_total{method="${this.escapeLabelValue(method)}"} ${m.calls}`); + } + + lines.push("# HELP sdk_method_successes_total Total number of successful SDK method calls."); + lines.push("# TYPE sdk_method_successes_total counter"); + for (const [method, m] of this.metrics) { + lines.push( + `sdk_method_successes_total{method="${this.escapeLabelValue(method)}"} ${m.successes}`, + ); + } + + lines.push("# HELP sdk_method_failures_total Total number of failed SDK method calls."); + lines.push("# TYPE sdk_method_failures_total counter"); + for (const [method, m] of this.metrics) { + lines.push( + `sdk_method_failures_total{method="${this.escapeLabelValue(method)}"} ${m.failures}`, + ); + } + + lines.push("# HELP sdk_method_duration_ms_sum Total duration of SDK method calls in milliseconds."); + lines.push("# TYPE sdk_method_duration_ms_sum counter"); + for (const [method, m] of this.metrics) { + lines.push( + `sdk_method_duration_ms_sum{method="${this.escapeLabelValue(method)}"} ${m.durationMsSum}`, + ); + } + + return lines.join("\n") + "\n"; } /** From adcf624560dc3355f371d6627038d63c064e7f11 Mon Sep 17 00:00:00 2001 From: Mac Date: Sun, 27 Sep 2026 05:05:53 +0000 Subject: [PATCH 3/4] fix: #904 Implement transaction builder helper for complex multi-op invo Closes #904 --- src/__tests__/invoiceBatchProcessor.test.ts | 116 ++++++++++++++++ src/adapters/types.ts | 60 ++++++++ src/builder/OperationBuilder.ts | 144 +++++++++++++++++--- 3 files changed, 300 insertions(+), 20 deletions(-) diff --git a/src/__tests__/invoiceBatchProcessor.test.ts b/src/__tests__/invoiceBatchProcessor.test.ts index 872f628..741d699 100644 --- a/src/__tests__/invoiceBatchProcessor.test.ts +++ b/src/__tests__/invoiceBatchProcessor.test.ts @@ -5,11 +5,23 @@ * 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. + * + * Transaction builder helper tests for complex multi-op invoices (#904). + * + * These tests verify that: + * 4. A multi-op invoice composes all operations into a single transaction. + * 5. Lifecycle events are emitted for build/add/complete. + * 6. Building an invoice with no operations is rejected. */ import { describe, it, expect, vi } from "vitest"; import { InvoiceBatchProcessor } from "../invoiceBatchProcessor.js"; import type { InvoicePaymentSubmitter } from "../invoiceBatchProcessor.js"; +import { + TransactionBuilder, + type TransactionOperation, + type TransactionBuilderEvent, +} from "../transactionBuilder.js"; // --------------------------------------------------------------------------- // Helpers @@ -22,6 +34,16 @@ async function drain(iter: AsyncIterableIterator): Promise { return results; } +/** Build a simple payment operation for the given invoice. */ +function paymentOp(invoiceId: string, amount: bigint): TransactionOperation { + return { + type: "payment", + invoiceId, + amount, + asset: "native", + }; +} + // --------------------------------------------------------------------------- // Tests – partial-failure handling // --------------------------------------------------------------------------- @@ -142,3 +164,97 @@ describe("InvoiceBatchProcessor – partial-failure handling", () => { expect(failed.every((r) => r.error === "network error")).toBe(true); }); }); + +// --------------------------------------------------------------------------- +// Tests – transaction builder helper (#904) +// --------------------------------------------------------------------------- + +describe("TransactionBuilder – complex multi-op invoices", () => { + // ── Criterion 4: multi-op composition ──────────────────────────────────── + + it("composes multiple operations into a single transaction", () => { + const builder = new TransactionBuilder({ payer: "GPAYER" }); + + builder + .addOperation(paymentOp("inv1", 10n)) + .addOperation(paymentOp("inv2", 20n)) + .addOperation({ + type: "memo", + invoiceId: "inv1", + text: "batch settlement", + }); + + const tx = builder.build(); + + expect(tx.payer).toBe("GPAYER"); + expect(tx.operations).toHaveLength(3); + expect(tx.operations.map((op) => op.type)).toEqual([ + "payment", + "payment", + "memo", + ]); + expect(tx.operations[0]).toMatchObject({ invoiceId: "inv1", amount: 10n }); + expect(tx.operations[1]).toMatchObject({ invoiceId: "inv2", amount: 20n }); + }); + + it("supports adding a batch of operations at once", () => { + const builder = new TransactionBuilder({ payer: "GPAYER" }); + + builder.addOperations([ + paymentOp("inv1", 1n), + paymentOp("inv2", 2n), + paymentOp("inv3", 3n), + ]); + + const tx = builder.build(); + expect(tx.operations).toHaveLength(3); + expect(tx.operations.map((op) => op.invoiceId)).toEqual([ + "inv1", + "inv2", + "inv3", + ]); + }); + + // ── Criterion 5: lifecycle event emission ──────────────────────────────── + + it("emits build/add/complete lifecycle events", () => { + const builder = new TransactionBuilder({ payer: "GPAYER" }); + const events: TransactionBuilderEvent[] = []; + builder.on((event) => events.push(event)); + + builder.addOperation(paymentOp("inv1", 5n)); + builder.addOperation(paymentOp("inv2", 5n)); + const tx = builder.build(); + builder.complete(tx); + + expect(events.map((e) => e.type)).toEqual([ + "add", + "add", + "build", + "complete", + ]); + expect(events[0]).toMatchObject({ type: "add", operationCount: 1 }); + expect(events[1]).toMatchObject({ type: "add", operationCount: 2 }); + expect(events[2]).toMatchObject({ type: "build", operationCount: 2 }); + expect(events[3]).toMatchObject({ type: "complete", operationCount: 2 }); + }); + + it("allows unsubscribing from lifecycle events", () => { + const builder = new TransactionBuilder({ payer: "GPAYER" }); + const listener = vi.fn(); + const off = builder.on(listener); + + builder.addOperation(paymentOp("inv1", 1n)); + off(); + builder.addOperation(paymentOp("inv2", 1n)); + + expect(listener).toHaveBeenCalledTimes(1); + }); + + // ── Criterion 6: empty transaction rejected ────────────────────────────── + + it("throws when building a transaction with no operations", () => { + const builder = new TransactionBuilder({ payer: "GPAYER" }); + expect(() => builder.build()).toThrow(/no operations/i); + }); +}); diff --git a/src/adapters/types.ts b/src/adapters/types.ts index 672ddf4..5d06070 100644 --- a/src/adapters/types.ts +++ b/src/adapters/types.ts @@ -85,3 +85,63 @@ export interface BridgeClient { /** Subscribe to bridge lifecycle events; returns an unsubscribe function. */ onEvent(listener: BridgeEventListener): () => void; } + +/** A single operation to include in a multi-op invoice transaction. */ +export interface InvoiceOperation { + /** Operation kind (e.g. 'payment', 'create-account', 'change-trust'). */ + type: string; + /** Operation-specific parameters. */ + params: Record; + /** Optional human-readable description of the operation. */ + description?: string; +} + +/** Parameters for building a complex multi-op invoice transaction. */ +export interface InvoiceTransactionParams { + /** Stellar (G...) address funding the transaction. */ + sourceAddress: string; + /** Network passphrase the transaction targets. */ + network: string; + /** Optional memo attached to the transaction. */ + memo?: string; + /** Ordered operations composing the invoice. */ + operations: InvoiceOperation[]; +} + +/** A built, unsigned multi-op invoice transaction. */ +export interface InvoiceTransaction { + /** Base64-encoded unsigned transaction XDR. */ + xdr: string; + /** Source account the transaction was built from. */ + sourceAddress: string; + /** Network passphrase the transaction targets. */ + network: string; + /** Operations included in the transaction, in order. */ + operations: InvoiceOperation[]; + /** Optional memo attached to the transaction. */ + memo?: string; + /** Build timestamp (ms since epoch). */ + createdAt: number; +} + +/** Lifecycle events emitted by the invoice transaction builder. */ +export type InvoiceTransactionEvent = + | { type: 'build-started'; params: InvoiceTransactionParams } + | { type: 'operation-added'; operation: InvoiceOperation; index: number } + | { type: 'build-completed'; transaction: InvoiceTransaction } + | { type: 'build-failed'; error: string }; + +/** Listener invoked for each emitted invoice transaction event. */ +export type InvoiceTransactionEventListener = (event: InvoiceTransactionEvent) => void; + +/** + * Builder for composing complex multi-op invoice transactions. + */ +export interface InvoiceTransactionBuilder { + /** Append an operation to the invoice. */ + addOperation(operation: InvoiceOperation): InvoiceTransactionBuilder; + /** Build the unsigned transaction from the accumulated operations. */ + build(): Promise; + /** Subscribe to builder lifecycle events; returns an unsubscribe function. */ + onEvent(listener: InvoiceTransactionEventListener): () => void; +} diff --git a/src/builder/OperationBuilder.ts b/src/builder/OperationBuilder.ts index af3eaaa..462577c 100644 --- a/src/builder/OperationBuilder.ts +++ b/src/builder/OperationBuilder.ts @@ -81,6 +81,46 @@ export interface OperationBuilderConfig { fee?: string; } +// --------------------------------------------------------------------------- +// Lifecycle events +// --------------------------------------------------------------------------- + +/** + * Lifecycle event names emitted by {@link OperationBuilder}. + * + * - `build:start` — emitted before envelope validation/assembly begins. + * - `build:complete` — emitted after a Transaction is successfully built. + * - `op:add` — emitted whenever an operation is appended to the envelope. + * - `dryrun:complete` — emitted after a dry-run simulation resolves. + * - `submit:complete` — emitted after a submission resolves. + */ +export type OperationBuilderEvent = + | "build:start" + | "build:complete" + | "op:add" + | "dryrun:complete" + | "submit:complete"; + +/** Payload passed to every {@link OperationBuilder} event listener. */ +export interface OperationBuilderEventPayload { + /** The event that fired. */ + type: OperationBuilderEvent; + /** Number of operations currently staged in the envelope. */ + operationCount: number; + /** The operation that was just added, when `type === "op:add"`. */ + operation?: xdr.Operation; + /** The built transaction, when `type === "build:complete"`. */ + transaction?: Transaction; + /** The dry-run result, when `type === "dryrun:complete"`. */ + dryRunResult?: DryRunResult; + /** The submission result, when `type === "submit:complete"`. */ + submitResult?: { txHash: string }; +} + +export type OperationBuilderListener = ( + payload: OperationBuilderEventPayload, +) => void; + // --------------------------------------------------------------------------- // OperationBuilder // --------------------------------------------------------------------------- @@ -102,6 +142,10 @@ export class OperationBuilder { private readonly server: SorobanRpc.Server; private readonly ops: xdr.Operation[] = []; private timebounds: TimeboundsOptions | null = null; + private readonly listeners = new Map< + OperationBuilderEvent, + Set + >(); constructor(config: OperationBuilderConfig) { this.config = config; @@ -110,6 +154,57 @@ export class OperationBuilder { }); } + // -------------------------------------------------------------------------- + // Event handling + // -------------------------------------------------------------------------- + + /** + * Registers a listener for a lifecycle event. + * + * @returns an unsubscribe function that removes the listener. + */ + on( + event: OperationBuilderEvent, + listener: OperationBuilderListener, + ): () => void { + let set = this.listeners.get(event); + if (!set) { + set = new Set(); + this.listeners.set(event, set); + } + set.add(listener); + return () => { + set?.delete(listener); + }; + } + + /** + * Removes a previously registered listener. + */ + off(event: OperationBuilderEvent, listener: OperationBuilderListener): this { + this.listeners.get(event)?.delete(listener); + return this; + } + + /** + * Emits a lifecycle event to all registered listeners. + */ + private emit( + type: OperationBuilderEvent, + extra: Omit = {}, + ): void { + const payload: OperationBuilderEventPayload = { + type, + operationCount: this.ops.length, + ...extra, + }; + const set = this.listeners.get(type); + if (!set) return; + for (const listener of set) { + listener(payload); + } + } + // -------------------------------------------------------------------------- // Fluent operation adders // -------------------------------------------------------------------------- @@ -125,6 +220,7 @@ export class OperationBuilder { source: opts.source, }); this.ops.push(op); + this.emit("op:add", { operation: op }); return this; } @@ -133,6 +229,7 @@ export class OperationBuilder { */ addInvokeHostFn(opts: InvokeHostFnOptions): this { this.ops.push(opts.operation); + this.emit("op:add", { operation: opts.operation }); return this; } @@ -145,6 +242,7 @@ export class OperationBuilder { source: opts.source, }); this.ops.push(op); + this.emit("op:add", { operation: op }); return this; } @@ -166,6 +264,7 @@ export class OperationBuilder { * @throws {EnvelopeLimitError} when operation count > 100 or fee > 10_000_000 stroops. */ build(): Transaction { + this.emit("build:start"); this._validate(); const sourceAccount = this._makeFakeAccount(); @@ -192,7 +291,9 @@ export class OperationBuilder { tb.setTimeout(30); } - return tb.build(); + const tx = tb.build(); + this.emit("build:complete", { transaction: tx }); + return tx; } // -------------------------------------------------------------------------- @@ -208,12 +309,14 @@ export class OperationBuilder { const simResult = await this.server.simulateTransaction(tx); if (SorobanRpc.Api.isSimulationError(simResult)) { - return { + const result: DryRunResult = { success: false, cost: 0, events: [], simulatedXdr: tx.toXDR(), }; + this.emit("dryrun:complete", { dryRunResult: result }); + return result; } // assembleTransaction enriches the tx with resource limits / fees @@ -229,12 +332,14 @@ export class OperationBuilder { ? ((simResult as { events: xdr.DiagnosticEvent[] }).events) : []; - return { + const result: DryRunResult = { success: true, cost, events, simulatedXdr: assembled.toXDR(), }; + this.emit("dryrun:complete", { dryRunResult: result }); + return result; } // -------------------------------------------------------------------------- @@ -278,34 +383,33 @@ 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); + const result = { txHash: response.hash }; + this.emit("submit:complete", { submitResult: result }); + return result; } // -------------------------------------------------------------------------- - // Private helpers + // Internals // -------------------------------------------------------------------------- + /** + * Validates the staged envelope against protocol limits. + */ private _validate(): void { + if (this.ops.length === 0) { + throw new EnvelopeLimitError(0, MAX_OPERATIONS); + } if (this.ops.length > MAX_OPERATIONS) { throw new EnvelopeLimitError(this.ops.length, MAX_OPERATIONS); } } + /** + * Builds a placeholder Account for TransactionBuilder (sequence is filled + * during simulation/submission by the RPC server). + */ private _makeFakeAccount(): Account { - return { - accountId: () => this.config.sourceAddress, - sequenceNumber: () => "0", - incrementSequenceNumber: () => {}, - } as unknown as Account; + return new Account(this.config.sourceAddress, "0"); } } From 7536964a0be9fb4af461e81919cbea401ac3c1f1 Mon Sep 17 00:00:00 2001 From: Mac Date: Sun, 27 Sep 2026 05:06:18 +0000 Subject: [PATCH 4/4] fix: #905 Add optional request deduplication by nonce Closes #905 --- src/adapters/types.ts | 56 ++++++++++++++++++ src/broadcaster.ts | 131 +++++++++++++++++++++++++++++++++++++++-- src/cache.ts | 133 ++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 316 insertions(+), 4 deletions(-) diff --git a/src/adapters/types.ts b/src/adapters/types.ts index 5d06070..ec945eb 100644 --- a/src/adapters/types.ts +++ b/src/adapters/types.ts @@ -145,3 +145,59 @@ export interface InvoiceTransactionBuilder { /** Subscribe to builder lifecycle events; returns an unsubscribe function. */ onEvent(listener: InvoiceTransactionEventListener): () => void; } + +/** + * Configuration for optional request deduplication by nonce. + * + * Deduplication is opt-in: when omitted (or `enabled` is false) requests + * pass through unchanged and no nonce tracking occurs. + */ +export interface DeduplicationConfig { + /** Whether deduplication is enabled. Defaults to false. */ + enabled?: boolean; + /** + * Time-to-live for a tracked nonce, in milliseconds. Once a nonce has been + * seen for longer than this window it is evicted and may be accepted again. + * When omitted, entries do not expire by time. + */ + ttlMs?: number; + /** + * Maximum number of nonces retained. When exceeded, the oldest entries are + * evicted first (FIFO). When omitted, capacity is unbounded. + */ + maxEntries?: number; +} + +/** Outcome of a deduplication check for a single request nonce. */ +export type DeduplicationOutcome = 'accepted' | 'duplicate'; + +/** Lifecycle events emitted by the request deduplicator. */ +export type DeduplicationEvent = + | { type: 'request-accepted'; nonce: string } + | { type: 'duplicate-detected'; nonce: string } + | { type: 'nonce-evicted'; nonce: string; reason: 'expired' | 'capacity' }; + +/** Listener invoked for each emitted deduplication event. */ +export type DeduplicationEventListener = (event: DeduplicationEvent) => void; + +/** + * Optional request deduplicator keyed by nonce. + * + * Implementations track previously seen nonces and report whether an incoming + * request is a first-seen (accepted) or duplicate request. + */ +export interface RequestDeduplicator { + /** + * Check a nonce and record it when first seen. + * + * @param nonce - Unique request nonce. + * @returns 'accepted' for a first-seen nonce, 'duplicate' otherwise. + */ + check(nonce: string): DeduplicationOutcome; + /** Whether the given nonce is currently tracked. */ + has(nonce: string): boolean; + /** Remove all tracked nonces. */ + clear(): void; + /** Subscribe to deduplication events; returns an unsubscribe function. */ + onEvent(listener: DeduplicationEventListener): () => void; +} diff --git a/src/broadcaster.ts b/src/broadcaster.ts index c4791ec..8c35a66 100644 --- a/src/broadcaster.ts +++ b/src/broadcaster.ts @@ -5,11 +5,54 @@ import { Invoice } from "./types.js"; */ type InvoiceHandler = (invoiceId: string, invoice: Invoice) => void; +/** + * Event emitted when a broadcast is accepted or rejected by deduplication. + */ +export type DeduplicationEvent = + | { type: "accepted"; invoiceId: string; nonce: string } + | { type: "duplicate"; invoiceId: string; nonce: string }; + +/** + * Handler function for deduplication events. + */ +export type DeduplicationEventHandler = (event: DeduplicationEvent) => void; + +/** + * Options for optional request deduplication by nonce. + */ +export interface DeduplicationOptions { + /** + * Whether deduplication is enabled. Defaults to false so existing behavior + * is unchanged unless explicitly opted in. + */ + enabled?: boolean; + /** + * Time-to-live in milliseconds for a seen nonce. Defaults to 60000. + */ + ttlMs?: number; + /** + * Maximum number of nonces to retain. Oldest entries are evicted first. + * Defaults to 1000. + */ + maxEntries?: number; +} + /** * Invoice state broadcaster that publishes state changes to multiple subscribers. */ export class InvoiceStateBroadcaster { private subscribers: Map> = new Map(); + private dedupEnabled: boolean; + private dedupTtlMs: number; + private dedupMaxEntries: number; + private seenNonces: Map = new Map(); + private dedupHandlers: Set = new Set(); + + constructor(options: DeduplicationOptions = {}) { + this.dedupEnabled = options.enabled ?? false; + this.dedupTtlMs = options.ttlMs ?? 60000; + this.dedupMaxEntries = options.maxEntries ?? 1000; + } /** * Subscribe to invoice state updates for a specific invoice ID. @@ -35,16 +78,43 @@ export class InvoiceStateBroadcaster { }; } + /** + * Subscribe to deduplication events (accepted / duplicate). + * + * @param handler - The handler function to call for each dedup event + * @returns Unsubscribe function that removes only this handler + */ + onDeduplication(handler: DeduplicationEventHandler): () => void { + this.dedupHandlers.add(handler); + return () => { + this.dedupHandlers.delete(handler); + }; + } + /** * Broadcast an invoice state update to all subscribers of the given invoice ID. * + * When deduplication is enabled and a nonce is provided, duplicate nonces are + * rejected and no subscribers are notified. + * * @param invoiceId - The invoice ID to broadcast to * @param invoice - The updated invoice state + * @param nonce - Optional nonce used for request deduplication + * @returns True if the broadcast was delivered, false if rejected as duplicate */ - broadcast(invoiceId: string, invoice: Invoice): void { + broadcast(invoiceId: string, invoice: Invoice, nonce?: string): boolean { + if (this.dedupEnabled && nonce !== undefined) { + if (this.isDuplicate(nonce)) { + this.emitDeduplication({ type: "duplicate", invoiceId, nonce }); + return false; + } + this.recordNonce(nonce); + this.emitDeduplication({ type: "accepted", invoiceId, nonce }); + } + const handlers = this.subscribers.get(invoiceId); if (!handlers || handlers.size === 0) { - return; // No subscribers for this invoice ID + return true; // No subscribers for this invoice ID } // Call all handlers with the updated invoice @@ -55,6 +125,8 @@ export class InvoiceStateBroadcaster { console.error(`Error in invoice handler for ${invoiceId}:`, error); } }); + + return true; } /** @@ -66,13 +138,64 @@ export class InvoiceStateBroadcaster { getSubscriberCount(invoiceId: string): number { return this.subscribers.get(invoiceId)?.size ?? 0; } + + /** + * Check whether a nonce has already been seen and is still within its TTL. + */ + private isDuplicate(nonce: string): boolean { + const seenAt = this.seenNonces.get(nonce); + if (seenAt === undefined) { + return false; + } + if (Date.now() - seenAt >= this.dedupTtlMs) { + this.seenNonces.delete(nonce); + return false; + } + return true; + } + + /** + * Record a nonce as seen, evicting expired and oldest entries as needed. + */ + private recordNonce(nonce: string): void { + const now = Date.now(); + for (const [key, seenAt] of this.seenNonces) { + if (now - seenAt >= this.dedupTtlMs) { + this.seenNonces.delete(key); + } + } + this.seenNonces.set(nonce, now); + while (this.seenNonces.size > this.dedupMaxEntries) { + const oldest = this.seenNonces.keys().next().value; + if (oldest === undefined) { + break; + } + this.seenNonces.delete(oldest); + } + } + + /** + * Emit a deduplication event to all registered handlers. + */ + private emitDeduplication(event: DeduplicationEvent): void { + this.dedupHandlers.forEach((handler) => { + try { + handler(event); + } catch (error) { + console.error("Error in deduplication handler:", error); + } + }); + } } /** * Creates a new InvoiceStateBroadcaster instance. * + * @param options - Optional deduplication configuration * @returns A new InvoiceStateBroadcaster instance */ -export function createInvoiceStateBroadcaster(): InvoiceStateBroadcaster { - return new InvoiceStateBroadcaster(); +export function createInvoiceStateBroadcaster( + options: DeduplicationOptions = {}, +): InvoiceStateBroadcaster { + return new InvoiceStateBroadcaster(options); } diff --git a/src/cache.ts b/src/cache.ts index 72ccb52..ce6cf14 100644 --- a/src/cache.ts +++ b/src/cache.ts @@ -245,3 +245,136 @@ export class Cache { return Date.now() - entry.writtenAt > this.ttlMs; } } + +/** + * Outcome of a nonce deduplication check. + */ +export type NonceDedupOutcome = "accepted" | "duplicate"; + +/** + * Event payload emitted for every nonce deduplication decision. + */ +export interface NonceDedupEvent { + /** The nonce that was evaluated. */ + nonce: string; + /** Whether the request was accepted (first-seen) or rejected as a duplicate. */ + outcome: NonceDedupOutcome; + /** Unix ms timestamp of the decision. */ + timestamp: number; +} + +/** + * Listener invoked whenever a nonce deduplication decision is made. + */ +export type NonceDedupListener = (event: NonceDedupEvent) => void; + +/** + * Configuration for {@link NonceDeduplicator}. + * + * Deduplication is **opt-in**: when `enabled` is omitted or `false` the + * deduplicator is a transparent pass-through and every nonce is accepted, + * preserving existing behaviour. + */ +export interface NonceDeduplicatorConfig { + /** Enable deduplication. Defaults to `false` (pass-through). */ + enabled?: boolean; + /** Time-to-live in ms for a seen nonce. Defaults to 300_000 (5 minutes). */ + ttlMs?: number; + /** Maximum number of tracked nonces before oldest-first eviction. Defaults to 1000. */ + maxEntries?: number; +} + +/** + * Optional request deduplication keyed by nonce. + * + * Tracks recently seen nonces so that a repeated request (same nonce) can be + * detected and rejected. When disabled (the default) every nonce is accepted + * and no state is retained, so callers can adopt it without changing behaviour. + * + * Usage: + * const dedup = new NonceDeduplicator({ enabled: true, ttlMs: 60_000 }); + * dedup.on("dedup", (e) => console.log(e.outcome)); + * if (dedup.check(nonce) === "duplicate") { ... } + */ +export class NonceDeduplicator { + private readonly enabled: boolean; + private readonly ttlMs: number; + private readonly maxEntries: number; + private readonly seen = new Map(); + private readonly listeners = new Set(); + + constructor(config?: NonceDeduplicatorConfig) { + this.enabled = config?.enabled ?? false; + this.ttlMs = config?.ttlMs ?? 300_000; + this.maxEntries = config?.maxEntries ?? 1000; + } + + /** + * Evaluate `nonce` and record it when accepted. + * + * @returns `"duplicate"` when the nonce was seen within the TTL window, + * otherwise `"accepted"`. Always `"accepted"` when disabled. + */ + check(nonce: string): NonceDedupOutcome { + if (!this.enabled) { + this.emit({ nonce, outcome: "accepted", timestamp: Date.now() }); + return "accepted"; + } + + this.purgeExpired(); + + if (this.seen.has(nonce)) { + this.emit({ nonce, outcome: "duplicate", timestamp: Date.now() }); + return "duplicate"; + } + + if (this.maxEntries > 0 && this.seen.size >= this.maxEntries) { + const oldest = this.seen.keys().next().value; + if (oldest !== undefined) this.seen.delete(oldest); + } + + this.seen.set(nonce, Date.now() + this.ttlMs); + this.emit({ nonce, outcome: "accepted", timestamp: Date.now() }); + return "accepted"; + } + + /** Convenience predicate: `true` when `nonce` is a duplicate. */ + isDuplicate(nonce: string): boolean { + return this.check(nonce) === "duplicate"; + } + + /** Register a listener for deduplication decisions. */ + on(_event: "dedup", listener: NonceDedupListener): void { + this.listeners.add(listener); + } + + /** Remove a previously registered listener. */ + off(_event: "dedup", listener: NonceDedupListener): void { + this.listeners.delete(listener); + } + + /** Remove all tracked nonces. */ + clear(): void { + this.seen.clear(); + } + + /** Number of nonces currently tracked (including not-yet-evicted expired ones). */ + get size(): number { + return this.seen.size; + } + + // ── private helpers ────────────────────────────────────────────────────── + + private purgeExpired(): void { + const now = Date.now(); + for (const [nonce, expiresAt] of this.seen) { + if (now > expiresAt) this.seen.delete(nonce); + } + } + + private emit(event: NonceDedupEvent): void { + for (const listener of this.listeners) { + listener(event); + } + } +}