1use std::cell::RefCell;
4use std::collections::{BTreeMap, BTreeSet};
5
6use chrono::NaiveDateTime;
7use qs_core::types::{Effect, PositionStatus};
8use qs_core::{ManagementProfile, RuleConfigDef, StoplossMode, TradeEngine};
9use qs_strategy::{
10 CommandFact, CommandFeedback, CommandTerminalStatus, ConfiguredActionKind, ConfiguredCommand,
11 ConfiguredStrategy, ConfiguredStrategyRequirements, DecisionKind, MAX_GENERATED_ID_BYTES,
12 MAX_ID_BYTES, MAX_NAMED_VALUES, MAX_OUTPUT_COMMANDS, MAX_OUTPUT_NOTES, MAX_TEXT_BYTES,
13 NamedValue, NoteKind, OutputScalar, SourceId, StrategyInput, TradeSlotFacts, TradeSlotState,
14 Value, ValueType,
15};
16
17use crate::future_executor::FutureExecutor;
18use crate::ledger::ActionDispositionStatus;
19use crate::portfolio::CampaignExcursion;
20use crate::profile::PreparedEntryProfiles;
21
22use super::{
23 BarSeriesSpec, ClosedBar, HistoricalObservationView, HistoricalSeriesView, JournalKind,
24 SeriesId, StrategyDecisionDraft, StrategyDecisionKind, StrategyDescriptor, StrategyDomainError,
25 StrategyFeedbackEvent, StrategyJournalDraft, StrategyJournalError, StrategyObservation,
26 StrategyRequirements, StrategyResearchLimits, StrategyRetentionLimits,
27};
28
29const MAX_EXACT_F64_INTEGER: u64 = 1_u64 << 53;
30
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
33pub enum HistoricalVolumeProjection {
34 TickCountExact,
35 OptionalTickCount,
36}
37
38#[derive(Debug, Clone, PartialEq, Eq)]
40pub struct ConfiguredSourceBinding {
41 source: SourceId,
42 series: BarSeriesSpec,
43}
44
45impl ConfiguredSourceBinding {
46 pub fn new(source: SourceId, series: BarSeriesSpec) -> Self {
47 Self { source, series }
48 }
49
50 pub fn source(&self) -> &SourceId {
51 &self.source
52 }
53
54 pub fn series(&self) -> &BarSeriesSpec {
55 &self.series
56 }
57
58 pub fn series_id(&self) -> &SeriesId {
59 self.series.requirement().id()
60 }
61}
62
63#[derive(Clone, Copy)]
65pub struct NamedInputProjectionContext<'a> {
66 pub observed_through: NaiveDateTime,
67 pub closed_bars: &'a [ClosedBar],
68 pub observations: &'a [StrategyObservation],
69 pub series: &'a dyn HistoricalSeriesView,
70 pub observation_history: &'a dyn HistoricalObservationView,
71}
72
73#[derive(Debug, Clone, PartialEq)]
75pub struct ProjectedNamedInput {
76 pub value: Value,
77 pub updated: bool,
78}
79
80pub trait HistoricalNamedInputProjector: Send {
84 fn output_type(&self) -> ValueType;
85
86 fn project(
87 &self,
88 context: NamedInputProjectionContext<'_>,
89 ) -> Result<ProjectedNamedInput, NamedInputProjectionError>;
90}
91
92#[derive(Debug, Clone, Copy, PartialEq, Eq)]
94pub enum SourceBarFactKind {
95 Ordinal,
96 OpenTime,
97 CloseTime,
98 AvailableAt,
99 GapBefore,
100}
101
102#[derive(Debug, Clone, Copy, Default)]
103struct SourceBarFactState {
104 last_open: Option<NaiveDateTime>,
105 last_close: Option<NaiveDateTime>,
106 ordinal: u64,
107}
108
109pub struct SourceBarFactProjector {
111 series_id: SeriesId,
112 kind: SourceBarFactKind,
113 state: RefCell<SourceBarFactState>,
114}
115
116impl SourceBarFactProjector {
117 pub fn new(series_id: SeriesId, kind: SourceBarFactKind) -> Self {
118 Self {
119 series_id,
120 kind,
121 state: RefCell::new(SourceBarFactState::default()),
122 }
123 }
124}
125
126impl HistoricalNamedInputProjector for SourceBarFactProjector {
127 fn output_type(&self) -> ValueType {
128 ValueType::optional(match self.kind {
129 SourceBarFactKind::Ordinal => qs_strategy::ScalarType::Integer,
130 SourceBarFactKind::OpenTime
131 | SourceBarFactKind::CloseTime
132 | SourceBarFactKind::AvailableAt => qs_strategy::ScalarType::Timestamp,
133 SourceBarFactKind::GapBefore => qs_strategy::ScalarType::Bool,
134 })
135 }
136
137 fn project(
138 &self,
139 context: NamedInputProjectionContext<'_>,
140 ) -> Result<ProjectedNamedInput, NamedInputProjectionError> {
141 let Some(bar) = context
142 .closed_bars
143 .iter()
144 .find(|bar| bar.series_id() == &self.series_id)
145 else {
146 return Ok(ProjectedNamedInput {
147 value: Value::Missing(self.output_type().scalar),
148 updated: false,
149 });
150 };
151 let mut state = self.state.borrow_mut();
152 let is_new = state.last_open != Some(bar.open_time());
153 let previous_close = state.last_close;
154 if is_new {
155 state.ordinal = state
156 .ordinal
157 .checked_add(1)
158 .ok_or_else(|| NamedInputProjectionError::new("source ordinal overflowed"))?;
159 state.last_open = Some(bar.open_time());
160 state.last_close = Some(bar.close_time());
161 }
162 let value = match self.kind {
163 SourceBarFactKind::Ordinal => Value::Integer(
164 i64::try_from(state.ordinal)
165 .map_err(|_| NamedInputProjectionError::new("source ordinal exceeds i64"))?,
166 ),
167 SourceBarFactKind::OpenTime => Value::Timestamp(bar.open_time()),
168 SourceBarFactKind::CloseTime => Value::Timestamp(bar.close_time()),
169 SourceBarFactKind::AvailableAt => Value::Timestamp(context.observed_through),
170 SourceBarFactKind::GapBefore => {
171 Value::Bool(previous_close.is_some_and(|close| close != bar.open_time()))
172 }
173 };
174 Ok(ProjectedNamedInput {
175 value,
176 updated: is_new,
177 })
178 }
179}
180
181#[derive(Debug, Clone, Copy, PartialEq, Eq)]
182pub enum ConfirmedSwingFactKind {
183 Price,
184 AnchorOpenTime,
185 AnchorCloseTime,
186 ConfirmedAt,
187}
188
189pub struct ConfirmedSwingFactProjector {
191 series_id: SeriesId,
192 swing_kind: super::SwingKind,
193 fact: ConfirmedSwingFactKind,
194 retained: RefCell<Option<(u64, Value)>>,
195}
196
197impl ConfirmedSwingFactProjector {
198 pub fn new(
199 series_id: SeriesId,
200 swing_kind: super::SwingKind,
201 fact: ConfirmedSwingFactKind,
202 ) -> Self {
203 Self {
204 series_id,
205 swing_kind,
206 fact,
207 retained: RefCell::new(None),
208 }
209 }
210}
211
212impl HistoricalNamedInputProjector for ConfirmedSwingFactProjector {
213 fn output_type(&self) -> ValueType {
214 ValueType::optional(match self.fact {
215 ConfirmedSwingFactKind::Price => qs_strategy::ScalarType::Price,
216 _ => qs_strategy::ScalarType::Timestamp,
217 })
218 }
219
220 fn project(
221 &self,
222 context: NamedInputProjectionContext<'_>,
223 ) -> Result<ProjectedNamedInput, NamedInputProjectionError> {
224 let newest = context
225 .observations
226 .iter()
227 .filter(|observation| observation.source_series().contains(&self.series_id))
228 .filter_map(|observation| {
229 observation
230 .value()
231 .swing()
232 .map(|swing| (observation.sequence(), swing))
233 })
234 .filter(|(_, swing)| swing.kind() == self.swing_kind)
235 .max_by_key(|(sequence, _)| *sequence);
236 let mut retained = self.retained.borrow_mut();
237 let updated = newest.is_some_and(|(sequence, _)| {
238 retained
239 .as_ref()
240 .is_none_or(|(previous, _)| sequence > *previous)
241 });
242 if let Some((sequence, swing)) = newest
243 && updated
244 {
245 let value = match self.fact {
246 ConfirmedSwingFactKind::Price => Value::Price(swing.price()),
247 ConfirmedSwingFactKind::AnchorOpenTime => {
248 Value::Timestamp(swing.anchor_open_time())
249 }
250 ConfirmedSwingFactKind::AnchorCloseTime => {
251 Value::Timestamp(swing.anchor_close_time())
252 }
253 ConfirmedSwingFactKind::ConfirmedAt => Value::Timestamp(swing.confirmed_at()),
254 };
255 *retained = Some((sequence, value));
256 }
257 Ok(ProjectedNamedInput {
258 value: retained
259 .as_ref()
260 .map(|(_, value)| value.clone())
261 .unwrap_or(Value::Missing(self.output_type().scalar)),
262 updated,
263 })
264 }
265}
266
267pub struct ConfiguredNamedInputBinding {
269 name: String,
270 projector: Box<dyn HistoricalNamedInputProjector>,
271}
272
273impl ConfiguredNamedInputBinding {
274 pub fn new(name: impl Into<String>, projector: Box<dyn HistoricalNamedInputProjector>) -> Self {
275 Self {
276 name: name.into(),
277 projector,
278 }
279 }
280
281 pub fn name(&self) -> &str {
282 &self.name
283 }
284
285 pub fn output_type(&self) -> ValueType {
286 self.projector.output_type()
287 }
288}
289
290pub struct ConfiguredHistoricalBindings {
292 sources: Vec<ConfiguredSourceBinding>,
293 named_inputs: Vec<ConfiguredNamedInputBinding>,
294 volume: HistoricalVolumeProjection,
295}
296
297impl ConfiguredHistoricalBindings {
298 pub fn new(
299 sources: Vec<ConfiguredSourceBinding>,
300 named_inputs: Vec<ConfiguredNamedInputBinding>,
301 volume: HistoricalVolumeProjection,
302 ) -> Self {
303 Self {
304 sources,
305 named_inputs,
306 volume,
307 }
308 }
309
310 pub fn sources(&self) -> &[ConfiguredSourceBinding] {
311 &self.sources
312 }
313
314 pub fn named_inputs(&self) -> &[ConfiguredNamedInputBinding] {
315 &self.named_inputs
316 }
317
318 pub fn volume(&self) -> HistoricalVolumeProjection {
319 self.volume
320 }
321
322 pub fn into_parts(
323 self,
324 ) -> (
325 Vec<ConfiguredSourceBinding>,
326 Vec<ConfiguredNamedInputBinding>,
327 HistoricalVolumeProjection,
328 ) {
329 (self.sources, self.named_inputs, self.volume)
330 }
331}
332
333#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
335#[error("{message}")]
336pub struct NamedInputProjectionError {
337 message: String,
338}
339
340impl NamedInputProjectionError {
341 pub fn new(message: impl Into<String>) -> Self {
342 Self {
343 message: message.into(),
344 }
345 }
346}
347
348#[derive(Debug, thiserror::Error)]
350pub enum ConfiguredStrategyAdapterBuildError {
351 #[error("configured source '{source_id}' has no historical binding")]
352 MissingSourceBinding { source_id: SourceId },
353 #[error("configured source '{source_id}' is bound more than once")]
354 DuplicateSourceBinding { source_id: SourceId },
355 #[error(
356 "historical series ID '{series_id}' cannot be bound to more than one configured source"
357 )]
358 DuplicateSeriesBinding { series_id: SeriesId },
359 #[error("source '{source_id}' is not declared by the configured strategy")]
360 UndeclaredSourceBinding { source_id: SourceId },
361 #[error(
362 "source '{source_id}' is bound to symbol '{series_symbol}', but the configured strategy primary symbol is '{primary_symbol}'"
363 )]
364 SourceSymbolMismatch {
365 source_id: SourceId,
366 primary_symbol: String,
367 series_symbol: String,
368 },
369 #[error(
370 "source '{source_id}' requires lookback {required}, but retained history is {retained}"
371 )]
372 RetentionBelowLookback {
373 source_id: SourceId,
374 required: usize,
375 retained: usize,
376 },
377 #[error("source '{source_id}' requires lookback {required}, but historical warmup is {warmup}")]
378 WarmupBelowLookback {
379 source_id: SourceId,
380 required: usize,
381 warmup: usize,
382 },
383 #[error("configured named input '{name}' has no projector")]
384 MissingNamedInputProjector { name: String },
385 #[error("configured named input '{name}' has more than one projector")]
386 DuplicateNamedInputProjector { name: String },
387 #[error("named input '{name}' expects {expected:?}, but its projector returns {actual:?}")]
388 NamedInputTypeMismatch {
389 name: String,
390 expected: ValueType,
391 actual: ValueType,
392 },
393 #[error("named input projector '{name}' is not required by the configured strategy")]
394 UndeclaredNamedInputProjector { name: String },
395 #[error(transparent)]
396 HistoricalRequirements(#[from] StrategyDomainError),
397}
398
399#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
401pub enum ConfiguredStrategyAdapterPreflightError {
402 #[error("decision reason capacity {actual} is below configured output capacity {required}")]
403 DecisionReasonCapacity { actual: usize, required: usize },
404 #[error("signal capacity {actual} is below configured output capacity {required}")]
405 SignalCapacity { actual: usize, required: usize },
406 #[error("journal callback capacity {actual} is below configured output capacity {required}")]
407 JournalCallbackCapacity { actual: usize, required: usize },
408 #[error("journal reason capacity {actual} is below configured output capacity {required}")]
409 JournalReasonCapacity { actual: usize, required: usize },
410 #[error("journal value capacity {actual} is below configured output capacity {required}")]
411 JournalValueCapacity { actual: usize, required: usize },
412 #[error("journal key capacity {actual} is below configured output capacity {required}")]
413 JournalKeyCapacity { actual: usize, required: usize },
414 #[error(
415 "historical trade identity capacity {actual} is below configured identity capacity {required}"
416 )]
417 TradeIdentityCapacity { actual: usize, required: usize },
418}
419
420#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
422pub enum ConfiguredEntryProfileError {
423 #[error("entry class `{entry_class}` on trade slot `{slot}` has no management-profile route")]
424 UnroutedEntryClass { entry_class: String, slot: String },
425 #[error(
426 "management profile `{profile}` manages the stoploss of trade slot `{slot}`, which the strategy also moves"
427 )]
428 StoplossOwnerConflict {
429 profile: String,
430 slot: String,
431 entry_class: Option<String>,
432 },
433}
434
435#[derive(Debug, thiserror::Error)]
437pub enum ConfiguredStrategyAdapterError {
438 #[error("source '{source_id}' produced more than one completed bar at {timestamp}")]
439 DuplicateSourceUpdate {
440 source_id: SourceId,
441 timestamp: NaiveDateTime,
442 },
443 #[error("source '{source_id}' requires a tick count but the bar count is unknown")]
444 MissingTickCount { source_id: SourceId },
445 #[error("tick count {tick_count} cannot be represented exactly as f64")]
446 TickCountNotExactlyRepresentable { tick_count: u64 },
447 #[error("named input '{name}' projection failed: {source}")]
448 NamedInput {
449 name: String,
450 source: NamedInputProjectionError,
451 },
452 #[error("named input '{name}' returned a value incompatible with {expected:?}")]
453 NamedInputValueType { name: String, expected: ValueType },
454 #[error("trade slot '{slot}' has inconsistent engine state: {reason}")]
455 TradeSlot { slot: String, reason: String },
456 #[error("configured command '{command_id}' received an incompatible committed effect")]
457 IncompatibleCommandEffect { command_id: String },
458 #[error("configured strategy evaluation failed: {0}")]
459 Evaluation(#[from] qs_strategy::EvaluationError),
460 #[error("configured decision mapping failed: {0}")]
461 Decision(#[from] StrategyDomainError),
462 #[error("configured note mapping failed: {0}")]
463 Journal(#[from] StrategyJournalError),
464 #[error("configured output integer cannot be represented exactly as f64")]
465 IntegerOutputPrecision,
466}
467
468#[derive(Debug, Clone)]
469struct CommandRoute {
470 action: ConfiguredActionKind,
471 slot: String,
472 fact_seen: bool,
473 terminal: Option<CommandTerminalStatus>,
474}
475
476pub(crate) struct ConfiguredBoundaryOutput {
477 pub decision: Option<StrategyDecisionDraft>,
478 pub journal: Vec<StrategyJournalDraft>,
479 pub commands: Vec<ConfiguredCommand>,
480}
481
482pub(crate) struct BoundaryPositionFacts<'a> {
486 excursions: &'a BTreeMap<String, CampaignExcursion>,
487 executor: &'a FutureExecutor,
488}
489
490impl<'a> BoundaryPositionFacts<'a> {
491 pub(crate) fn new(
492 excursions: &'a BTreeMap<String, CampaignExcursion>,
493 executor: &'a FutureExecutor,
494 ) -> Self {
495 Self {
496 excursions,
497 executor,
498 }
499 }
500
501 fn excursion(&self, position_id: &str) -> Option<CampaignExcursion> {
502 self.excursions
503 .get(position_id)
504 .copied()
505 .filter(|excursion| excursion.observations > 0)
506 }
507
508 fn initial_risk(&self, position_id: &str) -> Option<f64> {
509 self.executor.open_initial_risk(position_id)
510 }
511}
512
513pub struct BacktestConfiguredStrategyAdapter {
515 strategy: ConfiguredStrategy,
516 descriptor: StrategyDescriptor,
517 requirements: StrategyRequirements,
518 bindings: ConfiguredHistoricalBindings,
519 command_routes: BTreeMap<String, CommandRoute>,
520 evaluation_start: Option<NaiveDateTime>,
521 first_ready_at: Option<NaiveDateTime>,
522}
523
524impl BacktestConfiguredStrategyAdapter {
525 pub fn new(
526 strategy: ConfiguredStrategy,
527 descriptor: StrategyDescriptor,
528 bindings: ConfiguredHistoricalBindings,
529 decision_latency_ms: u64,
530 ) -> Result<Self, ConfiguredStrategyAdapterBuildError> {
531 validate_bindings(&strategy, &bindings)?;
532 let series = bindings
533 .sources
534 .iter()
535 .map(|binding| binding.series.requirement().clone())
536 .collect::<Vec<_>>();
537 let mut instruments = Vec::new();
538 for requirement in &series {
539 if !instruments
540 .iter()
541 .any(|symbol| symbol == requirement.symbol())
542 {
543 instruments.push(requirement.symbol().to_owned());
544 }
545 }
546 let needs_feedback = strategy.input_requirements().needs_command_feedback;
547 let requirements = StrategyRequirements::new(
548 instruments,
549 series,
550 decision_latency_ms,
551 true,
552 needs_feedback,
553 )?;
554 Ok(Self {
555 strategy,
556 descriptor,
557 requirements,
558 bindings,
559 command_routes: BTreeMap::new(),
560 evaluation_start: None,
561 first_ready_at: None,
562 })
563 }
564
565 pub fn descriptor(&self) -> &StrategyDescriptor {
566 &self.descriptor
567 }
568
569 pub fn requirements(&self) -> &StrategyRequirements {
570 &self.requirements
571 }
572
573 pub fn configured_requirements(&self) -> &ConfiguredStrategyRequirements {
574 self.strategy.input_requirements()
575 }
576
577 pub fn source_bindings(&self) -> &[ConfiguredSourceBinding] {
578 &self.bindings.sources
579 }
580
581 pub fn series_specs(&self) -> impl ExactSizeIterator<Item = &BarSeriesSpec> {
582 self.bindings.sources.iter().map(|binding| &binding.series)
583 }
584
585 pub fn first_ready_at(&self) -> Option<NaiveDateTime> {
586 self.first_ready_at
587 }
588
589 pub fn set_evaluation_start(&mut self, evaluation_start: Option<NaiveDateTime>) {
590 self.evaluation_start = evaluation_start;
591 }
592
593 pub fn configured_strategy(&self) -> &ConfiguredStrategy {
594 &self.strategy
595 }
596
597 pub fn into_configured_strategy(self) -> ConfiguredStrategy {
598 self.strategy
599 }
600
601 pub fn preflight_entry_profiles(
603 &self,
604 profiles: &PreparedEntryProfiles,
605 ) -> Result<(), ConfiguredEntryProfileError> {
606 let requirements = self.strategy.input_requirements();
607 for entry in &requirements.entries {
608 let profile = match entry.entry_class.as_ref() {
609 Some(entry_class) => Some(profiles.routes().get(entry_class).ok_or_else(|| {
610 ConfiguredEntryProfileError::UnroutedEntryClass {
611 entry_class: entry_class.clone(),
612 slot: entry.slot.clone(),
613 }
614 })?),
615 None => profiles.default_profile(),
616 };
617 if let Some(profile) = profile
618 && profile_manages_stoploss(profile)
619 && requirements.stop_managed_slots.contains(&entry.slot)
620 {
621 return Err(ConfiguredEntryProfileError::StoplossOwnerConflict {
622 profile: profile.name.clone(),
623 slot: entry.slot.clone(),
624 entry_class: entry.entry_class.clone(),
625 });
626 }
627 }
628 Ok(())
629 }
630
631 pub fn preflight(
632 &self,
633 retention: StrategyRetentionLimits,
634 research: StrategyResearchLimits,
635 ) -> Result<(), ConfiguredStrategyAdapterPreflightError> {
636 if retention.max_reason_bytes() < MAX_TEXT_BYTES {
637 return Err(
638 ConfiguredStrategyAdapterPreflightError::DecisionReasonCapacity {
639 actual: retention.max_reason_bytes(),
640 required: MAX_TEXT_BYTES,
641 },
642 );
643 }
644 if retention.max_signals_per_callback() < MAX_OUTPUT_COMMANDS {
645 return Err(ConfiguredStrategyAdapterPreflightError::SignalCapacity {
646 actual: retention.max_signals_per_callback(),
647 required: MAX_OUTPUT_COMMANDS,
648 });
649 }
650 if research.max_journal_per_callback() < MAX_OUTPUT_NOTES {
651 return Err(
652 ConfiguredStrategyAdapterPreflightError::JournalCallbackCapacity {
653 actual: research.max_journal_per_callback(),
654 required: MAX_OUTPUT_NOTES,
655 },
656 );
657 }
658 if research.max_reason_bytes() < MAX_TEXT_BYTES {
659 return Err(
660 ConfiguredStrategyAdapterPreflightError::JournalReasonCapacity {
661 actual: research.max_reason_bytes(),
662 required: MAX_TEXT_BYTES,
663 },
664 );
665 }
666 if research.max_values_per_record() < MAX_NAMED_VALUES {
667 return Err(
668 ConfiguredStrategyAdapterPreflightError::JournalValueCapacity {
669 actual: research.max_values_per_record(),
670 required: MAX_NAMED_VALUES,
671 },
672 );
673 }
674 if research.max_value_key_bytes() < MAX_ID_BYTES {
675 return Err(
676 ConfiguredStrategyAdapterPreflightError::JournalKeyCapacity {
677 actual: research.max_value_key_bytes(),
678 required: MAX_ID_BYTES,
679 },
680 );
681 }
682 if super::MAX_TRADE_ID_BYTES < MAX_GENERATED_ID_BYTES {
683 return Err(
684 ConfiguredStrategyAdapterPreflightError::TradeIdentityCapacity {
685 actual: super::MAX_TRADE_ID_BYTES,
686 required: MAX_GENERATED_ID_BYTES,
687 },
688 );
689 }
690 Ok(())
691 }
692
693 #[allow(clippy::too_many_arguments)]
694 pub(crate) fn evaluate_boundary(
695 &mut self,
696 observed_through: NaiveDateTime,
697 ready: bool,
698 closed_bars: &[ClosedBar],
699 observations: &[StrategyObservation],
700 series: &dyn HistoricalSeriesView,
701 observation_history: &dyn HistoricalObservationView,
702 engine: &TradeEngine,
703 positions: &BoundaryPositionFacts<'_>,
704 feedback_events: &[StrategyFeedbackEvent],
705 retention: StrategyRetentionLimits,
706 research: StrategyResearchLimits,
707 ) -> Result<ConfiguredBoundaryOutput, ConfiguredStrategyAdapterError> {
708 let ready = ready
709 && self
710 .evaluation_start
711 .is_none_or(|start| observed_through >= start);
712 if ready && self.first_ready_at.is_none() {
713 self.first_ready_at = Some(observed_through);
714 }
715 let feedback = self.project_feedback(feedback_events)?;
716 let input = StrategyInput {
717 time: observed_through,
718 ready,
719 completed_bars: self.project_bars(observed_through, closed_bars)?,
720 values: self.project_named_inputs(NamedInputProjectionContext {
721 observed_through,
722 closed_bars,
723 observations,
724 series,
725 observation_history,
726 })?,
727 trade_slots: self.project_trade_slots(engine, positions)?,
728 feedback,
729 };
730 let output = self.strategy.evaluate(&input)?;
731 for command in &output.commands {
732 self.command_routes.insert(
733 command.command_id.clone(),
734 CommandRoute {
735 action: command.action_kind,
736 slot: command.trade_slot.clone(),
737 fact_seen: false,
738 terminal: None,
739 },
740 );
741 }
742 let emitted_signals = output
743 .commands
744 .iter()
745 .map(|command| command.signal.clone())
746 .collect::<Vec<_>>();
747 let decision = output
748 .decision
749 .map(|decision| map_decision(decision, emitted_signals, retention))
750 .transpose()?;
751 let journal = output
752 .notes
753 .into_iter()
754 .map(|note| map_note(note, self.strategy.primary_symbol(), research))
755 .collect::<Result<Vec<_>, _>>()?;
756 Ok(ConfiguredBoundaryOutput {
757 decision,
758 journal,
759 commands: output.commands,
760 })
761 }
762
763 fn project_bars(
764 &self,
765 observed_through: NaiveDateTime,
766 closed_bars: &[ClosedBar],
767 ) -> Result<Vec<qs_strategy::CompletedBarUpdate>, ConfiguredStrategyAdapterError> {
768 self.strategy
769 .input_requirements()
770 .completed_bars
771 .iter()
772 .filter_map(|requirement| {
773 let binding = self
774 .bindings
775 .sources
776 .iter()
777 .find(|binding| binding.source == requirement.source)
778 .expect("bindings were validated at construction");
779 let mut matching = closed_bars
780 .iter()
781 .filter(|bar| bar.series_id() == binding.series_id());
782 let bar = matching.next()?;
783 Some(if matching.next().is_some() {
784 Err(ConfiguredStrategyAdapterError::DuplicateSourceUpdate {
785 source_id: requirement.source.clone(),
786 timestamp: observed_through,
787 })
788 } else {
789 Ok(qs_strategy::CompletedBarUpdate {
790 source: requirement.source.clone(),
791 bar: qs_strategy::CompletedBar {
792 open: bar.open(),
793 high: bar.high(),
794 low: bar.low(),
795 close: bar.close(),
796 volume: match self.bindings.volume {
797 HistoricalVolumeProjection::TickCountExact => {
798 let tick_count = match bar.tick_count() {
799 Some(count) => count,
800 None => return Some(Err(ConfiguredStrategyAdapterError::MissingTickCount { source_id: requirement.source.clone() })),
801 };
802 if tick_count > MAX_EXACT_F64_INTEGER {
803 return Some(Err(
804 ConfiguredStrategyAdapterError::TickCountNotExactlyRepresentable { tick_count },
805 ));
806 }
807 Some(tick_count as f64)
808 }
809 HistoricalVolumeProjection::OptionalTickCount => bar.tick_count().map(|count| count as f64)
810 },
811 },
812 })
813 })
814 })
815 .collect()
816 }
817
818 fn project_named_inputs(
819 &self,
820 context: NamedInputProjectionContext<'_>,
821 ) -> Result<Vec<NamedValue>, ConfiguredStrategyAdapterError> {
822 self.strategy
823 .input_requirements()
824 .named_inputs
825 .iter()
826 .map(|requirement| {
827 let binding = self
828 .bindings
829 .named_inputs
830 .iter()
831 .find(|binding| binding.name == requirement.name)
832 .expect("named input bindings were validated at construction");
833 let projected = binding.projector.project(context).map_err(|source| {
834 ConfiguredStrategyAdapterError::NamedInput {
835 name: requirement.name.clone(),
836 source,
837 }
838 })?;
839 if !value_matches_type(&projected.value, requirement.value_type) {
840 return Err(ConfiguredStrategyAdapterError::NamedInputValueType {
841 name: requirement.name.clone(),
842 expected: requirement.value_type,
843 });
844 }
845 Ok(NamedValue {
846 name: requirement.name.clone(),
847 value: projected.value,
848 updated: projected.updated,
849 })
850 })
851 .collect()
852 }
853
854 fn project_trade_slots(
855 &self,
856 engine: &TradeEngine,
857 positions: &BoundaryPositionFacts<'_>,
858 ) -> Result<Vec<TradeSlotFacts>, ConfiguredStrategyAdapterError> {
859 self.strategy
860 .input_requirements()
861 .trade_slots
862 .iter()
863 .map(|slot| {
864 let position = self
865 .strategy
866 .trade_id_for_slot(slot)
867 .and_then(|trade_id| engine.manager.id_by_trade_id(trade_id))
868 .and_then(|position_id| {
869 engine
870 .get_position(&position_id)
871 .map(|position| (position_id, position))
872 });
873 let state = match position {
874 None => TradeSlotState::Vacant,
875 Some((position_id, position)) => match position.data.status {
876 PositionStatus::Pending => TradeSlotState::Pending {
877 side: position.data.side,
878 requested_price: position.data.pending_price,
879 stoploss: position.current_stoploss(),
880 },
881 PositionStatus::Open => {
882 let opened_at = position.data.open_ts.ok_or_else(|| {
883 ConfiguredStrategyAdapterError::TradeSlot {
884 slot: slot.clone(),
885 reason: "open position has no entry fill time".into(),
886 }
887 })?;
888 let excursion = positions.excursion(&position_id);
889 TradeSlotState::Open {
890 side: position.data.side,
891 entry_price: position.data.average_entry(),
892 remaining_size: position.data.remaining_size(),
893 stoploss: position.current_stoploss(),
894 opened_at,
895 favorable_excursion: excursion.map(|excursion| excursion.mfe),
896 adverse_excursion: excursion.map(|excursion| excursion.mae),
897 initial_risk: positions.initial_risk(&position_id),
898 }
899 }
900 PositionStatus::Closed | PositionStatus::Cancelled => {
901 TradeSlotState::Vacant
902 }
903 },
904 };
905 Ok(TradeSlotFacts {
906 slot: slot.clone(),
907 state,
908 })
909 })
910 .collect()
911 }
912
913 fn project_feedback(
914 &mut self,
915 events: &[StrategyFeedbackEvent],
916 ) -> Result<Vec<CommandFeedback>, ConfiguredStrategyAdapterError> {
917 project_command_feedback(&mut self.command_routes, events)
918 }
919
920 pub(crate) fn finalize_feedback(
921 &mut self,
922 events: &[StrategyFeedbackEvent],
923 ) -> Result<(), ConfiguredStrategyAdapterError> {
924 let feedback = self.project_feedback(events)?;
925 self.strategy.finalize_command_feedback(&feedback)?;
926 self.command_routes.clear();
927 Ok(())
928 }
929}
930
931fn project_command_feedback(
932 routes: &mut BTreeMap<String, CommandRoute>,
933 events: &[StrategyFeedbackEvent],
934) -> Result<Vec<CommandFeedback>, ConfiguredStrategyAdapterError> {
935 let mut projected = Vec::new();
936 for event in events {
937 let Some(command_id) = event.action_id() else {
938 continue;
939 };
940 let Some(route) = routes.get_mut(command_id) else {
941 continue;
942 };
943 match event {
944 StrategyFeedbackEvent::Effect { effect, .. } => {
945 if let Some(fact) = map_effect(route.action, effect.effect()).map_err(|()| {
946 ConfiguredStrategyAdapterError::IncompatibleCommandEffect {
947 command_id: command_id.to_owned(),
948 }
949 })? {
950 route.fact_seen = true;
951 projected.push(CommandFeedback::Fact {
952 command_id: command_id.to_owned(),
953 fact,
954 });
955 }
956 }
957 StrategyFeedbackEvent::Disposition(disposition) => {
958 let status = match disposition.status {
959 ActionDispositionStatus::Applied => CommandTerminalStatus::Applied,
960 ActionDispositionStatus::Skipped => CommandTerminalStatus::Skipped,
961 ActionDispositionStatus::Rejected => CommandTerminalStatus::Rejected,
962 ActionDispositionStatus::Failed => CommandTerminalStatus::Failed,
963 };
964 route.terminal = Some(status);
965 projected.push(CommandFeedback::Terminal {
966 command_id: command_id.to_owned(),
967 status,
968 reason: disposition.reason.clone(),
969 });
970 }
971 }
972 let completed = route
973 .terminal
974 .is_some_and(|status| status != CommandTerminalStatus::Applied)
975 || (route.terminal == Some(CommandTerminalStatus::Applied) && route.fact_seen);
976 if completed {
977 let completed_route = routes
978 .remove(command_id)
979 .expect("completed route remains registered");
980 if completed_route.action == ConfiguredActionKind::CancelPending
981 && completed_route.terminal == Some(CommandTerminalStatus::Applied)
982 {
983 routes.retain(|_, route| {
984 !(route.action == ConfiguredActionKind::Entry
985 && route.slot == completed_route.slot)
986 });
987 }
988 }
989 }
990 Ok(projected)
991}
992
993fn validate_bindings(
994 strategy: &ConfiguredStrategy,
995 bindings: &ConfiguredHistoricalBindings,
996) -> Result<(), ConfiguredStrategyAdapterBuildError> {
997 let declared = strategy.declared_sources().iter().collect::<BTreeSet<_>>();
998 let mut sources = BTreeSet::new();
999 let mut series = BTreeSet::new();
1000 for binding in &bindings.sources {
1001 if !declared.contains(&binding.source) {
1002 return Err(
1003 ConfiguredStrategyAdapterBuildError::UndeclaredSourceBinding {
1004 source_id: binding.source.clone(),
1005 },
1006 );
1007 }
1008 if !sources.insert(binding.source.clone()) {
1009 return Err(
1010 ConfiguredStrategyAdapterBuildError::DuplicateSourceBinding {
1011 source_id: binding.source.clone(),
1012 },
1013 );
1014 }
1015 if !series.insert(binding.series_id().clone()) {
1016 return Err(
1017 ConfiguredStrategyAdapterBuildError::DuplicateSeriesBinding {
1018 series_id: binding.series_id().clone(),
1019 },
1020 );
1021 }
1022 let series_symbol = binding.series.requirement().symbol();
1023 if series_symbol != strategy.primary_symbol() {
1024 return Err(ConfiguredStrategyAdapterBuildError::SourceSymbolMismatch {
1025 source_id: binding.source.clone(),
1026 primary_symbol: strategy.primary_symbol().to_owned(),
1027 series_symbol: series_symbol.to_owned(),
1028 });
1029 }
1030 }
1031 for source in strategy.declared_sources() {
1032 if !sources.contains(source) {
1033 return Err(ConfiguredStrategyAdapterBuildError::MissingSourceBinding {
1034 source_id: source.clone(),
1035 });
1036 }
1037 }
1038 for requirement in &strategy.input_requirements().completed_bars {
1039 let binding = bindings
1040 .sources
1041 .iter()
1042 .find(|binding| binding.source == requirement.source)
1043 .expect("every declared source was checked above");
1044 if binding.series.retained_bars() < requirement.required_lookback {
1045 return Err(
1046 ConfiguredStrategyAdapterBuildError::RetentionBelowLookback {
1047 source_id: requirement.source.clone(),
1048 required: requirement.required_lookback,
1049 retained: binding.series.retained_bars(),
1050 },
1051 );
1052 }
1053 let warmup = binding.series.requirement().warmup().required_bars();
1054 if warmup < requirement.required_lookback {
1055 return Err(ConfiguredStrategyAdapterBuildError::WarmupBelowLookback {
1056 source_id: requirement.source.clone(),
1057 required: requirement.required_lookback,
1058 warmup,
1059 });
1060 }
1061 }
1062 let mut names = BTreeSet::new();
1063 for binding in &bindings.named_inputs {
1064 if !names.insert(binding.name.clone()) {
1065 return Err(
1066 ConfiguredStrategyAdapterBuildError::DuplicateNamedInputProjector {
1067 name: binding.name.clone(),
1068 },
1069 );
1070 }
1071 let Some(requirement) = strategy
1072 .input_requirements()
1073 .named_inputs
1074 .iter()
1075 .find(|requirement| requirement.name == binding.name)
1076 else {
1077 return Err(
1078 ConfiguredStrategyAdapterBuildError::UndeclaredNamedInputProjector {
1079 name: binding.name.clone(),
1080 },
1081 );
1082 };
1083 let actual = binding.output_type();
1084 if actual != requirement.value_type {
1085 return Err(
1086 ConfiguredStrategyAdapterBuildError::NamedInputTypeMismatch {
1087 name: binding.name.clone(),
1088 expected: requirement.value_type,
1089 actual,
1090 },
1091 );
1092 }
1093 }
1094 for requirement in &strategy.input_requirements().named_inputs {
1095 if !names.contains(&requirement.name) {
1096 return Err(
1097 ConfiguredStrategyAdapterBuildError::MissingNamedInputProjector {
1098 name: requirement.name.clone(),
1099 },
1100 );
1101 }
1102 }
1103 Ok(())
1104}
1105
1106fn value_matches_type(value: &Value, expected: ValueType) -> bool {
1107 if value.is_missing() {
1108 return expected.optional && value.scalar_type() == expected.scalar;
1109 }
1110 if value.scalar_type() != expected.scalar {
1111 return false;
1112 }
1113 match value {
1114 Value::Number(value)
1115 | Value::Price(value)
1116 | Value::Ratio(value)
1117 | Value::Percent(value)
1118 | Value::PricePerObservation(value)
1119 | Value::PricePerObservationSquared(value)
1120 | Value::RatioPerObservation(value)
1121 | Value::RatioPerObservationSquared(value)
1122 | Value::LogReturn(value)
1123 | Value::LogReturnVariance(value) => value.is_finite(),
1124 Value::Text(value) => !value.is_empty() && value.len() <= MAX_TEXT_BYTES,
1125 _ => true,
1126 }
1127}
1128
1129fn profile_manages_stoploss(profile: &ManagementProfile) -> bool {
1131 !matches!(profile.stoploss_mode, StoplossMode::FromSignal)
1132 || profile.rules.iter().any(|rule| {
1133 matches!(
1134 rule,
1135 RuleConfigDef::FixedStoploss { .. }
1136 | RuleConfigDef::TrailingStop { .. }
1137 | RuleConfigDef::BreakevenWhen { .. }
1138 | RuleConfigDef::BreakevenWhenOffset { .. }
1139 | RuleConfigDef::BreakevenAfterTargets { .. }
1140 )
1141 })
1142}
1143
1144fn map_effect(action: ConfiguredActionKind, effect: &Effect) -> Result<Option<CommandFact>, ()> {
1145 let mapped = match effect {
1146 Effect::PositionOpened { .. } => {
1147 Some((ConfiguredActionKind::Entry, CommandFact::EntryFilled))
1148 }
1149 Effect::PositionClosed { .. } => {
1150 Some((ConfiguredActionKind::Close, CommandFact::PositionClosed))
1151 }
1152 Effect::PartialClose { .. } => Some((
1153 ConfiguredActionKind::ClosePartial,
1154 CommandFact::PositionReduced,
1155 )),
1156 Effect::StoplossModified { .. } => match action {
1157 ConfiguredActionKind::MoveStoplossToEntry | ConfiguredActionKind::ModifyStoploss => {
1158 return Ok(Some(CommandFact::StoplossModified));
1159 }
1160 _ => return Err(()),
1161 },
1162 Effect::OrderCancelled { .. } => Some((
1163 ConfiguredActionKind::CancelPending,
1164 CommandFact::PendingCancelled,
1165 )),
1166 Effect::OrderPlaced { .. }
1167 | Effect::StoplossRemoved { .. }
1168 | Effect::ScaledIn { .. }
1169 | Effect::RuleTriggered { .. } => None,
1170 };
1171 match mapped {
1172 Some((expected, fact)) if expected == action => Ok(Some(fact)),
1173 Some(_) => Err(()),
1174 None => Ok(None),
1175 }
1176}
1177
1178fn map_decision(
1179 decision: qs_strategy::Decision,
1180 emitted_signals: Vec<qs_core::RawSignal>,
1181 limits: StrategyRetentionLimits,
1182) -> Result<StrategyDecisionDraft, StrategyDomainError> {
1183 let kind = match decision.kind {
1184 DecisionKind::Entry => StrategyDecisionKind::Entry,
1185 DecisionKind::Management => StrategyDecisionKind::Management,
1186 DecisionKind::Exit => StrategyDecisionKind::Exit,
1187 DecisionKind::Observation => StrategyDecisionKind::Annotation,
1188 };
1189 StrategyDecisionDraft::new(
1190 kind,
1191 decision.reason,
1192 decision.related_trade.map(|trade| trade.trade_id),
1193 emitted_signals,
1194 limits,
1195 )
1196}
1197
1198fn map_note(
1199 note: qs_strategy::Note,
1200 symbol: &str,
1201 limits: StrategyResearchLimits,
1202) -> Result<StrategyJournalDraft, ConfiguredStrategyAdapterError> {
1203 let kind = match note.kind {
1204 NoteKind::Observation | NoteKind::Risk => JournalKind::DecisionContext,
1205 NoteKind::Execution | NoteKind::Lifecycle => JournalKind::OutcomeReview,
1206 };
1207 let mut values = BTreeMap::new();
1208 for output in note.values {
1209 let value = match output.value {
1210 OutputScalar::Bool(value) => f64::from(value),
1211 OutputScalar::Integer(value) => {
1212 if value.unsigned_abs() > MAX_EXACT_F64_INTEGER {
1213 return Err(ConfiguredStrategyAdapterError::IntegerOutputPrecision);
1214 }
1215 value as f64
1216 }
1217 OutputScalar::Number(value)
1218 | OutputScalar::Price(value)
1219 | OutputScalar::Ratio(value)
1220 | OutputScalar::Percent(value)
1221 | OutputScalar::PricePerObservation(value)
1222 | OutputScalar::PricePerObservationSquared(value)
1223 | OutputScalar::RatioPerObservation(value)
1224 | OutputScalar::RatioPerObservationSquared(value)
1225 | OutputScalar::LogReturn(value)
1226 | OutputScalar::LogReturnVariance(value) => value,
1227 };
1228 values.insert(output.name, value);
1229 }
1230 Ok(StrategyJournalDraft::new(
1231 kind,
1232 symbol,
1233 note.related_trade.map(|trade| trade.trade_id),
1234 note.reason,
1235 None,
1236 values,
1237 limits,
1238 )?)
1239}
1240
1241#[cfg(test)]
1242mod tests {
1243 use super::*;
1244 use crate::ledger::ActionDisposition;
1245 use qs_core::types::FutureEffect;
1246
1247 #[test]
1248 fn applied_terminal_before_effect_completes_command_correlation() {
1249 let command_id = "opaque-command".to_owned();
1250 let mut routes = BTreeMap::from([(
1251 command_id.clone(),
1252 CommandRoute {
1253 action: ConfiguredActionKind::Entry,
1254 slot: "primary".into(),
1255 fact_seen: false,
1256 terminal: None,
1257 },
1258 )]);
1259 let terminal = project_command_feedback(
1260 &mut routes,
1261 &[StrategyFeedbackEvent::Disposition(
1262 ActionDisposition::applied(command_id.clone()),
1263 )],
1264 )
1265 .unwrap();
1266
1267 assert_eq!(
1268 terminal,
1269 vec![CommandFeedback::Terminal {
1270 command_id: command_id.clone(),
1271 status: CommandTerminalStatus::Applied,
1272 reason: None,
1273 }]
1274 );
1275 assert!(routes.contains_key(&command_id));
1276
1277 let fact = project_command_feedback(
1278 &mut routes,
1279 &[StrategyFeedbackEvent::Effect {
1280 action_id: Some(command_id.clone()),
1281 effect: FutureEffect::plain(Effect::PositionOpened {
1282 id: "position-1".into(),
1283 }),
1284 }],
1285 )
1286 .unwrap();
1287
1288 assert_eq!(
1289 fact,
1290 vec![CommandFeedback::Fact {
1291 command_id,
1292 fact: CommandFact::EntryFilled,
1293 }]
1294 );
1295 assert!(routes.is_empty());
1296 }
1297}