diff --git a/backend/src/services/crossChainAnomalyEngine.service.ts b/backend/src/services/crossChainAnomalyEngine.service.ts new file mode 100644 index 00000000..cb8f3359 --- /dev/null +++ b/backend/src/services/crossChainAnomalyEngine.service.ts @@ -0,0 +1,494 @@ +import * as StellarSdk from "@stellar/stellar-sdk"; +import { redis } from "../utils/redis.js"; +import { logger } from "../utils/logger.js"; +import { getMetricsService } from "./metrics.service.js"; +import { getCircuitBreakerService, PauseScope } from "./circuitBreaker.service.js"; +import type { FederatedEvent } from "./eventFederation/types.js"; +import { getDatabase } from "../database/connection.js"; + +export type AnomalyType = "double_spend" | "nonce_jump" | "reentrancy" | "threshold_breach"; + +export interface DetectedAnomaly { + id: string; + type: AnomalyType; + bridgeId: string; + chainId: string; + sequenceId?: number; + depositTxHash?: string; + details: Record; + timestamp: number; +} + +export interface FlashPauseResult { + triggered: boolean; + bridgeId: string; + anomalyCount: number; + reason: string; + timestamp: number; + contractPaused: boolean; +} + +export interface AnomalyEngineOptions { + windowSeconds?: number; + anomalyThreshold?: number; + nonceWindowSeconds?: number; + txHashWindowSeconds?: number; +} + +export class CrossChainAnomalyEngineService { + private readonly windowSeconds: number; + private readonly anomalyThreshold: number; + private readonly nonceWindowSeconds: number; + private readonly txHashWindowSeconds: number; + + // L1 In-Memory sliding window and state cache for sub-millisecond graph analysis + private readonly memoryStore = new Map(); + private readonly memoryAnomalies = new Map(); + private readonly memoryBreakers = new Map(); + + constructor(options: AnomalyEngineOptions = {}) { + this.windowSeconds = options.windowSeconds ?? 5; + this.anomalyThreshold = options.anomalyThreshold ?? 2; + this.nonceWindowSeconds = options.nonceWindowSeconds ?? 3600; + this.txHashWindowSeconds = options.txHashWindowSeconds ?? 3600; + } + + /** + * Main entry point for ingesting real-time federated stream events. + * Analyzes event for double-spend attempts, out-of-order sequence nonce jumps, and cross-chain re-entrancy. + */ + async processEvent(event: FederatedEvent): Promise { + const anomalies: DetectedAnomaly[] = []; + const now = Date.now(); + const bridgeId = this.extractBridgeId(event); + const chainId = event.chain || "unknown"; + + const depositTxHash = this.extractDepositTxHash(event); + const sequenceId = this.extractSequenceId(event); + + // 1. Double-Spend Anomaly Detection (Duplicate deposit tx hash across chains/relayers) + if (depositTxHash) { + const isDoubleSpend = await this.checkDoubleSpend(bridgeId, chainId, depositTxHash, now); + if (isDoubleSpend) { + const anomaly: DetectedAnomaly = { + id: `ds_${event.id}_${now}`, + type: "double_spend", + bridgeId, + chainId, + depositTxHash, + sequenceId, + details: { + message: `Double-spend deposit tx hash detected: ${depositTxHash}`, + eventId: event.id, + sourceId: event.sourceId, + }, + timestamp: now, + }; + anomalies.push(anomaly); + } + } + + // 2. Out-of-Order Sequence Nonce Jump Detection + if (sequenceId !== undefined && sequenceId !== null) { + const isNonceJump = await this.checkNonceJump(bridgeId, chainId, sequenceId, now); + if (isNonceJump) { + const anomaly: DetectedAnomaly = { + id: `nj_${event.id}_${now}`, + type: "nonce_jump", + bridgeId, + chainId, + sequenceId, + depositTxHash, + details: { + message: `Out-of-order sequence nonce jump detected: sequenceId ${sequenceId}`, + eventId: event.id, + sourceId: event.sourceId, + }, + timestamp: now, + }; + anomalies.push(anomaly); + } + } + + // 3. Cross-Chain Re-Entrancy Detection (rapid sub-second duplicate calls for same reference) + const isReentrancy = await this.checkReentrancy(bridgeId, chainId, event, now); + if (isReentrancy) { + const anomaly: DetectedAnomaly = { + id: `re_${event.id}_${now}`, + type: "reentrancy", + bridgeId, + chainId, + depositTxHash, + sequenceId, + details: { + message: `Cross-chain re-entrancy pattern detected within sub-second block window`, + eventId: event.id, + sourceId: event.sourceId, + }, + timestamp: now, + }; + anomalies.push(anomaly); + } + + // Record any detected anomalies and evaluate 5-second Flash-Pause threshold + if (anomalies.length > 0) { + for (const anomaly of anomalies) { + await this.recordAnomaly(anomaly); + } + + await this.evaluateFlashPauseThreshold(bridgeId, anomalies); + } + + return anomalies; + } + + /** + * Checks if deposit transaction hash was already processed or seen on another chain/relayer. + */ + private async checkDoubleSpend( + bridgeId: string, + chainId: string, + txHash: string, + now: number + ): Promise { + const key = `ccae:txhash:${bridgeId}:${txHash}`; + const memChain = this.memoryStore.get(key) as string | undefined; + let existingChain: string | null = memChain ?? null; + + if (!existingChain) { + try { + existingChain = await redis.get(key); + } catch { + existingChain = null; + } + } + + if (existingChain && existingChain !== chainId) { + return true; + } + + this.memoryStore.set(key, chainId); + try { + await redis.set(key, chainId, "EX", this.txHashWindowSeconds); + } catch { + // Redis optional L2 + } + + return false; + } + + /** + * Tracks sequence nonce per bridge/chain and flags jumps or duplicate/regressive nonces. + */ + private async checkNonceJump( + bridgeId: string, + chainId: string, + sequenceId: number, + now: number + ): Promise { + const key = `ccae:seq:${bridgeId}:${chainId}`; + const memSeq = this.memoryStore.get(key); + let lastSeqStr: string | null = memSeq !== undefined ? String(memSeq) : null; + + if (lastSeqStr === null) { + try { + lastSeqStr = await redis.get(key); + } catch { + lastSeqStr = null; + } + } + + let isJump = false; + if (lastSeqStr !== null && lastSeqStr !== undefined) { + const lastSeq = parseInt(lastSeqStr, 10); + if (!isNaN(lastSeq)) { + if (sequenceId > lastSeq + 1 || sequenceId <= lastSeq) { + isJump = true; + } + } + } + + this.memoryStore.set(key, sequenceId); + try { + await redis.set(key, String(sequenceId), "EX", this.nonceWindowSeconds); + } catch { + // Redis optional L2 + } + + return isJump; + } + + /** + * Detects rapid sub-second execution with identical key parameters (re-entrancy signature). + */ + private async checkReentrancy( + bridgeId: string, + chainId: string, + event: FederatedEvent, + now: number + ): Promise { + const refKey = event.sourceId || event.id; + const key = `ccae:reentrancy:${bridgeId}:${refKey}`; + + const memLastSeen = this.memoryStore.get(key) as number | undefined; + let lastSeen: number | null = memLastSeen ?? null; + + if (lastSeen === null) { + try { + const str = await redis.get(key); + if (str) lastSeen = parseInt(str, 10); + } catch { + lastSeen = null; + } + } + + let isReentrant = false; + if (lastSeen !== null && !isNaN(lastSeen) && now - lastSeen < 1000) { + isReentrant = true; + } + + this.memoryStore.set(key, now); + try { + await redis.set(key, String(now), "EX", 10); + } catch { + // Redis optional L2 + } + + return isReentrant; + } + + /** + * Records detected anomaly into L1 Memory + L2 Redis sliding window. + */ + private async recordAnomaly(anomaly: DetectedAnomaly): Promise { + const bridgeId = anomaly.bridgeId; + const list = this.memoryAnomalies.get(bridgeId) ?? []; + list.push(anomaly); + + const cutoff = anomaly.timestamp - (this.windowSeconds * 1000); + const filtered = list.filter((a) => a.timestamp >= cutoff); + this.memoryAnomalies.set(bridgeId, filtered); + + try { + const windowKey = `ccae:anomalies:${bridgeId}`; + await redis.zadd(windowKey, anomaly.timestamp, JSON.stringify(anomaly)); + await redis.zremrangebyscore(windowKey, "-inf", cutoff); + await redis.expire(windowKey, this.windowSeconds * 2); + } catch { + // Redis optional L2 + } + + logger.warn({ anomaly }, "Cross-chain anomaly recorded"); + } + + /** + * Evaluates total anomalies recorded within rolling 5-second window. + * If threshold is breached, triggers automated Flash-Pause. + */ + async evaluateFlashPauseThreshold( + bridgeId: string, + recentAnomalies: DetectedAnomaly[] + ): Promise { + const now = Date.now(); + const cutoff = now - (this.windowSeconds * 1000); + + const memList = (this.memoryAnomalies.get(bridgeId) ?? []).filter((a) => a.timestamp >= cutoff); + this.memoryAnomalies.set(bridgeId, memList); + + let anomalyCount = memList.length; + + try { + const windowKey = `ccae:anomalies:${bridgeId}`; + await redis.zremrangebyscore(windowKey, "-inf", cutoff); + const anomaliesInWindow = await redis.zrangebyscore(windowKey, cutoff, "+inf"); + if (Array.isArray(anomaliesInWindow) && anomaliesInWindow.length > anomalyCount) { + anomalyCount = anomaliesInWindow.length; + } + } catch { + // Redis optional L2 + } + + if (anomalyCount >= this.anomalyThreshold) { + const reason = `Automated Flash-Pause: ${anomalyCount} cross-chain anomalies detected within ${this.windowSeconds}s window`; + return this.triggerFlashPause(bridgeId, recentAnomalies, reason); + } + + return { + triggered: false, + bridgeId, + anomalyCount, + reason: "Below threshold", + timestamp: now, + contractPaused: false, + }; + } + + /** + * Triggers an automated Flash-Pause directly invoking Soroban pause_contract RPC and setting emergency breaker flags. + */ + async triggerFlashPause( + bridgeId: string, + anomalies: DetectedAnomaly[], + reason: string, + signer?: StellarSdk.Keypair + ): Promise { + const now = Date.now(); + const breakerKey = `ccae:breaker:${bridgeId}`; + + this.memoryBreakers.set(bridgeId, true); + + try { + await redis.set( + breakerKey, + JSON.stringify({ + active: true, + triggeredAt: now, + reason, + anomalyCount: anomalies.length, + }), + "EX", + 86400 + ); + } catch { + // Redis optional L2 + } + + let contractPaused = false; + const circuitBreaker = getCircuitBreakerService(); + + if (circuitBreaker) { + try { + const keypair = signer || StellarSdk.Keypair.random(); + await circuitBreaker.triggerPause(keypair, PauseScope.Bridge, bridgeId, reason); + contractPaused = true; + logger.info({ bridgeId, reason }, "Soroban contract pause_bridge invoked successfully"); + } catch (err) { + logger.error({ err, bridgeId }, "Failed invoking Soroban contract pause_bridge RPC"); + } + } + + try { + const metricsService = getMetricsService(); + metricsService.circuitBreakerTrips.inc({ + bridge_id: bridgeId, + reason: "flash_pause_anomaly", + }); + } catch (err) { + logger.warn({ err }, "Could not update metrics for flash pause"); + } + + try { + const db = getDatabase(); + const SYSTEM_RULE_ID = "00000000-0000-0000-0000-000000000000"; + await db("alert_events").insert({ + rule_id: SYSTEM_RULE_ID, + asset_code: bridgeId, + alert_type: "cross_chain_flash_pause", + priority: "critical", + triggered_value: anomalies.length, + threshold: this.anomalyThreshold, + metric: "cross_chain_anomaly_threshold", + webhook_delivered: false, + webhook_attempts: 0, + }); + } catch (err) { + logger.warn({ err }, "Could not persist flash pause alert event to DB"); + } + + logger.error({ bridgeId, reason, anomalyCount: anomalies.length }, "EMERGENCY FLASH-PAUSE TRIGGERED"); + + return { + triggered: true, + bridgeId, + anomalyCount: anomalies.length, + reason, + timestamp: now, + contractPaused, + }; + } + + /** + * Checks whether the emergency breaker is active for a given bridge. + */ + async isEmergencyBreakerActive(bridgeId: string): Promise { + if (this.memoryBreakers.get(bridgeId) === true) { + return true; + } + + const breakerKey = `ccae:breaker:${bridgeId}`; + try { + const data = await redis.get(breakerKey); + if (!data) return false; + const parsed = JSON.parse(data); + return Boolean(parsed.active); + } catch { + return false; + } + } + + /** + * Resets emergency breaker state flag. + */ + async resetEmergencyBreaker(bridgeId: string): Promise { + const breakerKey = `ccae:breaker:${bridgeId}`; + const windowKey = `ccae:anomalies:${bridgeId}`; + + this.memoryBreakers.delete(bridgeId); + this.memoryAnomalies.delete(bridgeId); + + try { + await redis.del(breakerKey); + await redis.del(windowKey); + } catch { + // Redis optional L2 + } + + logger.info({ bridgeId }, "Emergency breaker state reset"); + } + + /** + * Helper to extract bridge ID from event payload. + */ + private extractBridgeId(event: FederatedEvent): string { + const raw = event.raw ?? {}; + if (typeof raw.bridgeId === "string") return raw.bridgeId; + if (typeof raw.bridge_id === "string") return raw.bridge_id; + if (event.assetCode) return event.assetCode; + return "default-bridge"; + } + + /** + * Helper to extract deposit transaction hash from event payload. + */ + private extractDepositTxHash(event: FederatedEvent): string | undefined { + const raw = event.raw ?? {}; + if (typeof raw.depositTxHash === "string") return raw.depositTxHash; + if (typeof raw.deposit_tx_hash === "string") return raw.deposit_tx_hash; + if (typeof raw.txHash === "string") return raw.txHash; + if (event.type === "bridge_lock" || event.type === "bridge_release") { + return event.sourceId; + } + return undefined; + } + + /** + * Helper to extract sequence ID / nonce from event payload. + */ + private extractSequenceId(event: FederatedEvent): number | undefined { + const raw = event.raw ?? {}; + if (typeof raw.sequenceId === "number") return raw.sequenceId; + if (typeof raw.sequence === "number") return raw.sequence; + if (typeof raw.nonce === "number") return raw.nonce; + if (event.blockNumber > 0) return event.blockNumber; + return undefined; + } +} + +let _instance: CrossChainAnomalyEngineService | null = null; + +export function getCrossChainAnomalyEngineService(options?: AnomalyEngineOptions): CrossChainAnomalyEngineService { + if (!_instance || options) { + _instance = new CrossChainAnomalyEngineService(options); + } + return _instance; +} diff --git a/backend/src/services/eventFederation/EventFederationService.ts b/backend/src/services/eventFederation/EventFederationService.ts index 9f1d3453..98a82031 100644 --- a/backend/src/services/eventFederation/EventFederationService.ts +++ b/backend/src/services/eventFederation/EventFederationService.ts @@ -22,6 +22,7 @@ import type { ReplayRequest, } from "./types.js"; import { logger } from "../../utils/logger.js"; +import { getCrossChainAnomalyEngineService } from "../crossChainAnomalyEngine.service.js"; export const FEDERATION_EVENT = "event" as const; export const FEDERATION_HEALTH_EVENT = "health" as const; @@ -129,6 +130,11 @@ export class EventFederationService extends EventEmitter { this.totalProcessed++; this.replayBuffer.push(e); this.emit(FEDERATION_EVENT, e); + getCrossChainAnomalyEngineService() + .processEvent(e) + .catch((err) => { + logger.error({ err, eventId: e.id }, "CrossChainAnomalyEngine failed to process event"); + }); } } } diff --git a/backend/tests/integration/cross-chain/anomalyEngine.integration.test.ts b/backend/tests/integration/cross-chain/anomalyEngine.integration.test.ts new file mode 100644 index 00000000..447f0338 --- /dev/null +++ b/backend/tests/integration/cross-chain/anomalyEngine.integration.test.ts @@ -0,0 +1,91 @@ +import { describe, it, expect, beforeEach, afterEach, vi } from "vitest"; +import { EventFederationService } from "../../../src/services/eventFederation/EventFederationService.js"; +import { getCrossChainAnomalyEngineService } from "../../../src/services/crossChainAnomalyEngine.service.js"; +import type { FederatedEvent } from "../../../src/services/eventFederation/types.js"; +import * as circuitBreakerModule from "../../../src/services/circuitBreaker.service.js"; + +describe("Cross-Chain Anomaly Engine Integration Test", () => { + let federation: EventFederationService; + + beforeEach(async () => { + vi.clearAllMocks(); + federation = new EventFederationService(); + }); + + afterEach(async () => { + await federation.stop(); + }); + + it("ingests federated stream events and triggers automated Flash-Pause during double-mint exploit simulation", async () => { + const triggerPauseSpy = vi.fn().mockResolvedValue(undefined); + vi.spyOn(circuitBreakerModule, "getCircuitBreakerService").mockReturnValue({ + triggerPause: triggerPauseSpy, + } as any); + + const anomalyEngine = getCrossChainAnomalyEngineService({ + windowSeconds: 5, + anomalyThreshold: 2, + }); + + const bridgeId = "eth-stellar-usdc-bridge"; + const depositTxHash1 = "0xdeposit_exploit_hash_001"; + const depositTxHash2 = "0xdeposit_exploit_hash_002"; + + // Simulate cross-chain double mint exploit across 2 relayers + const event1: FederatedEvent = { + id: "evt_relayer_eth_1", + chain: "ethereum", + type: "bridge_lock", + blockNumber: 1000, + timestamp: new Date().toISOString(), + sourceId: "eth_tx_1000", + raw: { bridgeId, depositTxHash: depositTxHash1, sequenceId: 50 }, + }; + + const event2: FederatedEvent = { + id: "evt_relayer_stellar_1", + chain: "stellar", + type: "bridge_release", + blockNumber: 2000, + timestamp: new Date().toISOString(), + sourceId: "stl_tx_2000", + raw: { bridgeId, depositTxHash: depositTxHash1, sequenceId: 51 }, + }; + + const event3: FederatedEvent = { + id: "evt_relayer_eth_2", + chain: "ethereum", + type: "bridge_lock", + blockNumber: 1001, + timestamp: new Date().toISOString(), + sourceId: "eth_tx_1001", + raw: { bridgeId, depositTxHash: depositTxHash2, sequenceId: 52 }, + }; + + const event4: FederatedEvent = { + id: "evt_relayer_stellar_2", + chain: "stellar", + type: "bridge_release", + blockNumber: 2001, + timestamp: new Date().toISOString(), + sourceId: "stl_tx_2001", + raw: { bridgeId, depositTxHash: depositTxHash2, sequenceId: 53 }, + }; + + // Ingest events + await anomalyEngine.processEvent(event1); + await anomalyEngine.processEvent(event2); // Anomaly 1: Double-spend for depositTxHash1 + await anomalyEngine.processEvent(event3); + await anomalyEngine.processEvent(event4); // Anomaly 2: Double-spend for depositTxHash2 -> Threshold 2 breached! + + const isBreakerActive = await anomalyEngine.isEmergencyBreakerActive(bridgeId); + expect(isBreakerActive).toBe(true); + + expect(triggerPauseSpy).toHaveBeenCalledWith( + expect.anything(), + circuitBreakerModule.PauseScope.Bridge, + bridgeId, + expect.stringContaining("Automated Flash-Pause") + ); + }); +}); diff --git a/backend/tests/services/crossChainAnomalyEngine.service.test.ts b/backend/tests/services/crossChainAnomalyEngine.service.test.ts new file mode 100644 index 00000000..a74f0b1e --- /dev/null +++ b/backend/tests/services/crossChainAnomalyEngine.service.test.ts @@ -0,0 +1,223 @@ +import { describe, it, expect, beforeEach, vi } from "vitest"; +import { + CrossChainAnomalyEngineService, + getCrossChainAnomalyEngineService, +} from "../../src/services/crossChainAnomalyEngine.service.js"; +import type { FederatedEvent } from "../../src/services/eventFederation/types.js"; +import * as circuitBreakerModule from "../../src/services/circuitBreaker.service.js"; +import * as metricsModule from "../../src/services/metrics.service.js"; + +describe("CrossChainAnomalyEngineService", () => { + let engine: CrossChainAnomalyEngineService; + + beforeEach(() => { + vi.clearAllMocks(); + engine = new CrossChainAnomalyEngineService({ + windowSeconds: 5, + anomalyThreshold: 2, + nonceWindowSeconds: 3600, + txHashWindowSeconds: 3600, + }); + }); + + const createEvent = (overrides: Partial = {}): FederatedEvent => ({ + id: `evt_${Date.now()}_${Math.random()}`, + chain: "stellar", + type: "bridge_lock", + blockNumber: 100, + timestamp: new Date().toISOString(), + sourceId: `tx_${Math.random()}`, + raw: { + bridgeId: "usdc-stellar-eth", + depositTxHash: "0x123abc456def789", + sequenceId: 100, + }, + ...overrides, + }); + + describe("processEvent", () => { + it("processes legitimate sequential events cleanly without anomalies", async () => { + const event1 = createEvent({ + id: "evt_1", + chain: "stellar", + blockNumber: 100, + raw: { bridgeId: "usdc-bridge", depositTxHash: "0xhash1", sequenceId: 100 }, + }); + + const anomalies1 = await engine.processEvent(event1); + expect(anomalies1).toEqual([]); + + const event2 = createEvent({ + id: "evt_2", + chain: "stellar", + blockNumber: 101, + raw: { bridgeId: "usdc-bridge", depositTxHash: "0xhash2", sequenceId: 101 }, + }); + + const anomalies2 = await engine.processEvent(event2); + expect(anomalies2).toEqual([]); + }); + + it("detects double-spend attempt across chains/relayers", async () => { + const depositHash = "0xdouble_spend_hash_999"; + + const eventEthereum = createEvent({ + id: "evt_eth_1", + chain: "ethereum", + raw: { bridgeId: "usdc-bridge", depositTxHash: depositHash, sequenceId: 200 }, + }); + + await engine.processEvent(eventEthereum); + + const eventStellarDouble = createEvent({ + id: "evt_stl_1", + chain: "stellar", + raw: { bridgeId: "usdc-bridge", depositTxHash: depositHash, sequenceId: 201 }, + }); + + const anomalies = await engine.processEvent(eventStellarDouble); + + expect(anomalies.length).toBeGreaterThanOrEqual(1); + const dsAnomaly = anomalies.find((a) => a.type === "double_spend"); + expect(dsAnomaly).toBeDefined(); + expect(dsAnomaly?.depositTxHash).toBe(depositHash); + expect(dsAnomaly?.bridgeId).toBe("usdc-bridge"); + }); + + it("detects out-of-order sequence nonce jumps", async () => { + const event1 = createEvent({ + id: "evt_seq_1", + chain: "stellar", + raw: { bridgeId: "usdc-bridge", sequenceId: 50 }, + }); + + await engine.processEvent(event1); + + // Sequence jumps from 50 to 55 (jump > 1) + const eventJump = createEvent({ + id: "evt_seq_jump", + chain: "stellar", + raw: { bridgeId: "usdc-bridge", sequenceId: 55 }, + }); + + const anomalies = await engine.processEvent(eventJump); + + expect(anomalies.length).toBeGreaterThanOrEqual(1); + const jumpAnomaly = anomalies.find((a) => a.type === "nonce_jump"); + expect(jumpAnomaly).toBeDefined(); + expect(jumpAnomaly?.sequenceId).toBe(55); + }); + + it("detects cross-chain re-entrancy attack pattern", async () => { + const sharedSourceId = "reentrancy_source_tx_77"; + + const event1 = createEvent({ + id: "evt_re_1", + chain: "ethereum", + sourceId: sharedSourceId, + raw: { bridgeId: "reentrancy-bridge", depositTxHash: "0xre1" }, + }); + + await engine.processEvent(event1); + + // Rapid sub-second duplicate call with same sourceId + const eventReentrant = createEvent({ + id: "evt_re_2", + chain: "ethereum", + sourceId: sharedSourceId, + raw: { bridgeId: "reentrancy-bridge", depositTxHash: "0xre2" }, + }); + + const anomalies = await engine.processEvent(eventReentrant); + + expect(anomalies.length).toBeGreaterThanOrEqual(1); + const reAnomaly = anomalies.find((a) => a.type === "reentrancy"); + expect(reAnomaly).toBeDefined(); + }); + }); + + describe("Automated Flash-Pause", () => { + it("triggers automated Flash-Pause when anomaly threshold is breached within 5s window", async () => { + const triggerPauseSpy = vi.fn().mockResolvedValue(undefined); + vi.spyOn(circuitBreakerModule, "getCircuitBreakerService").mockReturnValue({ + triggerPause: triggerPauseSpy, + } as any); + + vi.spyOn(metricsModule, "getMetricsService").mockReturnValue({ + circuitBreakerTrips: { inc: vi.fn() }, + } as any); + + const bridgeId = "flash-pause-bridge"; + const depositHash1 = "0xhash_fp_1"; + const depositHash2 = "0xhash_fp_2"; + + // Trigger anomaly 1: double spend + await engine.processEvent( + createEvent({ + id: "evt_fp_1", + chain: "ethereum", + raw: { bridgeId, depositTxHash: depositHash1, sequenceId: 10 }, + }) + ); + + await engine.processEvent( + createEvent({ + id: "evt_fp_2", + chain: "stellar", + raw: { bridgeId, depositTxHash: depositHash1, sequenceId: 11 }, + }) + ); + + // Trigger anomaly 2: double spend for hash2 + await engine.processEvent( + createEvent({ + id: "evt_fp_3", + chain: "ethereum", + raw: { bridgeId, depositTxHash: depositHash2, sequenceId: 12 }, + }) + ); + + await engine.processEvent( + createEvent({ + id: "evt_fp_4", + chain: "stellar", + raw: { bridgeId, depositTxHash: depositHash2, sequenceId: 13 }, + }) + ); + + const isActive = await engine.isEmergencyBreakerActive(bridgeId); + expect(isActive).toBe(true); + + // Verify triggerPause on CircuitBreakerService was called + expect(triggerPauseSpy).toHaveBeenCalledWith( + expect.anything(), + circuitBreakerModule.PauseScope.Bridge, + bridgeId, + expect.stringContaining("Automated Flash-Pause") + ); + }); + + it("allows manual emergency breaker reset", async () => { + const bridgeId = "reset-bridge"; + + // Manually trigger flash pause + await engine.triggerFlashPause(bridgeId, [], "Test pause"); + + let isActive = await engine.isEmergencyBreakerActive(bridgeId); + expect(isActive).toBe(true); + + await engine.resetEmergencyBreaker(bridgeId); + + isActive = await engine.isEmergencyBreakerActive(bridgeId); + expect(isActive).toBe(false); + }); + }); + + describe("Singleton Factory", () => { + it("returns singleton instance of CrossChainAnomalyEngineService", () => { + const instance1 = getCrossChainAnomalyEngineService(); + const instance2 = getCrossChainAnomalyEngineService(); + expect(instance1).toBe(instance2); + }); + }); +});