Skip to main content

qs_backtest/strategy/
configured.rs

1//! Historical binding for reusable configured strategies.
2
3use 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/// Historical volume projection used for configured completed bars.
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum HistoricalVolumeProjection {
30    TickCountExact,
31}
32
33/// Complete historical binding for one logical configured source.
34#[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/// Immutable causal values available to one named-input projector.
59#[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/// One typed named-input value and its boundary update provenance.
69#[derive(Debug, Clone, PartialEq)]
70pub struct ProjectedNamedInput {
71    pub value: Value,
72    pub updated: bool,
73}
74
75/// Pure historical projection for one configured named input.
76pub trait HistoricalNamedInputProjector {
77    fn output_type(&self) -> ValueType;
78
79    fn project(
80        &self,
81        context: NamedInputProjectionContext<'_>,
82    ) -> Result<ProjectedNamedInput, NamedInputProjectionError>;
83}
84
85/// Binding from a configured input name to a historical projector.
86pub 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
108/// Complete caller-owned historical input binding.
109pub 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/// Named-input projector failure.
142#[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/// Binding failures detected before replay starts.
157#[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/// Static output compatibility failure detected before feed polling.
208#[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/// Runtime historical projection or configured evaluation failure.
229#[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
273/// Historical runtime adapter for one reusable configured strategy instance.
274pub 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}