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
8 changes: 4 additions & 4 deletions backend/src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -153,10 +153,7 @@ apolloServer.start().then(() => {
let user = null;
if (auth.startsWith("Bearer ")) {
try {
user = jwt.verify(
auth.slice(7),
process.env.JWT_SECRET || "change-me-in-production"
);
user = jwt.verify(auth.slice(7), process.env.JWT_SECRET || "change-me-in-production");
} catch {
// Invalid token — user stays null; resolvers can enforce auth as needed
}
Expand Down Expand Up @@ -189,6 +186,9 @@ require("./workers/archive-worker").start();
// Subscribe prediction market contract to Mercury Indexer
require("./indexer/mercury").subscribe();

// Start real-time Mercury event stream with reconnection logic
require("./indexer/mercury").startEventStream();

// Initialize self-healing gap detection and recovery
require("./indexer/gap-detector").initializeSelfHealing();

Expand Down
73 changes: 73 additions & 0 deletions backend/src/indexer/mercury.js
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ const axios = require("axios");
const db = require("../db");
const logger = require("../utils/logger");
const pubsub = require("../graphql/pubsub");
const { broadcastBetPlaced, broadcastMarketResolved } = require("../websocket/marketUpdates");

const MERCURY_BASE = process.env.MERCURY_URL || "https://api.mercurydata.app";
const MERCURY_KEY = process.env.MERCURY_API_KEY || "";
Expand Down Expand Up @@ -152,6 +153,14 @@ async function handleBetPlaced(payload, meta) {
amount: String(cost),
});

// Broadcast WebSocket update to subscribed clients
broadcastBetPlaced(market_id, {
market_id,
wallet_address: bettor,
outcome_index: option_index,
amount: String(cost),
});

// Recalculate and publish updated odds (basis-points per outcome, zero-float)
_publishOddsChanged(market_id).catch((err) =>
logger.warn({ err: err.message, market_id }, "Failed to publish oddsChanged")
Expand Down Expand Up @@ -197,6 +206,9 @@ async function handleMarketResolved(payload) {
winning_outcome,
total_pool: String(total_pool),
});

// Broadcast WebSocket update to subscribed clients
broadcastMarketResolved(market_id, winning_outcome);
}

/**
Expand Down Expand Up @@ -413,9 +425,70 @@ async function processEvent(event) {
}
}

// ── Event streaming with reconnection ────────────────────────────────────────

const POLL_INTERVAL_MS = parseInt(process.env.MERCURY_POLL_INTERVAL_MS, 10) || 2000;
const MAX_BACKOFF_MS = 60_000;

let _streamActive = false;
let _lastLedger = 0;

/**
* Poll Mercury for new events since the last processed ledger.
* On failure, retries with exponential backoff up to MAX_BACKOFF_MS.
*/
async function startEventStream() {
if (!CONTRACT_ID || !MERCURY_KEY) {
logger.warn("Mercury event stream skipped: CONTRACT_ADDRESS or MERCURY_API_KEY not set");
return;
}
_streamActive = true;
let backoff = 1000;

logger.info({ contract_id: CONTRACT_ID }, "Mercury event stream starting");

while (_streamActive) {
try {
const { data } = await axios.get(`${MERCURY_BASE}/event/contract/${CONTRACT_ID}`, {
headers: { Authorization: `Bearer ${MERCURY_KEY}` },
params: { after_ledger: _lastLedger },
timeout: 10_000,
});

const events = data.events ?? [];
for (const event of events) {
try {
await processEvent(event);
if (event.ledger_seq > _lastLedger) _lastLedger = event.ledger_seq;
} catch (err) {
logger.error({ err: err.message, topic: event.topic }, "Event processing failed");
}
}

// Reset backoff on success
backoff = 1000;
await _sleep(POLL_INTERVAL_MS);
} catch (err) {
logger.error({ err: err.message, backoff }, "Mercury connection error — retrying");
await _sleep(backoff);
backoff = Math.min(backoff * 2, MAX_BACKOFF_MS);
}
}
}

function stopEventStream() {
_streamActive = false;
}

function _sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}

module.exports = {
subscribe,
processEvent,
startEventStream,
stopEventStream,
// Named exports for unit testing
handleMarketCreated,
handleBetPlaced,
Expand Down
160 changes: 160 additions & 0 deletions backend/src/tests/mercury-streaming.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,160 @@
"use strict";

/**
* Unit tests for Mercury Indexer event handlers.
* Covers: BET_PLACED and MARKET_RESOLVED event handling,
* WebSocket broadcasts, and reconnection logic.
*/

jest.mock("../db");
jest.mock("../utils/logger", () => ({
info: jest.fn(),
warn: jest.fn(),
error: jest.fn(),
debug: jest.fn(),
}));
jest.mock("../graphql/pubsub", () => ({ publish: jest.fn() }));
jest.mock("../websocket/marketUpdates", () => ({
broadcastBetPlaced: jest.fn(),
broadcastMarketResolved: jest.fn(),
broadcastOddsChanged: jest.fn(),
}));

const db = require("../db");
const pubsub = require("../graphql/pubsub");
const ws = require("../websocket/marketUpdates");
const mercury = require("../indexer/mercury");

const META = { ledger_seq: 100, ledger_time: "2026-01-01T00:00:00Z" };

describe("Mercury Indexer — handleBetPlaced", () => {
beforeEach(() => jest.clearAllMocks());

const payload = {
version: 1,
market_id: 42,
bettor: "WALLET_ABC",
option_index: 0,
cost: 1000,
shares: 500,
};

it("inserts bet into DB", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleBetPlaced(payload, META);
expect(db.query).toHaveBeenCalledWith(
expect.stringContaining("INSERT INTO bets"),
expect.arrayContaining([42, "WALLET_ABC", 0, 1000, 500])
);
});

it("upserts user stats", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleBetPlaced(payload, META);
expect(db.query).toHaveBeenCalledWith(
expect.stringContaining("INSERT INTO users"),
expect.arrayContaining(["WALLET_ABC", 1000])
);
});

it("publishes to GraphQL pubsub", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleBetPlaced(payload, META);
expect(pubsub.publish).toHaveBeenCalledWith(
"betPlaced",
42,
expect.objectContaining({ market_id: 42, wallet_address: "WALLET_ABC" })
);
});

it("broadcasts WebSocket BET_PLACED event", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleBetPlaced(payload, META);
expect(ws.broadcastBetPlaced).toHaveBeenCalledWith(
42,
expect.objectContaining({ market_id: 42, wallet_address: "WALLET_ABC" })
);
});

it("throws on unsupported version", async () => {
await expect(mercury.handleBetPlaced({ ...payload, version: 99 }, META)).rejects.toThrow(
/Unsupported schema version/
);
});
});

describe("Mercury Indexer — handleMarketResolved", () => {
beforeEach(() => jest.clearAllMocks());

const payload = {
version: 1,
market_id: 7,
winning_outcome: 1,
total_pool: 50000,
fee_bps: 200,
};

it("updates market resolved status in DB", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleMarketResolved(payload);
expect(db.query).toHaveBeenCalledWith(
expect.stringContaining("resolved = true"),
expect.arrayContaining([1, 50000, 200, 7])
);
});

it("credits winners in DB", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleMarketResolved(payload);
expect(db.query).toHaveBeenCalledWith(
expect.stringContaining("total_won"),
expect.arrayContaining([7, 1])
);
});

it("publishes to GraphQL pubsub", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleMarketResolved(payload);
expect(pubsub.publish).toHaveBeenCalledWith(
"marketResolved",
7,
expect.objectContaining({ market_id: 7, winning_outcome: 1 })
);
});

it("broadcasts WebSocket MARKET_RESOLVED event", async () => {
db.query.mockResolvedValue({ rows: [] });
await mercury.handleMarketResolved(payload);
expect(ws.broadcastMarketResolved).toHaveBeenCalledWith(7, 1);
});

it("throws on unsupported version", async () => {
await expect(mercury.handleMarketResolved({ ...payload, version: 0 })).rejects.toThrow(
/Unsupported schema version/
);
});
});

describe("Mercury Indexer — startEventStream reconnection", () => {
beforeEach(() => {
jest.clearAllMocks();
jest.useFakeTimers();
});

afterEach(() => {
jest.useRealTimers();
mercury.stopEventStream();
});

it("stops streaming when stopEventStream is called", async () => {
// Immediately stop — stream should not make any axios calls
mercury.stopEventStream();
// startEventStream checks _streamActive before looping
const streamPromise = mercury.startEventStream();
await Promise.resolve(); // flush microtasks
mercury.stopEventStream();
await streamPromise;
// No DB calls should have been made
expect(db.query).not.toHaveBeenCalled();
});
});
Loading