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#[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#[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#[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#[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#[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#[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#[derive(Clone, Debug)]
312pub struct ResolvedInstrument {
313 pub reference: ResolvedInstrumentRef,
314 pub spec: Arc<InstrumentSpec>,
315}
316
317#[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#[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}