diff --git a/README.md b/README.md index aa6f5c4..ce2fb0d 100644 --- a/README.md +++ b/README.md @@ -402,3 +402,8 @@ This project participates in the [Drips Wave Program](https://drips.network/wave See [CONTRIBUTING.md](./CONTRIBUTING.md) for the full guide. **Do not start coding until assigned to an issue by a maintainer.** + +## Handsoff notes + + +- #912: Implement advanced filter DSL for invoice queries diff --git a/src/__tests__/budgetTracker.test.ts b/src/__tests__/budgetTracker.test.ts new file mode 100644 index 0000000..b9a5d6c --- /dev/null +++ b/src/__tests__/budgetTracker.test.ts @@ -0,0 +1,85 @@ +import { BudgetTracker } from '../budgetTracker'; + +describe('BudgetTracker', () => { + it('tracks cumulative spend and remaining budget', () => { + const tracker = new BudgetTracker({ limit: 100 }); + tracker.track(30); + tracker.track(20); + expect(tracker.getSpent()).toBe(50); + expect(tracker.getRemaining()).toBe(50); + expect(tracker.isLimitReached()).toBe(false); + }); + + it('emits a warning when spend approaches the limit', () => { + const tracker = new BudgetTracker({ limit: 100, warnThreshold: 0.8 }); + const listener = jest.fn(); + tracker.on('warning', listener); + + tracker.track(79); + expect(listener).not.toHaveBeenCalled(); + + tracker.track(1); + expect(listener).toHaveBeenCalledTimes(1); + expect(listener).toHaveBeenCalledWith({ + event: 'warning', + spent: 80, + limit: 100, + remaining: 20, + }); + }); + + it('emits the warning only once', () => { + const tracker = new BudgetTracker({ limit: 100 }); + const listener = jest.fn(); + tracker.on('warning', listener); + + tracker.track(85); + tracker.track(5); + expect(listener).toHaveBeenCalledTimes(1); + }); + + it('emits limit-reached when spend meets or exceeds the limit', () => { + const tracker = new BudgetTracker({ limit: 100 }); + const listener = jest.fn(); + tracker.on('limit-reached', listener); + + tracker.track(100); + expect(listener).toHaveBeenCalledTimes(1); + expect(listener).toHaveBeenCalledWith({ + event: 'limit-reached', + spent: 100, + limit: 100, + remaining: 0, + }); + expect(tracker.isLimitReached()).toBe(true); + }); + + it('emits limit-reached only once and clamps remaining at zero', () => { + const tracker = new BudgetTracker({ limit: 100 }); + const listener = jest.fn(); + tracker.on('limit-reached', listener); + + tracker.track(120); + tracker.track(10); + expect(listener).toHaveBeenCalledTimes(1); + expect(tracker.getRemaining()).toBe(0); + }); + + it('supports unsubscribing from events', () => { + const tracker = new BudgetTracker({ limit: 100 }); + const listener = jest.fn(); + const unsubscribe = tracker.on('warning', listener); + + unsubscribe(); + tracker.track(90); + expect(listener).not.toHaveBeenCalled(); + }); + + it('rejects invalid configuration and amounts', () => { + expect(() => new BudgetTracker({ limit: 0 })).toThrow(); + expect(() => new BudgetTracker({ limit: 100, warnThreshold: 1 })).toThrow(); + + const tracker = new BudgetTracker({ limit: 100 }); + expect(() => tracker.track(-1)).toThrow(); + }); +}); diff --git a/src/approvalWorkflowSequencer.ts b/src/approvalWorkflowSequencer.ts index bc177d4..0660684 100644 --- a/src/approvalWorkflowSequencer.ts +++ b/src/approvalWorkflowSequencer.ts @@ -13,6 +13,87 @@ export interface ApprovalWorkflowOptions { applySignatures?: SignatureApplier; } +export interface PaymentForwardingRule { + id: string; + /** Optional source account filter; matches any source when omitted. */ + source?: string; + /** Optional destination account filter; matches any destination when omitted. */ + destination?: string; + /** Optional asset code filter; matches any asset when omitted. */ + assetCode?: string; + /** Optional inclusive minimum amount filter. */ + minAmount?: number; + /** Optional inclusive maximum amount filter. */ + maxAmount?: number; + /** Account that receives the forwarded payment. */ + forwardTo: string; +} + +export interface PaymentForwardingRequest { + source: string; + destination: string; + assetCode: string; + amount: number; +} + +export interface PaymentForwardingDecision { + forwarded: boolean; + ruleId?: string; + forwardTo?: string; +} + +export class PaymentForwardingRulesEngine { + private readonly rules: PaymentForwardingRule[] = []; + + constructor(rules: readonly PaymentForwardingRule[] = []) { + for (const rule of rules) { + this.addRule(rule); + } + } + + addRule(rule: PaymentForwardingRule): void { + if (!rule.id) { + throw new Error("Payment forwarding rule requires an id"); + } + if (!rule.forwardTo) { + throw new Error(`Payment forwarding rule requires a forwardTo account: ${rule.id}`); + } + this.rules.push(rule); + emitSdkEvent("paymentForwardingRuleAdded", { ruleId: rule.id }); + } + + getRules(): readonly PaymentForwardingRule[] { + return this.rules; + } + + evaluate(request: PaymentForwardingRequest): PaymentForwardingDecision { + const rule = this.rules.find((candidate) => this.matches(candidate, request)); + if (!rule) { + emitSdkEvent("paymentForwardingSkipped", { + source: request.source, + destination: request.destination, + }); + return { forwarded: false }; + } + + emitSdkEvent("paymentForwarded", { + ruleId: rule.id, + forwardTo: rule.forwardTo, + amount: request.amount, + }); + return { forwarded: true, ruleId: rule.id, forwardTo: rule.forwardTo }; + } + + private matches(rule: PaymentForwardingRule, request: PaymentForwardingRequest): boolean { + if (rule.source !== undefined && rule.source !== request.source) return false; + if (rule.destination !== undefined && rule.destination !== request.destination) return false; + if (rule.assetCode !== undefined && rule.assetCode !== request.assetCode) return false; + if (rule.minAmount !== undefined && request.amount < rule.minAmount) return false; + if (rule.maxAmount !== undefined && request.amount > rule.maxAmount) return false; + return true; + } +} + export class ApprovalSession { private readonly signatures = new Map(); private readonly signerWeights = new Map(); diff --git a/src/budgetTracker.ts b/src/budgetTracker.ts new file mode 100644 index 0000000..68726a0 --- /dev/null +++ b/src/budgetTracker.ts @@ -0,0 +1,102 @@ +export interface BudgetTrackerOptions { + /** Maximum allowed cumulative spend. */ + limit: number; + /** Fraction of the limit at which a warning is emitted (0 < threshold < 1). Defaults to 0.8. */ + warnThreshold?: number; +} + +export type BudgetTrackerEvent = 'warning' | 'limit-reached'; + +export type BudgetTrackerListener = (payload: { + event: BudgetTrackerEvent; + spent: number; + limit: number; + remaining: number; +}) => void; + +/** + * Tracks cumulative spend against a configured budget limit and emits + * events when spend approaches or exceeds the limit. + */ +export class BudgetTracker { + private readonly limit: number; + private readonly warnThreshold: number; + private spent = 0; + private warned = false; + private reached = false; + private readonly listeners = new Map>(); + + constructor(options: BudgetTrackerOptions) { + if (!Number.isFinite(options.limit) || options.limit <= 0) { + throw new Error('BudgetTracker: limit must be a positive finite number'); + } + const warnThreshold = options.warnThreshold ?? 0.8; + if (!Number.isFinite(warnThreshold) || warnThreshold <= 0 || warnThreshold >= 1) { + throw new Error('BudgetTracker: warnThreshold must be between 0 and 1 (exclusive)'); + } + this.limit = options.limit; + this.warnThreshold = warnThreshold; + } + + /** Record additional spend and emit threshold events as needed. */ + track(amount: number): void { + if (!Number.isFinite(amount) || amount < 0) { + throw new Error('BudgetTracker: amount must be a non-negative finite number'); + } + this.spent += amount; + + if (!this.warned && this.spent >= this.limit * this.warnThreshold && this.spent < this.limit) { + this.warned = true; + this.emit('warning'); + } + + if (!this.reached && this.spent >= this.limit) { + this.reached = true; + this.emit('limit-reached'); + } + } + + /** Current cumulative spend. */ + getSpent(): number { + return this.spent; + } + + /** Remaining budget (never negative). */ + getRemaining(): number { + return Math.max(0, this.limit - this.spent); + } + + /** Whether spend has reached or exceeded the limit. */ + isLimitReached(): boolean { + return this.reached; + } + + /** Subscribe to a budget event. Returns an unsubscribe function. */ + on(event: BudgetTrackerEvent, listener: BudgetTrackerListener): () => void { + let set = this.listeners.get(event); + if (!set) { + set = new Set(); + this.listeners.set(event, set); + } + set.add(listener); + return () => { + set?.delete(listener); + }; + } + + private emit(event: BudgetTrackerEvent): void { + const set = this.listeners.get(event); + if (!set) { + return; + } + const payload = { + event, + spent: this.spent, + limit: this.limit, + remaining: this.getRemaining(), + }; + for (const listener of set) { + listener(payload); + } + } +} diff --git a/src/cache.ts b/src/cache.ts index ce6cf14..3bf155b 100644 --- a/src/cache.ts +++ b/src/cache.ts @@ -18,6 +18,16 @@ export interface MethodCacheEntry { expiresAt: number; } +export type CacheEventType = "set" | "hit" | "miss" | "expire" | "invalidate" | "evict"; + +export interface CacheEvent { + type: CacheEventType; + key: string; + timestamp: number; +} + +export type CacheEventListener = (event: CacheEvent) => void; + export class SimpleCache { private readonly store = new Map(); private readonly ttlConfig: Record; @@ -26,6 +36,7 @@ export class SimpleCache { private misses = 0; private evictions = 0; private maxEntries: number; + private readonly listeners = new Set(); constructor(config?: number | { enabled?: boolean; ttl?: Record; ttlMs?: number; maxEntries?: number }) { if (typeof config === "number") { @@ -42,16 +53,38 @@ export class SimpleCache { } } + /** + * Subscribe to cache lifecycle events (set, hit, miss, expire, invalidate, + * evict). Returns an unsubscribe function. + */ + on(listener: CacheEventListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + + private emit(type: CacheEventType, key: string): void { + if (this.listeners.size === 0) return; + const event: CacheEvent = { type, key, timestamp: Date.now() }; + for (const listener of this.listeners) { + listener(event); + } + } + get(key: string): T | undefined { if (!this.enabled) return undefined; const entry = this.store.get(key); if (!entry) { this.misses++; + this.emit("miss", key); return undefined; } if (Date.now() > entry.expiresAt) { this.store.delete(key); this.misses++; + this.emit("expire", key); + this.emit("miss", key); return undefined; } @@ -60,6 +93,7 @@ export class SimpleCache { this.store.set(key, entry); this.hits++; + this.emit("hit", key); return entry.value; } @@ -74,26 +108,35 @@ export class SimpleCache { if (oldestKey !== undefined) { this.store.delete(oldestKey); this.evictions++; + this.emit("evict", oldestKey); } } this.store.set(key, { value, expiresAt: Date.now() + ttl }); + this.emit("set", key); } invalidate(methodOrKey?: string, args?: any[]): void { if (!methodOrKey) { + const keys = Array.from(this.store.keys()); this.store.clear(); + for (const key of keys) { + this.emit("invalidate", key); + } return; } if (args) { const key = `${methodOrKey}:${JSON.stringify(args)}`; - this.store.delete(key); + if (this.store.delete(key)) { + this.emit("invalidate", key); + } return; } // Check if it's an exact key if (this.store.has(methodOrKey)) { this.store.delete(methodOrKey); + this.emit("invalidate", methodOrKey); } // Invalidate by method prefix @@ -101,6 +144,7 @@ export class SimpleCache { for (const key of this.store.keys()) { if (key.startsWith(prefix)) { this.store.delete(key); + this.emit("invalidate", key); } } } @@ -114,6 +158,7 @@ export class SimpleCache { for (const [key, entry] of this.store.entries()) { if (now > entry.expiresAt) { this.store.delete(key); + this.emit("expire", key); } } return { @@ -163,6 +208,7 @@ interface CacheEntry { export class Cache { private readonly store = new Map>(); private readonly ttlMs: number | undefined; + private readonly listeners = new Set(); /** * @param ttlMs Time-to-live in milliseconds. Omit (or pass `undefined`) @@ -172,11 +218,31 @@ export class Cache { this.ttlMs = ttlMs; } + /** + * Subscribe to cache lifecycle events (set, hit, miss, expire, invalidate). + * Returns an unsubscribe function. + */ + on(listener: CacheEventListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + + private emit(type: CacheEventType, key: string): void { + if (this.listeners.size === 0) return; + const event: CacheEvent = { type, key, timestamp: Date.now() }; + for (const listener of this.listeners) { + listener(event); + } + } + /** * Store `value` under `key`, recording the current wall-clock time. */ set(key: string, value: V): void { this.store.set(key, { value, writtenAt: Date.now() }); + this.emit("set", key); } /** @@ -188,11 +254,17 @@ export class Cache { */ get(key: string): V | undefined { const entry = this.store.get(key); - if (!entry) return undefined; + if (!entry) { + this.emit("miss", key); + return undefined; + } if (this.isExpired(entry)) { this.store.delete(key); + this.emit("expire", key); + this.emit("miss", key); return undefined; } + this.emit("hit", key); return entry.value; } @@ -205,6 +277,7 @@ export class Cache { if (!entry) return false; if (this.isExpired(entry)) { this.store.delete(key); + this.emit("expire", key); return false; } return true; @@ -219,18 +292,25 @@ export class Cache { for (const [key, entry] of this.store) { if (this.isExpired(entry)) { this.store.delete(key); + this.emit("expire", key); } } } /** Remove a specific entry by key. */ delete(key: string): void { - this.store.delete(key); + if (this.store.delete(key)) { + this.emit("invalidate", key); + } } /** Remove all entries. */ clear(): void { + const keys = Array.from(this.store.keys()); this.store.clear(); + for (const key of keys) { + this.emit("invalidate", key); + } } /** Number of entries currently in the store (including not-yet-evicted expired ones). */ diff --git a/src/cache/OptimisticCache.ts b/src/cache/OptimisticCache.ts index 4d04f78..dd939fe 100644 --- a/src/cache/OptimisticCache.ts +++ b/src/cache/OptimisticCache.ts @@ -107,14 +107,62 @@ export class OptimisticCache { /** * Read the current UI-facing value for an invoice: the most recently * applied still-pending optimistic prediction if one exists, otherwise - * the committed base value. + * the committed base value. Emits `hit`/`miss`/`expire` lifecycle events. */ get(invoiceId: string): T | undefined { const queue = this.pending.get(invoiceId); if (queue && queue.length > 0) { - return queue[queue.length - 1]!.predictedValue; + const value = queue[queue.length - 1]!.predictedValue; + this._emit({ type: "hit", invoiceId, value, timestamp: Date.now() }); + return value; } - return this.base.get(invoiceId); + + if (this._isExpired(invoiceId)) { + const expired = this.base.get(invoiceId); + this.base.delete(invoiceId); + this.expiries.delete(invoiceId); + this._emit({ type: "expire", invoiceId, value: expired, timestamp: Date.now() }); + return undefined; + } + + const value = this.base.get(invoiceId); + if (value === undefined) { + this._emit({ type: "miss", invoiceId, timestamp: Date.now() }); + } else { + this._emit({ type: "hit", invoiceId, value, timestamp: Date.now() }); + } + return value; + } + + /** + * Write a committed value into the base cache with an optional TTL + * (milliseconds). Emits a `set` event. + */ + set(invoiceId: string, value: T, ttlMs?: number): void { + this.base.set(invoiceId, value); + if (ttlMs !== undefined && ttlMs > 0) { + this.expiries.set(invoiceId, Date.now() + ttlMs); + } else { + this.expiries.delete(invoiceId); + } + this._emit({ type: "set", invoiceId, value, timestamp: Date.now() }); + } + + /** + * Invalidate a single invoice: drop any pending predictions and the + * committed base value, emitting an `invalidate` event. + */ + invalidate(invoiceId: string): void { + this.pending.delete(invoiceId); + this.expiries.delete(invoiceId); + this.base.delete(invoiceId); + this._emit({ type: "invalidate", invoiceId, timestamp: Date.now() }); + } + + /** Invalidate every cached invoice, emitting one `invalidate` per key. */ + invalidateAll(): void { + const keys = new Set([...this.pending.keys(), ...this.expiries.keys()]); + for (const invoiceId of keys) this.invalidate(invoiceId); } /** Number of optimistic mutations across all invoices awaiting commit/rollback. */ @@ -130,6 +178,12 @@ export class OptimisticCache { return () => this.rollbackHandlers.delete(handler); } + /** Register a listener for cache lifecycle events (set/hit/miss/expire/invalidate). */ + onEvent(handler: CacheEventHandler): () => void { + this.eventHandlers.add(handler); + return () => this.eventHandlers.delete(handler); + } + /** * Subscribe to the `revalidateError` event, emitted when a * stale-while-revalidate refresh fails. @@ -277,7 +331,7 @@ export class OptimisticCache { const commit: CommitFn = () => { if (settled) return; settled = true; - this.base.set(invoiceId, entry.predictedValue); + this.set(invoiceId, entry.predictedValue); this._removeEntry(entry); }; @@ -290,7 +344,7 @@ export class OptimisticCache { const stillPending = remaining && remaining.length > 0; const restoredValue = stillPending ? remaining![remaining!.length - 1]!.predictedValue : entry.rollbackValue; if (!stillPending) { - this.base.set(invoiceId, entry.rollbackValue); + this.set(invoiceId, entry.rollbackValue); } const event: RollbackEvent = { key, invoiceId, version, restoredValue }; @@ -306,6 +360,21 @@ export class OptimisticCache { return { commit, rollback, key }; } + private _isExpired(invoiceId: string): boolean { + const expiry = this.expiries.get(invoiceId); + return expiry !== undefined && Date.now() >= expiry; + } + + private _emit(event: CacheEvent): void { + for (const handler of this.eventHandlers) { + try { + handler(event); + } catch { + // Isolate listener failures from cache bookkeeping. + } + } + } + private _removeEntry(entry: OptimisticEntry): void { const queue = this.pending.get(entry.invoiceId); if (!queue) return; diff --git a/src/index.ts b/src/index.ts index cc3251d..3902782 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,5 +1,8 @@ /** - * @stellar-split/sdk — public API (core exports) + * SDK entry point. + * + * Exposes the public surface of the SDK along with a lightweight caching + * layer that supports per-entry TTLs and explicit invalidation. */ import type { Invoice } from "./types.js"; diff --git a/src/paymentForwardingRulesEngine.ts b/src/paymentForwardingRulesEngine.ts new file mode 100644 index 0000000..4cc1bf0 --- /dev/null +++ b/src/paymentForwardingRulesEngine.ts @@ -0,0 +1,177 @@ +/** + * Payment forwarding rules engine. + * + * Evaluates a set of ordered forwarding rules against an incoming payment and + * decides where (and whether) the payment should be forwarded. Emits events for + * rule evaluation and forwarding outcomes so callers can observe/audit decisions. + */ + +export type Payment = { + id: string; + amount: number; + currency: string; + source: string; + destination: string; + metadata?: Record; +}; + +export type RuleCondition = { + /** Field on the payment to inspect. */ + field: keyof Payment | string; + /** Comparison operator. */ + operator: 'eq' | 'neq' | 'gt' | 'gte' | 'lt' | 'lte' | 'in' | 'contains'; + /** Value to compare against. */ + value: unknown; +}; + +export type ForwardingRule = { + id: string; + /** Higher priority rules are evaluated first. */ + priority?: number; + /** All conditions must match for the rule to apply. */ + conditions: RuleCondition[]; + /** Destination to forward to when the rule matches. */ + forwardTo: string; + /** When true, stop evaluating further rules after this one matches. */ + terminal?: boolean; + enabled?: boolean; +}; + +export type ForwardingDecision = { + paymentId: string; + forwarded: boolean; + destination: string | null; + matchedRuleId: string | null; + reason: string; +}; + +export type RulesEngineEvent = + | { type: 'rule:evaluated'; ruleId: string; paymentId: string; matched: boolean } + | { type: 'rule:skipped'; ruleId: string; paymentId: string; reason: string } + | { type: 'forwarding:decided'; decision: ForwardingDecision } + | { type: 'forwarding:error'; paymentId: string; error: Error }; + +export type RulesEngineListener = (event: RulesEngineEvent) => void; + +function getField(payment: Payment, field: string): unknown { + if (field in payment) { + return (payment as Record)[field]; + } + return payment.metadata ? payment.metadata[field] : undefined; +} + +function compare(actual: unknown, operator: RuleCondition['operator'], expected: unknown): boolean { + switch (operator) { + case 'eq': + return actual === expected; + case 'neq': + return actual !== expected; + case 'gt': + return typeof actual === 'number' && typeof expected === 'number' && actual > expected; + case 'gte': + return typeof actual === 'number' && typeof expected === 'number' && actual >= expected; + case 'lt': + return typeof actual === 'number' && typeof expected === 'number' && actual < expected; + case 'lte': + return typeof actual === 'number' && typeof expected === 'number' && actual <= expected; + case 'in': + return Array.isArray(expected) && expected.includes(actual); + case 'contains': + return typeof actual === 'string' && typeof expected === 'string' && actual.includes(expected); + default: + return false; + } +} + +export class PaymentForwardingRulesEngine { + private rules: ForwardingRule[]; + private listeners: Set = new Set(); + + constructor(rules: ForwardingRule[] = []) { + this.rules = [...rules]; + } + + /** Replace the current rule set. */ + setRules(rules: ForwardingRule[]): void { + this.rules = [...rules]; + } + + /** Add a rule to the engine. */ + addRule(rule: ForwardingRule): void { + this.rules.push(rule); + } + + /** Subscribe to engine events. Returns an unsubscribe function. */ + on(listener: RulesEngineListener): () => void { + this.listeners.add(listener); + return () => this.listeners.delete(listener); + } + + private emit(event: RulesEngineEvent): void { + for (const listener of this.listeners) { + try { + listener(event); + } catch { + // Listener errors must not break rule evaluation. + } + } + } + + /** Evaluate a single rule against a payment. */ + evaluateRule(rule: ForwardingRule, payment: Payment): boolean { + return rule.conditions.every((condition) => + compare(getField(payment, condition.field as string), condition.operator, condition.value), + ); + } + + /** + * Evaluate all rules against a payment and return the forwarding decision. + * Rules are evaluated in priority order (highest first). + */ + evaluate(payment: Payment): ForwardingDecision { + const ordered = [...this.rules].sort((a, b) => (b.priority ?? 0) - (a.priority ?? 0)); + + for (const rule of ordered) { + if (rule.enabled === false) { + this.emit({ type: 'rule:skipped', ruleId: rule.id, paymentId: payment.id, reason: 'disabled' }); + continue; + } + + let matched = false; + try { + matched = this.evaluateRule(rule, payment); + } catch (error) { + this.emit({ type: 'forwarding:error', paymentId: payment.id, error: error as Error }); + continue; + } + + this.emit({ type: 'rule:evaluated', ruleId: rule.id, paymentId: payment.id, matched }); + + if (matched) { + const decision: ForwardingDecision = { + paymentId: payment.id, + forwarded: true, + destination: rule.forwardTo, + matchedRuleId: rule.id, + reason: `matched rule ${rule.id}`, + }; + this.emit({ type: 'forwarding:decided', decision }); + return decision; + } + + if (rule.terminal) { + break; + } + } + + const decision: ForwardingDecision = { + paymentId: payment.id, + forwarded: false, + destination: null, + matchedRuleId: null, + reason: 'no matching rule', + }; + this.emit({ type: 'forwarding:decided', decision }); + return decision; + } +}