Skip to content
Merged
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
78 changes: 78 additions & 0 deletions backend/src/api/routes/bftOracle.routes.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
import { z } from "zod";
import { logger } from "../../utils/logger.js";
import { bftOracleAggregatorService } from "../../services/bftOracleAggregator.service.js";

const reportItemSchema = z.object({
providerKey: z.string().min(1),
price: z.number().positive(),
healthScore: z.number().optional(),
timestamp: z.string().optional().default(() => new Date().toISOString()),
signature: z.string().optional(),
});

const aggregateSchema = z.object({
assetCode: z.string().min(1),
reports: z.array(reportItemSchema).min(1),
});

const registerNodeSchema = z.object({
providerKey: z.string().min(1),
displayName: z.string().min(1),
publicKey: z.string().min(1),
stakeWeight: z.number().positive().optional().default(1.0),
});

export async function bftOracleRoutes(server: FastifyInstance) {
server.post("/aggregate", async (request: FastifyRequest<{ Body: z.infer<typeof aggregateSchema> }>, reply: FastifyReply) => {
try {
const { assetCode, reports } = aggregateSchema.parse(request.body);
const result = await bftOracleAggregatorService.aggregateBftState(assetCode, reports);
return reply.code(200).send(result);
} catch (error) {
logger.error(error, "Failed to run BFT state aggregation");
return reply.code(400).send({ error: "Failed to run BFT state aggregation", details: String(error) });
}
});

server.post("/providers", async (request: FastifyRequest<{ Body: z.infer<typeof registerNodeSchema> }>, reply: FastifyReply) => {
try {
const body = registerNodeSchema.parse(request.body);
const provider = await bftOracleAggregatorService.registerProviderNode(body);
return reply.code(201).send(provider);
} catch (error) {
logger.error(error, "Failed to register BFT provider node");
return reply.code(400).send({ error: "Failed to register BFT provider node", details: String(error) });
}
});

server.get("/providers", async (_request: FastifyRequest, reply: FastifyReply) => {
try {
const providers = await bftOracleAggregatorService.getRegisteredProviders();
return reply.code(200).send({ providers, total: providers.length });
} catch (error) {
logger.error(error, "Failed to fetch BFT provider nodes");
return reply.code(500).send({ error: "Failed to fetch BFT provider nodes" });
}
});

server.get("/rounds/:assetCode", async (request: FastifyRequest<{ Params: { assetCode: string } }>, reply: FastifyReply) => {
try {
const rounds = await bftOracleAggregatorService.getPastRounds(request.params.assetCode);
return reply.code(200).send({ assetCode: request.params.assetCode, rounds });
} catch (error) {
logger.error(error, "Failed to fetch BFT rounds");
return reply.code(500).send({ error: "Failed to fetch BFT rounds" });
}
});

server.get("/slashing-events", async (_request: FastifyRequest, reply: FastifyReply) => {
try {
const events = await bftOracleAggregatorService.getSlashingHistory();
return reply.code(200).send({ events, total: events.length });
} catch (error) {
logger.error(error, "Failed to fetch slashing events");
return reply.code(500).send({ error: "Failed to fetch slashing events" });
}
});
}
5 changes: 5 additions & 0 deletions backend/src/api/routes/route-groups/provider-routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import type { FastifyInstance } from "fastify";
import { providerHealthRegistryRoutes } from "../providerHealthRegistry.routes.js";
import { providerAllowlistRoutes } from "../providerAllowlist.routes.js";
import { providerCircuitBreakerRoutes } from "../providerCircuitBreaker.routes.js";
import { bftOracleRoutes } from "../bftOracle.routes.js";

export async function registerProviderRoutes(server: FastifyInstance): Promise<void> {
server.register(providerHealthRegistryRoutes, {
Expand All @@ -13,4 +14,8 @@ export async function registerProviderRoutes(server: FastifyInstance): Promise<v
server.register(providerCircuitBreakerRoutes, {
prefix: "/api/v1/providers/circuit-breaker",
});
server.register(bftOracleRoutes, {
prefix: "/api/v1/bft-oracle",
});
}

53 changes: 53 additions & 0 deletions backend/src/database/migrations/049_bft_oracle_aggregator.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
import type { Knex } from "knex";

export async function up(knex: Knex): Promise<void> {
await knex.schema.createTable("bft_oracle_providers", (table) => {
table.string("provider_key", 120).primary();
table.string("display_name", 120).notNullable();
table.string("public_key", 120).notNullable();
table.double("stake_weight").notNullable().defaultTo(1.0);
table.string("status", 20).notNullable().defaultTo("active");
table.boolean("slashed").notNullable().defaultTo(false);
table.timestamp("slashed_at", { useTz: true }).nullable();
table.string("slash_reason", 255).nullable();
table.integer("total_submissions").notNullable().defaultTo(0);
table.integer("total_slashes").notNullable().defaultTo(0);
table.timestamps(true, true);
table.index(["status"], "idx_bft_providers_status");
});

await knex.schema.createTable("bft_consensus_rounds", (table) => {
table.uuid("id").primary().defaultTo(knex.raw("gen_random_uuid()"));
table.string("asset_code", 20).notNullable();
table.double("consensus_price").notNullable();
table.double("median_of_medians").notNullable();
table.double("mean").notNullable();
table.double("std_dev").notNullable();
table.integer("total_providers").notNullable();
table.integer("valid_providers").notNullable();
table.boolean("quorum_reached").notNullable();
table.text("aggregate_signature").nullable();
table.timestamp("created_at", { useTz: true }).notNullable().defaultTo(knex.fn.now());
table.index(["asset_code", "created_at"], "idx_bft_rounds_asset_time");
});

await knex.schema.createTable("bft_slashing_events", (table) => {
table.uuid("id").primary().defaultTo(knex.raw("gen_random_uuid()"));
table.string("provider_key", 120).notNullable();
table.uuid("round_id").nullable();
table.string("asset_code", 20).notNullable();
table.double("reported_value").notNullable();
table.double("consensus_value").notNullable();
table.double("deviation_sigma").notNullable();
table.double("slashed_stake").notNullable();
table.string("reason", 255).notNullable();
table.timestamp("created_at", { useTz: true }).notNullable().defaultTo(knex.fn.now());
table.index(["provider_key", "created_at"], "idx_bft_slashing_provider_time");
});
}

export async function down(knex: Knex): Promise<void> {
await knex.schema.dropTableIfExists("bft_slashing_events");
await knex.schema.dropTableIfExists("bft_consensus_rounds");
await knex.schema.dropTableIfExists("bft_oracle_providers");
}
45 changes: 45 additions & 0 deletions backend/src/database/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -532,3 +532,48 @@ export interface CircuitBreakerActionLog {
executed_at: Date;
}

// ─── bft_oracle_providers / bft_consensus_rounds / bft_slashing_events ───────

export interface BftOracleProvider {
provider_key: string;
display_name: string;
public_key: string;
stake_weight: number;
status: "active" | "slashed" | "degraded" | "suspended";
slashed: boolean;
slashed_at: Date | null;
slash_reason: string | null;
total_submissions: number;
total_slashes: number;
created_at: Date;
updated_at: Date;
}

export interface BftConsensusRound {
id: string;
asset_code: string;
consensus_price: number;
median_of_medians: number;
mean: number;
std_dev: number;
total_providers: number;
valid_providers: number;
quorum_reached: boolean;
aggregate_signature: string | null;
created_at: Date;
}

export interface BftSlashingEvent {
id: string;
provider_key: string;
round_id: string | null;
asset_code: string;
reported_value: number;
consensus_value: number;
deviation_sigma: number;
slashed_stake: number;
reason: string;
created_at: Date;
}


155 changes: 155 additions & 0 deletions backend/src/services/__tests__/bftOracleAggregator.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,155 @@
import { describe, it, expect, beforeEach, vi } from "vitest";
import { BftOracleAggregatorService, type OracleReport, type OracleProviderNode } from "../bftOracleAggregator.service.js";

vi.mock("../database/connection.js", () => {
const mockDb: any = vi.fn().mockImplementation(() => mockDb);
mockDb.schema = {
hasTable: vi.fn().mockResolvedValue(true),
};
mockDb.where = vi.fn().mockReturnValue(mockDb);
mockDb.update = vi.fn().mockResolvedValue(1);
mockDb.insert = vi.fn().mockResolvedValue([1]);
mockDb.select = vi.fn().mockResolvedValue([]);
mockDb.first = vi.fn().mockResolvedValue(null);
mockDb.orderBy = vi.fn().mockReturnValue(mockDb);
mockDb.limit = vi.fn().mockResolvedValue([]);
mockDb.raw = vi.fn((str) => str);
mockDb.fn = { now: () => new Date().toISOString() };
return { getDatabase: () => mockDb };
});

vi.mock("../providerHealthRegistry.service.js", () => {
return {
providerHealthRegistryService: {
flagAndSlashProvider: vi.fn().mockResolvedValue(true),
},
};
});

describe("BftOracleAggregatorService", () => {
let service: BftOracleAggregatorService;

beforeEach(() => {
service = new BftOracleAggregatorService("test-secret-key");
vi.clearAllMocks();
});

const generateNodes = (count: number): OracleProviderNode[] => {
return Array.from({ length: count }, (_, i) => ({
providerKey: `oracle_node_${i + 1}`,
displayName: `Oracle Node ${i + 1}`,
publicKey: `pubkey_oracle_node_${i + 1}`,
stakeWeight: 1.0,
status: "active",
slashed: false,
slashedAt: null,
slashReason: null,
totalSubmissions: 0,
totalSlashes: 0,
}));
};

it("computes consensus when 3f+1 honest nodes report (N=4, f=1, Quorum=3)", async () => {
const nodes = generateNodes(4);
const reports: OracleReport[] = [
{ providerKey: "oracle_node_1", price: 100.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_2", price: 100.2, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_3", price: 99.8, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_4", price: 100.1, timestamp: new Date().toISOString() },
];

const result = await service.aggregateBftState("USDC", reports, nodes);

expect(result.quorumReached).toBe(true);
expect(result.validProviders).toBe(4);
expect(result.requiredQuorum).toBe(3);
expect(result.consensusPrice).toBeCloseTo(100.025, 2);
expect(result.slashedProviders).toHaveLength(0);
expect(service.verifyAggregatePayload(result)).toBe(true);
});

it("detects and slashes > 5 sigma Byzantine outlier node", async () => {
const nodes = generateNodes(4);
const reports: OracleReport[] = [
{ providerKey: "oracle_node_1", price: 100.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_2", price: 100.1, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_3", price: 99.9, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_4", price: 500.0, timestamp: new Date().toISOString() }, // Malicious outlier (> 5 sigma)
];

const result = await service.aggregateBftState("USDC", reports, nodes);

expect(result.quorumReached).toBe(true);
expect(result.slashedProviders).toContain("oracle_node_4");
expect(result.validProviders).toBe(3);
expect(result.consensusPrice).toBeCloseTo(100.0, 1);
});

it("fails quorum when reporting node count is below 2f+1 threshold (network partition)", async () => {
const nodes = generateNodes(4); // N=4, f=1, Quorum=3
const reports: OracleReport[] = [
{ providerKey: "oracle_node_1", price: 100.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_2", price: 100.1, timestamp: new Date().toISOString() },
]; // Only 2 reports submitted

const result = await service.aggregateBftState("USDC", reports, nodes);

expect(result.quorumReached).toBe(false);
expect(result.reportingProviders).toBe(2);
expect(result.requiredQuorum).toBe(3);
expect(result.consensusPrice).toBe(0);
});

it("ignores reports from already slashed or suspended oracle nodes", async () => {
const nodes = generateNodes(4);
nodes[3].slashed = true;
nodes[3].status = "slashed";

const reports: OracleReport[] = [
{ providerKey: "oracle_node_1", price: 100.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_2", price: 100.1, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_3", price: 99.9, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_4", price: 999.9, timestamp: new Date().toISOString() },
];

const result = await service.aggregateBftState("USDC", reports, nodes);

expect(result.quorumReached).toBe(true);
expect(result.reportingProviders).toBe(3);
expect(result.validProviders).toBe(3);
expect(result.slashedProviders).not.toContain("oracle_node_4");
});

it("correctly weights stake median for consensus calculation", async () => {
const nodes = generateNodes(3);
nodes[0].stakeWeight = 10.0; // High stake
nodes[1].stakeWeight = 1.0;
nodes[2].stakeWeight = 1.0;

const reports: OracleReport[] = [
{ providerKey: "oracle_node_1", price: 105.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_2", price: 100.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_3", price: 98.0, timestamp: new Date().toISOString() },
];

const result = await service.aggregateBftState("ETH", reports, nodes);

expect(result.weightedMedianPrice).toBe(105.0);
expect(result.consensusPrice).toBeGreaterThan(100.0);
});

it("deduplicates duplicate report submissions from the same provider key", async () => {
const nodes = generateNodes(4);
const reports: OracleReport[] = [
{ providerKey: "oracle_node_1", price: 100.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_1", price: 100.0, timestamp: new Date().toISOString() },
{ providerKey: "oracle_node_1", price: 100.0, timestamp: new Date().toISOString() },
];

const result = await service.aggregateBftState("USDC", reports, nodes);

expect(result.quorumReached).toBe(false);
expect(result.reportingProviders).toBe(1);
});
});

Loading
Loading