import { defineConnector, raw, type NormalizedBatch } from "@market-atlas/connector-sdk"; import { parseTimestamp } from "@market-atlas/market-model"; import { pairObservations, planPair, splitVenueSymbol } from "../_shared/pairs.js"; const WS_URL = "wss://api.gemini.com/v2/marketdata"; const SYMBOLS = ["BTCUSD", "ETHUSD", "SOLUSD", "XRPUSD", "LTCUSD", "LINKUSD", "DOGEUSD", "AVAXUSD", "DOTUSD", "BTCGUSD"]; /** Gemini market data v2, `l2` subscription: initial snapshot with recent trades, then `trade` events. */ export const geminiWs = defineConnector({ metadata: { id: "gemini-ws", name: "Gemini — market data v2 (trades)", version: "1.0.0", sourceId: "gemini", organization: "Gemini Trust Company, LLC", sourceType: "WEBSOCKET", jurisdiction: "US", rightsStatus: "PUBLIC_ATTRIBUTED", realtimeStatus: "REALTIME", expectedLatencyMs: 400, supportsStreaming: true, supportsHistorical: false, assetClasses: ["CRYPTO"], exchanges: ["gemini"], homepage: "https://docs.gemini.com/websocket-api/#market-data-version-2", description: "Public Gemini market-data v2 feed (l2 subscription): last trades for the major USD pairs with exchange timestamps.", rightsNotes: "Public market data displayed with attribution to Gemini.", termsUrl: "https://www.gemini.com/legal/api-agreement", sourceFamily: "gemini", enabled: true, }, seeds: SYMBOLS.map((s) => { const [b, q] = splitVenueSymbol(s, null)!; return { symbol: s, hint: planPair(b, q, "gemini").hint, aliases: [`${b}-${q}`, `${b}/${q}`] }; }), defaultSymbols: SYMBOLS, async start(ctx) { const ws = ctx.openWebSocket(WS_URL, { label: "gemini", staleAfterMs: 180_000, heartbeat: { intervalMs: 25_000 }, onOpen: (sock) => sock.send({ type: "subscribe", subscriptions: [{ name: "l2", symbols: ctx.watchedSymbols() }] }), onMessage: (data) => { let msg: any; try { msg = JSON.parse(data); } catch { return; } if (msg?.type === "trade") ctx.emit(raw("gemini-ws", "gemini", "trade", msg)); else if (msg?.type === "l2_updates" && Array.isArray(msg.trades) && msg.trades.length) ctx.emit(raw("gemini-ws", "gemini", "trade", { ...msg.trades[msg.trades.length - 1], symbol: msg.symbol, type: "trade" })); }, }); ws.connect(); }, normalize(r): NormalizedBatch { const m = r.payload as Record; if (r.kind !== "trade" || typeof m.symbol !== "string") return { observations: [] }; const split = splitVenueSymbol(m.symbol, null); if (!split) return { observations: [] }; const [base, quote] = split; return { observations: pairObservations(m.symbol, base, quote, "gemini", { last: m.price }, { sourceTimestamp: parseTimestamp(m.timestamp), timestampTrust: "EXCHANGE", sequence: typeof m.event_id === "number" ? m.event_id : null, rightsStatus: "PUBLIC_ATTRIBUTED", realtimeStatus: "REALTIME", }), }; }, fixturesDir: "fixtures", });