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}