Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
117 changes: 96 additions & 21 deletions crates/common/src/cache/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2347,7 +2347,8 @@ impl Cache {
.collect();

if let Some(db) = &self.database {
self.index.order_position = db.load_index_order_position()?;
let order_position = db.load_index_order_position()?;
self.index.order_position = self.sanitize_order_position_index(order_position);
self.index.order_client = db.load_index_order_client()?;
}

Expand Down Expand Up @@ -2444,7 +2445,8 @@ impl Cache {
};

if let Some(db) = &self.database {
self.index.order_position = db.load_index_order_position()?;
let order_position = db.load_index_order_position()?;
self.index.order_position = self.sanitize_order_position_index(order_position);
self.index.order_client = db.load_index_order_client()?;
}

Expand All @@ -2454,6 +2456,23 @@ impl Cache {
Ok(())
}

fn sanitize_order_position_index(
&self,
mut order_position: AHashMap<ClientOrderId, PositionId>,
) -> AHashMap<ClientOrderId, PositionId> {
let original_len = order_position.len();
order_position.retain(|client_order_id, _| self.orders.contains_key(client_order_id));
let removed = original_len - order_position.len();

if removed > 0 {
log::warn!(
"Filtered {removed} stale order-position index entries without backing orders during cache load"
);
}

order_position
}

/// Clears and reloads the position cache from the database.
///
/// # Errors
Expand Down Expand Up @@ -2644,25 +2663,32 @@ impl Cache {
.insert(*position_id, position.strategy_id);

// 3: Build index.position_orders -> {PositionId, {ClientOrderId}}
self.index
.position_orders
.entry(*position_id)
.or_default()
.extend(position.client_order_ids());
let position_orders = self.index.position_orders.entry(*position_id).or_default();
position_orders.extend(
position
.client_order_ids()
.into_iter()
.filter(|client_order_id| self.orders.contains_key(client_order_id)),
);

// 4: Build index.instrument_positions -> {InstrumentId, {PositionId}}
self.index
.instrument_positions
.entry(instrument_id)
.or_default()
.insert(*position_id);
self.index
.instrument_orders
.entry(instrument_id)
.or_default();

// 5: Build index.strategy_positions -> {StrategyId, {PositionId}}
self.index
.strategy_positions
.entry(strategy_id)
.or_default()
.insert(*position_id);
self.index.strategy_orders.entry(strategy_id).or_default();

// 6: Build index.account_positions -> {AccountId, {PositionId}}
self.index
Expand Down Expand Up @@ -4632,18 +4658,30 @@ impl Cache {
database.index_order_position(*client_order_id, *position_id)?;
}

// Index: PositionId -> StrategyId
self.index
.position_strategy
.insert(*position_id, *strategy_id);

// Index: PositionId -> set[ClientOrderId]
self.index_position(position_id, venue, strategy_id);
self.index
.position_orders
.entry(*position_id)
.or_default()
.insert(*client_order_id);

Ok(())
}

fn index_position(
&mut self,
position_id: &PositionId,
venue: &Venue,
strategy_id: &StrategyId,
) {
// Index: PositionId -> StrategyId
self.index
.position_strategy
.insert(*position_id, *strategy_id);

// Every position has a reverse-order bucket, including orderless positions.
self.index.position_orders.entry(*position_id).or_default();

// Index: StrategyId -> set[PositionId]
self.index
.strategy_positions
Expand All @@ -4657,8 +4695,6 @@ impl Cache {
.entry(*venue)
.or_default()
.insert(*position_id);

Ok(())
}

// Propagates parent OTO `position_id` to contingent children that are missing one.
Expand Down Expand Up @@ -4722,21 +4758,56 @@ impl Cache {
///
/// Returns an error if persisting the position to the backing database fails.
pub fn add_position(&mut self, position: &Position, oms_type: OmsType) -> anyhow::Result<()> {
self.add_position_inner(position, oms_type, true)
}

/// Adds a position whose opening fill intentionally has no backing order.
///
/// # Errors
///
/// Returns an error if persisting the position to the backing database fails.
pub fn add_position_without_order(
&mut self,
position: &Position,
oms_type: OmsType,
) -> anyhow::Result<()> {
self.add_position_inner(position, oms_type, false)
}

fn add_position_inner(
&mut self,
position: &Position,
oms_type: OmsType,
index_order: bool,
) -> anyhow::Result<()> {
self.positions
.insert(position.id, SharedCell::new(position.clone()));
self.index.position_oms.insert(position.id, oms_type);
self.index.positions.insert(position.id);
self.index.positions_open.insert(position.id);
self.index.positions_closed.remove(&position.id); // Cleanup for NETTING reopen
self.index.strategies.insert(position.strategy_id);
self.index
.strategy_orders
.entry(position.strategy_id)
.or_default();

log::debug!("Adding {position}");

self.add_position_id(
&position.id,
&position.instrument_id.venue,
&position.opening_order_id,
&position.strategy_id,
)?;
if index_order {
self.add_position_id(
&position.id,
&position.instrument_id.venue,
&position.opening_order_id,
&position.strategy_id,
)?;
} else {
self.index_position(
&position.id,
&position.instrument_id.venue,
&position.strategy_id,
);
}

let venue = position.instrument_id.venue;
let venue_positions = self.index.venue_positions.entry(venue).or_default();
Expand All @@ -4750,6 +4821,10 @@ impl Cache {
.entry(instrument_id)
.or_default();
instrument_positions.insert(position.id);
self.index
.instrument_orders
.entry(instrument_id)
.or_default();

// Index: AccountId -> AHashSet<PositionId>
self.index
Expand Down
Loading
Loading