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
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#[derive(Debug, Clone, PartialEq, Eq)]
414pub enum ObservationOrigin {
415 Analyzer,
416 CausalAnnotation { annotation_id: AnnotationId },
417}
418
419#[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#[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#[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#[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#[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
650pub 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#[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
756pub 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#[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
811pub trait HistoricalAnalyzer {
813 fn on_bar(
814 &mut self,
815 bar: &ClosedBar,
816 context: AnalysisContext<'_>,
817 ) -> Result<Vec<StrategyObservationDraft>, AnalysisError>;
818}
819
820#[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
837pub 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#[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#[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#[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}