-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathtrades_stream.rs
More file actions
126 lines (115 loc) · 4.81 KB
/
Copy pathtrades_stream.rs
File metadata and controls
126 lines (115 loc) · 4.81 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
//! Subscribe to all-trades through `MoonClient` and print stream signals.
//!
//! This is a bounded CLI/demo loop. It uses `drain_events() + sleep(...)` because
//! the program only prints a stream for a few seconds. A real terminal should
//! wire `MoonEventSink` into its UI event loop and read snapshots from the UI's
//! normal update/render cadence instead of creating a polling feed thread.
//!
//! Run:
//! cargo run --example trades_stream --release -- "<key_base64>" [host:port] [market|all] [watch_seconds]
use std::env;
use std::time::{Duration, Instant};
use moonproto::state::{
MarketHandle, MarketHistoryReaders, SeqRingCursor, TradeHistoryRow, TradesEvent,
};
use moonproto::{Event, TradesStreamMode};
mod common;
fn main() {
let args: Vec<String> = env::args().collect();
if args.len() < 2 {
eprintln!("Usage: trades_stream <key_base64> [host:port] [market|all] [watch_seconds]");
std::process::exit(1);
}
let market_filter = match args.get(3).map(String::as_str) {
Some("all") | None => None,
Some(name) => Some(name.to_string()),
};
let watch_secs: u64 = args.get(4).and_then(|s| s.parse().ok()).unwrap_or(30);
let client = match common::connect(&args[1], args.get(2), common::init_config()) {
Ok(client) => client,
Err(err) => {
eprintln!("[connect/init] failed: {err}");
std::process::exit(2);
}
};
if let Some(market) = market_filter.as_ref() {
client
.streams()
.subscribe_trades_for(TradesStreamMode::TradesOnly, [market.as_str()])
.expect("runtime stopped");
println!("[subscribe] all-trades, retained market={market}");
} else {
client
.streams()
.subscribe_all_trades(TradesStreamMode::TradesOnly)
.expect("runtime stopped");
println!("[subscribe] all-trades, retained all markets");
}
let selected_market: Option<MarketHandle> = market_filter.as_deref().and_then(|name| {
client
.snapshot()
.and_then(|snapshot| snapshot.markets().get(name))
});
if market_filter.is_some() && selected_market.is_none() {
eprintln!(
"[warn] selected market was not found in the current market snapshot; \
waiting for stream signals only"
);
}
let mut signals = 0u64;
let mut trades = 0u64;
let mut printed = 0u64;
let mut readers: Option<MarketHistoryReaders> = None;
let mut cursor: Option<SeqRingCursor> = None;
let mut rows: Vec<TradeHistoryRow> = Vec::new();
let deadline = Instant::now() + Duration::from_secs(watch_secs);
while Instant::now() < deadline {
for event in client.drain_events() {
if let Event::Trade(TradesEvent::Applied { .. }) = event {
signals += 1;
if let Some(market) = selected_market.as_ref() {
let Some(snapshot) = client.snapshot() else {
continue;
};
if readers.is_none() {
readers = snapshot.market_history_readers_for(market);
}
let Some(reader) = readers.as_ref().and_then(|r| r.futures_trades.clone())
else {
continue;
};
let cursor = cursor.get_or_insert_with(|| reader.cursor_from_oldest());
rows.clear();
let meta = reader.copy_new_since(cursor, 4096, &mut rows);
if meta.clipped {
println!("[retained-gap] local cursor fell behind retained history");
}
trades += rows.len() as u64;
let remaining_to_print = 25u64.saturating_sub(printed) as usize;
for row in rows.iter().take(remaining_to_print) {
let side = if row.is_buy() { "buy" } else { "sell" };
println!(
"[retained-trade] {} {side} price={} qty={} time_ms={}",
market.name(),
row.price,
row.quantity(),
row.unix_millis()
);
printed += 1;
if printed >= 25 {
break;
}
}
} else {
trades += 1;
if printed < 25 {
printed += 1;
println!("[trade-signal] retained rows updated");
}
}
}
}
std::thread::sleep(Duration::from_millis(50));
}
println!("[done] update_signals={signals} visible_updates={trades}");
}