examples: move Rust examples into a top-level examples/rust/ crate
The three Rust examples (backtest, fetch_btcusdt, live_binance) used to
live each in their own crate's examples/ dir, splitting the example set
across crates and burying it inside the source tree. Move them into a new
workspace member crate at `examples/rust/` (package `wickra-examples`,
`publish = false`) so all language examples sit under one top-level
`examples/<lang>/` tree.
* `examples/rust/Cargo.toml` declares the per-binary deps (wickra,
wickra-data with the `live-binance` feature always on, serde_json, tokio
for the macro and current-thread runtime).
* `examples/rust/src/bin/{backtest,fetch_btcusdt,live_binance}.rs` are the
three migrated binaries; their doc-comments and the fetch_btcusdt output
path are updated for the new location and run command
(`cargo run -p wickra-examples --bin <name>`).
* Workspace `Cargo.toml` lists the new member; the now-empty
`[dev-dependencies]` extras (`wickra`, `tokio` in wickra-data and
`serde_json` in wickra) that existed only for these examples are dropped.
* The `[[example]] live_binance` table is removed from wickra-data's
manifest since the file moved out.
* README "Languages" + project-layout, examples/README.md, Quickstart-Rust
and Data-Layer are pointed at the new paths and commands.
`cargo build -p wickra-examples` and `cargo run --release -p wickra-examples
--bin backtest -- examples/data/btcusdt-1d.csv` both succeed; the rest of
the workspace (core, data, wickra) builds, clippies (`--all-targets -D
warnings`) and tests (508 core + 28 data + 1 integration + 74+3+1
doctests) all stay green.
This commit is contained in:
@@ -0,0 +1,95 @@
|
||||
//! Rust example: backtest a basket of indicators against a CSV file.
|
||||
//!
|
||||
//! Build with:
|
||||
//! ```text
|
||||
//! cargo run --release -p wickra-examples --bin backtest -- path/to/ohlcv.csv
|
||||
//! ```
|
||||
|
||||
use std::env;
|
||||
|
||||
use wickra::{Adx, Atr, BatchExt, BollingerBands, Ema, Indicator, MacdIndicator, Obv, Rsi};
|
||||
use wickra_data::csv::CandleReader;
|
||||
|
||||
fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let args: Vec<String> = env::args().collect();
|
||||
let path = args.get(1).ok_or("usage: backtest <ohlcv.csv>")?;
|
||||
|
||||
let mut reader = CandleReader::open(path)?;
|
||||
let candles = reader.read_all()?;
|
||||
if candles.is_empty() {
|
||||
return Err("CSV is empty".into());
|
||||
}
|
||||
let n = candles.len();
|
||||
let closes: Vec<f64> = candles.iter().map(|c| c.close).collect();
|
||||
|
||||
let rsi = Rsi::new(14)?.batch(&closes);
|
||||
let ema = Ema::new(20)?.batch(&closes);
|
||||
let bb = BollingerBands::classic().batch(&closes);
|
||||
let macd = MacdIndicator::classic().batch(&closes);
|
||||
|
||||
let mut atr = Atr::new(14)?;
|
||||
let atr_series: Vec<_> = candles.iter().map(|c| atr.update(*c)).collect();
|
||||
|
||||
let mut adx = Adx::new(14)?;
|
||||
let adx_series: Vec<_> = candles.iter().map(|c| adx.update(*c)).collect();
|
||||
|
||||
let mut obv = Obv::new();
|
||||
let obv_series: Vec<_> = candles.iter().map(|c| obv.update(*c)).collect();
|
||||
|
||||
let last_rsi = rsi
|
||||
.iter()
|
||||
.rev()
|
||||
.flatten()
|
||||
.next()
|
||||
.copied()
|
||||
.unwrap_or(f64::NAN);
|
||||
let last_ema = ema
|
||||
.iter()
|
||||
.rev()
|
||||
.flatten()
|
||||
.next()
|
||||
.copied()
|
||||
.unwrap_or(f64::NAN);
|
||||
let last_bb = bb.iter().rev().flatten().next().copied();
|
||||
let last_macd = macd.iter().rev().flatten().next().copied();
|
||||
let last_atr = atr_series
|
||||
.iter()
|
||||
.rev()
|
||||
.flatten()
|
||||
.next()
|
||||
.copied()
|
||||
.unwrap_or(f64::NAN);
|
||||
let last_adx = adx_series.iter().rev().flatten().next().copied();
|
||||
let last_obv = obv_series
|
||||
.iter()
|
||||
.rev()
|
||||
.flatten()
|
||||
.next()
|
||||
.copied()
|
||||
.unwrap_or(f64::NAN);
|
||||
|
||||
println!("backtest summary for {path} ({n} bars)");
|
||||
println!(" RSI(14) = {last_rsi:>9.4}");
|
||||
println!(" EMA(20) = {last_ema:>9.4}");
|
||||
if let Some(bb) = last_bb {
|
||||
println!(
|
||||
" BB(20,2) upper={:>9.4} middle={:>9.4} lower={:>9.4} sd={:>8.4}",
|
||||
bb.upper, bb.middle, bb.lower, bb.stddev
|
||||
);
|
||||
}
|
||||
if let Some(m) = last_macd {
|
||||
println!(
|
||||
" MACD macd={:>9.4} signal={:>9.4} hist={:>9.4}",
|
||||
m.macd, m.signal, m.histogram
|
||||
);
|
||||
}
|
||||
println!(" ATR(14) = {last_atr:>9.4}");
|
||||
if let Some(a) = last_adx {
|
||||
println!(
|
||||
" ADX(14) +DI={:>6.2} -DI={:>6.2} ADX={:>6.2}",
|
||||
a.plus_di, a.minus_di, a.adx
|
||||
);
|
||||
}
|
||||
println!(" OBV = {last_obv:>14.2}");
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,254 @@
|
||||
//! Rust example: download real BTCUSDT spot candles from the Binance REST API
|
||||
//! and write them as CSV datasets under `examples/data/`.
|
||||
//!
|
||||
//! These are the datasets the indicator benchmarks (`crates/wickra/benches/
|
||||
//! indicators.rs`) and the `example_data` integration test run against. They
|
||||
//! live at the workspace `examples/data/` directory. Re-run this example to
|
||||
//! refresh them with the latest market history:
|
||||
//!
|
||||
//! ```text
|
||||
//! cargo run -p wickra-examples --bin fetch_btcusdt
|
||||
//! ```
|
||||
//!
|
||||
//! HTTPS is handled by shelling out to the system `curl` — shipped with
|
||||
//! Windows 10+, macOS, and every Linux distribution — so this example adds no
|
||||
//! HTTP or TLS dependency to the crate (`serde_json` is the only extra, and
|
||||
//! only as a dev-dependency). Each Binance klines response is capped at 1000
|
||||
//! rows, so larger datasets are paginated backwards through `endTime`. Only
|
||||
//! fully closed candles are kept; the still-forming candle of the current
|
||||
//! bucket is dropped so the datasets stay reproducible.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::error::Error;
|
||||
use std::fmt::Write as _;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::Command;
|
||||
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
||||
|
||||
use serde_json::Value as Json;
|
||||
use wickra::Candle;
|
||||
|
||||
/// Binance Spot REST endpoint for historical klines.
|
||||
const KLINES_URL: &str = "https://api.binance.com/api/v3/klines";
|
||||
/// Trading pair to download.
|
||||
const SYMBOL: &str = "BTCUSDT";
|
||||
/// Binance caps a single klines response at 1000 rows.
|
||||
const PAGE_LIMIT: usize = 1000;
|
||||
/// Courtesy pause between paginated requests to stay well under the rate limit.
|
||||
const REQUEST_PAUSE: Duration = Duration::from_millis(200);
|
||||
|
||||
/// One dataset to produce: the Binance interval code, the output file name,
|
||||
/// and how many of the most-recent *closed* candles to collect.
|
||||
struct Dataset {
|
||||
interval: &'static str,
|
||||
file: &'static str,
|
||||
target: usize,
|
||||
}
|
||||
|
||||
/// The seven datasets, one per timeframe. `12h`/`1d`/`1M` simply collect all
|
||||
/// the history Binance offers — it is shorter than their `target`.
|
||||
///
|
||||
/// The monthly file is named `btcusdt-1month.csv`, not `btcusdt-1M.csv`: on
|
||||
/// case-insensitive filesystems (Windows, default macOS) the latter would
|
||||
/// collide with `btcusdt-1m.csv` and one dataset would silently overwrite the
|
||||
/// other.
|
||||
const DATASETS: &[Dataset] = &[
|
||||
Dataset {
|
||||
interval: "1m",
|
||||
file: "btcusdt-1m.csv",
|
||||
target: 50_000,
|
||||
},
|
||||
Dataset {
|
||||
interval: "5m",
|
||||
file: "btcusdt-5m.csv",
|
||||
target: 10_000,
|
||||
},
|
||||
Dataset {
|
||||
interval: "15m",
|
||||
file: "btcusdt-15m.csv",
|
||||
target: 10_000,
|
||||
},
|
||||
Dataset {
|
||||
interval: "1h",
|
||||
file: "btcusdt-1h.csv",
|
||||
target: 10_000,
|
||||
},
|
||||
Dataset {
|
||||
interval: "12h",
|
||||
file: "btcusdt-12h.csv",
|
||||
target: 5_000,
|
||||
},
|
||||
Dataset {
|
||||
interval: "1d",
|
||||
file: "btcusdt-1d.csv",
|
||||
target: 5_000,
|
||||
},
|
||||
Dataset {
|
||||
interval: "1M",
|
||||
file: "btcusdt-1month.csv",
|
||||
target: 5_000,
|
||||
},
|
||||
];
|
||||
|
||||
fn main() -> Result<(), Box<dyn Error>> {
|
||||
// A single reference instant for the whole run: any candle whose bucket has
|
||||
// not closed by this point is treated as in-progress and dropped.
|
||||
let now_ms = i64::try_from(SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis())
|
||||
.map_err(|_| "system clock is beyond the i64 millisecond range")?;
|
||||
|
||||
let data_dir = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
|
||||
.join("..")
|
||||
.join("data");
|
||||
std::fs::create_dir_all(&data_dir)?;
|
||||
|
||||
println!(
|
||||
"Fetching {SYMBOL} klines from Binance into {}",
|
||||
data_dir.display()
|
||||
);
|
||||
for ds in DATASETS {
|
||||
let candles = collect(ds.interval, ds.target, now_ms)?;
|
||||
if candles.is_empty() {
|
||||
return Err(format!(
|
||||
"Binance returned no closed candles for interval {}",
|
||||
ds.interval
|
||||
)
|
||||
.into());
|
||||
}
|
||||
let path = data_dir.join(ds.file);
|
||||
write_csv(&path, &candles)?;
|
||||
println!(
|
||||
" {:>3} {:>6} candles -> examples/data/{}",
|
||||
ds.interval,
|
||||
candles.len(),
|
||||
ds.file
|
||||
);
|
||||
}
|
||||
println!("Done — {} datasets written.", DATASETS.len());
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Paginate the Binance REST API backwards until `target` closed candles have
|
||||
/// been collected (or the exchange runs out of history), then return the most
|
||||
/// recent `target` of them in ascending time order.
|
||||
fn collect(interval: &str, target: usize, now_ms: i64) -> Result<Vec<Candle>, Box<dyn Error>> {
|
||||
// Keyed by open time: the map sorts chronologically and dedupes the small
|
||||
// overlap that can occur between adjacent pages.
|
||||
let mut by_open: BTreeMap<i64, Candle> = BTreeMap::new();
|
||||
let mut end_time: Option<i64> = None;
|
||||
let mut pages = 0usize;
|
||||
|
||||
loop {
|
||||
let page = fetch_page(interval, end_time)?;
|
||||
pages += 1;
|
||||
if page.is_empty() {
|
||||
break;
|
||||
}
|
||||
|
||||
let mut oldest_open = i64::MAX;
|
||||
for raw in &page {
|
||||
let Some((open_time, close_time, candle)) = parse_kline(raw) else {
|
||||
continue;
|
||||
};
|
||||
oldest_open = oldest_open.min(open_time);
|
||||
// Keep only fully closed candles — drop the in-progress bucket.
|
||||
if close_time < now_ms {
|
||||
by_open.insert(open_time, candle);
|
||||
}
|
||||
}
|
||||
eprint!(
|
||||
"\r {interval}: collected {} candles over {pages} page(s)…",
|
||||
by_open.len()
|
||||
);
|
||||
|
||||
// Stop once enough is collected or Binance has no older history left.
|
||||
if by_open.len() >= target || page.len() < PAGE_LIMIT {
|
||||
break;
|
||||
}
|
||||
if oldest_open == i64::MAX {
|
||||
return Err(format!("Binance page for {interval} held no parseable klines").into());
|
||||
}
|
||||
// Step the window strictly before the oldest open time seen so far.
|
||||
end_time = Some(oldest_open - 1);
|
||||
std::thread::sleep(REQUEST_PAUSE);
|
||||
}
|
||||
eprintln!();
|
||||
|
||||
let mut candles: Vec<Candle> = by_open.into_values().collect();
|
||||
if candles.len() > target {
|
||||
// Trim the oldest surplus, keeping the most recent `target` candles.
|
||||
candles.drain(..candles.len() - target);
|
||||
}
|
||||
Ok(candles)
|
||||
}
|
||||
|
||||
/// Fetch one page of klines. `end_time`, when set, caps the open time so the
|
||||
/// caller can walk backwards through history.
|
||||
fn fetch_page(interval: &str, end_time: Option<i64>) -> Result<Vec<Json>, Box<dyn Error>> {
|
||||
let mut url = format!("{KLINES_URL}?symbol={SYMBOL}&interval={interval}&limit={PAGE_LIMIT}");
|
||||
if let Some(end) = end_time {
|
||||
write!(url, "&endTime={end}").expect("writing to a String never fails");
|
||||
}
|
||||
|
||||
let output = Command::new("curl")
|
||||
.args([
|
||||
"--silent",
|
||||
"--show-error",
|
||||
"--fail",
|
||||
"--max-time",
|
||||
"30",
|
||||
&url,
|
||||
])
|
||||
.output()
|
||||
.map_err(|e| format!("could not run `curl` (install it / put it on PATH): {e}"))?;
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
return Err(format!("curl failed for {url}: {}", stderr.trim()).into());
|
||||
}
|
||||
|
||||
let body = String::from_utf8(output.stdout)?;
|
||||
match serde_json::from_str(&body)? {
|
||||
Json::Array(rows) => Ok(rows),
|
||||
other => Err(format!("expected a JSON array of klines from Binance, got: {other}").into()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Parse one raw Binance kline array into `(open_time, close_time, Candle)`.
|
||||
///
|
||||
/// A Binance kline is `[openTime, "open", "high", "low", "close", "volume",
|
||||
/// closeTime, …]`; the OHLCV fields arrive as decimal strings. Returns `None`
|
||||
/// for any row that is malformed or fails `Candle::new`'s OHLC validation.
|
||||
fn parse_kline(raw: &Json) -> Option<(i64, i64, Candle)> {
|
||||
let arr = raw.as_array()?;
|
||||
if arr.len() < 7 {
|
||||
return None;
|
||||
}
|
||||
let open_time = arr[0].as_i64()?;
|
||||
let close_time = arr[6].as_i64()?;
|
||||
let field = |i: usize| -> Option<f64> { arr[i].as_str()?.parse::<f64>().ok() };
|
||||
let candle = Candle::new(
|
||||
field(1)?,
|
||||
field(2)?,
|
||||
field(3)?,
|
||||
field(4)?,
|
||||
field(5)?,
|
||||
open_time,
|
||||
)
|
||||
.ok()?;
|
||||
Some((open_time, close_time, candle))
|
||||
}
|
||||
|
||||
/// Write candles as a standard `timestamp,open,high,low,close,volume` CSV that
|
||||
/// [`wickra_data::csv::CandleReader`] can read back.
|
||||
fn write_csv(path: &Path, candles: &[Candle]) -> Result<(), Box<dyn Error>> {
|
||||
let mut out = String::with_capacity(candles.len() * 80 + 32);
|
||||
out.push_str("timestamp,open,high,low,close,volume\n");
|
||||
for c in candles {
|
||||
writeln!(
|
||||
out,
|
||||
"{},{},{},{},{},{}",
|
||||
c.timestamp, c.open, c.high, c.low, c.close, c.volume
|
||||
)?;
|
||||
}
|
||||
std::fs::write(path, out)?;
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
//! Connect to Binance, run incremental RSI on the closed kline events.
|
||||
//!
|
||||
//! Build with:
|
||||
//! ```text
|
||||
//! cargo run --release -p wickra-examples --bin live_binance -- BTCUSDT
|
||||
//! ```
|
||||
//!
|
||||
//! The example prints a line per bar close. Hit Ctrl+C to stop.
|
||||
|
||||
use std::env;
|
||||
|
||||
use wickra::{Indicator, Rsi};
|
||||
use wickra_data::live::binance::{BinanceKlineStream, Interval};
|
||||
|
||||
#[tokio::main(flavor = "current_thread")]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let args: Vec<String> = env::args().collect();
|
||||
let symbol = args
|
||||
.get(1)
|
||||
.cloned()
|
||||
.unwrap_or_else(|| "BTCUSDT".to_string());
|
||||
|
||||
let mut stream =
|
||||
BinanceKlineStream::connect(std::slice::from_ref(&symbol), Interval::OneMinute).await?;
|
||||
let mut rsi = Rsi::new(14)?;
|
||||
println!("Listening for {symbol} 1m closes…");
|
||||
|
||||
while let Some(event) = stream.next_event().await? {
|
||||
if !event.is_closed {
|
||||
continue;
|
||||
}
|
||||
let close = event.candle.close;
|
||||
let v = rsi.update(close);
|
||||
if let Some(value) = v {
|
||||
println!("{} close={:.4} rsi={:.2}", event.symbol, close, value);
|
||||
} else {
|
||||
println!("{} close={:.4} rsi=…warmup", event.symbol, close);
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
Reference in New Issue
Block a user