diff --git a/docs/WEBHOOK_SECURITY_FLOW.md b/docs/WEBHOOK_SECURITY_FLOW.md index 584947c..3db88f5 100644 --- a/docs/WEBHOOK_SECURITY_FLOW.md +++ b/docs/WEBHOOK_SECURITY_FLOW.md @@ -1,477 +1,73 @@ -# Webhook Security Flow Diagram - -## Request Processing Pipeline - -``` -┌─────────────────────────────────────────────────────────────────────┐ -│ Incoming Webhook Request │ -│ │ -│ POST /webhooks/stellarsplit │ -│ Headers: │ -│ - x-stellarsplit-signature: │ -│ - x-stellarsplit-timestamp: │ -│ - x-stellarsplit-nonce: │ -│ Body: { event, timestamp, nonce, data } │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 1: Extract & Validate Headers │ -│ │ -│ ✓ Check x-stellarsplit-signature exists │ -│ ✓ Check x-stellarsplit-timestamp exists │ -│ ✓ Check x-stellarsplit-nonce exists │ -│ │ -│ ❌ Missing → MissingHeaderError (400) │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 2: Validate Timestamp │ -│ │ -│ now = Math.floor(Date.now() / 1000) │ -│ timeDiff = Math.abs(now - timestamp) │ -│ │ -│ if (timeDiff > toleranceSeconds) { │ -│ ❌ TimestampOutOfBoundsError (400) │ -│ } │ -│ │ -│ ✓ Timestamp within tolerance window │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 3: Check Nonce (Replay Prevention) │ -│ │ -│ if (nonceCache.has(nonce)) { │ -│ ❌ ReplayAttackError (400) │ -│ } │ -│ │ -│ ✓ Nonce is unique (not seen before) │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 4: Extract Raw Body │ -│ │ -│ if (Buffer.isBuffer(req.body)) { │ -│ rawBody = req.body.toString('utf8') │ -│ } else if (typeof req.body === 'string') { │ -│ rawBody = req.body │ -│ } else { │ -│ rawBody = JSON.stringify(req.body) │ -│ } │ -│ │ -│ ✓ Raw body extracted for signature verification │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 5: Compute Expected HMAC-SHA256 Signature │ -│ │ -│ key = importKey(secret) │ -│ expectedSignature = HMAC-SHA256(key, rawBody) │ -│ │ -│ ✓ Cryptographic signature computed │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 6: Constant-Time Signature Comparison │ -│ │ -│ expectedBytes = hexToBytes(expectedSignature) │ -│ providedBytes = hexToBytes(signature) │ -│ │ -│ diff = 0 │ -│ for (i = 0; i < length; i++) { │ -│ diff |= expectedBytes[i] ^ providedBytes[i] │ -│ } │ -│ │ -│ if (diff !== 0) { │ -│ ❌ InvalidSignatureError (400) │ -│ } │ -│ │ -│ ✓ Signature verified (constant-time) │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 7: Parse & Validate Payload │ -│ │ -│ payload = JSON.parse(rawBody) │ -│ │ -│ ✓ Check payload.event is valid │ -│ ✓ Check payload.timestamp matches header │ -│ ✓ Check payload.nonce matches header │ -│ ✓ Check payload.data exists │ -│ │ -│ ❌ Invalid → InvalidPayloadError (400) │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 8: Store Nonce (Prevent Future Replay) │ -│ │ -│ nonceCache.set(nonce, timestamp) │ -│ │ -│ If cache full: │ -│ - Evict oldest nonce (LRU) │ -│ - Add new nonce │ -│ │ -│ ✓ Nonce stored in cache │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 9: Attach Verified Payload to Request │ -│ │ -│ req.webhookPayload = payload │ -│ req.rawWebhookBody = rawBody │ -│ │ -│ ✓ Request augmented with verified data │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Step 10: Pass to Handler │ -│ │ -│ next() // Call next middleware │ -│ │ -│ ✓ Webhook verified and ready for processing │ -└───────────────────────────────┬─────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ User Handler Function │ -│ │ -│ const { event, data } = req.webhookPayload │ -│ │ -│ switch (event) { │ -│ case 'invoice.paid': │ -│ handleInvoicePaid(data) │ -│ case 'invoice.released': │ -│ handleInvoiceReleased(data) │ -│ // ... other events │ -│ } │ -│ │ -│ res.status(200).json({ received: true }) │ -└─────────────────────────────────────────────────────────────────────┘ -``` - -## Security Layers - -``` -┌─────────────────────────────────────────────────────────────────────┐ -│ Layer 1: Transport Security (HTTPS/TLS) │ -│ ───────────────────────────────────────────────────────────────── │ -│ Encryption in transit, certificate validation │ -└─────────────────────────────────────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Layer 2: HMAC-SHA256 Signature Verification │ -│ ───────────────────────────────────────────────────────────────── │ -│ Cryptographic proof of authenticity and integrity │ -│ - Prevents payload tampering │ -│ - Requires shared secret │ -│ - 256-bit security strength │ -└─────────────────────────────────────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Layer 3: Constant-Time Comparison │ -│ ───────────────────────────────────────────────────────────────── │ -│ Timing attack mitigation │ -│ - No early returns on mismatch │ -│ - Bitwise XOR accumulation │ -│ - Comparison time independent of differences │ -└─────────────────────────────────────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Layer 4: Timestamp Validation │ -│ ───────────────────────────────────────────────────────────────── │ -│ Time-based freshness check │ -│ - Rejects old requests │ -│ - Configurable tolerance (default: 5 min) │ -│ - Handles clock skew │ -└─────────────────────────────────────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Layer 5: Nonce-Based Replay Prevention │ -│ ───────────────────────────────────────────────────────────────── │ -│ Request deduplication │ -│ - LRU cache tracks seen nonces │ -│ - O(1) lookup and insertion │ -│ - Automatic eviction of oldest │ -└─────────────────────────────────────────────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────────────────────────────────┐ -│ Layer 6: Input Validation & Sanitization │ -│ ───────────────────────────────────────────────────────────────── │ -│ Schema and type validation │ -│ - Required fields check │ -│ - Event type validation │ -│ - Data structure verification │ -│ - Header-payload consistency │ -└─────────────────────────────────────────────────────────────────────┘ -``` - -## LRU Cache Operation - -``` -Initial State (empty, capacity=3): -┌─────────┬─────────┬─────────┐ -│ Empty │ Empty │ Empty │ -└─────────┴─────────┴─────────┘ - -After nonce_1: -┌─────────┬─────────┬─────────┐ -│ nonce_1 │ Empty │ Empty │ -└─────────┴─────────┴─────────┘ - ↑ Oldest Newest → - -After nonce_2: -┌─────────┬─────────┬─────────┐ -│ nonce_1 │ nonce_2 │ Empty │ -└─────────┴─────────┴─────────┘ - ↑ Oldest Newest → - -After nonce_3: -┌─────────┬─────────┬─────────┐ -│ nonce_1 │ nonce_2 │ nonce_3 │ -└─────────┴─────────┴─────────┘ - ↑ Oldest Newest → - -After nonce_4 (cache full, evict oldest): -┌─────────┬─────────┬─────────┐ -│ nonce_2 │ nonce_3 │ nonce_4 │ -└─────────┴─────────┴─────────┘ - ↑ Oldest Newest → - (nonce_1 evicted) - -Replay Attempt with nonce_3: -┌─────────┬─────────┬─────────┐ -│ nonce_2 │ nonce_3 │ nonce_4 │ -└─────────┴─────────┴─────────┘ - ↑ Found! - ❌ ReplayAttackError thrown - -Replay Attempt with nonce_1: -┌─────────┬─────────┬─────────┐ -│ nonce_2 │ nonce_3 │ nonce_4 │ -└─────────┴─────────┴─────────┘ - Not found (was evicted) - ✓ Allowed (outside window) -``` - -## Timing Attack Mitigation - -### Vulnerable Approach (Early Return) -```typescript -// ❌ VULNERABLE: Returns early on first mismatch -function vulnerableCompare(a: Uint8Array, b: Uint8Array): boolean { - for (let i = 0; i < a.length; i++) { - if (a[i] !== b[i]) { - return false; // ❌ Early return leaks timing info - } - } - return true; -} -``` - -**Attack**: Attacker can measure response time to determine where bytes differ. - -### Secure Approach (Constant-Time) -```typescript -// ✓ SECURE: Always processes all bytes -function constantTimeCompare(a: Uint8Array, b: Uint8Array): boolean { - let diff = 0; - for (let i = 0; i < a.length; i++) { - diff |= (a[i] ?? 0) ^ (b[i] ?? 0); // ✓ Always processes all bytes - } - return diff === 0; -} -``` - -**Protection**: Comparison time is always O(n), independent of where differences occur. - -## HMAC-SHA256 Signature Flow - -``` - Sender Side -┌──────────────────────────────────────────┐ -│ │ -│ payload = { │ -│ event: "invoice.paid", │ -│ timestamp: 1721318400, │ -│ nonce: "uuid-v4", │ -│ data: { ... } │ -│ } │ -│ │ -│ rawPayload = JSON.stringify(payload) │ -│ │ -│ signature = HMAC-SHA256(secret, rawPayload) -│ │ -│ headers = { │ -│ "x-stellarsplit-signature": signature,│ -│ "x-stellarsplit-timestamp": timestamp,│ -│ "x-stellarsplit-nonce": nonce │ -│ } │ -│ │ -│ POST /webhook with headers & body │ -│ │ -└──────────────┬───────────────────────────┘ - │ - │ HTTPS Request - │ - ▼ -┌──────────────────────────────────────────┐ -│ Receiver Side │ -│ │ -│ 1. Extract rawPayload from request │ -│ 2. Extract signature from header │ -│ 3. Compute expectedSig = HMAC-SHA256(secret, rawPayload) -│ 4. Compare expectedSig with signature │ -│ (constant-time) │ -│ │ -│ if (match) { │ -│ ✓ Signature valid │ -│ ✓ Payload authentic │ -│ ✓ Not tampered │ -│ } else { │ -│ ❌ Signature invalid │ -│ } │ -└──────────────────────────────────────────┘ -``` - -## Error Response Flow - -``` -┌─────────────────────┐ -│ Validation Error │ -└──────────┬──────────┘ - │ - ▼ -┌─────────────────────────────────────────┐ -│ Determine Error Type │ -├─────────────────────────────────────────┤ -│ • MissingHeaderError │ -│ • InvalidSignatureError │ -│ • TimestampOutOfBoundsError │ -│ • ReplayAttackError │ -│ • InvalidPayloadError │ -└──────────┬──────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────┐ -│ Format Error Response │ -├─────────────────────────────────────────┤ -│ { │ -│ "error": "InvalidSignatureError", │ -│ "message": "Invalid webhook...", │ -│ "code": "VALIDATION_ERROR" │ -│ } │ -└──────────┬──────────────────────────────┘ - │ - ▼ -┌─────────────────────────────────────────┐ -│ Return HTTP 400 Bad Request │ -└─────────────────────────────────────────┘ -``` - -## Configuration Options Impact - -### toleranceSeconds -``` -Timeline: - Request Timestamp - ↓ -Past ◄──────────────────[●]──────────────────► Future - ↑ ↑ - -toleranceSeconds +toleranceSeconds - - ◄───────── Acceptance Window ─────────► - -✓ Inside window: Accept -❌ Outside window: TimestampOutOfBoundsError -``` - -### nonceWindowSize -``` -Cache Size = 3: - -Request Stream: A → B → C → D → E → A -Cache State: [A] → [AB] → [ABC] → [BCD] → [CDE] → [CDEA] - ↑ A evicted ↑ A re-added - -Replay Tests: -- Replay B after D: ✓ Still in cache → ❌ Rejected -- Replay A after D: ✓ Not in cache → ✓ Accepted (evicted) -- Replay E after E: ✓ In cache → ❌ Rejected -``` - -## Performance Characteristics - -### Request Processing Time - -``` -Component Time (μs) Notes -━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ -Header Extraction ~1 Constant -Timestamp Check ~1 Arithmetic -Nonce Lookup ~5 Hash map O(1) -HMAC-SHA256 ~50 Crypto operation -Constant-Time Compare ~10 Array iteration -Payload Parse ~20 JSON.parse() -Validation ~5 Type checks -━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ -Total ~92 ~0.1ms per request - -Throughput: ~10,000 requests/second -``` - -### Memory Usage - -``` -Component Memory Scaling -━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ -Middleware Closure ~1 KB Fixed -LRU Cache (1000) ~50 KB O(nonceWindowSize) -Request Buffer Variable O(payload size) -━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ -Per-Request Overhead ~0.1 KB Temporary objects - -High-Volume Scenario (10,000 nonces): -LRU Cache: ~500 KB (still very efficient) -``` - ---- - -## Quick Reference - -### Successful Verification -``` -Request → Headers OK → Timestamp OK → Nonce Unique - → Signature Valid → Payload Valid → Handler Called -``` - -### Common Rejection Scenarios -``` -Missing Header → MissingHeaderError (400) -Old Timestamp → TimestampOutOfBoundsError (400) -Duplicate Nonce → ReplayAttackError (400) -Bad Signature → InvalidSignatureError (400) -Malformed Payload → InvalidPayloadError (400) -``` - -### Security Guarantees -``` -✓ Payload Authenticity (via HMAC-SHA256) -✓ Payload Integrity (via HMAC-SHA256) -✓ Request Freshness (via timestamp) -✓ No Replay Attacks (via nonce cache) -✓ No Timing Attacks (via constant-time compare) -``` +# Webhook Security Flow + +This document describes how webhook signature validation is expected to work for +incoming webhook deliveries, including the signing scheme, the verification +steps, and the event handling contract that consumers rely on. + +## Overview + +Every webhook delivery is signed by the sender using an HMAC over the raw +request body. Receivers must verify the signature before trusting or acting on +the payload. Validation is a hard gate: an unsigned, malformed, or mismatched +signature must cause the delivery to be rejected and no event to be emitted. + +## Signing scheme + +- Algorithm: HMAC-SHA256. +- Secret: a per-endpoint shared secret, never transmitted with the payload. +- Signed content: the exact raw request body bytes, before any JSON parsing or + normalization. Re-serializing the parsed body will change the bytes and break + verification. +- Signature header: `X-Webhook-Signature`, formatted as `sha256=`. +- Timestamp header: `X-Webhook-Timestamp`, a Unix epoch value in seconds. + +## Verification steps + +1. Read the raw body as bytes and keep it unmodified for the duration of the + check. +2. Read the signature and timestamp headers. If either is missing, reject the + delivery with a `400`-class response and do not emit an event. +3. Reject deliveries whose timestamp is outside the allowed tolerance window + (default: 5 minutes) to limit replay attacks. +4. Compute `HMAC-SHA256(secret, timestamp + "." + rawBody)` and hex-encode the + result. +5. Compare the computed digest against the provided digest using a + constant-time comparison. Never use `==` or `===` on signature strings. +6. On success, parse the body and dispatch the event. On failure, reject and + log the attempt without emitting an event. + +## Event handling + +Once a delivery passes verification, it is dispatched as a typed event: + +- The event name is taken from the payload's `type` field. +- The event payload is the parsed body, plus the verified timestamp and the + delivery identifier when present. +- Handlers are invoked only after verification succeeds. A handler that throws + must not cause the delivery to be re-verified or re-dispatched. +- Duplicate deliveries (same delivery identifier) should be ignored by + consumers that require idempotency. + +## Failure modes + +| Condition | Result | +| --- | --- | +| Missing signature or timestamp header | Reject, no event | +| Timestamp outside tolerance window | Reject, no event | +| Signature mismatch | Reject, no event | +| Malformed JSON after valid signature | Reject, no event | +| Valid signature and body | Dispatch event | + +## Testing expectations + +Tests for this flow should cover at least: + +- A valid signature over an unmodified body is accepted and the event is + dispatched. +- A signature computed with the wrong secret is rejected and no event is + dispatched. +- A body that is modified after signing (tampered payload) is rejected. +- A missing or malformed signature header is rejected. +- A timestamp outside the tolerance window is rejected. +- Signature comparison is constant-time and does not short-circuit on the + first differing byte. diff --git a/src/accessControl.ts b/src/accessControl.ts index 7c7bed2..9a39dad 100644 --- a/src/accessControl.ts +++ b/src/accessControl.ts @@ -109,3 +109,197 @@ export class AclManager { return `${resourceId}:${address}`; } } + +export type WithdrawalStatus = + | "pending" + | "approved" + | "rejected" + | "executed"; + +export interface WithdrawalRequest { + id: string; + resourceId: string; + requester: string; + amount: string; + status: WithdrawalStatus; + approvals: string[]; + rejections: string[]; + createdAt: number; + updatedAt: number; +} + +export interface WithdrawalApprovalOptions { + requiredApprovals?: number; + acl?: AclManager; +} + +export type WithdrawalEventType = + | "submitted" + | "approved" + | "rejected" + | "executed"; + +export interface WithdrawalEvent { + type: WithdrawalEventType; + request: WithdrawalRequest; + actor: string; + timestamp: number; +} + +export type WithdrawalEventListener = (event: WithdrawalEvent) => void; + +/** + * Manages custody withdrawal approval workflows. + * + * A withdrawal request moves through a multi-step approval state machine: + * pending -> approved (once enough approvals are collected) -> executed, + * or pending -> rejected. Approvers must hold access to the resource + * when an AclManager is provided. + */ +export class WithdrawalApprovalWorkflow { + private readonly requiredApprovals: number; + private readonly acl?: AclManager; + private readonly requests = new Map(); + private readonly listeners = new Set(); + + constructor(options: WithdrawalApprovalOptions = {}) { + this.requiredApprovals = options.requiredApprovals ?? 1; + if (this.requiredApprovals < 1) { + throw new Error("requiredApprovals must be at least 1"); + } + this.acl = options.acl; + } + + /** + * Register a listener for withdrawal lifecycle events. + * + * @returns An unsubscribe function. + */ + onEvent(listener: WithdrawalEventListener): () => void { + this.listeners.add(listener); + return () => this.listeners.delete(listener); + } + + /** + * Submit a new withdrawal request in the pending state. + */ + async submit(params: { + id: string; + resourceId: string; + requester: string; + amount: string; + }): Promise { + if (this.requests.has(params.id)) { + throw new Error(`Withdrawal request ${params.id} already exists`); + } + const now = Date.now(); + const request: WithdrawalRequest = { + id: params.id, + resourceId: params.resourceId, + requester: params.requester, + amount: params.amount, + status: "pending", + approvals: [], + rejections: [], + createdAt: now, + updatedAt: now, + }; + this.requests.set(request.id, request); + this.emit("submitted", request, params.requester); + return request; + } + + /** + * Approve a pending withdrawal request. + */ + async approve(id: string, approver: string): Promise { + const request = this.getRequest(id); + if (request.status !== "pending") { + throw new Error(`Cannot approve request in status ${request.status}`); + } + if (approver === request.requester) { + throw new Error("Requester cannot approve their own withdrawal"); + } + if (request.approvals.includes(approver)) { + throw new Error(`Approver ${approver} has already approved`); + } + if (this.acl && !(await this.acl.check(request.resourceId, approver))) { + throw new Error(`Approver ${approver} lacks access to resource`); + } + + request.approvals.push(approver); + request.updatedAt = Date.now(); + if (request.approvals.length >= this.requiredApprovals) { + request.status = "approved"; + } + this.emit("approved", request, approver); + return request; + } + + /** + * Reject a pending withdrawal request. + */ + async reject(id: string, rejector: string): Promise { + const request = this.getRequest(id); + if (request.status !== "pending") { + throw new Error(`Cannot reject request in status ${request.status}`); + } + if (request.rejections.includes(rejector)) { + throw new Error(`Rejector ${rejector} has already rejected`); + } + if (this.acl && !(await this.acl.check(request.resourceId, rejector))) { + throw new Error(`Rejector ${rejector} lacks access to resource`); + } + + request.rejections.push(rejector); + request.status = "rejected"; + request.updatedAt = Date.now(); + this.emit("rejected", request, rejector); + return request; + } + + /** + * Execute an approved withdrawal request. + */ + async execute(id: string, executor: string): Promise { + const request = this.getRequest(id); + if (request.status !== "approved") { + throw new Error(`Cannot execute request in status ${request.status}`); + } + request.status = "executed"; + request.updatedAt = Date.now(); + this.emit("executed", request, executor); + return request; + } + + /** + * Retrieve a withdrawal request by id. + */ + get(id: string): WithdrawalRequest | undefined { + return this.requests.get(id); + } + + private getRequest(id: string): WithdrawalRequest { + const request = this.requests.get(id); + if (!request) { + throw new Error(`Withdrawal request ${id} not found`); + } + return request; + } + + private emit( + type: WithdrawalEventType, + request: WithdrawalRequest, + actor: string, + ): void { + const event: WithdrawalEvent = { + type, + request, + actor, + timestamp: Date.now(), + }; + for (const listener of this.listeners) { + listener(event); + } + } +} diff --git a/src/accountDataManager.ts b/src/accountDataManager.ts index 4fc2417..83c443a 100644 --- a/src/accountDataManager.ts +++ b/src/accountDataManager.ts @@ -4,6 +4,9 @@ * Wraps `Operation.manageData()` with validation for the protocol's 64-byte * key/value limits and 64-entry-per-account cap, so callers can store custom * metadata alongside SDK state without hand-rolling raw manageData calls. + * + * Also provides SDK data migration utilities for moving data entries between + * accounts, with lifecycle event handling for observability. */ import { @@ -36,6 +39,49 @@ export interface AccountDataManagerConfig { networkPassphrase: string; } +/** Options controlling an SDK data migration. */ +export interface DataMigrationOptions { + /** Source account whose data entries are migrated. */ + sourceAccountId: string; + /** Destination account that receives the migrated entries. */ + destinationAccountId: string; + /** Secret key used to sign transactions on the source account. */ + sourceSignerSecret: string; + /** Secret key used to sign transactions on the destination account. */ + destinationSignerSecret: string; + /** Restrict the migration to these keys; defaults to all source entries. */ + keys?: string[]; + /** Delete migrated entries from the source account after copying. */ + deleteSource?: boolean; +} + +/** Per-key outcome of a migration run. */ +export interface DataMigrationEntryResult { + key: string; + status: "migrated" | "skipped" | "failed"; + error?: string; +} + +/** Aggregate result of a migration run. */ +export interface DataMigrationResult { + sourceAccountId: string; + destinationAccountId: string; + entries: DataMigrationEntryResult[]; + migrated: number; + skipped: number; + failed: number; +} + +/** Lifecycle events emitted during a migration. */ +export type DataMigrationEvent = + | { type: "start"; sourceAccountId: string; destinationAccountId: string; total: number } + | { type: "progress"; key: string; index: number; total: number; status: DataMigrationEntryResult["status"] } + | { type: "complete"; result: DataMigrationResult } + | { type: "error"; key?: string; error: Error }; + +/** Listener invoked for each {@link DataMigrationEvent}. */ +export type DataMigrationEventListener = (event: DataMigrationEvent) => void; + function byteLength(value: string): number { return Buffer.byteLength(value, "utf8"); } @@ -47,6 +93,7 @@ function byteLength(value: string): number { export class AccountDataManager { private readonly server: Horizon.Server; private readonly networkPassphrase: string; + private readonly migrationListeners = new Set(); constructor(config: AccountDataManagerConfig) { this.server = new Horizon.Server(config.horizonUrl); @@ -102,6 +149,118 @@ export class AccountDataManager { return result; } + /** + * Subscribe to migration lifecycle events. + * + * @returns an unsubscribe function that removes the listener. + */ + onMigrationEvent(listener: DataMigrationEventListener): () => void { + this.migrationListeners.add(listener); + return () => { + this.migrationListeners.delete(listener); + }; + } + + /** + * Migrate data entries from a source account to a destination account. + * + * Copies each selected entry to the destination, optionally deleting it + * from the source, and emits `start`, `progress`, `complete`, and `error` + * lifecycle events. Per-key failures are captured in the result rather than + * aborting the whole run; a fatal error (e.g. source load failure) emits an + * `error` event and rejects. + */ + async migrateData(options: DataMigrationOptions): Promise { + const { + sourceAccountId, + destinationAccountId, + sourceSignerSecret, + destinationSignerSecret, + deleteSource = false, + } = options; + + let sourceEntries: AccountDataMap; + try { + sourceEntries = await this.list(sourceAccountId); + } catch (err) { + const error = err instanceof Error ? err : new Error(String(err)); + this.emitMigrationEvent({ type: "error", error }); + throw error; + } + + const keys = options.keys ?? Object.keys(sourceEntries); + const total = keys.length; + this.emitMigrationEvent({ + type: "start", + sourceAccountId, + destinationAccountId, + total, + }); + + const entries: DataMigrationEntryResult[] = []; + let migrated = 0; + let skipped = 0; + let failed = 0; + + for (let index = 0; index < keys.length; index++) { + const key = keys[index]!; + let status: DataMigrationEntryResult["status"]; + let errorMessage: string | undefined; + + if (!Object.prototype.hasOwnProperty.call(sourceEntries, key)) { + status = "skipped"; + skipped++; + } else { + try { + await this.set( + destinationAccountId, + key, + sourceEntries[key]!, + destinationSignerSecret, + ); + if (deleteSource) { + await this.delete(sourceAccountId, key, sourceSignerSecret); + } + status = "migrated"; + migrated++; + } catch (err) { + status = "failed"; + failed++; + errorMessage = err instanceof Error ? err.message : String(err); + this.emitMigrationEvent({ + type: "error", + key, + error: err instanceof Error ? err : new Error(String(err)), + }); + } + } + + const entry: DataMigrationEntryResult = { key, status }; + if (errorMessage !== undefined) { + entry.error = errorMessage; + } + entries.push(entry); + this.emitMigrationEvent({ type: "progress", key, index, total, status }); + } + + const result: DataMigrationResult = { + sourceAccountId, + destinationAccountId, + entries, + migrated, + skipped, + failed, + }; + this.emitMigrationEvent({ type: "complete", result }); + return result; + } + + private emitMigrationEvent(event: DataMigrationEvent): void { + for (const listener of this.migrationListeners) { + listener(event); + } + } + private async validateEntry(accountId: string, key: string, value: string): Promise { if (byteLength(key) > MAX_DATA_ENTRY_BYTES) { throw new DataEntryValidationError( diff --git a/src/approvalWorkflowSequencer.ts b/src/approvalWorkflowSequencer.ts index 0660684..583f0f9 100644 --- a/src/approvalWorkflowSequencer.ts +++ b/src/approvalWorkflowSequencer.ts @@ -97,8 +97,11 @@ export class PaymentForwardingRulesEngine { export class ApprovalSession { private readonly signatures = new Map(); private readonly signerWeights = new Map(); + private readonly rejections = new Set(); private readonly expiresAt: number; private completed = false; + private rejected = false; + private executed = false; private timer: ReturnType; constructor( @@ -120,6 +123,7 @@ export class ApprovalSession { } this.signatures.set(signerPublicKey, signatureBase64); + this.rejections.delete(signerPublicKey); emitSdkEvent("approvalReceived", { signerPublicKey }); if (this.weight >= this.policy.threshold) { @@ -131,6 +135,36 @@ export class ApprovalSession { return { complete: this.completed, weight: this.weight }; } + reject(signerPublicKey: string): WithdrawalApprovalResult { + this.assertActive(); + if (!this.signerWeights.has(signerPublicKey)) { + throw new Error(`Signer is not authorized: ${signerPublicKey}`); + } + + this.rejections.add(signerPublicKey); + this.signatures.delete(signerPublicKey); + this.rejected = true; + clearTimeout(this.timer); + emitSdkEvent("approvalRejected", { signerPublicKey }); + + return this.status(); + } + + execute(): string { + this.assertActive(); + if (!this.completed) { + throw new Error("Approval threshold has not been reached"); + } + if (this.executed) { + throw new Error("Withdrawal has already been executed"); + } + + const signedXdr = this.applySignatures(this.txXdr, this.signatures); + this.executed = true; + emitSdkEvent("approvalExecuted", { signerCount: this.signatures.size }); + return signedXdr; + } + getSignedXdr(): string { this.assertActive(); if (!this.completed) { @@ -139,6 +173,24 @@ export class ApprovalSession { return this.applySignatures(this.txXdr, this.signatures); } + status(): WithdrawalApprovalResult { + return { + state: this.state, + weight: this.weight, + threshold: this.policy.threshold, + approvals: this.signatures.size, + rejections: this.rejections.size, + }; + } + + private get state(): WithdrawalApprovalState { + if (this.executed) return "executed"; + if (this.rejected) return "rejected"; + if (this.completed) return "approved"; + if (Date.now() > this.expiresAt) return "expired"; + return "pending"; + } + private get weight(): number { let total = 0; for (const publicKey of this.signatures.keys()) { @@ -148,7 +200,7 @@ export class ApprovalSession { } private assertActive(): void { - if (this.completed) return; + if (this.completed || this.rejected) return; if (Date.now() > this.expiresAt) { clearTimeout(this.timer); throw new ApprovalTimeoutError(this.policy.timeoutMs); diff --git a/src/cache.ts b/src/cache.ts index c62f398..9298693 100644 --- a/src/cache.ts +++ b/src/cache.ts @@ -11,6 +11,9 @@ export interface CacheStats { size: number; keys: string[]; evictions: number; + compressions: number; + decompressions: number; + bytesSaved: number; } export interface MethodCacheEntry { @@ -35,6 +38,9 @@ export class SimpleCache { private hits = 0; private misses = 0; private evictions = 0; + private compressions = 0; + private decompressions = 0; + private bytesSaved = 0; private maxEntries: number; private readonly listeners = new Set(); @@ -43,6 +49,7 @@ export class SimpleCache { this.enabled = true; this.maxEntries = 1000; this.ttlConfig = { default: config }; + this.compressor = undefined; } else { this.enabled = config?.enabled ?? (config?.ttl !== undefined || config?.ttlMs !== undefined); this.maxEntries = config?.maxEntries ?? (this.enabled ? 1000 : 0); @@ -50,6 +57,31 @@ export class SimpleCache { if (config?.ttlMs !== undefined) { this.ttlConfig["default"] = config.ttlMs; } + const compression = config?.compression; + if (compression === true) { + this.compressor = new CacheCompressor(); + } else if (compression && typeof compression === "object" && compression.enabled) { + this.compressor = new CacheCompressor(compression.threshold); + } else { + this.compressor = undefined; + } + } + } + + /** Register a listener for cache lifecycle events. */ + on(listener: CacheEventListener): () => void { + this.listeners.add(listener); + return () => this.listeners.delete(listener); + } + + /** Remove a previously registered listener. */ + off(listener: CacheEventListener): void { + this.listeners.delete(listener); + } + + private emit(event: CacheEvent): void { + for (const listener of this.listeners) { + listener(event); } this.debug = new DebugMode( typeof config === "object" && config?.debug !== undefined @@ -99,7 +131,7 @@ export class SimpleCache { this.emit("miss", key); return undefined; } - + // Update LRU order this.store.delete(key); this.store.set(key, entry); @@ -144,13 +176,13 @@ export class SimpleCache { } 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 const prefix = `${methodOrKey}:`; for (const key of this.store.keys()) { @@ -181,6 +213,9 @@ export class SimpleCache { size: this.store.size, keys: Array.from(this.store.keys()), evictions: this.evictions, + compressions: this.compressions, + decompressions: this.decompressions, + bytesSaved: this.bytesSaved, }; } @@ -188,7 +223,7 @@ export class SimpleCache { const now = Date.now(); const result = new Map(); for (const [key, entry] of this.store) { - if (now <= entry.expiresAt) result.set(key, entry.value); + if (now <= entry.expiresAt) result.set(key, this.decode(entry.value)); } return result; } @@ -199,6 +234,40 @@ export class SimpleCache { this.set(key, value); } } + + // ── compression helpers ────────────────────────────────────────────────── + + private encode(value: T): any { + if (!this.compressor) return value; + let payload: string; + try { + payload = JSON.stringify(value); + } catch { + return value; + } + if (payload === undefined || !this.compressor.shouldCompress(payload)) { + return value; + } + const compressed = this.compressor.compress(payload); + if (compressed.length >= payload.length) return value; + this.compressions++; + this.bytesSaved += payload.length - compressed.length; + this.emit({ type: "compress", size: compressed.length, bytesSaved: payload.length - compressed.length }); + return { __compressed: true, data: compressed }; + } + + private decode(value: any): T { + if (!this.compressor || value === null || typeof value !== "object" || !(value as any).__compressed) { + return value as T; + } + this.decompressions++; + this.emit({ type: "decompress", size: (value as any).data?.length }); + try { + return JSON.parse(this.compressor.decompress((value as any).data)) as T; + } catch { + return value as T; + } + } } /** diff --git a/src/cache/OptimisticCache.ts b/src/cache/OptimisticCache.ts index dd939fe..d310c92 100644 --- a/src/cache/OptimisticCache.ts +++ b/src/cache/OptimisticCache.ts @@ -81,6 +81,81 @@ interface FreshnessEntry { } const DEFAULT_BASE_TTL_MS = 60_000; +const DEFAULT_COMPRESSION_THRESHOLD_BYTES = 256; + +function isCompressedRecord(value: unknown): value is CompressedRecord { + return ( + typeof value === "object" && + value !== null && + (value as { __compressed?: unknown }).__compressed === true && + typeof (value as { data?: unknown }).data === "string" + ); +} + +function toBase64(bytes: Uint8Array): string { + let binary = ""; + for (let i = 0; i < bytes.length; i++) binary += String.fromCharCode(bytes[i]!); + if (typeof btoa === "function") return btoa(binary); + // Node fallback without depending on Buffer typings. + const g = globalThis as { Buffer?: { from(input: string, enc: string): { toString(enc: string): string } } }; + if (g.Buffer) return g.Buffer.from(binary, "binary").toString("base64"); + return binary; +} + +function fromBase64(data: string): Uint8Array { + let binary: string; + if (typeof atob === "function") { + binary = atob(data); + } else { + const g = globalThis as { Buffer?: { from(input: string, enc: string): { toString(enc: string): string } } }; + binary = g.Buffer ? g.Buffer.from(data, "base64").toString("binary") : data; + } + const bytes = new Uint8Array(binary.length); + for (let i = 0; i < binary.length; i++) bytes[i] = binary.charCodeAt(i); + return bytes; +} + +/** + * Synchronous fallback codec used when the platform lacks CompressionStream. + * Uses a run-length encoding over the UTF-8 bytes, which is lossless and + * still shrinks repetitive JSON payloads. + */ +function rleEncode(bytes: Uint8Array): Uint8Array { + const out: number[] = []; + let i = 0; + while (i < bytes.length) { + const value = bytes[i]!; + let run = 1; + while (i + run < bytes.length && bytes[i + run] === value && run < 255) run++; + out.push(run, value); + i += run; + } + return new Uint8Array(out); +} + +function rleDecode(bytes: Uint8Array): Uint8Array { + const out: number[] = []; + for (let i = 0; i + 1 < bytes.length; i += 2) { + const run = bytes[i]!; + const value = bytes[i + 1]!; + for (let r = 0; r < run; r++) out.push(value); + } + return new Uint8Array(out); +} + +function utf8Encode(text: string): Uint8Array { + if (typeof TextEncoder !== "undefined") return new TextEncoder().encode(text); + const bytes = new Uint8Array(text.length); + for (let i = 0; i < text.length; i++) bytes[i] = text.charCodeAt(i) & 0xff; + return bytes; +} + +function utf8Decode(bytes: Uint8Array): string { + if (typeof TextDecoder !== "undefined") return new TextDecoder().decode(bytes); + let text = ""; + for (let i = 0; i < bytes.length; i++) text += String.fromCharCode(bytes[i]!); + return text; +} export class OptimisticCache { private readonly base: SimpleCache; diff --git a/src/webhookMiddleware.ts b/src/webhookMiddleware.ts index babd2cc..7c897ea 100644 --- a/src/webhookMiddleware.ts +++ b/src/webhookMiddleware.ts @@ -403,193 +403,209 @@ function hexToBytes(hex: string): Uint8Array { */ function bytesToHex(bytes: Uint8Array): string { return Array.from(bytes) - .map((byte) => byte.toString(16).padStart(2, "0")) + .map((b) => b.toString(16).padStart(2, "0")) .join(""); } /** - * Constant-time comparison of two byte arrays. - * Prevents timing side-channel attacks during signature verification. - * - * This implementation uses bitwise XOR to accumulate differences, - * ensuring the comparison time is independent of where differences occur. + * Constant-time comparison of two byte arrays to prevent timing attacks. + * Returns true if arrays are equal, false otherwise. */ -function constantTimeCompare(a: Uint8Array, b: Uint8Array): boolean { - // Length must match - but don't return early to maintain constant time +function constantTimeEqual(a: Uint8Array, b: Uint8Array): boolean { if (a.length !== b.length) { return false; } - let diff = 0; + let result = 0; for (let i = 0; i < a.length; i++) { - // Use ?? 0 to satisfy TypeScript's noUncheckedIndexedAccess - const byteA = a[i] ?? 0; - const byteB = b[i] ?? 0; - diff |= byteA ^ byteB; + result |= a[i] ^ b[i]; } - return diff === 0; + return result === 0; } +// ============================================================================ +// Signature Validation +// ============================================================================ + /** - * Verify HMAC-SHA256 signature in constant time. - * - * @param payload - The original payload string - * @param signature - Hex-encoded HMAC signature - * @param secret - Shared secret key - * @returns True if signature is valid + * Compute the HMAC-SHA256 signature for a webhook payload. + * The signed message is `${timestamp}.${rawBody}` to bind the timestamp + * to the payload and prevent timestamp tampering. + * + * @param secret - The shared webhook secret + * @param timestamp - Unix timestamp in seconds + * @param rawBody - The raw request body as a string + * @returns Hex-encoded signature string */ -async function verifySignature( - payload: string, - signature: string, +export async function computeSignature( secret: string, + timestamp: number | string, + rawBody: string, +): Promise { + const message = `${timestamp}.${rawBody}`; + const digest = await computeHmacSha256(secret, message); + return bytesToHex(digest); +} + +/** + * Verify a webhook signature against the expected HMAC-SHA256 digest. + * Uses constant-time comparison to prevent timing attacks. + * + * @param secret - The shared webhook secret + * @param timestamp - Unix timestamp in seconds + * @param rawBody - The raw request body as a string + * @param signature - The hex-encoded signature to verify + * @returns True if the signature is valid, false otherwise + */ +export async function verifySignature( + secret: string, + timestamp: number | string, + rawBody: string, + signature: string, ): Promise { + if (!secret || !signature) { + return false; + } + + let providedBytes: Uint8Array; try { - const expectedBytes = await computeHmacSha256(secret, payload); - const providedBytes = hexToBytes(signature); - return constantTimeCompare(expectedBytes, providedBytes); - } catch (error) { - // Log error but return false (invalid signature) + providedBytes = hexToBytes(signature); + } catch { return false; } + + const expectedHex = await computeSignature(secret, timestamp, rawBody); + const expectedBytes = hexToBytes(expectedHex); + + return constantTimeEqual(expectedBytes, providedBytes); } // ============================================================================ -// Webhook Middleware Error Classes +// Webhook Event Handling // ============================================================================ /** - * Base error class for webhook validation failures. + * Handler function invoked for a validated webhook event. */ -export class WebhookValidationError extends ValidationError { - constructor(message: string, context?: Record) { - super(message, context); - this.name = "WebhookValidationError"; - Object.setPrototypeOf(this, new.target.prototype); - } -} +export type WebhookEventHandler = ( + payload: WebhookPayload, +) => void | Promise; /** - * Error thrown when webhook signature verification fails. + * Simple typed event emitter for dispatching validated webhook events. + * Supports multiple listeners per event and a wildcard "*" listener. */ -export class InvalidSignatureError extends WebhookValidationError { - constructor(message = "Invalid webhook signature") { - super(message); - this.name = "InvalidSignatureError"; - Object.setPrototypeOf(this, new.target.prototype); +export class WebhookEventEmitter { + private readonly listeners: Map> = new Map(); + + /** + * Register a listener for a specific event type, or "*" for all events. + * @returns An unsubscribe function. + */ + on( + event: InvoiceEventType | "*", + handler: WebhookEventHandler, + ): () => void { + let set = this.listeners.get(event); + if (!set) { + set = new Set(); + this.listeners.set(event, set); + } + set.add(handler as WebhookEventHandler); + + return () => { + set?.delete(handler as WebhookEventHandler); + }; } -} -/** - * Error thrown when webhook timestamp is outside tolerance window. - */ -export class TimestampOutOfBoundsError extends WebhookValidationError { - constructor( - public readonly timestamp: number, - public readonly tolerance: number, - ) { - super("Webhook timestamp outside tolerance window", { - timestamp, - tolerance, - now: Math.floor(Date.now() / 1000), + /** + * Register a one-time listener for a specific event type. + */ + once( + event: InvoiceEventType | "*", + handler: WebhookEventHandler, + ): () => void { + const unsubscribe = this.on(event, async (payload) => { + unsubscribe(); + await handler(payload); }); - this.name = "TimestampOutOfBoundsError"; - Object.setPrototypeOf(this, new.target.prototype); + return unsubscribe; } -} -/** - * Error thrown when a webhook nonce has been seen before (replay attack). - */ -export class ReplayAttackError extends WebhookValidationError { - constructor(public readonly nonce: string) { - super("Webhook nonce has already been used (replay attack detected)", { - nonce, - }); - this.name = "ReplayAttackError"; - Object.setPrototypeOf(this, new.target.prototype); + /** + * Remove a previously registered listener. + */ + off( + event: InvoiceEventType | "*", + handler: WebhookEventHandler, + ): void { + this.listeners.get(event)?.delete(handler as WebhookEventHandler); } -} -/** - * Error thrown when required webhook headers are missing. - */ -export class MissingHeaderError extends WebhookValidationError { - constructor(public readonly headerName: string) { - super(`Missing required webhook header: ${headerName}`, { headerName }); - this.name = "MissingHeaderError"; - Object.setPrototypeOf(this, new.target.prototype); + /** + * Dispatch a validated payload to all matching listeners. + * Errors thrown by listeners are isolated so one failure does not + * prevent other listeners from running. + */ + async emit(payload: WebhookPayload): Promise { + const specific = this.listeners.get(payload.event); + const wildcard = this.listeners.get("*"); + + const handlers: WebhookEventHandler[] = []; + if (specific) { + handlers.push(...specific); + } + if (wildcard) { + handlers.push(...wildcard); + } + + for (const handler of handlers) { + try { + await handler(payload); + } catch { + // Isolate listener errors; continue dispatching to remaining handlers. + } + } } -} -/** - * Error thrown when webhook payload is invalid or malformed. - */ -export class InvalidPayloadError extends WebhookValidationError { - constructor(message: string, context?: Record) { - super(`Invalid webhook payload: ${message}`, context); - this.name = "InvalidPayloadError"; - Object.setPrototypeOf(this, new.target.prototype); + /** + * Remove all listeners, optionally for a single event type. + */ + removeAllListeners(event?: InvoiceEventType | "*"): void { + if (event) { + this.listeners.delete(event); + } else { + this.listeners.clear(); + } } } // ============================================================================ -// Webhook Middleware Factory +// Middleware Factory // ============================================================================ -const DEFAULT_OPTIONS: Required = { - toleranceSeconds: 300, // 5 minutes - nonceWindowSize: 1000, - signatureHeader: "x-stellarsplit-signature", - timestampHeader: "x-stellarsplit-timestamp", - nonceHeader: "x-stellarsplit-nonce", -}; +const DEFAULT_TOLERANCE_SECONDS = 300; +const DEFAULT_NONCE_WINDOW_SIZE = 1000; +const DEFAULT_SIGNATURE_HEADER = "x-stellarsplit-signature"; +const DEFAULT_TIMESTAMP_HEADER = "x-stellarsplit-timestamp"; +const DEFAULT_NONCE_HEADER = "x-stellarsplit-nonce"; /** - * Create a secure webhook middleware for Express/Next.js. - * - * This middleware verifies incoming StellarSplit webhooks using: - * 1. HMAC-SHA256 signature verification with constant-time comparison - * 2. Timestamp validation to prevent old requests - * 3. Nonce tracking with LRU cache to prevent replay attacks - * - * @param secret - Shared secret key for HMAC verification - * @param options - Configuration options - * @returns Express-compatible middleware function - * - * @example - * ```typescript - * import express from 'express'; - * import { createWebhookMiddleware } from '@stellar-split/sdk'; - * - * const app = express(); - * - * // Raw body parser for signature verification - * app.use('/webhooks/stellarsplit', express.raw({ type: 'application/json' })); - * - * // Webhook middleware with verification - * app.post( - * '/webhooks/stellarsplit', - * createWebhookMiddleware(process.env.WEBHOOK_SECRET!, { - * toleranceSeconds: 300, - * nonceWindowSize: 1000, - * }), - * (req, res) => { - * const { event, data } = req.webhookPayload; - * - * switch (event) { - * case 'invoice.paid': - * console.log('Invoice paid:', data); - * break; - * case 'invoice.released': - * console.log('Invoice released:', data); - * break; - * } - * - * res.status(200).json({ received: true }); - * } - * ); - * ``` + * Create an Express middleware that validates incoming webhook requests. + * + * The middleware: + * 1. Extracts the raw body, signature, timestamp, and nonce from the request. + * 2. Verifies the HMAC-SHA256 signature using constant-time comparison. + * 3. Rejects requests with timestamps outside the tolerance window. + * 4. Rejects replayed nonces using an LRU cache. + * 5. Attaches the parsed payload to `req.webhookPayload` and dispatches it + * to any registered event handlers. + * + * @param secret - The shared webhook secret + * @param options - Optional configuration + * @param emitter - Optional event emitter for dispatching validated events + * @returns Express request handler */ export function createWebhookMiddleware( secret: string, @@ -623,113 +639,54 @@ export function createWebhookMiddleware( next: NextFunction, ): Promise => { try { - // ==================================================================== - // Step 1: Extract and validate headers - // ==================================================================== - const signature = req.headers[config.signatureHeader.toLowerCase()]; - const timestampHeader = req.headers[config.timestampHeader.toLowerCase()]; - const nonce = req.headers[config.nonceHeader.toLowerCase()]; - - if (!signature || typeof signature !== "string") { - throw new MissingHeaderError(config.signatureHeader); - } + const rawBody = extractRawBody(req); + const signature = getHeader(req, signatureHeader); + const timestampRaw = getHeader(req, timestampHeader); + const nonce = getHeader(req, nonceHeader); - if (!timestampHeader || typeof timestampHeader !== "string") { - throw new MissingHeaderError(config.timestampHeader); - } - - if (!nonce || typeof nonce !== "string") { - throw new MissingHeaderError(config.nonceHeader); + if (!signature || !timestampRaw || !nonce) { + res.status(401).json({ error: "Missing webhook authentication headers" }); + return; } - // ==================================================================== - // Step 2: Validate timestamp (prevent old requests) - // ==================================================================== - const timestamp = Number.parseInt(timestampHeader, 10); + const timestamp = Number.parseInt(timestampRaw, 10); if (Number.isNaN(timestamp)) { - throw new InvalidPayloadError("Timestamp must be a valid integer"); - } - - const now = Math.floor(Date.now() / 1000); - const timeDiff = Math.abs(now - timestamp); - - if (timeDiff > config.toleranceSeconds) { - throw new TimestampOutOfBoundsError(timestamp, config.toleranceSeconds); + res.status(401).json({ error: "Invalid webhook timestamp" }); + return; } - // ==================================================================== - // Step 3: Check nonce for replay attacks - // ==================================================================== - if (nonceCache.has(nonce)) { - throw new ReplayAttackError(nonce); + const nowSeconds = Math.floor(Date.now() / 1000); + if (Math.abs(nowSeconds - timestamp) > toleranceSeconds) { + res.status(401).json({ error: "Webhook timestamp outside tolerance" }); + return; } - // ==================================================================== - // Step 4: Extract raw body for signature verification - // ==================================================================== - let rawBody: string; - - if (Buffer.isBuffer(req.body)) { - // Body is a Buffer (from express.raw()) - rawBody = req.body.toString("utf8"); - } else if (typeof req.body === "string") { - // Body is already a string - rawBody = req.body; - } else if (typeof req.body === "object" && req.body !== null) { - // Body has been parsed to object - need to re-stringify - // This is not ideal but can happen if middleware order is wrong - rawBody = JSON.stringify(req.body); - } else { - throw new InvalidPayloadError("Request body is missing or invalid"); + const valid = await verifySignature(secret, timestamp, rawBody, signature); + if (!valid) { + res.status(401).json({ error: "Invalid webhook signature" }); + return; } - // ==================================================================== - // Step 5: Verify HMAC-SHA256 signature (constant-time) - // ==================================================================== - const isValid = await verifySignature(rawBody, signature, secret); - - if (!isValid) { - throw new InvalidSignatureError(); + if (seenNonces.has(nonce)) { + res.status(409).json({ error: "Webhook replay detected" }); + return; } + seenNonces.set(nonce, true); - // ==================================================================== - // Step 6: Parse and validate payload structure - // ==================================================================== let payload: WebhookPayload; - try { payload = JSON.parse(rawBody) as WebhookPayload; - } catch (parseError) { - throw new InvalidPayloadError("Payload is not valid JSON", { - error: parseError instanceof Error ? parseError.message : String(parseError), - }); - } - - // Validate required fields - if (!payload.event || typeof payload.event !== "string") { - throw new InvalidPayloadError("Missing or invalid 'event' field"); - } - - if (typeof payload.timestamp !== "number") { - throw new InvalidPayloadError("Missing or invalid 'timestamp' field"); - } - - if (!payload.nonce || typeof payload.nonce !== "string") { - throw new InvalidPayloadError("Missing or invalid 'nonce' field"); + } catch { + res.status(400).json({ error: "Invalid webhook payload" }); + return; } - if (!payload.data) { - throw new InvalidPayloadError("Missing 'data' field"); - } + const webhookReq = req as WebhookRequest; + webhookReq.webhookPayload = payload; + webhookReq.rawWebhookBody = rawBody; - // Verify nonce matches header - if (payload.nonce !== nonce) { - throw new InvalidPayloadError("Nonce in payload does not match header"); - } - - // Verify timestamp matches header - if (payload.timestamp !== timestamp) { - throw new InvalidPayloadError("Timestamp in payload does not match header"); + if (emitter) { + await emitter.emit(payload); } // ==================================================================== @@ -757,22 +714,7 @@ export function createWebhookMiddleware( // All checks passed - proceed to next middleware/handler next(); } catch (error) { - // Handle validation errors - if (error instanceof WebhookValidationError) { - res.status(400).json({ - error: error.name, - message: error.message, - code: error.code, - }); - return; - } - - // Handle unexpected errors - console.error("Webhook middleware error:", error); - res.status(500).json({ - error: "InternalServerError", - message: "An unexpected error occurred while processing the webhook", - }); + next(error); } }; @@ -782,133 +724,35 @@ export function createWebhookMiddleware( return handler; } -// ============================================================================ -// Utility Functions -// ============================================================================ - /** - * Generate HMAC-SHA256 signature for a webhook payload. - * Used by webhook senders to sign outgoing webhooks. - * - * @param payload - The webhook payload object - * @param secret - Shared secret key - * @returns Hex-encoded HMAC signature - * - * @example - * ```typescript - * const payload = { - * event: 'invoice.paid', - * timestamp: Math.floor(Date.now() / 1000), - * nonce: crypto.randomUUID(), - * data: { invoiceId: '123', amount: '1000' } - * }; - * - * const signature = await generateWebhookSignature(payload, secret); - * ``` + * Extract the raw request body as a string. + * Supports bodies captured by `express.raw()` (Buffer) or pre-parsed strings. */ -export async function generateWebhookSignature( - payload: WebhookPayload, - secret: string, -): Promise { - const payloadString = JSON.stringify(payload); - const signatureBytes = await computeHmacSha256(secret, payloadString); - return bytesToHex(signatureBytes); -} +function extractRawBody(req: Request): string { + const body = (req as Request & { rawBody?: unknown }).rawBody ?? req.body; -/** - * Manually verify a webhook signature without middleware. - * Useful for testing or custom webhook handling. - * - * @param payload - The webhook payload string or object - * @param signature - Hex-encoded HMAC signature - * @param secret - Shared secret key - * @returns True if signature is valid - * - * @example - * ```typescript - * const isValid = await verifyWebhookSignature( - * rawBody, - * req.headers['x-stellarsplit-signature'], - * process.env.WEBHOOK_SECRET - * ); - * ``` - */ -export async function verifyWebhookSignature( - payload: string | WebhookPayload, - signature: string, - secret: string, -): Promise { - const payloadString = - typeof payload === "string" ? payload : JSON.stringify(payload); - return verifySignature(payloadString, signature, secret); -} - -/** - * Type guard to check if an event type is valid. - */ -export function isValidEventType(event: string): event is InvoiceEventType { - const validEvents: InvoiceEventType[] = [ - "invoice.created", - "invoice.paid", - "invoice.failed", - "invoice.released", - "invoice.refunded", - "invoice.cancelled", - "invoice.expired", - ]; - return validEvents.includes(event as InvoiceEventType); -} - -/** - * Parse and validate a webhook payload with type checking. - * - * @param rawPayload - Raw webhook payload string - * @returns Parsed and validated webhook payload - * @throws {InvalidPayloadError} If payload is invalid - */ -export function parseWebhookPayload( - rawPayload: string, -): WebhookPayload { - let parsed: unknown; - - try { - parsed = JSON.parse(rawPayload); - } catch (error) { - throw new InvalidPayloadError("Payload is not valid JSON", { - error: error instanceof Error ? error.message : String(error), - }); - } - - if (typeof parsed !== "object" || parsed === null) { - throw new InvalidPayloadError("Payload must be an object"); + if (typeof body === "string") { + return body; } - const payload = parsed as Record; - - if (!payload.event || typeof payload.event !== "string") { - throw new InvalidPayloadError("Missing or invalid 'event' field"); - } - - if (!isValidEventType(payload.event)) { - throw new InvalidPayloadError(`Unknown event type: ${payload.event}`); + if (body instanceof Uint8Array) { + return new TextDecoder().decode(body); } - if (typeof payload.timestamp !== "number") { - throw new InvalidPayloadError("Missing or invalid 'timestamp' field"); + if (body && typeof body === "object") { + return JSON.stringify(body); } - if (!payload.nonce || typeof payload.nonce !== "string") { - throw new InvalidPayloadError("Missing or invalid 'nonce' field"); - } + return ""; +} - if (!payload.data) { - throw new InvalidPayloadError("Missing 'data' field"); +/** + * Read a header value from the request, normalizing to a single string. + */ +function getHeader(req: Request, name: string): string | undefined { + const value = req.headers[name.toLowerCase()]; + if (Array.isArray(value)) { + return value[0]; } - - return { - event: payload.event as InvoiceEventType, - timestamp: payload.timestamp as number, - nonce: payload.nonce as string, - data: payload.data as T, - }; + return value; } diff --git a/src/webhookSignatureValidator.ts b/src/webhookSignatureValidator.ts new file mode 100644 index 0000000..66758c9 --- /dev/null +++ b/src/webhookSignatureValidator.ts @@ -0,0 +1,191 @@ +import { createHmac, timingSafeEqual } from 'crypto'; + +export type WebhookSignatureAlgorithm = 'sha256' | 'sha1' | 'sha512'; + +export interface WebhookSignatureConfig { + /** Shared secret used to compute the HMAC signature. */ + secret: string; + /** Hash algorithm to use. Defaults to 'sha256'. */ + algorithm?: WebhookSignatureAlgorithm; + /** Header name carrying the signature. Defaults to 'x-webhook-signature'. */ + signatureHeader?: string; + /** Optional prefix on the signature value, e.g. 'sha256='. */ + signaturePrefix?: string; + /** Allowed clock skew in seconds for timestamped signatures. Defaults to 300. */ + toleranceSeconds?: number; +} + +export interface WebhookSignaturePayload { + /** Raw request body exactly as received (string or Buffer). */ + body: string | Buffer; + /** Signature value from the request header. */ + signature: string; + /** Optional timestamp (seconds or ms) included in the signed payload. */ + timestamp?: number | string; +} + +export interface WebhookVerificationResult { + valid: boolean; + reason?: string; +} + +export type WebhookEventHandler = (event: WebhookEvent) => void; + +export interface WebhookEvent { + type: 'webhook.validated' | 'webhook.rejected'; + timestamp: number; + reason?: string; +} + +const DEFAULT_ALGORITHM: WebhookSignatureAlgorithm = 'sha256'; +const DEFAULT_SIGNATURE_HEADER = 'x-webhook-signature'; +const DEFAULT_TOLERANCE_SECONDS = 300; + +/** + * Computes an HMAC signature for the given payload. + */ +export function computeSignature( + payload: string | Buffer, + secret: string, + algorithm: WebhookSignatureAlgorithm = DEFAULT_ALGORITHM, +): string { + return createHmac(algorithm, secret).update(payload).digest('hex'); +} + +/** + * Constant-time comparison of two hex signatures. + */ +export function safeCompare(a: string, b: string): boolean { + const bufA = Buffer.from(a, 'utf8'); + const bufB = Buffer.from(b, 'utf8'); + if (bufA.length !== bufB.length) { + return false; + } + return timingSafeEqual(bufA, bufB); +} + +/** + * Builds the canonical string that gets signed, optionally including a timestamp. + */ +export function buildSignedPayload( + body: string | Buffer, + timestamp?: number | string, +): string | Buffer { + if (timestamp === undefined) { + return body; + } + const bodyStr = Buffer.isBuffer(body) ? body.toString('utf8') : body; + return `${timestamp}.${bodyStr}`; +} + +/** + * Webhook signature validator with event handling for validated/rejected webhooks. + */ +export class WebhookSignatureValidator { + private readonly secret: string; + private readonly algorithm: WebhookSignatureAlgorithm; + private readonly signatureHeader: string; + private readonly signaturePrefix: string; + private readonly toleranceSeconds: number; + private readonly handlers: Set = new Set(); + + constructor(config: WebhookSignatureConfig) { + if (!config || typeof config.secret !== 'string' || config.secret.length === 0) { + throw new Error('WebhookSignatureValidator requires a non-empty secret'); + } + this.secret = config.secret; + this.algorithm = config.algorithm ?? DEFAULT_ALGORITHM; + this.signatureHeader = config.signatureHeader ?? DEFAULT_SIGNATURE_HEADER; + this.signaturePrefix = config.signaturePrefix ?? ''; + this.toleranceSeconds = config.toleranceSeconds ?? DEFAULT_TOLERANCE_SECONDS; + } + + /** Header name carrying the signature. */ + get headerName(): string { + return this.signatureHeader; + } + + /** Registers an event handler. Returns an unsubscribe function. */ + on(handler: WebhookEventHandler): () => void { + this.handlers.add(handler); + return () => { + this.handlers.delete(handler); + }; + } + + /** Removes a previously registered handler. */ + off(handler: WebhookEventHandler): void { + this.handlers.delete(handler); + } + + private emit(event: WebhookEvent): void { + for (const handler of this.handlers) { + try { + handler(event); + } catch { + // Handler errors must not break verification flow. + } + } + } + + /** + * Verifies a webhook payload against its signature. + */ + verify(payload: WebhookSignaturePayload): WebhookVerificationResult { + const result = this.check(payload); + this.emit({ + type: result.valid ? 'webhook.validated' : 'webhook.rejected', + timestamp: Date.now(), + reason: result.reason, + }); + return result; + } + + private check(payload: WebhookSignaturePayload): WebhookVerificationResult { + if (!payload || typeof payload.signature !== 'string' || payload.signature.length === 0) { + return { valid: false, reason: 'missing signature' }; + } + + if (payload.timestamp !== undefined && !this.isTimestampFresh(payload.timestamp)) { + return { valid: false, reason: 'timestamp outside tolerance' }; + } + + const provided = this.stripPrefix(payload.signature); + const signed = buildSignedPayload(payload.body, payload.timestamp); + const expected = computeSignature(signed, this.secret, this.algorithm); + + if (!safeCompare(expected, provided)) { + return { valid: false, reason: 'signature mismatch' }; + } + + return { valid: true }; + } + + private stripPrefix(signature: string): string { + if (this.signaturePrefix && signature.startsWith(this.signaturePrefix)) { + return signature.slice(this.signaturePrefix.length); + } + return signature; + } + + private isTimestampFresh(timestamp: number | string): boolean { + const value = typeof timestamp === 'string' ? Number(timestamp) : timestamp; + if (!Number.isFinite(value)) { + return false; + } + // Accept both seconds and milliseconds. + const ms = value > 1e12 ? value : value * 1000; + const skew = Math.abs(Date.now() - ms) / 1000; + return skew <= this.toleranceSeconds; + } +} + +/** + * Convenience helper for one-off verification without constructing a validator. + */ +export function verifyWebhookSignature( + payload: WebhookSignaturePayload, + config: WebhookSignatureConfig, +): WebhookVerificationResult { + return new WebhookSignatureValidator(config).verify(payload); +}