1use std::collections::{BTreeMap, BTreeSet};
4
5use chrono::NaiveDateTime;
6use qs_core::TradeEngine;
7use qs_core::types::{Effect, PositionStatus};
8use qs_strategy::{
9 CommandFact, CommandFeedback, CommandTerminalStatus, ConfiguredActionKind, ConfiguredCommand,
10 ConfiguredStrategy, ConfiguredStrategyRequirements, DecisionKind, MAX_GENERATED_ID_BYTES,
11 MAX_ID_BYTES, MAX_NAMED_VALUES, MAX_OUTPUT_COMMANDS, MAX_OUTPUT_NOTES, MAX_TEXT_BYTES,
12 NamedValue, NoteKind, OutputScalar, SourceId, StrategyInput, TradeSlotFacts, TradeSlotState,
13 Value, ValueType,
14};
15
16use crate::ledger::ActionDispositionStatus;
17
18use super::{
19 BarSeriesSpec, ClosedBar, HistoricalObservationView, HistoricalSeriesView, JournalKind,
20 SeriesId, StrategyDecisionDraft, StrategyDecisionKind, StrategyDescriptor, StrategyDomainError,
21 StrategyFeedbackEvent, StrategyJournalDraft, StrategyJournalError, StrategyObservation,
22 StrategyRequirements, StrategyResearchLimits, StrategyRetentionLimits,
23};
24
25const MAX_EXACT_F64_INTEGER: u64 = 1_u64 << 53;
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum HistoricalVolumeProjection {
30 TickCountExact,
31}
32
33#[derive(Debug, Clone, PartialEq, Eq)]
35pub struct ConfiguredSourceBinding {
36 source: SourceId,
37 series: BarSeriesSpec,
38}
39
40impl ConfiguredSourceBinding {
41 pub fn new(source: SourceId, series: BarSeriesSpec) -> Self {
42 Self { source, series }
43 }
44
45 pub fn source(&self) -> &SourceId {
46 &self.source
47 }
48
49 pub fn series(&self) -> &BarSeriesSpec {
50 &self.series
51 }
52
53 pub fn series_id(&self) -> &SeriesId {
54 self.series.requirement().id()
55 }
56}
57
58#[derive(Clone, Copy)]
60pub struct NamedInputProjectionContext<'a> {
61 pub observed_through: NaiveDateTime,
62 pub closed_bars: &'a [ClosedBar],
63 pub observations: &'a [StrategyObservation],
64 pub series: &'a dyn HistoricalSeriesView,
65 pub observation_history: &'a dyn HistoricalObservationView,
66}
67
68#[derive(Debug, Clone, PartialEq)]
70pub struct ProjectedNamedInput {
71 pub value: Value,
72 pub updated: bool,
73}
74
75pub trait HistoricalNamedInputProjector {
77 fn output_type(&self) -> ValueType;
78
79 fn project(
80 &self,
81 context: NamedInputProjectionContext<'_>,
82 ) -> Result<ProjectedNamedInput, NamedInputProjectionError>;
83}
84
85pub struct ConfiguredNamedInputBinding {
87 name: String,
88 projector: Box<dyn HistoricalNamedInputProjector>,
89}
90
91impl ConfiguredNamedInputBinding {
92 pub fn new(name: impl Into<String>, projector: Box<dyn HistoricalNamedInputProjector>) -> Self {
93 Self {
94 name: name.into(),
95 projector,
96 }
97 }
98
99 pub fn name(&self) -> &str {
100 &self.name
101 }
102
103 pub fn output_type(&self) -> ValueType {
104 self.projector.output_type()
105 }
106}
107
108pub struct ConfiguredHistoricalBindings {
110 sources: Vec<ConfiguredSourceBinding>,
111 named_inputs: Vec<ConfiguredNamedInputBinding>,
112 volume: HistoricalVolumeProjection,
113}
114
115impl ConfiguredHistoricalBindings {
116 pub fn new(
117 sources: Vec<ConfiguredSourceBinding>,
118 named_inputs: Vec<ConfiguredNamedInputBinding>,
119 volume: HistoricalVolumeProjection,
120 ) -> Self {
121 Self {
122 sources,
123 named_inputs,
124 volume,
125 }
126 }
127
128 pub fn sources(&self) -> &[ConfiguredSourceBinding] {
129 &self.sources
130 }
131
132 pub fn named_inputs(&self) -> &[ConfiguredNamedInputBinding] {
133 &self.named_inputs
134 }
135
136 pub fn volume(&self) -> HistoricalVolumeProjection {
137 self.volume
138 }
139}
140
141#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
143#[error("{message}")]
144pub struct NamedInputProjectionError {
145 message: String,
146}
147
148impl NamedInputProjectionError {
149 pub fn new(message: impl Into<String>) -> Self {
150 Self {
151 message: message.into(),
152 }
153 }
154}
155
156#[derive(Debug, thiserror::Error)]
158pub enum ConfiguredStrategyAdapterBuildError {
159 #[error("configured source '{source_id}' has no historical binding")]
160 MissingSourceBinding { source_id: SourceId },
161 #[error("configured source '{source_id}' is bound more than once")]
162 DuplicateSourceBinding { source_id: SourceId },
163 #[error(
164 "historical series ID '{series_id}' cannot be bound to more than one configured source"
165 )]
166 DuplicateSeriesBinding { series_id: SeriesId },
167 #[error("source '{source_id}' is not declared by the configured strategy")]
168 UndeclaredSourceBinding { source_id: SourceId },
169 #[error(
170 "source '{source_id}' is bound to symbol '{series_symbol}', but the configured strategy primary symbol is '{primary_symbol}'"
171 )]
172 SourceSymbolMismatch {
173 source_id: SourceId,
174 primary_symbol: String,
175 series_symbol: String,
176 },
177 #[error(
178 "source '{source_id}' requires lookback {required}, but retained history is {retained}"
179 )]
180 RetentionBelowLookback {
181 source_id: SourceId,
182 required: usize,
183 retained: usize,
184 },
185 #[error("source '{source_id}' requires lookback {required}, but historical warmup is {warmup}")]
186 WarmupBelowLookback {
187 source_id: SourceId,
188 required: usize,
189 warmup: usize,
190 },
191 #[error("configured named input '{name}' has no projector")]
192 MissingNamedInputProjector { name: String },
193 #[error("configured named input '{name}' has more than one projector")]
194 DuplicateNamedInputProjector { name: String },
195 #[error("named input '{name}' expects {expected:?}, but its projector returns {actual:?}")]
196 NamedInputTypeMismatch {
197 name: String,
198 expected: ValueType,
199 actual: ValueType,
200 },
201 #[error("named input projector '{name}' is not required by the configured strategy")]
202 UndeclaredNamedInputProjector { name: String },
203 #[error(transparent)]
204 HistoricalRequirements(#[from] StrategyDomainError),
205}
206
207#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
209pub enum ConfiguredStrategyAdapterPreflightError {
210 #[error("decision reason capacity {actual} is below configured output capacity {required}")]
211 DecisionReasonCapacity { actual: usize, required: usize },
212 #[error("signal capacity {actual} is below configured output capacity {required}")]
213 SignalCapacity { actual: usize, required: usize },
214 #[error("journal callback capacity {actual} is below configured output capacity {required}")]
215 JournalCallbackCapacity { actual: usize, required: usize },
216 #[error("journal reason capacity {actual} is below configured output capacity {required}")]
217 JournalReasonCapacity { actual: usize, required: usize },
218 #[error("journal value capacity {actual} is below configured output capacity {required}")]
219 JournalValueCapacity { actual: usize, required: usize },
220 #[error("journal key capacity {actual} is below configured output capacity {required}")]
221 JournalKeyCapacity { actual: usize, required: usize },
222 #[error(
223 "historical trade identity capacity {actual} is below configured identity capacity {required}"
224 )]
225 TradeIdentityCapacity { actual: usize, required: usize },
226}
227
228#[derive(Debug, thiserror::Error)]
230pub enum ConfiguredStrategyAdapterError {
231 #[error("source '{source_id}' produced more than one completed bar at {timestamp}")]
232 DuplicateSourceUpdate {
233 source_id: SourceId,
234 timestamp: NaiveDateTime,
235 },
236 #[error("tick count {tick_count} cannot be represented exactly as f64")]
237 TickCountNotExactlyRepresentable { tick_count: u64 },
238 #[error("named input '{name}' projection failed: {source}")]
239 NamedInput {
240 name: String,
241 source: NamedInputProjectionError,
242 },
243 #[error("named input '{name}' returned a value incompatible with {expected:?}")]
244 NamedInputValueType { name: String, expected: ValueType },
245 #[error("trade slot '{slot}' has inconsistent engine state: {reason}")]
246 TradeSlot { slot: String, reason: String },
247 #[error("configured command '{command_id}' received an incompatible committed effect")]
248 IncompatibleCommandEffect { command_id: String },
249 #[error("configured strategy evaluation failed: {0}")]
250 Evaluation(#[from] qs_strategy::EvaluationError),
251 #[error("configured decision mapping failed: {0}")]
252 Decision(#[from] StrategyDomainError),
253 #[error("configured note mapping failed: {0}")]
254 Journal(#[from] StrategyJournalError),
255 #[error("configured output integer cannot be represented exactly as f64")]
256 IntegerOutputPrecision,
257}
258
259#[derive(Debug, Clone)]
260struct CommandRoute {
261 action: ConfiguredActionKind,
262 slot: String,
263 fact_seen: bool,
264 terminal: Option<CommandTerminalStatus>,
265}
266
267pub(crate) struct ConfiguredBoundaryOutput {
268 pub decision: Option<StrategyDecisionDraft>,
269 pub journal: Vec<StrategyJournalDraft>,
270 pub commands: Vec<ConfiguredCommand>,
271}
272
273pub struct BacktestConfiguredStrategyAdapter {
275 strategy: ConfiguredStrategy,
276 descriptor: StrategyDescriptor,
277 requirements: StrategyRequirements,
278 bindings: ConfiguredHistoricalBindings,
279 command_routes: BTreeMap<String, CommandRoute>,
280}
281
282impl BacktestConfiguredStrategyAdapter {
283 pub fn new(
284 strategy: ConfiguredStrategy,
285 descriptor: StrategyDescriptor,
286 bindings: ConfiguredHistoricalBindings,
287 decision_latency_ms: u64,
288 ) -> Result<Self, ConfiguredStrategyAdapterBuildError> {
289 validate_bindings(&strategy, &bindings)?;
290 let series = bindings
291 .sources
292 .iter()
293 .map(|binding| binding.series.requirement().clone())
294 .collect::<Vec<_>>();
295 let mut instruments = Vec::new();
296 for requirement in &series {
297 if !instruments
298 .iter()
299 .any(|symbol| symbol == requirement.symbol())
300 {
301 instruments.push(requirement.symbol().to_owned());
302 }
303 }
304 let needs_feedback = strategy.input_requirements().needs_command_feedback;
305 let requirements = StrategyRequirements::new(
306 instruments,
307 series,
308 decision_latency_ms,
309 true,
310 needs_feedback,
311 )?;
312 Ok(Self {
313 strategy,
314 descriptor,
315 requirements,
316 bindings,
317 command_routes: BTreeMap::new(),
318 })
319 }
320
321 pub fn descriptor(&self) -> &StrategyDescriptor {
322 &self.descriptor
323 }
324
325 pub fn requirements(&self) -> &StrategyRequirements {
326 &self.requirements
327 }
328
329 pub fn configured_requirements(&self) -> &ConfiguredStrategyRequirements {
330 self.strategy.input_requirements()
331 }
332
333 pub fn source_bindings(&self) -> &[ConfiguredSourceBinding] {
334 &self.bindings.sources
335 }
336
337 pub fn series_specs(&self) -> impl ExactSizeIterator<Item = &BarSeriesSpec> {
338 self.bindings.sources.iter().map(|binding| &binding.series)
339 }
340
341 pub fn configured_strategy(&self) -> &ConfiguredStrategy {
342 &self.strategy
343 }
344
345 pub fn into_configured_strategy(self) -> ConfiguredStrategy {
346 self.strategy
347 }
348
349 pub fn preflight(
350 &self,
351 retention: StrategyRetentionLimits,
352 research: StrategyResearchLimits,
353 ) -> Result<(), ConfiguredStrategyAdapterPreflightError> {
354 if retention.max_reason_bytes() < MAX_TEXT_BYTES {
355 return Err(
356 ConfiguredStrategyAdapterPreflightError::DecisionReasonCapacity {
357 actual: retention.max_reason_bytes(),
358 required: MAX_TEXT_BYTES,
359 },
360 );
361 }
362 if retention.max_signals_per_callback() < MAX_OUTPUT_COMMANDS {
363 return Err(ConfiguredStrategyAdapterPreflightError::SignalCapacity {
364 actual: retention.max_signals_per_callback(),
365 required: MAX_OUTPUT_COMMANDS,
366 });
367 }
368 if research.max_journal_per_callback() < MAX_OUTPUT_NOTES {
369 return Err(
370 ConfiguredStrategyAdapterPreflightError::JournalCallbackCapacity {
371 actual: research.max_journal_per_callback(),
372 required: MAX_OUTPUT_NOTES,
373 },
374 );
375 }
376 if research.max_reason_bytes() < MAX_TEXT_BYTES {
377 return Err(
378 ConfiguredStrategyAdapterPreflightError::JournalReasonCapacity {
379 actual: research.max_reason_bytes(),
380 required: MAX_TEXT_BYTES,
381 },
382 );
383 }
384 if research.max_values_per_record() < MAX_NAMED_VALUES {
385 return Err(
386 ConfiguredStrategyAdapterPreflightError::JournalValueCapacity {
387 actual: research.max_values_per_record(),
388 required: MAX_NAMED_VALUES,
389 },
390 );
391 }
392 if research.max_value_key_bytes() < MAX_ID_BYTES {
393 return Err(
394 ConfiguredStrategyAdapterPreflightError::JournalKeyCapacity {
395 actual: research.max_value_key_bytes(),
396 required: MAX_ID_BYTES,
397 },
398 );
399 }
400 if super::MAX_TRADE_ID_BYTES < MAX_GENERATED_ID_BYTES {
401 return Err(
402 ConfiguredStrategyAdapterPreflightError::TradeIdentityCapacity {
403 actual: super::MAX_TRADE_ID_BYTES,
404 required: MAX_GENERATED_ID_BYTES,
405 },
406 );
407 }
408 Ok(())
409 }
410
411 #[allow(clippy::too_many_arguments)]
412 pub(crate) fn evaluate_boundary(
413 &mut self,
414 observed_through: NaiveDateTime,
415 ready: bool,
416 closed_bars: &[ClosedBar],
417 observations: &[StrategyObservation],
418 series: &dyn HistoricalSeriesView,
419 observation_history: &dyn HistoricalObservationView,
420 engine: &TradeEngine,
421 feedback_events: &[StrategyFeedbackEvent],
422 retention: StrategyRetentionLimits,
423 research: StrategyResearchLimits,
424 ) -> Result<ConfiguredBoundaryOutput, ConfiguredStrategyAdapterError> {
425 let feedback = self.project_feedback(feedback_events)?;
426 let input = StrategyInput {
427 time: observed_through,
428 ready,
429 completed_bars: self.project_bars(observed_through, closed_bars)?,
430 values: self.project_named_inputs(NamedInputProjectionContext {
431 observed_through,
432 closed_bars,
433 observations,
434 series,
435 observation_history,
436 })?,
437 trade_slots: self.project_trade_slots(engine)?,
438 feedback,
439 };
440 let output = self.strategy.evaluate(&input)?;
441 for command in &output.commands {
442 self.command_routes.insert(
443 command.command_id.clone(),
444 CommandRoute {
445 action: command.action_kind,
446 slot: command.trade_slot.clone(),
447 fact_seen: false,
448 terminal: None,
449 },
450 );
451 }
452 let emitted_signals = output
453 .commands
454 .iter()
455 .map(|command| command.signal.clone())
456 .collect::<Vec<_>>();
457 let decision = output
458 .decision
459 .map(|decision| map_decision(decision, emitted_signals, retention))
460 .transpose()?;
461 let journal = output
462 .notes
463 .into_iter()
464 .map(|note| map_note(note, self.strategy.primary_symbol(), research))
465 .collect::<Result<Vec<_>, _>>()?;
466 Ok(ConfiguredBoundaryOutput {
467 decision,
468 journal,
469 commands: output.commands,
470 })
471 }
472
473 fn project_bars(
474 &self,
475 observed_through: NaiveDateTime,
476 closed_bars: &[ClosedBar],
477 ) -> Result<Vec<qs_strategy::CompletedBarUpdate>, ConfiguredStrategyAdapterError> {
478 self.strategy
479 .input_requirements()
480 .completed_bars
481 .iter()
482 .filter_map(|requirement| {
483 let binding = self
484 .bindings
485 .sources
486 .iter()
487 .find(|binding| binding.source == requirement.source)
488 .expect("bindings were validated at construction");
489 let mut matching = closed_bars
490 .iter()
491 .filter(|bar| bar.series_id() == binding.series_id());
492 let bar = matching.next()?;
493 Some(if matching.next().is_some() {
494 Err(ConfiguredStrategyAdapterError::DuplicateSourceUpdate {
495 source_id: requirement.source.clone(),
496 timestamp: observed_through,
497 })
498 } else {
499 Ok(qs_strategy::CompletedBarUpdate {
500 source: requirement.source.clone(),
501 bar: qs_strategy::CompletedBar {
502 open: bar.open(),
503 high: bar.high(),
504 low: bar.low(),
505 close: bar.close(),
506 volume: match self.bindings.volume {
507 HistoricalVolumeProjection::TickCountExact => {
508 if bar.tick_count() > MAX_EXACT_F64_INTEGER {
509 return Some(Err(
510 ConfiguredStrategyAdapterError::TickCountNotExactlyRepresentable {
511 tick_count: bar.tick_count(),
512 },
513 ));
514 }
515 bar.tick_count() as f64
516 }
517 },
518 },
519 })
520 })
521 })
522 .collect()
523 }
524
525 fn project_named_inputs(
526 &self,
527 context: NamedInputProjectionContext<'_>,
528 ) -> Result<Vec<NamedValue>, ConfiguredStrategyAdapterError> {
529 self.strategy
530 .input_requirements()
531 .named_inputs
532 .iter()
533 .map(|requirement| {
534 let binding = self
535 .bindings
536 .named_inputs
537 .iter()
538 .find(|binding| binding.name == requirement.name)
539 .expect("named input bindings were validated at construction");
540 let projected = binding.projector.project(context).map_err(|source| {
541 ConfiguredStrategyAdapterError::NamedInput {
542 name: requirement.name.clone(),
543 source,
544 }
545 })?;
546 if !value_matches_type(&projected.value, requirement.value_type) {
547 return Err(ConfiguredStrategyAdapterError::NamedInputValueType {
548 name: requirement.name.clone(),
549 expected: requirement.value_type,
550 });
551 }
552 Ok(NamedValue {
553 name: requirement.name.clone(),
554 value: projected.value,
555 updated: projected.updated,
556 })
557 })
558 .collect()
559 }
560
561 fn project_trade_slots(
562 &self,
563 engine: &TradeEngine,
564 ) -> Result<Vec<TradeSlotFacts>, ConfiguredStrategyAdapterError> {
565 self.strategy
566 .input_requirements()
567 .trade_slots
568 .iter()
569 .map(|slot| {
570 let state = self
571 .strategy
572 .trade_id_for_slot(slot)
573 .and_then(|trade_id| engine.manager.id_by_trade_id(trade_id))
574 .and_then(|position_id| engine.get_position(&position_id))
575 .map(|position| match position.data.status {
576 PositionStatus::Pending => TradeSlotState::Pending {
577 side: position.data.side,
578 requested_price: position.data.pending_price,
579 stoploss: position.current_stoploss(),
580 },
581 PositionStatus::Open => TradeSlotState::Open {
582 side: position.data.side,
583 entry_price: position.data.average_entry(),
584 remaining_size: position.data.remaining_size(),
585 stoploss: position.current_stoploss(),
586 },
587 PositionStatus::Closed | PositionStatus::Cancelled => {
588 TradeSlotState::Vacant
589 }
590 })
591 .unwrap_or(TradeSlotState::Vacant);
592 Ok(TradeSlotFacts {
593 slot: slot.clone(),
594 state,
595 })
596 })
597 .collect()
598 }
599
600 fn project_feedback(
601 &mut self,
602 events: &[StrategyFeedbackEvent],
603 ) -> Result<Vec<CommandFeedback>, ConfiguredStrategyAdapterError> {
604 project_command_feedback(&mut self.command_routes, events)
605 }
606
607 pub(crate) fn finalize_feedback(
608 &mut self,
609 events: &[StrategyFeedbackEvent],
610 ) -> Result<(), ConfiguredStrategyAdapterError> {
611 let feedback = self.project_feedback(events)?;
612 self.strategy.finalize_command_feedback(&feedback)?;
613 self.command_routes.clear();
614 Ok(())
615 }
616}
617
618fn project_command_feedback(
619 routes: &mut BTreeMap<String, CommandRoute>,
620 events: &[StrategyFeedbackEvent],
621) -> Result<Vec<CommandFeedback>, ConfiguredStrategyAdapterError> {
622 let mut projected = Vec::new();
623 for event in events {
624 let Some(command_id) = event.action_id() else {
625 continue;
626 };
627 let Some(route) = routes.get_mut(command_id) else {
628 continue;
629 };
630 match event {
631 StrategyFeedbackEvent::Effect { effect, .. } => {
632 if let Some(fact) = map_effect(route.action, effect.effect()).map_err(|()| {
633 ConfiguredStrategyAdapterError::IncompatibleCommandEffect {
634 command_id: command_id.to_owned(),
635 }
636 })? {
637 route.fact_seen = true;
638 projected.push(CommandFeedback::Fact {
639 command_id: command_id.to_owned(),
640 fact,
641 });
642 }
643 }
644 StrategyFeedbackEvent::Disposition(disposition) => {
645 let status = match disposition.status {
646 ActionDispositionStatus::Applied => CommandTerminalStatus::Applied,
647 ActionDispositionStatus::Skipped => CommandTerminalStatus::Skipped,
648 ActionDispositionStatus::Rejected => CommandTerminalStatus::Rejected,
649 ActionDispositionStatus::Failed => CommandTerminalStatus::Failed,
650 };
651 route.terminal = Some(status);
652 projected.push(CommandFeedback::Terminal {
653 command_id: command_id.to_owned(),
654 status,
655 reason: disposition.reason.clone(),
656 });
657 }
658 }
659 let completed = route
660 .terminal
661 .is_some_and(|status| status != CommandTerminalStatus::Applied)
662 || (route.terminal == Some(CommandTerminalStatus::Applied) && route.fact_seen);
663 if completed {
664 let completed_route = routes
665 .remove(command_id)
666 .expect("completed route remains registered");
667 if completed_route.action == ConfiguredActionKind::CancelPending
668 && completed_route.terminal == Some(CommandTerminalStatus::Applied)
669 {
670 routes.retain(|_, route| {
671 !(route.action == ConfiguredActionKind::Entry
672 && route.slot == completed_route.slot)
673 });
674 }
675 }
676 }
677 Ok(projected)
678}
679
680fn validate_bindings(
681 strategy: &ConfiguredStrategy,
682 bindings: &ConfiguredHistoricalBindings,
683) -> Result<(), ConfiguredStrategyAdapterBuildError> {
684 let declared = strategy.declared_sources().iter().collect::<BTreeSet<_>>();
685 let mut sources = BTreeSet::new();
686 let mut series = BTreeSet::new();
687 for binding in &bindings.sources {
688 if !declared.contains(&binding.source) {
689 return Err(
690 ConfiguredStrategyAdapterBuildError::UndeclaredSourceBinding {
691 source_id: binding.source.clone(),
692 },
693 );
694 }
695 if !sources.insert(binding.source.clone()) {
696 return Err(
697 ConfiguredStrategyAdapterBuildError::DuplicateSourceBinding {
698 source_id: binding.source.clone(),
699 },
700 );
701 }
702 if !series.insert(binding.series_id().clone()) {
703 return Err(
704 ConfiguredStrategyAdapterBuildError::DuplicateSeriesBinding {
705 series_id: binding.series_id().clone(),
706 },
707 );
708 }
709 let series_symbol = binding.series.requirement().symbol();
710 if series_symbol != strategy.primary_symbol() {
711 return Err(ConfiguredStrategyAdapterBuildError::SourceSymbolMismatch {
712 source_id: binding.source.clone(),
713 primary_symbol: strategy.primary_symbol().to_owned(),
714 series_symbol: series_symbol.to_owned(),
715 });
716 }
717 }
718 for source in strategy.declared_sources() {
719 if !sources.contains(source) {
720 return Err(ConfiguredStrategyAdapterBuildError::MissingSourceBinding {
721 source_id: source.clone(),
722 });
723 }
724 }
725 for requirement in &strategy.input_requirements().completed_bars {
726 let binding = bindings
727 .sources
728 .iter()
729 .find(|binding| binding.source == requirement.source)
730 .expect("every declared source was checked above");
731 if binding.series.retained_bars() < requirement.required_lookback {
732 return Err(
733 ConfiguredStrategyAdapterBuildError::RetentionBelowLookback {
734 source_id: requirement.source.clone(),
735 required: requirement.required_lookback,
736 retained: binding.series.retained_bars(),
737 },
738 );
739 }
740 let warmup = binding.series.requirement().warmup().required_bars();
741 if warmup < requirement.required_lookback {
742 return Err(ConfiguredStrategyAdapterBuildError::WarmupBelowLookback {
743 source_id: requirement.source.clone(),
744 required: requirement.required_lookback,
745 warmup,
746 });
747 }
748 }
749 let mut names = BTreeSet::new();
750 for binding in &bindings.named_inputs {
751 if !names.insert(binding.name.clone()) {
752 return Err(
753 ConfiguredStrategyAdapterBuildError::DuplicateNamedInputProjector {
754 name: binding.name.clone(),
755 },
756 );
757 }
758 let Some(requirement) = strategy
759 .input_requirements()
760 .named_inputs
761 .iter()
762 .find(|requirement| requirement.name == binding.name)
763 else {
764 return Err(
765 ConfiguredStrategyAdapterBuildError::UndeclaredNamedInputProjector {
766 name: binding.name.clone(),
767 },
768 );
769 };
770 let actual = binding.output_type();
771 if actual != requirement.value_type {
772 return Err(
773 ConfiguredStrategyAdapterBuildError::NamedInputTypeMismatch {
774 name: binding.name.clone(),
775 expected: requirement.value_type,
776 actual,
777 },
778 );
779 }
780 }
781 for requirement in &strategy.input_requirements().named_inputs {
782 if !names.contains(&requirement.name) {
783 return Err(
784 ConfiguredStrategyAdapterBuildError::MissingNamedInputProjector {
785 name: requirement.name.clone(),
786 },
787 );
788 }
789 }
790 Ok(())
791}
792
793fn value_matches_type(value: &Value, expected: ValueType) -> bool {
794 if value.is_missing() {
795 return expected.optional && value.scalar_type() == expected.scalar;
796 }
797 if value.scalar_type() != expected.scalar {
798 return false;
799 }
800 match value {
801 Value::Number(value) | Value::Price(value) => value.is_finite(),
802 Value::Text(value) => !value.is_empty() && value.len() <= MAX_TEXT_BYTES,
803 _ => true,
804 }
805}
806
807fn map_effect(action: ConfiguredActionKind, effect: &Effect) -> Result<Option<CommandFact>, ()> {
808 let mapped = match effect {
809 Effect::PositionOpened { .. } => {
810 Some((ConfiguredActionKind::Entry, CommandFact::EntryFilled))
811 }
812 Effect::PositionClosed { .. } => {
813 Some((ConfiguredActionKind::Close, CommandFact::PositionClosed))
814 }
815 Effect::PartialClose { .. } => Some((
816 ConfiguredActionKind::ClosePartial,
817 CommandFact::PositionReduced,
818 )),
819 Effect::StoplossModified { .. } => match action {
820 ConfiguredActionKind::MoveStoplossToEntry | ConfiguredActionKind::ModifyStoploss => {
821 return Ok(Some(CommandFact::StoplossModified));
822 }
823 _ => return Err(()),
824 },
825 Effect::OrderCancelled { .. } => Some((
826 ConfiguredActionKind::CancelPending,
827 CommandFact::PendingCancelled,
828 )),
829 Effect::OrderPlaced { .. }
830 | Effect::StoplossRemoved { .. }
831 | Effect::ScaledIn { .. }
832 | Effect::RuleTriggered { .. } => None,
833 };
834 match mapped {
835 Some((expected, fact)) if expected == action => Ok(Some(fact)),
836 Some(_) => Err(()),
837 None => Ok(None),
838 }
839}
840
841fn map_decision(
842 decision: qs_strategy::Decision,
843 emitted_signals: Vec<qs_core::RawSignal>,
844 limits: StrategyRetentionLimits,
845) -> Result<StrategyDecisionDraft, StrategyDomainError> {
846 let kind = match decision.kind {
847 DecisionKind::Entry => StrategyDecisionKind::Entry,
848 DecisionKind::Management => StrategyDecisionKind::Management,
849 DecisionKind::Exit => StrategyDecisionKind::Exit,
850 DecisionKind::Observation => StrategyDecisionKind::Annotation,
851 };
852 StrategyDecisionDraft::new(
853 kind,
854 decision.reason,
855 decision.related_trade.map(|trade| trade.trade_id),
856 emitted_signals,
857 limits,
858 )
859}
860
861fn map_note(
862 note: qs_strategy::Note,
863 symbol: &str,
864 limits: StrategyResearchLimits,
865) -> Result<StrategyJournalDraft, ConfiguredStrategyAdapterError> {
866 let kind = match note.kind {
867 NoteKind::Observation | NoteKind::Risk => JournalKind::DecisionContext,
868 NoteKind::Execution | NoteKind::Lifecycle => JournalKind::OutcomeReview,
869 };
870 let mut values = BTreeMap::new();
871 for output in note.values {
872 let value = match output.value {
873 OutputScalar::Integer(value) => {
874 if value.unsigned_abs() > MAX_EXACT_F64_INTEGER {
875 return Err(ConfiguredStrategyAdapterError::IntegerOutputPrecision);
876 }
877 value as f64
878 }
879 OutputScalar::Number(value) | OutputScalar::Price(value) => value,
880 };
881 values.insert(output.name, value);
882 }
883 Ok(StrategyJournalDraft::new(
884 kind,
885 symbol,
886 note.related_trade.map(|trade| trade.trade_id),
887 note.reason,
888 None,
889 values,
890 limits,
891 )?)
892}
893
894#[cfg(test)]
895mod tests {
896 use super::*;
897 use crate::ledger::ActionDisposition;
898 use qs_core::types::FutureEffect;
899
900 #[test]
901 fn applied_terminal_before_effect_completes_command_correlation() {
902 let command_id = "opaque-command".to_owned();
903 let mut routes = BTreeMap::from([(
904 command_id.clone(),
905 CommandRoute {
906 action: ConfiguredActionKind::Entry,
907 slot: "primary".into(),
908 fact_seen: false,
909 terminal: None,
910 },
911 )]);
912 let terminal = project_command_feedback(
913 &mut routes,
914 &[StrategyFeedbackEvent::Disposition(
915 ActionDisposition::applied(command_id.clone()),
916 )],
917 )
918 .unwrap();
919
920 assert_eq!(
921 terminal,
922 vec![CommandFeedback::Terminal {
923 command_id: command_id.clone(),
924 status: CommandTerminalStatus::Applied,
925 reason: None,
926 }]
927 );
928 assert!(routes.contains_key(&command_id));
929
930 let fact = project_command_feedback(
931 &mut routes,
932 &[StrategyFeedbackEvent::Effect {
933 action_id: Some(command_id.clone()),
934 effect: FutureEffect::plain(Effect::PositionOpened {
935 id: "position-1".into(),
936 }),
937 }],
938 )
939 .unwrap();
940
941 assert_eq!(
942 fact,
943 vec![CommandFeedback::Fact {
944 command_id,
945 fact: CommandFact::EntryFilled,
946 }]
947 );
948 assert!(routes.is_empty());
949 }
950}