forked from Hahfyeex/Stellar-PolyMarket
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmercury.js
More file actions
504 lines (447 loc) Β· 16 KB
/
Copy pathmercury.js
File metadata and controls
504 lines (447 loc) Β· 16 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
/**
* indexer/mercury.js
*
* Mercury Indexer integration β versioned event parser.
*
* Subscribes the prediction market contract to Mercury and routes each
* incoming event to a typed handler. Every handler validates the event
* version field and upserts the relevant PostgreSQL rows.
*
* Supported topics (see docs/events.md for full schema):
* MktCreate β markets table
* BetPlace β bets + users tables
* MktResolv β markets.resolved, users.total_won
* MktVoid β markets.status = VOIDED
* MktPause β markets.is_paused
* Payout β payout_batches table
* LpSeed β lp_contributions table
* LpClaim β lp_claims table
* Dispute β disputes table
* FeeColl β fee_collections table
*/
"use strict";
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 || "";
const CONTRACT_ID = process.env.CONTRACT_ADDRESS || "";
// Minimum schema version this parser understands.
const MIN_SUPPORTED_VERSION = 1;
// Maximum schema version this parser understands.
const MAX_SUPPORTED_VERSION = 1;
// ββ Subscription ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
/**
* Register the prediction market contract with Mercury so it starts
* delivering events. Safe to call on every startup β Mercury deduplicates.
*/
async function subscribe() {
if (!CONTRACT_ID || !MERCURY_KEY) {
logger.warn("Mercury subscription skipped: CONTRACT_ADDRESS or MERCURY_API_KEY not set");
return;
}
try {
await axios.post(
`${MERCURY_BASE}/event/subscribe`,
{ contract_id: CONTRACT_ID },
{ headers: { Authorization: `Bearer ${MERCURY_KEY}` } }
);
logger.info({ contract_id: CONTRACT_ID }, "Mercury subscription registered");
} catch (err) {
logger.error({ err: err.message }, "Mercury subscription failed");
}
}
// ββ Version guard βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
/**
* Assert the event payload version is within the supported range.
* Throws if the version is unknown so the caller can dead-letter the event.
*
* @param {object} payload
* @param {string} topic
*/
function assertVersion(payload, topic) {
const v = payload.version;
if (typeof v !== "number" || v < MIN_SUPPORTED_VERSION || v > MAX_SUPPORTED_VERSION) {
throw new Error(
`Unsupported schema version ${v} for topic "${topic}". ` +
`Expected ${MIN_SUPPORTED_VERSION}β${MAX_SUPPORTED_VERSION}.`
);
}
}
// ββ Event handlers ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
/**
* MktCreate β market created.
*
* Payload (v1):
* version, market_id, creator, question, options_count,
* deadline, token, lmsr_b, creation_fee, ledger_timestamp
*/
async function handleMarketCreated(payload, meta) {
assertVersion(payload, "MktCreate");
const { market_id, creator, question, options_count, deadline, token, lmsr_b, creation_fee } =
payload;
await db.query(
`INSERT INTO markets
(id, question, options_count, deadline, contract_address,
creator, lmsr_b, creation_fee, status, created_at)
VALUES ($1, $2, $3, to_timestamp($4), $5, $6, $7, $8, 'ACTIVE', $9)
ON CONFLICT (id) DO NOTHING`,
[
market_id,
question,
options_count,
deadline,
token,
creator,
lmsr_b,
creation_fee,
meta.ledger_time,
]
);
}
/**
* BetPlace β bet placed.
*
* Payload (v1):
* version, market_id, bettor, option_index, cost, shares, ledger_timestamp
*
* `cost` is the LMSR cost delta in stroops (what the bettor actually paid).
* `shares` is the number of outcome shares purchased.
*/
async function handleBetPlaced(payload, meta) {
assertVersion(payload, "BetPlace");
const { market_id, bettor, option_index, cost, shares } = payload;
await db.query(
`INSERT INTO bets
(market_id, wallet_address, outcome_index, cost, shares, created_at)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT DO NOTHING`,
[market_id, bettor, option_index, cost, shares, meta.ledger_time]
);
// Upsert user aggregate stats (cost = actual spend in stroops)
await db.query(
`INSERT INTO users (wallet_address, total_staked, bet_count, last_seen)
VALUES ($1, $2, 1, $3)
ON CONFLICT (wallet_address) DO UPDATE SET
total_staked = users.total_staked + EXCLUDED.total_staked,
bet_count = users.bet_count + 1,
last_seen = EXCLUDED.last_seen`,
[bettor, cost, meta.ledger_time]
);
// Publish real-time subscription events (amounts as strings β zero-float)
pubsub.publish("betPlaced", market_id, {
market_id,
wallet_address: bettor,
outcome_index: option_index,
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")
);
}
/**
* MktResolv β market resolved.
*
* Payload (v1):
* version, market_id, winning_outcome, total_pool, fee_bps, ledger_timestamp
*/
async function handleMarketResolved(payload) {
assertVersion(payload, "MktResolv");
const { market_id, winning_outcome, total_pool, fee_bps } = payload;
await db.query(
`UPDATE markets
SET resolved = true,
winning_outcome = $1,
total_pool = $2,
fee_bps = $3,
status = 'RESOLVED'
WHERE id = $4`,
[winning_outcome, total_pool, fee_bps, market_id]
);
// Credit winners: total_won += their bet cost (simplified; full payout in distributor)
await db.query(
`UPDATE users u
SET total_won = u.total_won + b.cost,
win_count = u.win_count + 1
FROM bets b
WHERE b.wallet_address = u.wallet_address
AND b.market_id = $1
AND b.outcome_index = $2`,
[market_id, winning_outcome]
);
// Publish real-time subscription event (amounts as strings β zero-float)
pubsub.publish("marketResolved", market_id, {
market_id,
winning_outcome,
total_pool: String(total_pool),
});
// Broadcast WebSocket update to subscribed clients
broadcastMarketResolved(market_id, winning_outcome);
}
/**
* MktVoid β conditional market voided.
*
* Payload (v1):
* version, market_id, condition_market_id, condition_outcome_actual, ledger_timestamp
*/
async function handleMarketVoided(payload) {
assertVersion(payload, "MktVoid");
const { market_id, condition_market_id, condition_outcome_actual } = payload;
await db.query(
`UPDATE markets
SET status = 'VOIDED',
condition_market_id = $2,
condition_outcome_actual = $3
WHERE id = $1`,
[market_id, condition_market_id, condition_outcome_actual]
);
}
/**
* MktPause β market paused or unpaused.
*
* Payload (v1):
* version, market_id, paused, ledger_timestamp
*/
async function handleMarketPaused(payload) {
assertVersion(payload, "MktPause");
const { market_id, paused } = payload;
await db.query(`UPDATE markets SET is_paused = $1 WHERE id = $2`, [paused, market_id]);
}
/**
* Payout β batch payout processed.
*
* Payload (v1):
* version, market_id, recipients_paid, total_distributed, cursor, ledger_timestamp
*/
async function handlePayoutClaimed(payload, meta) {
assertVersion(payload, "Payout");
const { market_id, recipients_paid, total_distributed, cursor } = payload;
await db.query(
`INSERT INTO payout_batches
(market_id, recipients_paid, total_distributed, cursor, processed_at)
VALUES ($1, $2, $3, $4, $5)`,
[market_id, recipients_paid, total_distributed, cursor, meta.ledger_time]
);
}
/**
* LpSeed β liquidity provided.
*
* Payload (v1):
* version, market_id, provider, amount, ledger_timestamp
*/
async function handleLiquidityProvided(payload, meta) {
assertVersion(payload, "LpSeed");
const { market_id, provider, amount } = payload;
await db.query(
`INSERT INTO lp_contributions (market_id, provider, amount, contributed_at)
VALUES ($1, $2, $3, $4)
ON CONFLICT (market_id, provider) DO UPDATE SET
amount = lp_contributions.amount + EXCLUDED.amount`,
[market_id, provider, amount, meta.ledger_time]
);
}
/**
* LpClaim β LP reward claimed.
*
* Payload (v1):
* version, market_id, lp, reward, ledger_timestamp
*/
async function handleLpRewardClaimed(payload, meta) {
assertVersion(payload, "LpClaim");
const { market_id, lp, reward } = payload;
await db.query(
`INSERT INTO lp_claims (market_id, lp_address, reward, claimed_at)
VALUES ($1, $2, $3, $4)`,
[market_id, lp, reward, meta.ledger_time]
);
}
/**
* Dispute β dispute raised.
*
* Payload (v1):
* version, market_id, disputer, bond_amount, ledger_timestamp
*/
async function handleDisputeRaised(payload, meta) {
assertVersion(payload, "Dispute");
const { market_id, disputer, bond_amount } = payload;
await db.query(
`INSERT INTO disputes (market_id, disputer, bond_amount, raised_at, active)
VALUES ($1, $2, $3, $4, true)
ON CONFLICT (market_id) DO UPDATE SET
disputer = EXCLUDED.disputer,
bond_amount = EXCLUDED.bond_amount,
raised_at = EXCLUDED.raised_at,
active = true`,
[market_id, disputer, bond_amount, meta.ledger_time]
);
await db.query(`UPDATE markets SET status = 'DISPUTED' WHERE id = $1`, [market_id]);
}
/**
* FeeColl β creation fee collected.
*
* Payload (v1):
* version, market_id, payer, fee_destination, amount, ledger_timestamp
*/
async function handleFeeCollected(payload, meta) {
assertVersion(payload, "FeeColl");
const { market_id, payer, fee_destination, amount } = payload;
await db.query(
`INSERT INTO fee_collections
(market_id, payer, fee_destination, amount, collected_at)
VALUES ($1, $2, $3, $4, $5)`,
[market_id, payer, fee_destination, amount, meta.ledger_time]
);
}
// ββ Odds helper βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
/**
* Recalculate per-outcome odds in basis-points (integer, zero-float) and
* publish to the oddsChanged subscription channel.
*
* odds_bps[i] = (outcome_i_stake * 10_000) / total_stake
* All arithmetic uses BigInt to avoid floating-point.
*/
async function _publishOddsChanged(market_id) {
const { rows } = await db.query(
`SELECT outcome_index, COALESCE(SUM(cost), 0) AS stake
FROM bets WHERE market_id = $1
GROUP BY outcome_index ORDER BY outcome_index`,
[market_id]
);
if (rows.length === 0) return;
const total = rows.reduce((acc, r) => acc + BigInt(r.stake), 0n);
const odds_bps = rows.map((r) =>
total === 0n ? "0" : String((BigInt(r.stake) * 10_000n) / total)
);
pubsub.publish("oddsChanged", market_id, { market_id, odds_bps });
}
// ββ Dispatcher ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
/**
* Process a single Mercury event object.
* Called by the webhook handler or the polling loop.
*
* @param {object} event - Mercury event envelope
* @param {string} event.topic - Event topic symbol (e.g. "MktCreate")
* @param {object} event.payload - Deserialised XDR data struct
* @param {string} event.tx_hash
* @param {number} event.event_index
* @param {number} event.ledger_seq
* @param {string} event.ledger_time - ISO timestamp
*/
async function processEvent(event) {
const { topic, payload, tx_hash, event_index, ledger_seq, ledger_time } = event;
const meta = { ledger_seq, ledger_time };
// Persist raw event for audit / replay before any handler runs
await db.query(
`INSERT INTO events
(contract_id, topic, payload, ledger_seq, ledger_time, tx_hash, event_index)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (tx_hash, event_index) DO NOTHING`,
[CONTRACT_ID, topic, JSON.stringify(payload), ledger_seq, ledger_time, tx_hash, event_index]
);
try {
switch (topic) {
case "MktCreate":
return await handleMarketCreated(payload, meta);
case "BetPlace":
return await handleBetPlaced(payload, meta);
case "MktResolv":
return await handleMarketResolved(payload);
case "MktVoid":
return await handleMarketVoided(payload);
case "MktPause":
return await handleMarketPaused(payload);
case "Payout":
return await handlePayoutClaimed(payload, meta);
case "LpSeed":
return await handleLiquidityProvided(payload, meta);
case "LpClaim":
return await handleLpRewardClaimed(payload, meta);
case "Dispute":
return await handleDisputeRaised(payload, meta);
case "FeeColl":
return await handleFeeCollected(payload, meta);
default:
logger.warn({ topic }, "Unknown event topic β stored but not processed");
}
} catch (err) {
logger.error({ topic, tx_hash, err: err.message }, "Event handler failed");
throw err; // re-throw so the caller can dead-letter or retry
}
}
// ββ 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,
handleMarketResolved,
handleMarketVoided,
handleMarketPaused,
handlePayoutClaimed,
handleLiquidityProvided,
handleLpRewardClaimed,
handleDisputeRaised,
handleFeeCollected,
_publishOddsChanged,
};