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.
This commit is contained in:
@@ -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
|
||||||
|
)
|
||||||
+56
-1
@@ -341,6 +341,62 @@ pub fn register(env: &Environment) {
|
|||||||
Value::Object(Rc::new(stream_node))
|
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<Value>| {
|
||||||
|
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)]
|
#[cfg(test)]
|
||||||
@@ -379,7 +435,6 @@ mod tests {
|
|||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_barrier_sync() {
|
fn test_barrier_sync() {
|
||||||
let root = RootStream::new();
|
|
||||||
// Pipe with 2 inputs
|
// Pipe with 2 inputs
|
||||||
let pipe = Rc::new(RefCell::new(PipeStream::new("barrier-pipe".to_string(), 2, None)));
|
let pipe = Rc::new(RefCell::new(PipeStream::new("barrier-pipe".to_string(), 2, None)));
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user