diff --git a/ts-sdk/README.md b/ts-sdk/README.md index 798b98e4..47685fee 100644 --- a/ts-sdk/README.md +++ b/ts-sdk/README.md @@ -57,6 +57,9 @@ const ob = useSyncExternalStore( `.backstopTransfers(account)`, `.candles(id, resolution, range)` (a `SeriesResource` with `setWindow`/`loadOlder`), `.orders(account, query)`, `.bridgeConfig`, `.withdrawals(account)`. +- **Activity (ADR 0057):** `client.activity(account, query)` — one account's + orders and the money it moved in one cursor-paged, `onEvent`-observable feed; + `.orders(account)` is the narrower view it replaces. - **Charting:** `createPodDatafeed(client)` returns an `IDatafeedChartApi`-shaped object for the TradingView Charting Library (framework-agnostic, no React). diff --git a/ts-sdk/package-lock.json b/ts-sdk/package-lock.json index 7ed900f3..8e98bab7 100644 --- a/ts-sdk/package-lock.json +++ b/ts-sdk/package-lock.json @@ -1,12 +1,12 @@ { "name": "@pod-network/trade-sdk", - "version": "0.10.0", + "version": "0.11.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@pod-network/trade-sdk", - "version": "0.10.0", + "version": "0.11.0", "license": "MIT", "dependencies": { "viem": "^2.21.0" diff --git a/ts-sdk/package.json b/ts-sdk/package.json index df177ebf..5a4fc6cd 100644 --- a/ts-sdk/package.json +++ b/ts-sdk/package.json @@ -1,6 +1,6 @@ { "name": "@pod-network/trade-sdk", - "version": "0.10.0", + "version": "0.11.0", "description": "Read/stream TypeScript SDK for the pod trading indexer: cacheable REST seeds + one multiplexed WebSocket, kept in memory. Framework-agnostic.", "type": "module", "sideEffects": false, diff --git a/ts-sdk/src/client.test.ts b/ts-sdk/src/client.test.ts new file mode 100644 index 00000000..88ce74b1 --- /dev/null +++ b/ts-sdk/src/client.test.ts @@ -0,0 +1,36 @@ +// Resource memoisation. Two calls that ask for the same feed must hand back the +// same instance: a second one is a second REST seed and a second socket +// subscription for a list the app already has. + +import { describe, expect, it } from "vitest"; + +import { PodTradeClient } from "./client.js"; +import type { Address } from "./types/public.js"; + +const ACCOUNT = "0xa1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1" as Address; + +class FakeWebSocket { + constructor(public url: string) {} + close(): void {} + send(): void {} +} + +const client = () => new PodTradeClient({ + restUrl: "http://node.test/v1", + wsUrl: "ws://node.test/v1", + WebSocket: FakeWebSocket as never, +}); + +describe("PodTradeClient.activity", () => { + it("reads no filter and an empty filter as one feed", () => { + const c = client(); + try { + expect(c.activity(ACCOUNT, { types: [] })).toBe(c.activity(ACCOUNT)); + // Key order is how the caller typed the object, not part of the question. + expect(c.activity(ACCOUNT, { to: 2, from: 1 })).toBe(c.activity(ACCOUNT, { from: 1, to: 2 })); + expect(c.activity(ACCOUNT, { from: 1 })).not.toBe(c.activity(ACCOUNT)); + } finally { + c.close(); + } + }); +}); diff --git a/ts-sdk/src/client.ts b/ts-sdk/src/client.ts index 2822a176..6a487f95 100644 --- a/ts-sdk/src/client.ts +++ b/ts-sdk/src/client.ts @@ -1,5 +1,5 @@ import type { - Address, BackstopTransfer, Balances, Bar, BridgeConfig, Market, + ActivityQuery, Address, BackstopTransfer, Balances, Bar, BridgeConfig, Market, MarketId, PositionsSnapshot, Resolution, Status, TimeRange, Trigger, TriggersQuery, TxExplorer, OrdersQuery, Withdrawal, } from "./types/public.js"; @@ -13,6 +13,7 @@ import { import { withdrawalsSource } from "./sync/withdrawals.js"; import { CandleSeries, candleTailFrom, fetchCandleHistory } from "./sync/candles.js"; import { OrderHistory } from "./sync/orders.js"; +import { ActivityHistory, normalizeActivityQuery } from "./sync/activity.js"; import { enrichPositions } from "./sync/positions-live.js"; import { fetchPnlHistory, streamPnlHistory, PnlHistoryCache, type PnlHistory, type PnlHistoryChunk, type PnlHistoryQuery } from "./sync/pnl-history.js"; @@ -306,6 +307,24 @@ export class PodTradeClient { } } + /** + * One account's whole activity, newest first (ADR 0057): its orders, the + * backstop legs it was swept into, and every bridge transfer and transfer that + * moved its money. + * + * `orders` is a separate stream (`pod_orders_v2`), with its own cursor and its + * own row shape — not this feed narrowed to orders. It is to be retired in + * favour of this one, so new code should start here. + */ + activity(account: Address, query?: ActivityQuery): ActivityHistory { + // The constructor normalises too, so the key must — otherwise `{}` and + // `{ types: [] }` ask for the same feed and get two of them, each with its own + // socket subscription. Sorted, since key order is an argument's accident. + const q = normalizeActivityQuery(query); + const key = `activity:${account.toLowerCase()}:${JSON.stringify(q, Object.keys(q).sort())}`; + return this.memo(key, () => new ActivityHistory(this.ctx, account, query)); + } + orders(account: Address, query?: OrdersQuery): OrderHistory { const key = `orders:${account}:${query ? JSON.stringify(query) : ""}`; return this.memo(key, () => new OrderHistory(this.ctx, account, query)); diff --git a/ts-sdk/src/codec/activity.test.ts b/ts-sdk/src/codec/activity.test.ts new file mode 100644 index 00000000..86255fc8 --- /dev/null +++ b/ts-sdk/src/codec/activity.test.ts @@ -0,0 +1,281 @@ +// The activity seed's entry union and the `pod_activity` frame fold. None of +// it is reachable from `typecheck`: an entry is `unknown` off REST and a frame is +// `unknown` off the socket, so every field mapping is only ever checked here. +// +// Shapes are written from the node: `ActivityEntry` in `node/src/rpc/types.rs`, +// `ActivityFrame`/`MoneyEvent` in `node/src/rpc/activity.rs`, and the JSON its +// own tests assert (`node/tests/activity_frame.rs`, `clob_indexer::rest::tests`). +// Encodings follow from those types: a `U256` is `0x` hex, an `I256` a signed +// decimal string, a `WireDec` a signed decimal string, and an absent field is +// absent rather than null. + +import { describe, expect, it } from "vitest"; + +import type { ActivityEntry, Address, MarketId } from "../types/public.js"; +import type { WireActivityEntry, WireActivityFrame } from "../types/wire.js"; +import { decodeActivityEntry } from "./decode.js"; +import { applyActivityFrame } from "./activity.js"; +import { WAD } from "./units.js"; + +const ALICE = "0xa1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1" as Address; +const TOKEN = "0x7e7e7e7e7e7e7e7e7e7e7e7e7e7e7e7e7e7e7e7e" as Address; +const BOOK = "0x000000000000000000000000000000000000000000000000000000000000000a" as MarketId; +const TICK = 5_000_000; + +const hex = (n: bigint) => `0x${n.toString(16)}`; +const units = (n: bigint) => n * WAD; + +/** `{ activity_type: "order", ts, ...OrderResponse }`. */ +const ORDER_ENTRY: WireActivityEntry = { + activity_type: "order", + ts: 8_000_000, + orderbook_id: BOOK, + market_type: "perpetual", + kind: "user_signed", + order_id: "0x0101010101010101010101010101010101010101010101010101010101010101", + tx_hash: "0x6565656565656565656565656565656565656565656565656565656565656565", + bidder: ALICE, + nonce: 1, + order_type: "limit", + status: "active", + side: "buy", + price: hex(units(100n)), + initial_size: units(2n).toString(), + filled_base_amount: hex(units(1n)), + filled_quote_amount: hex(units(100n)), + fee: "0x7", + deadline: 8_000_000, + end: 3_605_000_000, + included_batch: TICK, + effective_price: hex(units(100n)), + fills: [{ + base_amount: hex(units(1n)), + quote_amount: hex(units(100n)), + timestamp: TICK, + price: hex(units(100n)), + }], + reduce_only: false, + ioc: false, +}; + +/** The money entries name their fields exactly as the stream's money events do. */ +const BACKSTOP_ENTRY: WireActivityEntry = { + activity_type: "backstop", + ts: TICK, + book: BOOK, + size: (-units(2n)).toString(), + cash: "0", + mark: units(100n).toString(), + equity: (-units(5n)).toString(), + pnl: (-units(1n)).toString(), +}; + +const BRIDGE_ENTRY: WireActivityEntry = { + activity_type: "bridge_transfer", + ts: TICK, + tx: "0x0000000000000000000000000000000000000000000000000000000000000001", + idx: 1, + token: TOKEN, + amount: "-900", + error: "insufficient_balance", +}; + +const TRANSFER_ENTRY: WireActivityEntry = { + activity_type: "transfer", + ts: TICK, + id: "0x0000000000000000000000000000000000000000000000000000000000000003", + token: TOKEN, + amount: "1700", +}; + +describe("decodeActivityEntry", () => { + it("carries an order through with its fills", () => { + const entry = decodeActivityEntry(ORDER_ENTRY); + if (entry?.activityType !== "order") throw new Error("expected an order"); + // The node timestamps an order entry by its SIGNED deadline, not the batch it + // landed in, so `timeMs` and `order.includedMs` are deliberately different. + expect(entry.timeMs).toBe(8_000); + expect(entry.order.includedMs).toBe(5_000); + expect(entry.order.marketType).toBe("perp"); + expect(entry.order.price).toBe(units(100n)); + expect(entry.order.initialSize).toBe(units(2n)); + expect(entry.order.fills).toHaveLength(1); + }); + + it("carries a backstop leg with its realized PnL", () => { + const entry = decodeActivityEntry(BACKSTOP_ENTRY); + if (entry?.activityType !== "backstop") throw new Error("expected a backstop"); + expect(entry.timeMs).toBe(5_000); + expect(entry.orderbookId).toBe(BOOK); + expect(entry.size).toBe(-units(2n)); + expect(entry.markPrice).toBe(units(100n)); + expect(entry.equity).toBe(-units(5n)); + expect(entry.realizedPnl).toBe(-units(1n)); + }); + + it("keeps money signed from the account's side, with its refusal", () => { + const bridge = decodeActivityEntry(BRIDGE_ENTRY); + if (bridge?.activityType !== "bridge_transfer") throw new Error("expected a bridge transfer"); + expect(bridge.txHash).toBe(BRIDGE_ENTRY.tx); + // One tx can carry several deposits, so the hash alone does not identify a row. + expect(bridge.idx).toBe(1); + expect(bridge.amount).toBe(-900n); + expect(bridge.error).toBe("insufficient_balance"); + + const transfer = decodeActivityEntry(TRANSFER_ENTRY); + if (transfer?.activityType !== "transfer") throw new Error("expected a transfer"); + expect(transfer.transferId).toBe(TRANSFER_ENTRY.id); + expect(transfer.amount).toBe(1700n); + expect(transfer.error).toBeUndefined(); + }); +}); + +/** The shape `the_wire_is_tagged_stringly_and_never_null` asserts: one frame per + * tick, no `book` on the frame, books first and then the tick's money. */ +const FRAME: WireActivityFrame = { + batch: TICK, + orders: [{ + id: "0x0101010101010101010101010101010101010101010101010101010101010101", + tx: "0x6565656565656565656565656565656565656565656565656565656565656565", + book: BOOK, + n: 1, + px: units(100n).toString(), + sz: units(2n).toString(), + end: 3_605_000_000, + }], + events: [ + { k: "new", o: 0 }, + { + k: "fill", + o: 0, + b: units(1n).toString(), + q: units(100n).toString(), + tb: units(1n).toString(), + tq: units(100n).toString(), + tf: "7", + }, + { + k: "backstop", + book: BOOK, + size: (-units(2n)).toString(), + cash: "0", + mark: units(100n).toString(), + equity: (-units(5n)).toString(), + pnl: (-units(1n)).toString(), + }, + { + k: "bridge_transfer", + tx: "0x0000000000000000000000000000000000000000000000000000000000000001", + idx: 1, + token: TOKEN, + amount: units(10n).toString(), + }, + { + k: "transfer", + id: "0x0000000000000000000000000000000000000000000000000000000000000005", + token: TOKEN, + amount: (-units(5n)).toString(), + error: "insufficient_balance", + }, + ], +}; + +const emptyState = () => ({ orders: new Map(), entries: [] as ActivityEntry[] }); + +describe("applyActivityFrame", () => { + it("folds the order flow and collects the money the tick moved", () => { + const state = emptyState(); + const events = applyActivityFrame(FRAME, state, { account: ALICE }); + + expect(events.map((e) => ("activityType" in e ? e.activityType : e.kind))) + .toEqual(["new", "fill", "backstop", "bridge_transfer", "transfer"]); + + const order = state.orders.get(FRAME.orders[0]!.id); + expect(order?.status).toBe("active"); + expect(order?.filledBase).toBe(units(1n)); + // `book` moved from the frame onto the entity; there is no frame-level one. + expect(order?.orderbookId).toBe(BOOK); + // A frame with no `accts` covers exactly one account. + expect(order?.bidder).toBe(ALICE); + + expect(state.entries.map((e) => e.activityType)) + .toEqual(["backstop", "bridge_transfer", "transfer"]); + // Every money entry is timed by the batch that moved it. + expect(state.entries.every((e) => e.timeMs === 5_000)).toBe(true); + }); + + it("decodes each money kind into the entry the seed would have served", () => { + const state = emptyState(); + applyActivityFrame(FRAME, state, { account: ALICE }); + const [backstop, bridge, transfer] = state.entries; + + if (backstop?.activityType !== "backstop") throw new Error("expected a backstop entry"); + expect(backstop.orderbookId).toBe(BOOK); + expect(backstop.size).toBe(-units(2n)); + expect(backstop.markPrice).toBe(units(100n)); + expect(backstop.realizedPnl).toBe(-units(1n)); + + if (bridge?.activityType !== "bridge_transfer") throw new Error("expected a bridge entry"); + expect(bridge.txHash).toBe("0x0000000000000000000000000000000000000000000000000000000000000001"); + expect(bridge.idx).toBe(1); + expect(bridge.token).toBe(TOKEN); + expect(bridge.amount).toBe(units(10n)); + expect(bridge.error).toBeUndefined(); + + if (transfer?.activityType !== "transfer") throw new Error("expected a transfer entry"); + expect(transfer.amount).toBe(-units(5n)); + expect(transfer.error).toBe("insufficient_balance"); + }); + + it("ignores a kind it does not know, on either side of the union", () => { + const state = emptyState(); + const events = applyActivityFrame( + { ...FRAME, events: [{ k: "liquidation_notice", o: 0 }, { k: "airdrop", amount: "1" } as never] }, + state, + { account: ALICE }, + ); + expect(events).toEqual([]); + expect(state.entries).toEqual([]); + // The entity still lands: the frame said the order exists. + expect(state.orders.size).toBe(1); + }); +}); + +const BOOK_B = "0x000000000000000000000000000000000000000000000000000000000000000b" as MarketId; + +const sweep = (book: MarketId) => ({ + k: "backstop" as const, + book, + size: (-units(2n)).toString(), + cash: "0", + mark: units(100n).toString(), + equity: (-units(5n)).toString(), + pnl: (-units(1n)).toString(), +}); + +describe("applyActivityFrame interleaving", () => { + it("emits order and money events in the node's wire order", () => { + const state = emptyState(); + const events = applyActivityFrame({ + batch: TICK, + orders: [ + { id: "0x0a", tx: "0xaa", book: BOOK, n: 1, px: units(100n).toString(), sz: units(1n).toString() }, + { id: "0x0b", tx: "0xbb", book: BOOK_B, n: 2, px: units(100n).toString(), sz: units(1n).toString() }, + ], + events: [{ k: "new", o: 0 }, sweep(BOOK), { k: "new", o: 1 }, sweep(BOOK_B)], + }, state, { account: ALICE }); + + // Folding the orders first and appending the money would report the tick as + // two placements followed by two sweeps, which is not what happened. + expect(events.map((e) => ("activityType" in e ? e.activityType : e.kind))) + .toEqual(["new", "backstop", "new", "backstop"]); + expect(state.entries.map((e) => e.activityType === "backstop" ? e.orderbookId : e.activityType)) + .toEqual([BOOK, BOOK_B]); + }); +}); + +describe("decodeActivityEntry on an unknown kind", () => { + it("returns undefined rather than an untagged row", () => { + expect(decodeActivityEntry({ activity_type: "airdrop", ts: TICK } as never)).toBeUndefined(); + }); +}); diff --git a/ts-sdk/src/codec/activity.ts b/ts-sdk/src/codec/activity.ts new file mode 100644 index 00000000..98d469d6 --- /dev/null +++ b/ts-sdk/src/codec/activity.ts @@ -0,0 +1,46 @@ +// `pod_activity` frames -> the order map plus the account's money rows (ADR 0057). +// +// A frame is `pod_orders_v2` for one whole account: the order half folds through +// `applyOrdersFrame` unchanged, and the money half is a set of extra event kinds +// whose `k` tags are disjoint from the order ones. Money has no entity to +// reference and no state to mutate, so each event is decoded straight into the +// entry the REST seed would have served for it — in one pass over `frame.events`, +// so what a consumer sees is the node's own interleaving of the tick. + +import type { ActivityEntry, ActivityEvent, Address, MoneyActivity, Order } from "../types/public.js"; +import type { WireActivityFrame, WireMoneyEvent, WireOrderEvent, WireOrdersFrame } from "../types/wire.js"; +import { decodeMoneyEvent } from "./decode.js"; +import { applyOrdersFrame } from "./orders-v2.js"; +import { usToMs } from "./units.js"; + +const MONEY_KINDS: Record = { + backstop: true, + bridge_transfer: true, + transfer: true, +}; + +export const isMoney = (k: string): k is WireMoneyEvent["k"] => Object.hasOwn(MONEY_KINDS, k); + +export interface ActivityState { + orders: Map; + entries: ActivityEntry[]; +} + +export function applyActivityFrame( + frame: WireActivityFrame, + state: ActivityState, + ctx: { account: Address }, +): ActivityEvent[] { + const timeMs = usToMs(frame.batch); + return applyOrdersFrame(frame as WireOrdersFrame, state.orders, { + account: ctx.account, + // Order kinds, and money kinds this version does not know, are ignored rather + // than handed on — exactly as `applyOrdersFrame` does. + foreign: (event: WireOrderEvent) => { + if (!isMoney(event.k)) return undefined; + const entry = decodeMoneyEvent(event as unknown as WireMoneyEvent, timeMs); + state.entries.push(entry); + return entry; + }, + }); +} diff --git a/ts-sdk/src/codec/decode.ts b/ts-sdk/src/codec/decode.ts index 6f6b8192..70be07f0 100644 --- a/ts-sdk/src/codec/decode.ts +++ b/ts-sdk/src/codec/decode.ts @@ -2,14 +2,14 @@ // representations are normalized into bigint + millisecond numbers. import type { - Bar, BackstopTransfer, Balances, BridgeConfig, Market, Order, + ActivityEntry, Bar, BackstopTransfer, Balances, BridgeConfig, Market, Order, Orderbook, PartialFill, PerpPosition, Position, PositionsSnapshot, SpotHolding, - SpotPosition, Status, Trigger, MarketType, OrderDirection, OrderKind, OrderStatus, TriggerType, - Withdrawal, + SpotPosition, Status, Trigger, MarketType, MoneyActivity, OrderDirection, OrderKind, + OrderStatus, TriggerType, Withdrawal, } from "../types/public.js"; import type { - WireBackstopTransfer, WireBalances, WireBridgeConfig, WireCandle, - WireMarketDynamics, WireMarketStatic, WireOrder, WireOrderbook, WirePartialFill, + WireActivityEntry, WireBackstopTransfer, WireBalances, WireBridgeConfig, WireCandle, + WireMarketDynamics, WireMarketStatic, WireMoneyEvent, WireOrder, WireOrderbook, WirePartialFill, WirePerpPosition, WirePosition, WirePositionsSnapshot, WireSpotHolding, WireSpotPosition, WireStatus, WireTrigger, WireWithdrawal, } from "../types/wire.js"; @@ -263,10 +263,70 @@ export function decodeBackstopTransfer(w: WireBackstopTransfer): BackstopTransfe cash: dec(w.cash), markPrice: dec(w.mark_price), equity: dec(w.equity), + realizedPnl: dec(w.realized_pnl), time: usToMs(w.timestamp_us), }; } +/** One money row (ADR 0057), from a seed entry or a stream event alike — the two + * carry the same fields under different tags. */ +export function decodeMoneyEvent(event: WireMoneyEvent, timeMs: number): MoneyActivity { + switch (event.k) { + case "backstop": + return { + activityType: "backstop", + timeMs, + time: timeMs, + orderbookId: event.book, + size: dec(event.size), + cash: dec(event.cash), + markPrice: dec(event.mark), + equity: dec(event.equity), + realizedPnl: dec(event.pnl), + }; + case "bridge_transfer": + return { + activityType: "bridge_transfer", + timeMs, + txHash: event.tx, + idx: event.idx, + token: event.token, + amount: dec(event.amount), + // `||`, as `decodeWithdrawal` does: an empty string is the wire saying + // "no reason", and it would otherwise read as a failure. + error: event.error || undefined, + }; + case "transfer": + return { + activityType: "transfer", + timeMs, + transferId: event.id, + token: event.token, + amount: dec(event.amount), + error: event.error || undefined, + }; + } +} + +/** + * One activity row (ADR 0057), reusing the decoder its kind already had, or + * `undefined` for a kind this version does not know — the same silence the frame + * path keeps, rather than a row with no tag a consumer could act on. + * + * `timeMs` is the node's own sort key, which for an order is its **signed + * deadline** rather than the batch it landed in — `order.includedMs` is that. + */ +export function decodeActivityEntry(w: WireActivityEntry): ActivityEntry | undefined { + const timeMs = usToMs(w.ts); + switch (w.activity_type) { + case "order": return { activityType: "order", timeMs, order: decodeOrder(w) }; + case "backstop": return decodeMoneyEvent({ ...w, k: "backstop" }, timeMs); + case "bridge_transfer": return decodeMoneyEvent({ ...w, k: "bridge_transfer" }, timeMs); + case "transfer": return decodeMoneyEvent({ ...w, k: "transfer" }, timeMs); + default: return undefined; + } +} + export function decodeBridgeConfig(w: WireBridgeConfig): BridgeConfig { return { claimChainId: w.claim_chain_id, diff --git a/ts-sdk/src/codec/orders-v2.ts b/ts-sdk/src/codec/orders-v2.ts index e806e5c8..e7aef5f4 100644 --- a/ts-sdk/src/codec/orders-v2.ts +++ b/ts-sdk/src/codec/orders-v2.ts @@ -15,9 +15,15 @@ import { dec, endMsFromUs, usToMs } from "./units.js"; import { div } from "./fixed.js"; import { classifyPerpDirection } from "./direction.js"; -export interface OrdersFrameContext { +export interface OrdersFrameContext { /** The account the subscription is filtered to, used when the frame omits `accts`. */ account: Address; + /** + * Events this codec applied nothing for, offered back in wire order — how a + * superset channel (`pod_activity`) folds its own kinds without a second pass + * that would report the tick's money after all of its order flow. + */ + foreign?(event: WireOrderEvent): X | undefined; } /** Per-frame facts the events need. */ @@ -37,13 +43,17 @@ interface FrameFacts { * mid-frame: they are collected while applying but the rows are mutated in place, so a * consumer never sees a half-applied order. */ -export function applyOrdersFrame(frame: WireOrdersFrame, byId: Map, ctx: OrdersFrameContext): OrderEvent[] { +export function applyOrdersFrame( + frame: WireOrdersFrame, + byId: Map, + ctx: OrdersFrameContext, +): (OrderEvent | X)[] { // The book and the batch are frame constants: resolved once, not per entity. const batchMs = usToMs(frame.batch); const facts: FrameFacts = { batchMs, accts: frame.accts }; const created = frame.orders.map((e) => decodeEntity(e, frame, batchMs, ctx.account)); - const applied: OrderEvent[] = []; + const applied: (OrderEvent | X)[] = []; for (const event of frame.events) { // `o` indexes this frame's entities and `id` names one resting from an earlier // batch — but an event kind added later (ADR 0029 §6 reserves several) may name @@ -53,12 +63,11 @@ export function applyOrdersFrame(frame: WireOrdersFrame, byId: Map o.id === event.id); - if (!target) continue; // `undefined` means a kind this version does not know, which is ignored rather than // handed on as an event a consumer cannot act on (ADR 0029 §6). The switch is the // only list of what is recognised — a separate one could disagree with it. - const outcome = applyEvent(event, target, facts); - if (outcome) { + const outcome = target ? applyEvent(event, target, facts) : undefined; + if (outcome && target) { applied.push({ kind: event.k as OrderEventKind, order: target, @@ -66,7 +75,10 @@ export function applyOrdersFrame(frame: WireOrdersFrame, byId: Map `0x${n.toString(16)}`; +const units = (n: bigint) => n * WAD; +const id = (n: number) => `0x${n.toString(16).padStart(64, "0")}` as Hex; +const flush = () => new Promise((r) => setTimeout(r, 0)); + +const wireOrder = (n: number, at: number, timestampUs = at): WireActivityEntry => ({ + activity_type: "order", + ts: timestampUs, + orderbook_id: BOOK, + market_type: "perpetual", + kind: "user_signed", + order_id: id(n), + tx_hash: id(n + 100), + bidder: ACCOUNT, + nonce: n, + order_type: "limit", + status: "active", + side: "buy", + price: hex(units(100n)), + initial_size: units(2n).toString(), + filled_base_amount: "0x0", + filled_quote_amount: "0x0", + fee: "0x0", + deadline: at, + end: at + 600_000, + included_batch: at, + fills: [], +}); + +const wireBackstop = (at: number): WireActivityEntry => ({ + activity_type: "backstop", + ts: at, + book: BOOK, + size: (-units(2n)).toString(), + cash: "0", + mark: units(100n).toString(), + equity: (-units(5n)).toString(), + pnl: (-units(1n)).toString(), +}); + +const wireBridge = (n: number, at: number, idx = 0): WireActivityEntry => ({ + activity_type: "bridge_transfer", + ts: at, + tx: id(n), + idx, + token: TOKEN, + amount: "-900", + error: "insufficient_balance", +}); + +const wireTransfer = (n: number, at: number): WireActivityEntry => ({ + activity_type: "transfer", + ts: at, + id: id(n), + token: TOKEN, + amount: "1700", +}); + +/** The page the node serves, newest first: `(timestamp, ordinal, key)` descending. */ +const SEED: WireActivityEntry[] = [ + wireTransfer(3, SEED_TICK), + wireBridge(2, SEED_TICK), + wireBackstop(SEED_TICK), + wireOrder(1, SEED_TICK), +]; + +function page(entries: WireActivityEntry[], nextCursor: string | null = null): ActivityPage { + const activity = entries.map(decodeActivityEntry).filter((e): e is ActivityEntry => e !== undefined); + return { activity, nextCursor, solutionNow: SEED_TICK / 1000 }; +} + +/** + * `ActivityHistory` over a stub context: scripted REST pages, and a websocket + * whose `subscribe` hands us the frame and error callbacks so a test can deliver + * what the transport would. + */ +function harness(opts?: { pages?: (ActivityPage | Error)[]; query?: ActivityQuery; defer?: boolean }) { + const pages = [...(opts?.pages ?? [page(SEED)])]; + const pending: ((p: ActivityPage) => void)[] = []; + let onOpen: (() => void) | undefined; + let deliver: ((r: unknown) => void) | undefined; + let refuse: ((e: unknown) => void) | undefined; + const subscribed: SubParams[] = []; + const updates: SubParams[] = []; + let resubscribes = 0; + const ws = { + state: "open", + on: (_event: string, handler: () => void) => { onOpen = handler; return () => { onOpen = undefined; }; }, + subscribe: ( + _channel: string, + params: SubParams, + onMessage: (r: unknown) => void, + onError: (e: unknown) => void, + ) => { + subscribed.push(params); + deliver = onMessage; + refuse = onError; + return { + unsubscribe: () => {}, + update: (p: SubParams) => updates.push(p), + resubscribe: () => { resubscribes++; }, + }; + }, + }; + const queries: unknown[] = []; + const rest = { + activity: vi.fn((_account: Address, q?: unknown) => { + queries.push(q); + if (opts?.defer) return new Promise((resolve) => { pending.push(resolve); }); + const next = pages.shift() ?? page([]); + return next instanceof Error ? Promise.reject(next) : Promise.resolve(next); + }), + }; + const history = new ActivityHistory( + { rest, ws, positionResyncMs: 0, marketResyncMs: 0 } as unknown as SyncContext, + ACCOUNT, + opts?.query, + ); + return { + history, + queries, + subscribed, + updates, + frame: (f: unknown) => deliver?.(f), + close: (e: unknown) => refuse?.(e), + open: () => onOpen?.(), + pending, + settle: (i: number, p: ActivityPage) => pending[i]!(p), + restCalls: () => rest.activity.mock.calls.length, + resubscribes: () => resubscribes, + }; +} + +const FRAME: WireActivityFrame = { + batch: SEED_TICK + 500_000, + orders: [{ id: id(9), tx: id(109), book: BOOK, n: 9, px: units(101n).toString(), sz: units(1n).toString() }], + events: [ + { k: "new", o: 0 }, + { + k: "backstop", + book: BOOK, + size: (-units(2n)).toString(), + cash: "0", + mark: units(100n).toString(), + equity: (-units(5n)).toString(), + pnl: (-units(1n)).toString(), + }, + { k: "bridge_transfer", tx: id(21), idx: 0, token: TOKEN, amount: units(10n).toString() }, + { k: "transfer", id: id(22), token: TOKEN, amount: (-units(5n)).toString() }, + ], +}; + +describe("ActivityHistory seed", () => { + it("decodes every kind the page carries", async () => { + const { history } = harness(); + const entries = await history.ready(); + + expect(entries.map((e) => e.activityType)).toEqual(["transfer", "bridge_transfer", "backstop", "order"]); + const [transfer, bridge, backstop] = entries; + if (transfer?.activityType !== "transfer") throw new Error("expected a transfer"); + expect(transfer.amount).toBe(1700n); + if (bridge?.activityType !== "bridge_transfer") throw new Error("expected a bridge transfer"); + expect(bridge.amount).toBe(-900n); + expect(bridge.error).toBe("insufficient_balance"); + if (backstop?.activityType !== "backstop") throw new Error("expected a backstop"); + expect(backstop.realizedPnl).toBe(-units(1n)); + }); + + it("subscribes from the page's watermark, in micros", async () => { + const { history, subscribed } = harness(); + await history.ready(); + await flush(); + expect(subscribed).toEqual([{ account: ACCOUNT, since: SEED_TICK }]); + }); + + it("passes the query's types and window to REST, in micros", async () => { + const { history, queries } = harness({ + query: { types: ["transfer", "order"], from: 1_000, to: 2_000, limit: 25 }, + }); + await history.ready(); + expect(queries[0]).toMatchObject({ + types: ["transfer", "order"], + from: 1_000, + to: 2_000, + limit: 25, + }); + }); + + it("reads an empty `types` as no filter, on both REST and the frame", async () => { + const { history, frame, queries } = harness({ query: { types: [] } }); + await history.ready(); + await flush(); + + expect((queries[0] as { types?: unknown }).types).toBeUndefined(); + expect(history.get()).toHaveLength(4); + frame(FRAME); + expect(history.get()).toHaveLength(8); + }); +}); + +describe("ActivityHistory stream", () => { + it("applies a frame's orders and money, and emits every event", async () => { + const { history, frame } = harness(); + const seen: ActivityEvent[][] = []; + const off = history.onEvent((events) => seen.push(events)); + await history.ready(); + await flush(); + + frame(FRAME); + expect(seen).toHaveLength(1); + expect(seen[0]!.map((e) => ("activityType" in e ? e.activityType : e.kind))) + .toEqual(["new", "backstop", "bridge_transfer", "transfer"]); + + const entries = history.get()!; + // The new order, three money rows, and the four the seed carried. + expect(entries).toHaveLength(8); + expect(entries.filter((e) => e.activityType === "order")).toHaveLength(2); + off(); + }); + + it("drops a frame the seed already settled and applies the next one", async () => { + const { history, frame, updates } = harness(); + await history.ready(); + await flush(); + + // One frame per tick, so a frame at `since` is already delivered. + frame({ ...FRAME, batch: SEED_TICK }); + expect(history.get()).toHaveLength(4); + expect(updates).toHaveLength(0); + + frame(FRAME); + expect(history.get()).toHaveLength(8); + expect(updates).toEqual([{ since: FRAME.batch }]); + }); + + it("resumes a closed stream, then re-seeds once resuming stops working", async () => { + vi.useFakeTimers(); + try { + const { history, close, resubscribes, restCalls } = harness(); + history.subscribe(() => {}); + await vi.advanceTimersByTimeAsync(0); + const seeds = restCalls(); + + const closed = (resumable: boolean) => new PodSubscriptionClosedError({ + code: resumable ? -32020 : -32023, + data: { resumable, resume_since: SEED_TICK + 1_000_000 }, + }); + + close(closed(true)); + expect(resubscribes()).toBe(1); + expect(restCalls(), "a resumable close needs no re-seed").toBe(seeds); + + close(closed(false)); + await vi.advanceTimersByTimeAsync(1_000); + expect(restCalls(), "and a non-resumable one re-seeds over REST").toBe(seeds + 1); + expect(resubscribes()).toBe(2); + } finally { + vi.useRealTimers(); + } + }); + + it("retries a failed initial seed until one lands, then subscribes once", async () => { + vi.useFakeTimers(); + try { + const { history, subscribed, restCalls } = harness({ + pages: [new Error("indexer unavailable"), page(SEED)], + }); + history.subscribe(() => {}); + await vi.advanceTimersByTimeAsync(0); + // Nothing is subscribed yet, so no close or rejection would ever arrive to + // schedule another attempt — the seed has to re-arm itself. + expect(restCalls()).toBe(1); + expect(subscribed).toHaveLength(0); + + await vi.advanceTimersByTimeAsync(500); + expect(restCalls()).toBe(2); + expect(subscribed).toEqual([{ account: ACCOUNT, since: SEED_TICK }]); + expect(history.get()).toHaveLength(4); + + await vi.advanceTimersByTimeAsync(60_000); + expect(subscribed).toHaveLength(1); + } finally { + vi.useRealTimers(); + } + }); + + it("times an order by the node's row time, and a streamed one by its batch", async () => { + // The node times an order row by its SIGNED deadline, which is neither the + // batch it landed in nor anything a frame carries. + const signed = SEED_TICK + 2_000_000; + const { history, frame } = harness({ pages: [page([wireOrder(1, SEED_TICK, signed)])] }); + await history.ready(); + await flush(); + + frame({ + ...FRAME, + // The seeded order again (the frame re-delivers its terms), plus a new one. + orders: [ + { id: id(1), tx: id(101), book: BOOK, n: 1, px: units(100n).toString(), sz: units(2n).toString() }, + FRAME.orders[0]!, + ], + events: [{ k: "new", o: 1 }], + }); + + const times = new Map(history.get()!.flatMap((e) => + e.activityType === "order" ? [[e.order.id, e.timeMs] as const] : [])); + expect(times.get(id(1))).toBe(signed / 1000); + expect(times.get(id(9))).toBe(FRAME.batch / 1000); + }); + + it("drops the events of a type the query excludes", async () => { + // REST filters server-side, so the seed is already narrow; the frame is not. + const { history, frame } = harness({ + pages: [page([wireTransfer(3, SEED_TICK)])], + query: { types: ["transfer"] }, + }); + const seen: ActivityEvent[][] = []; + const off = history.onEvent((events) => seen.push(events)); + await history.ready(); + await flush(); + + frame(FRAME); + expect(seen[0]!.map((e) => ("activityType" in e ? e.activityType : e.kind))).toEqual(["transfer"]); + // The excluded kinds never reach the list either — including the order the + // frame's `new` would otherwise have created. + expect(history.get()!.map((e) => e.activityType)).toEqual(["transfer", "transfer"]); + off(); + }); +}); + +describe("ActivityHistory paging", () => { + it("merges an older page without duplicating what it already holds", async () => { + const older = [wireTransfer(3, SEED_TICK), wireOrder(1, SEED_TICK), wireOrder(4, SEED_TICK - 500_000)]; + const { history } = harness({ pages: [page(SEED, "3000000:4:00"), page(older)] }); + await history.ready(); + expect(history.hasMore()).toBe(true); + + await history.loadOlder(); + // Only the one row the first page did not carry. + expect(history.get()).toHaveLength(5); + expect(history.hasMore()).toBe(false); + }); + + it("keeps every same-tick money row when a page boundary splits the run", async () => { + // Three transfers and one two-deposit tx, all in one tick, split across two + // pages: nothing on the wire orders them, so only the rows' own ids can key them. + const seed = [wireTransfer(31, SEED_TICK), wireTransfer(32, SEED_TICK)]; + const older = [ + wireTransfer(33, SEED_TICK), + wireBridge(41, SEED_TICK, 1), + wireBridge(41, SEED_TICK, 0), + ]; + const { history } = harness({ pages: [page(seed, "3000000:4:00"), page(older)] }); + await history.ready(); + await history.loadOlder(); + + const rows = history.get()!; + const keys = rows.map((e) => e.activityType === "transfer" ? e.transferId + : e.activityType === "bridge_transfer" ? `${e.txHash}:${e.idx}` : e.activityType); + expect(new Set(keys).size).toBe(5); + expect(rows).toHaveLength(5); + }); + + it("sorts newest first, ties by the node's ordinal", async () => { + const { history } = harness({ + pages: [page([wireOrder(1, SEED_TICK), wireBackstop(SEED_TICK), wireTransfer(3, SEED_TICK)])], + }); + const entries = await history.ready(); + expect(entries.map((e) => e.activityType)).toEqual(["transfer", "backstop", "order"]); + }); +}); + +describe("ActivityHistory.onEvent", () => { + it("starts the stream on its own, and stops delivering once released", async () => { + const { history, frame } = harness(); + const seen: string[] = []; + const off = history.onEvent((events) => seen.push(...events.map((e) => ("activityType" in e ? e.activityType : e.kind)))); + await flush(); + + frame(FRAME); + expect(seen).toEqual(["new", "backstop", "bridge_transfer", "transfer"]); + + off(); + frame({ ...FRAME, batch: FRAME.batch + 500_000 }); + expect(seen).toEqual(["new", "backstop", "bridge_transfer", "transfer"]); + }); + + it("says nothing for the REST seed", async () => { + const { history } = harness(); + const seen: ActivityEvent[][] = []; + const off = history.onEvent((events) => seen.push(events)); + await flush(); + expect(seen).toEqual([]); + off(); + }); +}); + +const moneyId = (e: ActivityEntry) => e.activityType === "transfer" ? e.transferId : e.activityType; + +describe("ActivityHistory window", () => { + it("drops a live frame outside the query's window", async () => { + const { history, frame } = harness({ query: { from: 0, to: 3_000 } }); + await history.ready(); + await flush(); + const before = history.get()!.length; + + // REST filters server-side; nothing stops the stream from pushing a later tick. + frame({ ...FRAME, batch: 3_500_000 }); + expect(history.get()).toHaveLength(before); + }); +}); + +describe("ActivityHistory re-seed", () => { + it("keeps a streamed order's time when a re-seed reports a later one", async () => { + const later = SEED_TICK + 9_000_000; + const { history, frame, open } = harness({ + pages: [page([]), page([wireOrder(9, SEED_TICK, later)])], + }); + history.subscribe(() => {}); + await flush(); + frame(FRAME); + + open(); + await flush(); + const row = history.get()!.find((e) => e.activityType === "order"); + expect(row?.timeMs).toBe(FRAME.batch / 1000); + }); + + it("keeps the paging state across a reconnect re-seed", async () => { + const { history, open } = harness({ + pages: [page(SEED, "c1"), page([wireOrder(4, SEED_TICK - 500_000)]), page(SEED, "c1")], + }); + await history.ready(); + await history.loadOlder(); + expect(history.hasMore()).toBe(false); + + open(); + await flush(); + expect(history.hasMore(), "a re-seed re-paints the first page, not the paging state").toBe(false); + }); + + it("ignores a seed that lands after a fresher one", async () => { + const { history, open, pending, settle, subscribed } = harness({ defer: true }); + history.subscribe(() => {}); + await flush(); + open(); + await flush(); + expect(pending).toHaveLength(2); + + settle(1, page([wireTransfer(3, SEED_TICK)])); + await flush(); + settle(0, { + activity: [decodeActivityEntry(wireTransfer(9, SEED_TICK))!], + nextCursor: null, + solutionNow: (SEED_TICK + 5_000_000) / 1000, + }); + await flush(); + + expect(history.get()!.map(moneyId)).toEqual([id(3)]); + expect(subscribed).toEqual([{ account: ACCOUNT, since: SEED_TICK }]); + }); + + it("does not report an error once an empty first page has painted", async () => { + vi.useFakeTimers(); + try { + const { history, close } = harness({ pages: [page([]), new Error("indexer down")] }); + history.subscribe(() => {}); + await vi.advanceTimersByTimeAsync(0); + expect(history.get()).toEqual([]); + + close(new PodSubscriptionClosedError({ code: -32023, data: { resumable: false } })); + await vi.advanceTimersByTimeAsync(1_000); + // An account with no activity is a painted resource, not a failed one. + expect(history.error).toBeUndefined(); + } finally { + vi.useRealTimers(); + } + }); +}); + +describe("ActivityHistory ordering", () => { + it("breaks a same-tick tie the same way from a page and from a frame", async () => { + const seeded = harness({ + pages: [page([wireTransfer(31, SEED_TICK), wireTransfer(32, SEED_TICK)])], + }); + const fromPage = (await seeded.history.ready()).map(moneyId); + + const streamed = harness({ pages: [page([])] }); + streamed.history.subscribe(() => {}); + await flush(); + streamed.frame({ + batch: SEED_TICK + 500_000, + orders: [], + events: [ + { k: "transfer", id: id(32), token: TOKEN, amount: "1700" }, + { k: "transfer", id: id(31), token: TOKEN, amount: "1700" }, + ], + }); + + expect(streamed.history.get()!.map(moneyId)).toEqual(fromPage); + }); +}); diff --git a/ts-sdk/src/sync/activity.ts b/ts-sdk/src/sync/activity.ts new file mode 100644 index 00000000..18f0fd56 --- /dev/null +++ b/ts-sdk/src/sync/activity.ts @@ -0,0 +1,266 @@ +// ActivityHistory: one account's whole activity (ADR 0057) — `OrderHistory` +// widened to the money the account moved. Seeded from the first REST page, kept +// live by `pod_activity`, paged backwards by cursor. +// +// The cursor is the batch alone. `pod_activity` sends one frame per tick, +// covering every book, so a frame is already delivered exactly when its batch is +// at or below `since` — there is no `sinceBook` to resume inside a tick with. + +import type { + ActivityEntry, ActivityEvent, ActivityQuery, Address, MoneyActivity, Order, +} from "../types/public.js"; +import type { WireActivityFrame } from "../types/wire.js"; +import { applyActivityFrame, isMoney } from "../codec/activity.js"; +import { msToUs, usToMs } from "../codec/units.js"; +import { BaseResource, type ResourceHandle } from "../stores/resource.js"; +import type { SubParams } from "../transport/ws.js"; +import type { SeriesResource } from "./candles.js"; +import type { SyncContext } from "./sources.js"; +import { compareCursor, ResumableStream } from "./stream.js"; + +/** + * No types is no filter, and an empty list says the same thing. One form, so REST, + * `keep()` and the client's memo key cannot read the same query three ways. + */ +export function normalizeActivityQuery(query: ActivityQuery = {}): ActivityQuery { + return query.types?.length ? query : { ...query, types: undefined }; +} + +/** The node's page order within a tick: `(timestamp, ordinal, key)` descending. */ +const ORDINAL = { order: 1, backstop: 2, bridge_transfer: 3, transfer: 4 } as const; + +/** + * A money row's identity, read off the row itself — so the same row keys the same + * whichever page it arrives on and however many of its kind share its tick. A + * backstop leg has no id of its own, but a sweep touches each market once. + */ +function moneyKey(entry: MoneyActivity): string { + switch (entry.activityType) { + case "transfer": return `transfer:${entry.transferId}`; + case "bridge_transfer": return `bridge:${entry.txHash}:${entry.idx}`; + case "backstop": return `backstop:${entry.timeMs}:${entry.orderbookId ?? "cash"}`; + } +} + +const cmp = (a: string, b: string): number => a < b ? -1 : a > b ? 1 : 0; + +/** + * Separate two rows the tick and the ordinal cannot. Insertion order would do it + * for one source alone, but a tick that arrives as a frame and a tick that arrives + * as a page are built in different orders, and a reload must not reshuffle the list. + */ +function tiebreak(a: ActivityEntry, b: ActivityEntry): number { + if (a.activityType === "order" && b.activityType === "order") { + return b.order.nonce - a.order.nonce || cmp(a.order.id, b.order.id); + } + // Equal ordinals, so neither is an order. + return cmp(moneyKey(a as MoneyActivity), moneyKey(b as MoneyActivity)); +} + +export class ActivityHistory implements SeriesResource { + private readonly base: BaseResource; + private handle: ResourceHandle | undefined; + private readonly orders = new Map(); + /** + * When each order was reported, by order id. Kept beside the row rather than on + * it: the node times an order by its signed deadline, which no frame carries, so + * neither a later frame nor a re-seed must be able to retime a row already placed. + */ + private readonly orderTimeMs = new Map(); + /** The money rows, by {@link moneyKey}. */ + private readonly money = new Map(); + private nextCursor: string | null = null; + private painted = false; + private seedGen = 0; + private _hasMore = true; + private _loading = false; + private readonly eventListeners = new Set<(events: ActivityEvent[]) => void>(); + private readonly stream: ResumableStream; + + private readonly query: ActivityQuery; + + constructor( + private readonly ctx: SyncContext, + private readonly account: Address, + query: ActivityQuery = {}, + ) { + this.query = normalizeActivityQuery(query); + this.stream = new ResumableStream({ + ws: ctx.ws, + channel: "pod_activity", + params: { account }, + reseed: () => this.fetchFirstPage(), + onFrame: (r) => this.onFrame(r), + }); + this.base = new BaseResource((h) => { + this.handle = h; + this.stream.start(); + return () => { + this.stream.stop(); + this.handle = undefined; + }; + }); + } + + get(): ActivityEntry[] | undefined { return this.base.get(); } + subscribe(listener: () => void): () => void { return this.base.subscribe(listener); } + ready(): Promise { return this.base.ready(); } + get error(): Error | undefined { return this.base.error; } + hasMore(): boolean { return this._hasMore; } + loading(): boolean { return this._loading; } + destroy(): void { this.base.destroy(); } + + /** + * Observe what the stream reported, a tick at a time: the order transitions and + * the money that moved. Nothing is emitted for the REST seed. + */ + onEvent(listener: (events: ActivityEvent[]) => void): () => void { + this.eventListeners.add(listener); + // Listening starts the stream, like subscribing does — the resource is + // ref-counted from `subscribe`/`ready` alone. + const release = this.base.subscribe(() => {}); + return () => { this.eventListeners.delete(listener); release(); }; + } + + setWindow(): void { /* activity pages by cursor, not by time window */ } + + async loadOlder(): Promise { + if (!this.handle || !this.nextCursor || this._loading) return; + this._loading = true; + this.rebuild(); + try { + const page = await this.ctx.rest.activity(this.account, { + ...this.restQuery(), + cursor: this.nextCursor, + }); + this.absorb(page.activity, false); + this.nextCursor = page.nextCursor; + this._hasMore = page.nextCursor !== null; + } finally { + this._loading = false; + this.rebuild(); + } + } + + // --- internals --- + + private restQuery() { + return { ...this.query, limit: this.query.limit ?? 100 }; + } + + private async fetchFirstPage(): Promise { + const gen = ++this.seedGen; + const before = this.stream.cursor; + try { + const page = await this.ctx.rest.activity(this.account, this.restQuery()); + // A seed overtaken by a fresher one is dropped whole: it would absorb rows the + // newer page has already corrected, and raise the cursor to an older watermark. + if (!this.stream.running || gen !== this.seedGen) return false; + // Replayed frames can land while this is in flight, and the indexer trails the + // stream by design. Overwriting a row the stream has already advanced would + // revert it permanently: the cursor refuses to rewind and `onFrame` drops a + // re-delivery, so nothing would repair it. + this.absorb(page.activity, compareCursor(this.stream.cursor, before) <= 0); + // Paging is the consumer's position in history, not the seed's: a reconnect + // re-paints the first page, and taking its cursor again would hand back pages + // `loadOlder` has already walked past. + if (!this.painted) { + this.painted = true; + this.nextCursor = page.nextCursor; + this._hasMore = page.nextCursor !== null; + } + this.stream.raise({ since: msToUs(page.solutionNow) }); + this.rebuild(); + return true; + } catch (e) { + // Only before the first paint: an account whose history is genuinely empty has + // a resource, and reporting a later transient failure on it would blank the view. + if (this.handle && this.handle.current() === undefined) this.handle.fail(e as Error); + throw e; + } + } + + private absorb(entries: ActivityEntry[], overwriteOrders: boolean): void { + for (const entry of entries) { + if (entry.activityType === "order") { + if (!overwriteOrders && this.orders.has(entry.order.id)) continue; + this.orders.set(entry.order.id, entry.order); + // Only the source that first reported the order gets to time it. A re-seed + // carries the node's row time for an order the stream already placed by its + // batch, and taking it would move the row on every reconnect. + if (!this.orderTimeMs.has(entry.order.id)) this.orderTimeMs.set(entry.order.id, entry.timeMs); + continue; + } + const key = moneyKey(entry); + if (!this.money.has(key)) this.money.set(key, entry); + } + } + + /** One frame: everything that happened to the account in one auction batch. */ + private onFrame(result: unknown): void { + const frame = result as WireActivityFrame; + if (!frame || !Array.isArray(frame.orders) || !Array.isArray(frame.events)) return; + // Applying a frame is not idempotent — a fill appends to `order.fills` — and + // re-delivery is designed in: a resumed subscription replays from the cursor. + const at: SubParams = { since: frame.batch }; + if (this.stream.delivered(at)) return; + const kept = this.keep(frame); + if (!kept) return; + + const entries: ActivityEntry[] = []; + const events = applyActivityFrame(kept, { orders: this.orders, entries }, { + account: this.account, + }); + // An order the stream is the first to report is timed by its batch; one the + // seed already placed keeps the time the node gave it. + const batchMs = usToMs(kept.batch); + for (const o of kept.orders) if (!this.orderTimeMs.has(o.id)) this.orderTimeMs.set(o.id, batchMs); + this.absorb(entries, false); + this.stream.advance(at); + this.rebuild(); + // Strictly after `rebuild()`: a listener that reads the resource in response to an + // event must see the state that event produced. + if (events.length) { + for (const listener of this.eventListeners) { + try { listener(events); } catch { /* a listener's failure is not the stream's */ } + } + } + } + + /** + * Narrow a frame to what the query asked for, or drop it. REST filters both the + * kinds and the window server-side; the stream filters neither. + */ + private keep(frame: WireActivityFrame): WireActivityFrame | undefined { + const { types, from, to } = this.query; + if (from !== undefined && frame.batch < msToUs(from)) return undefined; + if (to !== undefined && frame.batch >= msToUs(to)) return undefined; + if (!types) return frame; + const want = new Set(types); + return { + ...frame, + // The order events index `orders`, so the two are kept or dropped together. + orders: want.has("order") ? frame.orders : [], + events: frame.events.filter((e) => want.has(isMoney(e.k) ? e.k : "order")), + }; + } + + private rebuild(): void { + if (!this.handle) return; + const arr: ActivityEntry[] = new Array(this.orders.size + this.money.size); + let i = 0; + for (const [id, order] of this.orders) { + arr[i++] = { + activityType: "order", + timeMs: this.orderTimeMs.get(id) ?? Number.MAX_SAFE_INTEGER, + order, + }; + } + for (const entry of this.money.values()) arr[i++] = entry; + arr.sort((a, b) => + b.timeMs - a.timeMs + || ORDINAL[b.activityType] - ORDINAL[a.activityType] + || tiebreak(a, b)); + this.handle.set(arr); + } +} diff --git a/ts-sdk/src/sync/orders.ts b/ts-sdk/src/sync/orders.ts index 58947dc9..18dbe9c6 100644 --- a/ts-sdk/src/sync/orders.ts +++ b/ts-sdk/src/sync/orders.ts @@ -1,109 +1,47 @@ // OrderHistory: a SeriesResource seeded from the warm first page, kept // live by the pod_orders_v2 stream (bidder-filtered, resumed from a (batch, book) -// cursor), and paged backwards by cursor for deep history. -// -// Reconnect: on every (re)connect we re-seed the first page (authoritative open -// orders) and refresh the cursor from its watermark; the WS auto-resubscribes, -// and if the cursor is too old (down too long) `onError` re-seeds and -// resubscribes. A server-initiated close is not the same as a rejection: it -// reports where delivery stopped, so resuming from one is a bare `resubscribe()` -// with no re-seed. +// cursor by `ResumableStream`), and paged backwards by cursor for deep history. -import type { Address, MarketId, Order, OrderEvent, OrdersQuery } from "../types/public.js"; +import type { Address, Order, OrderEvent, OrdersQuery } from "../types/public.js"; import type { WireOrdersFrame } from "../types/wire.js"; import { applyOrdersFrame } from "../codec/orders-v2.js"; import { BaseResource, type ResourceHandle } from "../stores/resource.js"; -import { PodSubscriptionClosedError, type SubParams, type Subscription } from "../transport/ws.js"; +import type { SubParams } from "../transport/ws.js"; import type { SeriesResource } from "./candles.js"; import type { SyncContext } from "./sources.js"; +import { compareCursor, ResumableStream } from "./stream.js"; -/** - * Consecutive server closes we resume from before falling back to the re-seed - * path. A lagged close is recoverable by resubscribing, so the fast path is the - * right default — but a stream that keeps closing is one we are not keeping up - * with, and re-seeding beats replaying a growing backlog on every attempt. - */ -const FAST_RESUMES_BEFORE_RESEED = 3; - -/** - * Re-seed attempts that keep the cursor before we conclude the cursor is what the - * server is rejecting, drop it, and settle for streaming live. - */ -const RETRIES_BEFORE_DROPPING_CURSOR = 2; - -/** The largest book id, which is what an absent `sinceBook` means. */ -const WHOLE_BATCH = `0x${"ff".repeat(32)}` as MarketId; - -/** - * Order two positions in the stream, the way the server does. - * - * A batch is delivered as one frame per book, so a position is the pair - * `(batch, book)` — and an absent book means the whole batch, which sorts *above* - * every book in it. This mirrors `already_delivered` in - * `node/src/rpc/orders_v2.rs`: `(frame_batch, frame_book) <= (since, - * since_book.unwrap_or(0xff…))`. Getting the absent case wrong turns "all of batch - * N" into "up to book B of batch N", which asks the server to re-send the rest. - * - * Book ids are fixed-width lowercase hex, so comparing them as strings is - * comparing their bytes. - */ -export function compareCursor(a: SubParams, b: SubParams): number { - const aBatch = a.since ?? 0; - const bBatch = b.since ?? 0; - if (aBatch !== bBatch) return aBatch < bBatch ? -1 : 1; - const aBook = a.sinceBook ?? WHOLE_BATCH; - const bBook = b.sinceBook ?? WHOLE_BATCH; - return aBook === bBook ? 0 : aBook < bBook ? -1 : 1; -} +export { compareCursor } from "./stream.js"; export class OrderHistory implements SeriesResource { private readonly base: BaseResource; private handle: ResourceHandle | undefined; private readonly byId = new Map(); private nextCursor: string | null = null; - private sub: Subscription | undefined; - private alive = false; + private painted = false; + private seedGen = 0; private _hasMore = true; private _loading = false; - private subRetries = 0; - private fastResumes = 0; - private retryTimer?: ReturnType; private readonly eventListeners = new Set<(events: OrderEvent[]) => void>(); - /** - * Where the stream is: the last frame accepted, as the `(batch, book)` pair the - * channel resumes from. Replaced whole rather than patched half at a time — - * `sinceBook` names a book *within* `since`, so the two are one fact and a - * mismatched pair asks the server to skip books we never saw. - */ - private cursor: SubParams = { since: 0, sinceBook: undefined }; - - /** - * Push the cursor to the transport, both halves. - * - * `update` merges, so sending `{ since }` alone would leave whatever `sinceBook` was - * there before — a book from an older batch beside a newer `since`, which asks the - * server to skip less than it should and re-send the difference. - */ - private pushCursor(): void { - this.sub?.update({ since: this.cursor.since, sinceBook: this.cursor.sinceBook }); - } + private readonly stream: ResumableStream; constructor( private readonly ctx: SyncContext, private readonly account: Address, private readonly query: OrdersQuery = {}, ) { + this.stream = new ResumableStream({ + ws: ctx.ws, + channel: "pod_orders_v2", + params: { account }, + reseed: () => this.fetchFirstPage(), + onFrame: (r) => this.onFrame(r), + }); this.base = new BaseResource((h) => { this.handle = h; - this.alive = true; - const offOpen = this.ctx.ws.on("open", () => { if (this.alive) this.seed(); }); - this.seed(); // initial paint (REST is independent of the socket being open) + this.stream.start(); return () => { - this.alive = false; - offOpen(); - if (this.retryTimer) clearTimeout(this.retryTimer); - this.sub?.unsubscribe(); - this.sub = undefined; + this.stream.stop(); this.handle = undefined; }; }); @@ -160,143 +98,59 @@ export class OrderHistory implements SeriesResource { // --- internals --- private async fetchFirstPage(): Promise { - const before = this.cursor; + const gen = ++this.seedGen; + const before = this.stream.cursor; try { const page = await this.ctx.rest.orders(this.account, { limit: this.query.limit ?? 100 }); - if (!this.alive) return false; + // A seed overtaken by a fresher one is dropped whole: it would absorb rows the + // newer page has already corrected, and raise the cursor to an older watermark. + if (!this.stream.running || gen !== this.seedGen) return false; // The transport resubscribes synchronously right after the `open` event that // starts this fetch, so replayed frames can land while it is still in flight — // and the indexer trails the stream by design. Overwriting a row the stream has - // already advanced would revert it, permanently: the cursor below refuses to - // rewind and `onFrame` drops a re-delivery, so nothing would repair it. - const streamMovedOn = compareCursor(this.cursor, before) > 0; + // already advanced would revert it, permanently: the cursor refuses to rewind + // and `onFrame` drops a re-delivery, so nothing would repair it. + const streamMovedOn = compareCursor(this.stream.cursor, before) > 0; for (const o of page.orders) { if (streamMovedOn && this.byId.has(o.id)) continue; this.byId.set(o.id, o); } - this.nextCursor = page.nextCursor; - this._hasMore = page.nextCursor !== null; + // Paging is the consumer's position in history, not the seed's: a reconnect + // re-paints the first page, and taking its cursor again would hand back pages + // `loadOlder` has already walked past. + if (!this.painted) { + this.painted = true; + this.nextCursor = page.nextCursor; + this._hasMore = page.nextCursor !== null; + } // Only ever forward, and a page settles whole batches, so its watermark // carries no book — which makes it *ahead* of a stream position in the same // batch, not equal to it. - const settled: SubParams = { since: page.solutionNow * 1000, sinceBook: undefined }; - if (compareCursor(settled, this.cursor) > 0) this.cursor = settled; + this.stream.raise({ since: page.solutionNow * 1000, sinceBook: undefined }); this.rebuild(); return true; } catch (e) { - if (this.byId.size === 0) this.handle?.fail(e as Error); - return false; - } - } - - private seed(): void { - void this.fetchFirstPage().then((ok) => { - if (!ok || !this.alive) return; - if (!this.sub) { - this.sub = this.ctx.ws.subscribe( - "pod_orders_v2", - { account: this.account, ...this.cursor }, - (r) => this.onFrame(r), - (e) => this.onSubError(e), - ); - } else { - this.pushCursor(); // refresh for the next reconnect - } - }); - } - - /** - * The subscription is not running: `eth_subscribe` was rejected, or the server - * closed it. - * - * A resumable close needs neither a re-seed nor a delay: the server reports where - * it stopped, so resubscribing from that delivers exactly the frames we never got. - * That fast path is only taken while the socket is actually open — `resubscribe()` - * is a no-op otherwise, which would spend the budget without an attempt and leave - * nothing scheduled, and `-32021` (node shutting down) arrives exactly as the - * socket goes away. - * - * Everything else (a rejection, a server bug, a close with the socket already - * gone, or closes that keep coming) takes the slow path: backed off and capped, so - * a server that keeps refusing cannot spin this into a tight re-seed loop, and - * eventually the cursor is dropped (the likely culprit) to just stream live. Both - * counters reset once live data flows (`onFrame`). - */ - private onSubError(err: unknown): void { - if (!this.alive) return; - const canFastResume = err instanceof PodSubscriptionClosedError && err.resumable - && this.fastResumes < FAST_RESUMES_BEFORE_RESEED && this.ctx.ws.state === "open"; - if (canFastResume) { - this.fastResumes++; - // Adopt the server's watermark only when it is ahead of ours: it knows which - // frames it handed over, but a re-seed may already have carried us past it, - // and rewinding would re-deliver frames we have applied. - const reported: SubParams = { since: err.resumeSince, sinceBook: err.resumeSinceBook }; - if (err.resumeSince !== undefined && compareCursor(reported, this.cursor) > 0) this.cursor = reported; - // The transport rewrote `sub.params` from the close before this ran, so without - // pushing our own decision back the wire resumes from the server's point - // regardless and the guard above protects nothing. - this.pushCursor(); - this.sub?.resubscribe(); - return; + // Only before the first paint: an account whose history is genuinely empty has + // a resource, and reporting a later transient failure on it would blank the view. + if (this.handle && this.handle.current() === undefined) this.handle.fail(e as Error); + throw e; } - // The slow path, and the one that answers "what if we fell behind the server's - // replay buffer": `eth_subscribe` rejects a `since` older than what the buffer - // retains, and the prescribed recovery is to backfill over REST and resubscribe. - // That is what this is. `fetchFirstPage` also *replaces the cursor* with the - // page's watermark, so the position that was too old is gone by the first retry - // and the resubscribe is accepted with no gap — the server replays from the page - // forward. Dropping the cursor below is the backstop for when even that fresh - // watermark is refused (the indexer further behind than the buffer retains): - // live-only resubscribe, trading the unreplayable window for a working stream. - this.scheduleReseed(); - } - - /** - * Back off, re-seed over REST, then resubscribe — and keep trying. - * - * The re-seed can fail too (the same node is usually behind both the stream and the - * indexer), and a failure has to re-arm here: not resubscribing means no further - * close or rejection arrives, so nothing else would ever schedule another attempt and - * the stream would stay down with no error surfaced. Capped, so a node that keeps - * refusing cannot spin this. - */ - private scheduleReseed(): void { - this.subRetries++; - const delay = Math.min(30_000, 500 * 2 ** (this.subRetries - 1)); - if (this.retryTimer) clearTimeout(this.retryTimer); - this.retryTimer = setTimeout(() => { - void this.fetchFirstPage().then((ok) => { - if (!this.alive) return; - if (!ok) { this.scheduleReseed(); return; } - const tooOld = this.subRetries > RETRIES_BEFORE_DROPPING_CURSOR; - if (tooOld) this.sub?.update({ since: undefined, sinceBook: undefined }); - else this.pushCursor(); - this.sub?.resubscribe(); - }); - }, delay); } /** One frame: everything that happened to one book in one auction batch. */ private onFrame(result: unknown): void { - // Live data flowing → the subscription is healthy on both paths. - this.subRetries = 0; - this.fastResumes = 0; const frame = result as WireOrdersFrame; if (!frame || !Array.isArray(frame.orders) || !Array.isArray(frame.events)) return; // Drop what we already hold. Applying a frame is not idempotent — a fill event // appends to `order.fills`, and re-creating an entity resets its totals — and // re-delivery is designed in: a resumed subscription replays from a cursor, and - // the replay boundary is a whole batch. Same predicate as the server's - // `already_delivered`, so client and server agree on what "already sent" means. + // the replay boundary is a whole batch. const at: SubParams = { since: frame.batch, sinceBook: frame.book }; - if (compareCursor(at, this.cursor) <= 0) return; + if (this.stream.delivered(at)) return; const events = applyOrdersFrame(frame, this.byId, { account: this.account }); - // A socket-level reconnect resubscribes from whatever is stored here. - this.cursor = at; - this.sub?.update(at); + this.stream.advance(at); this.rebuild(); // Strictly after `rebuild()`: a listener that reads the resource in response to an // event must see the state that event produced, not the state before it. diff --git a/ts-sdk/src/sync/stream.test.ts b/ts-sdk/src/sync/stream.test.ts new file mode 100644 index 00000000..b8f03319 --- /dev/null +++ b/ts-sdk/src/sync/stream.test.ts @@ -0,0 +1,226 @@ +// `ResumableStream` on its own: the subscribe / resume / re-seed loop both +// account feeds share. Every rule here is a recovery path the happy-path feed +// tests never reach, and none of it is reachable from `typecheck` — a close +// frame is `unknown` off the socket. + +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +import { ResumableStream } from "./stream.js"; +import { PodHttpError } from "../transport/rest.js"; +import { PodSubscriptionClosedError, type PodWsClient, type SubParams } from "../transport/ws.js"; +import type { Address, MarketId } from "../types/public.js"; + +const ACCOUNT = "0xa1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1" as Address; +const book = (n: number) => `0x${n.toString(16).padStart(64, "0")}` as MarketId; +const B1 = book(1); +const B9 = book(9); +const SEED = 3_000_000; + +/** + * A stream over a stub socket. `seeds` scripts the re-seeds in order — a number + * is a page whose watermark raises the cursor to it, an `Error` is a failed + * fetch; the last entry repeats for every attempt after it. + */ +function harness(seeds: Array = [SEED]) { + const queue = [...seeds]; + let onOpen: (() => void) | undefined; + let deliver: ((r: unknown) => void) | undefined; + let refuse: ((e: unknown) => void) | undefined; + const subscribed: SubParams[] = []; + const updates: SubParams[] = []; + const frames: unknown[] = []; + let resubscribes = 0; + let unsubscribes = 0; + let seeds_ = 0; + + const ws = { + state: "open", + on: (_event: string, handler: () => void) => { + onOpen = handler; + return () => { onOpen = undefined; }; + }, + subscribe: ( + _channel: string, + params: SubParams, + onMessage: (r: unknown) => void, + onError: (e: unknown) => void, + ) => { + subscribed.push({ ...params }); + deliver = onMessage; + refuse = onError; + return { + unsubscribe: () => { unsubscribes++; }, + update: (p: SubParams) => { updates.push(p); }, + resubscribe: () => { resubscribes++; }, + }; + }, + }; + + const stream = new ResumableStream({ + ws: ws as unknown as PodWsClient, + channel: "pod_orders_v2", + params: { account: ACCOUNT }, + reseed: async () => { + seeds_++; + const next = queue.length > 1 ? queue.shift()! : queue[0] ?? SEED; + if (next instanceof Error) throw next; + stream.raise({ since: next, sinceBook: undefined }); + return true; + }, + onFrame: (r) => { frames.push(r); }, + }); + + return { + stream, + subscribed, + updates, + frames, + open: () => onOpen?.(), + frame: (f: unknown) => deliver?.(f), + close: (e: unknown) => refuse?.(e), + seeds: () => seeds_, + resubscribes: () => resubscribes, + unsubscribes: () => unsubscribes, + }; +} + +const closed = (resumable: boolean, since?: number, atBook?: MarketId) => + new PodSubscriptionClosedError({ + code: resumable ? -32020 : -32023, + data: { resumable, resume_since: since, resume_since_book: atBook }, + }); + +beforeEach(() => { vi.useFakeTimers(); }); +afterEach(() => { vi.useRealTimers(); }); + +describe("ResumableStream cursor", () => { + it("subscribes from the seeded watermark and carries both halves forward", async () => { + const h = harness(); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + expect(h.subscribed).toEqual([{ account: ACCOUNT, since: SEED }]); + + h.stream.advance({ since: SEED + 1, sinceBook: B1 }); + expect(h.updates.at(-1)).toEqual({ since: SEED + 1, sinceBook: B1 }); + + // Never rewinds, and an absent book is the whole batch — above every book in it. + h.stream.raise({ since: SEED, sinceBook: undefined }); + expect(h.stream.cursor).toEqual({ since: SEED + 1, sinceBook: B1 }); + h.stream.raise({ since: SEED + 1, sinceBook: undefined }); + expect(h.stream.cursor.sinceBook).toBeUndefined(); + expect(h.stream.delivered({ since: SEED + 1, sinceBook: B9 })).toBe(true); + }); +}); + +describe("ResumableStream resume", () => { + it("resubscribes from the server's watermark on a resumable close, without re-seeding", async () => { + const h = harness(); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + + h.close(closed(true, SEED + 500_000, B1)); + expect(h.resubscribes()).toBe(1); + expect(h.seeds()).toBe(1); + // Both halves pushed back: `update` merges, so a lone `since` would leave a + // book from an older batch beside a newer batch. + expect(h.updates.at(-1)).toEqual({ since: SEED + 500_000, sinceBook: B1 }); + }); + + it("falls back to a re-seed once the fast resumes are spent", async () => { + const h = harness([SEED, SEED + 9_000_000]); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + + for (let i = 0; i < 3; i++) h.close(closed(true, SEED + i, undefined)); + expect(h.resubscribes()).toBe(3); + expect(h.seeds()).toBe(1); + + h.close(closed(true, SEED + 9, undefined)); + await vi.advanceTimersByTimeAsync(1_000); + expect(h.seeds(), "the fourth close re-seeds instead").toBe(2); + expect(h.resubscribes()).toBe(4); + }); + + it("drops a cursor the server keeps refusing", async () => { + const h = harness(); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + + for (let i = 0; i < 3; i++) { + h.close(closed(false)); + await vi.advanceTimersByTimeAsync(5_000); + } + // Two attempts keep the cursor; the third concludes the cursor is the problem. + expect(h.updates.at(-1)).toEqual({ since: undefined, sinceBook: undefined }); + expect(h.updates.slice(0, -1).every((u) => u.since !== undefined)).toBe(true); + }); + + it("resets the retry budget once a retried first paint subscribes", async () => { + const h = harness([new Error("down"), new Error("down"), SEED]); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + await vi.advanceTimersByTimeAsync(500); + await vi.advanceTimersByTimeAsync(1_000); + expect(h.subscribed).toHaveLength(1); + + // A live stream's first refusal must not inherit the seed's attempts. + h.close(closed(false)); + await vi.advanceTimersByTimeAsync(5_000); + expect(h.updates.at(-1)?.since).toBe(SEED); + }); + + it("cancels a pending retry when a later seed succeeds on its own", async () => { + const h = harness([new Error("down"), SEED]); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + expect(h.seeds()).toBe(1); + + // The socket came back before the retry was due; its seed is the one that lands. + h.open(); + await vi.advanceTimersByTimeAsync(0); + expect(h.seeds()).toBe(2); + expect(h.subscribed).toHaveLength(1); + + await vi.advanceTimersByTimeAsync(60_000); + expect(h.seeds(), "the orphaned retry must not fire").toBe(2); + expect(h.resubscribes()).toBe(0); + }); + + it("gives up on a seed the server refused on its own terms", async () => { + const h = harness([new PodHttpError(400, "http://node.test", "bad request")]); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + + await vi.advanceTimersByTimeAsync(60_000); + expect(h.seeds()).toBe(1); + expect(h.subscribed).toHaveLength(0); + }); + + it("treats a live frame as the all-clear", async () => { + const h = harness(); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + + h.close(closed(false)); + h.frame({ batch: SEED + 1 }); + + await vi.advanceTimersByTimeAsync(60_000); + expect(h.frames).toHaveLength(1); + expect(h.seeds(), "there is nothing for the armed retry to recover").toBe(1); + }); + + it("tears down on stop, and stays down", async () => { + const h = harness(); + h.stream.start(); + await vi.advanceTimersByTimeAsync(0); + + h.stream.stop(); + expect(h.unsubscribes()).toBe(1); + expect(h.stream.running).toBe(false); + + h.open(); + await vi.advanceTimersByTimeAsync(60_000); + expect(h.seeds()).toBe(1); + expect(h.subscribed).toHaveLength(1); + }); +}); diff --git a/ts-sdk/src/sync/stream.ts b/ts-sdk/src/sync/stream.ts new file mode 100644 index 00000000..b3ffa114 --- /dev/null +++ b/ts-sdk/src/sync/stream.ts @@ -0,0 +1,270 @@ +// The subscribe / resume / re-seed loop shared by the account feeds. +// +// Reconnect: on every (re)connect the owner re-seeds its first page (the +// authoritative snapshot) and refreshes the cursor from its watermark; the WS +// auto-resubscribes, and if the cursor is too old (down too long) `onSubError` +// re-seeds and resubscribes. A server-initiated close is not the same as a +// rejection: it reports where delivery stopped, so resuming from one is a bare +// `resubscribe()` with no re-seed. + +import type { MarketId } from "../types/public.js"; +import { PodHttpError } from "../transport/rest.js"; +import type { Channel, PodWsClient, SubParams, Subscription } from "../transport/ws.js"; +import { PodSubscriptionClosedError } from "../transport/ws.js"; + +/** + * Consecutive server closes we resume from before falling back to the re-seed + * path. A lagged close is recoverable by resubscribing, so the fast path is the + * right default — but a stream that keeps closing is one we are not keeping up + * with, and re-seeding beats replaying a growing backlog on every attempt. + */ +const FAST_RESUMES_BEFORE_RESEED = 3; + +/** + * Re-seed attempts that keep the cursor before we conclude the cursor is what the + * server is rejecting, drop it, and settle for streaming live. + */ +const RETRIES_BEFORE_DROPPING_CURSOR = 2; + +/** The largest book id, which is what an absent `sinceBook` means. */ +const WHOLE_BATCH = `0x${"ff".repeat(32)}` as MarketId; + +/** + * Order two positions in the stream, the way the server does. + * + * A batch is delivered as one frame per book, so a position is the pair + * `(batch, book)` — and an absent book means the whole batch, which sorts *above* + * every book in it. This mirrors `already_delivered` in + * `node/src/rpc/orders_v2.rs`: `(frame_batch, frame_book) <= (since, + * since_book.unwrap_or(0xff…))`. Getting the absent case wrong turns "all of batch + * N" into "up to book B of batch N", which asks the server to re-send the rest. + * + * `pod_activity` sends one frame per tick and names no book, so the pair collapses + * to the batch there — the same comparison, with both books absent. + * + * Book ids are fixed-width lowercase hex, so comparing them as strings is + * comparing their bytes. + */ +export function compareCursor(a: SubParams, b: SubParams): number { + const aBatch = a.since ?? 0; + const bBatch = b.since ?? 0; + if (aBatch !== bBatch) return aBatch < bBatch ? -1 : 1; + const aBook = a.sinceBook ?? WHOLE_BATCH; + const bBook = b.sinceBook ?? WHOLE_BATCH; + return aBook === bBook ? 0 : aBook < bBook ? -1 : 1; +} + +/** + * A seed the server refused on its own terms. Retrying reproduces it, so the + * backoff would be an infinite loop against an answer that will not change. + */ +const permanent = (err: unknown): boolean => + err instanceof PodHttpError && err.status >= 400 && err.status < 500; + +export interface ResumableStreamOptions { + ws: PodWsClient; + channel: Channel; + /** Everything that identifies the subscription except the cursor. */ + params: SubParams; + /** + * Re-seed over REST, and raise the cursor to the page's watermark. `true` when + * the page was applied, `false` when it was abandoned (the stream stopped, or a + * fresher seed overtook it); a rejection is a failed fetch, which is what decides + * between backing off and giving up. + */ + reseed(): Promise; + onFrame(result: unknown): void; +} + +export class ResumableStream { + /** + * Where the stream is: the last frame accepted. Replaced whole rather than + * patched half at a time — on `pod_orders_v2` `sinceBook` names a book *within* + * `since`, so the two are one fact and a mismatched pair asks the server to skip + * books we never saw. + */ + cursor: SubParams = {}; + private sub: Subscription | undefined; + private alive = false; + private subRetries = 0; + private fastResumes = 0; + private retryTimer?: ReturnType; + private offOpen?: () => void; + + constructor(private readonly opts: ResumableStreamOptions) {} + + /** Whether the owner's resource is still running; guards its async seeds. */ + get running(): boolean { return this.alive; } + + start(): void { + this.alive = true; + this.offOpen = this.opts.ws.on("open", () => { if (this.alive) this.seed(); }); + this.seed(); // initial paint (REST is independent of the socket being open) + } + + stop(): void { + this.alive = false; + this.offOpen?.(); + this.offOpen = undefined; + this.clearRetry(); + this.sub?.unsubscribe(); + this.sub = undefined; + } + + /** The server's `already_delivered`: is this frame at or behind where we are? */ + delivered(at: SubParams): boolean { + return compareCursor(at, this.cursor) <= 0; + } + + /** Move to `at` and tell the transport, so a reconnect resumes from there. */ + advance(at: SubParams): void { + this.cursor = at; + this.sub?.update(at); + } + + /** Move to `at` only when it is ahead. Never rewinds — a re-delivered frame is dropped. */ + raise(at: SubParams): void { + if (compareCursor(at, this.cursor) > 0) this.cursor = at; + } + + seed(): void { + // A seed that lands supersedes whatever the last failure armed; leaving the + // timer would re-seed and resubscribe a stream that is already healthy. + const armed = this.retryTimer !== undefined; + this.clearRetry(); + void this.opts.reseed().then( + (ok) => { if (ok) this.resume(); }, + (err) => { + // Nothing is subscribed yet, so no close or rejection can arrive to schedule + // another attempt — a failed first paint has to re-arm itself or the resource + // stays empty for good. With a subscription up, the live stream is unaffected + // and its own error path owns the recovery; the exception is a retry this call + // just disarmed, which nothing else would re-arm. + if (this.alive && (!this.sub || armed) && !permanent(err)) this.scheduleReseed(); + }, + ); + } + + private resume(): void { + if (!this.alive) return; + if (!this.sub) this.subscribe(); + else this.pushCursor(); // refresh for the next reconnect + } + + private subscribe(): void { + this.sub = this.opts.ws.subscribe( + this.opts.channel, + { ...this.opts.params, ...this.cursor }, + (r) => this.onFrame(r), + (e) => this.onSubError(e), + ); + } + + /** + * Push the cursor to the transport, both halves. + * + * `update` merges, so sending `{ since }` alone would leave whatever `sinceBook` + * was there before — a book from an older batch beside a newer `since`, which asks + * the server to skip less than it should and re-send the difference. + */ + private pushCursor(): void { + this.sub?.update({ since: this.cursor.since, sinceBook: this.cursor.sinceBook }); + } + + private clearRetry(): void { + if (this.retryTimer) { clearTimeout(this.retryTimer); this.retryTimer = undefined; } + } + + private onFrame(result: unknown): void { + // Live data flowing → the subscription is healthy on both paths, and a retry + // armed by whatever went wrong before has nothing left to recover. + this.subRetries = 0; + this.fastResumes = 0; + this.clearRetry(); + this.opts.onFrame(result); + } + + /** + * The subscription is not running: `eth_subscribe` was rejected, or the server + * closed it. + * + * A resumable close needs neither a re-seed nor a delay: the server reports where + * it stopped, so resubscribing from that delivers exactly the frames we never got. + * That fast path is only taken while the socket is actually open — `resubscribe()` + * is a no-op otherwise, which would spend the budget without an attempt and leave + * nothing scheduled, and `-32021` (node shutting down) arrives exactly as the + * socket goes away. + * + * Everything else (a rejection, a server bug, a close with the socket already + * gone, or closes that keep coming) takes the slow path: backed off and capped, so + * a server that keeps refusing cannot spin this into a tight re-seed loop, and + * eventually the cursor is dropped (the likely culprit) to just stream live. Both + * counters reset once live data flows (`onFrame`). + */ + private onSubError(err: unknown): void { + if (!this.alive) return; + const canFastResume = err instanceof PodSubscriptionClosedError && err.resumable + && this.fastResumes < FAST_RESUMES_BEFORE_RESEED && this.opts.ws.state === "open"; + if (canFastResume) { + this.fastResumes++; + // Adopt the server's watermark only when it is ahead of ours: it knows which + // frames it handed over, but a re-seed may already have carried us past it, + // and rewinding would re-deliver frames we have applied. + const reported: SubParams = { since: err.resumeSince, sinceBook: err.resumeSinceBook }; + if (err.resumeSince !== undefined) this.raise(reported); + // The transport rewrote `sub.params` from the close before this ran, so without + // pushing our own decision back the wire resumes from the server's point + // regardless and the guard above protects nothing. + this.pushCursor(); + this.sub?.resubscribe(); + return; + } + // The slow path, and the one that answers "what if we fell behind the server's + // replay buffer": `eth_subscribe` rejects a `since` older than what the buffer + // retains, and the prescribed recovery is to backfill over REST and resubscribe. + // That is what this is. The re-seed also *raises the cursor* to the page's + // watermark, so the position that was too old is gone by the first retry and the + // resubscribe is accepted with no gap — the server replays from the page forward. + // Dropping the cursor below is the backstop for when even that fresh watermark is + // refused (the indexer further behind than the buffer retains): live-only + // resubscribe, trading the unreplayable window for a working stream. + this.scheduleReseed(); + } + + /** + * Back off, re-seed over REST, then resubscribe — and keep trying. + * + * The re-seed can fail too (the same node is usually behind both the stream and the + * indexer), and a failure has to re-arm here: not resubscribing means no further + * close or rejection arrives, so nothing else would ever schedule another attempt and + * the stream would stay down with no error surfaced. Capped, so a node that keeps + * refusing cannot spin this. + */ + private scheduleReseed(): void { + this.subRetries++; + const delay = Math.min(30_000, 500 * 2 ** (this.subRetries - 1)); + this.clearRetry(); + this.retryTimer = setTimeout(() => { + this.retryTimer = undefined; + void this.opts.reseed().then((ok) => { + if (!this.alive || !ok) return; + if (!this.sub) { + // Retrying the first paint: there is no subscription to resume, and no + // cursor the server has refused — open one from the page we just landed, + // and let the live stream start on a full budget rather than inheriting + // the attempts the indexer's outage cost. + this.subRetries = 0; + this.fastResumes = 0; + this.subscribe(); + return; + } + const tooOld = this.subRetries > RETRIES_BEFORE_DROPPING_CURSOR; + if (tooOld) this.sub.update({ since: undefined, sinceBook: undefined }); + else this.pushCursor(); + this.sub.resubscribe(); + }, (err) => { + if (this.alive && !permanent(err)) this.scheduleReseed(); + }); + }, delay); + } +} diff --git a/ts-sdk/src/transport/rest.test.ts b/ts-sdk/src/transport/rest.test.ts index a29b2188..31ce9c91 100644 --- a/ts-sdk/src/transport/rest.test.ts +++ b/ts-sdk/src/transport/rest.test.ts @@ -32,3 +32,48 @@ describe("PodRestClient request timeout", () => { await expect(rest.status()).rejects.toThrow(); }); }); + +describe("PodRestClient.activity", () => { + const ACCOUNT = "0xa1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1a1"; + const served = (body: unknown) => { + let url = ""; + const rest = new PodRestClient({ + restUrl: "http://node.test/v1", + fetch: ((u: string) => { + url = u; + return Promise.resolve(new Response(JSON.stringify(body), { status: 200 })); + }) as unknown as typeof fetch, + }); + return { rest, query: () => new URL(url).searchParams }; + }; + + it("joins the activity types and sends the window in micros", async () => { + const { rest, query } = served({ activity: [], next_cursor: null, solution_now: 1 }); + await rest.activity(ACCOUNT as never, { + types: ["transfer", "order"], from: 1_000, to: 2_000, limit: 25, + }); + expect(query().get("activity_types")).toBe("transfer,order"); + expect(query().get("from")).toBe("1000000"); + expect(query().get("to")).toBe("2000000"); + expect(query().get("limit")).toBe("25"); + }); + + it("omits an empty type list rather than asking for nothing", async () => { + const { rest, query } = served({ activity: [], next_cursor: null, solution_now: 1 }); + await rest.activity(ACCOUNT as never, { types: [] }); + expect(query().get("activity_types")).toBeNull(); + }); + + it("drops a row whose kind this version does not know", async () => { + const { rest } = served({ + activity: [ + { activity_type: "transfer", ts: 5_000_000, id: "0x03", token: "0x7e", amount: "1700" }, + { activity_type: "airdrop", ts: 5_000_000, amount: "1" }, + ], + next_cursor: null, + solution_now: 5_000_000, + }); + const page = await rest.activity(ACCOUNT as never); + expect(page.activity.map((e) => e.activityType)).toEqual(["transfer"]); + }); +}); diff --git a/ts-sdk/src/transport/rest.ts b/ts-sdk/src/transport/rest.ts index 47f276e4..8d66c645 100644 --- a/ts-sdk/src/transport/rest.ts +++ b/ts-sdk/src/transport/rest.ts @@ -3,21 +3,21 @@ // ms in, decoded (bigint / ms) out. import type { - Address, BackstopTransfer, Balances, Bar, BridgeConfig, CandleQuery, + ActivityEntry, ActivityType, Address, BackstopTransfer, Balances, Bar, BridgeConfig, CandleQuery, Market, MarketId, Order, Orderbook, PositionsSnapshot, Resolution, Status, Trigger, TxExplorer, Withdrawal, WithdrawalsQuery, } from "../types/public.js"; import type { - WireBackstopPage, WireBalances, WireBridgeConfig, WireCandlesEnvelope, + WireActivityPage, WireBackstopPage, WireBalances, WireBridgeConfig, WireCandlesEnvelope, WireMarketStatic, WireMarketStatsPage, WireOrderbook, WireOrdersPage, WirePositionsSnapshot, WireStatus, WireTriggersPage, WireWithdrawal, } from "../types/wire.js"; import { - decodeBackstopTransfer, decodeBalances, decodeBridgeConfig, decodeCandle, + decodeActivityEntry, decodeBackstopTransfer, decodeBalances, decodeBridgeConfig, decodeCandle, decodeMarketDynamics, decodeMarketStatic, decodeOrder, decodeOrderbook, decodePositions, decodeStatus, decodeTrigger, decodeWithdrawal, } from "../codec/decode.js"; -import { msToSecs, usToMs } from "../codec/units.js"; +import { msToSecs, msToUs, usToMs } from "../codec/units.js"; export type MarketDynamicsPatch = Partial & { id: string }; @@ -31,6 +31,11 @@ export interface OrdersPage { totalCount: number; solutionNow: number; } +export interface ActivityPage { + activity: ActivityEntry[]; + nextCursor: string | null; + solutionNow: number; // ms +} export interface BackstopPage { transfers: BackstopTransfer[]; totalCount: number; @@ -53,6 +58,13 @@ export interface OrdersQueryRest { limit?: number; cursor?: string; } +export interface ActivityQueryRest { + limit?: number; + cursor?: string; + types?: ActivityType[]; + from?: number; // ms + to?: number; // ms +} export interface TriggersQueryRest { orderbook?: MarketId; limit?: number; @@ -177,6 +189,26 @@ export class PodRestClient { }; } + /** + * One account's activity, newest first: orders at their placement tick, and the + * money that moved (ADR 0057 §5). The seed behind `pod_activity`; page it + * with `cursor`. + */ + async activity(account: Address, q?: ActivityQueryRest): Promise { + const w = await this.get(`/clob/activity/${account}`, { + limit: q?.limit, + cursor: q?.cursor, + activity_types: q?.types?.length ? q.types.join(",") : undefined, + from: q?.from !== undefined ? msToUs(q.from) : undefined, + to: q?.to !== undefined ? msToUs(q.to) : undefined, + }); + return { + activity: w.activity.map(decodeActivityEntry).filter((e): e is ActivityEntry => e !== undefined), + nextCursor: w.next_cursor, + solutionNow: usToMs(w.solution_now), + }; + } + async backstopTransfers(account: Address): Promise { const w = await this.get(`/clob/backstop-transfers/${account}`); return { diff --git a/ts-sdk/src/transport/ws.ts b/ts-sdk/src/transport/ws.ts index 82506bf8..ef29c7eb 100644 --- a/ts-sdk/src/transport/ws.ts +++ b/ts-sdk/src/transport/ws.ts @@ -10,7 +10,11 @@ export type Channel = | "pod_markets" | "pod_positions" | "pod_triggers" /** Terminal withdrawal outcomes; one array per tick, `account`-filtered on the * debited account. ADR 0033 §6. */ - | "pod_withdrawals"; + | "pod_withdrawals" + /** One account's whole activity: `pod_orders_v2` for every book the tick + * cleared, plus the money it moved. One frame per tick, so `since` alone + * resumes it. ADR 0057 §6. */ + | "pod_activity"; export interface SubParams { clobIds?: MarketId[]; diff --git a/ts-sdk/src/types/public.ts b/ts-sdk/src/types/public.ts index fb821d14..a684cb1d 100644 --- a/ts-sdk/src/types/public.ts +++ b/ts-sdk/src/types/public.ts @@ -359,6 +359,8 @@ export interface BackstopTransfer { cash: bigint; markPrice: bigint; equity: bigint; + /** Crystallized by the forced close; `0n` on the terminal cash sweep. */ + realizedPnl: bigint; time: number; } @@ -449,6 +451,66 @@ export interface CandleQuery { limit?: number; } +export type ActivityType = "order" | "backstop" | "bridge_transfer" | "transfer"; + +/** + * One row of an account's activity (ADR 0057): an order at its placement tick, + * or money that moved. + * + * `timeMs` is what the row is ordered by. Fills, cancels and amendments are not + * rows — they are events on the stream, and the order an `order` row carries + * already reflects them. + * + * Money amounts are signed from the account's side — negative left, positive + * arrived — and say nothing about the other end. + */ +export type ActivityEntry = + | { activityType: "order"; timeMs: number; order: Order } + | (BackstopTransfer & { activityType: "backstop"; timeMs: number }) + | { + activityType: "bridge_transfer"; + timeMs: number; + txHash: Hash; + /** Position within the transaction: one tx can bridge several amounts, so + * the hash alone does not identify the row. */ + idx: number; + token: Address; + amount: bigint; + /** Absent when the movement settled. */ + error?: string; + } + | { + activityType: "transfer"; + timeMs: number; + transferId: Hash; + token: Address; + amount: bigint; + error?: string; + }; + +/** The activity rows that are not orders. */ +export type MoneyActivity = Exclude; + +/** + * What the stream reported in one tick: the order transitions, and the money + * that moved, which is reported as the entry it produced. + * + * Discriminate with `"activityType" in event` — an order event is tagged by + * `kind` and a money one by `activityType`. + */ +export type ActivityEvent = OrderEvent | MoneyActivity; + +/** + * `from` and `to` are milliseconds; `limit` is a row count, not a span of time. + * `types` names the kinds to keep — absent or empty keeps every kind. + */ +export interface ActivityQuery { + types?: ActivityType[]; + from?: number; + to?: number; + limit?: number; +} + export interface OrdersQuery { status?: OrderStatus; orderbookId?: MarketId; diff --git a/ts-sdk/src/types/wire.ts b/ts-sdk/src/types/wire.ts index 192bd00c..350e94ea 100644 --- a/ts-sdk/src/types/wire.ts +++ b/ts-sdk/src/types/wire.ts @@ -249,6 +249,8 @@ export interface WireBackstopTransfer { cash: WireDecimal; mark_price: WireDecimal; equity: WireDecimal; + /** Crystallized by the forced close; `"0"` on the terminal cash sweep. */ + realized_pnl?: WireDecimal; timestamp_us: number; } @@ -267,8 +269,9 @@ export interface WireBackstopPage { // means something specific — noted per field. export interface WireOrdersFrame { - /** The orderbook these actions happened on; constant for the frame. */ - book: Hex; + /** The orderbook these actions happened on; constant for the frame. Absent on + * `pod_activity`, which is one frame per tick and names the book per entity. */ + book?: Hex; /** Deadline (micros) of the batch the actions **landed in**. Half of the resume cursor; `book` is the other half. */ batch: number; /** @@ -286,6 +289,9 @@ export interface WireOrderEntity { id: Hex; /** Creating transaction, or its parent `submitBatch` envelope. */ tx: Hex; + /** The order's book. Sent here by `pod_activity`, which covers a whole + * account; `pod_orders_v2` names it once on the frame instead. */ + book?: Hex; /** Index into the frame's `accts`; present iff `accts` is. */ a?: number; n: number; @@ -432,3 +438,80 @@ export interface WireWithdrawalDetail { status: "claimable" | "pending" | "refused"; proof?: { claim_hash?: Hex | null }; } + +// --- account activity (ADR 0057) --- +// +// `GET /clob/activity/{account}` serves the seed and `pod_activity` streams it. +// The seed is a union tagged by `activity_type`: the order variant flattens the +// orders route's `WireOrder`, and every money variant carries exactly the fields +// its stream event does, so the two decode through one decoder. The channel is +// `pod_orders_v2`'s frame for a whole account — one frame per tick, with the +// money the tick moved as extra event kinds. + +export interface WireActivityPage { + activity: WireActivityEntry[]; + next_cursor: string | null; + solution_now: number; // micros +} + +/** What each money kind carries. Shared because the seed row and the stream + * event differ only in how they are tagged. */ +export interface WireBackstopMoney { + /** Absent on the terminal cash sweep, which belongs to no market. */ + book?: Hex; + size: WireDecimal; // signed + cash: WireDecimal; // signed + mark: WireDecimal; + equity: WireDecimal; // signed + pnl: WireDecimal; // signed +} + +export interface WireBridgeMoney { + tx: Hex; + /** Position within the transaction: one tx can bridge several amounts, so the + * hash alone does not identify the row. */ + idx: number; + token: Hex; + amount: WireDecimal; // signed + error?: string; +} + +export interface WireTransferMoney { + id: Hex; + token: Hex; + amount: WireDecimal; // signed + error?: string; +} + +/** + * One activity row. `ts` is what the node sorted the page by — for an order that + * is its **signed deadline**, not the batch it landed in. + * + * Money amounts are signed from the account's side: negative left, positive + * arrived, for the bridge and for a transfer alike. + */ +export type WireActivityEntry = + | ({ activity_type: "order"; ts: number } & WireOrder) + | ({ activity_type: "backstop"; ts: number } & WireBackstopMoney) + | ({ activity_type: "bridge_transfer"; ts: number } & WireBridgeMoney) + | ({ activity_type: "transfer"; ts: number } & WireTransferMoney); + +/** + * One tick of an account's activity. No `book` — a tick covers every book the + * account traded on, so each entity names its own — and no `accts`, since the + * channel takes exactly one account. + */ +export interface WireActivityFrame { + /** Deadline (micros) of the batch. One frame per tick, so this alone is the + * resume cursor: a frame is already delivered exactly when `batch <= since`. */ + batch: number; + orders: WireOrderEntity[]; + events: (WireOrderEvent | WireMoneyEvent)[]; +} + +/** Money the tick moved, discriminated by `k` — disjoint from the order kinds, + * which is what lets the two share one untagged union on the wire. */ +export type WireMoneyEvent = + | ({ k: "backstop" } & WireBackstopMoney) + | ({ k: "bridge_transfer" } & WireBridgeMoney) + | ({ k: "transfer" } & WireTransferMoney);