From 5598aa9433651af39c7f5355d2d974629acda0f3 Mon Sep 17 00:00:00 2001 From: James21 Date: Sun, 27 Sep 2026 05:29:17 +0000 Subject: [PATCH 1/4] fix: #918 Implement custody account management helpers Closes #918 --- src/accounts/AccountMergeDetector.ts | 98 +++++++++ src/accounts/AccountSignerWeightCalculator.ts | 194 ++++++++++++++++++ 2 files changed, 292 insertions(+) diff --git a/src/accounts/AccountMergeDetector.ts b/src/accounts/AccountMergeDetector.ts index 092de34..89e19da 100644 --- a/src/accounts/AccountMergeDetector.ts +++ b/src/accounts/AccountMergeDetector.ts @@ -36,9 +36,32 @@ export interface MergeEventPayload { mergedAt: Date; } +/** + * A custody account managed by the detector. Custody accounts are watched + * accounts whose funds are held on behalf of a recipient and which may be + * merged into a destination account. + */ +export interface CustodyAccount { + /** The custody account address */ + address: string; + /** Optional human-readable label */ + label?: string; + /** Optional asset the custody account is expected to hold */ + asset?: { code: string; issuer: string }; + /** When the custody account was registered */ + registeredAt: Date; +} + +/** Payload emitted on custody account lifecycle events. */ +export interface CustodyAccountEventPayload { + account: CustodyAccount; + at: Date; +} + export class AccountMergeDetector extends EventEmitter { private watchedAccounts = new Set(); private mergeCache = new Map(); // source -> destination mapping + private custodyAccounts = new Map(); private streamActive = false; private checkInterval: NodeJS.Timeout | null = null; @@ -89,6 +112,72 @@ export class AccountMergeDetector extends EventEmitter { this.watchedAccounts.delete(accountId); } + /** + * Register a custody account and begin watching it for merges. + * Emits "custody:registered" with the created custody account. + */ + registerCustodyAccount( + address: string, + options: { label?: string; asset?: { code: string; issuer: string } } = {}, + ): CustodyAccount { + const existing = this.custodyAccounts.get(address); + if (existing) { + return existing; + } + + const account: CustodyAccount = { + address, + label: options.label, + asset: options.asset, + registeredAt: new Date(), + }; + + this.custodyAccounts.set(address, account); + this.watchAccount(address); + + this.emit("custody:registered", { + account, + at: account.registeredAt, + } satisfies CustodyAccountEventPayload); + + return account; + } + + /** + * Remove a custody account from management and stop watching it. + * Emits "custody:removed" when an account was actually removed. + */ + removeCustodyAccount(address: string): boolean { + const account = this.custodyAccounts.get(address); + if (!account) { + return false; + } + + this.custodyAccounts.delete(address); + this.unwatchAccount(address); + + this.emit("custody:removed", { + account, + at: new Date(), + } satisfies CustodyAccountEventPayload); + + return true; + } + + /** + * Retrieve a managed custody account by address. + */ + getCustodyAccount(address: string): CustodyAccount | undefined { + return this.custodyAccounts.get(address); + } + + /** + * List all managed custody accounts. + */ + listCustodyAccounts(): CustodyAccount[] { + return Array.from(this.custodyAccounts.values()); + } + /** * Check if an account has been merged and resolve the final destination. * Supports recursive merge chains up to depth 5. @@ -200,6 +289,15 @@ export class AccountMergeDetector extends EventEmitter { mergedAt: event.timestamp, } satisfies MergeEventPayload); + // If the merged source was a managed custody account, emit a lifecycle event + const custodyAccount = this.custodyAccounts.get(sourceAccount); + if (custodyAccount) { + this.emit("custody:merged", { + account: custodyAccount, + at: event.timestamp, + } satisfies CustodyAccountEventPayload); + } + // Notify the client to reroute recipients try { // The client will handle rerouting via rerouteRecipient method diff --git a/src/accounts/AccountSignerWeightCalculator.ts b/src/accounts/AccountSignerWeightCalculator.ts index 09f9a4b..1c479cc 100644 --- a/src/accounts/AccountSignerWeightCalculator.ts +++ b/src/accounts/AccountSignerWeightCalculator.ts @@ -31,6 +31,64 @@ export interface SignerWeightResult { missingWeight: number; } +/** + * A single custody account signer entry, as returned by the custody account + * management helpers. + */ +export interface CustodySigner { + /** The signer's public key (G… address, pre-auth tx, or hash(x)). */ + key: string; + /** The signing weight assigned to this signer. */ + weight: number; +} + +/** + * A snapshot of a custody account's signer configuration and thresholds. + */ +export interface CustodyAccount { + /** The Stellar account G… address. */ + accountId: string; + /** The account's current signers. */ + signers: CustodySigner[]; + /** The account's threshold configuration. */ + thresholds: { + low: number; + medium: number; + high: number; + }; +} + +/** + * Event names emitted by the custody account management helpers during + * lifecycle operations. + */ +export type CustodyAccountEvent = + | "signer:added" + | "signer:removed" + | "signer:updated" + | "threshold:updated" + | "account:loaded"; + +/** + * Payload delivered to custody account event listeners. + */ +export interface CustodyAccountEventPayload { + /** The account the event pertains to. */ + accountId: string; + /** The lifecycle event that occurred. */ + event: CustodyAccountEvent; + /** The signer affected by the event, when applicable. */ + signer?: CustodySigner; + /** The threshold level affected by the event, when applicable. */ + thresholdLevel?: ThresholdLevel; + /** The previous value before the change, when applicable. */ + previousValue?: number; + /** The new value after the change, when applicable. */ + newValue?: number; +} + +export type CustodyAccountEventListener = (payload: CustodyAccountEventPayload) => void; + // --------------------------------------------------------------------------- // Cache entry // --------------------------------------------------------------------------- @@ -49,6 +107,7 @@ export class AccountSignerWeightCalculator { /** Cache TTL in milliseconds (default 30 seconds). */ private readonly cacheTtlMs: number; private readonly cache = new Map(); + private readonly listeners = new Set(); constructor(horizonUrl: string, cacheTtlMs = 30_000) { this.server = new Horizon.Server(horizonUrl, { allowHttp: horizonUrl.startsWith("http://") }); @@ -124,10 +183,145 @@ export class AccountSignerWeightCalculator { return result.sufficient; } + // -------------------------------------------------------------------------- + // Custody account management helpers + // -------------------------------------------------------------------------- + + /** + * Load a custody account snapshot (signers + thresholds) from Horizon. + * Emits an `account:loaded` event on success. + */ + async loadCustodyAccount(accountId: string): Promise { + const record = await this._loadAccount(accountId); + const account = this._toCustodyAccount(accountId, record); + this._emit({ accountId, event: "account:loaded" }); + return account; + } + + /** + * Add a signer to a custody account snapshot. If the signer already exists + * its weight is updated instead. Emits `signer:added` or `signer:updated`. + */ + addSigner(account: CustodyAccount, signer: CustodySigner): CustodyAccount { + const existing = account.signers.find((s) => s.key === signer.key); + let next: CustodyAccount; + + if (existing) { + next = { + ...account, + signers: account.signers.map((s) => (s.key === signer.key ? { ...signer } : s)), + }; + this._emit({ + accountId: account.accountId, + event: "signer:updated", + signer: { ...signer }, + previousValue: existing.weight, + newValue: signer.weight, + }); + } else { + next = { ...account, signers: [...account.signers, { ...signer }] }; + this._emit({ + accountId: account.accountId, + event: "signer:added", + signer: { ...signer }, + newValue: signer.weight, + }); + } + + return next; + } + + /** + * Remove a signer from a custody account snapshot by public key. + * Emits `signer:removed` when a signer was actually removed. + */ + removeSigner(account: CustodyAccount, signerKey: string): CustodyAccount { + const existing = account.signers.find((s) => s.key === signerKey); + if (!existing) { + return account; + } + + const next: CustodyAccount = { + ...account, + signers: account.signers.filter((s) => s.key !== signerKey), + }; + + this._emit({ + accountId: account.accountId, + event: "signer:removed", + signer: { ...existing }, + previousValue: existing.weight, + }); + + return next; + } + + /** + * Update a threshold level on a custody account snapshot. + * Emits `threshold:updated` when the value changes. + */ + updateThreshold( + account: CustodyAccount, + level: ThresholdLevel, + value: number, + ): CustodyAccount { + const previousValue = account.thresholds[level]; + if (previousValue === value) { + return account; + } + + const next: CustodyAccount = { + ...account, + thresholds: { ...account.thresholds, [level]: value }, + }; + + this._emit({ + accountId: account.accountId, + event: "threshold:updated", + thresholdLevel: level, + previousValue, + newValue: value, + }); + + return next; + } + + /** + * Register a listener for custody account lifecycle events. + * Returns an unsubscribe function. + */ + onCustodyAccountEvent(listener: CustodyAccountEventListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + // -------------------------------------------------------------------------- // Private helpers // -------------------------------------------------------------------------- + private _emit(payload: CustodyAccountEventPayload): void { + for (const listener of this.listeners) { + listener(payload); + } + } + + private _toCustodyAccount( + accountId: string, + record: Horizon.AccountResponse, + ): CustodyAccount { + return { + accountId, + signers: record.signers.map((s) => ({ key: s.key, weight: s.weight })), + thresholds: { + low: record.thresholds.low_threshold, + medium: record.thresholds.med_threshold, + high: record.thresholds.high_threshold, + }, + }; + } + private async _loadAccount(accountId: string): Promise { const now = Date.now(); const entry = this.cache.get(accountId); From 95e9dc1d2f14c08ac8efbf6840545df7d15187ff Mon Sep 17 00:00:00 2001 From: James21 Date: Sun, 27 Sep 2026 05:29:27 +0000 Subject: [PATCH 2/4] fix: #919 Add SDK fee estimation with historical analysis Closes #919 --- src/broadcaster.ts | 142 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 142 insertions(+) diff --git a/src/broadcaster.ts b/src/broadcaster.ts index c4791ec..57f6f5e 100644 --- a/src/broadcaster.ts +++ b/src/broadcaster.ts @@ -5,6 +5,42 @@ import { Invoice } from "./types.js"; */ type InvoiceHandler = (invoiceId: string, invoice: Invoice) => void; +/** + * A single historical fee observation used to inform fee estimates. + */ +export interface FeeHistoryEntry { + /** Fee rate observed, in the smallest fee unit (e.g. sat/vB). */ + feeRate: number; + /** Timestamp (ms since epoch) when the observation was recorded. */ + timestamp: number; +} + +/** + * Result of an SDK fee estimation, including the historical analysis used. + */ +export interface FeeEstimate { + /** Estimated fee rate, in the smallest fee unit (e.g. sat/vB). */ + feeRate: number; + /** Minimum fee rate observed in the analyzed history. */ + minFeeRate: number; + /** Maximum fee rate observed in the analyzed history. */ + maxFeeRate: number; + /** Average fee rate across the analyzed history. */ + averageFeeRate: number; + /** Number of historical samples that informed the estimate. */ + sampleCount: number; +} + +/** + * Event names emitted during the fee estimation lifecycle. + */ +export type FeeEstimationEvent = "estimate" | "error"; + +/** + * Handler invoked when a fee estimation event is emitted. + */ +type FeeEstimationHandler = (event: FeeEstimationEvent, payload: FeeEstimate | Error) => void; + /** * Invoice state broadcaster that publishes state changes to multiple subscribers. */ @@ -76,3 +112,109 @@ export class InvoiceStateBroadcaster { export function createInvoiceStateBroadcaster(): InvoiceStateBroadcaster { return new InvoiceStateBroadcaster(); } + +/** + * Estimates SDK fees using historical fee observations. + * + * The estimate is derived from the provided history: the most recent + * observation is weighted against the historical average so that recent + * network conditions inform the result without discarding past data. + * Emits lifecycle events so callers can react to estimates and errors. + */ +export class FeeEstimator { + private history: FeeHistoryEntry[] = []; + private handlers: Set = new Set(); + + /** + * Subscribe to fee estimation lifecycle events. + * + * @param handler - Handler invoked on "estimate" and "error" events + * @returns Unsubscribe function that removes only this handler + */ + on(handler: FeeEstimationHandler): () => void { + this.handlers.add(handler); + return () => { + this.handlers.delete(handler); + }; + } + + /** + * Record a historical fee observation. + * + * @param entry - The fee history entry to record + */ + record(entry: FeeHistoryEntry): void { + this.history.push(entry); + } + + /** + * Get a copy of the recorded fee history. + * + * @returns The recorded fee history entries + */ + getHistory(): FeeHistoryEntry[] { + return [...this.history]; + } + + /** + * Estimate the current fee rate using historical analysis. + * + * @param windowSize - Optional number of most recent samples to analyze + * @returns The fee estimate, or null when no history is available + */ + estimate(windowSize?: number): FeeEstimate | null { + try { + const samples = + windowSize && windowSize > 0 + ? this.history.slice(-windowSize) + : this.history; + + if (samples.length === 0) { + return null; + } + + const rates = samples.map((entry) => entry.feeRate); + const minFeeRate = Math.min(...rates); + const maxFeeRate = Math.max(...rates); + const averageFeeRate = + rates.reduce((sum, rate) => sum + rate, 0) / rates.length; + + // Weight the most recent observation against the historical average. + const latest = samples[samples.length - 1].feeRate; + const feeRate = Math.round((latest + averageFeeRate) / 2); + + const estimate: FeeEstimate = { + feeRate, + minFeeRate, + maxFeeRate, + averageFeeRate, + sampleCount: samples.length, + }; + + this.emit("estimate", estimate); + return estimate; + } catch (error) { + this.emit("error", error instanceof Error ? error : new Error(String(error))); + return null; + } + } + + private emit(event: FeeEstimationEvent, payload: FeeEstimate | Error): void { + this.handlers.forEach((handler) => { + try { + handler(event, payload); + } catch (error) { + console.error(`Error in fee estimation handler for ${event}:`, error); + } + }); + } +} + +/** + * Creates a new FeeEstimator instance. + * + * @returns A new FeeEstimator instance + */ +export function createFeeEstimator(): FeeEstimator { + return new FeeEstimator(); +} From 89added6d888e46c4857f64419b42711da5e4a4e Mon Sep 17 00:00:00 2001 From: James21 Date: Sun, 27 Sep 2026 05:29:39 +0000 Subject: [PATCH 3/4] fix: #920 Implement payment pathway optimization Closes #920 --- src/broadcaster.ts | 132 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 132 insertions(+) diff --git a/src/broadcaster.ts b/src/broadcaster.ts index 57f6f5e..4cca016 100644 --- a/src/broadcaster.ts +++ b/src/broadcaster.ts @@ -41,6 +41,53 @@ export type FeeEstimationEvent = "estimate" | "error"; */ type FeeEstimationHandler = (event: FeeEstimationEvent, payload: FeeEstimate | Error) => void; +/** + * A candidate payment pathway with its associated cost and reliability. + */ +export interface PaymentPathway { + /** Identifier of the pathway (e.g. channel or route id). */ + id: string; + /** Estimated fee rate for routing through this pathway, in sat/vB. */ + feeRate: number; + /** Estimated probability (0..1) that the payment succeeds via this pathway. */ + successProbability: number; + /** Optional available liquidity along the pathway, in the smallest unit. */ + liquidity?: number; +} + +/** + * A scored payment pathway produced by the optimizer. + */ +export interface ScoredPaymentPathway extends PaymentPathway { + /** Composite score; higher is better. */ + score: number; +} + +/** + * Result of a payment pathway optimization run. + */ +export interface PaymentPathwayOptimization { + /** Pathways ordered from best to worst by score. */ + pathways: ScoredPaymentPathway[]; + /** The recommended pathway, or null when no candidates were provided. */ + recommended: ScoredPaymentPathway | null; + /** Number of candidate pathways that were evaluated. */ + evaluatedCount: number; +} + +/** + * Event names emitted during the payment pathway optimization lifecycle. + */ +export type PaymentPathwayEvent = "optimized" | "error"; + +/** + * Handler invoked when a payment pathway optimization event is emitted. + */ +type PaymentPathwayHandler = ( + event: PaymentPathwayEvent, + payload: PaymentPathwayOptimization | Error, +) => void; + /** * Invoice state broadcaster that publishes state changes to multiple subscribers. */ @@ -218,3 +265,88 @@ export class FeeEstimator { export function createFeeEstimator(): FeeEstimator { return new FeeEstimator(); } + +/** + * Optimizes payment pathways by scoring candidates on cost and reliability. + * + * Each candidate is scored so that cheaper fees and higher success + * probabilities rank higher. The optimizer emits lifecycle events so callers + * can react to optimization results and errors. + */ +export class PaymentPathwayOptimizer { + private handlers: Set = new Set(); + + /** + * Subscribe to payment pathway optimization lifecycle events. + * + * @param handler - Handler invoked on "optimized" and "error" events + * @returns Unsubscribe function that removes only this handler + */ + on(handler: PaymentPathwayHandler): () => void { + this.handlers.add(handler); + return () => { + this.handlers.delete(handler); + }; + } + + /** + * Score a single payment pathway. + * + * The score rewards higher success probability and penalizes higher fees. + * A zero or negative fee rate is treated as the cheapest possible pathway. + * + * @param pathway - The candidate pathway to score + * @returns The pathway annotated with its composite score + */ + scorePathway(pathway: PaymentPathway): ScoredPaymentPathway { + const probability = Math.min(Math.max(pathway.successProbability, 0), 1); + const feeRate = pathway.feeRate > 0 ? pathway.feeRate : 1; + const score = probability / feeRate; + return { ...pathway, score }; + } + + /** + * Optimize a set of candidate payment pathways. + * + * @param pathways - The candidate pathways to evaluate + * @returns The optimization result, ordered from best to worst + */ + optimize(pathways: PaymentPathway[]): PaymentPathwayOptimization { + try { + const scored = pathways + .map((pathway) => this.scorePathway(pathway)) + .sort((a, b) => b.score - a.score); + + const result: PaymentPathwayOptimization = { + pathways: scored, + recommended: scored.length > 0 ? scored[0] : null, + evaluatedCount: scored.length, + }; + + this.emit("optimized", result); + return result; + } catch (error) { + this.emit("error", error instanceof Error ? error : new Error(String(error))); + return { pathways: [], recommended: null, evaluatedCount: 0 }; + } + } + + private emit(event: PaymentPathwayEvent, payload: PaymentPathwayOptimization | Error): void { + this.handlers.forEach((handler) => { + try { + handler(event, payload); + } catch (error) { + console.error(`Error in payment pathway handler for ${event}:`, error); + } + }); + } +} + +/** + * Creates a new PaymentPathwayOptimizer instance. + * + * @returns A new PaymentPathwayOptimizer instance + */ +export function createPaymentPathwayOptimizer(): PaymentPathwayOptimizer { + return new PaymentPathwayOptimizer(); +} From 1445c4af51af30092cfa6716aa053fe761086b44 Mon Sep 17 00:00:00 2001 From: James21 Date: Sun, 27 Sep 2026 05:29:53 +0000 Subject: [PATCH 4/4] fix: #921 Add SDK audit trail logging for compliance Closes #921 --- src/audit/AuditTrailHasher.ts | 31 +++++++++++++++++ src/auditLogger.ts | 64 +++++++++++++++++++++++++++++++++++ 2 files changed, 95 insertions(+) diff --git a/src/audit/AuditTrailHasher.ts b/src/audit/AuditTrailHasher.ts index 3ed5b75..978e7d5 100644 --- a/src/audit/AuditTrailHasher.ts +++ b/src/audit/AuditTrailHasher.ts @@ -4,8 +4,11 @@ import * as crypto from 'crypto'; // Use node crypto webcrypto subtle const subtle = crypto.webcrypto.subtle; +export type AuditTrailListener = (entry: AuditChainEntry) => void; + export class AuditTrailHasher { private entries: AuditChainEntry[] = []; + private listeners: Set = new Set(); constructor(entries: AuditChainEntry[] = []) { this.entries = [...entries]; @@ -22,6 +25,33 @@ export class AuditTrailHasher { return hashArray.map(b => b.toString(16).padStart(2, '0')).join(''); } + /** + * Subscribes to audit trail events. Returns an unsubscribe function. + */ + onAppend(listener: AuditTrailListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + + /** + * Removes a previously registered audit trail listener. + */ + offAppend(listener: AuditTrailListener): void { + this.listeners.delete(listener); + } + + private emitAppend(entry: AuditChainEntry): void { + for (const listener of this.listeners) { + try { + listener(entry); + } catch { + // Listener errors must not break the audit chain. + } + } + } + /** * Appends a new event to the audit trail */ @@ -35,6 +65,7 @@ export class AuditTrailHasher { const entry: AuditChainEntry = { event, hash, prevHash, index }; this.entries.push(entry); + this.emitAppend(entry); return entry; } diff --git a/src/auditLogger.ts b/src/auditLogger.ts index 063713a..1ab5c48 100644 --- a/src/auditLogger.ts +++ b/src/auditLogger.ts @@ -12,6 +12,28 @@ export interface AuditEntry { decodedXdr?: DecodedXDR; } +/** + * A single audit event emitted by the SDK for compliance tracking. + * + * Unlike {@link AuditEntry}, which is a low-level sink record, an + * `AuditEvent` carries a stable `type` discriminator and a monotonically + * increasing `sequence` so downstream consumers can order and reconcile + * events reliably. + */ +export interface AuditEvent { + /** Stable event type discriminator, e.g. `"audit.log"`. */ + type: string; + /** Monotonically increasing sequence number, starting at 1. */ + sequence: number; + /** Wall-clock time the event was emitted (ms since epoch). */ + timestamp: number; + /** The audit entry associated with this event. */ + entry: AuditEntry; +} + +/** Handler invoked for every emitted {@link AuditEvent}. */ +export type AuditEventListener = (event: AuditEvent) => void; + const STELLAR_ADDRESS_RE = /^G[A-Z0-9]{55}$/; /** Detect if a string value looks like base64-encoded XDR. */ @@ -23,13 +45,55 @@ const MIN_XDR_LENGTH = 40; export class AuditLogger { private readonly sink: (entry: AuditEntry) => void; private readonly splitAuditTrails = new Map(); + private readonly listeners = new Set(); + private sequence = 0; constructor(sink: (entry: AuditEntry) => void) { this.sink = sink; } + /** + * Subscribe to audit events. Returns an unsubscribe function. + * + * Listeners are invoked synchronously after the entry has been written to + * the configured sink, so a throwing listener can never prevent the audit + * record from being persisted. + */ + on(listener: AuditEventListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + + /** Remove a previously registered listener. */ + off(listener: AuditEventListener): void { + this.listeners.delete(listener); + } + + /** Emit an audit event to all registered listeners. */ + private emit(entry: AuditEntry): void { + if (this.listeners.size === 0) { + return; + } + const event: AuditEvent = { + type: "audit.log", + sequence: ++this.sequence, + timestamp: entry.timestamp, + entry, + }; + for (const listener of this.listeners) { + try { + listener(event); + } catch { + // A misbehaving listener must never break audit logging. + } + } + } + log(entry: AuditEntry): void { this.sink(entry); + this.emit(entry); } sanitize(params: Record): Record {