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
375#[derive(Deserialize)]
376#[serde(
377    tag = "kind",
378    content = "value",
379    rename_all = "snake_case",
380    deny_unknown_fields
381)]
382enum StrategyObservationValueDef {
383    Zone(PriceZone),
384    Swing(SwingPoint),
385    Rejection {
386        pattern: RejectionPattern,
387        anchor_open_time: NaiveDateTime,
388        anchor_close_time: NaiveDateTime,
389    },
390    Momentum(MomentumState),
391}
392
393impl<'de> Deserialize<'de> for StrategyObservationValue {
394    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
395    where
396        D: Deserializer<'de>,
397    {
398        match StrategyObservationValueDef::deserialize(deserializer)? {
399            StrategyObservationValueDef::Zone(value) => Ok(Self::Zone(value)),
400            StrategyObservationValueDef::Swing(value) => Ok(Self::Swing(value)),
401            StrategyObservationValueDef::Rejection {
402                pattern,
403                anchor_open_time,
404                anchor_close_time,
405            } => Self::rejection(pattern, anchor_open_time, anchor_close_time)
406                .map_err(serde::de::Error::custom),
407            StrategyObservationValueDef::Momentum(value) => Ok(Self::Momentum(value)),
408        }
409    }
410}
411
412/// Authoritative source category for a committed observation.
413#[derive(Debug, Clone, PartialEq, Eq)]
414pub enum ObservationOrigin {
415    Analyzer,
416    CausalAnnotation { annotation_id: AnnotationId },
417}
418
419/// One immutable committed causal observation.
420#[derive(Debug, Clone, PartialEq)]
421pub struct StrategyObservation {
422    sequence: u64,
423    observed_through: NaiveDateTime,
424    valid_from: NaiveDateTime,
425    symbol: String,
426    source_series: Vec<SeriesId>,
427    origin: ObservationOrigin,
428    value: StrategyObservationValue,
429}
430
431impl StrategyObservation {
432    pub fn sequence(&self) -> u64 {
433        self.sequence
434    }
435
436    pub fn observed_through(&self) -> NaiveDateTime {
437        self.observed_through
438    }
439
440    pub fn valid_from(&self) -> NaiveDateTime {
441        self.valid_from
442    }
443
444    pub fn symbol(&self) -> &str {
445        &self.symbol
446    }
447
448    pub fn source_series(&self) -> &[SeriesId] {
449        &self.source_series
450    }
451
452    pub fn origin(&self) -> &ObservationOrigin {
453        &self.origin
454    }
455
456    pub fn value(&self) -> &StrategyObservationValue {
457        &self.value
458    }
459
460    pub(crate) fn from_annotation(
461        sequence: u64,
462        observed_through: NaiveDateTime,
463        annotation: &StrategyAnnotation,
464    ) -> Result<Self, AnalysisError> {
465        let valid_from = annotation
466            .valid_from()
467            .expect("only causal annotations become observations");
468        Self::validated(
469            sequence,
470            observed_through,
471            valid_from,
472            annotation.symbol().to_string(),
473            annotation.source_series().to_vec(),
474            ObservationOrigin::CausalAnnotation {
475                annotation_id: annotation.annotation_id().clone(),
476            },
477            annotation.value().clone(),
478        )
479    }
480
481    fn from_draft(
482        sequence: u64,
483        observed_through: NaiveDateTime,
484        draft: StrategyObservationDraft,
485    ) -> Result<Self, AnalysisError> {
486        Self::validated(
487            sequence,
488            observed_through,
489            observed_through,
490            draft.symbol,
491            draft.source_series,
492            ObservationOrigin::Analyzer,
493            draft.value,
494        )
495    }
496
497    #[allow(clippy::too_many_arguments)]
498    fn validated(
499        sequence: u64,
500        observed_through: NaiveDateTime,
501        valid_from: NaiveDateTime,
502        symbol: String,
503        source_series: Vec<SeriesId>,
504        origin: ObservationOrigin,
505        value: StrategyObservationValue,
506    ) -> Result<Self, AnalysisError> {
507        validate_symbol(&symbol)?;
508        validate_source_series(&source_series, MAX_OBSERVATION_SOURCE_SERIES)?;
509        if valid_from > observed_through {
510            return Err(AnalysisError::ObservationNotYetValid {
511                valid_from,
512                observed_through,
513            });
514        }
515        value.validate_at(observed_through)?;
516        Ok(Self {
517            sequence,
518            observed_through,
519            valid_from,
520            symbol,
521            source_series,
522            origin,
523            value,
524        })
525    }
526}
527
528/// Metadata-free analyzer output validated and stamped by the pipeline.
529#[derive(Debug, Clone, PartialEq)]
530pub struct StrategyObservationDraft {
531    symbol: String,
532    source_series: Vec<SeriesId>,
533    value: StrategyObservationValue,
534}
535
536impl StrategyObservationDraft {
537    pub fn new(
538        symbol: impl Into<String>,
539        source_series: Vec<SeriesId>,
540        value: StrategyObservationValue,
541    ) -> Result<Self, AnalysisError> {
542        let symbol = symbol.into();
543        validate_symbol(&symbol)?;
544        validate_source_series(&source_series, MAX_OBSERVATION_SOURCE_SERIES)?;
545        Ok(Self {
546            symbol,
547            source_series,
548            value,
549        })
550    }
551
552    pub fn symbol(&self) -> &str {
553        &self.symbol
554    }
555
556    pub fn source_series(&self) -> &[SeriesId] {
557        &self.source_series
558    }
559
560    pub fn value(&self) -> &StrategyObservationValue {
561        &self.value
562    }
563}
564
565/// Caller-visible bounds for strategy-visible observation history.
566#[derive(Debug, Clone, Copy, PartialEq, Eq)]
567pub struct ObservationStoreLimits {
568    max_retained: usize,
569    max_per_boundary: usize,
570}
571
572impl ObservationStoreLimits {
573    pub fn new(max_retained: usize, max_per_boundary: usize) -> Result<Self, AnalysisError> {
574        validate_nonzero_limit("max_retained", max_retained, MAX_RETAINED_OBSERVATIONS)?;
575        validate_nonzero_limit(
576            "max_per_boundary",
577            max_per_boundary,
578            MAX_OBSERVATIONS_PER_BOUNDARY,
579        )?;
580        Ok(Self {
581            max_retained,
582            max_per_boundary,
583        })
584    }
585
586    pub fn max_retained(self) -> usize {
587        self.max_retained
588    }
589
590    pub fn max_per_boundary(self) -> usize {
591        self.max_per_boundary
592    }
593}
594
595impl Default for ObservationStoreLimits {
596    fn default() -> Self {
597        Self {
598            max_retained: 10_000,
599            max_per_boundary: 256,
600        }
601    }
602}
603
604/// Allocation-free suffix view over a possibly wrapped observation deque.
605#[derive(Debug, Clone, Copy)]
606pub struct ObservationWindow<'a> {
607    older: &'a [StrategyObservation],
608    newer: &'a [StrategyObservation],
609}
610
611impl<'a> ObservationWindow<'a> {
612    pub fn len(&self) -> usize {
613        self.older.len() + self.newer.len()
614    }
615
616    pub fn is_empty(&self) -> bool {
617        self.older.is_empty() && self.newer.is_empty()
618    }
619
620    pub fn latest(&self) -> Option<&'a StrategyObservation> {
621        self.newer.last().or_else(|| self.older.last())
622    }
623
624    pub fn iter(&self) -> impl DoubleEndedIterator<Item = &'a StrategyObservation> {
625        self.older.iter().chain(self.newer.iter())
626    }
627}
628
629/// Borrowed deterministic symbol-filtered selection.
630#[derive(Debug, Clone, Copy)]
631pub struct ObservationSelection<'a> {
632    window: ObservationWindow<'a>,
633    symbol: &'a str,
634    count: usize,
635}
636
637impl<'a> ObservationSelection<'a> {
638    pub fn iter(&self) -> impl DoubleEndedIterator<Item = &'a StrategyObservation> {
639        self.window
640            .iter()
641            .rev()
642            .filter(move |observation| observation.symbol() == self.symbol)
643            .take(self.count)
644            .collect::<Vec<_>>()
645            .into_iter()
646            .rev()
647    }
648}
649
650/// Object-safe read-only view over committed causal observations.
651pub trait HistoricalObservationView {
652    fn observations(&self, count: usize) -> ObservationWindow<'_>;
653    fn for_symbol<'a>(&'a self, symbol: &'a str, count: usize) -> ObservationSelection<'a>;
654    fn latest_zone(&self, zone_id: &ZoneId) -> Option<&StrategyObservation>;
655    fn omitted(&self) -> u64;
656}
657
658/// Bounded strategy-visible working history.
659#[derive(Debug, Clone)]
660pub struct ObservationStore {
661    retained: VecDeque<StrategyObservation>,
662    omitted: u64,
663    limits: ObservationStoreLimits,
664}
665
666impl ObservationStore {
667    pub fn new(limits: ObservationStoreLimits) -> Self {
668        Self {
669            retained: VecDeque::with_capacity(limits.max_retained),
670            omitted: 0,
671            limits,
672        }
673    }
674
675    pub fn limits(&self) -> ObservationStoreLimits {
676        self.limits
677    }
678
679    pub fn len(&self) -> usize {
680        self.retained.len()
681    }
682
683    pub fn is_empty(&self) -> bool {
684        self.retained.is_empty()
685    }
686
687    pub fn omitted(&self) -> u64 {
688        self.omitted
689    }
690
691    pub fn observations(&self, count: usize) -> ObservationWindow<'_> {
692        HistoricalObservationView::observations(self, count)
693    }
694
695    pub fn for_symbol<'a>(&'a self, symbol: &'a str, count: usize) -> ObservationSelection<'a> {
696        HistoricalObservationView::for_symbol(self, symbol, count)
697    }
698
699    pub fn latest_zone(&self, zone_id: &ZoneId) -> Option<&StrategyObservation> {
700        HistoricalObservationView::latest_zone(self, zone_id)
701    }
702
703    fn push(&mut self, observation: StrategyObservation) -> Result<(), AnalysisError> {
704        if self.retained.len() == self.limits.max_retained {
705            self.omitted = self
706                .omitted
707                .checked_add(1)
708                .ok_or(AnalysisError::OmittedCountOverflow)?;
709            self.retained.pop_front();
710        }
711        self.retained.push_back(observation);
712        Ok(())
713    }
714}
715
716impl HistoricalObservationView for ObservationStore {
717    fn observations(&self, count: usize) -> ObservationWindow<'_> {
718        let (older, newer) = self.retained.as_slices();
719        let available = count.min(self.retained.len());
720        let skip = self.retained.len() - available;
721        if skip < older.len() {
722            ObservationWindow {
723                older: &older[skip..],
724                newer,
725            }
726        } else {
727            ObservationWindow {
728                older: &older[older.len()..],
729                newer: &newer[skip - older.len()..],
730            }
731        }
732    }
733
734    fn for_symbol<'a>(&'a self, symbol: &'a str, count: usize) -> ObservationSelection<'a> {
735        ObservationSelection {
736            window: self.observations(usize::MAX),
737            symbol,
738            count,
739        }
740    }
741
742    fn latest_zone(&self, zone_id: &ZoneId) -> Option<&StrategyObservation> {
743        self.retained.iter().rev().find(|observation| {
744            observation
745                .value()
746                .zone()
747                .is_some_and(|zone| zone.zone_id() == zone_id)
748        })
749    }
750
751    fn omitted(&self) -> u64 {
752        self.omitted
753    }
754}
755
756/// Complete historical boundary supplied after a closed-bar series commits one timestamp batch.
757pub struct AnalysisBoundary<'a> {
758    observed_through: NaiveDateTime,
759    closed_bars: &'a [ClosedBar],
760    series: &'a dyn HistoricalSeriesView,
761}
762
763impl<'a> AnalysisBoundary<'a> {
764    pub fn new(
765        observed_through: NaiveDateTime,
766        closed_bars: &'a [ClosedBar],
767        series: &'a dyn HistoricalSeriesView,
768    ) -> Self {
769        Self {
770            observed_through,
771            closed_bars,
772            series,
773        }
774    }
775
776    pub fn observed_through(&self) -> NaiveDateTime {
777        self.observed_through
778    }
779
780    pub fn closed_bars(&self) -> &'a [ClosedBar] {
781        self.closed_bars
782    }
783
784    pub fn series(&self) -> &'a dyn HistoricalSeriesView {
785        self.series
786    }
787}
788
789/// Read-only causal state visible during one analyzer callback.
790#[derive(Clone, Copy)]
791pub struct AnalysisContext<'a> {
792    observed_through: NaiveDateTime,
793    series: &'a dyn HistoricalSeriesView,
794    observations: &'a dyn HistoricalObservationView,
795}
796
797impl<'a> AnalysisContext<'a> {
798    pub fn observed_through(self) -> NaiveDateTime {
799        self.observed_through
800    }
801
802    pub fn series(self) -> &'a dyn HistoricalSeriesView {
803        self.series
804    }
805
806    pub fn observations(self) -> &'a dyn HistoricalObservationView {
807        self.observations
808    }
809}
810
811/// Synchronous extension point for causal bar analysis.
812pub trait HistoricalAnalyzer {
813    fn on_bar(
814        &mut self,
815        bar: &ClosedBar,
816        context: AnalysisContext<'_>,
817    ) -> Result<Vec<StrategyObservationDraft>, AnalysisError>;
818}
819
820/// Observations atomically committed by one complete boundary.
821#[derive(Debug, Clone, PartialEq)]
822pub struct AnalysisBoundaryOutput {
823    observed_through: NaiveDateTime,
824    observations: Vec<StrategyObservation>,
825}
826
827impl AnalysisBoundaryOutput {
828    pub fn observed_through(&self) -> NaiveDateTime {
829        self.observed_through
830    }
831
832    pub fn observations(&self) -> &[StrategyObservation] {
833        &self.observations
834    }
835}
836
837/// Complete-boundary analyzer pipeline with bounded committed history.
838pub struct AnalysisPipeline {
839    analyzers: Vec<Box<dyn HistoricalAnalyzer>>,
840    observations: ObservationStore,
841    annotations: AnnotationTimeline,
842    next_sequence: u64,
843    last_boundary: Option<NaiveDateTime>,
844    failed: bool,
845}
846
847impl AnalysisPipeline {
848    pub fn new(
849        analyzers: Vec<Box<dyn HistoricalAnalyzer>>,
850        observation_limits: ObservationStoreLimits,
851        annotation_limits: AnnotationLimits,
852    ) -> Result<Self, AnalysisError> {
853        if analyzers.len() > MAX_ANALYZERS {
854            return Err(AnalysisError::TooManyAnalyzers {
855                actual: analyzers.len(),
856                maximum: MAX_ANALYZERS,
857            });
858        }
859        Ok(Self {
860            analyzers,
861            observations: ObservationStore::new(observation_limits),
862            annotations: AnnotationTimeline::new(annotation_limits),
863            next_sequence: 0,
864            last_boundary: None,
865            failed: false,
866        })
867    }
868
869    pub fn add_annotation(&mut self, annotation: StrategyAnnotation) -> Result<(), AnalysisError> {
870        if self.failed {
871            return Err(AnalysisError::PipelineFailed);
872        }
873        self.annotations
874            .add(annotation, self.last_boundary)
875            .map_err(AnalysisError::Annotation)
876    }
877
878    pub fn observations(&self) -> &ObservationStore {
879        &self.observations
880    }
881
882    pub fn annotations(&self) -> &AnnotationTimeline {
883        &self.annotations
884    }
885
886    pub(crate) fn into_research_annotations(self) -> Vec<StrategyAnnotation> {
887        self.annotations.into_research_only()
888    }
889
890    pub fn is_failed(&self) -> bool {
891        self.failed
892    }
893
894    pub fn on_boundary(
895        &mut self,
896        boundary: AnalysisBoundary<'_>,
897    ) -> Result<AnalysisBoundaryOutput, AnalysisError> {
898        if self.failed {
899            return Err(AnalysisError::PipelineFailed);
900        }
901        self.validate_boundary(&boundary)?;
902
903        let activates_annotation = self
904            .annotations
905            .pending_causal()
906            .first()
907            .and_then(|annotation| annotation.valid_from())
908            .is_some_and(|valid_from| valid_from <= boundary.observed_through);
909        if boundary.closed_bars.is_empty() && !activates_annotation {
910            self.last_boundary = Some(boundary.observed_through);
911            return Ok(AnalysisBoundaryOutput {
912                observed_through: boundary.observed_through,
913                observations: Vec::new(),
914            });
915        }
916
917        let mut staged_store = self.observations.clone();
918        let mut staged_annotations = self.annotations.clone();
919        let mut staged_sequence = self.next_sequence;
920        let mut committed = Vec::new();
921        let eligible = staged_annotations
922            .activate(boundary.observed_through)
923            .map_err(AnalysisError::Annotation)?;
924        self.ensure_boundary_capacity(eligible.len())?;
925        for annotation in eligible {
926            let sequence = take_sequence(&mut staged_sequence)?;
927            let observation = StrategyObservation::from_annotation(
928                sequence,
929                boundary.observed_through,
930                &annotation,
931            )?;
932            staged_store.push(observation.clone())?;
933            committed.push(observation);
934        }
935
936        for bar in boundary.closed_bars {
937            for analyzer_index in 0..self.analyzers.len() {
938                let context = AnalysisContext {
939                    observed_through: boundary.observed_through,
940                    series: boundary.series,
941                    observations: &staged_store,
942                };
943                let drafts = match self.analyzers[analyzer_index].on_bar(bar, context) {
944                    Ok(drafts) => drafts,
945                    Err(source) => {
946                        self.failed = true;
947                        return Err(AnalysisError::AnalyzerFailure {
948                            analyzer_index,
949                            source: Box::new(source),
950                        });
951                    }
952                };
953                let next_count = match committed.len().checked_add(drafts.len()) {
954                    Some(count) => count,
955                    None => {
956                        self.failed = true;
957                        return Err(AnalysisError::BoundaryOutputOverflow);
958                    }
959                };
960                if let Err(error) = self.ensure_boundary_capacity(next_count) {
961                    self.failed = true;
962                    return Err(error);
963                }
964                for draft in drafts {
965                    let sequence = match take_sequence(&mut staged_sequence) {
966                        Ok(sequence) => sequence,
967                        Err(error) => {
968                            self.failed = true;
969                            return Err(error);
970                        }
971                    };
972                    let observation = match StrategyObservation::from_draft(
973                        sequence,
974                        boundary.observed_through,
975                        draft,
976                    ) {
977                        Ok(observation) => observation,
978                        Err(error) => {
979                            self.failed = true;
980                            return Err(error);
981                        }
982                    };
983                    if let Err(error) = staged_store.push(observation.clone()) {
984                        self.failed = true;
985                        return Err(error);
986                    }
987                    committed.push(observation);
988                }
989            }
990        }
991
992        self.observations = staged_store;
993        self.annotations = staged_annotations;
994        self.next_sequence = staged_sequence;
995        self.last_boundary = Some(boundary.observed_through);
996        Ok(AnalysisBoundaryOutput {
997            observed_through: boundary.observed_through,
998            observations: committed,
999        })
1000    }
1001
1002    fn validate_boundary(&self, boundary: &AnalysisBoundary<'_>) -> Result<(), AnalysisError> {
1003        if let Some(previous) = self.last_boundary
1004            && boundary.observed_through < previous
1005        {
1006            return Err(AnalysisError::BoundaryRegression {
1007                previous,
1008                current: boundary.observed_through,
1009            });
1010        }
1011        for bar in boundary.closed_bars {
1012            if bar.close_time() > boundary.observed_through {
1013                return Err(AnalysisError::BarAfterBoundary {
1014                    series_id: bar.series_id().clone(),
1015                    close_time: bar.close_time(),
1016                    observed_through: boundary.observed_through,
1017                });
1018            }
1019        }
1020        Ok(())
1021    }
1022
1023    fn ensure_boundary_capacity(&self, actual: usize) -> Result<(), AnalysisError> {
1024        let maximum = self.observations.limits.max_per_boundary;
1025        if actual > maximum {
1026            Err(AnalysisError::TooManyBoundaryObservations { actual, maximum })
1027        } else {
1028            Ok(())
1029        }
1030    }
1031}
1032
1033/// Validated configuration for the reference confirmed-pivot analyzer.
1034#[derive(Debug, Clone, PartialEq, Eq)]
1035pub struct PivotConfig {
1036    series_id: SeriesId,
1037    left_bars: usize,
1038    right_bars: usize,
1039    required_bars: usize,
1040}
1041
1042impl PivotConfig {
1043    pub fn new(
1044        series_id: SeriesId,
1045        left_bars: usize,
1046        right_bars: usize,
1047    ) -> Result<Self, AnalysisError> {
1048        validate_nonzero_limit("left_bars", left_bars, MAX_PIVOT_SIDE_BARS)?;
1049        validate_nonzero_limit("right_bars", right_bars, MAX_PIVOT_SIDE_BARS)?;
1050        let required_bars = left_bars
1051            .checked_add(1)
1052            .and_then(|value| value.checked_add(right_bars))
1053            .ok_or(AnalysisError::PivotHistoryOverflow)?;
1054        Ok(Self {
1055            series_id,
1056            left_bars,
1057            right_bars,
1058            required_bars,
1059        })
1060    }
1061
1062    pub fn series_id(&self) -> &SeriesId {
1063        &self.series_id
1064    }
1065
1066    pub fn left_bars(&self) -> usize {
1067        self.left_bars
1068    }
1069
1070    pub fn right_bars(&self) -> usize {
1071        self.right_bars
1072    }
1073
1074    pub fn required_bars(&self) -> usize {
1075        self.required_bars
1076    }
1077}
1078
1079/// Exact delayed-confirmation high and low pivot analyzer.
1080#[derive(Debug, Clone)]
1081pub struct ConfirmedPivotAnalyzer {
1082    config: PivotConfig,
1083}
1084
1085impl ConfirmedPivotAnalyzer {
1086    pub fn new(config: PivotConfig) -> Self {
1087        Self { config }
1088    }
1089
1090    pub fn config(&self) -> &PivotConfig {
1091        &self.config
1092    }
1093}
1094
1095impl HistoricalAnalyzer for ConfirmedPivotAnalyzer {
1096    fn on_bar(
1097        &mut self,
1098        bar: &ClosedBar,
1099        context: AnalysisContext<'_>,
1100    ) -> Result<Vec<StrategyObservationDraft>, AnalysisError> {
1101        context
1102            .series()
1103            .latest_bar(self.config.series_id())
1104            .map_err(AnalysisError::SeriesView)?;
1105        if bar.series_id() != self.config.series_id() {
1106            return Ok(Vec::new());
1107        }
1108        let history = context
1109            .series()
1110            .bars(self.config.series_id(), self.config.required_bars())
1111            .map_err(AnalysisError::SeriesView)?;
1112        if history.len() < self.config.required_bars() {
1113            return Ok(Vec::new());
1114        }
1115        let bars = history.iter().collect::<Vec<_>>();
1116        let candidate = bars[self.config.left_bars];
1117        let left = &bars[..self.config.left_bars];
1118        let right = &bars[self.config.left_bars + 1..];
1119        let is_high = left
1120            .iter()
1121            .chain(right.iter())
1122            .all(|neighbor| candidate.high() > neighbor.high());
1123        let is_low = left
1124            .iter()
1125            .chain(right.iter())
1126            .all(|neighbor| candidate.low() < neighbor.low());
1127        let mut drafts = Vec::with_capacity(usize::from(is_high) + usize::from(is_low));
1128        if is_high {
1129            drafts.push(pivot_draft(
1130                candidate,
1131                SwingKind::High,
1132                context.observed_through(),
1133            )?);
1134        }
1135        if is_low {
1136            drafts.push(pivot_draft(
1137                candidate,
1138                SwingKind::Low,
1139                context.observed_through(),
1140            )?);
1141        }
1142        Ok(drafts)
1143    }
1144}
1145
1146fn pivot_draft(
1147    candidate: &ClosedBar,
1148    kind: SwingKind,
1149    confirmed_at: NaiveDateTime,
1150) -> Result<StrategyObservationDraft, AnalysisError> {
1151    let price = match kind {
1152        SwingKind::High => candidate.high(),
1153        SwingKind::Low => candidate.low(),
1154    };
1155    StrategyObservationDraft::new(
1156        candidate.symbol(),
1157        vec![candidate.series_id().clone()],
1158        StrategyObservationValue::Swing(SwingPoint::new(
1159            kind,
1160            price,
1161            candidate.open_time(),
1162            candidate.close_time(),
1163            confirmed_at,
1164        )?),
1165    )
1166}
1167
1168/// Typed construction, causality, ordering, lookup, and pipeline failures.
1169#[derive(Debug, thiserror::Error)]
1170pub enum AnalysisError {
1171    #[error("zone ID must contain 1 to {MAX_ZONE_ID_BYTES} ASCII identifier bytes")]
1172    InvalidZoneId,
1173    #[error("{field} must be finite and positive")]
1174    InvalidPrice { field: &'static str },
1175    #[error("zone lower {lower} must be strictly below upper {upper}")]
1176    InvalidZoneGeometry { lower: f64, upper: f64 },
1177    #[error("anchor open {open} must be before anchor close {close}")]
1178    InvalidAnchorRange {
1179        open: NaiveDateTime,
1180        close: NaiveDateTime,
1181    },
1182    #[error("confirmation {confirmed_at} cannot precede anchor close {anchor_close}")]
1183    ConfirmationBeforeAnchorClose {
1184        confirmed_at: NaiveDateTime,
1185        anchor_close: NaiveDateTime,
1186    },
1187    #[error("invalid observation symbol '{symbol}'")]
1188    InvalidSymbol { symbol: String },
1189    #[error("observation source series must not contain duplicates")]
1190    DuplicateSourceSeries,
1191    #[error("observation source-series count {actual} exceeds maximum {maximum}")]
1192    TooManySourceSeries { actual: usize, maximum: usize },
1193    #[error("observation valid time {valid_from} is after boundary {observed_through}")]
1194    ObservationNotYetValid {
1195        valid_from: NaiveDateTime,
1196        observed_through: NaiveDateTime,
1197    },
1198    #[error("observation value time {value_time} is after boundary {observed_through}")]
1199    ValueAfterObservation {
1200        value_time: NaiveDateTime,
1201        observed_through: NaiveDateTime,
1202    },
1203    #[error("{field} must be greater than zero")]
1204    ZeroLimit { field: &'static str },
1205    #[error("{field} {actual} exceeds maximum {maximum}")]
1206    LimitTooLarge {
1207        field: &'static str,
1208        actual: usize,
1209        maximum: usize,
1210    },
1211    #[error("analyzer count {actual} exceeds maximum {maximum}")]
1212    TooManyAnalyzers { actual: usize, maximum: usize },
1213    #[error("analysis boundary moved backwards from {previous} to {current}")]
1214    BoundaryRegression {
1215        previous: NaiveDateTime,
1216        current: NaiveDateTime,
1217    },
1218    #[error("bar for '{series_id}' closes at {close_time} after boundary {observed_through}")]
1219    BarAfterBoundary {
1220        series_id: SeriesId,
1221        close_time: NaiveDateTime,
1222        observed_through: NaiveDateTime,
1223    },
1224    #[error("boundary emitted {actual} observations, exceeding maximum {maximum}")]
1225    TooManyBoundaryObservations { actual: usize, maximum: usize },
1226    #[error("boundary observation count overflowed")]
1227    BoundaryOutputOverflow,
1228    #[error("observation sequence overflowed")]
1229    SequenceOverflow,
1230    #[error("observation omitted counter overflowed")]
1231    OmittedCountOverflow,
1232    #[error("pivot history requirement overflowed")]
1233    PivotHistoryOverflow,
1234    #[error(transparent)]
1235    SeriesView(#[from] SeriesViewError),
1236    #[error(transparent)]
1237    Annotation(#[from] super::annotation::AnnotationError),
1238    #[error("analyzer {analyzer_index} failed: {source}")]
1239    AnalyzerFailure {
1240        analyzer_index: usize,
1241        source: Box<AnalysisError>,
1242    },
1243    #[error("analysis pipeline is terminally failed")]
1244    PipelineFailed,
1245    #[error("analyzer failed: {message}")]
1246    Analyzer { message: String },
1247}
1248
1249pub(crate) fn validate_symbol(symbol: &str) -> Result<(), AnalysisError> {
1250    if symbol.is_empty()
1251        || symbol.len() > super::MAX_INSTRUMENT_BYTES
1252        || !symbol.bytes().all(|byte| {
1253            byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-' | b'/' | b':')
1254        })
1255    {
1256        return Err(AnalysisError::InvalidSymbol {
1257            symbol: symbol.to_string(),
1258        });
1259    }
1260    Ok(())
1261}
1262
1263pub(crate) fn validate_source_series(
1264    source_series: &[SeriesId],
1265    maximum: usize,
1266) -> Result<(), AnalysisError> {
1267    if source_series.len() > maximum {
1268        return Err(AnalysisError::TooManySourceSeries {
1269            actual: source_series.len(),
1270            maximum,
1271        });
1272    }
1273    let mut unique = BTreeSet::new();
1274    if source_series.iter().any(|series| !unique.insert(series)) {
1275        return Err(AnalysisError::DuplicateSourceSeries);
1276    }
1277    Ok(())
1278}
1279
1280fn validate_price(value: f64, field: &'static str) -> Result<(), AnalysisError> {
1281    if value.is_finite() && value > 0.0 {
1282        Ok(())
1283    } else {
1284        Err(AnalysisError::InvalidPrice { field })
1285    }
1286}
1287
1288fn validate_nonzero_limit(
1289    field: &'static str,
1290    actual: usize,
1291    maximum: usize,
1292) -> Result<(), AnalysisError> {
1293    if actual == 0 {
1294        return Err(AnalysisError::ZeroLimit { field });
1295    }
1296    if actual > maximum {
1297        return Err(AnalysisError::LimitTooLarge {
1298            field,
1299            actual,
1300            maximum,
1301        });
1302    }
1303    Ok(())
1304}
1305
1306fn valid_identifier(value: &str, maximum: usize) -> bool {
1307    !value.is_empty()
1308        && value.len() <= maximum
1309        && value
1310            .bytes()
1311            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-'))
1312}
1313
1314fn take_sequence(next: &mut u64) -> Result<u64, AnalysisError> {
1315    let sequence = *next;
1316    *next = next.checked_add(1).ok_or(AnalysisError::SequenceOverflow)?;
1317    Ok(sequence)
1318}
1319
1320#[cfg(test)]
1321mod tests {
1322    use super::*;
1323
1324    fn timestamp() -> NaiveDateTime {
1325        chrono::NaiveDate::from_ymd_opt(2026, 1, 2)
1326            .unwrap()
1327            .and_hms_opt(0, 0, 0)
1328            .unwrap()
1329    }
1330
1331    fn observation(sequence: u64) -> StrategyObservation {
1332        StrategyObservation::validated(
1333            sequence,
1334            timestamp(),
1335            timestamp(),
1336            "EURUSD".to_string(),
1337            Vec::new(),
1338            ObservationOrigin::Analyzer,
1339            StrategyObservationValue::Momentum(MomentumState::Advancing),
1340        )
1341        .unwrap()
1342    }
1343
1344    #[test]
1345    fn sequence_overflow_does_not_advance_sequence() {
1346        let mut next = u64::MAX;
1347        assert!(matches!(
1348            take_sequence(&mut next),
1349            Err(AnalysisError::SequenceOverflow)
1350        ));
1351        assert_eq!(next, u64::MAX);
1352    }
1353
1354    #[test]
1355    fn omitted_overflow_does_not_mutate_retained_history() {
1356        let mut store = ObservationStore::new(ObservationStoreLimits::new(1, 1).unwrap());
1357        store.push(observation(0)).unwrap();
1358        store.omitted = u64::MAX;
1359        let snapshot = store
1360            .observations(usize::MAX)
1361            .iter()
1362            .cloned()
1363            .collect::<Vec<_>>();
1364        assert!(matches!(
1365            store.push(observation(1)),
1366            Err(AnalysisError::OmittedCountOverflow)
1367        ));
1368        assert_eq!(
1369            store
1370                .observations(usize::MAX)
1371                .iter()
1372                .cloned()
1373                .collect::<Vec<_>>(),
1374            snapshot
1375        );
1376        assert_eq!(store.omitted(), u64::MAX);
1377    }
1378}