1use std::collections::{BTreeMap, BTreeSet, VecDeque};
4
5use chrono::{DateTime, NaiveDateTime};
6
7use crate::data_feed::{MarketEvent, TimestampBatch};
8
9use super::{PriceBasis, SeriesId, SeriesRequirement, StrategyRequirements, WarmupRequirement};
10
11pub const MAX_RETAINED_BARS: usize = 1_000_000;
12
13#[derive(Debug, Clone, Copy, PartialEq, Eq)]
15pub enum MissingIntervalPolicy {
16 Skip,
17 Reject,
18}
19
20#[derive(Debug, Clone, PartialEq, Eq)]
22pub struct BarSeriesSpec {
23 requirement: SeriesRequirement,
24 retained_bars: usize,
25 alignment_offset_seconds: i64,
26 missing_interval: MissingIntervalPolicy,
27}
28
29impl BarSeriesSpec {
30 pub fn new(
31 requirement: SeriesRequirement,
32 retained_bars: usize,
33 alignment_offset_seconds: i32,
34 missing_interval: MissingIntervalPolicy,
35 ) -> Result<Self, SeriesError> {
36 if retained_bars == 0 {
37 return Err(SeriesError::ZeroRetention {
38 series_id: requirement.id().clone(),
39 });
40 }
41 if retained_bars > MAX_RETAINED_BARS {
42 return Err(SeriesError::RetentionTooLarge {
43 series_id: requirement.id().clone(),
44 retained_bars,
45 maximum: MAX_RETAINED_BARS,
46 });
47 }
48 let required_bars = requirement.warmup().required_bars();
49 if retained_bars < required_bars {
50 return Err(SeriesError::RetentionBelowWarmup {
51 series_id: requirement.id().clone(),
52 retained_bars,
53 required_bars,
54 });
55 }
56 let duration = i64::try_from(requirement.timeframe().duration_seconds())
57 .expect("fixed timeframe duration always fits i64");
58 let alignment_offset_seconds = i64::from(alignment_offset_seconds).rem_euclid(duration);
59 Ok(Self {
60 requirement,
61 retained_bars,
62 alignment_offset_seconds,
63 missing_interval,
64 })
65 }
66
67 pub fn requirement(&self) -> &SeriesRequirement {
68 &self.requirement
69 }
70
71 pub fn retained_bars(&self) -> usize {
72 self.retained_bars
73 }
74
75 pub fn alignment_offset_seconds(&self) -> i64 {
76 self.alignment_offset_seconds
77 }
78
79 pub fn missing_interval(&self) -> MissingIntervalPolicy {
80 self.missing_interval
81 }
82}
83
84#[derive(Debug, Clone, PartialEq)]
86pub struct ClosedBar {
87 series_id: SeriesId,
88 symbol: String,
89 open_time: NaiveDateTime,
90 close_time: NaiveDateTime,
91 open: f64,
92 high: f64,
93 low: f64,
94 close: f64,
95 tick_count: u64,
96}
97
98impl ClosedBar {
99 pub fn series_id(&self) -> &SeriesId {
100 &self.series_id
101 }
102
103 pub fn symbol(&self) -> &str {
104 &self.symbol
105 }
106
107 pub fn open_time(&self) -> NaiveDateTime {
108 self.open_time
109 }
110
111 pub fn close_time(&self) -> NaiveDateTime {
112 self.close_time
113 }
114
115 pub fn open(&self) -> f64 {
116 self.open
117 }
118
119 pub fn high(&self) -> f64 {
120 self.high
121 }
122
123 pub fn low(&self) -> f64 {
124 self.low
125 }
126
127 pub fn close(&self) -> f64 {
128 self.close
129 }
130
131 pub fn tick_count(&self) -> u64 {
132 self.tick_count
133 }
134}
135
136#[derive(Debug, Clone, Copy)]
138pub struct BarWindow<'a> {
139 older: &'a [ClosedBar],
140 newer: &'a [ClosedBar],
141}
142
143impl<'a> BarWindow<'a> {
144 pub fn len(&self) -> usize {
145 self.older.len() + self.newer.len()
146 }
147
148 pub fn is_empty(&self) -> bool {
149 self.older.is_empty() && self.newer.is_empty()
150 }
151
152 pub fn latest(&self) -> Option<&'a ClosedBar> {
153 self.newer.last().or_else(|| self.older.last())
154 }
155
156 pub fn iter(&self) -> impl Iterator<Item = &'a ClosedBar> {
157 self.older.iter().chain(self.newer.iter())
158 }
159}
160
161#[derive(Debug, Clone, Copy, PartialEq, Eq)]
163pub struct SeriesWarmupState {
164 required: WarmupRequirement,
165 available_bars: usize,
166}
167
168impl SeriesWarmupState {
169 pub fn required(self) -> WarmupRequirement {
170 self.required
171 }
172
173 pub fn available_bars(self) -> usize {
174 self.available_bars
175 }
176
177 pub fn is_ready(self) -> bool {
178 self.available_bars >= self.required.required_bars()
179 }
180}
181
182pub trait HistoricalSeriesView {
184 fn latest_bar(&self, id: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError>;
185 fn bars(&self, id: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError>;
186 fn warmup(&self, id: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError>;
187}
188
189#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
191pub enum SeriesError {
192 #[error("series '{series_id}' retained-bar capacity must be greater than zero")]
193 ZeroRetention { series_id: SeriesId },
194 #[error("series '{series_id}' retained-bar capacity {retained_bars} exceeds maximum {maximum}")]
195 RetentionTooLarge {
196 series_id: SeriesId,
197 retained_bars: usize,
198 maximum: usize,
199 },
200 #[error(
201 "series '{series_id}' retained-bar capacity {retained_bars} is below warmup {required_bars}"
202 )]
203 RetentionBelowWarmup {
204 series_id: SeriesId,
205 retained_bars: usize,
206 required_bars: usize,
207 },
208 #[error("series ID '{series_id}' is configured more than once")]
209 DuplicateSeriesId { series_id: SeriesId },
210 #[error("batch at {batch_ts} contains an event at {event_ts}")]
211 BatchTimestampMismatch {
212 batch_ts: NaiveDateTime,
213 event_ts: NaiveDateTime,
214 },
215 #[error("duplicate event ordering metadata ({series_rank}, {row_sequence}) at {timestamp}")]
216 DuplicateOrderingMetadata {
217 timestamp: NaiveDateTime,
218 series_rank: u32,
219 row_sequence: u64,
220 },
221 #[error("primary tick for '{symbol}' moved backwards from {previous} to {current}")]
222 TimestampRegression {
223 symbol: String,
224 previous: NaiveDateTime,
225 current: NaiveDateTime,
226 },
227 #[error("series '{series_id}' cannot represent a bucket boundary for {timestamp}")]
228 BoundaryOverflow {
229 series_id: SeriesId,
230 timestamp: NaiveDateTime,
231 },
232 #[error("series '{series_id}' has one or more empty intervals before {next_open}")]
233 MissingInterval {
234 series_id: SeriesId,
235 previous_close: NaiveDateTime,
236 next_open: NaiveDateTime,
237 },
238 #[error("series '{series_id}' tick count overflowed")]
239 TickCountOverflow { series_id: SeriesId },
240 #[error("series '{series_id}' completed-bar count overflowed")]
241 CompletedBarCountOverflow { series_id: SeriesId },
242}
243
244#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
246pub enum SeriesViewError {
247 #[error("unknown historical series '{series_id}'")]
248 UnknownSeries { series_id: SeriesId },
249}
250
251#[derive(Debug, Clone)]
252struct OpenBar {
253 open_time: NaiveDateTime,
254 close_time: NaiveDateTime,
255 open: f64,
256 high: f64,
257 low: f64,
258 close: f64,
259 tick_count: u64,
260}
261
262impl OpenBar {
263 fn new(open_time: NaiveDateTime, close_time: NaiveDateTime, price: f64) -> Self {
264 Self {
265 open_time,
266 close_time,
267 open: price,
268 high: price,
269 low: price,
270 close: price,
271 tick_count: 1,
272 }
273 }
274
275 fn update(&mut self, series_id: &SeriesId, price: f64) -> Result<(), SeriesError> {
276 self.high = self.high.max(price);
277 self.low = self.low.min(price);
278 self.close = price;
279 self.tick_count =
280 self.tick_count
281 .checked_add(1)
282 .ok_or_else(|| SeriesError::TickCountOverflow {
283 series_id: series_id.clone(),
284 })?;
285 Ok(())
286 }
287}
288
289#[derive(Debug, Clone)]
290struct SeriesState {
291 spec: BarSeriesSpec,
292 open: Option<OpenBar>,
293 closed: VecDeque<ClosedBar>,
294 completed_bars: usize,
295}
296
297impl SeriesState {
298 fn new(spec: BarSeriesSpec) -> Self {
299 Self {
300 closed: VecDeque::with_capacity(spec.retained_bars),
301 spec,
302 open: None,
303 completed_bars: 0,
304 }
305 }
306
307 fn duration_seconds(&self) -> i64 {
308 i64::try_from(self.spec.requirement.timeframe().duration_seconds())
309 .expect("fixed timeframe duration always fits i64")
310 }
311}
312
313#[derive(Debug)]
314struct BatchSeriesState {
315 open: Option<OpenBar>,
316 completed_bars: usize,
317 emitted: Vec<ClosedBar>,
318}
319
320impl BatchSeriesState {
321 fn from_committed(state: &SeriesState) -> Self {
322 Self {
323 open: state.open.clone(),
324 completed_bars: state.completed_bars,
325 emitted: Vec::new(),
326 }
327 }
328
329 fn apply_tick(
330 &mut self,
331 spec: &BarSeriesSpec,
332 timestamp: NaiveDateTime,
333 bid: f64,
334 ask: f64,
335 ) -> Result<(), SeriesError> {
336 let price = match spec.requirement.price_basis() {
337 PriceBasis::Bid => bid,
338 PriceBasis::Ask => ask,
339 PriceBasis::Mid => bid + (ask - bid) / 2.0,
340 };
341 let duration = i64::try_from(spec.requirement.timeframe().duration_seconds())
342 .expect("fixed timeframe duration always fits i64");
343 let (open_time, close_time) = bucket_bounds(
344 spec.requirement.id(),
345 timestamp,
346 duration,
347 spec.alignment_offset_seconds,
348 )?;
349
350 let Some(current) = self.open.as_mut() else {
351 self.open = Some(OpenBar::new(open_time, close_time, price));
352 return Ok(());
353 };
354 if open_time == current.open_time {
355 current.update(spec.requirement.id(), price)?;
356 return Ok(());
357 }
358
359 if spec.missing_interval == MissingIntervalPolicy::Reject && open_time > current.close_time
360 {
361 return Err(SeriesError::MissingInterval {
362 series_id: spec.requirement.id().clone(),
363 previous_close: current.close_time,
364 next_open: open_time,
365 });
366 }
367
368 let completed = self
369 .open
370 .take()
371 .expect("open bar exists after transition validation");
372 self.completed_bars = self.completed_bars.checked_add(1).ok_or_else(|| {
373 SeriesError::CompletedBarCountOverflow {
374 series_id: spec.requirement.id().clone(),
375 }
376 })?;
377 self.emitted.push(ClosedBar {
378 series_id: spec.requirement.id().clone(),
379 symbol: spec.requirement.symbol().to_string(),
380 open_time: completed.open_time,
381 close_time: completed.close_time,
382 open: completed.open,
383 high: completed.high,
384 low: completed.low,
385 close: completed.close,
386 tick_count: completed.tick_count,
387 });
388 self.open = Some(OpenBar::new(open_time, close_time, price));
389 Ok(())
390 }
391}
392
393#[derive(Debug)]
395pub struct MultiTimeframeSeries {
396 series: BTreeMap<SeriesId, SeriesState>,
397 last_source_ts: BTreeMap<String, NaiveDateTime>,
398}
399
400impl MultiTimeframeSeries {
401 pub fn new(specs: Vec<BarSeriesSpec>) -> Result<Self, SeriesError> {
402 let mut series = BTreeMap::new();
403 for spec in specs {
404 let id = spec.requirement.id().clone();
405 if series.insert(id.clone(), SeriesState::new(spec)).is_some() {
406 return Err(SeriesError::DuplicateSeriesId { series_id: id });
407 }
408 }
409 Ok(Self {
410 series,
411 last_source_ts: BTreeMap::new(),
412 })
413 }
414
415 pub fn on_batch(&mut self, batch: &TimestampBatch) -> Result<Vec<ClosedBar>, SeriesError> {
416 validate_batch(batch)?;
417 let mut staged_series = self
418 .series
419 .iter()
420 .map(|(id, state)| (id.clone(), BatchSeriesState::from_committed(state)))
421 .collect::<BTreeMap<_, _>>();
422 let mut staged_source_ts = self.last_source_ts.clone();
423 self.preflight_batch(batch, &mut staged_series, &mut staged_source_ts)?;
424
425 let mut emitted = Vec::new();
426 for (id, mut staged) in staged_series {
427 let state = self
428 .series
429 .get_mut(&id)
430 .expect("staged series originates from committed state");
431 state.open = staged.open;
432 state.completed_bars = staged.completed_bars;
433 for closed in staged.emitted.drain(..) {
434 if state.closed.len() == state.spec.retained_bars {
435 state.closed.pop_front();
436 }
437 state.closed.push_back(closed.clone());
438 emitted.push(closed);
439 }
440 }
441 self.last_source_ts = staged_source_ts;
442 emitted.sort_by(|left, right| {
443 left.close_time
444 .cmp(&right.close_time)
445 .then_with(|| {
446 self.series[left.series_id()]
447 .duration_seconds()
448 .cmp(&self.series[right.series_id()].duration_seconds())
449 })
450 .then_with(|| left.series_id.cmp(&right.series_id))
451 });
452 Ok(emitted)
453 }
454
455 pub fn latest_bar(&self, series: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError> {
456 HistoricalSeriesView::latest_bar(self, series)
457 }
458
459 pub fn bars(&self, series: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError> {
460 HistoricalSeriesView::bars(self, series, count)
461 }
462
463 pub fn warmup(&self, series: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError> {
464 HistoricalSeriesView::warmup(self, series)
465 }
466
467 pub fn warmup_complete(
468 &self,
469 requirements: &StrategyRequirements,
470 ) -> Result<bool, SeriesViewError> {
471 for requirement in requirements.series() {
472 self.state(requirement.id())?;
473 }
474 Ok(requirements
475 .warmup_complete(|id| self.series.get(id).map_or(0, |state| state.completed_bars)))
476 }
477
478 fn preflight_batch(
479 &self,
480 batch: &TimestampBatch,
481 staged_series: &mut BTreeMap<SeriesId, BatchSeriesState>,
482 staged_source_ts: &mut BTreeMap<String, NaiveDateTime>,
483 ) -> Result<(), SeriesError> {
484 let mut ordered = batch.events.iter().collect::<Vec<_>>();
485 ordered.sort_by_key(|event| (event.metadata.series_rank, event.metadata.row_sequence));
486
487 for feed_event in ordered {
488 if !feed_event.metadata.roles.primary {
489 continue;
490 }
491 let MarketEvent::Tick {
492 symbol,
493 ts,
494 bid,
495 ask,
496 } = &feed_event.event
497 else {
498 continue;
499 };
500 if let Some(previous) = staged_source_ts.get(symbol)
501 && *ts < *previous
502 {
503 return Err(SeriesError::TimestampRegression {
504 symbol: symbol.clone(),
505 previous: *previous,
506 current: *ts,
507 });
508 }
509 staged_source_ts.insert(symbol.clone(), *ts);
510
511 if feed_event.event.to_valid_quote().is_none() {
512 continue;
513 }
514 for (id, state) in self
515 .series
516 .iter()
517 .filter(|(_, state)| state.spec.requirement.symbol() == symbol)
518 {
519 staged_series
520 .get_mut(id)
521 .expect("staged series originates from committed state")
522 .apply_tick(&state.spec, *ts, *bid, *ask)?;
523 }
524 }
525 Ok(())
526 }
527
528 fn state(&self, id: &SeriesId) -> Result<&SeriesState, SeriesViewError> {
529 self.series
530 .get(id)
531 .ok_or_else(|| SeriesViewError::UnknownSeries {
532 series_id: id.clone(),
533 })
534 }
535}
536
537impl HistoricalSeriesView for MultiTimeframeSeries {
538 fn latest_bar(&self, id: &SeriesId) -> Result<Option<&ClosedBar>, SeriesViewError> {
539 Ok(self.state(id)?.closed.back())
540 }
541
542 fn bars(&self, id: &SeriesId, count: usize) -> Result<BarWindow<'_>, SeriesViewError> {
543 let state = self.state(id)?;
544 let (older, newer) = state.closed.as_slices();
545 let available = count.min(state.closed.len());
546 let skip = state.closed.len() - available;
547 if skip < older.len() {
548 Ok(BarWindow {
549 older: &older[skip..],
550 newer,
551 })
552 } else {
553 Ok(BarWindow {
554 older: &older[older.len()..],
555 newer: &newer[skip - older.len()..],
556 })
557 }
558 }
559
560 fn warmup(&self, id: &SeriesId) -> Result<SeriesWarmupState, SeriesViewError> {
561 let state = self.state(id)?;
562 Ok(SeriesWarmupState {
563 required: state.spec.requirement.warmup(),
564 available_bars: state.completed_bars,
565 })
566 }
567}
568
569fn validate_batch(batch: &TimestampBatch) -> Result<(), SeriesError> {
570 let mut ordering = BTreeSet::new();
571 for feed_event in &batch.events {
572 let event_ts = feed_event.event.ts();
573 if event_ts != batch.ts {
574 return Err(SeriesError::BatchTimestampMismatch {
575 batch_ts: batch.ts,
576 event_ts,
577 });
578 }
579 let key = (
580 feed_event.metadata.series_rank,
581 feed_event.metadata.row_sequence,
582 );
583 if !ordering.insert(key) {
584 return Err(SeriesError::DuplicateOrderingMetadata {
585 timestamp: batch.ts,
586 series_rank: key.0,
587 row_sequence: key.1,
588 });
589 }
590 }
591 Ok(())
592}
593
594fn bucket_bounds(
595 series_id: &SeriesId,
596 timestamp: NaiveDateTime,
597 duration: i64,
598 offset: i64,
599) -> Result<(NaiveDateTime, NaiveDateTime), SeriesError> {
600 let timestamp_seconds = i128::from(timestamp.and_utc().timestamp());
601 let duration = i128::from(duration);
602 let offset = i128::from(offset);
603 let bucket_index = (timestamp_seconds - offset).div_euclid(duration);
604 let open_seconds = bucket_index
605 .checked_mul(duration)
606 .and_then(|value| value.checked_add(offset))
607 .ok_or_else(|| SeriesError::BoundaryOverflow {
608 series_id: series_id.clone(),
609 timestamp,
610 })?;
611 let close_seconds =
612 open_seconds
613 .checked_add(duration)
614 .ok_or_else(|| SeriesError::BoundaryOverflow {
615 series_id: series_id.clone(),
616 timestamp,
617 })?;
618 let open_seconds = i64::try_from(open_seconds).map_err(|_| SeriesError::BoundaryOverflow {
619 series_id: series_id.clone(),
620 timestamp,
621 })?;
622 let close_seconds =
623 i64::try_from(close_seconds).map_err(|_| SeriesError::BoundaryOverflow {
624 series_id: series_id.clone(),
625 timestamp,
626 })?;
627 let open_time = DateTime::from_timestamp(open_seconds, 0)
628 .map(|value| value.naive_utc())
629 .ok_or_else(|| SeriesError::BoundaryOverflow {
630 series_id: series_id.clone(),
631 timestamp,
632 })?;
633 let close_time = DateTime::from_timestamp(close_seconds, 0)
634 .map(|value| value.naive_utc())
635 .ok_or_else(|| SeriesError::BoundaryOverflow {
636 series_id: series_id.clone(),
637 timestamp,
638 })?;
639 Ok((open_time, close_time))
640}