Files
RustAst/src/ast/rtl/data_streams.rs
T
Brummel fad8bc3471 Extract data_server into external data-server crate
The data_server module had no dependency on the rest of myc (only
chrono/regex/zip), so it was moved verbatim into a standalone leaf
crate at its own Gitea repo and is now consumed as a git dependency.
A 'pub use ::data_server;' re-export in src/ast/mod.rs keeps all
existing myc::ast::data_server::… paths valid; no call sites changed.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-18 11:04:03 +02:00

220 lines
8.7 KiB
Rust

//! RTL registration for file-based market data streams.
//!
//! Bridges the `Arc`-based [`::data_server::DataServer`] cache with the `Rc`-based VM world
//! by registering generator closures that read pre-parsed chunks and convert
//! them into `Value::Record` items for `RootStream::tick()`.
use crate::ast::data_server::get_data_server;
use crate::ast::environment::Environment;
use crate::ast::rtl::streams::{RootStream, StreamNode};
use crate::ast::types::{Keyword, Purity, RecordLayout, Signature, StaticType, Value};
use std::rc::Rc;
/// Registers `create-m1-stream` and `create-tick-stream` as RTL functions.
pub fn register(env: &Environment) {
register_m1_stream(env);
register_tick_stream(env);
}
// -- M1 Stream ---------------------------------------------------------------
fn register_m1_stream(env: &Environment) {
let generators = env.pipeline_generators.clone();
let m1_layout = RecordLayout::get_or_create(vec![
(Keyword::intern("time"), StaticType::DateTime),
(Keyword::intern("open"), StaticType::Float),
(Keyword::intern("high"), StaticType::Float),
(Keyword::intern("low"), StaticType::Float),
(Keyword::intern("close"), StaticType::Float),
(Keyword::intern("spread"), StaticType::Float),
(Keyword::intern("volume"), StaticType::Int),
]);
let ret_type = StaticType::Stream(Box::new(StaticType::Record(m1_layout.clone())));
env.register_native_fn(
"create-m1-stream",
StaticType::Function(Box::new(Signature {
params: StaticType::Any,
ret: ret_type,
})),
Purity::Impure,
move |args: &[Value]| {
let symbol = match &args[0] {
Value::Text(s) => s.as_ref(),
_ => panic!("create-m1-stream: first argument must be a string"),
};
let from_ms = args.get(1).map(|v| match v {
Value::DateTime(ms) => *ms,
_ => panic!("create-m1-stream: 'from' must be a DateTime"),
});
let to_ms = args.get(2).map(|v| match v {
Value::DateTime(ms) => *ms,
_ => panic!("create-m1-stream: 'to' must be a DateTime"),
});
let server = get_data_server()
.expect("create-m1-stream: DataServer not initialized. Call init_data_server() first.");
// Unknown symbol → Void (no stream created)
if !server.has_symbol(symbol) {
return Value::Void;
}
let mut iter = server.stream_m1_windowed(symbol, from_ms, to_ms);
// Create the pipeline entry point
let root_stream = Rc::new(RootStream::new());
let layout = m1_layout.clone();
let stream_node = StreamNode {
inner: root_stream.clone(),
element_type: StaticType::Record(layout.clone()),
};
// Pre-fetch the first chunk (None if time window is empty)
let mut current_chunk = iter.as_mut().and_then(|it| it.next_chunk());
let mut pos = 0;
let generator = move || -> bool {
loop {
let chunk = match &current_chunk {
Some(c) => c,
None => return false,
};
if pos < chunk.len() {
let rec = &chunk[pos];
pos += 1;
let value = Value::Record(
layout.clone(),
Rc::new(vec![
Value::DateTime(rec.time_ms),
Value::Float(rec.open),
Value::Float(rec.high),
Value::Float(rec.low),
Value::Float(rec.close),
Value::Float(rec.spread),
Value::Int(rec.volume),
]),
);
root_stream.tick(value);
return true;
}
// Advance to next chunk (triggers prefetch internally)
current_chunk = iter.as_mut().and_then(|it| it.next_chunk());
pos = 0;
}
};
generators.borrow_mut().push(Box::new(generator));
Value::Stream(Rc::new(stream_node))
},
)
.doc("Creates a stream of M1 (minute) OHLCV bars from tick data files.")
.description("Reads pre-cached binary data from the DataServer. Optional from/to DateTime arguments restrict the time window (millisecond precision). Returns Void if the symbol is unknown. Returns an empty (immediately exhausted) stream if the symbol exists but the time window contains no data.")
.examples(&[
"(def ohlc (create-m1-stream \"EURUSD\"))",
"(def ohlc (create-m1-stream \"EURUSD\" (date \"2020-01-01\") (date \"2020-12-31\")))",
]);
}
// -- Tick Stream -------------------------------------------------------------
fn register_tick_stream(env: &Environment) {
let generators = env.pipeline_generators.clone();
let tick_layout = RecordLayout::get_or_create(vec![
(Keyword::intern("time"), StaticType::DateTime),
(Keyword::intern("ask"), StaticType::Float),
(Keyword::intern("bid"), StaticType::Float),
]);
let ret_type = StaticType::Stream(Box::new(StaticType::Record(tick_layout.clone())));
env.register_native_fn(
"create-tick-stream",
StaticType::Function(Box::new(Signature {
params: StaticType::Any,
ret: ret_type,
})),
Purity::Impure,
move |args: &[Value]| {
let symbol = match &args[0] {
Value::Text(s) => s.as_ref(),
_ => panic!("create-tick-stream: first argument must be a string"),
};
let from_ms = args.get(1).map(|v| match v {
Value::DateTime(ms) => *ms,
_ => panic!("create-tick-stream: 'from' must be a DateTime"),
});
let to_ms = args.get(2).map(|v| match v {
Value::DateTime(ms) => *ms,
_ => panic!("create-tick-stream: 'to' must be a DateTime"),
});
let server = get_data_server()
.expect("create-tick-stream: DataServer not initialized.");
// Unknown symbol → Void (no stream created)
if !server.has_symbol(symbol) {
return Value::Void;
}
let mut iter = server.stream_tick_windowed(symbol, from_ms, to_ms);
let root_stream = Rc::new(RootStream::new());
let layout = tick_layout.clone();
let stream_node = StreamNode {
inner: root_stream.clone(),
element_type: StaticType::Record(layout.clone()),
};
// Pre-fetch the first chunk (None if time window is empty)
let mut current_chunk = iter.as_mut().and_then(|it| it.next_chunk());
let mut pos = 0;
let generator = move || -> bool {
loop {
let chunk = match &current_chunk {
Some(c) => c,
None => return false,
};
if pos < chunk.len() {
let rec = &chunk[pos];
pos += 1;
let value = Value::Record(
layout.clone(),
Rc::new(vec![
Value::DateTime(rec.time_ms),
Value::Float(rec.ask),
Value::Float(rec.bid),
]),
);
root_stream.tick(value);
return true;
}
current_chunk = iter.as_mut().and_then(|it| it.next_chunk());
pos = 0;
}
};
generators.borrow_mut().push(Box::new(generator));
Value::Stream(Rc::new(stream_node))
},
)
.doc("Creates a stream of tick data (ask/bid) from tick data files.")
.description("Reads pre-cached binary data from the DataServer. Optional from/to DateTime arguments restrict the time window (millisecond precision). Returns Void if the symbol is unknown. Returns an empty (immediately exhausted) stream if the symbol exists but the time window contains no data.")
.examples(&[
"(def ticks (create-tick-stream \"EURUSD\"))",
"(def ticks (create-tick-stream \"EURUSD\" (date \"2020-01-01\") (date \"2020-12-31\")))",
]);
}