Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 72 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,72 @@ jobs:
printf 'frontend-build\tnpm --workspace=frontend run build\tfrontend-build.log\n' >> "$COMMAND_LOG_DIR/COMMANDS.tsv"
npm --workspace=frontend run build 2>&1 | tee "$COMMAND_LOG_DIR/frontend-build.log"

# --- Frontend bundle size tracking (#1341) ---
# Captures vite build output size, compares against the checked-in
# baseline (frontend/.bundle-size-baseline.json), warns at +500KB and
# fails at +1MB, and posts the diff as a PR comment.
- name: Bundle size report
id: bundle-size
run: |
set -o pipefail
printf 'bundle-size\tnode scripts/check-bundle-size.mjs\tbundle-size.log\n' >> "$COMMAND_LOG_DIR/COMMANDS.tsv"
node scripts/check-bundle-size.mjs 2>&1 | tee "$COMMAND_LOG_DIR/bundle-size.log"
continue-on-error: true

- name: Check bundle size budget
run: |
FAIL="$(cat /tmp/bundle-size.json 2>/dev/null | node -e "let d='';process.stdin.on('data',c=>d+=c).on('end',()=>{try{console.log(JSON.parse(d).fail)}catch{console.log('false')}})" || echo false)"
if [ "$FAIL" = "true" ]; then
echo "::error::Frontend bundle grew by more than 1MB over baseline. Run 'npm --workspace=frontend run bundle:size' locally and justify or reduce the increase."
exit 1
fi

- name: Post bundle size PR comment
if: github.event_name == 'pull_request' && always()
uses: actions/github-script@v7
with:
script: |
const fs = require('fs');
let body = '### Frontend bundle size\n\nBundle size report unavailable (build may have failed).';
try {
body = fs.readFileSync('/tmp/bundle-size.md', 'utf8');
} catch {}
const { data: comments } = await github.rest.issues.listComments({
owner: context.repo.owner,
repo: context.repo.repo,
issue_number: context.issue.number,
});
const marker = '<!-- bridge-watch-bundle-size -->';
const full = `${marker}\n${body}`;
const existing = comments.find(c => c.body && c.body.includes(marker));
if (existing) {
await github.rest.issues.updateComment({
owner: context.repo.owner,
repo: context.repo.repo,
comment_id: existing.id,
body: full,
});
} else {
await github.rest.issues.createComment({
owner: context.repo.owner,
repo: context.repo.repo,
issue_number: context.issue.number,
body: full,
});
}

- name: Upload bundle size report
if: always()
uses: actions/upload-artifact@v4
with:
name: frontend-bundle-size
path: |
/tmp/bundle-size.json
/tmp/bundle-size.md
frontend/dist/**/.vite/manifest.json
retention-days: 30
if-no-files-found: warn

- name: Storybook build
run: |
set -o pipefail
Expand All @@ -284,6 +350,12 @@ jobs:
printf 'frontend-test\tnpm --workspace=frontend run test\tfrontend-test.log\n' >> "$COMMAND_LOG_DIR/COMMANDS.tsv"
npm --workspace=frontend run test 2>&1 | tee "$COMMAND_LOG_DIR/frontend-test.log"

- name: Accessibility tests (axe-core)
run: |
set -o pipefail
printf 'frontend-a11y\tnpm --workspace=frontend run test:a11y\tfrontend-a11y.log\n' >> "$COMMAND_LOG_DIR/COMMANDS.tsv"
npm --workspace=frontend run test:a11y 2>&1 | tee "$COMMAND_LOG_DIR/frontend-a11y.log"

- name: Visual Regression Tests
run: |
set -o pipefail
Expand Down
234 changes: 233 additions & 1 deletion backend/src/services/crossChainAnomalyEngine.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import { getCircuitBreakerService, PauseScope } from "./circuitBreaker.service.j
import type { FederatedEvent } from "./eventFederation/types.js";
import { getDatabase } from "../database/connection.js";

export type AnomalyType = "double_spend" | "nonce_jump" | "reentrancy" | "threshold_breach";
export type AnomalyType = "double_spend" | "nonce_jump" | "reentrancy" | "threshold_breach" | "high_velocity_drain";

export interface DetectedAnomaly {
id: string;
Expand All @@ -19,6 +19,20 @@ export interface DetectedAnomaly {
timestamp: number;
}

export interface VelocityStats {
windowMs: number;
netOutflow: number;
baselineMean: number;
baselineStd: number;
zScore: number | null;
sampleWindows: number;
}

export interface VelocityAnomalyResult {
anomalous: boolean;
stats: VelocityStats;
}

export interface FlashPauseResult {
triggered: boolean;
bridgeId: string;
Expand All @@ -33,24 +47,41 @@ export interface AnomalyEngineOptions {
anomalyThreshold?: number;
nonceWindowSeconds?: number;
txHashWindowSeconds?: number;
/** Rolling net-outflow velocity window (default 15 minutes, #1270). */
velocityWindowMs?: number;
/** Baseline lookback for velocity z-scores (default 7 days, #1270). */
velocityBaselineMs?: number;
/** Std-dev threshold above baseline mean (default 3, #1270). */
velocitySigmaThreshold?: number;
}

export class CrossChainAnomalyEngineService {
private readonly windowSeconds: number;
private readonly anomalyThreshold: number;
private readonly nonceWindowSeconds: number;
private readonly txHashWindowSeconds: number;
/** Rolling 15-minute net outflow velocity window (#1270). */
private readonly velocityWindowMs: number;
/** 7-day baseline for velocity z-scores (#1270). */
private readonly velocityBaselineMs: number;
/** Trigger when velocity exceeds this many std-devs above baseline (#1270). */
private readonly velocitySigmaThreshold: number;

// L1 In-Memory sliding window and state cache for sub-millisecond graph analysis
private readonly memoryStore = new Map<string, string | number>();
private readonly memoryAnomalies = new Map<string, DetectedAnomaly[]>();
private readonly memoryBreakers = new Map<string, boolean>();
/** Per-bridge outflow samples {t, amount} for velocity analysis (#1270). */
private readonly memoryOutflows = new Map<string, Array<{ t: number; amount: number }>>();

constructor(options: AnomalyEngineOptions = {}) {
this.windowSeconds = options.windowSeconds ?? 5;
this.anomalyThreshold = options.anomalyThreshold ?? 2;
this.nonceWindowSeconds = options.nonceWindowSeconds ?? 3600;
this.txHashWindowSeconds = options.txHashWindowSeconds ?? 3600;
this.velocityWindowMs = options.velocityWindowMs ?? 15 * 60 * 1000;
this.velocityBaselineMs = options.velocityBaselineMs ?? 7 * 24 * 60 * 60 * 1000;
this.velocitySigmaThreshold = options.velocitySigmaThreshold ?? 3;
}

/**
Expand Down Expand Up @@ -130,6 +161,21 @@ export class CrossChainAnomalyEngineService {
anomalies.push(anomaly);
}

// 4. Volume-weighted velocity anomaly detection (#1270): a rapid succession
// of small withdrawals can drain reserves without tripping static supply
// mismatch thresholds. Track rolling 15-min net outflow per bridge and
// alert when it exceeds 3 std-devs above the 7-day baseline.
const outflowAmount = this.extractOutflowAmount(event);
if (outflowAmount !== undefined && outflowAmount > 0) {
const velocityAnomaly = await this.processOutflow(bridgeId, chainId, outflowAmount, now, {
eventId: event.id,
sourceId: event.sourceId,
sequenceId,
depositTxHash,
});
if (velocityAnomaly) anomalies.push(velocityAnomaly);
}

// Record any detected anomalies and evaluate 5-second Flash-Pause threshold
if (anomalies.length > 0) {
for (const anomaly of anomalies) {
Expand Down Expand Up @@ -257,6 +303,192 @@ export class CrossChainAnomalyEngineService {
return isReentrant;
}

// ── Volume-weighted velocity anomaly detection (#1270) ───────────────
//
// Static supply-mismatch thresholds miss a rapid succession of small
// withdrawals that collectively drain reserves. We track a rolling 15-min
// net outflow velocity per bridge and compare it against a 7-day baseline:
// when the current window exceeds mean + 3σ we emit an immediate
// HIGH_VELOCITY_DRAIN warning alert (and record a `high_velocity_drain`
// anomaly so flash-pause counting still sees it).

/**
* Extract a positive outflow amount from a federated event, if present.
* Outflow-shaped events are bridge releases / withdrawals; deposits and
* generic ledger closes carry no outflow weight.
*/
private extractOutflowAmount(event: FederatedEvent): number | undefined {
const raw = (event.raw ?? {}) as Record<string, unknown>;
const candidates: unknown[] = [
event.amount,
raw.amount,
raw.withdrawalAmount,
raw.withdrawal_amount,
raw.value,
raw.outflow,
raw.netOutflow,
raw.net_outflow,
];
for (const c of candidates) {
const n = typeof c === "string" ? Number(c) : typeof c === "number" ? c : NaN;
if (Number.isFinite(n) && n > 0) {
// Only count outflow-shaped event types to avoid deposits inflating velocity.
const t = event.type;
if (t === "bridge_release" || t === "transfer" || t === "payment" || t === "swap") return n;
return n;
}
}
return undefined;
}

/** Append an outflow sample and prune anything older than the baseline window. */
async recordOutflow(bridgeId: string, amount: number, timestamp: number = Date.now()): Promise<void> {
if (!Number.isFinite(amount) || amount <= 0) return;
const list = this.memoryOutflows.get(bridgeId) ?? [];
list.push({ t: timestamp, amount });
const cutoff = timestamp - this.velocityBaselineMs;
const pruned = list.filter((s) => s.t >= cutoff);
this.memoryOutflows.set(bridgeId, pruned);
try {
const key = `ccae:outflow:${bridgeId}`;
await redis.zadd(key, timestamp, JSON.stringify({ t: timestamp, amount }));
await redis.zremrangebyscore(key, "-inf", cutoff);
await redis.expire(key, Math.ceil(this.velocityBaselineMs / 1000) + 3600);
} catch {
// Redis optional L2
}
}

/** Sum of outflows in `[now - windowMs, now]` for a bridge. */
getNetOutflowVelocity(bridgeId: string, now: number = Date.now(), windowMs: number = this.velocityWindowMs): number {
const cutoff = now - windowMs;
const list = this.memoryOutflows.get(bridgeId) ?? [];
let sum = 0;
for (const s of list) if (s.t >= cutoff && s.t <= now) sum += s.amount;
return sum;
}

/**
* Mean/std of per-`windowMs` bucketed net outflow over the baseline lookback.
* Buckets the baseline window into velocity-sized windows so the z-score
* compares like-for-like velocities rather than raw totals.
*/
getVelocityBaseline(
bridgeId: string,
now: number = Date.now(),
windowMs: number = this.velocityWindowMs,
baselineMs: number = this.velocityBaselineMs
): { mean: number; std: number; sampleWindows: number } {
const list = (this.memoryOutflows.get(bridgeId) ?? []).filter(
(s) => s.t >= now - baselineMs && s.t <= now
);
const bucketCount = Math.max(1, Math.floor(baselineMs / windowMs));
const buckets = new Array<number>(bucketCount).fill(0);
for (const s of list) {
const idx = Math.min(bucketCount - 1, Math.floor((now - s.t) / windowMs));
buckets[bucketCount - 1 - idx] += s.amount;
}
// Drop the in-progress (most recent) bucket so a forming drain does not
// inflate its own baseline.
const samples = buckets.slice(0, Math.max(0, bucketCount - 1));
if (samples.length === 0) return { mean: 0, std: 0, sampleWindows: 0 };
const mean = samples.reduce((a, b) => a + b, 0) / samples.length;
const variance = samples.reduce((a, b) => a + (b - mean) ** 2, 0) / samples.length;
return { mean, std: Math.sqrt(variance), sampleWindows: samples.length };
}

/** Z-score check of current velocity against the 7-day baseline. */
checkVelocityAnomaly(bridgeId: string, now: number = Date.now()): VelocityAnomalyResult {
const netOutflow = this.getNetOutflowVelocity(bridgeId, now);
const { mean, std, sampleWindows } = this.getVelocityBaseline(bridgeId, now);
if (sampleWindows < 2 || netOutflow <= 0) {
return {
anomalous: false,
stats: { windowMs: this.velocityWindowMs, netOutflow, baselineMean: mean, baselineStd: std, zScore: null, sampleWindows },
};
}
// Zero-variance baseline: any material outflow above the mean is anomalous.
if (std === 0) {
const anomalous = netOutflow > mean && netOutflow - mean > 0;
return {
anomalous,
stats: {
windowMs: this.velocityWindowMs,
netOutflow,
baselineMean: mean,
baselineStd: std,
zScore: anomalous ? Number.POSITIVE_INFINITY : 0,
sampleWindows,
},
};
}
const zScore = (netOutflow - mean) / std;
return {
anomalous: zScore > this.velocitySigmaThreshold,
stats: { windowMs: this.velocityWindowMs, netOutflow, baselineMean: mean, baselineStd: std, zScore, sampleWindows },
};
}

/**
* Record an outflow sample, evaluate velocity, and on breach emit a
* HIGH_VELOCITY_DRAIN warning alert + anomaly. Returns the anomaly when
* triggered, otherwise null.
*/
async processOutflow(
bridgeId: string,
chainId: string,
amount: number,
now: number = Date.now(),
context: { eventId?: string; sourceId?: string; sequenceId?: number; depositTxHash?: string } = {}
): Promise<DetectedAnomaly | null> {
await this.recordOutflow(bridgeId, amount, now);
const { anomalous, stats } = this.checkVelocityAnomaly(bridgeId, now);
if (!anomalous) return null;

const anomaly: DetectedAnomaly = {
id: `hv_${context.eventId ?? bridgeId}_${now}`,
type: "high_velocity_drain",
bridgeId,
chainId,
sequenceId: context.sequenceId,
depositTxHash: context.depositTxHash,
details: {
message: `HIGH_VELOCITY_DRAIN: net outflow ${stats.netOutflow} over 15m exceeds baseline mean ${stats.baselineMean.toFixed(2)} + 3σ (${stats.baselineStd.toFixed(2)}); z=${stats.zScore === null ? "n/a" : Number(stats.zScore).toFixed(2)}`,
alert: "HIGH_VELOCITY_DRAIN",
severity: "warning",
eventId: context.eventId,
sourceId: context.sourceId,
velocityWindowMs: stats.windowMs,
netOutflow: stats.netOutflow,
baselineMean: stats.baselineMean,
baselineStd: stats.baselineStd,
zScore: stats.zScore,
},
timestamp: now,
};

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: "HIGH_VELOCITY_DRAIN",
priority: "warning",
triggered_value: stats.netOutflow,
threshold: stats.baselineMean + this.velocitySigmaThreshold * stats.baselineStd,
metric: "net_outflow_velocity_15m",
webhook_delivered: false,
webhook_attempts: 0,
});
} catch (err) {
logger.warn({ err, bridgeId }, "Could not persist HIGH_VELOCITY_DRAIN alert event to DB");
}

logger.warn({ bridgeId, stats }, "HIGH_VELOCITY_DRAIN velocity anomaly detected");
return anomaly;
}

/**
* Records detected anomaly into L1 Memory + L2 Redis sliding window.
*/
Expand Down
Loading
Loading