diff --git a/examples/telemetry-hooks-example.ts b/examples/telemetry-hooks-example.ts index 93bffdc..c1b0ec5 100644 --- a/examples/telemetry-hooks-example.ts +++ b/examples/telemetry-hooks-example.ts @@ -281,4 +281,59 @@ client.setTelemetryHooks({ console.log("\nEven if hooks throw exceptions, SDK operations continue normally."); console.log("Hook exceptions are logged to console but don't propagate to your code."); +// Example 8: Built-in SDK performance profiling +console.log("\n=== Example 8: Built-in Performance Profiling ==="); + +// The SDK ships with a built-in profiler that records operation timings +// and emits lifecycle events (start/stop/mark/measure) without any extra setup. +client.startProfiling(); + +// Subscribe to profiling lifecycle events +const unsubscribeProfiling = client.onProfilingEvent((event) => { + switch (event.type) { + case "start": + console.log(`🟢 Profiling started at ${new Date(event.timestamp).toISOString()}`); + break; + case "mark": + console.log(`šŸ“ Mark "${event.name}" at ${event.timestamp}`); + break; + case "measure": + console.log(`ā±ļø Measure "${event.name}": ${event.durationMs}ms`); + break; + case "stop": + console.log(`šŸ”“ Profiling stopped at ${new Date(event.timestamp).toISOString()}`); + break; + } +}); + +async function demonstrateProfiling() { + console.log("\n=== Running Profiled SDK Operations ===\n"); + + client.mark("before-getInvoice"); + try { + await client.getInvoice("789"); + } catch (error) { + console.log("Expected error caught in main code"); + } + client.mark("after-getInvoice"); + client.measure("getInvoice-roundtrip", "before-getInvoice", "after-getInvoice"); + + // Retrieve aggregated profiling metrics + const metrics = client.getProfilingMetrics(); + console.log("\nšŸ“Š Built-in Profiling Metrics:"); + for (const [name, stats] of Object.entries(metrics)) { + console.log(` ${name}:`); + console.log(` Calls: ${stats.calls}`); + console.log(` Avg: ${stats.avgMs.toFixed(2)}ms`); + console.log(` Min: ${stats.minMs}ms, Max: ${stats.maxMs}ms`); + } + + // Stop profiling and clean up the event subscription + client.stopProfiling(); + unsubscribeProfiling(); + console.log("āœ… Profiling stopped and listener removed"); +} + +demonstrateProfiling().catch(console.error); + export { client, perfMonitor, errorTracker, analytics }; diff --git a/src/audit/AuditTrailHasher.ts b/src/audit/AuditTrailHasher.ts index 978e7d5..38ecb96 100644 --- a/src/audit/AuditTrailHasher.ts +++ b/src/audit/AuditTrailHasher.ts @@ -69,6 +69,106 @@ export class AuditTrailHasher { return entry; } + /** + * Records an invoice audit entry scoped to a tenant, enabling cross-tenant auditing. + * Emits an 'invoice.audited' lifecycle event. + */ + async auditInvoice(tenantId: string, invoiceId: string, event: AuditEvent): Promise { + const entry = await this.append(event); + const record: CrossTenantAuditRecord = { + tenantId, + invoiceId, + event, + entry, + recordedAt: Date.now(), + }; + this.crossTenantRecords.push(record); + this.emit({ + type: 'invoice.audited', + tenantId, + invoiceId, + entry, + timestamp: record.recordedAt, + }); + return record; + } + + /** + * Records a cross-tenant access attempt against an invoice and emits a + * 'invoice.cross-tenant-access' event so consumers can react to it. + */ + async recordCrossTenantAccess( + accessingTenantId: string, + invoiceTenantId: string, + invoiceId: string, + event: AuditEvent, + ): Promise { + const entry = await this.append(event); + const record: CrossTenantAuditRecord = { + tenantId: accessingTenantId, + invoiceId, + event, + entry, + recordedAt: Date.now(), + }; + this.crossTenantRecords.push(record); + this.emit({ + type: 'invoice.cross-tenant-access', + tenantId: accessingTenantId, + invoiceId, + entry, + timestamp: record.recordedAt, + }); + return record; + } + + /** + * Queries recorded cross-tenant audit entries, optionally filtered by tenant + * and/or invoice. Returns a defensive copy to preserve isolation. + */ + queryCrossTenantAudits(query: CrossTenantAuditQuery = {}): CrossTenantAuditRecord[] { + return this.crossTenantRecords + .filter(r => (query.tenantId === undefined || r.tenantId === query.tenantId)) + .filter(r => (query.invoiceId === undefined || r.invoiceId === query.invoiceId)) + .map(r => ({ ...r })); + } + + /** + * Verifies that a tenant's recorded audit entries are intact and match the + * expected chain root, enforcing cross-tenant isolation. + */ + async verifyTenantAudit( + tenantId: string, + expectedRoot: AuditTrailRoot, + ): Promise<{ valid: boolean; mismatchAt?: number; length?: number }> { + const tenantEntries = this.crossTenantRecords + .filter(r => r.tenantId === tenantId) + .map(r => r.entry); + const scoped = new AuditTrailHasher(tenantEntries); + return scoped.verify(expectedRoot); + } + + /** + * Registers a listener for cross-tenant audit lifecycle events. + */ + on(listener: CrossTenantAuditListener): () => void { + this.listeners.add(listener); + return () => this.listeners.delete(listener); + } + + /** + * Removes a previously registered listener. + */ + off(listener: CrossTenantAuditListener): void { + this.listeners.delete(listener); + } + + private emit(event: CrossTenantAuditEvent): void { + for (const listener of this.listeners) { + listener(event); + } + } + /** * Computes a Merkle root over all current chain entry hashes using pairwise SHA-256 combining */ diff --git a/src/auditLogger.ts b/src/auditLogger.ts index 1ab5c48..d6ec606 100644 --- a/src/auditLogger.ts +++ b/src/auditLogger.ts @@ -92,6 +92,11 @@ export class AuditLogger { } log(entry: AuditEntry): void { + this.debugLog("log", { + method: entry.method, + success: entry.success, + durationMs: entry.durationMs, + }); this.sink(entry); this.emit(entry); } @@ -202,4 +207,104 @@ export class AuditLogger { async exportSplitAuditTrail(invoiceId: string): Promise { return [...(this.splitAuditTrails.get(invoiceId) ?? [])]; } + + /** + * Subscribe to cross-tenant invoice audit lifecycle events. + * + * @returns an unsubscribe function. + */ + onCrossTenantAudit(listener: CrossTenantAuditEventListener): () => void { + this.crossTenantListeners.add(listener); + return () => { + this.crossTenantListeners.delete(listener); + }; + } + + /** + * Record a cross-tenant invoice audit entry. + * + * Persists the entry to the in-memory cross-tenant trail, writes a + * sanitized `AuditEntry` to the configured sink, and emits the appropriate + * lifecycle event: + * - `cross_tenant_access_detected` when the actor differs from the owner, + * - `cross_tenant_access_denied` when such access is unauthorized, + * - `invoice_audited` for every recorded entry. + */ + recordCrossTenantInvoiceAudit( + entry: CrossTenantInvoiceAuditEntry, + ): void { + this.crossTenantAudits.push(entry); + + this.log({ + timestamp: entry.timestamp, + method: "cross_tenant_invoice_audit", + params: this.sanitize({ + ownerTenantId: entry.ownerTenantId, + actorTenantId: entry.actorTenantId, + invoiceId: entry.invoiceId, + action: entry.action, + authorized: entry.authorized, + ...(entry.metadata ?? {}), + }), + success: entry.authorized, + durationMs: 0, + }); + + const isCrossTenant = entry.actorTenantId !== entry.ownerTenantId; + if (isCrossTenant) { + this.emitCrossTenantAudit({ + type: "cross_tenant_access_detected", + entry, + }); + if (!entry.authorized) { + this.emitCrossTenantAudit({ + type: "cross_tenant_access_denied", + entry, + }); + } + } + this.emitCrossTenantAudit({ type: "invoice_audited", entry }); + } + + /** + * Query recorded cross-tenant invoice audit entries. + * + * All filters are optional and combined with AND semantics. Results are + * returned in the order they were recorded. + */ + queryCrossTenantInvoiceAudits(filter?: { + ownerTenantId?: string; + actorTenantId?: string; + invoiceId?: string; + action?: string; + authorized?: boolean; + }): CrossTenantInvoiceAuditEntry[] { + return this.crossTenantAudits.filter((entry) => { + if (filter?.ownerTenantId && entry.ownerTenantId !== filter.ownerTenantId) { + return false; + } + if (filter?.actorTenantId && entry.actorTenantId !== filter.actorTenantId) { + return false; + } + if (filter?.invoiceId && entry.invoiceId !== filter.invoiceId) { + return false; + } + if (filter?.action && entry.action !== filter.action) { + return false; + } + if ( + filter?.authorized !== undefined && + entry.authorized !== filter.authorized + ) { + return false; + } + return true; + }); + } + + private emitCrossTenantAudit(event: CrossTenantAuditEvent): void { + for (const listener of this.crossTenantListeners) { + listener(event); + } + } } diff --git a/src/cache.ts b/src/cache.ts index 3bf155b..c62f398 100644 --- a/src/cache.ts +++ b/src/cache.ts @@ -38,7 +38,7 @@ export class SimpleCache { private maxEntries: number; private readonly listeners = new Set(); - constructor(config?: number | { enabled?: boolean; ttl?: Record; ttlMs?: number; maxEntries?: number }) { + constructor(config?: number | { enabled?: boolean; ttl?: Record; ttlMs?: number; maxEntries?: number; debug?: boolean | DebugModeOptions }) { if (typeof config === "number") { this.enabled = true; this.maxEntries = 1000; @@ -51,6 +51,18 @@ export class SimpleCache { this.ttlConfig["default"] = config.ttlMs; } } + this.debug = new DebugMode( + typeof config === "object" && config?.debug !== undefined + ? typeof config.debug === "boolean" + ? { enabled: config.debug } + : config.debug + : undefined + ); + } + + /** Access the debug-mode controller for this cache instance. */ + getDebugMode(): DebugMode { + return this.debug; } /** @@ -147,10 +159,12 @@ export class SimpleCache { this.emit("invalidate", key); } } + this.debug.log(`[cache] invalidate ${methodOrKey}`); } clear(): void { this.store.clear(); + this.debug.log("[cache] clear"); } getStats(): CacheStats { @@ -213,9 +227,18 @@ export class Cache { /** * @param ttlMs Time-to-live in milliseconds. Omit (or pass `undefined`) * for no-expiry behaviour. + * @param debug Optional debug-mode configuration for verbose logging. */ - constructor(ttlMs?: number) { + constructor(ttlMs?: number, debug?: boolean | DebugModeOptions) { this.ttlMs = ttlMs; + this.debug = new DebugMode( + typeof debug === "boolean" ? { enabled: debug } : debug + ); + } + + /** Access the debug-mode controller for this cache instance. */ + getDebugMode(): DebugMode { + return this.debug; } /** diff --git a/src/channelReconciler.ts b/src/channelReconciler.ts index fcbb7bd..3063a68 100644 --- a/src/channelReconciler.ts +++ b/src/channelReconciler.ts @@ -38,6 +38,37 @@ export type ChannelStateFetcher = ( payer: string ) => Promise; +/** Lifecycle events emitted during invoice reconciliation. */ +export type ReconciliationEvent = + | { type: "reconciliation:started"; invoiceId: string; payer: string } + | { type: "reconciliation:matched"; invoiceId: string; payer: string; result: ChannelReconciliationResult } + | { type: "reconciliation:discrepancy"; invoiceId: string; payer: string; result: ChannelReconciliationResult } + | { type: "reconciliation:completed"; invoiceId: string; payer: string; result: ChannelReconciliationResult } + | { type: "reconciliation:error"; invoiceId: string; payer: string; error: unknown }; + +/** Listener invoked for every reconciliation lifecycle event. */ +export type ReconciliationEventListener = (event: ReconciliationEvent) => void; + +const _listeners = new Set(); + +/** + * Subscribe to reconciliation lifecycle events. + * + * @returns An unsubscribe function that removes the listener. + */ +export function onReconciliationEvent(listener: ReconciliationEventListener): () => void { + _listeners.add(listener); + return () => { + _listeners.delete(listener); + }; +} + +function emit(event: ReconciliationEvent): void { + for (const listener of _listeners) { + listener(event); + } +} + let _fetcher: ChannelStateFetcher | null = null; /** Register (or clear) the function that reads on-chain channel state. */ @@ -64,23 +95,39 @@ export async function reconcileChannel( localPayments: bigint[], fetcher?: ChannelStateFetcher ): Promise { - const resolveFetcher = fetcher ?? _fetcher; - if (!resolveFetcher) { - throw new ChannelReconciliationError( - "No channel state fetcher registered. Call registerChannelStateFetcher() first." - ); - } + emit({ type: "reconciliation:started", invoiceId, payer }); - const { deposited, balance: onChainBalance } = await resolveFetcher(invoiceId, payer); + try { + const resolveFetcher = fetcher ?? _fetcher; + if (!resolveFetcher) { + throw new ChannelReconciliationError( + "No channel state fetcher registered. Call registerChannelStateFetcher() first." + ); + } - const totalPaid = localPayments.reduce((sum, amt) => sum + amt, 0n); - const expectedBalance = deposited - totalPaid; - const delta = onChainBalance - expectedBalance; + const { deposited, balance: onChainBalance } = await resolveFetcher(invoiceId, payer); - return { - inSync: delta === 0n, - onChainBalance, - expectedBalance, - delta, - }; + const totalPaid = localPayments.reduce((sum, amt) => sum + amt, 0n); + const expectedBalance = deposited - totalPaid; + const delta = onChainBalance - expectedBalance; + + const result: ChannelReconciliationResult = { + inSync: delta === 0n, + onChainBalance, + expectedBalance, + delta, + }; + + if (result.inSync) { + emit({ type: "reconciliation:matched", invoiceId, payer, result }); + } else { + emit({ type: "reconciliation:discrepancy", invoiceId, payer, result }); + } + + emit({ type: "reconciliation:completed", invoiceId, payer, result }); + return result; + } catch (error) { + emit({ type: "reconciliation:error", invoiceId, payer, error }); + throw error; + } } \ No newline at end of file diff --git a/src/telemetryHooks.ts b/src/telemetryHooks.ts index c2ac810..d573c69 100644 --- a/src/telemetryHooks.ts +++ b/src/telemetryHooks.ts @@ -55,6 +55,50 @@ export interface TelemetryCallEndParams { traceId?: string; } +/** + * A single recorded performance measurement produced by the built-in profiler. + */ +export interface ProfileMeasurement { + /** The SDK method or operation name that was measured. */ + name: string; + /** Duration of the operation in milliseconds. */ + durationMs: number; + /** Timestamp when the measurement was recorded (milliseconds since epoch). */ + timestamp: number; + /** Optional trace ID correlating this measurement with an SDK call. */ + traceId?: string; + /** Optional arbitrary metadata attached to the measurement. */ + metadata?: Record; +} + +/** + * Aggregated statistics for a profiled operation name. + */ +export interface ProfileStats { + /** The operation name these stats describe. */ + name: string; + /** Number of recorded measurements. */ + count: number; + /** Total accumulated duration in milliseconds. */ + totalMs: number; + /** Minimum observed duration in milliseconds. */ + minMs: number; + /** Maximum observed duration in milliseconds. */ + maxMs: number; + /** Mean duration in milliseconds. */ + avgMs: number; +} + +/** + * Parameters passed to the onProfile hook when a measurement is recorded. + */ +export interface TelemetryProfileParams { + /** The recorded measurement. */ + measurement: ProfileMeasurement; + /** Aggregated stats for the measured operation name. */ + stats: ProfileStats; +} + /** * Telemetry hooks that can be registered with the SDK. * All hooks are optional and fire-and-forget. @@ -81,6 +125,13 @@ export interface TelemetryHooks { * @param params - Call results including method name, duration, success status, and optional error. */ onCallEnd?(params: TelemetryCallEndParams): void; + + /** + * Called whenever the built-in profiler records a measurement. + * + * @param params - The measurement and its aggregated stats. + */ + onProfile?(params: TelemetryProfileParams): void; } /** @@ -165,10 +216,228 @@ export class TelemetryHookManager { } } + /** + * Invoke the onProfile hook if registered. + * Exceptions within the hook are caught and logged but do not propagate. + * + * @param params - Profile measurement and stats. + */ + fireOnProfile(params: TelemetryProfileParams): void { + if (!this.hooks.onProfile) { + return; + } + + try { + this.hooks.onProfile(params); + } catch (hookError) { + // Fire-and-forget: hook errors must not propagate + console.error("[TelemetryHook] onProfile hook threw an exception:", hookError); + } + } + /** * Check if any hooks are registered. */ hasHooks(): boolean { - return !!(this.hooks.onError || this.hooks.onCallStart || this.hooks.onCallEnd); + return !!( + this.hooks.onError || + this.hooks.onCallStart || + this.hooks.onCallEnd || + this.hooks.onProfile + ); + } +} + +/** + * Built-in SDK performance profiler. + * + * Records operation timings, aggregates per-operation statistics, and emits + * lifecycle events (start/stop/mark/measure) through the telemetry hook manager. + * + * The profiler is disabled by default and must be explicitly enabled via + * {@link enable} so it adds zero overhead unless opted into. + */ +export class SdkProfiler { + private enabled = false; + private readonly measurements: ProfileMeasurement[] = []; + private readonly stats = new Map(); + private readonly activeMarks = new Map(); + + constructor(private readonly hookManager?: TelemetryHookManager) {} + + /** + * Enable profiling. Subsequent {@link measure} and {@link mark}/{@link endMark} + * calls will record measurements. + */ + enable(): void { + this.enabled = true; + } + + /** + * Disable profiling. Existing recorded measurements are retained until cleared. + */ + disable(): void { + this.enabled = false; + } + + /** + * Whether profiling is currently enabled. + */ + isEnabled(): boolean { + return this.enabled; + } + + /** + * Record a completed measurement for the given operation name. + * No-op when profiling is disabled. + * + * @param name - The operation name being measured. + * @param durationMs - Duration of the operation in milliseconds. + * @param options - Optional trace ID and metadata. + * @returns The recorded measurement, or undefined when disabled. + */ + measure( + name: string, + durationMs: number, + options?: { traceId?: string; metadata?: Record } + ): ProfileMeasurement | undefined { + if (!this.enabled) { + return undefined; + } + + const measurement: ProfileMeasurement = { + name, + durationMs, + timestamp: Date.now(), + traceId: options?.traceId, + metadata: options?.metadata, + }; + + this.measurements.push(measurement); + const stats = this.updateStats(measurement); + this.hookManager?.fireOnProfile({ measurement, stats }); + + return measurement; + } + + /** + * Start timing an operation. Pairs with {@link endMark}. + * No-op when profiling is disabled. + * + * @param name - The operation name to start timing. + */ + mark(name: string): void { + if (!this.enabled) { + return; + } + this.activeMarks.set(name, Date.now()); + } + + /** + * Finish timing an operation started with {@link mark} and record the measurement. + * No-op when profiling is disabled or no matching mark exists. + * + * @param name - The operation name to finish timing. + * @param options - Optional trace ID and metadata. + * @returns The recorded measurement, or undefined when disabled/unmatched. + */ + endMark( + name: string, + options?: { traceId?: string; metadata?: Record } + ): ProfileMeasurement | undefined { + if (!this.enabled) { + return undefined; + } + + const start = this.activeMarks.get(name); + if (start === undefined) { + return undefined; + } + + this.activeMarks.delete(name); + return this.measure(name, Date.now() - start, options); + } + + /** + * Convenience helper that times an async operation and records a measurement. + * No-op passthrough when profiling is disabled. + * + * @param name - The operation name being measured. + * @param fn - The async function to execute and time. + * @param options - Optional trace ID and metadata. + * @returns The resolved value of the wrapped function. + */ + async profile( + name: string, + fn: () => Promise, + options?: { traceId?: string; metadata?: Record } + ): Promise { + if (!this.enabled) { + return fn(); + } + + const start = Date.now(); + try { + return await fn(); + } finally { + this.measure(name, Date.now() - start, options); + } + } + + /** + * Get all recorded measurements (a defensive copy). + */ + getMeasurements(): ProfileMeasurement[] { + return [...this.measurements]; + } + + /** + * Get aggregated stats for a single operation name, if any. + */ + getStats(name: string): ProfileStats | undefined { + const stats = this.stats.get(name); + return stats ? { ...stats } : undefined; + } + + /** + * Get aggregated stats for all profiled operation names. + */ + getAllStats(): ProfileStats[] { + return Array.from(this.stats.values(), (stats) => ({ ...stats })); + } + + /** + * Clear all recorded measurements, stats, and active marks. + */ + clear(): void { + this.measurements.length = 0; + this.stats.clear(); + this.activeMarks.clear(); + } + + private updateStats(measurement: ProfileMeasurement): ProfileStats { + const existing = this.stats.get(measurement.name); + + const next: ProfileStats = existing + ? { + name: measurement.name, + count: existing.count + 1, + totalMs: existing.totalMs + measurement.durationMs, + minMs: Math.min(existing.minMs, measurement.durationMs), + maxMs: Math.max(existing.maxMs, measurement.durationMs), + avgMs: 0, + } + : { + name: measurement.name, + count: 1, + totalMs: measurement.durationMs, + minMs: measurement.durationMs, + maxMs: measurement.durationMs, + avgMs: 0, + }; + + next.avgMs = next.totalMs / next.count; + this.stats.set(measurement.name, next); + return next; } }