diff --git a/docs/binance-orderbook-trade-development.md b/docs/binance-orderbook-trade-development.md index 0366dbc..57516ab 100644 --- a/docs/binance-orderbook-trade-development.md +++ b/docs/binance-orderbook-trade-development.md @@ -55,6 +55,7 @@ src/binance-orderbook-trade/ continuous-ladder.js decimal.js depth-profile-book.js + binance-native-depth-source.js depth-profile-session.js ladder-plan.js orderbook.js @@ -107,7 +108,7 @@ scripts/ The optional depth profile is a userscript-owned canvas beside the TradingView price axis. It does not clone Binance's native Depth React component. A feature-gated adapter reads the active TradingView main pane height and maps each depth price through the pane's current price scale, including logarithmic and inverted modes. The divider uses Binance's latest visible trade price rather than the order-book midpoint so it follows TradingView's live-price line. If that contract is unavailable or invalid, the profile fails closed instead of falling back to an approximate scale. Binance Basic and native Depth modes do not expose the verified coordinate contract, so the profile is hidden in those modes. -The data session opens the official USD-M public `{symbol}@depth@100ms` stream before requesting `/fapi/v1/depth?limit=1000`. `core/depth-profile-book.js` applies the official `lastUpdateId`, `U`, `u`, and `pu` sequence contract and treats quantities as absolute values; zero removes a price level. `core/depth-profile-session.js` owns one inflight snapshot, at most three resynchronizations per connection, and at most five reconnect attempts. A terminal failure remains visible until the user collapses and reopens the profile or changes symbols. +The userscript installs `core/binance-native-depth-source.js` at `document-start` and passively observes the native `/fapi/v1/rpiDepth?limit=1000` response and `{symbol}@rpiDepth@500ms` messages. It preserves Binance's original `fetch` result and WebSocket instances and never opens a second depth connection. `core/depth-profile-book.js` applies the observed `lastUpdateId`, `U`, `u`, and `pu` sequence contract and treats quantities as absolute values; zero removes a price level. `core/depth-profile-session.js` only subscribes the active symbol to that page-owned source. A sequence gap waits for Binance's native resynchronization instead of issuing a userscript-owned retry request, while a changed private RPI contract fails the profile explicitly without blocking Binance's own request. The session must stop and invalidate old work on symbol change, non-trading routes, hidden documents, and `pagehide`. The overlay canvas uses `pointer-events: none`; only its compact collapse control may receive pointer input. Do not connect this visualization book to ladder pricing or any trading decision. diff --git a/scripts/binance-orderbook-trade.user.js b/scripts/binance-orderbook-trade.user.js index b5d53c0..a231e1d 100644 --- a/scripts/binance-orderbook-trade.user.js +++ b/scripts/binance-orderbook-trade.user.js @@ -3,7 +3,7 @@ // @namespace binance.orderbook.trade // @icon data:image/svg+xml,%3Csvg%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%20viewBox%3D%220%200%2064%2064%22%3E%3Crect%20width%3D%2264%22%20height%3D%2264%22%20rx%3D%2214%22%20fill%3D%22%23f0b90b%22%2F%3E%3Ctext%20x%3D%2232%22%20y%3D%2249%22%20text-anchor%3D%22middle%22%20font-family%3D%22Arial%2C%20sans-serif%22%20font-size%3D%2242%22%20font-weight%3D%22800%22%20fill%3D%22%23111827%22%3EJ%3C%2Ftext%3E%3C%2Fsvg%3E // @icon64 data:image/svg+xml,%3Csvg%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%20viewBox%3D%220%200%2064%2064%22%3E%3Crect%20width%3D%2264%22%20height%3D%2264%22%20rx%3D%2214%22%20fill%3D%22%23f0b90b%22%2F%3E%3Ctext%20x%3D%2232%22%20y%3D%2249%22%20text-anchor%3D%22middle%22%20font-family%3D%22Arial%2C%20sans-serif%22%20font-size%3D%2242%22%20font-weight%3D%22800%22%20fill%3D%22%23111827%22%3EJ%3C%2Ftext%3E%3C%2Fsvg%3E -// @version 2.7.184 +// @version 2.7.185 // @author jackhai9 // @description 单击订单簿价格,按当前开仓/平仓 tab 自动填数量并执行下单,内置数量倍率面板 // @match https://www.binance.com/*/futures/* @@ -3598,194 +3598,294 @@ }; } - // src/binance-orderbook-trade/core/depth-profile-session.js - var MAX_RESYNCS_PER_CONNECTION = 3; - var MAX_RECONNECT_ATTEMPTS = 5; - var RESYNC_DELAY_MS = 1e3; - var RECONNECT_DELAYS_MS = [1e3, 2e3, 5e3, 1e4, 3e4]; + // src/binance-orderbook-trade/core/binance-native-depth-source.js + var NATIVE_RPI_DEPTH_PATH = "/fapi/v1/rpiDepth"; + var NATIVE_RPI_DEPTH_LIMIT = "1000"; + var NATIVE_RPI_STREAM_PATTERN = /^([a-z0-9_]+)@rpiDepth@500ms$/; function assertFunction(value, field) { - if (typeof value !== "function") throw new Error(`Invalid depth profile ${field}`); + if (typeof value !== "function") throw new Error(`Invalid native depth ${field}`); return value; } - function createDepthProfileSession(options) { - const { + function assertSymbol2(value) { + if (typeof value !== "string" || !/^[A-Z0-9_]+$/.test(value)) { + throw new Error("Invalid native depth symbol"); + } + return value; + } + function resolveRequestUrl(input, baseUrl) { + if (typeof input === "string") return new URL(input, baseUrl); + if (input && typeof input === "object" && typeof input.url === "string") { + return new URL(input.url, baseUrl); + } + return null; + } + function resolveNativeSnapshotSymbol(input, baseUrl) { + const url = resolveRequestUrl(input, baseUrl); + if (!url || url.pathname !== NATIVE_RPI_DEPTH_PATH) return null; + const symbol = assertSymbol2(url.searchParams.get("symbol")); + if (url.searchParams.get("limit") !== NATIVE_RPI_DEPTH_LIMIT) { + return { + symbol, + error: new Error("Invalid Binance native RPI depth snapshot limit") + }; + } + return { symbol, error: null }; + } + function createRecord(symbol) { + return { symbol, - fetchFn, - WebSocketCtor, - onProfile, - onStatus, - setTimer = setTimeout, - clearTimer = clearTimeout - } = options || {}; - if (typeof symbol !== "string" || !/^[A-Z0-9_]+$/.test(symbol)) { - throw new Error("Invalid depth profile session symbol"); + book: createDepthProfileBook(symbol), + profile: null, + status: { symbol, status: "connecting", detail: "" }, + subscribers: /* @__PURE__ */ new Set() + }; + } + function installBinanceNativeDepthSource(globalObject) { + if (!globalObject || typeof globalObject !== "object") { + throw new Error("Invalid native depth global object"); } - assertFunction(fetchFn, "fetch function"); - assertFunction(WebSocketCtor, "WebSocket constructor"); - assertFunction(onProfile, "profile listener"); - assertFunction(onStatus, "status listener"); - assertFunction(setTimer, "timer scheduler"); - assertFunction(clearTimer, "timer clearer"); - let active = false; - let epoch = 0; - let socket = null; - let snapshotController = null; - let reconnectTimer = 0; - let resyncTimer = 0; - let reconnectAttempts = 0; - let resyncAttempts = 0; - let book = createDepthProfileBook(symbol); - const snapshotUrl = `https://fapi.binance.com/fapi/v1/depth?symbol=${symbol}&limit=${DEPTH_PROFILE_LEVELS_PER_SIDE}`; - const streamUrl = `wss://fstream.binance.com/public/ws/${symbol.toLowerCase()}@depth@100ms`; - function emitStatus(status, detail = "") { - if (!active) return; - onStatus({ symbol, status, detail }); - } - function clearScheduledWork() { - if (reconnectTimer) clearTimer(reconnectTimer); - if (resyncTimer) clearTimer(resyncTimer); - reconnectTimer = 0; - resyncTimer = 0; - } - function abortSnapshot() { - snapshotController?.abort(); - snapshotController = null; - } - function closeSocket() { - const currentSocket = socket; - socket = null; - if (currentSocket && currentSocket.readyState < 2) currentSocket.close(); - } - function emitProfile(profile) { - reconnectAttempts = 0; - onProfile(profile); - emitStatus("ready"); - } - function failSession(error) { - clearScheduledWork(); - abortSnapshot(); - closeSocket(); - emitStatus("failed", error?.message || String(error)); - active = false; - epoch += 1; - } - function scheduleReconnect(error) { - if (!active || reconnectTimer) return; - abortSnapshot(); - closeSocket(); - if (reconnectAttempts >= MAX_RECONNECT_ATTEMPTS) { - failSession(error); - return; - } - const delay = RECONNECT_DELAYS_MS[reconnectAttempts]; - reconnectAttempts += 1; - emitStatus("reconnecting", error?.message || String(error)); - const scheduledEpoch = epoch; - reconnectTimer = setTimer(() => { - reconnectTimer = 0; - if (!active || scheduledEpoch !== epoch) return; - connect(); - }, delay); - } - function scheduleResync(error) { - if (!active || resyncTimer) return; - abortSnapshot(); - book = createDepthProfileBook(symbol); - if (resyncAttempts >= MAX_RESYNCS_PER_CONNECTION) { - scheduleReconnect(error); - return; - } - resyncAttempts += 1; - emitStatus("resyncing", error?.message || String(error)); - const scheduledEpoch = epoch; - resyncTimer = setTimer(() => { - resyncTimer = 0; - if (!active || scheduledEpoch !== epoch) return; - requestSnapshot(); - }, RESYNC_DELAY_MS); - } - async function requestSnapshot() { - if (!active || snapshotController) return; - const requestEpoch = epoch; - const controller = new AbortController(); - snapshotController = controller; - let profile = null; + const nativeFetch = assertFunction(globalObject.fetch, "fetch function"); + const NativeWebSocket = assertFunction(globalObject.WebSocket, "WebSocket constructor"); + const nativeSocketAddEventListener = assertFunction( + NativeWebSocket.prototype?.addEventListener, + "WebSocket event listener" + ); + const baseUrl = globalObject.location?.href; + if (typeof baseUrl !== "string") throw new Error("Invalid native depth page URL"); + const records = /* @__PURE__ */ new Map(); + const socketSymbols = /* @__PURE__ */ new WeakMap(); + let restored = false; + function ensureRecord(symbol) { + const normalizedSymbol = assertSymbol2(symbol); + let record = records.get(normalizedSymbol); + if (!record) { + record = createRecord(normalizedSymbol); + records.set(normalizedSymbol, record); + } + return record; + } + function notifyStatus(record, status, detail = "") { + record.status = { symbol: record.symbol, status, detail }; + for (const subscriber of record.subscribers) subscriber.onStatus(record.status); + } + function notifyProfile(record, profile) { + record.profile = profile; + for (const subscriber of record.subscribers) subscriber.onProfile(profile); + notifyStatus(record, "ready"); + } + function failRecord(record, error) { + record.profile = null; + notifyStatus(record, "failed", error?.message || String(error)); + } + function beginSnapshot(symbol) { + const record = ensureRecord(symbol); + record.book = createDepthProfileBook(symbol); + record.profile = null; + notifyStatus(record, "synchronizing"); + for (const [otherSymbol, otherRecord] of records) { + if (otherSymbol !== symbol && otherRecord.subscribers.size === 0) records.delete(otherSymbol); + } + return record; + } + function acceptSnapshot(symbol, payload) { + if (restored) return; + const record = ensureRecord(symbol); try { - const response = await fetchFn(snapshotUrl, { signal: controller.signal }); - if (!active || requestEpoch !== epoch || snapshotController !== controller) return; - if (!response.ok) throw new Error(`Depth snapshot HTTP ${response.status}`); - const payload = await response.json(); - if (!active || requestEpoch !== epoch || snapshotController !== controller) return; - const ready = applyDepthProfileSnapshot(book, payload); - snapshotController = null; - if (ready) profile = buildDepthProfile(book); + const ready = applyDepthProfileSnapshot(record.book, payload); + if (ready) notifyProfile(record, buildDepthProfile(record.book)); + else notifyStatus(record, "synchronizing"); } catch (error) { - if (snapshotController === controller) snapshotController = null; - if (controller.signal.aborted || !active || requestEpoch !== epoch) return; - scheduleResync(error); - return; + failRecord(record, error); } - if (profile) emitProfile(profile); - else emitStatus("synchronizing"); } - function handleMessage(currentSocket, event) { - if (!active || socket !== currentSocket) return; - let profile = null; + function acceptUpdate(symbol, payload) { + if (restored) return; + const record = ensureRecord(symbol); try { - const payload = JSON.parse(event.data); - const ready = pushDepthProfileUpdate(book, payload); - if (ready) profile = buildDepthProfile(book); + const ready = pushDepthProfileUpdate(record.book, payload); + if (ready) notifyProfile(record, buildDepthProfile(record.book)); + else notifyStatus(record, "synchronizing"); } catch (error) { - if (error instanceof DepthProfileSequenceError) scheduleResync(error); - else scheduleReconnect(error); - return; + record.profile = null; + if (error instanceof DepthProfileSequenceError) { + record.book = createDepthProfileBook(symbol); + notifyStatus(record, "resyncing", error.message); + } else { + failRecord(record, error); + } } - if (profile) emitProfile(profile); } - function connect() { - if (!active || socket) return; - book = createDepthProfileBook(symbol); - resyncAttempts = 0; - emitStatus("connecting"); - let currentSocket; - try { - currentSocket = new WebSocketCtor(streamUrl); - } catch (error) { - scheduleReconnect(error); - return; - } - socket = currentSocket; - currentSocket.addEventListener("open", () => { - if (!active || socket !== currentSocket) return; - emitStatus("synchronizing"); - requestSnapshot(); - }, { once: true }); - currentSocket.addEventListener("message", (event) => handleMessage(currentSocket, event)); - currentSocket.addEventListener("error", () => { - if (!active || socket !== currentSocket) return; - currentSocket.close(); - }); - currentSocket.addEventListener("close", () => { - if (!active || socket !== currentSocket) return; - socket = null; - scheduleReconnect(new Error("Depth WebSocket closed")); - }, { once: true }); + function failAll(error) { + for (const record of records.values()) failRecord(record, error); } + function observeSocket(socket) { + const symbols = /* @__PURE__ */ new Set(); + socketSymbols.set(socket, symbols); + Reflect.apply(nativeSocketAddEventListener, socket, ["message", (event) => { + if (restored || typeof event?.data !== "string" || !event.data.includes("@rpiDepth@500ms")) { + return; + } + let envelope; + try { + envelope = JSON.parse(event.data); + const match = NATIVE_RPI_STREAM_PATTERN.exec(envelope?.stream); + if (!match || !envelope.data || typeof envelope.data !== "object") { + throw new Error("Invalid Binance native RPI depth message"); + } + const symbol = assertSymbol2(match[1].toUpperCase()); + symbols.add(symbol); + acceptUpdate(symbol, envelope.data); + } catch (error) { + failAll(error); + } + }]); + Reflect.apply(nativeSocketAddEventListener, socket, ["close", () => { + if (restored) return; + for (const symbol of socketSymbols.get(socket) || []) { + const record = records.get(symbol); + if (record) notifyStatus(record, "reconnecting", "Binance native depth socket closed"); + } + }]); + } + const observedFetch = new Proxy(nativeFetch, { + apply(target, receiver, args) { + let observation = null; + try { + observation = resolveNativeSnapshotSymbol(args[0], baseUrl); + } catch (error) { + failAll(error); + } + const symbol = observation?.symbol || null; + if (observation?.error) failRecord(ensureRecord(symbol), observation.error); + else if (symbol) beginSnapshot(symbol); + let result; + try { + result = Reflect.apply(target, receiver, args); + } catch (error) { + if (symbol) failRecord(ensureRecord(symbol), error); + throw error; + } + if (symbol && !observation.error) { + Promise.resolve(result).then((response) => { + if (!response?.ok) { + throw new Error(`Binance native RPI depth snapshot HTTP ${response?.status}`); + } + return response.clone().json(); + }).then( + (payload) => acceptSnapshot(symbol, payload), + (error) => failRecord(ensureRecord(symbol), error) + ); + } + return result; + } + }); + const ObservedWebSocket = new Proxy(NativeWebSocket, { + construct(target, args, newTarget) { + const socket = Reflect.construct(target, args, newTarget); + observeSocket(socket); + return socket; + } + }); + globalObject.fetch = observedFetch; + globalObject.WebSocket = ObservedWebSocket; return { + subscribe(options) { + const { + symbol, + onProfile, + onStatus + } = options || {}; + const record = ensureRecord(symbol); + const subscriber = { + onProfile: assertFunction(onProfile, "profile listener"), + onStatus: assertFunction(onStatus, "status listener") + }; + if (restored) throw new Error("Binance native depth source has been restored"); + record.subscribers.add(subscriber); + subscriber.onStatus(record.status); + if (record.profile) subscriber.onProfile(record.profile); + return () => record.subscribers.delete(subscriber); + }, + getState(symbol) { + const record = records.get(assertSymbol2(symbol)); + if (!record) return null; + return { + status: record.status, + bidCount: record.profile?.bids.length || 0, + askCount: record.profile?.asks.length || 0, + minPrice: record.profile?.minPrice ?? null, + maxPrice: record.profile?.maxPrice ?? null + }; + }, + restore() { + if (restored) return; + restored = true; + if (globalObject.fetch === observedFetch) globalObject.fetch = nativeFetch; + if (globalObject.WebSocket === ObservedWebSocket) globalObject.WebSocket = NativeWebSocket; + records.clear(); + } + }; + } + + // src/binance-orderbook-trade/core/depth-profile-session.js + function assertFunction2(value, field) { + if (typeof value !== "function") throw new Error(`Invalid depth profile ${field}`); + return value; + } + function assertSymbol3(value) { + if (typeof value !== "string" || !/^[A-Z0-9_]+$/.test(value)) { + throw new Error("Invalid depth profile session symbol"); + } + return value; + } + function createDepthProfileSession(options) { + const { symbol, + source, + onProfile, + onStatus + } = options || {}; + const normalizedSymbol = assertSymbol3(symbol); + if (!source || typeof source !== "object") throw new Error("Invalid depth profile source"); + const subscribe = assertFunction2(source.subscribe, "source subscriber"); + const profileListener = assertFunction2(onProfile, "profile listener"); + const statusListener = assertFunction2(onStatus, "status listener"); + let active = false; + let unsubscribe = null; + function assertMatchingSymbol(value, field) { + if (value?.symbol !== normalizedSymbol) { + throw new Error( + `Depth profile ${field} symbol mismatch: expected ${normalizedSymbol}, received ${value?.symbol}` + ); + } + } + return { + symbol: normalizedSymbol, start() { if (active) throw new Error("Depth profile session already started"); active = true; - epoch += 1; - connect(); + statusListener({ symbol: normalizedSymbol, status: "connecting", detail: "" }); + unsubscribe = Reflect.apply(subscribe, source, [{ + symbol: normalizedSymbol, + onProfile(profile) { + if (!active) return; + assertMatchingSymbol(profile, "profile"); + profileListener(profile); + }, + onStatus(status) { + if (!active) return; + assertMatchingSymbol(status, "status"); + statusListener(status); + } + }]); + assertFunction2(unsubscribe, "unsubscribe function"); }, stop() { if (!active) return; active = false; - epoch += 1; - clearScheduledWork(); - abortSnapshot(); - closeSocket(); + const stopSubscription = unsubscribe; + unsubscribe = null; + stopSubscription?.(); }, isActive() { return active; @@ -4655,6 +4755,7 @@ return isFuturesTradingPathname(location.pathname); } if (!isFuturesTradingPage()) return; + const nativeDepthSource = installBinanceNativeDepthSource(window); const CFG = { // true=只填数量;false=填数量并自动点“开多/开空/平多/平空” SAFE_MODE: false, @@ -11305,8 +11406,7 @@ let session = null; session = createDepthProfileSession({ symbol, - fetchFn: window.fetch.bind(window), - WebSocketCtor: window.WebSocket, + source: nativeDepthSource, onProfile(profile) { if (depthProfileSession !== session || getCurrentSymbol() !== symbol) return; depthProfileData = profile; @@ -11922,6 +12022,10 @@ get continuousChartSaveStats() { return continuousChartSaveController?.getStats() || null; }, + get nativeDepthState() { + const symbol = getCurrentSymbol(); + return symbol ? nativeDepthSource.getState(symbol) : null; + }, get cachedCloseState() { return getCachedCloseState(getCurrentSymbol()); }, diff --git a/src/binance-orderbook-trade/core/binance-native-depth-source.js b/src/binance-orderbook-trade/core/binance-native-depth-source.js new file mode 100644 index 0000000..10e2b55 --- /dev/null +++ b/src/binance-orderbook-trade/core/binance-native-depth-source.js @@ -0,0 +1,259 @@ +import { + applyDepthProfileSnapshot, + buildDepthProfile, + createDepthProfileBook, + DepthProfileSequenceError, + pushDepthProfileUpdate, +} from './depth-profile-book.js'; + +const NATIVE_RPI_DEPTH_PATH = '/fapi/v1/rpiDepth'; +const NATIVE_RPI_DEPTH_LIMIT = '1000'; +const NATIVE_RPI_STREAM_PATTERN = /^([a-z0-9_]+)@rpiDepth@500ms$/; + +function assertFunction(value, field) { + if (typeof value !== 'function') throw new Error(`Invalid native depth ${field}`); + return value; +} + +function assertSymbol(value) { + if (typeof value !== 'string' || !/^[A-Z0-9_]+$/.test(value)) { + throw new Error('Invalid native depth symbol'); + } + return value; +} + +function resolveRequestUrl(input, baseUrl) { + if (typeof input === 'string') return new URL(input, baseUrl); + if (input && typeof input === 'object' && typeof input.url === 'string') { + return new URL(input.url, baseUrl); + } + return null; +} + +function resolveNativeSnapshotSymbol(input, baseUrl) { + const url = resolveRequestUrl(input, baseUrl); + if (!url || url.pathname !== NATIVE_RPI_DEPTH_PATH) return null; + const symbol = assertSymbol(url.searchParams.get('symbol')); + if (url.searchParams.get('limit') !== NATIVE_RPI_DEPTH_LIMIT) { + return { + symbol, + error: new Error('Invalid Binance native RPI depth snapshot limit'), + }; + } + return { symbol, error: null }; +} + +function createRecord(symbol) { + return { + symbol, + book: createDepthProfileBook(symbol), + profile: null, + status: { symbol, status: 'connecting', detail: '' }, + subscribers: new Set(), + }; +} + +/** + * Observes Binance's existing RPI depth transport. The wrappers preserve the page's + * fetch result and WebSocket instances; this module never opens a network connection. + */ +export function installBinanceNativeDepthSource(globalObject) { + if (!globalObject || typeof globalObject !== 'object') { + throw new Error('Invalid native depth global object'); + } + const nativeFetch = assertFunction(globalObject.fetch, 'fetch function'); + const NativeWebSocket = assertFunction(globalObject.WebSocket, 'WebSocket constructor'); + const nativeSocketAddEventListener = assertFunction( + NativeWebSocket.prototype?.addEventListener, + 'WebSocket event listener', + ); + const baseUrl = globalObject.location?.href; + if (typeof baseUrl !== 'string') throw new Error('Invalid native depth page URL'); + + const records = new Map(); + const socketSymbols = new WeakMap(); + let restored = false; + + function ensureRecord(symbol) { + const normalizedSymbol = assertSymbol(symbol); + let record = records.get(normalizedSymbol); + if (!record) { + record = createRecord(normalizedSymbol); + records.set(normalizedSymbol, record); + } + return record; + } + + function notifyStatus(record, status, detail = '') { + record.status = { symbol: record.symbol, status, detail }; + for (const subscriber of record.subscribers) subscriber.onStatus(record.status); + } + + function notifyProfile(record, profile) { + record.profile = profile; + for (const subscriber of record.subscribers) subscriber.onProfile(profile); + notifyStatus(record, 'ready'); + } + + function failRecord(record, error) { + record.profile = null; + notifyStatus(record, 'failed', error?.message || String(error)); + } + + function beginSnapshot(symbol) { + const record = ensureRecord(symbol); + record.book = createDepthProfileBook(symbol); + record.profile = null; + notifyStatus(record, 'synchronizing'); + for (const [otherSymbol, otherRecord] of records) { + if (otherSymbol !== symbol && otherRecord.subscribers.size === 0) records.delete(otherSymbol); + } + return record; + } + + function acceptSnapshot(symbol, payload) { + if (restored) return; + const record = ensureRecord(symbol); + try { + const ready = applyDepthProfileSnapshot(record.book, payload); + if (ready) notifyProfile(record, buildDepthProfile(record.book)); + else notifyStatus(record, 'synchronizing'); + } catch (error) { + failRecord(record, error); + } + } + + function acceptUpdate(symbol, payload) { + if (restored) return; + const record = ensureRecord(symbol); + try { + const ready = pushDepthProfileUpdate(record.book, payload); + if (ready) notifyProfile(record, buildDepthProfile(record.book)); + else notifyStatus(record, 'synchronizing'); + } catch (error) { + record.profile = null; + if (error instanceof DepthProfileSequenceError) { + record.book = createDepthProfileBook(symbol); + notifyStatus(record, 'resyncing', error.message); + } else { + failRecord(record, error); + } + } + } + + function failAll(error) { + for (const record of records.values()) failRecord(record, error); + } + + function observeSocket(socket) { + const symbols = new Set(); + socketSymbols.set(socket, symbols); + Reflect.apply(nativeSocketAddEventListener, socket, ['message', (event) => { + if (restored || typeof event?.data !== 'string' || !event.data.includes('@rpiDepth@500ms')) { + return; + } + let envelope; + try { + envelope = JSON.parse(event.data); + const match = NATIVE_RPI_STREAM_PATTERN.exec(envelope?.stream); + if (!match || !envelope.data || typeof envelope.data !== 'object') { + throw new Error('Invalid Binance native RPI depth message'); + } + const symbol = assertSymbol(match[1].toUpperCase()); + symbols.add(symbol); + acceptUpdate(symbol, envelope.data); + } catch (error) { + failAll(error); + } + }]); + Reflect.apply(nativeSocketAddEventListener, socket, ['close', () => { + if (restored) return; + for (const symbol of socketSymbols.get(socket) || []) { + const record = records.get(symbol); + if (record) notifyStatus(record, 'reconnecting', 'Binance native depth socket closed'); + } + }]); + } + + const observedFetch = new Proxy(nativeFetch, { + apply(target, receiver, args) { + let observation = null; + try { + observation = resolveNativeSnapshotSymbol(args[0], baseUrl); + } catch (error) { + failAll(error); + } + const symbol = observation?.symbol || null; + if (observation?.error) failRecord(ensureRecord(symbol), observation.error); + else if (symbol) beginSnapshot(symbol); + let result; + try { + result = Reflect.apply(target, receiver, args); + } catch (error) { + if (symbol) failRecord(ensureRecord(symbol), error); + throw error; + } + if (symbol && !observation.error) { + Promise.resolve(result).then((response) => { + if (!response?.ok) { + throw new Error(`Binance native RPI depth snapshot HTTP ${response?.status}`); + } + return response.clone().json(); + }).then( + (payload) => acceptSnapshot(symbol, payload), + (error) => failRecord(ensureRecord(symbol), error), + ); + } + return result; + }, + }); + + const ObservedWebSocket = new Proxy(NativeWebSocket, { + construct(target, args, newTarget) { + const socket = Reflect.construct(target, args, newTarget); + observeSocket(socket); + return socket; + }, + }); + + globalObject.fetch = observedFetch; + globalObject.WebSocket = ObservedWebSocket; + + return { + subscribe(options) { + const { + symbol, + onProfile, + onStatus, + } = options || {}; + const record = ensureRecord(symbol); + const subscriber = { + onProfile: assertFunction(onProfile, 'profile listener'), + onStatus: assertFunction(onStatus, 'status listener'), + }; + if (restored) throw new Error('Binance native depth source has been restored'); + record.subscribers.add(subscriber); + subscriber.onStatus(record.status); + if (record.profile) subscriber.onProfile(record.profile); + return () => record.subscribers.delete(subscriber); + }, + getState(symbol) { + const record = records.get(assertSymbol(symbol)); + if (!record) return null; + return { + status: record.status, + bidCount: record.profile?.bids.length || 0, + askCount: record.profile?.asks.length || 0, + minPrice: record.profile?.minPrice ?? null, + maxPrice: record.profile?.maxPrice ?? null, + }; + }, + restore() { + if (restored) return; + restored = true; + if (globalObject.fetch === observedFetch) globalObject.fetch = nativeFetch; + if (globalObject.WebSocket === ObservedWebSocket) globalObject.WebSocket = NativeWebSocket; + records.clear(); + }, + }; +} diff --git a/src/binance-orderbook-trade/core/depth-profile-session.js b/src/binance-orderbook-trade/core/depth-profile-session.js index 391227e..1630047 100644 --- a/src/binance-orderbook-trade/core/depth-profile-session.js +++ b/src/binance-orderbook-trade/core/depth-profile-session.js @@ -1,215 +1,66 @@ -import { - applyDepthProfileSnapshot, - buildDepthProfile, - createDepthProfileBook, - DEPTH_PROFILE_LEVELS_PER_SIDE, - DepthProfileSequenceError, - pushDepthProfileUpdate, -} from './depth-profile-book.js'; - -const MAX_RESYNCS_PER_CONNECTION = 3; -const MAX_RECONNECT_ATTEMPTS = 5; -const RESYNC_DELAY_MS = 1000; -const RECONNECT_DELAYS_MS = [1000, 2000, 5000, 10000, 30000]; - function assertFunction(value, field) { if (typeof value !== 'function') throw new Error(`Invalid depth profile ${field}`); return value; } +function assertSymbol(value) { + if (typeof value !== 'string' || !/^[A-Z0-9_]+$/.test(value)) { + throw new Error('Invalid depth profile session symbol'); + } + return value; +} + export function createDepthProfileSession(options) { const { symbol, - fetchFn, - WebSocketCtor, + source, onProfile, onStatus, - setTimer = setTimeout, - clearTimer = clearTimeout, } = options || {}; - if (typeof symbol !== 'string' || !/^[A-Z0-9_]+$/.test(symbol)) { - throw new Error('Invalid depth profile session symbol'); - } - assertFunction(fetchFn, 'fetch function'); - assertFunction(WebSocketCtor, 'WebSocket constructor'); - assertFunction(onProfile, 'profile listener'); - assertFunction(onStatus, 'status listener'); - assertFunction(setTimer, 'timer scheduler'); - assertFunction(clearTimer, 'timer clearer'); + const normalizedSymbol = assertSymbol(symbol); + if (!source || typeof source !== 'object') throw new Error('Invalid depth profile source'); + const subscribe = assertFunction(source.subscribe, 'source subscriber'); + const profileListener = assertFunction(onProfile, 'profile listener'); + const statusListener = assertFunction(onStatus, 'status listener'); let active = false; - let epoch = 0; - let socket = null; - let snapshotController = null; - let reconnectTimer = 0; - let resyncTimer = 0; - let reconnectAttempts = 0; - let resyncAttempts = 0; - let book = createDepthProfileBook(symbol); - - const snapshotUrl = `https://fapi.binance.com/fapi/v1/depth?symbol=${symbol}&limit=${DEPTH_PROFILE_LEVELS_PER_SIDE}`; - const streamUrl = `wss://fstream.binance.com/public/ws/${symbol.toLowerCase()}@depth@100ms`; - - function emitStatus(status, detail = '') { - if (!active) return; - onStatus({ symbol, status, detail }); - } - - function clearScheduledWork() { - if (reconnectTimer) clearTimer(reconnectTimer); - if (resyncTimer) clearTimer(resyncTimer); - reconnectTimer = 0; - resyncTimer = 0; - } - - function abortSnapshot() { - snapshotController?.abort(); - snapshotController = null; - } - - function closeSocket() { - const currentSocket = socket; - socket = null; - if (currentSocket && currentSocket.readyState < 2) currentSocket.close(); - } + let unsubscribe = null; - function emitProfile(profile) { - reconnectAttempts = 0; - onProfile(profile); - emitStatus('ready'); - } - - function failSession(error) { - clearScheduledWork(); - abortSnapshot(); - closeSocket(); - emitStatus('failed', error?.message || String(error)); - active = false; - epoch += 1; - } - - function scheduleReconnect(error) { - if (!active || reconnectTimer) return; - abortSnapshot(); - closeSocket(); - if (reconnectAttempts >= MAX_RECONNECT_ATTEMPTS) { - failSession(error); - return; + function assertMatchingSymbol(value, field) { + if (value?.symbol !== normalizedSymbol) { + throw new Error( + `Depth profile ${field} symbol mismatch: expected ${normalizedSymbol}, received ${value?.symbol}`, + ); } - const delay = RECONNECT_DELAYS_MS[reconnectAttempts]; - reconnectAttempts += 1; - emitStatus('reconnecting', error?.message || String(error)); - const scheduledEpoch = epoch; - reconnectTimer = setTimer(() => { - reconnectTimer = 0; - if (!active || scheduledEpoch !== epoch) return; - connect(); - }, delay); - } - - function scheduleResync(error) { - if (!active || resyncTimer) return; - abortSnapshot(); - book = createDepthProfileBook(symbol); - if (resyncAttempts >= MAX_RESYNCS_PER_CONNECTION) { - scheduleReconnect(error); - return; - } - resyncAttempts += 1; - emitStatus('resyncing', error?.message || String(error)); - const scheduledEpoch = epoch; - resyncTimer = setTimer(() => { - resyncTimer = 0; - if (!active || scheduledEpoch !== epoch) return; - requestSnapshot(); - }, RESYNC_DELAY_MS); - } - - async function requestSnapshot() { - if (!active || snapshotController) return; - const requestEpoch = epoch; - const controller = new AbortController(); - snapshotController = controller; - let profile = null; - try { - const response = await fetchFn(snapshotUrl, { signal: controller.signal }); - if (!active || requestEpoch !== epoch || snapshotController !== controller) return; - if (!response.ok) throw new Error(`Depth snapshot HTTP ${response.status}`); - const payload = await response.json(); - if (!active || requestEpoch !== epoch || snapshotController !== controller) return; - const ready = applyDepthProfileSnapshot(book, payload); - snapshotController = null; - if (ready) profile = buildDepthProfile(book); - } catch (error) { - if (snapshotController === controller) snapshotController = null; - if (controller.signal.aborted || !active || requestEpoch !== epoch) return; - scheduleResync(error); - return; - } - if (profile) emitProfile(profile); - else emitStatus('synchronizing'); - } - - function handleMessage(currentSocket, event) { - if (!active || socket !== currentSocket) return; - let profile = null; - try { - const payload = JSON.parse(event.data); - const ready = pushDepthProfileUpdate(book, payload); - if (ready) profile = buildDepthProfile(book); - } catch (error) { - if (error instanceof DepthProfileSequenceError) scheduleResync(error); - else scheduleReconnect(error); - return; - } - if (profile) emitProfile(profile); - } - - function connect() { - if (!active || socket) return; - book = createDepthProfileBook(symbol); - resyncAttempts = 0; - emitStatus('connecting'); - let currentSocket; - try { - currentSocket = new WebSocketCtor(streamUrl); - } catch (error) { - scheduleReconnect(error); - return; - } - socket = currentSocket; - currentSocket.addEventListener('open', () => { - if (!active || socket !== currentSocket) return; - emitStatus('synchronizing'); - requestSnapshot(); - }, { once: true }); - currentSocket.addEventListener('message', (event) => handleMessage(currentSocket, event)); - currentSocket.addEventListener('error', () => { - if (!active || socket !== currentSocket) return; - currentSocket.close(); - }); - currentSocket.addEventListener('close', () => { - if (!active || socket !== currentSocket) return; - socket = null; - scheduleReconnect(new Error('Depth WebSocket closed')); - }, { once: true }); } return { - symbol, + symbol: normalizedSymbol, start() { if (active) throw new Error('Depth profile session already started'); active = true; - epoch += 1; - connect(); + statusListener({ symbol: normalizedSymbol, status: 'connecting', detail: '' }); + unsubscribe = Reflect.apply(subscribe, source, [{ + symbol: normalizedSymbol, + onProfile(profile) { + if (!active) return; + assertMatchingSymbol(profile, 'profile'); + profileListener(profile); + }, + onStatus(status) { + if (!active) return; + assertMatchingSymbol(status, 'status'); + statusListener(status); + }, + }]); + assertFunction(unsubscribe, 'unsubscribe function'); }, stop() { if (!active) return; active = false; - epoch += 1; - clearScheduledWork(); - abortSnapshot(); - closeSocket(); + const stopSubscription = unsubscribe; + unsubscribe = null; + stopSubscription?.(); }, isActive() { return active; diff --git a/src/binance-orderbook-trade/index.user.js b/src/binance-orderbook-trade/index.user.js index 3edb9f8..4b68dbd 100644 --- a/src/binance-orderbook-trade/index.user.js +++ b/src/binance-orderbook-trade/index.user.js @@ -3,7 +3,7 @@ // @namespace binance.orderbook.trade // @icon data:image/svg+xml,%3Csvg%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%20viewBox%3D%220%200%2064%2064%22%3E%3Crect%20width%3D%2264%22%20height%3D%2264%22%20rx%3D%2214%22%20fill%3D%22%23f0b90b%22%2F%3E%3Ctext%20x%3D%2232%22%20y%3D%2249%22%20text-anchor%3D%22middle%22%20font-family%3D%22Arial%2C%20sans-serif%22%20font-size%3D%2242%22%20font-weight%3D%22800%22%20fill%3D%22%23111827%22%3EJ%3C%2Ftext%3E%3C%2Fsvg%3E // @icon64 data:image/svg+xml,%3Csvg%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%20viewBox%3D%220%200%2064%2064%22%3E%3Crect%20width%3D%2264%22%20height%3D%2264%22%20rx%3D%2214%22%20fill%3D%22%23f0b90b%22%2F%3E%3Ctext%20x%3D%2232%22%20y%3D%2249%22%20text-anchor%3D%22middle%22%20font-family%3D%22Arial%2C%20sans-serif%22%20font-size%3D%2242%22%20font-weight%3D%22800%22%20fill%3D%22%23111827%22%3EJ%3C%2Ftext%3E%3C%2Fsvg%3E -// @version 2.7.184 +// @version 2.7.185 // @author jackhai9 // @description 单击订单簿价格,按当前开仓/平仓 tab 自动填数量并执行下单,内置数量倍率面板 // @match https://www.binance.com/*/futures/* @@ -204,6 +204,7 @@ import { createTradingViewRemovalSaveController, } from './core/chart-save-coalescer.js'; import { findBinanceTradingViewTarget } from './dom/tradingview-target.js'; +import { installBinanceNativeDepthSource } from './core/binance-native-depth-source.js'; import { createDepthProfileSession } from './core/depth-profile-session.js'; import { clearDepthProfile, @@ -244,6 +245,8 @@ import { showUsdtRebalanceDialog } from './dom/usdt-rebalance-dialog.js'; if (!isFuturesTradingPage()) return; + const nativeDepthSource = installBinanceNativeDepthSource(window); + const CFG = { // true=只填数量;false=填数量并自动点“开多/开空/平多/平空” SAFE_MODE: false, @@ -8024,8 +8027,7 @@ import { showUsdtRebalanceDialog } from './dom/usdt-rebalance-dialog.js'; let session = null; session = createDepthProfileSession({ symbol, - fetchFn: window.fetch.bind(window), - WebSocketCtor: window.WebSocket, + source: nativeDepthSource, onProfile(profile) { if (depthProfileSession !== session || getCurrentSymbol() !== symbol) return; depthProfileData = profile; @@ -8726,6 +8728,10 @@ import { showUsdtRebalanceDialog } from './dom/usdt-rebalance-dialog.js'; get continuousChartSaveStats() { return continuousChartSaveController?.getStats() || null; }, + get nativeDepthState() { + const symbol = getCurrentSymbol(); + return symbol ? nativeDepthSource.getState(symbol) : null; + }, get cachedCloseState() { return getCachedCloseState(getCurrentSymbol()); }, get displayCloseState() { return lastDisplayCloseState; }, get closeGuard() { return closeGuard; }, diff --git a/test/unit/binance-orderbook-trade/binance-native-depth-source.test.js b/test/unit/binance-orderbook-trade/binance-native-depth-source.test.js new file mode 100644 index 0000000..f53c12f --- /dev/null +++ b/test/unit/binance-orderbook-trade/binance-native-depth-source.test.js @@ -0,0 +1,280 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; + +import { installBinanceNativeDepthSource } from '../../../src/binance-orderbook-trade/core/binance-native-depth-source.js'; + +class FakeResponse { + constructor(payload, { ok = true, status = 200 } = {}) { + this.payload = payload; + this.ok = ok; + this.status = status; + } + + clone() { + return new FakeResponse(this.payload, { ok: this.ok, status: this.status }); + } + + async json() { + return this.payload; + } +} + +class FakeNativeSocket { + static CONNECTING = 0; + static OPEN = 1; + static CLOSING = 2; + static CLOSED = 3; + static instances = []; + + constructor(...args) { + this.args = args; + this.listeners = new Map(); + FakeNativeSocket.instances.push(this); + } + + addEventListener(type, listener) { + const listeners = this.listeners.get(type) || []; + listeners.push(listener); + this.listeners.set(type, listeners); + } + + emit(type, event = {}) { + for (const listener of this.listeners.get(type) || []) listener(event); + } + + message(payload) { + this.emit('message', { data: JSON.stringify(payload) }); + } +} + +function snapshot(overrides = {}) { + return { + lastUpdateId: 101, + bids: [['100', '1'], ['99', '2']], + asks: [['101', '3'], ['102', '4']], + ...overrides, + }; +} + +function update(overrides = {}) { + return { + e: 'depthUpdate', + s: 'BTCUSDT', + st: 1, + U: 100, + u: 102, + pu: 99, + b: [['100', '2'], ['99', '3']], + a: [['101', '4'], ['102', '5']], + ...overrides, + }; +} + +function rpiMessage(payload = update()) { + return { stream: 'btcusdt@rpiDepth@500ms', data: payload }; +} + +function createHarness(response = new FakeResponse(snapshot())) { + FakeNativeSocket.instances = []; + const fetchCalls = []; + const originalFetch = function originalFetch(...args) { + const result = Promise.resolve(typeof response === 'function' ? response() : response); + fetchCalls.push({ receiver: this, args, result }); + return result; + }; + const globalObject = { + fetch: originalFetch, + WebSocket: FakeNativeSocket, + location: { href: 'https://www.binance.com/zh-CN/futures/BTCUSDT' }, + }; + const source = installBinanceNativeDepthSource(globalObject); + return { fetchCalls, globalObject, originalFetch, source }; +} + +test('observes the native RPI snapshot and stream without creating another socket', async () => { + const { fetchCalls, globalObject, source } = createHarness(); + const profiles = []; + const statuses = []; + source.subscribe({ + symbol: 'BTCUSDT', + onProfile: (profile) => profiles.push(profile), + onStatus: (status) => statuses.push(status), + }); + + const nativeSocket = new globalObject.WebSocket('wss://native-binance-stream.example/ws'); + const responsePromise = globalObject.fetch('/fapi/v1/rpiDepth?symbol=BTCUSDT&limit=1000'); + nativeSocket.message(rpiMessage()); + const response = await responsePromise; + await new Promise((resolve) => setImmediate(resolve)); + + assert.equal(response.ok, true); + assert.equal(FakeNativeSocket.instances.length, 1); + assert.equal(nativeSocket.args[0], 'wss://native-binance-stream.example/ws'); + assert.equal(fetchCalls.length, 1); + assert.equal(profiles.length, 1, JSON.stringify({ statuses, state: source.getState('BTCUSDT') })); + assert.equal(profiles[0].symbol, 'BTCUSDT'); + assert.equal(profiles[0].bids.at(-1).price, 99); + assert.equal(profiles[0].asks.at(-1).price, 102); + assert.equal(statuses.at(-1).status, 'ready'); + + source.restore(); +}); + +test('preserves native fetch, WebSocket prototype, instanceof, and static constants', async () => { + const { + fetchCalls, + globalObject, + originalFetch, + source, + } = createHarness(); + const wrappedFetch = globalObject.fetch; + const WrappedWebSocket = globalObject.WebSocket; + const receiver = { marker: 'receiver' }; + + const fetchResult = Reflect.apply(wrappedFetch, receiver, ['/unrelated']); + await fetchResult; + const socket = new WrappedWebSocket('wss://native-binance-stream.example/ws', ['json']); + + assert.equal(socket instanceof FakeNativeSocket, true); + assert.equal(socket instanceof WrappedWebSocket, true); + assert.equal(wrappedFetch.name, originalFetch.name); + assert.equal(wrappedFetch.length, originalFetch.length); + assert.equal(WrappedWebSocket.name, FakeNativeSocket.name); + assert.equal(WrappedWebSocket.OPEN, FakeNativeSocket.OPEN); + assert.deepEqual(socket.args, ['wss://native-binance-stream.example/ws', ['json']]); + assert.equal(fetchResult, fetchCalls[0].result); + + source.restore(); + assert.equal(globalObject.fetch, originalFetch); + assert.equal(globalObject.WebSocket, FakeNativeSocket); +}); + +test('reports a changed native snapshot contract without blocking the Binance fetch', async () => { + const { fetchCalls, globalObject, source } = createHarness(); + const statuses = []; + source.subscribe({ + symbol: 'BTCUSDT', + onProfile: () => {}, + onStatus: (status) => statuses.push(status), + }); + + const response = await globalObject.fetch('/fapi/v1/rpiDepth?symbol=BTCUSDT&limit=500'); + + assert.equal(response.ok, true); + assert.equal(fetchCalls.length, 1); + assert.equal(statuses.at(-1).status, 'failed'); + assert.match(statuses.at(-1).detail, /snapshot limit/); + source.restore(); +}); + +test('does not block a malformed native snapshot request', async () => { + const { fetchCalls, globalObject, source } = createHarness(); + const statuses = []; + source.subscribe({ + symbol: 'BTCUSDT', + onProfile: () => {}, + onStatus: (status) => statuses.push(status), + }); + + const response = await globalObject.fetch('/fapi/v1/rpiDepth?limit=1000'); + + assert.equal(response.ok, true); + assert.equal(fetchCalls.length, 1); + assert.equal(statuses.at(-1).status, 'failed'); + assert.match(statuses.at(-1).detail, /symbol/); + source.restore(); +}); + +test('keeps the latest native profile for subscribers that start after page initialization', async () => { + const { globalObject, source } = createHarness(); + const nativeSocket = new globalObject.WebSocket('wss://native-binance-stream.example/ws'); + const responsePromise = globalObject.fetch('/fapi/v1/rpiDepth?symbol=BTCUSDT&limit=1000'); + nativeSocket.message(rpiMessage()); + await responsePromise; + await new Promise((resolve) => setImmediate(resolve)); + + const profiles = []; + const statuses = []; + source.subscribe({ + symbol: 'BTCUSDT', + onProfile: (profile) => profiles.push(profile), + onStatus: (status) => statuses.push(status), + }); + + assert.equal(profiles.length, 1, JSON.stringify({ statuses, state: source.getState('BTCUSDT') })); + assert.equal(profiles[0].symbol, 'BTCUSDT'); + assert.equal(statuses.at(-1).status, 'ready'); + source.restore(); +}); + +test('ignores unrelated native fetches and WebSocket streams', async () => { + const { globalObject, source } = createHarness(); + const profiles = []; + source.subscribe({ + symbol: 'BTCUSDT', + onProfile: (profile) => profiles.push(profile), + onStatus: () => {}, + }); + const nativeSocket = new globalObject.WebSocket('wss://native-binance-stream.example/ws'); + + nativeSocket.message({ stream: 'btcusdt@depth@100ms', data: update() }); + await globalObject.fetch('/fapi/v1/depth?symbol=BTCUSDT&limit=1000'); + await Promise.resolve(); + + assert.equal(profiles.length, 0); + source.restore(); +}); + +test('waits for Binance native resynchronization after a sequence gap', async () => { + const { globalObject, source } = createHarness(); + const profiles = []; + const statuses = []; + source.subscribe({ + symbol: 'BTCUSDT', + onProfile: (profile) => profiles.push(profile), + onStatus: (status) => statuses.push(status), + }); + const nativeSocket = new globalObject.WebSocket('wss://native-binance-stream.example/ws'); + const responsePromise = globalObject.fetch('/fapi/v1/rpiDepth?symbol=BTCUSDT&limit=1000'); + nativeSocket.message(rpiMessage()); + await responsePromise; + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(profiles.length, 1, JSON.stringify({ statuses, state: source.getState('BTCUSDT') })); + + nativeSocket.message(rpiMessage(update({ U: 105, u: 106, pu: 104 }))); + assert.equal(statuses.at(-1).status, 'resyncing'); + + nativeSocket.message(rpiMessage(update({ U: 106, u: 107, pu: 106 }))); + assert.equal(profiles.length, 1); + source.restore(); +}); + +test('recovers when Binance performs its next native snapshot synchronization', async () => { + let currentResponse = new FakeResponse(snapshot()); + const { globalObject, source } = createHarness(() => currentResponse); + const profiles = []; + const statuses = []; + source.subscribe({ + symbol: 'BTCUSDT', + onProfile: (profile) => profiles.push(profile), + onStatus: (status) => statuses.push(status), + }); + const nativeSocket = new globalObject.WebSocket('wss://native-binance-stream.example/ws'); + + let responsePromise = globalObject.fetch('/fapi/v1/rpiDepth?symbol=BTCUSDT&limit=1000'); + nativeSocket.message(rpiMessage()); + await responsePromise; + await new Promise((resolve) => setImmediate(resolve)); + nativeSocket.message(rpiMessage(update({ U: 105, u: 106, pu: 104 }))); + assert.equal(statuses.at(-1).status, 'resyncing'); + + currentResponse = new FakeResponse(snapshot({ lastUpdateId: 106 })); + responsePromise = globalObject.fetch('/fapi/v1/rpiDepth?symbol=BTCUSDT&limit=1000'); + nativeSocket.message(rpiMessage(update({ U: 106, u: 107, pu: 106 }))); + await responsePromise; + await new Promise((resolve) => setImmediate(resolve)); + + assert.equal(profiles.length, 2); + assert.equal(statuses.at(-1).status, 'ready'); + source.restore(); +}); diff --git a/test/unit/binance-orderbook-trade/depth-profile-session.test.js b/test/unit/binance-orderbook-trade/depth-profile-session.test.js index 14daec9..bfe48f8 100644 --- a/test/unit/binance-orderbook-trade/depth-profile-session.test.js +++ b/test/unit/binance-orderbook-trade/depth-profile-session.test.js @@ -3,202 +3,91 @@ import assert from 'node:assert/strict'; import { createDepthProfileSession } from '../../../src/binance-orderbook-trade/core/depth-profile-session.js'; -class FakeSocket { - static instances = []; - - constructor(url) { - this.url = url; - this.readyState = 0; - this.listeners = new Map(); - FakeSocket.instances.push(this); - } - - addEventListener(type, listener) { - const listeners = this.listeners.get(type) || []; - listeners.push(listener); - this.listeners.set(type, listeners); - } - - emit(type, event = {}) { - for (const listener of this.listeners.get(type) || []) listener(event); +class FakeNativeDepthSource { + constructor() { + this.subscriptions = []; } - open() { - this.readyState = 1; - this.emit('open'); + subscribe(subscription) { + this.subscriptions.push(subscription); + return () => { + this.subscriptions = this.subscriptions.filter((item) => item !== subscription); + }; } - message(payload) { - this.emit('message', { data: JSON.stringify(payload) }); + profile(profile) { + for (const subscription of this.subscriptions) subscription.onProfile(profile); } - close() { - if (this.readyState >= 2) return; - this.readyState = 3; - this.emit('close'); + status(status) { + for (const subscription of this.subscriptions) subscription.onStatus(status); } } -function depthUpdate(overrides = {}) { - return { - e: 'depthUpdate', - s: 'BTCUSDT', - st: 1, - U: 100, - u: 102, - pu: 99, - b: [['100', '2'], ['99', '3']], - a: [['101', '4'], ['102', '5']], - ...overrides, - }; -} - -function depthSnapshot(overrides = {}) { - return { - lastUpdateId: 101, - bids: [['100', '1'], ['99', '2']], - asks: [['101', '3'], ['102', '4']], - ...overrides, - }; -} - -function createTimerHarness() { - const timers = []; - return { - timers, - setTimer(callback, delay) { - const timer = { callback, delay, cleared: false }; - timers.push(timer); - return timer; - }, - clearTimer(timer) { - timer.cleared = true; - }, - }; -} - -test('opens the official public stream, synchronizes a snapshot, and emits a profile', async () => { - FakeSocket.instances = []; +test('subscribes to the existing Binance native depth source', () => { const statuses = []; const profiles = []; - const timer = createTimerHarness(); - const fetchCalls = []; + const source = new FakeNativeDepthSource(); const session = createDepthProfileSession({ symbol: 'BTCUSDT', - WebSocketCtor: FakeSocket, - fetchFn: async (url) => { - fetchCalls.push(url); - return { ok: true, json: async () => depthSnapshot() }; - }, + source, onProfile: (profile) => profiles.push(profile), onStatus: (status) => statuses.push(status), - setTimer: timer.setTimer, - clearTimer: timer.clearTimer, }); session.start(); - const socket = FakeSocket.instances[0]; - assert.equal(socket.url, 'wss://fstream.binance.com/public/ws/btcusdt@depth@100ms'); - socket.message(depthUpdate()); - socket.open(); - await Promise.resolve(); - await Promise.resolve(); + source.status({ symbol: 'BTCUSDT', status: 'synchronizing', detail: '' }); + source.profile({ symbol: 'BTCUSDT', bids: [], asks: [] }); - assert.equal(fetchCalls[0], 'https://fapi.binance.com/fapi/v1/depth?symbol=BTCUSDT&limit=1000'); assert.equal(profiles.length, 1); assert.equal(profiles[0].symbol, 'BTCUSDT'); - assert.equal(statuses.at(-1).status, 'ready'); - + assert.equal(statuses.at(-1).status, 'synchronizing'); session.stop(); - assert.equal(socket.readyState, 3); + assert.equal(source.subscriptions.length, 0); assert.equal(session.isActive(), false); }); -test('stopped sessions ignore a late snapshot', async () => { - FakeSocket.instances = []; - let resolveSnapshot; +test('stopped sessions ignore later native source events', () => { const profiles = []; + const source = new FakeNativeDepthSource(); const session = createDepthProfileSession({ symbol: 'BTCUSDT', - WebSocketCtor: FakeSocket, - fetchFn: () => new Promise((resolve) => { resolveSnapshot = resolve; }), + source, onProfile: (profile) => profiles.push(profile), onStatus: () => {}, }); session.start(); - const socket = FakeSocket.instances[0]; - socket.message(depthUpdate()); - socket.open(); session.stop(); - resolveSnapshot({ ok: true, json: async () => depthSnapshot() }); - await Promise.resolve(); - await Promise.resolve(); - + source.profile({ symbol: 'BTCUSDT', bids: [], asks: [] }); assert.equal(profiles.length, 0); }); -test('a sequence gap resets the book and performs one delayed resync', async () => { - FakeSocket.instances = []; - const statuses = []; - const profiles = []; - const timer = createTimerHarness(); - let snapshotCount = 0; +test('rejects mismatched native profile symbols', () => { + const source = new FakeNativeDepthSource(); const session = createDepthProfileSession({ symbol: 'BTCUSDT', - WebSocketCtor: FakeSocket, - fetchFn: async () => { - snapshotCount += 1; - return { ok: true, json: async () => depthSnapshot() }; - }, - onProfile: (profile) => profiles.push(profile), - onStatus: (status) => statuses.push(status), - setTimer: timer.setTimer, - clearTimer: timer.clearTimer, + source, + onProfile: () => {}, + onStatus: () => {}, }); session.start(); - const socket = FakeSocket.instances[0]; - socket.message(depthUpdate()); - socket.open(); - await Promise.resolve(); - await Promise.resolve(); - assert.equal(profiles.length, 1); - - socket.message(depthUpdate({ U: 105, u: 106, pu: 104 })); - assert.equal(statuses.at(-1).status, 'resyncing'); - assert.equal(timer.timers.at(-1).delay, 1000); - - socket.message(depthUpdate({ U: 100, u: 103, pu: 99 })); - timer.timers.at(-1).callback(); - await Promise.resolve(); - await Promise.resolve(); - - assert.equal(snapshotCount, 2); - assert.equal(profiles.length, 2); - assert.equal(statuses.at(-1).status, 'ready'); - session.stop(); + assert.throws( + () => source.profile({ symbol: 'ETHUSDT', bids: [], asks: [] }), + /symbol mismatch/, + ); }); -test('unexpected close uses bounded reconnect scheduling', () => { - FakeSocket.instances = []; - const statuses = []; - const timer = createTimerHarness(); +test('rejects duplicate starts', () => { + const source = new FakeNativeDepthSource(); const session = createDepthProfileSession({ symbol: 'BTCUSDT', - WebSocketCtor: FakeSocket, - fetchFn: async () => ({ ok: true, json: async () => depthSnapshot() }), + source, onProfile: () => {}, - onStatus: (status) => statuses.push(status), - setTimer: timer.setTimer, - clearTimer: timer.clearTimer, + onStatus: () => {}, }); session.start(); - FakeSocket.instances[0].close(); - - assert.equal(statuses.at(-1).status, 'reconnecting'); - assert.equal(timer.timers.at(-1).delay, 1000); - session.stop(); - assert.equal(timer.timers.at(-1).cleared, true); + assert.throws(() => session.start(), /already started/); });