1use 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#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
57#[serde(rename_all = "snake_case")]
58pub enum ZoneSource {
59 CausalAnnotation,
60 DeterministicAnalyzer,
61}
62
63#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(rename_all = "snake_case")]
66pub enum ZoneSide {
67 Support,
68 Resistance,
69}
70
71#[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#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
194#[serde(rename_all = "snake_case")]
195pub enum SwingKind {
196 High,
197 Low,
198}
199
200#[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#[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#[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#[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#[derive(Debug, Clone, PartialEq, Eq)]
421pub enum ObservationOrigin {
422 Analyzer,
423 CausalAnnotation { annotation_id: AnnotationId },
424}
425
426#[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#[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#[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#[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#[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
657pub 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#[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
763pub 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#[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
818pub trait HistoricalAnalyzer: Send {
822 fn on_bar(
823 &mut self,
824 bar: &ClosedBar,
825 context: AnalysisContext<'_>,
826 ) -> Result<Vec<StrategyObservationDraft>, AnalysisError>;
827}
828
829#[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
846pub 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#[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#[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#[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}