From 52bca15b09e3d91e9b87942b1879b79423dd5058 Mon Sep 17 00:00:00 2001 From: kodaodife-dev Date: Sun, 27 Sep 2026 07:49:37 +0000 Subject: [PATCH 1/4] fix: #926 Implement payment batch scheduling Closes #926 --- src/paymentBatchScheduler.ts | 199 +++++++++++++++++++++++++++++++++++ 1 file changed, 199 insertions(+) create mode 100644 src/paymentBatchScheduler.ts diff --git a/src/paymentBatchScheduler.ts b/src/paymentBatchScheduler.ts new file mode 100644 index 0000000..78d0339 --- /dev/null +++ b/src/paymentBatchScheduler.ts @@ -0,0 +1,199 @@ +import { EventEmitter } from 'events'; + +export interface Payment { + id: string; + amount: number; + currency: string; + recipient: string; + metadata?: Record; +} + +export type BatchStatus = + | 'scheduled' + | 'executing' + | 'executed' + | 'cancelled' + | 'failed'; + +export interface PaymentBatch { + id: string; + payments: Payment[]; + scheduledAt: number; + status: BatchStatus; + createdAt: number; + executedAt?: number; + error?: string; +} + +export interface BatchSchedulerOptions { + /** Maximum number of payments allowed in a single batch. */ + maxBatchSize?: number; + /** Injectable clock for deterministic scheduling and tests. */ + now?: () => number; + /** Injectable timer functions for deterministic tests. */ + setTimeoutFn?: (fn: () => void, ms: number) => unknown; + clearTimeoutFn?: (handle: unknown) => void; +} + +export type BatchExecutor = (payments: Payment[]) => Promise | void; + +export interface BatchSchedulerEvents { + scheduled: (batch: PaymentBatch) => void; + executed: (batch: PaymentBatch) => void; + cancelled: (batch: PaymentBatch) => void; + failed: (batch: PaymentBatch, error: Error) => void; +} + +const DEFAULT_MAX_BATCH_SIZE = 100; + +/** + * Schedules and executes batches of payments. + * + * Batches are queued by their scheduled time and executed in order once their + * scheduled time is reached. Lifecycle transitions emit typed events so callers + * can observe scheduling, execution, cancellation, and failure. + */ +export class PaymentBatchScheduler extends EventEmitter { + private readonly batches = new Map(); + private readonly timers = new Map(); + private readonly maxBatchSize: number; + private readonly now: () => number; + private readonly setTimeoutFn: (fn: () => void, ms: number) => unknown; + private readonly clearTimeoutFn: (handle: unknown) => void; + private sequence = 0; + + constructor( + private readonly executor: BatchExecutor, + options: BatchSchedulerOptions = {}, + ) { + super(); + this.maxBatchSize = options.maxBatchSize ?? DEFAULT_MAX_BATCH_SIZE; + this.now = options.now ?? (() => Date.now()); + this.setTimeoutFn = + options.setTimeoutFn ?? + ((fn, ms) => setTimeout(fn, ms) as unknown); + this.clearTimeoutFn = + options.clearTimeoutFn ?? ((handle) => clearTimeout(handle as never)); + } + + /** + * Schedule a batch of payments for execution at the given time. + * Returns the created batch. Throws if the batch is invalid. + */ + schedule( + payments: Payment[], + scheduledAt: number, + id?: string, + ): PaymentBatch { + if (!Array.isArray(payments) || payments.length === 0) { + throw new Error('Cannot schedule an empty payment batch'); + } + if (payments.length > this.maxBatchSize) { + throw new Error( + `Batch size ${payments.length} exceeds maximum of ${this.maxBatchSize}`, + ); + } + if (!Number.isFinite(scheduledAt)) { + throw new Error('scheduledAt must be a finite timestamp'); + } + + const batchId = id ?? this.nextId(); + if (this.batches.has(batchId)) { + throw new Error(`Batch ${batchId} already exists`); + } + + const batch: PaymentBatch = { + id: batchId, + payments: payments.map((p) => ({ ...p })), + scheduledAt, + status: 'scheduled', + createdAt: this.now(), + }; + + this.batches.set(batchId, batch); + this.armTimer(batch); + this.emit('scheduled', batch); + return batch; + } + + /** Cancel a scheduled batch. Returns true if it was cancelled. */ + cancel(id: string): boolean { + const batch = this.batches.get(id); + if (!batch || batch.status !== 'scheduled') { + return false; + } + this.disarmTimer(id); + batch.status = 'cancelled'; + this.emit('cancelled', batch); + return true; + } + + /** Retrieve a batch by id. */ + getBatch(id: string): PaymentBatch | undefined { + return this.batches.get(id); + } + + /** List all known batches, ordered by scheduled time then creation order. */ + listBatches(): PaymentBatch[] { + return Array.from(this.batches.values()).sort((a, b) => { + if (a.scheduledAt !== b.scheduledAt) { + return a.scheduledAt - b.scheduledAt; + } + return a.createdAt - b.createdAt; + }); + } + + /** Cancel all pending timers. Useful for teardown. */ + dispose(): void { + for (const id of Array.from(this.timers.keys())) { + this.disarmTimer(id); + } + } + + private nextId(): string { + this.sequence += 1; + return `batch_${this.now()}_${this.sequence}`; + } + + private armTimer(batch: PaymentBatch): void { + const delay = Math.max(0, batch.scheduledAt - this.now()); + const handle = this.setTimeoutFn(() => { + this.timers.delete(batch.id); + void this.execute(batch.id); + }, delay); + this.timers.set(batch.id, handle); + } + + private disarmTimer(id: string): void { + const handle = this.timers.get(id); + if (handle !== undefined) { + this.clearTimeoutFn(handle); + this.timers.delete(id); + } + } + + /** Execute a batch immediately, regardless of its scheduled time. */ + async execute(id: string): Promise { + const batch = this.batches.get(id); + if (!batch || batch.status !== 'scheduled') { + return batch; + } + + this.disarmTimer(id); + batch.status = 'executing'; + + try { + await this.executor(batch.payments); + batch.status = 'executed'; + batch.executedAt = this.now(); + this.emit('executed', batch); + } catch (err) { + const error = err instanceof Error ? err : new Error(String(err)); + batch.status = 'failed'; + batch.error = error.message; + this.emit('failed', batch, error); + } + + return batch; + } +} From 9397d24df26e5443b0b36bfc7e324d2e1b657505 Mon Sep 17 00:00:00 2001 From: kodaodife-dev Date: Sun, 27 Sep 2026 07:49:55 +0000 Subject: [PATCH 2/4] fix: #927 Add SDK transaction rollback simulation Closes #927 --- src/autoRecovery.ts | 118 ++++++++++++++++++++++++++++++++++++ src/autoResolveSimulator.ts | 114 ++++++++++++++++++++++++++++++++++ src/broadcaster.ts | 83 +++++++++++++++++++++++++ 3 files changed, 315 insertions(+) diff --git a/src/autoRecovery.ts b/src/autoRecovery.ts index 9b61d44..80e5bb4 100644 --- a/src/autoRecovery.ts +++ b/src/autoRecovery.ts @@ -74,3 +74,121 @@ export class AutoRecoveryMonitor { } } } + +/** + * Lifecycle phase of a simulated transaction rollback. + */ +export type RollbackPhase = "start" | "success" | "failure"; + +export interface RollbackEvent { + phase: RollbackPhase; + transactionId: string; + reason?: string; + error?: Error; + timestamp: number; +} + +export interface RollbackSimulationOptions { + /** Maximum number of simulated attempts before giving up. Defaults to 3. */ + maxAttempts?: number; + /** Delay in ms between simulated attempts. Defaults to 0. */ + retryDelayMs?: number; + /** Optional predicate deciding whether a given attempt should fail. */ + shouldFail?: (attempt: number, transactionId: string) => boolean; + /** Optional hook invoked for every rollback lifecycle event. */ + onEvent?: (event: RollbackEvent) => void; +} + +export interface RollbackSimulationResult { + transactionId: string; + success: boolean; + attempts: number; + rolledBack: boolean; + error?: Error; +} + +/** + * Simulates the rollback of an SDK transaction, emitting lifecycle events + * (start/success/failure) and reporting whether the transaction was rolled + * back after exhausting its retry budget. + */ +export class TransactionRollbackSimulator { + private readonly maxAttempts: number; + private readonly retryDelayMs: number; + private readonly shouldFail: (attempt: number, transactionId: string) => boolean; + private readonly onEvent: (event: RollbackEvent) => void; + + constructor(options: RollbackSimulationOptions = {}) { + this.maxAttempts = options.maxAttempts ?? 3; + this.retryDelayMs = options.retryDelayMs ?? 0; + this.shouldFail = options.shouldFail ?? (() => false); + this.onEvent = options.onEvent ?? (() => {}); + } + + private emit(event: RollbackEvent): void { + try { + this.onEvent(event); + } catch { + // Listener errors must not break the simulation. + } + } + + private async delay(): Promise { + if (this.retryDelayMs <= 0) { + return; + } + await new Promise((resolve) => setTimeout(resolve, this.retryDelayMs)); + } + + /** + * Run a rollback simulation for the given transaction id. + * + * @param transactionId - Identifier of the transaction being simulated + * @param execute - Optional executor invoked per attempt; a thrown error + * marks the attempt as failed. + */ + async simulate( + transactionId: string, + execute?: (attempt: number) => Promise | void + ): Promise { + this.emit({ phase: "start", transactionId, timestamp: Date.now() }); + + let attempts = 0; + let lastError: Error | undefined; + + while (attempts < this.maxAttempts) { + attempts += 1; + try { + if (this.shouldFail(attempts, transactionId)) { + throw new Error(`Simulated failure on attempt ${attempts}`); + } + if (execute) { + await execute(attempts); + } + this.emit({ phase: "success", transactionId, timestamp: Date.now() }); + return { transactionId, success: true, attempts, rolledBack: false }; + } catch (err) { + lastError = err instanceof Error ? err : new Error(String(err)); + if (attempts < this.maxAttempts) { + await this.delay(); + } + } + } + + this.emit({ + phase: "failure", + transactionId, + reason: lastError?.message, + error: lastError, + timestamp: Date.now(), + }); + + return { + transactionId, + success: false, + attempts, + rolledBack: true, + error: lastError, + }; + } +} diff --git a/src/autoResolveSimulator.ts b/src/autoResolveSimulator.ts index 7e9c79e..a30bd53 100644 --- a/src/autoResolveSimulator.ts +++ b/src/autoResolveSimulator.ts @@ -38,3 +38,117 @@ export function simulateAutoResolve(invoice: Invoice): AutoResolveSimulation { return { wouldResolve: false, action: null, matchedRule: null }; } + +/** + * Lifecycle phase of a simulated transaction rollback. + */ +export type RollbackPhase = "start" | "success" | "failure"; + +/** + * Event emitted as a simulated transaction rollback progresses through its + * lifecycle. `start` is emitted before the simulated transaction is applied, + * followed by exactly one terminal event (`success` or `failure`). + */ +export interface RollbackEvent { + phase: RollbackPhase; + /** Human-readable description of the phase. */ + message: string; + /** Error that caused a `failure` event, when applicable. */ + error?: Error; +} + +/** + * A single step in a simulated transaction. Each step is applied in order and + * may throw to signal that the transaction should be rolled back. + */ +export interface RollbackStep { + /** Label used in emitted events and error messages. */ + name: string; + /** Pure function that produces the next state from the current state. */ + apply: (state: TState) => TState; +} + +/** + * Result of simulating a transaction rollback. + */ +export interface RollbackSimulationResult { + /** Whether every step applied without throwing. */ + committed: boolean; + /** State after the simulation: the committed state or the original state. */ + state: TState; + /** Name of the step that failed, or `null` when the transaction committed. */ + failedStep: string | null; + /** Error thrown by the failing step, or `null` when the transaction committed. */ + error: Error | null; + /** Ordered lifecycle events emitted during the simulation. */ + events: RollbackEvent[]; +} + +/** + * Simulate a transaction against an initial state, rolling back to that state + * if any step throws. + * + * Steps are applied in order to a working copy of the state. If a step throws, + * the working copy is discarded and the original state is returned, mirroring + * the all-or-nothing semantics of an on-chain transaction. Lifecycle events are + * emitted for the start of the simulation and for its terminal outcome. + * + * Pure function — performs no RPC calls and never mutates `initialState`. + * + * @param initialState - State the transaction starts from and rolls back to. + * @param steps - Ordered steps to apply. + * @param onEvent - Optional listener invoked for each lifecycle event. + * @returns The simulation result, including the emitted events. + */ +export function simulateTransactionRollback( + initialState: TState, + steps: ReadonlyArray>, + onEvent?: (event: RollbackEvent) => void, +): RollbackSimulationResult { + const events: RollbackEvent[] = []; + + const emit = (event: RollbackEvent): void => { + events.push(event); + onEvent?.(event); + }; + + emit({ + phase: "start", + message: `Simulating transaction with ${steps.length} step(s)`, + }); + + let workingState = initialState; + + for (const step of steps) { + try { + workingState = step.apply(workingState); + } catch (cause) { + const error = cause instanceof Error ? cause : new Error(String(cause)); + emit({ + phase: "failure", + message: `Step "${step.name}" failed; rolling back`, + error, + }); + return { + committed: false, + state: initialState, + failedStep: step.name, + error, + events, + }; + } + } + + emit({ + phase: "success", + message: `Transaction committed after ${steps.length} step(s)`, + }); + + return { + committed: true, + state: workingState, + failedStep: null, + error: null, + events, + }; +} diff --git a/src/broadcaster.ts b/src/broadcaster.ts index c4791ec..f3889f6 100644 --- a/src/broadcaster.ts +++ b/src/broadcaster.ts @@ -76,3 +76,86 @@ export class InvoiceStateBroadcaster { export function createInvoiceStateBroadcaster(): InvoiceStateBroadcaster { return new InvoiceStateBroadcaster(); } + +/** + * Lifecycle phase of a simulated transaction rollback. + */ +export type RollbackPhase = "start" | "success" | "failure"; + +/** + * Event emitted during a transaction rollback simulation. + */ +export interface RollbackEvent { + /** The transaction identifier being rolled back. */ + transactionId: string; + /** The lifecycle phase this event represents. */ + phase: RollbackPhase; + /** Optional error when the rollback fails. */ + error?: Error; +} + +/** + * Handler invoked for each rollback lifecycle event. + */ +export type RollbackEventHandler = (event: RollbackEvent) => void; + +/** + * Simulates SDK transaction rollbacks, emitting lifecycle events for the + * start, success, and failure phases of each rollback. + */ +export class TransactionRollbackSimulator { + private handlers: Set = new Set(); + + /** + * Register a handler for rollback lifecycle events. + * + * @param handler - The handler to invoke on each event + * @returns Unsubscribe function that removes only this handler + */ + onRollback(handler: RollbackEventHandler): () => void { + this.handlers.add(handler); + return () => { + this.handlers.delete(handler); + }; + } + + /** + * Simulate rolling back a transaction. Emits a "start" event, then either a + * "success" event or a "failure" event depending on the outcome. + * + * @param transactionId - The transaction identifier to roll back + * @param shouldFail - When true, the rollback fails and emits a failure event + * @returns True when the rollback succeeded, false otherwise + */ + simulateRollback(transactionId: string, shouldFail = false): boolean { + this.emit({ transactionId, phase: "start" }); + + if (shouldFail) { + const error = new Error(`Rollback failed for transaction ${transactionId}`); + this.emit({ transactionId, phase: "failure", error }); + return false; + } + + this.emit({ transactionId, phase: "success" }); + return true; + } + + private emit(event: RollbackEvent): void { + this.handlers.forEach((handler) => { + try { + handler(event); + } catch (error) { + console.error(`Error in rollback handler for ${event.transactionId}:`, error); + } + }); + } +} + +/** + * Creates a new TransactionRollbackSimulator instance. + * + * @returns A new TransactionRollbackSimulator instance + */ +export function createTransactionRollbackSimulator(): TransactionRollbackSimulator { + return new TransactionRollbackSimulator(); +} From d2038cc934d41dd813277d7aef76426576760402 Mon Sep 17 00:00:00 2001 From: kodaodife-dev Date: Sun, 27 Sep 2026 07:50:07 +0000 Subject: [PATCH 3/4] fix: #928 Implement invoice forensics tools Closes #928 --- src/invoiceForensics.ts | 300 ++++++++++++++++++++++++++++++++++++++++ 1 file changed, 300 insertions(+) create mode 100644 src/invoiceForensics.ts diff --git a/src/invoiceForensics.ts b/src/invoiceForensics.ts new file mode 100644 index 0000000..3705b82 --- /dev/null +++ b/src/invoiceForensics.ts @@ -0,0 +1,300 @@ +/** + * Invoice forensics tools. + * + * Provides analysis of invoice fields to detect tampering, mismatches and + * anomalies, plus typed event handling for forensic findings. + */ + +export type InvoiceField = 'invoiceNumber' | 'amount' | 'currency' | 'issuedAt' | 'dueAt' | 'vendor' | 'buyer' | 'lineItems'; + +export type FindingSeverity = 'info' | 'warning' | 'critical'; + +export interface InvoiceLineItem { + description: string; + quantity: number; + unitPrice: number; +} + +export interface Invoice { + invoiceNumber: string; + amount: number; + currency: string; + issuedAt: string; + dueAt: string; + vendor: string; + buyer: string; + lineItems: InvoiceLineItem[]; +} + +export interface ForensicFinding { + code: string; + severity: FindingSeverity; + field: InvoiceField; + message: string; + expected?: unknown; + actual?: unknown; +} + +export interface ForensicReport { + invoiceNumber: string; + findings: ForensicFinding[]; + riskScore: number; + passed: boolean; +} + +export type ForensicEventType = 'finding' | 'report' | 'error'; + +export interface ForensicEventMap { + finding: ForensicFinding; + report: ForensicReport; + error: { message: string; error?: unknown }; +} + +export type ForensicEventHandler = (payload: ForensicEventMap[T]) => void; + +export interface InvoiceForensicsOptions { + /** Absolute tolerance when comparing monetary amounts. */ + amountTolerance?: number; + /** Maximum allowed gap (ms) between issue and due dates. */ + maxTermMs?: number; + /** Risk score at or above which the report is considered failed. */ + riskThreshold?: number; +} + +const DEFAULT_AMOUNT_TOLERANCE = 0.01; +const DEFAULT_MAX_TERM_MS = 365 * 24 * 60 * 60 * 1000; +const DEFAULT_RISK_THRESHOLD = 50; + +const SEVERITY_WEIGHT: Record = { + info: 1, + warning: 10, + critical: 40, +}; + +/** + * Analyzes invoices for tampering, mismatches and anomalies and emits typed + * events for each finding and for the final report. + */ +export class InvoiceForensics { + private readonly amountTolerance: number; + private readonly maxTermMs: number; + private readonly riskThreshold: number; + private readonly handlers: { [K in ForensicEventType]: Set> } = { + finding: new Set(), + report: new Set(), + error: new Set(), + }; + + constructor(options: InvoiceForensicsOptions = {}) { + this.amountTolerance = options.amountTolerance ?? DEFAULT_AMOUNT_TOLERANCE; + this.maxTermMs = options.maxTermMs ?? DEFAULT_MAX_TERM_MS; + this.riskThreshold = options.riskThreshold ?? DEFAULT_RISK_THRESHOLD; + } + + on(event: T, handler: ForensicEventHandler): () => void { + this.handlers[event].add(handler as ForensicEventHandler); + return () => this.off(event, handler); + } + + off(event: T, handler: ForensicEventHandler): void { + this.handlers[event].delete(handler as ForensicEventHandler); + } + + private emit(event: T, payload: ForensicEventMap[T]): void { + for (const handler of this.handlers[event]) { + try { + (handler as ForensicEventHandler)(payload); + } catch (error) { + if (event !== 'error') { + this.emit('error', { message: `Handler for "${event}" threw`, error }); + } + } + } + } + + /** Runs all forensic checks against a single invoice. */ + analyze(invoice: Invoice): ForensicReport { + const findings: ForensicFinding[] = []; + + try { + findings.push(...this.checkRequiredFields(invoice)); + findings.push(...this.checkDates(invoice)); + findings.push(...this.checkAmounts(invoice)); + findings.push(...this.checkLineItems(invoice)); + } catch (error) { + this.emit('error', { message: 'Invoice analysis failed', error }); + } + + for (const finding of findings) { + this.emit('finding', finding); + } + + const riskScore = findings.reduce((sum, f) => sum + SEVERITY_WEIGHT[f.severity], 0); + const report: ForensicReport = { + invoiceNumber: invoice.invoiceNumber, + findings, + riskScore, + passed: riskScore < this.riskThreshold, + }; + + this.emit('report', report); + return report; + } + + /** Analyzes a batch of invoices and returns their reports. */ + analyzeAll(invoices: Invoice[]): ForensicReport[] { + return invoices.map((invoice) => this.analyze(invoice)); + } + + private checkRequiredFields(invoice: Invoice): ForensicFinding[] { + const findings: ForensicFinding[] = []; + const required: InvoiceField[] = ['invoiceNumber', 'currency', 'vendor', 'buyer']; + + for (const field of required) { + const value = invoice[field]; + if (typeof value !== 'string' || value.trim() === '') { + findings.push({ + code: 'MISSING_FIELD', + severity: 'critical', + field, + message: `Required field "${field}" is missing or empty`, + actual: value, + }); + } + } + + return findings; + } + + private checkDates(invoice: Invoice): ForensicFinding[] { + const findings: ForensicFinding[] = []; + const issued = Date.parse(invoice.issuedAt); + const due = Date.parse(invoice.dueAt); + + if (Number.isNaN(issued)) { + findings.push({ + code: 'INVALID_DATE', + severity: 'critical', + field: 'issuedAt', + message: 'issuedAt is not a valid date', + actual: invoice.issuedAt, + }); + } + + if (Number.isNaN(due)) { + findings.push({ + code: 'INVALID_DATE', + severity: 'critical', + field: 'dueAt', + message: 'dueAt is not a valid date', + actual: invoice.dueAt, + }); + } + + if (!Number.isNaN(issued) && !Number.isNaN(due)) { + if (due < issued) { + findings.push({ + code: 'DUE_BEFORE_ISSUED', + severity: 'critical', + field: 'dueAt', + message: 'Due date precedes issue date', + expected: invoice.issuedAt, + actual: invoice.dueAt, + }); + } else if (due - issued > this.maxTermMs) { + findings.push({ + code: 'EXCESSIVE_TERM', + severity: 'warning', + field: 'dueAt', + message: 'Payment term exceeds the configured maximum', + expected: this.maxTermMs, + actual: due - issued, + }); + } + } + + return findings; + } + + private checkAmounts(invoice: Invoice): ForensicFinding[] { + const findings: ForensicFinding[] = []; + + if (typeof invoice.amount !== 'number' || !Number.isFinite(invoice.amount)) { + findings.push({ + code: 'INVALID_AMOUNT', + severity: 'critical', + field: 'amount', + message: 'Invoice amount is not a finite number', + actual: invoice.amount, + }); + return findings; + } + + if (invoice.amount <= 0) { + findings.push({ + code: 'NON_POSITIVE_AMOUNT', + severity: 'critical', + field: 'amount', + message: 'Invoice amount must be positive', + actual: invoice.amount, + }); + } + + const lineTotal = invoice.lineItems.reduce( + (sum, item) => sum + item.quantity * item.unitPrice, + 0, + ); + + if (Math.abs(lineTotal - invoice.amount) > this.amountTolerance) { + findings.push({ + code: 'AMOUNT_MISMATCH', + severity: 'critical', + field: 'amount', + message: 'Invoice amount does not match the sum of line items', + expected: lineTotal, + actual: invoice.amount, + }); + } + + return findings; + } + + private checkLineItems(invoice: Invoice): ForensicFinding[] { + const findings: ForensicFinding[] = []; + + if (!Array.isArray(invoice.lineItems) || invoice.lineItems.length === 0) { + findings.push({ + code: 'NO_LINE_ITEMS', + severity: 'warning', + field: 'lineItems', + message: 'Invoice has no line items', + actual: invoice.lineItems, + }); + return findings; + } + + invoice.lineItems.forEach((item, index) => { + if (typeof item.quantity !== 'number' || item.quantity <= 0) { + findings.push({ + code: 'INVALID_QUANTITY', + severity: 'warning', + field: 'lineItems', + message: `Line item ${index} has an invalid quantity`, + actual: item.quantity, + }); + } + + if (typeof item.unitPrice !== 'number' || item.unitPrice < 0) { + findings.push({ + code: 'INVALID_UNIT_PRICE', + severity: 'warning', + field: 'lineItems', + message: `Line item ${index} has an invalid unit price`, + actual: item.unitPrice, + }); + } + }); + + return findings; + } +} From 479c3b7306557dd168807c7958693da902ff3465 Mon Sep 17 00:00:00 2001 From: kodaodife-dev Date: Sun, 27 Sep 2026 07:50:38 +0000 Subject: [PATCH 4/4] fix: #929 Add SDK price oracle integration helpers Closes #929 --- src/__tests__/priceOracle.test.ts | 61 ++++++++ src/ammCalculator.ts | 224 +++++++++++++++++++++++++++++- src/priceOracle.ts | 154 ++++++++++++-------- 3 files changed, 379 insertions(+), 60 deletions(-) create mode 100644 src/__tests__/priceOracle.test.ts diff --git a/src/__tests__/priceOracle.test.ts b/src/__tests__/priceOracle.test.ts new file mode 100644 index 0000000..05c3d7f --- /dev/null +++ b/src/__tests__/priceOracle.test.ts @@ -0,0 +1,61 @@ +import { + PriceOracle, + parseOraclePrice, + type OraclePrice, +} from "../priceOracle"; + +describe("parseOraclePrice", () => { + it("parses a numeric price and timestamp", () => { + const result = parseOraclePrice("XLM", { price: 0.42, timestamp: 1700000000000 }); + expect(result).toEqual({ symbol: "XLM", price: 0.42, timestamp: 1700000000000 }); + }); + + it("parses string prices and aliases", () => { + const result = parseOraclePrice("USDC", { value: "1.01", time: 1700000000000 }); + expect(result.price).toBe(1.01); + expect(result.timestamp).toBe(1700000000000); + }); + + it("falls back to the current time when no timestamp is present", () => { + const before = Date.now(); + const result = parseOraclePrice("XLM", { amount: 2 }); + expect(result.timestamp).toBeGreaterThanOrEqual(before); + }); + + it("throws on an invalid price", () => { + expect(() => parseOraclePrice("XLM", { price: "not-a-number" })).toThrow( + /Invalid oracle price/, + ); + expect(() => parseOraclePrice("XLM", {})).toThrow(/Invalid oracle price/); + }); +}); + +describe("PriceOracle", () => { + it("fetches, parses, and caches a price", async () => { + const oracle = new PriceOracle(async () => ({ price: 3.5, timestamp: 1 })); + const price = await oracle.fetchPrice("XLM"); + expect(price).toEqual({ symbol: "XLM", price: 3.5, timestamp: 1 }); + expect(oracle.getCachedPrice("XLM")).toEqual(price); + }); + + it("emits price updates to listeners", async () => { + const oracle = new PriceOracle(async () => ({ price: 7 })); + const received: OraclePrice[] = []; + const unsubscribe = oracle.onPriceUpdate((p) => received.push(p)); + + await oracle.fetchPrice("XLM"); + expect(received).toHaveLength(1); + expect(received[0].price).toBe(7); + + unsubscribe(); + await oracle.fetchPrice("XLM"); + expect(received).toHaveLength(1); + }); + + it("fetches multiple prices in order", async () => { + const oracle = new PriceOracle(async (symbol) => ({ price: symbol.length })); + const prices = await oracle.fetchPrices(["XLM", "USDC"]); + expect(prices.map((p) => p.symbol)).toEqual(["XLM", "USDC"]); + expect(prices.map((p) => p.price)).toEqual([3, 4]); + }); +}); diff --git a/src/ammCalculator.ts b/src/ammCalculator.ts index 4e5af3a..4ae3d42 100644 --- a/src/ammCalculator.ts +++ b/src/ammCalculator.ts @@ -183,10 +183,232 @@ export function calculatePoolShare( }; } +// --------------------------------------------------------------------------- +// Price oracle integration helpers +// --------------------------------------------------------------------------- + +/** + * A single price observation returned by a price oracle source. + */ +export interface OraclePriceObservation { + /** Asset identifier the price refers to (e.g. "native" or "USDC:GA..."). */ + asset: string; + /** Price expressed as a decimal string, in the oracle's quote asset. */ + price: string; + /** Unix timestamp (seconds) at which the price was observed. */ + timestamp: number; +} + +/** + * A price oracle source that can be queried for the latest price of an asset. + * Implementations may wrap an on-chain oracle contract, an HTTP feed, or a + * cached in-memory source. + */ +export interface PriceOracleSource { + /** Fetches the latest observation for the given asset. */ + getPrice(asset: string): Promise; +} + +/** + * A parsed, validated oracle price ready for use in AMM calculations. + */ +export interface ParsedOraclePrice { + asset: string; + /** Price as a decimal string, normalized to a fixed precision. */ + price: string; + /** Price scaled to an integer string (price * 10^decimals). */ + scaledPrice: string; + /** Number of decimal places used for scaledPrice. */ + decimals: number; + timestamp: number; + /** Age of the observation in seconds relative to the provided `now`. */ + ageSeconds: number; + /** Whether the observation is older than the staleness threshold. */ + stale: boolean; +} + +/** + * Options controlling how oracle prices are parsed and validated. + */ +export interface OraclePriceOptions { + /** Decimal places used when scaling the price. Defaults to 7 (Stellar). */ + decimals?: number; + /** Maximum acceptable age in seconds before a price is flagged stale. */ + maxAgeSeconds?: number; + /** Reference time (unix seconds) used to compute age. Defaults to now. */ + now?: number; +} + +const DEFAULT_ORACLE_DECIMALS = 7; +const DEFAULT_ORACLE_MAX_AGE_SECONDS = 300; + +/** + * Fetches the latest price for `asset` from the given oracle source and parses + * it into a validated {@link ParsedOraclePrice}. + * + * @param source - The oracle source to query. + * @param asset - The asset identifier to fetch a price for. + * @param options - Parsing/validation options. + * @throws Error when the source returns a malformed or non-positive price. + */ +export async function fetchOraclePrice( + source: PriceOracleSource, + asset: string, + options: OraclePriceOptions = {} +): Promise { + const observation = await source.getPrice(asset); + return parseOraclePrice(observation, options); +} + +/** + * Parses and validates a raw {@link OraclePriceObservation} into a + * {@link ParsedOraclePrice}, scaling the price to an integer string and + * computing staleness relative to `options.now`. + * + * @param observation - The raw observation to parse. + * @param options - Parsing/validation options. + * @throws Error when the price is missing, malformed, or non-positive. + */ +export function parseOraclePrice( + observation: OraclePriceObservation, + options: OraclePriceOptions = {} +): ParsedOraclePrice { + const decimals = options.decimals ?? DEFAULT_ORACLE_DECIMALS; + const maxAgeSeconds = + options.maxAgeSeconds ?? DEFAULT_ORACLE_MAX_AGE_SECONDS; + const now = options.now ?? Math.floor(Date.now() / 1000); + + if (!observation || typeof observation.price !== "string") { + throw new Error("Oracle price observation is missing a price"); + } + + const price = observation.price.trim(); + if (!/^\d+(\.\d+)?$/.test(price)) { + throw new Error(`Malformed oracle price: ${observation.price}`); + } + + const scaledPrice = scaleDecimal(price, decimals); + if (BigInt(scaledPrice) <= 0n) { + throw new Error(`Oracle price must be positive: ${observation.price}`); + } + + const ageSeconds = Math.max(0, now - observation.timestamp); + + return { + asset: observation.asset, + price, + scaledPrice, + decimals, + timestamp: observation.timestamp, + ageSeconds, + stale: ageSeconds > maxAgeSeconds, + }; +} + +/** + * Computes the cross price between two assets using their oracle prices. + * + * Given `base` priced in quote units and `counter` priced in the same quote + * units, returns how many `counter` units one `base` unit is worth, as a + * decimal string with `decimals` places of precision. + * + * @param base - Parsed price for the base asset. + * @param counter - Parsed price for the counter asset. + * @param decimals - Output precision. Defaults to the base price decimals. + * @throws Error when the counter price is zero. + */ +export function computeCrossPrice( + base: ParsedOraclePrice, + counter: ParsedOraclePrice, + decimals: number = base.decimals +): string { + const counterScaled = BigInt(counter.scaledPrice); + if (counterScaled === 0n) { + throw new Error("Cannot compute cross price with a zero counter price"); + } + + const baseScaled = BigInt(base.scaledPrice); + const SCALE = 10n ** BigInt(decimals); + const result = (baseScaled * SCALE) / counterScaled; + + return formatScaled(result, decimals); +} + +/** + * Subscribes to oracle price updates for a set of assets, invoking `onUpdate` + * whenever a new observation is produced. Returns an unsubscribe function. + * + * The source is polled at `intervalMs`; each poll fetches every asset and + * emits only observations whose price or timestamp changed since the previous + * poll. Errors from the source are forwarded to `onError` (if provided) and + * do not stop the subscription. + * + * @param source - The oracle source to poll. + * @param assets - Asset identifiers to watch. + * @param onUpdate - Callback invoked with each changed parsed price. + * @param options - Polling interval, parse options, and error handler. + * @returns A function that stops the subscription when called. + */ +export function subscribeToOraclePrices( + source: PriceOracleSource, + assets: string[], + onUpdate: (price: ParsedOraclePrice) => void, + options: OraclePriceOptions & { + intervalMs?: number; + onError?: (error: unknown) => void; + } = {} +): () => void { + const intervalMs = options.intervalMs ?? 15_000; + const lastSeen = new Map(); + let stopped = false; + + const poll = async (): Promise => { + for (const asset of assets) { + if (stopped) return; + try { + const parsed = await fetchOraclePrice(source, asset, options); + const fingerprint = `${parsed.price}@${parsed.timestamp}`; + if (lastSeen.get(asset) !== fingerprint) { + lastSeen.set(asset, fingerprint); + onUpdate(parsed); + } + } catch (error) { + options.onError?.(error); + } + } + }; + + void poll(); + const timer = setInterval(() => { + void poll(); + }, intervalMs); + + return () => { + stopped = true; + clearInterval(timer); + }; +} + // --------------------------------------------------------------------------- // Internal helpers // --------------------------------------------------------------------------- +function scaleDecimal(value: string, decimals: number): string { + const dot = value.indexOf("."); + const intPart = dot === -1 ? value : value.slice(0, dot); + const fracPart = dot === -1 ? "" : value.slice(dot + 1); + const paddedFrac = fracPart.padEnd(decimals, "0").slice(0, decimals); + return `${intPart}${paddedFrac}`.replace(/^0+(?=\d)/, "") || "0"; +} + +function formatScaled(value: bigint, decimals: number): string { + const SCALE = 10n ** BigInt(decimals); + const intPart = value / SCALE; + const fracPart = value % SCALE; + if (decimals === 0) return intPart.toString(); + return `${intPart}.${fracPart.toString().padStart(decimals, "0")}`; +} + function computeSpotPrice(reserveIn: bigint, reserveOut: bigint): string { // spotPrice = reserveOut / reserveIn as a decimal string if (reserveIn === 0n) return "0"; @@ -246,8 +468,6 @@ function computeOwnershipPercent(owned: bigint, total: bigint): string { function formatRatio(numerator: bigint, denominator: bigint): string { if (denominator === 0n) return "0"; - // Use BigInt-safe decimal division with up to 12 decimal places. - // Multiply numerator by 10^12 before division, then insert decimal point. const SCALE = 10n ** 12n; const scaled = (numerator * SCALE) / denominator; const intPart = scaled / SCALE; diff --git a/src/priceOracle.ts b/src/priceOracle.ts index 1a73de9..6db84e4 100644 --- a/src/priceOracle.ts +++ b/src/priceOracle.ts @@ -1,72 +1,110 @@ /** - * Cross-asset price oracle for settlement calculations (invoice normalisation, - * multi-asset line items). Unlike `currencyConverter.ts` (display-only, never - * feeds into real amounts), rates returned here are used to compute amounts - * that settle on-chain, so callers should treat failures as fatal rather than - * falling back to a stale/cached display value. + * SDK price oracle integration helpers. + * + * Provides a small, dependency-free abstraction for fetching and parsing + * oracle prices, plus an event emitter for price updates. */ -import { - Contract, - rpc as SorobanRpc, - TransactionBuilder, - BASE_FEE, - nativeToScVal, - scValToNative, -} from "@stellar/stellar-sdk"; -import { OraclePriceError, NoReturnValueError } from "./errors.js"; +export interface OraclePrice { + /** Asset symbol, e.g. "XLM" or "USDC". */ + symbol: string; + /** Price expressed in the oracle's quote currency. */ + price: number; + /** Unix timestamp (ms) when the price was observed. */ + timestamp: number; +} -/** Resolves a conversion rate between two on-chain assets. */ -export interface PriceOracle { - /** - * Fixed-point rate (1e18 = 1.0) to convert 1 unit of `fromAsset` into - * `toAsset`, or `undefined` when no price is available for that pair. - */ - getRate(fromAsset: string, toAsset: string): Promise; +/** + * Minimal shape of a raw oracle response. Real oracles vary, so we accept a + * permissive record and normalize it in {@link parseOraclePrice}. + */ +export type RawOracleResponse = Record; + +/** + * Fetches a raw price payload for a symbol. Implementations may hit an HTTP + * endpoint, a contract, or a cache. + */ +export type OracleFetcher = (symbol: string) => Promise; + +export type PriceUpdateListener = (price: OraclePrice) => void; + +/** + * Parse a raw oracle response into a normalized {@link OraclePrice}. + * + * Accepts common field aliases (`price`/`value`/`amount`, `timestamp`/`time`) + * and throws when a usable numeric price cannot be found. + */ +export function parseOraclePrice( + symbol: string, + raw: RawOracleResponse, +): OraclePrice { + const rawPrice = raw.price ?? raw.value ?? raw.amount; + const price = typeof rawPrice === "string" ? Number(rawPrice) : rawPrice; + + if (typeof price !== "number" || !Number.isFinite(price)) { + throw new Error(`Invalid oracle price for ${symbol}`); + } + + const rawTimestamp = raw.timestamp ?? raw.time ?? raw.updatedAt; + const timestamp = + typeof rawTimestamp === "number" + ? rawTimestamp + : typeof rawTimestamp === "string" + ? Date.parse(rawTimestamp) + : Date.now(); + + return { + symbol, + price, + timestamp: Number.isFinite(timestamp) ? timestamp : Date.now(), + }; } -/** Soroban contract-backed price oracle. */ -export class ContractPriceOracle implements PriceOracle { - constructor( - private readonly server: SorobanRpc.Server, - private readonly oracleAddress: string, - private readonly networkPassphrase: string - ) {} +/** + * Price oracle client that fetches, parses, and emits price updates. + */ +export class PriceOracle { + private readonly fetcher: OracleFetcher; + private readonly listeners = new Set(); + private readonly cache = new Map(); - async getRate(fromAsset: string, toAsset: string): Promise { - const contract = new Contract(this.oracleAddress); - const operation = contract.call( - "get_price", - nativeToScVal(fromAsset, { type: "symbol" }), - nativeToScVal(toAsset, { type: "symbol" }) - ); + constructor(fetcher: OracleFetcher) { + this.fetcher = fetcher; + } - const sourceAccount = { - accountId: () => this.oracleAddress, - sequenceNumber: () => "0", - incrementSequenceNumber: () => {}, - } as any; + /** Subscribe to price updates. Returns an unsubscribe function. */ + onPriceUpdate(listener: PriceUpdateListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } - const tx = new TransactionBuilder(sourceAccount, { - fee: BASE_FEE, - networkPassphrase: this.networkPassphrase, - }) - .addOperation(operation) - .setTimeout(30) - .build(); + /** Return the last cached price for a symbol, if any. */ + getCachedPrice(symbol: string): OraclePrice | undefined { + return this.cache.get(symbol); + } - const simResult = await this.server.simulateTransaction(tx); - if (SorobanRpc.Api.isSimulationError(simResult)) { - if (/no price|not found|unsupported/i.test(simResult.error)) { - return undefined; - } - throw new OraclePriceError(`Oracle simulation failed: ${simResult.error}`); - } + /** + * Fetch and parse the current price for a symbol, caching the result and + * notifying listeners. + */ + async fetchPrice(symbol: string): Promise { + const raw = await this.fetcher(symbol); + const price = parseOraclePrice(symbol, raw); + this.cache.set(symbol, price); + this.emit(price); + return price; + } - const returnVal = (simResult as SorobanRpc.Api.SimulateTransactionSuccessResponse).result - ?.retval; - if (!returnVal) throw new NoReturnValueError("oracle get_price"); + /** Fetch prices for multiple symbols, preserving input order. */ + async fetchPrices(symbols: string[]): Promise { + return Promise.all(symbols.map((symbol) => this.fetchPrice(symbol))); + } - return BigInt(scValToNative(returnVal)); + private emit(price: OraclePrice): void { + for (const listener of this.listeners) { + listener(price); + } } }