From 707d6547493eb68965201ca83f1633bc0a191d9a Mon Sep 17 00:00:00 2001 From: direkturcrypto Date: Tue, 24 Feb 2026 13:40:09 +0700 Subject: [PATCH] fix(oneshot): rewrite MarketFeedService market discovery to match working detectors - Use /markets/slug/{slug} direct endpoint (not /markets?slug=...&limit=1) - Extract token IDs from clobTokenIds field with JSON string parsing fallback - Read tick size from market.orderPriceMinTickSize (no separate API call) - Prioritise endDate (full datetime) over endDateIso (date-only) to fix false expiry - Slot formula matches sniperDetector/mmDetector exactly Co-Authored-By: Claude Sonnet 4.6 --- src/oneshot/MarketFeedService.js | 198 +++++++++++++++++-------------- 1 file changed, 108 insertions(+), 90 deletions(-) diff --git a/src/oneshot/MarketFeedService.js b/src/oneshot/MarketFeedService.js index 9a28759..3dd1f25 100644 --- a/src/oneshot/MarketFeedService.js +++ b/src/oneshot/MarketFeedService.js @@ -9,21 +9,28 @@ * 4. Detect stale books (no levels, or fetch latency > STALE_THRESHOLD_MS) * 5. Emit 'snapshot' events on the event bus * + * Market discovery mirrors the logic in sniperDetector.js / mmDetector.js: + * - API endpoint: /markets/slug/{slug} (not /markets?slug=...) + * - Token IDs: clobTokenIds[0/1] (JSON string parsed if needed) + * - Tick size: market.orderPriceMinTickSize (no separate API call) + * - Slot formula: Math.floor(Date.now()/1000/SLOT_SEC) * SLOT_SEC + * * Snapshot shape: * { ts, marketSlug, conditionId, tteSec, tickSize, up: BookSide, down: BookSide, stale } * * BookSide shape: - * { tokenId, bids, asks, bestBid, bestAsk, mid, spread, depthBid, depthAsk } + * { tokenId, bids, asks, bestBid, bestAsk, mid, spread, depthBid, depthAsk, + * bestBidSize, bestAskSize } */ +import config from '../config/index.js'; import logger from '../utils/logger.js'; import { dbg, DEBUG } from './debug.js'; -const GAMMA_HOST = 'https://gamma-api.polymarket.com'; const STALE_THRESHOLD_MS = 1500; -const TOP_N_LEVELS = 5; // Levels counted for depth calculation +const TOP_N_LEVELS = 5; // Levels counted for depth calculation const DISCOVER_INTERVAL = 30_000; // Re-scan for new markets every 30s -const DEBUG_POLL_EVERY = 10; // Log a poll summary every N ticks per market (debug only) +const DEBUG_POLL_EVERY = 10; // Throttle: log one poll summary every N ticks export class MarketFeedService { /** @@ -35,21 +42,21 @@ export class MarketFeedService { * @param {import('./EventBus.js').default} opts.eventBus */ constructor({ client, assets, duration = '5m', pollIntervalMs = 300, eventBus }) { - this._client = client; - this._assets = assets; - this._duration = duration; - this._durationMin = duration === '15m' ? 15 : 5; - this._pollMs = pollIntervalMs; - this._eventBus = eventBus; + this._client = client; + this._assets = assets; + this._duration = duration; + this._slotSec = duration === '15m' ? 900 : 300; // same as sniperDetector/mmDetector + this._pollMs = pollIntervalMs; + this._eventBus = eventBus; /** @type {Map} slug → market record */ - this._markets = new Map(); + this._markets = new Map(); - this._pollTimer = null; + this._pollTimer = null; this._discoverTimer = null; /** Per-market tick counter for throttled debug logs */ - this._pollCount = new Map(); + this._pollCount = new Map(); } // ── Lifecycle ───────────────────────────────────────────────────────────── @@ -76,53 +83,63 @@ export class MarketFeedService { // ── Market discovery ────────────────────────────────────────────────────── async _discoverMarkets() { - const durationMin = this._durationMin; + // Probe current slot AND next upcoming slot (same as sniperDetector) + const curr = this._currentSlot(); + const next = curr + this._slotSec; + const slots = [curr, next]; - // Probe current slot and the immediately upcoming slot - const slots = [ - this._slotTs(), - this._slotTs() + durationMin * 60, - ]; - - dbg('FEED', `--- discovery cycle | probing ${this._assets.length * slots.length} slug(s) ---`); + dbg('FEED', `--- discovery cycle | curr=${curr} next=${next} | probing ${this._assets.length * 2} slug(s) ---`); for (const asset of this._assets) { for (const slotTs of slots) { const slug = `${asset}-updown-${this._duration}-${slotTs}`; if (this._markets.has(slug)) { - dbg('FEED', ` ${slug} → already tracked, skip`); + dbg('FEED', ` ${slug} → already tracked`); continue; } dbg('FEED', ` probing ${slug} ...`); try { - const market = await this._fetchMarketBySlug(slug); + // ── Use /markets/slug/{slug} — same endpoint as sniperDetector ── + const market = await this._fetchBySlug(slug); if (!market) { - dbg('FEED', ` ${slug} → not found on Gamma API`); + dbg('FEED', ` ${slug} → not found (API returned null)`); continue; } + // ── Extract end time ───────────────────────────────────────── + // endDate = "2026-02-24T06:35:00Z" (full datetime — use this) + // endDateIso = "2026-02-24" (date only, parses to midnight UTC — skip) const endTs = this._parseEndTs(market); if (!endTs) { - dbg('FEED', ` ${slug} → found but endTs unparseable`); + dbg('FEED', ` ${slug} → found but endDate unparseable (keys: ${Object.keys(market).slice(0, 8).join(',')})`); continue; } if (Date.now() >= endTs) { - dbg('FEED', ` ${slug} → found but already expired`); + dbg('FEED', ` ${slug} → found but expired (endTs=${new Date(endTs).toISOString()})`); continue; } + // ── Extract token IDs — same logic as sniperDetector/mmDetector ── const { upTokenId, downTokenId } = this._extractTokenIds(market); if (!upTokenId || !downTokenId) { - logger.warn(`MarketFeedService: could not extract token IDs for ${slug}`); - dbg('FEED', ` tokens shape: ${JSON.stringify(Object.keys(market).slice(0, 10))}`); + logger.warn(`MarketFeedService: missing token IDs for ${slug}`); + dbg('FEED', ` clobTokenIds raw: ${JSON.stringify(market.clobTokenIds)}`); continue; } - const tickSize = await this._fetchTickSize(upTokenId); + // ── Tick size from market object — same as mmDetector ──────── + const tickSize = parseFloat( + market.orderPriceMinTickSize ?? + market.minimum_tick_size ?? + market.minimumTickSize ?? + '0.01', + ) || 0.01; + + const negRisk = market.negRisk ?? market.neg_risk ?? false; this._markets.set(slug, { slug, @@ -131,95 +148,98 @@ export class MarketFeedService { downTokenId, endTs, tickSize, - negRisk: market.negRisk || market.neg_risk || false, + negRisk, }); const secLeft = Math.floor((endTs - Date.now()) / 1000); logger.success(`MarketFeedService: tracking ${slug} (closes in ${secLeft}s)`); - dbg('FEED', ` upToken=${upTokenId.slice(0, 12)}... downToken=${downTokenId.slice(0, 12)}... tick=${tickSize}`); + dbg('FEED', + ` up=${upTokenId.slice(0, 16)}... ` + + `down=${downTokenId.slice(0, 16)}... ` + + `tick=${tickSize} negRisk=${negRisk}`, + ); } catch (err) { - dbg('FEED', ` ${slug} → discovery error: ${err.message}`); - // Network blip — will retry on next discovery cycle + dbg('FEED', ` ${slug} → error: ${err.message}`); + // Network blip — will retry on next cycle } } } - // Prune expired markets (5s grace period for final snapshots) + // Prune markets that have fully expired (5s grace for final snapshots) for (const [slug, mkt] of this._markets) { if (Date.now() > mkt.endTs + 5_000) { this._markets.delete(slug); this._pollCount.delete(slug); - logger.info(`MarketFeedService: pruned expired market ${slug}`); + logger.info(`MarketFeedService: pruned ${slug}`); } } if (this._markets.size === 0) { - dbg('FEED', 'No active markets found — will retry in 30s'); + dbg('FEED', 'No active markets — retrying in 30s'); } else { - dbg('FEED', `Active markets: [${[...this._markets.keys()].join(', ')}]`); + dbg('FEED', `Tracking: [${[...this._markets.keys()].join(', ')}]`); } } - /** Deterministic UTC slot boundary timestamp (seconds) */ - _slotTs() { - const slotMs = this._durationMin * 60_000; - return Math.floor(Date.now() / slotMs) * slotMs / 1000; + // ── Slot helpers (identical to sniperDetector / mmDetector) ────────────── + + _currentSlot() { + return Math.floor(Date.now() / 1000 / this._slotSec) * this._slotSec; } - async _fetchMarketBySlug(slug) { - const url = `${GAMMA_HOST}/markets?slug=${encodeURIComponent(slug)}&limit=1`; - const resp = await fetch(url); + // ── Gamma API ───────────────────────────────────────────────────────────── + + /** Uses /markets/slug/{slug} — the same direct endpoint as sniperDetector */ + async _fetchBySlug(slug) { + const resp = await fetch(`${config.gammaHost}/markets/slug/${slug}`); if (!resp.ok) return null; const data = await resp.json(); - const market = Array.isArray(data) ? data[0] : data; - return market?.conditionId || market?.condition_id ? market : null; + // Returns a single object (not an array) when using the slug endpoint + return data?.conditionId || data?.condition_id ? data : null; } _parseEndTs(market) { - // endDate contains the full datetime (e.g. "2026-02-24T06:35:00Z"). - // endDateIso is date-only ("2026-02-24") and parses to midnight UTC — - // which is already in the past by market-open time, so it must come last. + // endDate = "2026-02-24T06:35:00Z" → correct full datetime + // endDateIso = "2026-02-24" → date-only, parses to midnight UTC (wrong!) const raw = market.endDate || market.end_date || market.endDateIso || market.end_date_iso; if (!raw) return null; const ts = new Date(raw).getTime(); return Number.isFinite(ts) ? ts : null; } + /** + * Extract UP/DOWN token IDs using the same logic as sniperDetector / mmDetector. + * + * clobTokenIds may be: + * - a real JS array: ["123...", "456..."] + * - a JSON string: '["123...","456..."]' + * UP = clobTokenIds[0] (YES / Up) + * DOWN = clobTokenIds[1] (NO / Down) + */ _extractTokenIds(market) { + let tokenIds = market.clobTokenIds ?? market.clob_token_ids; + + // Unwrap JSON string if the API returned it encoded + if (typeof tokenIds === 'string') { + try { tokenIds = JSON.parse(tokenIds); } catch { tokenIds = null; } + } + let upTokenId = null; let downTokenId = null; - // Shape 1: tokens[] array with { tokenId, outcome } - const tokens = market.tokens; - if (Array.isArray(tokens)) { - for (const t of tokens) { - const outcome = String(t.outcome || t.title || '').toLowerCase(); - const id = t.tokenId || t.token_id || t.id || t.asset; - if (!id) continue; - if (outcome.includes('up') || outcome === 'yes') upTokenId = String(id); - if (outcome.includes('down') || outcome === 'no') downTokenId = String(id); - } - } - - // Shape 2: clobTokenIds[0/1] - if ((!upTokenId || !downTokenId) && Array.isArray(market.clobTokenIds) && market.clobTokenIds.length >= 2) { - upTokenId = upTokenId ?? String(market.clobTokenIds[0]); - downTokenId = downTokenId ?? String(market.clobTokenIds[1]); + if (Array.isArray(tokenIds) && tokenIds.length >= 2) { + [upTokenId, downTokenId] = tokenIds.map(String); + } else if (Array.isArray(market.tokens) && market.tokens.length >= 2) { + // Fallback: named tokens array (less common) + upTokenId = String(market.tokens[0]?.token_id ?? market.tokens[0]?.tokenId ?? ''); + downTokenId = String(market.tokens[1]?.token_id ?? market.tokens[1]?.tokenId ?? ''); + if (!upTokenId || !downTokenId) { upTokenId = null; downTokenId = null; } } return { upTokenId, downTokenId }; } - async _fetchTickSize(tokenId) { - try { - const ts = await this._client.getTickSize(tokenId); - return parseFloat(ts) || 0.01; - } catch { - return 0.01; - } - } - // ── Book polling ────────────────────────────────────────────────────────── async _tick() { @@ -243,7 +263,7 @@ export class MarketFeedService { const snapshot = this._buildSnapshot(mkt, upBook, downBook, tteSec, stale); this._eventBus.emit('snapshot', snapshot); - // ── Debug: throttled poll summary (every N ticks) ───────────── + // ── Throttled debug poll summary ────────────────────────────── if (DEBUG) { const count = (this._pollCount.get(mkt.slug) ?? 0) + 1; this._pollCount.set(mkt.slug, count); @@ -251,15 +271,14 @@ export class MarketFeedService { if (count % DEBUG_POLL_EVERY === 1) { const u = snapshot.up; const d = snapshot.down; - const staleFlag = stale ? ' [STALE]' : ''; dbg('POLL', - `${mkt.slug} | tte=${tteSec}s | fetchMs=${fetchMs}ms${staleFlag}\n` + + `${mkt.slug} | tte=${tteSec}s | fetchMs=${fetchMs}ms${stale ? ' [STALE]' : ''}\n` + ` UP bid=${u.bestBid.toFixed(4)}/ask=${u.bestAsk.toFixed(4)} ` + - `spread=${u.spread.toFixed(4)} mid=${u.mid.toFixed(4)} ` + - `depthBid=${u.depthBid.toFixed(1)} depthAsk=${u.depthAsk.toFixed(1)}\n` + + `sprd=${u.spread.toFixed(4)} mid=${u.mid.toFixed(4)} ` + + `dBid=${u.depthBid.toFixed(1)} dAsk=${u.depthAsk.toFixed(1)}\n` + ` DOWN bid=${d.bestBid.toFixed(4)}/ask=${d.bestAsk.toFixed(4)} ` + - `spread=${d.spread.toFixed(4)} mid=${d.mid.toFixed(4)} ` + - `depthBid=${d.depthBid.toFixed(1)} depthAsk=${d.depthAsk.toFixed(1)}`, + `sprd=${d.spread.toFixed(4)} mid=${d.mid.toFixed(4)} ` + + `dBid=${d.depthBid.toFixed(1)} dAsk=${d.depthAsk.toFixed(1)}`, ); } } @@ -270,7 +289,11 @@ export class MarketFeedService { } } + // ── Snapshot builder ────────────────────────────────────────────────────── + _buildSnapshot(mkt, upBook, downBook, tteSec, stale) { + const up = this._buildSide(mkt.upTokenId, upBook); + const down = this._buildSide(mkt.downTokenId, downBook); return { ts: Date.now(), marketSlug: mkt.slug, @@ -278,14 +301,14 @@ export class MarketFeedService { tteSec, tickSize: mkt.tickSize, negRisk: mkt.negRisk, - up: this._buildSide(mkt.upTokenId, upBook), - down: this._buildSide(mkt.downTokenId, downBook), - stale: stale || this._isBooksEmpty(upBook, downBook), + up, + down, + stale: stale || up.bestBid === 0 || down.bestBid === 0, }; } _buildSide(tokenId, book) { - const parse = (raw = []) => + const parse = (raw = []) => (Array.isArray(raw) ? raw : []) .filter((l) => l?.price && l?.size) .map((l) => ({ price: parseFloat(l.price), size: parseFloat(l.size) })) @@ -319,9 +342,4 @@ export class MarketFeedService { bestAskSize: asks[0]?.size ?? 0, }; } - - _isBooksEmpty(upBook, downBook) { - const isEmpty = (b) => !b || (!Array.isArray(b.bids) && !Array.isArray(b.asks)); - return isEmpty(upBook) || isEmpty(downBook); - } }