Skip to main content

qs_instruments/
catalog.rs

1use std::collections::{BTreeMap, BTreeSet};
2use std::sync::Arc;
3
4use chrono::{DateTime, Utc};
5use serde::{Deserialize, Serialize};
6
7use crate::{
8    AssetId, AssetSpec, EffectiveInterval, InstrumentAlias, InstrumentId, InstrumentSpec,
9    ListingVenueId, MarketKind, SpecValidationError,
10};
11
12/// Strict authoring document compiled into an immutable catalog snapshot.
13#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
14#[serde(deny_unknown_fields)]
15pub struct CatalogDocument {
16    pub schema_version: u32,
17    pub version: String,
18    pub assets: Vec<AssetSpec>,
19    pub instruments: Vec<InstrumentSpec>,
20}
21
22/// Human-readable identity of one immutable catalog snapshot.
23#[derive(Clone, Debug, Eq, PartialEq, Ord, PartialOrd, Hash, Serialize, Deserialize)]
24#[serde(deny_unknown_fields)]
25pub struct CatalogSnapshotId {
26    pub version: String,
27}
28
29/// Immutable, indexed catalog revision safe to share across in-flight operations.
30#[derive(Clone, Debug)]
31pub struct InstrumentCatalogSnapshot {
32    id: CatalogSnapshotId,
33    assets: BTreeMap<AssetId, AssetSpec>,
34    history: BTreeMap<InstrumentId, Vec<InstrumentSpec>>,
35    aliases: BTreeMap<InstrumentAlias, BTreeSet<InstrumentId>>,
36}
37
38impl InstrumentCatalogSnapshot {
39    pub fn compile(mut document: CatalogDocument) -> Result<Self, CatalogCompileError> {
40        if document.schema_version != 1 {
41            return Err(CatalogCompileError::UnsupportedSchema(
42                document.schema_version,
43            ));
44        }
45        if document.version.trim().is_empty()
46            || document.version.trim() != document.version
47            || document.version.len() > 64
48            || !document.version.is_ascii()
49        {
50            return Err(CatalogCompileError::InvalidVersion);
51        }
52
53        document
54            .assets
55            .sort_by(|left, right| left.asset.cmp(&right.asset));
56        document.instruments.sort_by(|left, right| {
57            left.instrument
58                .cmp(&right.instrument)
59                .then(left.effective.valid_from.cmp(&right.effective.valid_from))
60                .then(left.revision.cmp(&right.revision))
61        });
62
63        let mut assets = BTreeMap::new();
64        for asset in &document.assets {
65            asset.validate()?;
66            if assets.insert(asset.asset.clone(), asset.clone()).is_some() {
67                return Err(CatalogCompileError::DuplicateAsset(asset.asset.clone()));
68            }
69        }
70
71        let mut history: BTreeMap<InstrumentId, Vec<InstrumentSpec>> = BTreeMap::new();
72        let mut aliases: BTreeMap<InstrumentAlias, BTreeSet<InstrumentId>> = BTreeMap::new();
73        for spec in &document.instruments {
74            spec.validate()?;
75            validate_assets(spec, &assets)?;
76            for alias in &spec.aliases {
77                aliases
78                    .entry(alias.clone())
79                    .or_default()
80                    .insert(spec.instrument.clone());
81            }
82            history
83                .entry(spec.instrument.clone())
84                .or_default()
85                .push(spec.clone());
86        }
87
88        for (instrument, specifications) in &mut history {
89            specifications.sort_by(|left, right| {
90                left.effective
91                    .valid_from
92                    .cmp(&right.effective.valid_from)
93                    .then(left.revision.cmp(&right.revision))
94            });
95            for pair in specifications.windows(2) {
96                if pair[0].effective.overlaps(&pair[1].effective) {
97                    return Err(CatalogCompileError::OverlappingInterval {
98                        instrument: instrument.clone(),
99                        left: pair[0].effective,
100                        right: pair[1].effective,
101                    });
102                }
103            }
104        }
105
106        Ok(Self {
107            id: CatalogSnapshotId {
108                version: document.version,
109            },
110            assets,
111            history,
112            aliases,
113        })
114    }
115
116    pub fn id(&self) -> &CatalogSnapshotId {
117        &self.id
118    }
119
120    pub fn asset(&self, asset: &AssetId) -> Option<&AssetSpec> {
121        self.assets.get(asset)
122    }
123
124    pub fn instrument_ids(&self) -> impl Iterator<Item = &InstrumentId> {
125        self.history.keys()
126    }
127
128    pub fn spec_at(
129        &self,
130        instrument: &InstrumentId,
131        at: DateTime<Utc>,
132    ) -> Result<ResolvedInstrument, InstrumentResolutionError> {
133        let history = self
134            .history
135            .get(instrument)
136            .ok_or(InstrumentResolutionError::Unknown)?;
137        let resolved = history
138            .iter()
139            .find(|candidate| candidate.effective.contains(at))
140            .ok_or_else(|| InstrumentResolutionError::Inactive {
141                instrument: instrument.clone(),
142                intervals: history.iter().map(|entry| entry.effective).collect(),
143            })?;
144        Ok(self.resolved(resolved))
145    }
146
147    pub fn resolve(
148        &self,
149        selector: &InstrumentSelector,
150        context: &InstrumentResolutionContext,
151        at: DateTime<Utc>,
152    ) -> Result<ResolvedInstrument, InstrumentResolutionError> {
153        match selector {
154            InstrumentSelector::Exact { instrument } => {
155                if !context.allowed_instruments.contains(instrument) {
156                    return Err(InstrumentResolutionError::Disallowed {
157                        instrument: instrument.clone(),
158                    });
159                }
160                self.spec_at(instrument, at)
161            }
162            InstrumentSelector::Alias {
163                alias,
164                listing_venue,
165                market_kind,
166            } => {
167                let candidates = self
168                    .aliases
169                    .get(alias)
170                    .ok_or(InstrumentResolutionError::Unknown)?;
171                let listing_venue = listing_venue
172                    .as_ref()
173                    .or(context.default_listing_venue.as_ref());
174                let market_kind = market_kind
175                    .as_ref()
176                    .or(context.default_market_kind.as_ref());
177
178                let filtered = candidates
179                    .iter()
180                    .filter(|instrument| {
181                        listing_venue.is_none_or(|venue| &instrument.listing_venue == venue)
182                            && market_kind.is_none_or(|kind| &instrument.market_kind == kind)
183                    })
184                    .cloned()
185                    .collect::<Vec<_>>();
186                if filtered.is_empty() {
187                    return Err(InstrumentResolutionError::Unknown);
188                }
189
190                let allowed = filtered
191                    .iter()
192                    .filter(|instrument| context.allowed_instruments.contains(*instrument))
193                    .cloned()
194                    .collect::<Vec<_>>();
195                if allowed.is_empty() {
196                    return Err(InstrumentResolutionError::Disallowed {
197                        instrument: filtered[0].clone(),
198                    });
199                }
200
201                let mut active = Vec::new();
202                let mut inactive = Vec::new();
203                for instrument in allowed {
204                    let history = self
205                        .history
206                        .get(&instrument)
207                        .ok_or(InstrumentResolutionError::Unknown)?;
208                    let matching = history
209                        .iter()
210                        .filter(|entry| entry.aliases.contains(alias))
211                        .collect::<Vec<_>>();
212                    if let Some(resolved) =
213                        matching.iter().find(|entry| entry.effective.contains(at))
214                    {
215                        active.push(self.resolved(resolved));
216                    } else {
217                        inactive.push((
218                            instrument,
219                            matching.into_iter().map(|entry| entry.effective).collect(),
220                        ));
221                    }
222                }
223                if active.len() == 1 {
224                    return Ok(active.remove(0));
225                }
226                if active.len() > 1 {
227                    return Err(InstrumentResolutionError::Ambiguous {
228                        candidates: active
229                            .into_iter()
230                            .map(|resolved| resolved.reference.instrument)
231                            .collect(),
232                    });
233                }
234                let (instrument, intervals) = inactive.remove(0);
235                Err(InstrumentResolutionError::Inactive {
236                    instrument,
237                    intervals,
238                })
239            }
240        }
241    }
242
243    fn resolved(&self, resolved: &InstrumentSpec) -> ResolvedInstrument {
244        ResolvedInstrument {
245            reference: ResolvedInstrumentRef {
246                instrument: resolved.instrument.clone(),
247                catalog: self.id.clone(),
248                spec_revision: resolved.revision.clone(),
249            },
250            spec: Arc::new(resolved.clone()),
251        }
252    }
253}
254
255fn validate_assets(
256    spec: &InstrumentSpec,
257    assets: &BTreeMap<AssetId, AssetSpec>,
258) -> Result<(), CatalogCompileError> {
259    let mut referenced = BTreeSet::new();
260    referenced.extend(spec.assets.base.iter().cloned());
261    referenced.extend(spec.assets.quote.iter().cloned());
262    referenced.insert(spec.assets.settlement.clone());
263    referenced.extend(spec.assets.fee_assets.iter().cloned());
264    referenced.insert(spec.economics.settlement_asset.clone());
265    if let Some(notional) = &spec.notional {
266        referenced.insert(notional.asset.clone());
267    }
268    for asset in referenced {
269        if !assets.contains_key(&asset) {
270            return Err(CatalogCompileError::UnknownAsset {
271                instrument: spec.instrument.clone(),
272                asset,
273            });
274        }
275    }
276    Ok(())
277}
278
279/// Selector for exact identity or an explicitly scoped alias.
280#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
281#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
282pub enum InstrumentSelector {
283    Exact {
284        instrument: InstrumentId,
285    },
286    Alias {
287        alias: InstrumentAlias,
288        listing_venue: Option<ListingVenueId>,
289        market_kind: Option<MarketKind>,
290    },
291}
292
293/// Immutable deployment or account constraints applied during resolution.
294#[derive(Clone, Debug, Eq, PartialEq)]
295pub struct InstrumentResolutionContext {
296    pub allowed_instruments: BTreeSet<InstrumentId>,
297    pub default_listing_venue: Option<ListingVenueId>,
298    pub default_market_kind: Option<MarketKind>,
299}
300
301/// Persistable identity of one specification resolved from one catalog snapshot.
302#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
303#[serde(deny_unknown_fields)]
304pub struct ResolvedInstrumentRef {
305    pub instrument: InstrumentId,
306    pub catalog: CatalogSnapshotId,
307    pub spec_revision: crate::SpecRevision,
308}
309
310/// Resolved reference and immutable specification value.
311#[derive(Clone, Debug)]
312pub struct ResolvedInstrument {
313    pub reference: ResolvedInstrumentRef,
314    pub spec: Arc<InstrumentSpec>,
315}
316
317/// Catalog compilation failures.
318#[derive(Debug, thiserror::Error)]
319pub enum CatalogCompileError {
320    #[error("unsupported catalog schema version {0}")]
321    UnsupportedSchema(u32),
322    #[error("catalog version must contain between 1 and 64 bytes")]
323    InvalidVersion,
324    #[error("duplicate asset {0}")]
325    DuplicateAsset(AssetId),
326    #[error("instrument {instrument} references unknown asset {asset}")]
327    UnknownAsset {
328        instrument: InstrumentId,
329        asset: AssetId,
330    },
331    #[error("instrument {instrument} has overlapping effective intervals")]
332    OverlappingInterval {
333        instrument: InstrumentId,
334        left: EffectiveInterval,
335        right: EffectiveInterval,
336    },
337    #[error(transparent)]
338    InvalidSpec(#[from] SpecValidationError),
339}
340
341/// Deterministic instrument-resolution failures.
342#[derive(Clone, Debug, Eq, PartialEq, thiserror::Error)]
343pub enum InstrumentResolutionError {
344    #[error("instrument is unknown")]
345    Unknown,
346    #[error("instrument {instrument} is outside the resolution allowlist")]
347    Disallowed { instrument: InstrumentId },
348    #[error("instrument selector is ambiguous")]
349    Ambiguous { candidates: Vec<InstrumentId> },
350    #[error("instrument {instrument} is inactive at the requested time")]
351    Inactive {
352        instrument: InstrumentId,
353        intervals: Vec<EffectiveInterval>,
354    },
355}