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(¤cies.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(®istry()).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(®istry()).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, ®istry()).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, ®istry()) {
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}