Skip to main content

qs_backtest/strategy/
analysis.rs

1//! Causal historical observations and complete-boundary analysis.
2
3use std::collections::{BTreeSet, VecDeque};
4use std::fmt;
5
6use chrono::NaiveDateTime;
7use serde::{Deserialize, Deserializer, Serialize};
8
9use super::annotation::{AnnotationId, AnnotationLimits, AnnotationTimeline, StrategyAnnotation};
10use super::{ClosedBar, HistoricalSeriesView, SeriesId, SeriesViewError};
11
12pub const MAX_ANALYZERS: usize = 256;
13pub const MAX_RETAINED_OBSERVATIONS: usize = 1_000_000;
14pub const MAX_OBSERVATIONS_PER_BOUNDARY: usize = 4096;
15pub const MAX_OBSERVATION_SOURCE_SERIES: usize = 64;
16pub const MAX_ZONE_ID_BYTES: usize = 64;
17pub const MAX_PIVOT_SIDE_BARS: usize = 1_000_000;
18
19/// Stable caller-supplied identity for one price zone.
20#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
21#[serde(transparent)]
22pub struct ZoneId(String);
23
24impl ZoneId {
25    pub fn new(value: impl Into<String>) -> Result<Self, AnalysisError> {
26        let value = value.into();
27        if valid_identifier(&value, MAX_ZONE_ID_BYTES) {
28            Ok(Self(value))
29        } else {
30            Err(AnalysisError::InvalidZoneId)
31        }
32    }
33
34    pub fn as_str(&self) -> &str {
35        &self.0
36    }
37}
38
39impl fmt::Display for ZoneId {
40    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
41        formatter.write_str(&self.0)
42    }
43}
44
45impl<'de> Deserialize<'de> for ZoneId {
46    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
47    where
48        D: Deserializer<'de>,
49    {
50        let value = String::deserialize(deserializer)?;
51        Self::new(value).map_err(serde::de::Error::custom)
52    }
53}
54
55/// Provenance category for a descriptive price zone.
56#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
57#[serde(rename_all = "snake_case")]
58pub enum ZoneSource {
59    CausalAnnotation,
60    DeterministicAnalyzer,
61}
62
63/// Descriptive side of a price zone.
64#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(rename_all = "snake_case")]
66pub enum ZoneSide {
67    Support,
68    Resistance,
69}
70
71/// Descriptive lifecycle state of an immutable zone snapshot.
72#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
73#[serde(rename_all = "snake_case")]
74pub enum ZoneState {
75    Active,
76    Broken,
77    RetestPending,
78    Retested,
79    Choppy,
80    Degraded,
81    Invalid,
82}
83
84/// One validated immutable price-zone snapshot.
85#[derive(Debug, Clone, PartialEq, Serialize)]
86pub struct PriceZone {
87    zone_id: ZoneId,
88    side: ZoneSide,
89    lower: f64,
90    upper: f64,
91    created_at: NaiveDateTime,
92    touch_count: u32,
93    state: ZoneState,
94    source: ZoneSource,
95}
96
97impl PriceZone {
98    #[allow(clippy::too_many_arguments)]
99    pub fn new(
100        zone_id: ZoneId,
101        side: ZoneSide,
102        lower: f64,
103        upper: f64,
104        created_at: NaiveDateTime,
105        touch_count: u32,
106        state: ZoneState,
107        source: ZoneSource,
108    ) -> Result<Self, AnalysisError> {
109        validate_price(lower, "zone lower")?;
110        validate_price(upper, "zone upper")?;
111        if lower >= upper {
112            return Err(AnalysisError::InvalidZoneGeometry { lower, upper });
113        }
114        Ok(Self {
115            zone_id,
116            side,
117            lower,
118            upper,
119            created_at,
120            touch_count,
121            state,
122            source,
123        })
124    }
125
126    pub fn zone_id(&self) -> &ZoneId {
127        &self.zone_id
128    }
129
130    pub fn side(&self) -> ZoneSide {
131        self.side
132    }
133
134    pub fn lower(&self) -> f64 {
135        self.lower
136    }
137
138    pub fn upper(&self) -> f64 {
139        self.upper
140    }
141
142    pub fn created_at(&self) -> NaiveDateTime {
143        self.created_at
144    }
145
146    pub fn touch_count(&self) -> u32 {
147        self.touch_count
148    }
149
150    pub fn state(&self) -> ZoneState {
151        self.state
152    }
153
154    pub fn source(&self) -> ZoneSource {
155        self.source
156    }
157}
158
159#[derive(Deserialize)]
160#[serde(deny_unknown_fields)]
161struct PriceZoneDef {
162    zone_id: ZoneId,
163    side: ZoneSide,
164    lower: f64,
165    upper: f64,
166    created_at: NaiveDateTime,
167    touch_count: u32,
168    state: ZoneState,
169    source: ZoneSource,
170}
171
172impl<'de> Deserialize<'de> for PriceZone {
173    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
174    where
175        D: Deserializer<'de>,
176    {
177        let value = PriceZoneDef::deserialize(deserializer)?;
178        Self::new(
179            value.zone_id,
180            value.side,
181            value.lower,
182            value.upper,
183            value.created_at,
184            value.touch_count,
185            value.state,
186            value.source,
187        )
188        .map_err(serde::de::Error::custom)
189    }
190}
191
192/// High or low classification for a confirmed swing point.
193#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
194#[serde(rename_all = "snake_case")]
195pub enum SwingKind {
196    High,
197    Low,
198}
199
200/// One immutable swing point with explicit anchor and confirmation times.
201#[derive(Debug, Clone, PartialEq, Serialize)]
202pub struct SwingPoint {
203    kind: SwingKind,
204    price: f64,
205    anchor_open_time: NaiveDateTime,
206    anchor_close_time: NaiveDateTime,
207    confirmed_at: NaiveDateTime,
208}
209
210impl SwingPoint {
211    pub fn new(
212        kind: SwingKind,
213        price: f64,
214        anchor_open_time: NaiveDateTime,
215        anchor_close_time: NaiveDateTime,
216        confirmed_at: NaiveDateTime,
217    ) -> Result<Self, AnalysisError> {
218        validate_price(price, "swing price")?;
219        if anchor_open_time >= anchor_close_time {
220            return Err(AnalysisError::InvalidAnchorRange {
221                open: anchor_open_time,
222                close: anchor_close_time,
223            });
224        }
225        if confirmed_at < anchor_close_time {
226            return Err(AnalysisError::ConfirmationBeforeAnchorClose {
227                confirmed_at,
228                anchor_close: anchor_close_time,
229            });
230        }
231        Ok(Self {
232            kind,
233            price,
234            anchor_open_time,
235            anchor_close_time,
236            confirmed_at,
237        })
238    }
239
240    pub fn kind(&self) -> SwingKind {
241        self.kind
242    }
243
244    pub fn price(&self) -> f64 {
245        self.price
246    }
247
248    pub fn anchor_open_time(&self) -> NaiveDateTime {
249        self.anchor_open_time
250    }
251
252    pub fn anchor_close_time(&self) -> NaiveDateTime {
253        self.anchor_close_time
254    }
255
256    pub fn confirmed_at(&self) -> NaiveDateTime {
257        self.confirmed_at
258    }
259}
260
261#[derive(Deserialize)]
262#[serde(deny_unknown_fields)]
263struct SwingPointDef {
264    kind: SwingKind,
265    price: f64,
266    anchor_open_time: NaiveDateTime,
267    anchor_close_time: NaiveDateTime,
268    confirmed_at: NaiveDateTime,
269}
270
271impl<'de> Deserialize<'de> for SwingPoint {
272    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
273    where
274        D: Deserializer<'de>,
275    {
276        let value = SwingPointDef::deserialize(deserializer)?;
277        Self::new(
278            value.kind,
279            value.price,
280            value.anchor_open_time,
281            value.anchor_close_time,
282            value.confirmed_at,
283        )
284        .map_err(serde::de::Error::custom)
285    }
286}
287
288/// Common descriptive rejection patterns.
289#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
290#[serde(rename_all = "snake_case")]
291pub enum RejectionPattern {
292    LongWick,
293    Engulfing,
294    DoubleTouch,
295    SnapBackInside,
296}
297
298/// Common descriptive momentum states.
299#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
300#[serde(rename_all = "snake_case")]
301pub enum MomentumState {
302    Advancing,
303    Stalling,
304    Sideways,
305    Reversing,
306}
307
308/// Small common observation vocabulary shared by historical strategies.
309#[derive(Debug, Clone, PartialEq, Serialize)]
310#[serde(tag = "kind", content = "value", rename_all = "snake_case")]
311pub enum StrategyObservationValue {
312    Zone(PriceZone),
313    Swing(SwingPoint),
314    Rejection {
315        pattern: RejectionPattern,
316        anchor_open_time: NaiveDateTime,
317        anchor_close_time: NaiveDateTime,
318    },
319    Momentum(MomentumState),
320}
321
322impl StrategyObservationValue {
323    pub fn rejection(
324        pattern: RejectionPattern,
325        anchor_open_time: NaiveDateTime,
326        anchor_close_time: NaiveDateTime,
327    ) -> Result<Self, AnalysisError> {
328        if anchor_open_time >= anchor_close_time {
329            return Err(AnalysisError::InvalidAnchorRange {
330                open: anchor_open_time,
331                close: anchor_close_time,
332            });
333        }
334        Ok(Self::Rejection {
335            pattern,
336            anchor_open_time,
337            anchor_close_time,
338        })
339    }
340
341    pub(crate) fn validate_at(&self, observed_through: NaiveDateTime) -> Result<(), AnalysisError> {
342        match self {
343            Self::Zone(zone) if zone.created_at() > observed_through => {
344                Err(AnalysisError::ValueAfterObservation {
345                    value_time: zone.created_at(),
346                    observed_through,
347                })
348            }
349            Self::Swing(swing) if swing.confirmed_at() > observed_through => {
350                Err(AnalysisError::ValueAfterObservation {
351                    value_time: swing.confirmed_at(),
352                    observed_through,
353                })
354            }
355            Self::Rejection {
356                anchor_close_time, ..
357            } if *anchor_close_time > observed_through => {
358                Err(AnalysisError::ValueAfterObservation {
359                    value_time: *anchor_close_time,
360                    observed_through,
361                })
362            }
363            _ => Ok(()),
364        }
365    }
366
367    pub fn zone(&self) -> Option<&PriceZone> {
368        match self {
369            Self::Zone(zone) => Some(zone),
370            _ => None,
371        }
372    }
373
374    pub fn swing(&self) -> Option<&SwingPoint> {
375        match self {
376            Self::Swing(swing) => Some(swing),
377            _ => None,
378        }
379    }
380}
381
382#[derive(Deserialize)]
383#[serde(
384    tag = "kind",
385    content = "value",
386    rename_all = "snake_case",
387    deny_unknown_fields
388)]
389enum StrategyObservationValueDef {
390    Zone(PriceZone),
391    Swing(SwingPoint),
392    Rejection {
393        pattern: RejectionPattern,
394        anchor_open_time: NaiveDateTime,
395        anchor_close_time: NaiveDateTime,
396    },
397    Momentum(MomentumState),
398}
399
400impl<'de> Deserialize<'de> for StrategyObservationValue {
401    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
402    where
403        D: Deserializer<'de>,
404    {
405        match StrategyObservationValueDef::deserialize(deserializer)? {
406            StrategyObservationValueDef::Zone(value) => Ok(Self::Zone(value)),
407            StrategyObservationValueDef::Swing(value) => Ok(Self::Swing(value)),
408            StrategyObservationValueDef::Rejection {
409                pattern,
410                anchor_open_time,
411                anchor_close_time,
412            } => Self::rejection(pattern, anchor_open_time, anchor_close_time)
413                .map_err(serde::de::Error::custom),
414            StrategyObservationValueDef::Momentum(value) => Ok(Self::Momentum(value)),
415        }
416    }
417}
418
419/// Authoritative source category for a committed observation.
420#[derive(Debug, Clone, PartialEq, Eq)]
421pub enum ObservationOrigin {
422    Analyzer,
423    CausalAnnotation { annotation_id: AnnotationId },
424}
425
426/// One immutable committed causal observation.
427#[derive(Debug, Clone, PartialEq)]
428pub struct StrategyObservation {
429    sequence: u64,
430    observed_through: NaiveDateTime,
431    valid_from: NaiveDateTime,
432    symbol: String,
433    source_series: Vec<SeriesId>,
434    origin: ObservationOrigin,
435    value: StrategyObservationValue,
436}
437
438impl StrategyObservation {
439    pub fn sequence(&self) -> u64 {
440        self.sequence
441    }
442
443    pub fn observed_through(&self) -> NaiveDateTime {
444        self.observed_through
445    }
446
447    pub fn valid_from(&self) -> NaiveDateTime {
448        self.valid_from
449    }
450
451    pub fn symbol(&self) -> &str {
452        &self.symbol
453    }
454
455    pub fn source_series(&self) -> &[SeriesId] {
456        &self.source_series
457    }
458
459    pub fn origin(&self) -> &ObservationOrigin {
460        &self.origin
461    }
462
463    pub fn value(&self) -> &StrategyObservationValue {
464        &self.value
465    }
466
467    pub(crate) fn from_annotation(
468        sequence: u64,
469        observed_through: NaiveDateTime,
470        annotation: &StrategyAnnotation,
471    ) -> Result<Self, AnalysisError> {
472        let valid_from = annotation
473            .valid_from()
474            .expect("only causal annotations become observations");
475        Self::validated(
476            sequence,
477            observed_through,
478            valid_from,
479            annotation.symbol().to_string(),
480            annotation.source_series().to_vec(),
481            ObservationOrigin::CausalAnnotation {
482                annotation_id: annotation.annotation_id().clone(),
483            },
484            annotation.value().clone(),
485        )
486    }
487
488    fn from_draft(
489        sequence: u64,
490        observed_through: NaiveDateTime,
491        draft: StrategyObservationDraft,
492    ) -> Result<Self, AnalysisError> {
493        Self::validated(
494            sequence,
495            observed_through,
496            observed_through,
497            draft.symbol,
498            draft.source_series,
499            ObservationOrigin::Analyzer,
500            draft.value,
501        )
502    }
503
504    #[allow(clippy::too_many_arguments)]
505    fn validated(
506        sequence: u64,
507        observed_through: NaiveDateTime,
508        valid_from: NaiveDateTime,
509        symbol: String,
510        source_series: Vec<SeriesId>,
511        origin: ObservationOrigin,
512        value: StrategyObservationValue,
513    ) -> Result<Self, AnalysisError> {
514        validate_symbol(&symbol)?;
515        validate_source_series(&source_series, MAX_OBSERVATION_SOURCE_SERIES)?;
516        if valid_from > observed_through {
517            return Err(AnalysisError::ObservationNotYetValid {
518                valid_from,
519                observed_through,
520            });
521        }
522        value.validate_at(observed_through)?;
523        Ok(Self {
524            sequence,
525            observed_through,
526            valid_from,
527            symbol,
528            source_series,
529            origin,
530            value,
531        })
532    }
533}
534
535/// Metadata-free analyzer output validated and stamped by the pipeline.
536#[derive(Debug, Clone, PartialEq)]
537pub struct StrategyObservationDraft {
538    symbol: String,
539    source_series: Vec<SeriesId>,
540    value: StrategyObservationValue,
541}
542
543impl StrategyObservationDraft {
544    pub fn new(
545        symbol: impl Into<String>,
546        source_series: Vec<SeriesId>,
547        value: StrategyObservationValue,
548    ) -> Result<Self, AnalysisError> {
549        let symbol = symbol.into();
550        validate_symbol(&symbol)?;
551        validate_source_series(&source_series, MAX_OBSERVATION_SOURCE_SERIES)?;
552        Ok(Self {
553            symbol,
554            source_series,
555            value,
556        })
557    }
558
559    pub fn symbol(&self) -> &str {
560        &self.symbol
561    }
562
563    pub fn source_series(&self) -> &[SeriesId] {
564        &self.source_series
565    }
566
567    pub fn value(&self) -> &StrategyObservationValue {
568        &self.value
569    }
570}
571
572/// Caller-visible bounds for strategy-visible observation history.
573#[derive(Debug, Clone, Copy, PartialEq, Eq)]
574pub struct ObservationStoreLimits {
575    max_retained: usize,
576    max_per_boundary: usize,
577}
578
579impl ObservationStoreLimits {
580    pub fn new(max_retained: usize, max_per_boundary: usize) -> Result<Self, AnalysisError> {
581        validate_nonzero_limit("max_retained", max_retained, MAX_RETAINED_OBSERVATIONS)?;
582        validate_nonzero_limit(
583            "max_per_boundary",
584            max_per_boundary,
585            MAX_OBSERVATIONS_PER_BOUNDARY,
586        )?;
587        Ok(Self {
588            max_retained,
589            max_per_boundary,
590        })
591    }
592
593    pub fn max_retained(self) -> usize {
594        self.max_retained
595    }
596
597    pub fn max_per_boundary(self) -> usize {
598        self.max_per_boundary
599    }
600}
601
602impl Default for ObservationStoreLimits {
603    fn default() -> Self {
604        Self {
605            max_retained: 10_000,
606            max_per_boundary: 256,
607        }
608    }
609}
610
611/// Allocation-free suffix view over a possibly wrapped observation deque.
612#[derive(Debug, Clone, Copy)]
613pub struct ObservationWindow<'a> {
614    older: &'a [StrategyObservation],
615    newer: &'a [StrategyObservation],
616}
617
618impl<'a> ObservationWindow<'a> {
619    pub fn len(&self) -> usize {
620        self.older.len() + self.newer.len()
621    }
622
623    pub fn is_empty(&self) -> bool {
624        self.older.is_empty() && self.newer.is_empty()
625    }
626
627    pub fn latest(&self) -> Option<&'a StrategyObservation> {
628        self.newer.last().or_else(|| self.older.last())
629    }
630
631    pub fn iter(&self) -> impl DoubleEndedIterator<Item = &'a StrategyObservation> {
632        self.older.iter().chain(self.newer.iter())
633    }
634}
635
636/// Borrowed deterministic symbol-filtered selection.
637#[derive(Debug, Clone, Copy)]
638pub struct ObservationSelection<'a> {
639    window: ObservationWindow<'a>,
640    symbol: &'a str,
641    count: usize,
642}
643
644impl<'a> ObservationSelection<'a> {
645    pub fn iter(&self) -> impl DoubleEndedIterator<Item = &'a StrategyObservation> {
646        self.window
647            .iter()
648            .rev()
649            .filter(move |observation| observation.symbol() == self.symbol)
650            .take(self.count)
651            .collect::<Vec<_>>()
652            .into_iter()
653            .rev()
654    }
655}
656
657/// Object-safe read-only view over committed causal observations.
658pub trait HistoricalObservationView {
659    fn observations(&self, count: usize) -> ObservationWindow<'_>;
660    fn for_symbol<'a>(&'a self, symbol: &'a str, count: usize) -> ObservationSelection<'a>;
661    fn latest_zone(&self, zone_id: &ZoneId) -> Option<&StrategyObservation>;
662    fn omitted(&self) -> u64;
663}
664
665/// Bounded strategy-visible working history.
666#[derive(Debug, Clone)]
667pub struct ObservationStore {
668    retained: VecDeque<StrategyObservation>,
669    omitted: u64,
670    limits: ObservationStoreLimits,
671}
672
673impl ObservationStore {
674    pub fn new(limits: ObservationStoreLimits) -> Self {
675        Self {
676            retained: VecDeque::with_capacity(limits.max_retained),
677            omitted: 0,
678            limits,
679        }
680    }
681
682    pub fn limits(&self) -> ObservationStoreLimits {
683        self.limits
684    }
685
686    pub fn len(&self) -> usize {
687        self.retained.len()
688    }
689
690    pub fn is_empty(&self) -> bool {
691        self.retained.is_empty()
692    }
693
694    pub fn omitted(&self) -> u64 {
695        self.omitted
696    }
697
698    pub fn observations(&self, count: usize) -> ObservationWindow<'_> {
699        HistoricalObservationView::observations(self, count)
700    }
701
702    pub fn for_symbol<'a>(&'a self, symbol: &'a str, count: usize) -> ObservationSelection<'a> {
703        HistoricalObservationView::for_symbol(self, symbol, count)
704    }
705
706    pub fn latest_zone(&self, zone_id: &ZoneId) -> Option<&StrategyObservation> {
707        HistoricalObservationView::latest_zone(self, zone_id)
708    }
709
710    fn push(&mut self, observation: StrategyObservation) -> Result<(), AnalysisError> {
711        if self.retained.len() == self.limits.max_retained {
712            self.omitted = self
713                .omitted
714                .checked_add(1)
715                .ok_or(AnalysisError::OmittedCountOverflow)?;
716            self.retained.pop_front();
717        }
718        self.retained.push_back(observation);
719        Ok(())
720    }
721}
722
723impl HistoricalObservationView for ObservationStore {
724    fn observations(&self, count: usize) -> ObservationWindow<'_> {
725        let (older, newer) = self.retained.as_slices();
726        let available = count.min(self.retained.len());
727        let skip = self.retained.len() - available;
728        if skip < older.len() {
729            ObservationWindow {
730                older: &older[skip..],
731                newer,
732            }
733        } else {
734            ObservationWindow {
735                older: &older[older.len()..],
736                newer: &newer[skip - older.len()..],
737            }
738        }
739    }
740
741    fn for_symbol<'a>(&'a self, symbol: &'a str, count: usize) -> ObservationSelection<'a> {
742        ObservationSelection {
743            window: self.observations(usize::MAX),
744            symbol,
745            count,
746        }
747    }
748
749    fn latest_zone(&self, zone_id: &ZoneId) -> Option<&StrategyObservation> {
750        self.retained.iter().rev().find(|observation| {
751            observation
752                .value()
753                .zone()
754                .is_some_and(|zone| zone.zone_id() == zone_id)
755        })
756    }
757
758    fn omitted(&self) -> u64 {
759        self.omitted
760    }
761}
762
763/// Complete historical boundary supplied after a closed-bar series commits one timestamp batch.
764pub struct AnalysisBoundary<'a> {
765    observed_through: NaiveDateTime,
766    closed_bars: &'a [ClosedBar],
767    series: &'a dyn HistoricalSeriesView,
768}
769
770impl<'a> AnalysisBoundary<'a> {
771    pub fn new(
772        observed_through: NaiveDateTime,
773        closed_bars: &'a [ClosedBar],
774        series: &'a dyn HistoricalSeriesView,
775    ) -> Self {
776        Self {
777            observed_through,
778            closed_bars,
779            series,
780        }
781    }
782
783    pub fn observed_through(&self) -> NaiveDateTime {
784        self.observed_through
785    }
786
787    pub fn closed_bars(&self) -> &'a [ClosedBar] {
788        self.closed_bars
789    }
790
791    pub fn series(&self) -> &'a dyn HistoricalSeriesView {
792        self.series
793    }
794}
795
796/// Read-only causal state visible during one analyzer callback.
797#[derive(Clone, Copy)]
798pub struct AnalysisContext<'a> {
799    observed_through: NaiveDateTime,
800    series: &'a dyn HistoricalSeriesView,
801    observations: &'a dyn HistoricalObservationView,
802}
803
804impl<'a> AnalysisContext<'a> {
805    pub fn observed_through(self) -> NaiveDateTime {
806        self.observed_through
807    }
808
809    pub fn series(self) -> &'a dyn HistoricalSeriesView {
810        self.series
811    }
812
813    pub fn observations(self) -> &'a dyn HistoricalObservationView {
814        self.observations
815    }
816}
817
818/// Synchronous extension point for causal bar analysis.
819///
820/// An analyzer must be movable between threads, for the same reason a named-input projector must: it is a deterministic accumulation over borrowed bars, and the same implementation has to serve a replay running on a worker thread and a live instance running in its own task.
821pub trait HistoricalAnalyzer: Send {
822    fn on_bar(
823        &mut self,
824        bar: &ClosedBar,
825        context: AnalysisContext<'_>,
826    ) -> Result<Vec<StrategyObservationDraft>, AnalysisError>;
827}
828
829/// Observations atomically committed by one complete boundary.
830#[derive(Debug, Clone, PartialEq)]
831pub struct AnalysisBoundaryOutput {
832    observed_through: NaiveDateTime,
833    observations: Vec<StrategyObservation>,
834}
835
836impl AnalysisBoundaryOutput {
837    pub fn observed_through(&self) -> NaiveDateTime {
838        self.observed_through
839    }
840
841    pub fn observations(&self) -> &[StrategyObservation] {
842        &self.observations
843    }
844}
845
846/// Complete-boundary analyzer pipeline with bounded committed history.
847pub struct AnalysisPipeline {
848    analyzers: Vec<Box<dyn HistoricalAnalyzer>>,
849    observations: ObservationStore,
850    annotations: AnnotationTimeline,
851    next_sequence: u64,
852    last_boundary: Option<NaiveDateTime>,
853    failed: bool,
854}
855
856impl AnalysisPipeline {
857    pub fn new(
858        analyzers: Vec<Box<dyn HistoricalAnalyzer>>,
859        observation_limits: ObservationStoreLimits,
860        annotation_limits: AnnotationLimits,
861    ) -> Result<Self, AnalysisError> {
862        if analyzers.len() > MAX_ANALYZERS {
863            return Err(AnalysisError::TooManyAnalyzers {
864                actual: analyzers.len(),
865                maximum: MAX_ANALYZERS,
866            });
867        }
868        Ok(Self {
869            analyzers,
870            observations: ObservationStore::new(observation_limits),
871            annotations: AnnotationTimeline::new(annotation_limits),
872            next_sequence: 0,
873            last_boundary: None,
874            failed: false,
875        })
876    }
877
878    pub fn add_annotation(&mut self, annotation: StrategyAnnotation) -> Result<(), AnalysisError> {
879        if self.failed {
880            return Err(AnalysisError::PipelineFailed);
881        }
882        self.annotations
883            .add(annotation, self.last_boundary)
884            .map_err(AnalysisError::Annotation)
885    }
886
887    pub fn observations(&self) -> &ObservationStore {
888        &self.observations
889    }
890
891    pub fn annotations(&self) -> &AnnotationTimeline {
892        &self.annotations
893    }
894
895    pub(crate) fn into_research_annotations(self) -> Vec<StrategyAnnotation> {
896        self.annotations.into_research_only()
897    }
898
899    pub fn is_failed(&self) -> bool {
900        self.failed
901    }
902
903    pub fn on_boundary(
904        &mut self,
905        boundary: AnalysisBoundary<'_>,
906    ) -> Result<AnalysisBoundaryOutput, AnalysisError> {
907        if self.failed {
908            return Err(AnalysisError::PipelineFailed);
909        }
910        self.validate_boundary(&boundary)?;
911
912        let activates_annotation = self
913            .annotations
914            .pending_causal()
915            .first()
916            .and_then(|annotation| annotation.valid_from())
917            .is_some_and(|valid_from| valid_from <= boundary.observed_through);
918        if boundary.closed_bars.is_empty() && !activates_annotation {
919            self.last_boundary = Some(boundary.observed_through);
920            return Ok(AnalysisBoundaryOutput {
921                observed_through: boundary.observed_through,
922                observations: Vec::new(),
923            });
924        }
925
926        let mut staged_store = self.observations.clone();
927        let mut staged_annotations = self.annotations.clone();
928        let mut staged_sequence = self.next_sequence;
929        let mut committed = Vec::new();
930        let eligible = staged_annotations
931            .activate(boundary.observed_through)
932            .map_err(AnalysisError::Annotation)?;
933        self.ensure_boundary_capacity(eligible.len())?;
934        for annotation in eligible {
935            let sequence = take_sequence(&mut staged_sequence)?;
936            let observation = StrategyObservation::from_annotation(
937                sequence,
938                boundary.observed_through,
939                &annotation,
940            )?;
941            staged_store.push(observation.clone())?;
942            committed.push(observation);
943        }
944
945        for bar in boundary.closed_bars {
946            for analyzer_index in 0..self.analyzers.len() {
947                let context = AnalysisContext {
948                    observed_through: boundary.observed_through,
949                    series: boundary.series,
950                    observations: &staged_store,
951                };
952                let drafts = match self.analyzers[analyzer_index].on_bar(bar, context) {
953                    Ok(drafts) => drafts,
954                    Err(source) => {
955                        self.failed = true;
956                        return Err(AnalysisError::AnalyzerFailure {
957                            analyzer_index,
958                            source: Box::new(source),
959                        });
960                    }
961                };
962                let next_count = match committed.len().checked_add(drafts.len()) {
963                    Some(count) => count,
964                    None => {
965                        self.failed = true;
966                        return Err(AnalysisError::BoundaryOutputOverflow);
967                    }
968                };
969                if let Err(error) = self.ensure_boundary_capacity(next_count) {
970                    self.failed = true;
971                    return Err(error);
972                }
973                for draft in drafts {
974                    let sequence = match take_sequence(&mut staged_sequence) {
975                        Ok(sequence) => sequence,
976                        Err(error) => {
977                            self.failed = true;
978                            return Err(error);
979                        }
980                    };
981                    let observation = match StrategyObservation::from_draft(
982                        sequence,
983                        boundary.observed_through,
984                        draft,
985                    ) {
986                        Ok(observation) => observation,
987                        Err(error) => {
988                            self.failed = true;
989                            return Err(error);
990                        }
991                    };
992                    if let Err(error) = staged_store.push(observation.clone()) {
993                        self.failed = true;
994                        return Err(error);
995                    }
996                    committed.push(observation);
997                }
998            }
999        }
1000
1001        self.observations = staged_store;
1002        self.annotations = staged_annotations;
1003        self.next_sequence = staged_sequence;
1004        self.last_boundary = Some(boundary.observed_through);
1005        Ok(AnalysisBoundaryOutput {
1006            observed_through: boundary.observed_through,
1007            observations: committed,
1008        })
1009    }
1010
1011    fn validate_boundary(&self, boundary: &AnalysisBoundary<'_>) -> Result<(), AnalysisError> {
1012        if let Some(previous) = self.last_boundary
1013            && boundary.observed_through < previous
1014        {
1015            return Err(AnalysisError::BoundaryRegression {
1016                previous,
1017                current: boundary.observed_through,
1018            });
1019        }
1020        for bar in boundary.closed_bars {
1021            if bar.close_time() > boundary.observed_through {
1022                return Err(AnalysisError::BarAfterBoundary {
1023                    series_id: bar.series_id().clone(),
1024                    close_time: bar.close_time(),
1025                    observed_through: boundary.observed_through,
1026                });
1027            }
1028        }
1029        Ok(())
1030    }
1031
1032    fn ensure_boundary_capacity(&self, actual: usize) -> Result<(), AnalysisError> {
1033        let maximum = self.observations.limits.max_per_boundary;
1034        if actual > maximum {
1035            Err(AnalysisError::TooManyBoundaryObservations { actual, maximum })
1036        } else {
1037            Ok(())
1038        }
1039    }
1040}
1041
1042/// Validated configuration for the reference confirmed-pivot analyzer.
1043#[derive(Debug, Clone, PartialEq, Eq)]
1044pub struct PivotConfig {
1045    series_id: SeriesId,
1046    left_bars: usize,
1047    right_bars: usize,
1048    required_bars: usize,
1049}
1050
1051impl PivotConfig {
1052    pub fn new(
1053        series_id: SeriesId,
1054        left_bars: usize,
1055        right_bars: usize,
1056    ) -> Result<Self, AnalysisError> {
1057        validate_nonzero_limit("left_bars", left_bars, MAX_PIVOT_SIDE_BARS)?;
1058        validate_nonzero_limit("right_bars", right_bars, MAX_PIVOT_SIDE_BARS)?;
1059        let required_bars = left_bars
1060            .checked_add(1)
1061            .and_then(|value| value.checked_add(right_bars))
1062            .ok_or(AnalysisError::PivotHistoryOverflow)?;
1063        Ok(Self {
1064            series_id,
1065            left_bars,
1066            right_bars,
1067            required_bars,
1068        })
1069    }
1070
1071    pub fn series_id(&self) -> &SeriesId {
1072        &self.series_id
1073    }
1074
1075    pub fn left_bars(&self) -> usize {
1076        self.left_bars
1077    }
1078
1079    pub fn right_bars(&self) -> usize {
1080        self.right_bars
1081    }
1082
1083    pub fn required_bars(&self) -> usize {
1084        self.required_bars
1085    }
1086}
1087
1088/// Exact delayed-confirmation high and low pivot analyzer.
1089#[derive(Debug, Clone)]
1090pub struct ConfirmedPivotAnalyzer {
1091    config: PivotConfig,
1092}
1093
1094impl ConfirmedPivotAnalyzer {
1095    pub fn new(config: PivotConfig) -> Self {
1096        Self { config }
1097    }
1098
1099    pub fn config(&self) -> &PivotConfig {
1100        &self.config
1101    }
1102}
1103
1104impl HistoricalAnalyzer for ConfirmedPivotAnalyzer {
1105    fn on_bar(
1106        &mut self,
1107        bar: &ClosedBar,
1108        context: AnalysisContext<'_>,
1109    ) -> Result<Vec<StrategyObservationDraft>, AnalysisError> {
1110        context
1111            .series()
1112            .latest_bar(self.config.series_id())
1113            .map_err(AnalysisError::SeriesView)?;
1114        if bar.series_id() != self.config.series_id() {
1115            return Ok(Vec::new());
1116        }
1117        let history = context
1118            .series()
1119            .bars(self.config.series_id(), self.config.required_bars())
1120            .map_err(AnalysisError::SeriesView)?;
1121        if history.len() < self.config.required_bars() {
1122            return Ok(Vec::new());
1123        }
1124        let bars = history.iter().collect::<Vec<_>>();
1125        let candidate = bars[self.config.left_bars];
1126        let left = &bars[..self.config.left_bars];
1127        let right = &bars[self.config.left_bars + 1..];
1128        let is_high = left
1129            .iter()
1130            .chain(right.iter())
1131            .all(|neighbor| candidate.high() > neighbor.high());
1132        let is_low = left
1133            .iter()
1134            .chain(right.iter())
1135            .all(|neighbor| candidate.low() < neighbor.low());
1136        let mut drafts = Vec::with_capacity(usize::from(is_high) + usize::from(is_low));
1137        if is_high {
1138            drafts.push(pivot_draft(
1139                candidate,
1140                SwingKind::High,
1141                context.observed_through(),
1142            )?);
1143        }
1144        if is_low {
1145            drafts.push(pivot_draft(
1146                candidate,
1147                SwingKind::Low,
1148                context.observed_through(),
1149            )?);
1150        }
1151        Ok(drafts)
1152    }
1153}
1154
1155fn pivot_draft(
1156    candidate: &ClosedBar,
1157    kind: SwingKind,
1158    confirmed_at: NaiveDateTime,
1159) -> Result<StrategyObservationDraft, AnalysisError> {
1160    let price = match kind {
1161        SwingKind::High => candidate.high(),
1162        SwingKind::Low => candidate.low(),
1163    };
1164    StrategyObservationDraft::new(
1165        candidate.symbol(),
1166        vec![candidate.series_id().clone()],
1167        StrategyObservationValue::Swing(SwingPoint::new(
1168            kind,
1169            price,
1170            candidate.open_time(),
1171            candidate.close_time(),
1172            confirmed_at,
1173        )?),
1174    )
1175}
1176
1177/// Typed construction, causality, ordering, lookup, and pipeline failures.
1178#[derive(Debug, thiserror::Error)]
1179pub enum AnalysisError {
1180    #[error("zone ID must contain 1 to {MAX_ZONE_ID_BYTES} ASCII identifier bytes")]
1181    InvalidZoneId,
1182    #[error("{field} must be finite and positive")]
1183    InvalidPrice { field: &'static str },
1184    #[error("zone lower {lower} must be strictly below upper {upper}")]
1185    InvalidZoneGeometry { lower: f64, upper: f64 },
1186    #[error("anchor open {open} must be before anchor close {close}")]
1187    InvalidAnchorRange {
1188        open: NaiveDateTime,
1189        close: NaiveDateTime,
1190    },
1191    #[error("confirmation {confirmed_at} cannot precede anchor close {anchor_close}")]
1192    ConfirmationBeforeAnchorClose {
1193        confirmed_at: NaiveDateTime,
1194        anchor_close: NaiveDateTime,
1195    },
1196    #[error("invalid observation symbol '{symbol}'")]
1197    InvalidSymbol { symbol: String },
1198    #[error("observation source series must not contain duplicates")]
1199    DuplicateSourceSeries,
1200    #[error("observation source-series count {actual} exceeds maximum {maximum}")]
1201    TooManySourceSeries { actual: usize, maximum: usize },
1202    #[error("observation valid time {valid_from} is after boundary {observed_through}")]
1203    ObservationNotYetValid {
1204        valid_from: NaiveDateTime,
1205        observed_through: NaiveDateTime,
1206    },
1207    #[error("observation value time {value_time} is after boundary {observed_through}")]
1208    ValueAfterObservation {
1209        value_time: NaiveDateTime,
1210        observed_through: NaiveDateTime,
1211    },
1212    #[error("{field} must be greater than zero")]
1213    ZeroLimit { field: &'static str },
1214    #[error("{field} {actual} exceeds maximum {maximum}")]
1215    LimitTooLarge {
1216        field: &'static str,
1217        actual: usize,
1218        maximum: usize,
1219    },
1220    #[error("analyzer count {actual} exceeds maximum {maximum}")]
1221    TooManyAnalyzers { actual: usize, maximum: usize },
1222    #[error("analysis boundary moved backwards from {previous} to {current}")]
1223    BoundaryRegression {
1224        previous: NaiveDateTime,
1225        current: NaiveDateTime,
1226    },
1227    #[error("bar for '{series_id}' closes at {close_time} after boundary {observed_through}")]
1228    BarAfterBoundary {
1229        series_id: SeriesId,
1230        close_time: NaiveDateTime,
1231        observed_through: NaiveDateTime,
1232    },
1233    #[error("boundary emitted {actual} observations, exceeding maximum {maximum}")]
1234    TooManyBoundaryObservations { actual: usize, maximum: usize },
1235    #[error("boundary observation count overflowed")]
1236    BoundaryOutputOverflow,
1237    #[error("observation sequence overflowed")]
1238    SequenceOverflow,
1239    #[error("observation omitted counter overflowed")]
1240    OmittedCountOverflow,
1241    #[error("pivot history requirement overflowed")]
1242    PivotHistoryOverflow,
1243    #[error(transparent)]
1244    SeriesView(#[from] SeriesViewError),
1245    #[error(transparent)]
1246    Annotation(#[from] super::annotation::AnnotationError),
1247    #[error("analyzer {analyzer_index} failed: {source}")]
1248    AnalyzerFailure {
1249        analyzer_index: usize,
1250        source: Box<AnalysisError>,
1251    },
1252    #[error("analysis pipeline is terminally failed")]
1253    PipelineFailed,
1254    #[error("analyzer failed: {message}")]
1255    Analyzer { message: String },
1256}
1257
1258pub(crate) fn validate_symbol(symbol: &str) -> Result<(), AnalysisError> {
1259    if symbol.is_empty()
1260        || symbol.len() > super::MAX_INSTRUMENT_BYTES
1261        || !symbol.bytes().all(|byte| {
1262            byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-' | b'/' | b':')
1263        })
1264    {
1265        return Err(AnalysisError::InvalidSymbol {
1266            symbol: symbol.to_string(),
1267        });
1268    }
1269    Ok(())
1270}
1271
1272pub(crate) fn validate_source_series(
1273    source_series: &[SeriesId],
1274    maximum: usize,
1275) -> Result<(), AnalysisError> {
1276    if source_series.len() > maximum {
1277        return Err(AnalysisError::TooManySourceSeries {
1278            actual: source_series.len(),
1279            maximum,
1280        });
1281    }
1282    let mut unique = BTreeSet::new();
1283    if source_series.iter().any(|series| !unique.insert(series)) {
1284        return Err(AnalysisError::DuplicateSourceSeries);
1285    }
1286    Ok(())
1287}
1288
1289fn validate_price(value: f64, field: &'static str) -> Result<(), AnalysisError> {
1290    if value.is_finite() && value > 0.0 {
1291        Ok(())
1292    } else {
1293        Err(AnalysisError::InvalidPrice { field })
1294    }
1295}
1296
1297fn validate_nonzero_limit(
1298    field: &'static str,
1299    actual: usize,
1300    maximum: usize,
1301) -> Result<(), AnalysisError> {
1302    if actual == 0 {
1303        return Err(AnalysisError::ZeroLimit { field });
1304    }
1305    if actual > maximum {
1306        return Err(AnalysisError::LimitTooLarge {
1307            field,
1308            actual,
1309            maximum,
1310        });
1311    }
1312    Ok(())
1313}
1314
1315fn valid_identifier(value: &str, maximum: usize) -> bool {
1316    !value.is_empty()
1317        && value.len() <= maximum
1318        && value
1319            .bytes()
1320            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-'))
1321}
1322
1323fn take_sequence(next: &mut u64) -> Result<u64, AnalysisError> {
1324    let sequence = *next;
1325    *next = next.checked_add(1).ok_or(AnalysisError::SequenceOverflow)?;
1326    Ok(sequence)
1327}
1328
1329#[cfg(test)]
1330mod tests {
1331    use super::*;
1332
1333    fn timestamp() -> NaiveDateTime {
1334        chrono::NaiveDate::from_ymd_opt(2026, 1, 2)
1335            .unwrap()
1336            .and_hms_opt(0, 0, 0)
1337            .unwrap()
1338    }
1339
1340    fn observation(sequence: u64) -> StrategyObservation {
1341        StrategyObservation::validated(
1342            sequence,
1343            timestamp(),
1344            timestamp(),
1345            "EURUSD".to_string(),
1346            Vec::new(),
1347            ObservationOrigin::Analyzer,
1348            StrategyObservationValue::Momentum(MomentumState::Advancing),
1349        )
1350        .unwrap()
1351    }
1352
1353    #[test]
1354    fn sequence_overflow_does_not_advance_sequence() {
1355        let mut next = u64::MAX;
1356        assert!(matches!(
1357            take_sequence(&mut next),
1358            Err(AnalysisError::SequenceOverflow)
1359        ));
1360        assert_eq!(next, u64::MAX);
1361    }
1362
1363    #[test]
1364    fn omitted_overflow_does_not_mutate_retained_history() {
1365        let mut store = ObservationStore::new(ObservationStoreLimits::new(1, 1).unwrap());
1366        store.push(observation(0)).unwrap();
1367        store.omitted = u64::MAX;
1368        let snapshot = store
1369            .observations(usize::MAX)
1370            .iter()
1371            .cloned()
1372            .collect::<Vec<_>>();
1373        assert!(matches!(
1374            store.push(observation(1)),
1375            Err(AnalysisError::OmittedCountOverflow)
1376        ));
1377        assert_eq!(
1378            store
1379                .observations(usize::MAX)
1380                .iter()
1381                .cloned()
1382                .collect::<Vec<_>>(),
1383            snapshot
1384        );
1385        assert_eq!(store.omitted(), u64::MAX);
1386    }
1387}