From f3459baf431665a9dccbace534d99acc78de75ac Mon Sep 17 00:00:00 2001 From: Michael Schimmel Date: Mon, 2 Mar 2026 13:15:56 +0100 Subject: [PATCH] Add create-ticker function Implements the `create-ticker` function, which allows for the creation of a ticker stream based on a provided condition closure. This function is crucial for synchronizing pipeline execution based on specific events or conditions. --- examples/pipeline_script_ticker.myc | 48 ++++++++++++++++++++++++ src/ast/rtl/streams.rs | 57 ++++++++++++++++++++++++++++- 2 files changed, 104 insertions(+), 1 deletion(-) create mode 100644 examples/pipeline_script_ticker.myc diff --git a/examples/pipeline_script_ticker.myc b/examples/pipeline_script_ticker.myc new file mode 100644 index 0000000..6eb2a42 --- /dev/null +++ b/examples/pipeline_script_ticker.myc @@ -0,0 +1,48 @@ +;; Benchmark: 2.0us +;; Benchmark-Repeat: 1015 +;; Output: PipelineNode[last: Some(110.45243843391206)] + +(do + ;; Set the random seed to match the existing test output exactly + (seed! 42) + + ;; Use our new generic create-ticker to pulse 3 times + (def cnt 3) + (def ticker (create-ticker (fn [] (> (assign cnt (- cnt 1)) -1)))) + + ;; The candle generator reacting to the ticker + (def last-close 100.0) + + (def src1 + (pipe [ticker] + (fn [_] + (do + (def change (* (- (random) 0.5) 2.0)) ; Matches Rust: (rng.f64() - 0.5) * 2.0 + (def open last-close) + (def high (+ open (abs (* (random) 2.0)))) + (def low (- open (abs (* (random) 2.0)))) + (assign last-close (+ open change)) + {:open open :high high :low low :close last-close} + ) + ) + ) + ) + + (def combined + (pipe [src1] + (fn [t1] + (.close t1) + ) + ) + ) + + (def my_indicator + (pipe [combined] + (fn [close_price] + (+ close_price 10.0) + ) + ) + ) + + my_indicator +) diff --git a/src/ast/rtl/streams.rs b/src/ast/rtl/streams.rs index f12df26..50ba5e8 100644 --- a/src/ast/rtl/streams.rs +++ b/src/ast/rtl/streams.rs @@ -341,6 +341,62 @@ pub fn register(env: &Environment) { Value::Object(Rc::new(stream_node)) }, ); + + // (create-ticker condition-closure) -> StreamNode + let ticker_generators = env.pipeline_generators.clone(); + let globals = env.global_values.clone(); + + env.register_native_fn( + "create-ticker", + StaticType::Function(Box::new(Signature { + params: StaticType::Tuple(vec![StaticType::Any]), // Expects a closure + ret: StaticType::Any, // Returns StreamNode + })), + Purity::Impure, + move |args: std::vec::Vec| { + if args.len() != 1 { + panic!("create-ticker expects exactly 1 argument (the condition closure)"); + } + + let closure_obj = if let Value::Object(obj) = &args[0] { + obj.clone() + } else { + panic!("create-ticker expects a closure as its argument"); + }; + + // 1. Create the RootStream + let root_stream = Rc::new(RootStream::new()); + let stream_node = StreamNode { inner: root_stream.clone() }; + + // 2. Setup isolated VM + let mut ticker_vm = crate::ast::vm::VM::new(globals.clone()); + let my_closure = closure_obj.clone(); + + // 3. Create generator + let generator = move || -> bool { + match ticker_vm.run_with_args(my_closure.clone(), vec![]) { + Ok(Value::Bool(b)) => { + if b { + // Ticker pulses with a simple `true` value or `Void` + // We use true here so the pipe receives something tangible. + root_stream.tick(Value::Bool(true)); + true + } else { + false // Exhausted + } + } + Ok(_) => panic!("create-ticker closure must return a boolean"), + Err(e) => panic!("create-ticker closure execution failed: {}", e), + } + }; + + // 4. Register the generator + ticker_generators.borrow_mut().push(Box::new(generator)); + + // 5. Return stream + Value::Object(Rc::new(stream_node)) + } + ); } #[cfg(test)] @@ -379,7 +435,6 @@ mod tests { #[test] fn test_barrier_sync() { - let root = RootStream::new(); // Pipe with 2 inputs let pipe = Rc::new(RefCell::new(PipeStream::new("barrier-pipe".to_string(), 2, None)));