Skip to main content

backtest_server/
instrument_catalog.rs

1use std::collections::{BTreeMap, BTreeSet};
2use std::fs;
3use std::sync::Arc;
4
5use chrono::{DateTime, NaiveDateTime, Utc};
6use qs_backtest::{
7    ReplayInstrumentArtifact, ReplayInstrumentManifest, guarded_instrument_spec,
8    resolve_legacy_economics,
9};
10use qs_instruments::{
11    AssetId, AssetKind, AssetSpec, CatalogDocument, Decimal, DecimalGrid, EconomicsModelId,
12    EffectiveInterval, InstrumentAlias, InstrumentAssets, InstrumentCatalogSnapshot,
13    InstrumentEconomics, InstrumentId, InstrumentResolutionContext, InstrumentResolutionError,
14    InstrumentSelector, InstrumentSpec, ListingStatus, ListingVenueId, MarketDataSourceId,
15    MarketKind, PriceRules, QuantityRules, QuantityUnit, StoredSeriesBinding,
16};
17use qs_symbols::{SymbolCurrencyMetadata, SymbolRegistry, SymbolSpec};
18
19use crate::config::{InstrumentsSection, LinearInstrumentConfig};
20use crate::error::{BacktestServerError, Result};
21use crate::rpc_types::InstrumentExclusionReasonMsg;
22use qs_market_loader::SymbolPartitionResolver;
23
24const CATALOG_SCHEMA_VERSION: u32 = 1;
25const COMPATIBILITY_LISTING_VENUE: &str = "repository-default";
26const REGISTRY_CATALOG_VERSION: &str = "symbol-registry-1";
27const SPEC_VALID_FROM: &str = "1970-01-01T00:00:00Z";
28
29#[derive(Clone)]
30pub struct InstrumentDomain {
31    snapshot: Arc<InstrumentCatalogSnapshot>,
32    default_listing_venue: Option<ListingVenueId>,
33    data_source: MarketDataSourceId,
34    source_symbols: BTreeMap<String, String>,
35}
36
37impl InstrumentDomain {
38    pub fn load(config: &InstrumentsSection, registry: &SymbolRegistry) -> Result<Self> {
39        let configured_listing_venue = config
40            .default_listing_venue
41            .as_deref()
42            .map(ListingVenueId::new)
43            .transpose()
44            .map_err(|error| BacktestServerError::Config(error.to_string()))?;
45        let data_source = MarketDataSourceId::new(&config.market_data_source)
46            .map_err(|error| BacktestServerError::Config(error.to_string()))?;
47        if config.catalog_path.is_some() && !config.linear_instruments.is_empty() {
48            return Err(BacktestServerError::Config(
49                "catalog_path and linear_instruments are alternative specification authorities"
50                    .into(),
51            ));
52        }
53        let mut source_symbols = BTreeMap::new();
54        let mut physical = BTreeSet::new();
55        for (symbol, source) in &config.source_symbols {
56            let canonical = registry.normalize(symbol).ok_or_else(|| {
57                BacktestServerError::Config(format!("unknown source binding symbol '{symbol}'"))
58            })?;
59            if source.is_empty()
60                || source.contains(['/', '\\'])
61                || source == "."
62                || source == ".."
63                || registry
64                    .normalize(source)
65                    .is_some_and(|target| target != canonical)
66                || !physical.insert(source.to_ascii_lowercase())
67                || source_symbols
68                    .insert(canonical.into(), source.clone())
69                    .is_some()
70            {
71                return Err(BacktestServerError::Config(format!(
72                    "invalid or duplicate source binding for '{symbol}'"
73                )));
74            }
75        }
76        let (snapshot, default_listing_venue) = match &config.catalog_path {
77            Some(path) => {
78                let content = fs::read_to_string(path).map_err(|error| {
79                    BacktestServerError::Config(format!(
80                        "failed to read instrument catalog '{path}': {error}"
81                    ))
82                })?;
83                let document = toml::from_str::<CatalogDocument>(&content).map_err(|error| {
84                    BacktestServerError::Config(format!(
85                        "failed to parse instrument catalog '{path}': {error}"
86                    ))
87                })?;
88                let snapshot = InstrumentCatalogSnapshot::compile(document).map_err(|error| {
89                    BacktestServerError::Config(format!(
90                        "invalid instrument catalog '{path}': {error}"
91                    ))
92                })?;
93                (snapshot, configured_listing_venue)
94            }
95            None => {
96                let listing_venue = configured_listing_venue.unwrap_or_else(|| {
97                    ListingVenueId::new(COMPATIBILITY_LISTING_VENUE)
98                        .expect("valid built-in compatibility listing venue")
99                });
100                (
101                    compatibility_snapshot(registry, &listing_venue, &config.linear_instruments)?,
102                    Some(listing_venue),
103                )
104            }
105        };
106        Ok(Self {
107            snapshot: Arc::new(snapshot),
108            default_listing_venue,
109            data_source,
110            source_symbols,
111        })
112    }
113
114    pub fn compatibility(registry: &SymbolRegistry) -> Result<Self> {
115        Self::load(&InstrumentsSection::default(), registry)
116    }
117
118    pub fn snapshot(&self) -> &InstrumentCatalogSnapshot {
119        &self.snapshot
120    }
121
122    pub fn data_source(&self) -> &MarketDataSourceId {
123        &self.data_source
124    }
125
126    pub fn symbol_resolver<'a>(
127        &'a self,
128        registry: &'a SymbolRegistry,
129    ) -> SymbolPartitionResolver<'a> {
130        SymbolPartitionResolver::with_source_symbols(registry, &self.source_symbols)
131    }
132
133    pub fn resolve_manifest(
134        &self,
135        symbols: &[String],
136        at: NaiveDateTime,
137        through: Option<NaiveDateTime>,
138    ) -> Result<ReplayInstrumentManifest> {
139        let allowed_instruments = self
140            .snapshot
141            .instrument_ids()
142            .cloned()
143            .collect::<BTreeSet<_>>();
144        let context = InstrumentResolutionContext {
145            allowed_instruments,
146            default_listing_venue: self.default_listing_venue.clone(),
147            default_market_kind: None,
148        };
149        let at = at.and_utc();
150        let through = through.map(|value| value.and_utc());
151        let mut instruments = BTreeMap::new();
152        for symbol in symbols {
153            let selector = InstrumentSelector::Alias {
154                alias: InstrumentAlias::new(symbol).map_err(|error| {
155                    BacktestServerError::InvalidRequest(format!(
156                        "invalid instrument selector for '{symbol}': {error}"
157                    ))
158                })?,
159                listing_venue: None,
160                market_kind: None,
161            };
162            let resolved = self
163                .snapshot
164                .resolve(&selector, &context, at)
165                .map_err(|error| resolution_error(symbol, error, at))?;
166            validate_replay_spec(symbol, &resolved.spec)?;
167            if let Some(through) = through {
168                let end = self
169                    .snapshot
170                    .resolve(&selector, &context, through)
171                    .map_err(|error| resolution_error(symbol, error, through))?;
172                if end.reference != resolved.reference {
173                    return Err(BacktestServerError::InstrumentUnavailable {
174                        symbol: symbol.clone(),
175                        reason: InstrumentExclusionReasonMsg::UnsupportedEconomics,
176                        details:
177                            "instrument changes specification during the requested replay range"
178                                .into(),
179                    });
180                }
181            }
182            instruments.insert(
183                symbol.clone(),
184                ReplayInstrumentArtifact {
185                    resolved: resolved.reference,
186                    spec: resolved.spec.as_ref().clone(),
187                },
188            );
189        }
190        Ok(ReplayInstrumentManifest {
191            instruments,
192            stored_series: Vec::new(),
193        })
194    }
195
196    pub fn attach_stored_series(
197        &self,
198        manifest: &mut ReplayInstrumentManifest,
199        coordinates: impl IntoIterator<Item = (String, String, String)>,
200    ) -> Result<()> {
201        let mut bindings = Vec::new();
202        for (symbol, source_partition, source_symbol) in coordinates {
203            let artifact = manifest.instruments.get(&symbol).ok_or_else(|| {
204                BacktestServerError::InvalidRequest(format!(
205                    "stored series for '{symbol}' has no resolved instrument"
206                ))
207            })?;
208            bindings.push(StoredSeriesBinding {
209                data_source: self.data_source.clone(),
210                source_partition,
211                source_symbol,
212                instrument: artifact.resolved.clone(),
213                effective: artifact.spec.effective,
214            });
215        }
216        bindings.sort_by(|left, right| {
217            left.source_partition
218                .cmp(&right.source_partition)
219                .then(left.source_symbol.cmp(&right.source_symbol))
220                .then(left.instrument.instrument.cmp(&right.instrument.instrument))
221        });
222        bindings.dedup();
223        manifest.stored_series = bindings;
224        Ok(())
225    }
226}
227
228fn resolution_error(
229    symbol: &str,
230    error: InstrumentResolutionError,
231    at: DateTime<Utc>,
232) -> BacktestServerError {
233    let details = format!("cannot resolve instrument '{symbol}': {error}");
234    match error {
235        InstrumentResolutionError::Inactive { .. } => BacktestServerError::InactiveInstrument {
236            symbol: symbol.into(),
237            at,
238            details,
239        },
240        InstrumentResolutionError::Ambiguous { .. } => BacktestServerError::InstrumentUnavailable {
241            symbol: symbol.into(),
242            reason: InstrumentExclusionReasonMsg::AmbiguousMapping,
243            details,
244        },
245        _ => BacktestServerError::InstrumentUnavailable {
246            symbol: symbol.into(),
247            reason: InstrumentExclusionReasonMsg::UnknownInstrument,
248            details,
249        },
250    }
251}
252
253fn compatibility_snapshot(
254    registry: &SymbolRegistry,
255    listing_venue: &ListingVenueId,
256    linear_instruments: &[LinearInstrumentConfig],
257) -> Result<InstrumentCatalogSnapshot> {
258    let effective = EffectiveInterval::new(
259        SPEC_VALID_FROM
260            .parse::<DateTime<Utc>>()
261            .expect("valid built-in instrument epoch"),
262        None,
263    )
264    .expect("valid built-in instrument interval");
265    let mut assets = BTreeMap::<AssetId, AssetSpec>::new();
266    let mut instruments = Vec::new();
267    for (symbol, currencies) in registry.entries() {
268        let Ok(economics) = resolve_legacy_economics(symbol) else {
269            continue;
270        };
271        register_assets(&mut assets, symbol, currencies)?;
272        let instrument = InstrumentId::new(
273            listing_venue.clone(),
274            market_kind(symbol)?,
275            symbol.canonical.parse().map_err(|error| {
276                BacktestServerError::Config(format!(
277                    "invalid listing ID for '{}': {error}",
278                    symbol.canonical
279                ))
280            })?,
281        );
282        instruments.push(
283            guarded_instrument_spec(symbol, currencies, economics, instrument, effective)
284                .map_err(|error| BacktestServerError::Config(error.to_string()))?,
285        );
286    }
287    let mut configured = BTreeSet::new();
288    for rules in linear_instruments {
289        let canonical = registry.normalize(&rules.symbol).ok_or_else(|| {
290            BacktestServerError::Config(format!("unknown linear instrument '{}'", rules.symbol))
291        })?;
292        if !configured.insert(canonical.to_owned()) {
293            return Err(BacktestServerError::Config(format!(
294                "duplicate linear instrument '{canonical}'"
295            )));
296        }
297        let symbol = registry
298            .spec(canonical)
299            .expect("normalized registered symbol");
300        let currencies = registry.currency_metadata(canonical).ok_or_else(|| {
301            BacktestServerError::Config(format!("missing currencies for '{canonical}'"))
302        })?;
303        register_assets(&mut assets, symbol, currencies)?;
304        let settlement = AssetId::new(&currencies.pnl_currency)
305            .map_err(|error| BacktestServerError::Config(error.to_string()))?;
306        let spec = InstrumentSpec {
307            revision: "1.0.0"
308                .parse()
309                .map_err(|error: qs_instruments::IdentifierError| {
310                    BacktestServerError::Config(error.to_string())
311                })?,
312            instrument: InstrumentId::new(
313                listing_venue.clone(),
314                MarketKind::new(MarketKind::LINEAR_EXPOSURE)
315                    .map_err(|error| BacktestServerError::Config(error.to_string()))?,
316                canonical
317                    .parse()
318                    .map_err(|error: qs_instruments::IdentifierError| {
319                        BacktestServerError::Config(error.to_string())
320                    })?,
321            ),
322            effective,
323            status: ListingStatus::Trading,
324            assets: InstrumentAssets {
325                base: currencies
326                    .base_currency
327                    .as_deref()
328                    .map(AssetId::new)
329                    .transpose()
330                    .map_err(|error| BacktestServerError::Config(error.to_string()))?,
331                quote: currencies
332                    .quote_currency
333                    .as_deref()
334                    .map(AssetId::new)
335                    .transpose()
336                    .map_err(|error| BacktestServerError::Config(error.to_string()))?,
337                settlement: settlement.clone(),
338                fee_assets: BTreeSet::new(),
339            },
340            price: PriceRules {
341                grid: DecimalGrid::new(Decimal::ZERO, rules.price_step),
342                display_scale: rules.display_scale,
343            },
344            quantity: QuantityRules {
345                grid: DecimalGrid::new(Decimal::ZERO, rules.quantity_step),
346                minimum: rules.minimum,
347                maximum: Some(rules.maximum),
348                storage_scale: rules
349                    .quantity_step
350                    .get()
351                    .scale()
352                    .max(rules.minimum.get().scale())
353                    .max(rules.maximum.get().scale()),
354            },
355            notional: None,
356            economics: InstrumentEconomics {
357                pnl_model: EconomicsModelId::new(EconomicsModelId::CFD_QUOTE_LINEAR_V1)
358                    .map_err(|error| BacktestServerError::Config(error.to_string()))?,
359                quantity_unit: QuantityUnit::StandardLot,
360                contract_multiplier: rules.contract_multiplier,
361                settlement_asset: settlement,
362                fee_model: None,
363                funding_model: None,
364                margin_model: None,
365            },
366            aliases: BTreeSet::from([InstrumentAlias::new(canonical)
367                .map_err(|error| BacktestServerError::Config(error.to_string()))?]),
368        };
369        spec.validate()
370            .map_err(|error| BacktestServerError::Config(error.to_string()))?;
371        instruments.retain(|existing| {
372            !existing
373                .aliases
374                .contains(&InstrumentAlias::new(canonical).expect("validated alias"))
375        });
376        instruments.push(spec);
377    }
378    InstrumentCatalogSnapshot::compile(CatalogDocument {
379        schema_version: CATALOG_SCHEMA_VERSION,
380        version: REGISTRY_CATALOG_VERSION.into(),
381        assets: assets.into_values().collect(),
382        instruments,
383    })
384    .map_err(|error| BacktestServerError::Config(error.to_string()))
385}
386
387fn register_assets(
388    assets: &mut BTreeMap<AssetId, AssetSpec>,
389    symbol: &SymbolSpec,
390    currencies: &SymbolCurrencyMetadata,
391) -> Result<()> {
392    let values = [
393        currencies.base_currency.as_deref(),
394        currencies.quote_currency.as_deref(),
395        Some(currencies.pnl_currency.as_str()),
396    ];
397    for value in values.into_iter().flatten() {
398        let asset =
399            AssetId::new(value).map_err(|error| BacktestServerError::Config(error.to_string()))?;
400        let kind = match symbol.category.as_str() {
401            "crypto" if currencies.base_currency.as_deref() == Some(value) => AssetKind::Crypto,
402            "metal" | "commodity" if currencies.base_currency.as_deref() == Some(value) => {
403                AssetKind::Commodity
404            }
405            _ => AssetKind::Fiat,
406        };
407        assets.entry(asset.clone()).or_insert(AssetSpec {
408            asset,
409            kind,
410            display_code: value.to_ascii_uppercase(),
411            storage_scale: None,
412        });
413    }
414    Ok(())
415}
416
417fn market_kind(symbol: &SymbolSpec) -> Result<MarketKind> {
418    let kind = match symbol.category.as_str() {
419        "forex" => MarketKind::FX_CFD,
420        "metal" => MarketKind::METAL_CFD,
421        "commodity" => MarketKind::COMMODITY_CFD,
422        "index" => MarketKind::INDEX_CFD,
423        category => {
424            return Err(BacktestServerError::Config(format!(
425                "unsupported registry category '{category}' for {}",
426                symbol.canonical
427            )));
428        }
429    };
430    MarketKind::new(kind).map_err(|error| BacktestServerError::Config(error.to_string()))
431}
432
433fn validate_replay_spec(symbol: &str, spec: &qs_instruments::InstrumentSpec) -> Result<()> {
434    if spec.status != ListingStatus::Trading {
435        return Err(BacktestServerError::InstrumentUnavailable {
436            symbol: symbol.into(),
437            reason: InstrumentExclusionReasonMsg::UnsupportedEconomics,
438            details: "instrument is not in trading status".into(),
439        });
440    }
441    if spec.economics.quantity_unit != QuantityUnit::StandardLot {
442        return Err(BacktestServerError::InstrumentUnavailable {
443            symbol: symbol.into(),
444            reason: InstrumentExclusionReasonMsg::UnsupportedEconomics,
445            details: "unsupported quantity unit".into(),
446        });
447    }
448    let model = spec.economics.pnl_model.as_str();
449    if model != EconomicsModelId::FX_QUOTE_LINEAR_V1
450        && model != EconomicsModelId::CFD_QUOTE_LINEAR_V1
451    {
452        return Err(BacktestServerError::InstrumentUnavailable {
453            symbol: symbol.into(),
454            reason: InstrumentExclusionReasonMsg::UnsupportedEconomics,
455            details: format!("unsupported P&L model: {model}"),
456        });
457    }
458    Ok(())
459}
460
461#[cfg(test)]
462mod tests {
463    use std::collections::BTreeSet;
464    use std::time::{SystemTime, UNIX_EPOCH};
465
466    use qs_instruments::{
467        Decimal, DecimalGrid, InstrumentAssets, InstrumentEconomics, InstrumentSpec, PriceRules,
468        QuantityRules,
469    };
470
471    use super::*;
472
473    fn registry() -> SymbolRegistry {
474        SymbolRegistry::from_toml(
475            r#"
476[[symbol]]
477canonical = "eurusd"
478aliases = ["eur/usd"]
479pip_position = 4
480digits = 5
481category = "forex"
482base_currency = "EUR"
483quote_currency = "USD"
484pnl_currency = "USD"
485lot_base_units = 100000
486lot_step_units = 1000
487
488[[symbol]]
489canonical = "btcusd"
490aliases = ["btc/usd"]
491pip_position = 1
492digits = 2
493category = "crypto"
494base_currency = "BTC"
495quote_currency = "USD"
496pnl_currency = "USD"
497lot_base_units = 100000000
498lot_step_units = 100000
499"#,
500        )
501        .unwrap()
502    }
503
504    #[test]
505    fn compatibility_catalog_uses_owned_listing_namespace_not_data_platform() {
506        let domain = InstrumentDomain::compatibility(&registry()).unwrap();
507        let manifest = domain
508            .resolve_manifest(
509                &["eurusd".into()],
510                "2026-01-01T00:00:00".parse().unwrap(),
511                None,
512            )
513            .unwrap();
514        let instrument = &manifest.instruments["eurusd"].resolved.instrument;
515        assert_eq!(instrument.listing_venue.as_str(), "repository-default");
516        assert_ne!(instrument.listing_venue.as_str(), "ctrader");
517        assert_eq!(instrument.market_kind.as_str(), MarketKind::FX_CFD);
518    }
519
520    #[test]
521    fn unsupported_registry_rows_do_not_enter_the_compatibility_catalog() {
522        let domain = InstrumentDomain::compatibility(&registry()).unwrap();
523        let error = domain
524            .resolve_manifest(
525                &["btcusd".into()],
526                "2026-01-01T00:00:00".parse().unwrap(),
527                None,
528            )
529            .unwrap_err();
530        assert!(error.to_string().contains("cannot resolve instrument"));
531    }
532
533    #[test]
534    fn explicit_catalog_loader_preserves_broker_listing_and_rejects_unknown_fields() {
535        let usd: AssetId = "USD".parse().unwrap();
536        let document = CatalogDocument {
537            schema_version: 1,
538            version: "broker-catalog-1".into(),
539            assets: vec![
540                AssetSpec {
541                    asset: "EUR".parse().unwrap(),
542                    kind: AssetKind::Fiat,
543                    display_code: "EUR".into(),
544                    storage_scale: Some(2),
545                },
546                AssetSpec {
547                    asset: usd.clone(),
548                    kind: AssetKind::Fiat,
549                    display_code: "USD".into(),
550                    storage_scale: Some(2),
551                },
552            ],
553            instruments: vec![InstrumentSpec {
554                revision: "1.0.0".parse().unwrap(),
555                instrument: InstrumentId::new(
556                    "ic-markets".parse().unwrap(),
557                    MarketKind::new(MarketKind::FX_CFD).unwrap(),
558                    "EURUSD".parse().unwrap(),
559                ),
560                effective: EffectiveInterval::new("2026-01-01T00:00:00Z".parse().unwrap(), None)
561                    .unwrap(),
562                status: ListingStatus::Trading,
563                assets: InstrumentAssets {
564                    base: Some("EUR".parse().unwrap()),
565                    quote: Some(usd.clone()),
566                    settlement: usd.clone(),
567                    fee_assets: BTreeSet::new(),
568                },
569                price: PriceRules {
570                    grid: DecimalGrid::new(Decimal::ZERO, "0.00001".parse().unwrap()),
571                    display_scale: 5,
572                },
573                quantity: QuantityRules {
574                    grid: DecimalGrid::new(Decimal::ZERO, "0.01".parse().unwrap()),
575                    minimum: "0.01".parse().unwrap(),
576                    maximum: Some("100".parse().unwrap()),
577                    storage_scale: 2,
578                },
579                notional: None,
580                economics: InstrumentEconomics {
581                    pnl_model: EconomicsModelId::new(EconomicsModelId::FX_QUOTE_LINEAR_V1).unwrap(),
582                    quantity_unit: QuantityUnit::StandardLot,
583                    contract_multiplier: "100000".parse().unwrap(),
584                    settlement_asset: usd,
585                    fee_model: None,
586                    funding_model: None,
587                    margin_model: None,
588                },
589                aliases: BTreeSet::from(["EURUSD".parse().unwrap()]),
590            }],
591        };
592        let unique = SystemTime::now()
593            .duration_since(UNIX_EPOCH)
594            .unwrap()
595            .as_nanos();
596        let path = std::env::temp_dir().join(format!(
597            "qs_instrument_catalog_{}_{}.toml",
598            std::process::id(),
599            unique
600        ));
601        let content = toml::to_string(&document).unwrap();
602        fs::write(&path, &content).unwrap();
603        let config = InstrumentsSection {
604            catalog_path: Some(path.to_string_lossy().into_owned()),
605            default_listing_venue: None,
606            market_data_source: "test-parquet".into(),
607            ..InstrumentsSection::default()
608        };
609
610        let domain = InstrumentDomain::load(&config, &registry()).unwrap();
611        let manifest = domain
612            .resolve_manifest(
613                &["eurusd".into()],
614                "2026-02-01T00:00:00".parse().unwrap(),
615                None,
616            )
617            .unwrap();
618        let instrument = &manifest.instruments["eurusd"].resolved.instrument;
619        assert_eq!(instrument.listing_venue.as_str(), "ic-markets");
620        assert_ne!(instrument.listing_venue.as_str(), "ctrader");
621
622        fs::write(&path, format!("unexpected = true\n{content}")).unwrap();
623        let error = match InstrumentDomain::load(&config, &registry()) {
624            Ok(_) => panic!("catalog with an unknown field must be rejected"),
625            Err(error) => error,
626        };
627        assert!(error.to_string().contains("unknown field"));
628        fs::remove_file(path).unwrap();
629    }
630}