Skip to main content

qs_backtest/strategy/
configured.rs

1//! Historical binding for reusable configured strategies.
2
3use 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/// Historical volume projection used for configured completed bars.
32#[derive(Debug, Clone, Copy, PartialEq, Eq)]
33pub enum HistoricalVolumeProjection {
34    TickCountExact,
35    OptionalTickCount,
36}
37
38/// Complete historical binding for one logical configured source.
39#[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/// Immutable causal values available to one named-input projector.
64#[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/// One typed named-input value and its boundary update provenance.
74#[derive(Debug, Clone, PartialEq)]
75pub struct ProjectedNamedInput {
76    pub value: Value,
77    pub updated: bool,
78}
79
80/// Pure historical projection for one configured named input.
81///
82/// A projector must be movable between threads. It is a deterministic transformation over a borrowed context, so an implementation that could not move holds shared state it has no reason to hold. The bound also keeps one projector implementation usable by both the historical adapter and a later live adapter, which must run inside its own task.
83pub 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/// One explicitly selected causal fact from a completed historical source bar.
93#[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
109/// Projects source-bar timing and explicit gap facts without using host time or callback count.
110pub 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
189/// Retains the latest confirmed high or low swing while preserving its actual confirmation update.
190pub 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
267/// Binding from a configured input name to a historical projector.
268pub 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
290/// Complete caller-owned historical input binding.
291pub 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/// Named-input projector failure.
334#[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/// Binding failures detected before replay starts.
349#[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/// Static output compatibility failure detected before feed polling.
400#[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/// Management-profile routing failure for a configured strategy, detected before feed polling.
421#[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/// Runtime historical projection or configured evaluation failure.
436#[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
482/// Open-position economics a configured boundary may observe.
483///
484/// Excursion is the campaign's extreme as of the start of the boundary's batch, before that batch's own quotes are marked, so a stored bar's close cannot reach the strategy before the bar closes. Initial risk is the basis completed positions report for R normalization.
485pub(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
513/// Historical runtime adapter for one reusable configured strategy instance.
514pub 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    /// Check that every Entry the strategy can emit resolves to a profile the same way replay will select it, and that no selected profile manages a stop the strategy also moves.
602    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
1129/// A profile owns the stop when it replaces the signal stop or attaches a rule that moves the stop while the position is open.
1130fn 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}