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