-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmonitor_client.py
More file actions
344 lines (296 loc) · 11.9 KB
/
Copy pathmonitor_client.py
File metadata and controls
344 lines (296 loc) · 11.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
# File: monitor_client.py
from __future__ import annotations
import asyncio
import json
import os
import sys
import time
from collections import deque
from dataclasses import dataclass, field
from typing import Any
import websockets
from websockets.exceptions import ConnectionClosed
FEED_URI = "ws://127.0.0.1:9010"
STARTING_CAPITAL = 10_000.0
RECONNECT_DELAY_SECONDS = 1.0
REFRESH_SECONDS = 0.5
MAX_TRADES = 10
BOOK_DEPTH = 5
def round4(value: float) -> float:
rounded = round(float(value), 4)
if rounded == 0:
return 0.0
return rounded
def supports_color() -> bool:
return sys.stdout.isatty() and os.getenv("NO_COLOR") is None
class Palette:
def __init__(self, enabled: bool) -> None:
if enabled:
self.reset = "\033[0m"
self.green = "\033[32m"
self.red = "\033[31m"
self.yellow = "\033[33m"
self.cyan = "\033[36m"
self.bold = "\033[1m"
else:
self.reset = ""
self.green = ""
self.red = ""
self.yellow = ""
self.cyan = ""
self.bold = ""
def colorize(self, text: str, color: str) -> str:
return f"{color}{text}{self.reset}" if color else text
@dataclass(slots=True)
class TraderState:
position: float = 0.0
cash: float = STARTING_CAPITAL
realized_pnl: float = 0.0
unrealized_pnl: float = 0.0
total_equity: float = STARTING_CAPITAL
@dataclass(slots=True)
class MonitorState:
order_book: dict[str, list[tuple[float, float]]] = field(
default_factory=lambda: {"bids": [], "asks": []}
)
trades: deque[dict[str, Any]] = field(default_factory=lambda: deque(maxlen=MAX_TRADES))
traders: dict[str, TraderState] = field(default_factory=dict)
connected: bool = False
last_error: str = ""
last_event_ts: int | None = None
@property
def best_bid(self) -> float | None:
bids = self.order_book["bids"]
return bids[0][0] if bids else None
@property
def best_ask(self) -> float | None:
asks = self.order_book["asks"]
return asks[0][0] if asks else None
def mid_price(self) -> float | None:
if self.best_bid is None or self.best_ask is None:
return None
return round4((self.best_bid + self.best_ask) / 2.0)
def recalc_trader_metrics(self) -> None:
mid = self.mid_price()
if mid is None:
for row in self.traders.values():
row.unrealized_pnl = 0.0
row.total_equity = round4(row.cash)
return
for row in self.traders.values():
unrealized = round4(row.position * mid)
total_equity = round4(row.cash + unrealized)
row.unrealized_pnl = unrealized
row.total_equity = total_equity
class MonitorDashboard:
def __init__(self, uri: str = FEED_URI) -> None:
self._uri = uri
self._state = MonitorState()
self._palette = Palette(supports_color())
self._shutdown = asyncio.Event()
async def run(self) -> None:
receiver_task = asyncio.create_task(self._receiver_loop(), name="receiver")
render_task = asyncio.create_task(self._render_loop(), name="renderer")
try:
await self._shutdown.wait()
finally:
receiver_task.cancel()
render_task.cancel()
await asyncio.gather(receiver_task, render_task, return_exceptions=True)
async def _receiver_loop(self) -> None:
while not self._shutdown.is_set():
try:
self._state.connected = False
async with websockets.connect(self._uri) as websocket:
self._state.connected = True
self._state.last_error = ""
async for raw in websocket:
payload = self._safe_json(raw)
if payload is None:
continue
self._apply_event(payload)
except ConnectionClosed as exc:
self._state.connected = False
self._state.last_error = f"connection closed: {exc.code} {exc.reason}"
except Exception as exc:
self._state.connected = False
self._state.last_error = f"connection error: {exc}"
await asyncio.sleep(RECONNECT_DELAY_SECONDS)
async def _render_loop(self) -> None:
while not self._shutdown.is_set():
self._state.recalc_trader_metrics()
self._render()
await asyncio.sleep(REFRESH_SECONDS)
def _safe_json(self, raw: str) -> dict[str, Any] | None:
try:
payload = json.loads(raw)
except json.JSONDecodeError:
self._state.last_error = "received invalid JSON payload"
return None
if not isinstance(payload, dict):
self._state.last_error = "received non-object JSON payload"
return None
return payload
def _apply_event(self, payload: dict[str, Any]) -> None:
event_type = payload.get("type")
if not isinstance(event_type, str):
return
if event_type == "book_update":
self._handle_book_update(payload)
return
if event_type == "trade":
self._handle_trade(payload)
return
if event_type == "position_update":
self._handle_position_update(payload)
return
def _handle_book_update(self, payload: dict[str, Any]) -> None:
bids = self._parse_levels(payload.get("bids"))
asks = self._parse_levels(payload.get("asks"))
bids.sort(key=lambda x: x[0], reverse=True)
asks.sort(key=lambda x: x[0])
self._state.order_book["bids"] = bids
self._state.order_book["asks"] = asks
ts = payload.get("timestamp")
if isinstance(ts, int):
self._state.last_event_ts = ts
def _handle_trade(self, payload: dict[str, Any]) -> None:
price = payload.get("price")
qty = payload.get("qty")
ts = payload.get("timestamp")
if not isinstance(price, (int, float)) or not isinstance(qty, (int, float)):
return
trade = {
"price": round4(price),
"qty": round4(qty),
"timestamp": ts if isinstance(ts, int) else None,
"buy_trader_id": payload.get("buy_trader_id"),
"sell_trader_id": payload.get("sell_trader_id"),
}
self._state.trades.appendleft(trade)
if isinstance(ts, int):
self._state.last_event_ts = ts
def _handle_position_update(self, payload: dict[str, Any]) -> None:
trader_id = payload.get("trader_id")
if not isinstance(trader_id, str) or not trader_id.strip():
return
trader = self._state.traders.get(trader_id)
if trader is None:
trader = TraderState()
self._state.traders[trader_id] = trader
trader.position = float(payload.get("position", trader.position))
trader.cash = round4(float(payload.get("cash", trader.cash)))
trader.realized_pnl = round4(float(payload.get("realized_pnl", trader.realized_pnl)))
ts = payload.get("timestamp")
if isinstance(ts, int):
self._state.last_event_ts = ts
def _parse_levels(self, raw_levels: Any) -> list[tuple[float, float]]:
levels: list[tuple[float, float]] = []
if not isinstance(raw_levels, list):
return levels
for level in raw_levels:
if not isinstance(level, (list, tuple)) or len(level) != 2:
continue
px, qty = level
if not isinstance(px, (int, float)) or not isinstance(qty, (int, float)):
continue
qty_val = round4(qty)
if qty_val <= 0:
continue
levels.append((round4(px), qty_val))
return levels
def _render(self) -> None:
print("\033[H\033[J", end="")
p = self._palette
status = p.colorize("CONNECTED", p.green) if self._state.connected else p.colorize("DISCONNECTED", p.red)
mid = self._state.mid_price()
print("===============================")
print("OpenMarketSim Monitor")
print("===============================")
print(f"Feed: {self._uri}")
print(f"Status: {status}")
if self._state.last_error:
print(f"Last Error: {self._state.last_error}")
print(f"Best Bid: {self._fmt_px(self._state.best_bid)} Best Ask: {self._fmt_px(self._state.best_ask)} Mid: {self._fmt_px(mid)}")
print("")
print("ORDER BOOK (top 5)")
print("-------------------------------")
print(" BID_QTY BID_PX | ASK_PX ASK_QTY")
bids = self._state.order_book["bids"][:BOOK_DEPTH]
asks = self._state.order_book["asks"][:BOOK_DEPTH]
for i in range(BOOK_DEPTH):
bid_px, bid_qty = ("", "")
ask_px, ask_qty = ("", "")
if i < len(bids):
bid_px = f"{bids[i][0]:.2f}"
bid_qty = f"{bids[i][1]:.4f}".rstrip("0").rstrip(".")
if i < len(asks):
ask_px = f"{asks[i][0]:.2f}"
ask_qty = f"{asks[i][1]:.4f}".rstrip("0").rstrip(".")
print(f"{bid_qty:>11} {bid_px:>11} | {ask_px:>9} {ask_qty:>11}")
print("")
print("RECENT TRADES")
print("-------------------------------")
print(" TS(ms) PRICE QTY")
for trade in list(self._state.trades)[:MAX_TRADES]:
ts = trade.get("timestamp")
ts_txt = str(ts) if isinstance(ts, int) else "-"
px_txt = f"{float(trade['price']):.2f}"
qty_txt = f"{float(trade['qty']):.4f}".rstrip("0").rstrip(".")
print(f"{ts_txt:>12} {px_txt:>12} {qty_txt:>8}")
if not self._state.trades:
print(" (no trades yet)")
print("")
rows = self._leaderboard_rows()
print("BOT PERFORMANCE")
print("Trader | Pos | Cash | Unreal | Total | PnL")
print("-------------------------------------------")
for row in rows:
pnl_color = p.green if row["pnl"] >= 0 else p.red
pnl_text = p.colorize(f"{row['pnl']:.2f}", pnl_color)
print(
f"{row['trader_id']:<10} {row['position']:>6.2f} {row['cash']:>10.2f} "
f"{row['unrealized']:>10.2f} {row['total_equity']:>10.2f} {pnl_text:>10}"
)
if not rows:
print("(no trader state yet; waiting for position_update events)")
print("")
print("LEADERBOARD")
print("-------------------------------------------")
for i, row in enumerate(rows, 1):
print(f"{i:>2}. {row['trader_id']:<12} PnL: {row['pnl']:>10.2f}")
if not rows:
print("(leaderboard unavailable)")
sys.stdout.flush()
def _leaderboard_rows(self) -> list[dict[str, Any]]:
rows: list[dict[str, Any]] = []
for trader_id, state in self._state.traders.items():
unrealized = state.unrealized_pnl
total_equity = state.total_equity
pnl = round4(total_equity - STARTING_CAPITAL)
rows.append(
{
"trader_id": trader_id,
"position": state.position,
"cash": state.cash,
"realized_pnl": state.realized_pnl,
"unrealized": unrealized,
"total_equity": total_equity,
"pnl": pnl,
}
)
rows.sort(key=lambda x: (-x["pnl"], x["trader_id"]))
return rows
@staticmethod
def _fmt_px(value: float | None) -> str:
if value is None:
return "-"
return f"{value:.2f}"
async def main() -> None:
dashboard = MonitorDashboard(uri=FEED_URI)
await dashboard.run()
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
pass