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, EconomicsModelId, EffectiveInterval,
12    InstrumentAlias, InstrumentCatalogSnapshot, InstrumentId, InstrumentResolutionContext,
13    InstrumentSelector, ListingStatus, ListingVenueId, MarketDataSourceId, MarketKind,
14    QuantityUnit, StoredSeriesBinding,
15};
16use qs_symbols::{SymbolCurrencyMetadata, SymbolRegistry, SymbolSpec};
17
18use crate::config::InstrumentsSection;
19use crate::error::{BacktestServerError, Result};
20
21const CATALOG_SCHEMA_VERSION: u32 = 1;
22const COMPATIBILITY_LISTING_VENUE: &str = "repository-default";
23const REGISTRY_CATALOG_VERSION: &str = "symbol-registry-1";
24const SPEC_VALID_FROM: &str = "1970-01-01T00:00:00Z";
25
26#[derive(Clone)]
27pub struct InstrumentDomain {
28    snapshot: Arc<InstrumentCatalogSnapshot>,
29    default_listing_venue: Option<ListingVenueId>,
30    data_source: MarketDataSourceId,
31}
32
33impl InstrumentDomain {
34    pub fn load(config: &InstrumentsSection, registry: &SymbolRegistry) -> Result<Self> {
35        let configured_listing_venue = config
36            .default_listing_venue
37            .as_deref()
38            .map(ListingVenueId::new)
39            .transpose()
40            .map_err(|error| BacktestServerError::Config(error.to_string()))?;
41        let data_source = MarketDataSourceId::new(&config.market_data_source)
42            .map_err(|error| BacktestServerError::Config(error.to_string()))?;
43        let (snapshot, default_listing_venue) = match &config.catalog_path {
44            Some(path) => {
45                let content = fs::read_to_string(path).map_err(|error| {
46                    BacktestServerError::Config(format!(
47                        "failed to read instrument catalog '{path}': {error}"
48                    ))
49                })?;
50                let document = toml::from_str::<CatalogDocument>(&content).map_err(|error| {
51                    BacktestServerError::Config(format!(
52                        "failed to parse instrument catalog '{path}': {error}"
53                    ))
54                })?;
55                let snapshot = InstrumentCatalogSnapshot::compile(document).map_err(|error| {
56                    BacktestServerError::Config(format!(
57                        "invalid instrument catalog '{path}': {error}"
58                    ))
59                })?;
60                (snapshot, configured_listing_venue)
61            }
62            None => {
63                let listing_venue = configured_listing_venue.unwrap_or_else(|| {
64                    ListingVenueId::new(COMPATIBILITY_LISTING_VENUE)
65                        .expect("valid built-in compatibility listing venue")
66                });
67                (
68                    compatibility_snapshot(registry, &listing_venue)?,
69                    Some(listing_venue),
70                )
71            }
72        };
73        Ok(Self {
74            snapshot: Arc::new(snapshot),
75            default_listing_venue,
76            data_source,
77        })
78    }
79
80    pub fn compatibility(registry: &SymbolRegistry) -> Result<Self> {
81        Self::load(&InstrumentsSection::default(), registry)
82    }
83
84    pub fn snapshot(&self) -> &InstrumentCatalogSnapshot {
85        &self.snapshot
86    }
87
88    pub fn data_source(&self) -> &MarketDataSourceId {
89        &self.data_source
90    }
91
92    pub fn resolve_manifest(
93        &self,
94        symbols: &[String],
95        at: NaiveDateTime,
96        through: Option<NaiveDateTime>,
97    ) -> Result<ReplayInstrumentManifest> {
98        let allowed_instruments = self
99            .snapshot
100            .instrument_ids()
101            .cloned()
102            .collect::<BTreeSet<_>>();
103        let context = InstrumentResolutionContext {
104            allowed_instruments,
105            default_listing_venue: self.default_listing_venue.clone(),
106            default_market_kind: None,
107        };
108        let at = at.and_utc();
109        let through = through.map(|value| value.and_utc());
110        let mut instruments = BTreeMap::new();
111        for symbol in symbols {
112            let selector = InstrumentSelector::Alias {
113                alias: InstrumentAlias::new(symbol).map_err(|error| {
114                    BacktestServerError::InvalidRequest(format!(
115                        "invalid instrument selector for '{symbol}': {error}"
116                    ))
117                })?,
118                listing_venue: None,
119                market_kind: None,
120            };
121            let resolved = self
122                .snapshot
123                .resolve(&selector, &context, at)
124                .map_err(|error| {
125                    BacktestServerError::InvalidRequest(format!(
126                        "cannot resolve instrument '{symbol}': {error}"
127                    ))
128                })?;
129            validate_replay_spec(symbol, &resolved.spec)?;
130            if let Some(through) = through {
131                let end = self
132                    .snapshot
133                    .resolve(&selector, &context, through)
134                    .map_err(|error| {
135                        BacktestServerError::InvalidRequest(format!(
136                            "cannot resolve instrument '{symbol}' at replay end: {error}"
137                        ))
138                    })?;
139                if end.reference != resolved.reference {
140                    return Err(BacktestServerError::InvalidRequest(format!(
141                        "instrument '{symbol}' changes specification during the requested replay range"
142                    )));
143                }
144            }
145            instruments.insert(
146                symbol.clone(),
147                ReplayInstrumentArtifact {
148                    resolved: resolved.reference,
149                    spec: resolved.spec.as_ref().clone(),
150                },
151            );
152        }
153        Ok(ReplayInstrumentManifest {
154            instruments,
155            stored_series: Vec::new(),
156        })
157    }
158
159    pub fn attach_stored_series(
160        &self,
161        manifest: &mut ReplayInstrumentManifest,
162        coordinates: impl IntoIterator<Item = (String, String, String)>,
163    ) -> Result<()> {
164        let mut bindings = Vec::new();
165        for (symbol, source_partition, source_symbol) in coordinates {
166            let artifact = manifest.instruments.get(&symbol).ok_or_else(|| {
167                BacktestServerError::InvalidRequest(format!(
168                    "stored series for '{symbol}' has no resolved instrument"
169                ))
170            })?;
171            bindings.push(StoredSeriesBinding {
172                data_source: self.data_source.clone(),
173                source_partition,
174                source_symbol,
175                instrument: artifact.resolved.clone(),
176                effective: artifact.spec.effective,
177            });
178        }
179        bindings.sort_by(|left, right| {
180            left.source_partition
181                .cmp(&right.source_partition)
182                .then(left.source_symbol.cmp(&right.source_symbol))
183                .then(left.instrument.instrument.cmp(&right.instrument.instrument))
184        });
185        bindings.dedup();
186        manifest.stored_series = bindings;
187        Ok(())
188    }
189}
190
191fn compatibility_snapshot(
192    registry: &SymbolRegistry,
193    listing_venue: &ListingVenueId,
194) -> Result<InstrumentCatalogSnapshot> {
195    let effective = EffectiveInterval::new(
196        SPEC_VALID_FROM
197            .parse::<DateTime<Utc>>()
198            .expect("valid built-in instrument epoch"),
199        None,
200    )
201    .expect("valid built-in instrument interval");
202    let mut assets = BTreeMap::<AssetId, AssetSpec>::new();
203    let mut instruments = Vec::new();
204    for (symbol, currencies) in registry.entries() {
205        let Ok(economics) = resolve_legacy_economics(symbol) else {
206            continue;
207        };
208        register_assets(&mut assets, symbol, currencies)?;
209        let instrument = InstrumentId::new(
210            listing_venue.clone(),
211            market_kind(symbol)?,
212            symbol.canonical.parse().map_err(|error| {
213                BacktestServerError::Config(format!(
214                    "invalid listing ID for '{}': {error}",
215                    symbol.canonical
216                ))
217            })?,
218        );
219        instruments.push(
220            guarded_instrument_spec(symbol, currencies, economics, instrument, effective)
221                .map_err(|error| BacktestServerError::Config(error.to_string()))?,
222        );
223    }
224    InstrumentCatalogSnapshot::compile(CatalogDocument {
225        schema_version: CATALOG_SCHEMA_VERSION,
226        version: REGISTRY_CATALOG_VERSION.into(),
227        assets: assets.into_values().collect(),
228        instruments,
229    })
230    .map_err(|error| BacktestServerError::Config(error.to_string()))
231}
232
233fn register_assets(
234    assets: &mut BTreeMap<AssetId, AssetSpec>,
235    symbol: &SymbolSpec,
236    currencies: &SymbolCurrencyMetadata,
237) -> Result<()> {
238    let values = [
239        currencies.base_currency.as_deref(),
240        currencies.quote_currency.as_deref(),
241        Some(currencies.pnl_currency.as_str()),
242    ];
243    for value in values.into_iter().flatten() {
244        let asset =
245            AssetId::new(value).map_err(|error| BacktestServerError::Config(error.to_string()))?;
246        let kind = match symbol.category.as_str() {
247            "metal" | "commodity" if currencies.base_currency.as_deref() == Some(value) => {
248                AssetKind::Commodity
249            }
250            _ => AssetKind::Fiat,
251        };
252        assets.entry(asset.clone()).or_insert(AssetSpec {
253            asset,
254            kind,
255            display_code: value.to_ascii_uppercase(),
256            storage_scale: None,
257        });
258    }
259    Ok(())
260}
261
262fn market_kind(symbol: &SymbolSpec) -> Result<MarketKind> {
263    let kind = match symbol.category.as_str() {
264        "forex" => MarketKind::FX_CFD,
265        "metal" => MarketKind::METAL_CFD,
266        "commodity" => MarketKind::COMMODITY_CFD,
267        "index" => MarketKind::INDEX_CFD,
268        category => {
269            return Err(BacktestServerError::Config(format!(
270                "unsupported registry category '{category}' for {}",
271                symbol.canonical
272            )));
273        }
274    };
275    MarketKind::new(kind).map_err(|error| BacktestServerError::Config(error.to_string()))
276}
277
278fn validate_replay_spec(symbol: &str, spec: &qs_instruments::InstrumentSpec) -> Result<()> {
279    if spec.status != ListingStatus::Trading {
280        return Err(BacktestServerError::InvalidRequest(format!(
281            "instrument '{symbol}' is not in trading status"
282        )));
283    }
284    if spec.economics.quantity_unit != QuantityUnit::StandardLot {
285        return Err(BacktestServerError::InvalidRequest(format!(
286            "unsupported quantity unit for instrument '{symbol}'"
287        )));
288    }
289    let model = spec.economics.pnl_model.as_str();
290    if model != EconomicsModelId::FX_QUOTE_LINEAR_V1
291        && model != EconomicsModelId::CFD_QUOTE_LINEAR_V1
292    {
293        return Err(BacktestServerError::InvalidRequest(format!(
294            "unsupported P&L model for instrument '{symbol}': {model}"
295        )));
296    }
297    Ok(())
298}
299
300#[cfg(test)]
301mod tests {
302    use std::collections::BTreeSet;
303    use std::time::{SystemTime, UNIX_EPOCH};
304
305    use qs_instruments::{
306        Decimal, DecimalGrid, InstrumentAssets, InstrumentEconomics, InstrumentSpec, PriceRules,
307        QuantityRules,
308    };
309
310    use super::*;
311
312    fn registry() -> SymbolRegistry {
313        SymbolRegistry::from_toml(
314            r#"
315[[symbol]]
316canonical = "eurusd"
317aliases = ["eur/usd"]
318pip_position = 4
319digits = 5
320category = "forex"
321base_currency = "EUR"
322quote_currency = "USD"
323pnl_currency = "USD"
324lot_base_units = 100000
325lot_step_units = 1000
326
327[[symbol]]
328canonical = "btcusd"
329aliases = ["btc/usd"]
330pip_position = 1
331digits = 2
332category = "crypto"
333base_currency = "BTC"
334quote_currency = "USD"
335pnl_currency = "USD"
336lot_base_units = 100000000
337lot_step_units = 100000
338"#,
339        )
340        .unwrap()
341    }
342
343    #[test]
344    fn compatibility_catalog_uses_owned_listing_namespace_not_data_platform() {
345        let domain = InstrumentDomain::compatibility(&registry()).unwrap();
346        let manifest = domain
347            .resolve_manifest(
348                &["eurusd".into()],
349                "2026-01-01T00:00:00".parse().unwrap(),
350                None,
351            )
352            .unwrap();
353        let instrument = &manifest.instruments["eurusd"].resolved.instrument;
354        assert_eq!(instrument.listing_venue.as_str(), "repository-default");
355        assert_ne!(instrument.listing_venue.as_str(), "ctrader");
356        assert_eq!(instrument.market_kind.as_str(), MarketKind::FX_CFD);
357    }
358
359    #[test]
360    fn unsupported_registry_rows_do_not_enter_the_compatibility_catalog() {
361        let domain = InstrumentDomain::compatibility(&registry()).unwrap();
362        let error = domain
363            .resolve_manifest(
364                &["btcusd".into()],
365                "2026-01-01T00:00:00".parse().unwrap(),
366                None,
367            )
368            .unwrap_err();
369        assert!(error.to_string().contains("cannot resolve instrument"));
370    }
371
372    #[test]
373    fn explicit_catalog_loader_preserves_broker_listing_and_rejects_unknown_fields() {
374        let usd: AssetId = "USD".parse().unwrap();
375        let document = CatalogDocument {
376            schema_version: 1,
377            version: "broker-catalog-1".into(),
378            assets: vec![
379                AssetSpec {
380                    asset: "EUR".parse().unwrap(),
381                    kind: AssetKind::Fiat,
382                    display_code: "EUR".into(),
383                    storage_scale: Some(2),
384                },
385                AssetSpec {
386                    asset: usd.clone(),
387                    kind: AssetKind::Fiat,
388                    display_code: "USD".into(),
389                    storage_scale: Some(2),
390                },
391            ],
392            instruments: vec![InstrumentSpec {
393                revision: "1.0.0".parse().unwrap(),
394                instrument: InstrumentId::new(
395                    "ic-markets".parse().unwrap(),
396                    MarketKind::new(MarketKind::FX_CFD).unwrap(),
397                    "EURUSD".parse().unwrap(),
398                ),
399                effective: EffectiveInterval::new("2026-01-01T00:00:00Z".parse().unwrap(), None)
400                    .unwrap(),
401                status: ListingStatus::Trading,
402                assets: InstrumentAssets {
403                    base: Some("EUR".parse().unwrap()),
404                    quote: Some(usd.clone()),
405                    settlement: usd.clone(),
406                    fee_assets: BTreeSet::new(),
407                },
408                price: PriceRules {
409                    grid: DecimalGrid::new(Decimal::ZERO, "0.00001".parse().unwrap()),
410                    display_scale: 5,
411                },
412                quantity: QuantityRules {
413                    grid: DecimalGrid::new(Decimal::ZERO, "0.01".parse().unwrap()),
414                    minimum: "0.01".parse().unwrap(),
415                    maximum: Some("100".parse().unwrap()),
416                    storage_scale: 2,
417                },
418                notional: None,
419                economics: InstrumentEconomics {
420                    pnl_model: EconomicsModelId::new(EconomicsModelId::FX_QUOTE_LINEAR_V1).unwrap(),
421                    quantity_unit: QuantityUnit::StandardLot,
422                    contract_multiplier: "100000".parse().unwrap(),
423                    settlement_asset: usd,
424                    fee_model: None,
425                    funding_model: None,
426                    margin_model: None,
427                },
428                aliases: BTreeSet::from(["EURUSD".parse().unwrap()]),
429            }],
430        };
431        let unique = SystemTime::now()
432            .duration_since(UNIX_EPOCH)
433            .unwrap()
434            .as_nanos();
435        let path = std::env::temp_dir().join(format!(
436            "qs_instrument_catalog_{}_{}.toml",
437            std::process::id(),
438            unique
439        ));
440        let content = toml::to_string(&document).unwrap();
441        fs::write(&path, &content).unwrap();
442        let config = InstrumentsSection {
443            catalog_path: Some(path.to_string_lossy().into_owned()),
444            default_listing_venue: None,
445            market_data_source: "test-parquet".into(),
446        };
447
448        let domain = InstrumentDomain::load(&config, &registry()).unwrap();
449        let manifest = domain
450            .resolve_manifest(
451                &["eurusd".into()],
452                "2026-02-01T00:00:00".parse().unwrap(),
453                None,
454            )
455            .unwrap();
456        let instrument = &manifest.instruments["eurusd"].resolved.instrument;
457        assert_eq!(instrument.listing_venue.as_str(), "ic-markets");
458        assert_ne!(instrument.listing_venue.as_str(), "ctrader");
459
460        fs::write(&path, format!("unexpected = true\n{content}")).unwrap();
461        let error = match InstrumentDomain::load(&config, &registry()) {
462            Ok(_) => panic!("catalog with an unknown field must be rejected"),
463            Err(error) => error,
464        };
465        assert!(error.to_string().contains("unknown field"));
466        fs::remove_file(path).unwrap();
467    }
468}