Skip to main content

wickra_core/
traits.rs

1//! Core traits: the [`Indicator`] state machine and the [`BatchExt`] blanket extension.
2
3use crate::ohlcv::Candle;
4
5/// A streaming technical indicator.
6///
7/// Every indicator in Wickra implements this trait. The contract is:
8///
9/// - [`update`](Indicator::update) is called once per input point and must be O(1) in
10///   the input length. Pre-existing buffered state may be touched, but no full
11///   recomputation over the entire series is permitted.
12/// - The returned `Option<Output>` is `None` while the indicator is still in its
13///   *warmup* phase (insufficient inputs to produce a defined value), and `Some`
14///   once it is ready.
15/// - [`reset`](Indicator::reset) clears all state, returning the indicator to the
16///   exact configuration it had immediately after construction.
17///
18/// Implementors that consume scalar prices use `Input = f64` so they automatically
19/// gain access to chaining via [`Chain`].
20pub trait Indicator {
21    /// Type of one input data point (typically `f64` for a price, or `Candle` / `Tick`).
22    type Input;
23    /// Type of one output value.
24    type Output;
25
26    /// Feed one new data point into the indicator and return the freshly
27    /// computed output, or `None` if there is no value for this input.
28    ///
29    /// `None` covers exactly two cases:
30    ///
31    /// * the indicator is still warming up, and
32    /// * the input was rejected as non-finite.
33    ///
34    /// A rejected input is *skipped*: it does not enter the indicator's state,
35    /// so a single bad tick cannot corrupt the values that follow it. The
36    /// alternative — repeating the last computed value — was rejected because
37    /// it hands the caller a stale number that looks exactly like a fresh one.
38    fn update(&mut self, input: Self::Input) -> Option<Self::Output>;
39
40    /// Reset all internal state, leaving the indicator equivalent to a freshly
41    /// constructed instance with the same parameters.
42    fn reset(&mut self);
43
44    /// Number of inputs required before the first non-`None` output can be produced.
45    fn warmup_period(&self) -> usize;
46
47    /// Whether the indicator has emitted at least one value since the last reset.
48    fn is_ready(&self) -> bool;
49
50    /// Stable, human-readable indicator name. Used by chaining and diagnostics.
51    fn name(&self) -> &'static str;
52
53    /// Run the indicator over `inputs`, writing one output per input into the
54    /// caller-owned `out` buffer (`NaN` where [`update`](Indicator::update)
55    /// returns `None`). It exists for every indicator whose output converts to an
56    /// `f64`; the batch pipelines use it for the scalar `f64 -> f64` ones.
57    ///
58    /// This is the exact batch: every value is *bit-for-bit* the one a replay of
59    /// `update` produces, and the indicator is left in the state that replay
60    /// leaves it in. The default replays `update`; indicators with a faster exact
61    /// formulation override it, and every binding and extension method routes
62    /// through here, so the override reaches all of them. Writing into a buffer
63    /// the caller owns skips the output allocation, which on a large series costs
64    /// more than the arithmetic of a simple indicator.
65    ///
66    /// The bounds name only the associated types, so `Indicator` stays usable as
67    /// a trait object and the method, called through a `dyn Indicator`, reaches
68    /// the concrete indicator's override.
69    ///
70    /// # Panics
71    ///
72    /// Panics if `out.len() != inputs.len()`.
73    fn batch_nan_into(&mut self, inputs: &[Self::Input], out: &mut [f64])
74    where
75        Self::Input: Copy,
76        Self::Output: Into<f64>,
77    {
78        assert_eq!(
79            inputs.len(),
80            out.len(),
81            "batch output length must equal input length"
82        );
83        for (slot, &x) in out.iter_mut().zip(inputs) {
84            *slot = self.update(x).map_or(f64::NAN, Into::into);
85        }
86    }
87
88    /// Opt-in fast batch: like [`batch_nan_into`](Indicator::batch_nan_into), but
89    /// an indicator with a vectorised kernel may reassociate its arithmetic to
90    /// run it in SIMD lanes. Each value then agrees with the exact batch to within
91    /// the tolerance the indicator documents (a few units in the last place), not
92    /// bit for bit; warmup positions, `NaN` placement and the output length are
93    /// identical. The kernels are deterministic: the same input produces the same
94    /// bits on every platform, with or without SIMD hardware.
95    ///
96    /// The kernel only runs from a fresh (just constructed or reset) indicator
97    /// over an all-finite slice; any other call is served by the exact batch.
98    /// Afterwards the indicator continues streaming from the kernel's final
99    /// state. The default is the exact batch, so every scalar indicator offers
100    /// this method and the ones without a kernel simply return exact values.
101    ///
102    /// # Panics
103    ///
104    /// Panics if `out.len() != inputs.len()`.
105    fn batch_fast_into(&mut self, inputs: &[Self::Input], out: &mut [f64])
106    where
107        Self::Input: Copy,
108        Self::Output: Into<f64>,
109    {
110        self.batch_nan_into(inputs, out);
111    }
112}
113
114/// Blanket extension that adds batch evaluation to every [`Indicator`].
115///
116/// The naive `batch` simply replays `update` over a slice, which is always correct
117/// because `update` is the only state transition. Concrete indicators may override
118/// `batch` if they have a faster vectorized path; the default keeps the contract
119/// `batch == repeated update`.
120pub trait BatchExt: Indicator {
121    /// Run the indicator over a slice of inputs in order, returning one output (or
122    /// `None` during warmup) per input.
123    fn batch(&mut self, inputs: &[Self::Input]) -> Vec<Option<Self::Output>>
124    where
125        Self::Input: Clone,
126    {
127        let mut out = Vec::with_capacity(inputs.len());
128        for x in inputs {
129            out.push(self.update(x.clone()));
130        }
131        out
132    }
133
134    /// Run an independent copy of the indicator over each input series in parallel.
135    ///
136    /// Each asset is processed by its own fresh instance built via `make`, so state
137    /// never leaks across assets. Requires the `parallel` feature (enabled by
138    /// default), which pulls in `rayon`.
139    #[cfg(feature = "parallel")]
140    fn batch_parallel<F>(
141        inputs_per_asset: &[Vec<Self::Input>],
142        make: F,
143    ) -> Vec<Vec<Option<Self::Output>>>
144    where
145        Self: Sized + Send,
146        Self::Input: Sync + Clone,
147        Self::Output: Send,
148        F: Fn() -> Self + Sync + Send,
149    {
150        use rayon::prelude::*;
151        inputs_per_asset
152            .par_iter()
153            .map(|series| {
154                let mut ind = make();
155                ind.batch(series)
156            })
157            .collect()
158    }
159}
160
161impl<T: Indicator> BatchExt for T {}
162
163/// Fast batch for scalar `f64 -> f64` indicators.
164///
165/// The generic [`BatchExt::batch`] returns `Vec<Option<f64>>` — 16 bytes per
166/// element (no niche fits an arbitrary `f64`), which a caller wanting a dense
167/// `f64` series then has to walk a second time to map warmup `None`s to `NaN`.
168/// This skips both the wide intermediate and the second pass: one allocation,
169/// one pass, warmup encoded as `NaN`. Both methods allocate the result and fill
170/// it through [`Indicator::batch_nan_into`] / [`Indicator::batch_fast_into`],
171/// so an indicator's faster formulation reaches them too.
172pub trait BatchNanExt: Indicator<Input = f64, Output = f64> {
173    /// One `f64` per input, warmup positions filled with `NaN`, bit-for-bit equal
174    /// to replaying `update`.
175    fn batch_nan(&mut self, inputs: &[f64]) -> Vec<f64> {
176        let mut out = vec![0.0; inputs.len()];
177        self.batch_nan_into(inputs, &mut out);
178        out
179    }
180
181    /// The opt-in fast batch ([`Indicator::batch_fast_into`]) into a fresh
182    /// vector: within the indicator's documented tolerance of
183    /// [`batch_nan`](BatchNanExt::batch_nan), deterministic across platforms.
184    fn batch_fast(&mut self, inputs: &[f64]) -> Vec<f64> {
185        let mut out = vec![0.0; inputs.len()];
186        self.batch_fast_into(inputs, &mut out);
187        out
188    }
189}
190
191impl<T: Indicator<Input = f64, Output = f64>> BatchNanExt for T {}
192
193/// A streaming *bar builder* — an alternative-chart constructor (Renko, Kagi,
194/// Point-and-Figure) that turns a candle stream into a stream of price-driven
195/// bars.
196///
197/// Bar builders are deliberately **not** [`Indicator`]s: a single input candle
198/// may complete zero, one, or many bars (a large move can print several Renko
199/// bricks at once), which breaks the `update -> Option<Output>` one-in-one-out
200/// contract and the `batch == repeated update` length invariant. They get their
201/// own trait instead, returning a `Vec` of freshly completed bars per candle.
202///
203/// The contract is:
204///
205/// - [`update`](BarBuilder::update) ingests one candle and returns every bar it
206///   *completed* on that candle, in chronological order. An empty vector means
207///   the move was not large enough to finish a bar yet.
208/// - [`reset`](BarBuilder::reset) clears all state, returning the builder to the
209///   configuration it had immediately after construction.
210/// - [`batch`](BarBuilder::batch) concatenates the bars from replaying `update`
211///   over a slice; the flattened length is data-dependent, not the input length.
212///
213/// Bar builders cannot participate in [`Chain`] (which requires
214/// `Indicator<Input = f64, Output = f64>`); feed a downstream indicator from the
215/// bars' close prices manually if you need to chain off them.
216///
217/// ```text
218/// let mut renko = RenkoBars::new(1.0).unwrap();
219/// let bricks = renko.update(candle); // Vec<RenkoBrick>: 0..n completed bricks
220/// ```
221pub trait BarBuilder {
222    /// Type of one completed bar.
223    type Bar;
224
225    /// Feed one candle and return every bar completed on it (possibly none).
226    fn update(&mut self, candle: Candle) -> Vec<Self::Bar>;
227
228    /// Reset all internal state to the freshly-constructed configuration.
229    fn reset(&mut self);
230
231    /// Stable, human-readable builder name.
232    fn name(&self) -> &'static str;
233
234    /// Replay `update` over a slice, concatenating all completed bars. The
235    /// result length is data-dependent (not the input length).
236    fn batch(&mut self, candles: &[Candle]) -> Vec<Self::Bar> {
237        let mut out = Vec::new();
238        for candle in candles {
239            out.extend(self.update(*candle));
240        }
241        out
242    }
243}
244
245/// Chain two indicators so the output of the first becomes the input of the second.
246///
247/// Both indicators must agree on `f64` as the bridging type, which is the common
248/// case for price-in/value-out indicators. The chain itself is an indicator, so
249/// chains can be nested arbitrarily.
250///
251/// # Example
252///
253/// ```
254/// use wickra_core::{Chain, Ema, Indicator, Rsi};
255///
256/// // RSI(7) on top of EMA(14). EMA seeds at input 14, then RSI needs 7+1 more
257/// // valid inputs to emit, so the chain becomes ready at input 21.
258/// let mut chain = Chain::new(Ema::new(14).unwrap(), Rsi::new(7).unwrap());
259/// for i in 1..=21 {
260///     chain.update(f64::from(i));
261/// }
262/// assert!(chain.is_ready());
263/// ```
264#[derive(Debug, Clone)]
265pub struct Chain<A, B>
266where
267    A: Indicator<Input = f64, Output = f64>,
268    B: Indicator<Input = f64>,
269{
270    first: A,
271    second: B,
272}
273
274impl<A, B> Chain<A, B>
275where
276    A: Indicator<Input = f64, Output = f64>,
277    B: Indicator<Input = f64>,
278{
279    /// Construct a chain whose inputs flow through `first` and then `second`.
280    pub const fn new(first: A, second: B) -> Self {
281        Self { first, second }
282    }
283
284    /// Add a third stage on top.
285    pub fn then<C>(self, third: C) -> Chain<Self, C>
286    where
287        C: Indicator<Input = f64>,
288        Self: Indicator<Input = f64, Output = f64>,
289    {
290        Chain::new(self, third)
291    }
292
293    /// Borrow the upstream indicator.
294    pub const fn first(&self) -> &A {
295        &self.first
296    }
297
298    /// Borrow the downstream indicator.
299    pub const fn second(&self) -> &B {
300        &self.second
301    }
302}
303
304impl<A, B> Indicator for Chain<A, B>
305where
306    A: Indicator<Input = f64, Output = f64>,
307    B: Indicator<Input = f64>,
308{
309    type Input = f64;
310    type Output = B::Output;
311
312    fn update(&mut self, input: f64) -> Option<Self::Output> {
313        self.first.update(input).and_then(|v| self.second.update(v))
314    }
315
316    fn reset(&mut self) {
317        self.first.reset();
318        self.second.reset();
319    }
320
321    fn warmup_period(&self) -> usize {
322        // Not an upper bound: this method promises the input count before the
323        // first value, so over-declaring it is as wrong as under-declaring it.
324        // The second stage receives its first input on the bar the first stage
325        // emits, so the two warmups overlap by exactly one.
326        // A stage declaring 0 still needs its first input to produce anything,
327        // so each side counts as at least one bar before the overlap is taken
328        // off -- otherwise two pass-through stages underflow.
329        self.first.warmup_period().max(1) + self.second.warmup_period().max(1) - 1
330    }
331
332    fn is_ready(&self) -> bool {
333        self.first.is_ready() && self.second.is_ready()
334    }
335
336    fn name(&self) -> &'static str {
337        "Chain"
338    }
339}
340
341#[cfg(test)]
342mod tests {
343    use super::*;
344
345    /// A trivial test indicator: identity (passes input through).
346    #[derive(Debug, Default)]
347    struct Identity {
348        seen: bool,
349    }
350
351    impl Indicator for Identity {
352        type Input = f64;
353        type Output = f64;
354        fn update(&mut self, input: f64) -> Option<f64> {
355            self.seen = true;
356            Some(input)
357        }
358        fn reset(&mut self) {
359            self.seen = false;
360        }
361        fn warmup_period(&self) -> usize {
362            0
363        }
364        fn is_ready(&self) -> bool {
365            self.seen
366        }
367        fn name(&self) -> &'static str {
368            "Identity"
369        }
370    }
371
372    /// Another trivial test indicator: scales input by 2.
373    #[derive(Debug, Default)]
374    struct Doubler {
375        seen: bool,
376    }
377
378    impl Indicator for Doubler {
379        type Input = f64;
380        type Output = f64;
381        fn update(&mut self, input: f64) -> Option<f64> {
382            self.seen = true;
383            Some(input * 2.0)
384        }
385        fn reset(&mut self) {
386            self.seen = false;
387        }
388        fn warmup_period(&self) -> usize {
389            0
390        }
391        fn is_ready(&self) -> bool {
392            self.seen
393        }
394        fn name(&self) -> &'static str {
395            "Doubler"
396        }
397    }
398
399    #[test]
400    fn batch_replays_update() {
401        let mut id = Identity::default();
402        let out = id.batch(&[1.0, 2.0, 3.0]);
403        assert_eq!(out, vec![Some(1.0), Some(2.0), Some(3.0)]);
404    }
405
406    /// The blanket [`BatchNanExt::batch_nan`] default (used by every scalar
407    /// indicator without an inherent fast path) maps `update` outputs to a dense
408    /// `f64` series, warmup `None` becoming `NaN`. `Identity` is always ready, so
409    /// the result is just the inputs back.
410    #[test]
411    fn batch_nan_default_maps_none_to_nan() {
412        let mut id = Identity::default();
413        let out = id.batch_nan(&[1.0, 2.0, 3.0]);
414        assert_eq!(out, vec![1.0, 2.0, 3.0]);
415    }
416
417    /// The default `batch_nan_into` writes one value per input into the
418    /// caller's buffer, overwriting whatever was there.
419    #[test]
420    fn batch_nan_into_default_fills_caller_buffer() {
421        let mut id = Identity::default();
422        let mut out = [f64::INFINITY; 3];
423        id.batch_nan_into(&[4.0, 5.0, 6.0], &mut out);
424        assert_eq!(out, [4.0, 5.0, 6.0]);
425    }
426
427    /// A length mismatch is a caller bug, not something to truncate silently.
428    #[test]
429    #[should_panic(expected = "batch output length must equal input length")]
430    fn batch_nan_into_rejects_mismatched_lengths() {
431        let mut id = Identity::default();
432        let mut out = [0.0; 2];
433        id.batch_nan_into(&[1.0, 2.0, 3.0], &mut out);
434    }
435
436    /// An indicator without a vectorised kernel serves `batch_fast` from the
437    /// exact batch, so the two agree bit for bit.
438    #[test]
439    fn batch_fast_default_is_the_exact_batch() {
440        let exact = Doubler::default().batch_nan(&[1.5, -2.0, 3.25]);
441        let fast = Doubler::default().batch_fast(&[1.5, -2.0, 3.25]);
442        assert_eq!(exact, fast);
443        let mut out = [0.0; 3];
444        Doubler::default().batch_fast_into(&[1.5, -2.0, 3.25], &mut out);
445        assert_eq!(out.to_vec(), exact);
446    }
447
448    /// `Indicator` and `BatchNanExt` stay usable as trait objects, and a batch
449    /// called through one reaches the concrete indicator's implementation.
450    #[test]
451    fn batch_traits_are_dyn_compatible() {
452        let mut boxed: Box<dyn BatchNanExt<Input = f64, Output = f64>> =
453            Box::new(Doubler::default());
454        assert_eq!(boxed.batch_nan(&[1.0, 2.0]), vec![2.0, 4.0]);
455        assert_eq!(boxed.batch_fast(&[3.0]), vec![6.0]);
456        let series: Vec<f64> = (0..64).map(|i| f64::from(i % 7) * 1.25 + 2.0).collect();
457        let mut plain: Box<dyn Indicator<Input = f64, Output = f64>> =
458            Box::new(crate::Sma::new(5).unwrap());
459        let mut exact = vec![0.0; series.len()];
460        plain.batch_nan_into(&series, &mut exact);
461        let bits = |v: &[f64]| v.iter().map(|x| x.to_bits()).collect::<Vec<_>>();
462        let want = crate::Sma::new(5).unwrap().batch_nan(&series);
463        assert_eq!(bits(&exact), bits(&want));
464        plain.reset();
465        let mut fast = vec![0.0; series.len()];
466        plain.batch_fast_into(&series, &mut fast);
467        let want = crate::Sma::new(5).unwrap().batch_fast(&series);
468        assert_eq!(bits(&fast), bits(&want));
469    }
470
471    #[test]
472    fn chain_pipes_first_into_second() {
473        let mut c = Chain::new(Doubler::default(), Doubler::default());
474        // 5 -> 10 -> 20
475        assert_eq!(c.update(5.0), Some(20.0));
476    }
477
478    #[test]
479    fn chain_is_ready_only_after_both_stages_emit() {
480        let mut c = Chain::new(Doubler::default(), Doubler::default());
481        assert!(!c.is_ready());
482        c.update(1.0);
483        assert!(c.is_ready());
484    }
485
486    #[test]
487    fn chain_reset_propagates() {
488        let mut c = Chain::new(Doubler::default(), Doubler::default());
489        c.update(1.0);
490        assert!(c.is_ready());
491        c.reset();
492        assert!(!c.is_ready());
493    }
494
495    #[test]
496    fn chain_three_levels_via_then() {
497        let c = Chain::new(Doubler::default(), Doubler::default()).then(Doubler::default());
498        let mut c = c;
499        // 1 -> 2 -> 4 -> 8
500        assert_eq!(c.update(1.0), Some(8.0));
501    }
502
503    /// Cover the `Chain::first` / `Chain::second` borrow accessors and the
504    /// `Chain::warmup_period` + `Chain::name` Indicator-impl bodies.
505    ///
506    /// Existing chain tests only invoked the Indicator surface (`update`,
507    /// `reset`, `is_ready`) on the wrapped `Chain`. The const borrow accessors
508    /// and the `warmup_period` / `name` impls were never traversed, so Codecov
509    /// flagged traits.rs lines 140-142, 145-147, 167-170, 176-178 as missed.
510    /// `chain.warmup_period()` also reaches `Doubler::warmup_period`
511    /// (228-230), and `chain.first().name()` reaches `Doubler::name`
512    /// (234-236) — both helper methods were uncovered for the same reason.
513    #[test]
514    fn chain_accessors_and_metadata() {
515        let chain = Chain::new(Doubler::default(), Doubler::default());
516        // Borrow accessors return the wrapped stages; query each via .name()
517        // so Doubler::name (lines 234-236) is also exercised.
518        assert_eq!(chain.first().name(), "Doubler");
519        assert_eq!(chain.second().name(), "Doubler");
520        // Doubler::warmup_period (lines 228-230) is 0, meaning it emits on its
521        // first input; chaining two of them still emits on the first input.
522        assert_eq!(chain.first().warmup_period(), 0);
523        assert_eq!(chain.second().warmup_period(), 0);
524        assert_eq!(chain.warmup_period(), 1);
525        // Chain::name returns the literal "Chain" (line 177).
526        assert_eq!(chain.name(), "Chain");
527    }
528
529    /// Cover the full Indicator surface of the `Identity` test helper:
530    /// `reset` (198-200), `warmup_period` (201-203), `is_ready` (204-206),
531    /// and `name` (207-209). The only other test using `Identity`
532    /// (`batch_replays_update`) calls `batch`, which exercises `update`
533    /// alone, leaving the remaining four trait methods uncovered.
534    #[test]
535    fn identity_helper_full_indicator_surface() {
536        let mut id = Identity::default();
537        // warmup_period is the literal 0; name is the literal "Identity".
538        assert_eq!(id.warmup_period(), 0);
539        assert_eq!(id.name(), "Identity");
540        // is_ready exercises the `self.seen` return with seen=false first…
541        assert!(!id.is_ready());
542        // …then with seen=true after a single update.
543        let out = id.update(42.0);
544        assert_eq!(out, Some(42.0));
545        assert!(id.is_ready());
546        // reset() flips seen back to false; is_ready reflects it.
547        id.reset();
548        assert!(!id.is_ready());
549    }
550
551    #[cfg(feature = "parallel")]
552    #[test]
553    fn batch_parallel_runs_independent_instances() {
554        let series: Vec<Vec<f64>> = vec![vec![1.0, 2.0, 3.0], vec![4.0, 5.0, 6.0]];
555        let out = Doubler::batch_parallel(&series, Doubler::default);
556        assert_eq!(out.len(), 2);
557        assert_eq!(out[0], vec![Some(2.0), Some(4.0), Some(6.0)]);
558        assert_eq!(out[1], vec![Some(8.0), Some(10.0), Some(12.0)]);
559    }
560}