Merge remote-tracking branch 'origin/main'

This commit is contained in:
boris
2026-09-14 09:15:58 +08:00
17 changed files with 4973 additions and 403 deletions
+510 -72
View File
@@ -216,6 +216,9 @@ struct OpenOrder {
commission_remaining: Option<f64>,
execution_cursor: Option<NaiveDateTime>,
reason: String,
algo_request: Option<AlgoExecutionRequest>,
value_budget: Option<f64>,
reserved_cash: Option<f64>,
}
#[derive(Debug, Clone, Copy)]
@@ -225,6 +228,13 @@ struct RestingOrderOrigin {
accepted_date: NaiveDate,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum BrokerCallbackPhase {
Normal,
ControlsOnly,
BeforeStrategy,
}
#[derive(Debug, Default)]
struct BrokerExecutionSession {
date: Option<NaiveDate>,
@@ -420,6 +430,15 @@ struct AlgoExecutionRequest {
style: AlgoExecutionStyle,
start_time: Option<NaiveTime>,
end_time: Option<NaiveTime>,
total_quantity: Option<u32>,
filled_quantity: u32,
commission_remaining: Option<f64>,
order_id: Option<u64>,
}
struct RestoreCell<'a, T: Copy>(&'a Cell<T>, T);
impl<T: Copy> Drop for RestoreCell<'_, T> {
fn drop(&mut self) { self.0.set(self.1); }
}
pub struct BrokerSimulator<C, R> {
@@ -450,6 +469,10 @@ pub struct BrokerSimulator<C, R> {
intraday_execution_start_time: Option<NaiveTime>,
runtime_intraday_start_time: Cell<Option<NaiveTime>>,
runtime_intraday_end_time: Cell<Option<NaiveTime>>,
runtime_execution_clock: Cell<Option<NaiveTime>>,
runtime_callback_phase: Cell<BrokerCallbackPhase>,
runtime_algo_schedule: Cell<Option<AlgoExecutionRequest>>,
runtime_unprocessed_algorithm_cash: Cell<FixedMoney>,
runtime_decision_date: Cell<Option<NaiveDate>>,
runtime_buy_denials: RefCell<BTreeMap<String, String>>,
runtime_auto_buy_denials: RefCell<BTreeMap<String, String>>,
@@ -494,6 +517,10 @@ impl<C, R> BrokerSimulator<C, R> {
intraday_execution_start_time: None,
runtime_intraday_start_time: Cell::new(None),
runtime_intraday_end_time: Cell::new(None),
runtime_execution_clock: Cell::new(None),
runtime_callback_phase: Cell::new(BrokerCallbackPhase::Normal),
runtime_algo_schedule: Cell::new(None),
runtime_unprocessed_algorithm_cash: Cell::new(FixedMoney::ZERO),
runtime_decision_date: Cell::new(None),
runtime_buy_denials: RefCell::new(BTreeMap::new()),
runtime_auto_buy_denials: RefCell::new(BTreeMap::new()),
@@ -542,6 +569,10 @@ impl<C, R> BrokerSimulator<C, R> {
intraday_execution_start_time: None,
runtime_intraday_start_time: Cell::new(None),
runtime_intraday_end_time: Cell::new(None),
runtime_execution_clock: Cell::new(None),
runtime_callback_phase: Cell::new(BrokerCallbackPhase::Normal),
runtime_algo_schedule: Cell::new(None),
runtime_unprocessed_algorithm_cash: Cell::new(FixedMoney::ZERO),
runtime_decision_date: Cell::new(None),
runtime_buy_denials: RefCell::new(BTreeMap::new()),
runtime_auto_buy_denials: RefCell::new(BTreeMap::new()),
@@ -726,6 +757,10 @@ impl<C, R> BrokerSimulator<C, R> {
.or(self.intraday_execution_start_time)
}
fn execution_clock(&self) -> Option<NaiveTime> {
self.runtime_execution_clock.get().or(self.runtime_intraday_start_time.get())
}
fn order_origin(&self) -> (Option<NaiveDate>, Option<NaiveTime>) {
self.runtime_resting_order_origin.get().map_or(
(self.runtime_order_created_date.get(), self.submission_time()),
@@ -898,6 +933,7 @@ impl<C, R> BrokerSimulator<C, R> {
avg_price: 0.0,
transaction_cost: 0.0,
limit_price: order.limit_price,
reserved_cash: order.reserved_cash,
reason: order.reason.clone(),
})
.collect()
@@ -916,11 +952,12 @@ impl<C, R> BrokerSimulator<C, R> {
fn resting_order_session_close(&self, date: NaiveDate, order: &OpenOrder) -> NaiveTime {
let post_close = self.execution_phase_for_submission(date, order.order_created_date, order.submission_time)
== EquityExecutionPhase::PostCloseFixedPrice;
NaiveTime::from_hms_opt(15, if post_close { 30 } else { 0 }, 0).expect("session end")
let close=NaiveTime::from_hms_opt(15, if post_close { 30 } else { 0 }, 0).expect("session end");
order.algo_request.and_then(|request|request.end_time).map_or(close,|end|end.min(close))
}
pub(crate) fn next_day_order_expiry(&self, date: NaiveDate) -> Option<NaiveTime> {
self.open_orders.borrow().iter().filter(|order| order.time_in_force == OrderTimeInForce::Day)
self.open_orders.borrow().iter().filter(|order| order.time_in_force == OrderTimeInForce::Day || order.algo_request.is_some())
.map(|order| self.resting_order_session_close(date, order)).min()
}
}
@@ -1616,17 +1653,21 @@ where
self.deferred_stock_pools.borrow_mut().remove(&contract.pool_id);
}
}
self.process_open_orders(
date,
portfolio,
data,
&mut session.intraday_turnover,
&mut session.execution_cursors,
&mut session.global_execution_cursor,
&mut session.commission_state,
&mut report,
)?;
self.resume_stock_pool_executions(date, portfolio, data, session, &mut report)?;
if self.runtime_callback_phase.get() != BrokerCallbackPhase::ControlsOnly {
self.process_open_orders(
date,
portfolio,
data,
&mut session.intraday_turnover,
&mut session.execution_cursors,
&mut session.global_execution_cursor,
&mut session.commission_state,
&mut report,
)?;
if self.runtime_callback_phase.get() == BrokerCallbackPhase::Normal {
self.resume_stock_pool_executions(date, portfolio, data, session, &mut report)?;
}
}
if !decision.order_intents.is_empty() {
let mut ordered_intents = decision.order_intents.iter().collect::<Vec<_>>();
if self.effective_rebalance_cash_mode() != RebalanceCashMode::PreOpenCash
@@ -1803,6 +1844,110 @@ where
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn execute_controls_without_matching(
&self,
date: NaiveDate,
decision_date: NaiveDate,
portfolio: &mut PortfolioState,
data: &DataSet,
decision: &StrategyDecision,
clock: Option<NaiveTime>,
) -> Result<BrokerExecutionReport, BacktestError> {
if decision.rebalance
|| !decision.target_weights.is_empty()
|| !decision.exit_symbols.is_empty()
|| decision.order_intents.iter().any(|intent| {
!matches!(
intent.unwrapped(),
OrderIntent::CancelOrder { .. }
| OrderIntent::CancelSymbol { .. }
| OrderIntent::CancelAll { .. }
| OrderIntent::ModifyOrder { .. }
)
})
{
return Err(BacktestError::Execution(
"non-matching control phase only accepts cancel or modify requests".into(),
));
}
let _guard = RestoreCell(
&self.runtime_callback_phase,
self.runtime_callback_phase
.replace(BrokerCallbackPhase::ControlsOnly),
);
self.execute_between_with_event_dates(
date,
decision_date,
decision_date,
portfolio,
data,
decision,
clock,
clock,
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn execute_coarse_at_clock(
&self,
date: NaiveDate,
decision_date: NaiveDate,
order_created_date: NaiveDate,
decision_total_equity: Option<f64>,
portfolio: &mut PortfolioState,
data: &DataSet,
decision: &StrategyDecision,
clock: Option<NaiveTime>,
) -> Result<BrokerExecutionReport, BacktestError> {
// Advancing the engine clock must not turn a daily closing-bar order
// into an explicitly submitted post-close order.
let _clock_guard = RestoreCell(
&self.runtime_execution_clock,
self.runtime_execution_clock.replace(clock),
);
self.execute_between_with_event_dates_and_decision_equity(
date,
decision_date,
order_created_date,
decision_total_equity,
portfolio,
data,
decision,
None,
clock,
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn execute_before_strategy_at_clock(
&self,
date: NaiveDate,
decision_date: NaiveDate,
order_created_date: NaiveDate,
decision_total_equity: Option<f64>,
portfolio: &mut PortfolioState,
data: &DataSet,
decision: &StrategyDecision,
clock: Option<NaiveTime>,
) -> Result<BrokerExecutionReport, BacktestError> {
let _guard = RestoreCell(
&self.runtime_callback_phase,
self.runtime_callback_phase
.replace(BrokerCallbackPhase::BeforeStrategy),
);
self.execute_coarse_at_clock(
date,
decision_date,
order_created_date,
decision_total_equity,
portfolio,
data,
decision,
clock,
)
}
pub fn execute_between_with_event_dates(
&self,
date: NaiveDate,
@@ -2682,18 +2827,26 @@ where
let mut open_orders = self.open_orders.borrow_mut();
std::mem::take(&mut *open_orders)
};
let reserved=FixedMoney::checked_sum_f64(pending_orders.iter().filter_map(|order|order.reserved_cash))
.ok_or_else(||BacktestError::Execution("working order cash reservation is invalid".into()))?;
let _reservation_guard=RestoreCell(&self.runtime_unprocessed_algorithm_cash,
self.runtime_unprocessed_algorithm_cash.replace(reserved));
for order in pending_orders {
if let Some(reserved)=order.reserved_cash {
self.runtime_unprocessed_algorithm_cash.set(self.runtime_unprocessed_algorithm_cash.get()
.checked_sub(FixedMoney::from_f64(reserved).expect("validated reservation")).expect("reserved cash subset"));
}
if self.matching_type == MatchingType::NextBarOpen && self.runtime_intraday_start_time.get().is_none()
&& order.accepted_date == date {
&& order.accepted_date == date && order.algo_request.is_none() {
self.open_orders.borrow_mut().push(order);
continue;
}
let close = self.resting_order_session_close(date, &order);
let clock = self.submission_time();
let past_day = order.time_in_force == OrderTimeInForce::Day
let clock = self.execution_clock().or(self.submission_time());
let past_day = (order.time_in_force == OrderTimeInForce::Day || order.algo_request.is_some())
&& order.accepted_date < date;
if past_day || clock.is_some_and(|time| time > close) {
if order.time_in_force == OrderTimeInForce::Day {
if order.time_in_force == OrderTimeInForce::Day || order.algo_request.is_some() {
Self::emit_resting_day_expiry(report, date, &order, order.filled_quantity);
} else {
self.open_orders.borrow_mut().push(order);
@@ -2730,7 +2883,18 @@ where
accepted_date: order.accepted_date,
}));
let previous_decision_date = self.runtime_decision_date.replace(order.decision_date);
let execution_result = self.process_limit_shares_internal(
let execution_result = if let Some(mut algorithm)=order.algo_request {
algorithm.total_quantity=Some(order.requested_quantity);
algorithm.filled_quantity=order.filled_quantity;
algorithm.commission_remaining=order.commission_remaining;
if order.side==OrderSide::Buy {
self.process_buy(date,portfolio,data,&order.symbol,order.remaining_quantity,order.order_id,&order.reason,
intraday_turnover,execution_cursors,global_execution_cursor,commission_state,order.value_budget,None,false,false,Some(&algorithm),report)
} else {
self.process_sell(date,portfolio,data,&order.symbol,order.remaining_quantity,order.order_id,&order.reason,
intraday_turnover,execution_cursors,global_execution_cursor,commission_state,None,false,false,Some(&algorithm),report)
}
} else { self.process_limit_shares_internal(
date,
portfolio,
data,
@@ -2745,7 +2909,7 @@ where
global_execution_cursor,
commission_state,
report,
);
) };
self.runtime_time_in_force.set(previous_time_in_force);
self.runtime_resting_order_origin.set(previous_origin);
self.runtime_decision_date.set(previous_decision_date);
@@ -2843,7 +3007,8 @@ where
}
fn emit_resting_day_expiry(report: &mut BrokerExecutionReport, date: NaiveDate, order: &OpenOrder, filled: u32) {
let detail = format!("DAY order expired at market close: {} remaining_quantity={}", order.symbol, order.requested_quantity.saturating_sub(filled));
let label=if order.algo_request.is_some() {"algorithm execution window expired"} else {"DAY order expired at market close"};
let detail = format!("{label}: {} remaining_quantity={}", order.symbol, order.requested_quantity.saturating_sub(filled));
report.order_events.push(OrderEvent {
date, decision_date: order.decision_date, order_created_date: order.order_created_date,
execution_date: Some(date), order_id: Some(order.order_id), symbol: order.symbol.clone(),
@@ -2929,6 +3094,11 @@ where
let target_total_quantity = new_total_quantity.unwrap_or(existing.requested_quantity);
let target_limit_price = new_limit_price.unwrap_or(existing.limit_price);
if existing.algo_request.is_some() {
Self::emit_open_order_update_rejected(report,date,order_id,Some(&existing.symbol),Some(existing.side),reason,
"algorithm schedule is immutable; cancel it before submitting a different schedule");
return;
}
if target_total_quantity == existing.requested_quantity
&& target_limit_price.to_bits() == existing.limit_price.to_bits()
{
@@ -3898,6 +4068,10 @@ where
},
start_time: *start_time,
end_time: *end_time,
total_quantity: None,
filled_quantity: 0,
commission_remaining: None,
order_id: None,
}),
_ => None,
};
@@ -4173,9 +4347,8 @@ where
return self.execution_limit_check_price(snapshot, side);
}
let matching_type = self.matching_type_for_algo_request(algo_request);
let start_cursor = algo_request
.and_then(|request| request.start_time)
.or(self.runtime_intraday_start_time.get())
let start_cursor = self.execution_clock()
.or_else(||algo_request.and_then(|request| request.start_time))
.or(self.intraday_execution_start_time)
.map(|start_time| date.and_time(start_time));
self.latest_known_quote_at_or_before(
@@ -4187,7 +4360,9 @@ where
false,
)
.and_then(|quote| self.select_quote_reference_price(snapshot, quote, side, matching_type))
.unwrap_or_else(|| self.execution_limit_check_price(snapshot, side))
.unwrap_or_else(|| if algo_request.is_some() && self.execution_clock().is_some() {
f64::NAN
} else {self.execution_limit_check_price(snapshot, side)})
}
#[cfg(test)]
@@ -4534,6 +4709,8 @@ where
algo_request: Option<&AlgoExecutionRequest>,
report: &mut BrokerExecutionReport,
) -> Result<(), BacktestError> {
let algorithm = self.normalized_algorithm(date, requested_qty, order_id, commission_state.get(&order_id).copied(), algo_request);
let algo_request = algorithm.as_ref();
// Existing accepted orders are not canceled by a subsequently enabled lock.
if emit_creation_events && self.runtime_auto_sell_denials.borrow().contains_key(symbol) {
return Ok(());
@@ -4768,6 +4945,9 @@ where
time_in_force: Self::pending_time_in_force(remainder_policy),
commission_remaining: commission_state.get(&order_id).copied(),
execution_cursor: execution_cursors.get(symbol).copied(),
algo_request: None,
value_budget: None,
reserved_cash: None,
reason: reason.to_string(),
});
// Waiting without a fill is not a new order-state transition.
@@ -4859,6 +5039,9 @@ where
time_in_force: Self::pending_time_in_force(remainder_policy),
commission_remaining: commission_state.get(&order_id).copied(),
execution_cursor: execution_cursors.get(symbol).copied(),
algo_request: None,
value_budget: None,
reserved_cash: None,
reason: reason.to_string(),
});
// Waiting without a fill is not a new order-state transition.
@@ -4976,8 +5159,8 @@ where
price: execution_price,
mark_price: self.snapshot_mark_price(snapshot, OrderSide::Sell),
quantity: fillable_qty,
execution_start_timestamp: None,
execution_timestamp: None,
execution_start_timestamp: self.runtime_execution_clock.get().map(|time|date.and_time(time)),
execution_timestamp: self.runtime_execution_clock.get().map(|time|date.and_time(time)),
}],
None,
Vec::new(),
@@ -5014,8 +5197,9 @@ where
let detail = partial_fill_reason
.as_deref()
.unwrap_or("limit price not marketable yet");
if Self::keeps_remainder_open(remainder_policy)
&& Self::limit_order_can_remain_open(Some(detail))
if (Self::keeps_remainder_open(remainder_policy)
&& Self::limit_order_can_remain_open(Some(detail)))
|| self.algorithm_still_working(algo_request, Some(detail))
{
self.upsert_open_order(OpenOrder {
order_id,
@@ -5028,10 +5212,13 @@ where
requested_quantity: requested_qty,
filled_quantity: 0,
remaining_quantity: requested_qty,
limit_price: limit_price.expect("limit price for pending limit sell"),
time_in_force: Self::pending_time_in_force(remainder_policy),
limit_price: if algo_request.is_some() {limit_price.unwrap_or(0.0)} else {limit_price.expect("limit price for pending limit sell")},
time_in_force: if algo_request.is_some() {self.runtime_time_in_force.get().unwrap_or(OrderTimeInForce::Day)} else {Self::pending_time_in_force(remainder_policy)},
commission_remaining: commission_state.get(&order_id).copied(),
execution_cursor: execution_cursors.get(symbol).copied(),
algo_request: Self::progressed_algorithm(algo_request, 0, commission_state.get(&order_id).copied()),
value_budget: None,
reserved_cash: None,
reason: reason.to_string(),
});
// Waiting without a fill is not a new order-state transition.
@@ -5072,7 +5259,7 @@ where
side: OrderSide::Sell,
requested_quantity: requested_qty,
filled_quantity: 0,
status: zero_fill_status_for_reason(detail),
status: self.unfilled_algorithm_status(algo_request, detail),
reason: format!("{reason}: {detail}"),
});
Self::emit_order_process_event(
@@ -5084,7 +5271,7 @@ where
OrderSide::Sell,
format!(
"status={:?} reason={detail}",
zero_fill_status_for_reason(detail)
self.unfilled_algorithm_status(algo_request, detail)
),
);
self.clear_open_order(order_id);
@@ -5185,9 +5372,10 @@ where
*intraday_turnover.entry(symbol.to_string()).or_default() += filled_qty;
let remaining_qty = requested_qty.saturating_sub(filled_qty);
let keep_open = Self::keeps_remainder_open(remainder_policy)
let keep_open = (Self::keeps_remainder_open(remainder_policy)
&& remaining_qty > 0
&& Self::limit_order_can_remain_open(partial_fill_reason.as_deref());
&& Self::limit_order_can_remain_open(partial_fill_reason.as_deref()))
|| (remaining_qty > 0 && self.algorithm_still_working(algo_request,partial_fill_reason.as_deref()));
if keep_open {
self.upsert_open_order(OpenOrder {
order_id,
@@ -5200,10 +5388,13 @@ where
requested_quantity: requested_qty,
filled_quantity: filled_qty,
remaining_quantity: remaining_qty,
limit_price: limit_price.expect("limit price for pending limit sell"),
time_in_force: Self::pending_time_in_force(remainder_policy),
limit_price: if algo_request.is_some() {limit_price.unwrap_or(0.0)} else {limit_price.expect("limit price for pending limit sell")},
time_in_force: if algo_request.is_some() {self.runtime_time_in_force.get().unwrap_or(OrderTimeInForce::Day)} else {Self::pending_time_in_force(remainder_policy)},
commission_remaining: commission_state.get(&order_id).copied(),
execution_cursor: execution_cursors.get(symbol).copied(),
algo_request: Self::progressed_algorithm(algo_request, filled_qty, commission_state.get(&order_id).copied()),
value_budget: None,
reserved_cash: None,
reason: reason.to_string(),
});
} else {
@@ -5213,7 +5404,7 @@ where
let status = if keep_open {
OrderStatus::PartiallyFilled
} else if filled_qty < requested_qty {
OrderStatus::Canceled
if self.algorithm_window_expired(algo_request, partial_fill_reason.as_deref().unwrap_or("")) {OrderStatus::Expired} else {OrderStatus::Canceled}
} else {
OrderStatus::Filled
};
@@ -5250,7 +5441,7 @@ where
status,
reason: order_reason,
});
if matches!(status, OrderStatus::Canceled | OrderStatus::Rejected) {
if matches!(status, OrderStatus::Canceled | OrderStatus::Rejected | OrderStatus::Expired) {
Self::emit_order_process_event(
report,
date,
@@ -5399,6 +5590,10 @@ where
},
start_time,
end_time,
total_quantity: None,
filled_quantity: 0,
commission_remaining: None,
order_id: None,
};
if target_value <= f64::EPSILON {
@@ -6080,12 +6275,19 @@ where
},
start_time,
end_time,
total_quantity: None,
filled_quantity: 0,
commission_remaining: None,
order_id: None,
};
if value > 0.0 {
let round_lot = self.round_lot(data, symbol);
let minimum_order_quantity = self.minimum_order_quantity(data, symbol);
let order_step_size = self.order_step_size(data, symbol);
let price = self.sizing_price(snapshot);
let price = self.execution_order_limit_check_price(date, data, symbol, snapshot, OrderSide::Buy, Some(&algo_request));
if !price.is_finite() || price <= 0.0 {
return Err(BacktestError::MissingPrice {date, symbol:symbol.to_string(), field:"algorithm_submission_price"});
}
let snapshot_requested_qty = self.value_buy_quantity(
date,
value.abs(),
@@ -6126,7 +6328,10 @@ where
report,
)
} else {
let price = self.sizing_price(snapshot);
let price = self.execution_order_limit_check_price(date, data, symbol, snapshot, OrderSide::Sell, Some(&algo_request));
if !price.is_finite() || price <= 0.0 {
return Err(BacktestError::MissingPrice {date, symbol:symbol.to_string(), field:"algorithm_submission_price"});
}
let requested_qty = self.round_buy_quantity(
(value.abs() / price).floor() as u32,
self.minimum_order_quantity(data, symbol),
@@ -6337,6 +6542,9 @@ where
algo_request: Option<&AlgoExecutionRequest>,
report: &mut BrokerExecutionReport,
) -> Result<(), BacktestError> {
let algorithm = self.normalized_algorithm(date, requested_qty, order_id, commission_state.get(&order_id).copied(), algo_request);
let algo_request = algorithm.as_ref();
let fill_start = report.fill_events.len();
if emit_creation_events && self.runtime_auto_buy_denials.borrow().contains_key(symbol) {
return Ok(());
}
@@ -6592,6 +6800,9 @@ where
time_in_force: Self::pending_time_in_force(remainder_policy),
commission_remaining: commission_state.get(&order_id).copied(),
execution_cursor: execution_cursors.get(symbol).copied(),
algo_request: None,
value_budget: None,
reserved_cash: None,
reason: reason.to_string(),
});
// Waiting without a fill is not a new order-state transition.
@@ -6651,13 +6862,14 @@ where
}
};
let value_gross_limit = self.value_buy_gross_limit(value_budget);
let available_cash=self.cash_after_algorithm_reservations(portfolio.cash(),Some(order_id))?;
let buy_cash_limit = if self.strict_value_budget {
value_budget
.filter(|budget| budget.is_finite() && *budget > 0.0)
.map(|budget| portfolio.cash().min(budget))
.unwrap_or_else(|| portfolio.cash())
.map(|budget| available_cash.min(budget))
.unwrap_or(available_cash)
} else {
portfolio.cash()
available_cash
};
let fill = self.resolve_execution_fill(
@@ -6779,8 +6991,8 @@ where
price: execution_price,
mark_price: self.snapshot_mark_price(snapshot, OrderSide::Buy),
quantity: filled_qty,
execution_start_timestamp: None,
execution_timestamp: None,
execution_start_timestamp: self.runtime_execution_clock.get().map(|time|date.and_time(time)),
execution_timestamp: self.runtime_execution_clock.get().map(|time|date.and_time(time)),
}],
None,
Vec::new(),
@@ -6814,8 +7026,9 @@ where
let detail = partial_fill_reason
.as_deref()
.unwrap_or("insufficient cash after fees");
if Self::keeps_remainder_open(remainder_policy)
&& Self::limit_order_can_remain_open(Some(detail))
if (Self::keeps_remainder_open(remainder_policy)
&& Self::limit_order_can_remain_open(Some(detail)))
|| self.algorithm_still_working(algo_request,Some(detail))
{
self.upsert_open_order(OpenOrder {
order_id,
@@ -6828,10 +7041,13 @@ where
requested_quantity: requested_qty,
filled_quantity: 0,
remaining_quantity: requested_qty,
limit_price: limit_price.expect("limit price for pending limit buy"),
time_in_force: Self::pending_time_in_force(remainder_policy),
limit_price: if algo_request.is_some() {limit_price.unwrap_or(0.0)} else {limit_price.expect("limit price for pending limit buy")},
time_in_force: if algo_request.is_some() {self.runtime_time_in_force.get().unwrap_or(OrderTimeInForce::Day)} else {Self::pending_time_in_force(remainder_policy)},
commission_remaining: commission_state.get(&order_id).copied(),
execution_cursor: execution_cursors.get(symbol).copied(),
algo_request: Self::progressed_algorithm(algo_request, 0, commission_state.get(&order_id).copied()),
value_budget: if algo_request.is_some() {value_budget} else {None},
reserved_cash: if algo_request.is_some() {Some(self.algorithm_cash_reservation(date,value_budget,requested_qty,size_check_price,order_id,commission_state.get(&order_id).copied(),data.instruments().get(symbol),portfolio.cash())?)} else {None},
reason: reason.to_string(),
});
// Waiting without a fill is not a new order-state transition.
@@ -6872,7 +7088,7 @@ where
side: OrderSide::Buy,
requested_quantity: requested_qty,
filled_quantity: 0,
status: zero_fill_status_for_reason(detail),
status: self.unfilled_algorithm_status(algo_request, detail),
reason: format!("{reason}: {detail}"),
});
Self::emit_order_process_event(
@@ -6884,7 +7100,7 @@ where
OrderSide::Buy,
format!(
"status={:?} reason={detail}",
zero_fill_status_for_reason(detail)
self.unfilled_algorithm_status(algo_request, detail)
),
);
self.clear_open_order(order_id);
@@ -6987,9 +7203,10 @@ where
*intraday_turnover.entry(symbol.to_string()).or_default() += filled_qty;
let remaining_qty = requested_qty.saturating_sub(filled_qty);
let keep_open = Self::keeps_remainder_open(remainder_policy)
let keep_open = (Self::keeps_remainder_open(remainder_policy)
&& remaining_qty > 0
&& Self::limit_order_can_remain_open(partial_fill_reason.as_deref());
&& Self::limit_order_can_remain_open(partial_fill_reason.as_deref()))
|| (remaining_qty > 0 && self.algorithm_still_working(algo_request,partial_fill_reason.as_deref()));
if keep_open {
self.upsert_open_order(OpenOrder {
order_id,
@@ -7002,10 +7219,13 @@ where
requested_quantity: requested_qty,
filled_quantity: filled_qty,
remaining_quantity: remaining_qty,
limit_price: limit_price.expect("limit price for pending limit buy"),
time_in_force: Self::pending_time_in_force(remainder_policy),
limit_price: if algo_request.is_some() {limit_price.unwrap_or(0.0)} else {limit_price.expect("limit price for pending limit buy")},
time_in_force: if algo_request.is_some() {self.runtime_time_in_force.get().unwrap_or(OrderTimeInForce::Day)} else {Self::pending_time_in_force(remainder_policy)},
commission_remaining: commission_state.get(&order_id).copied(),
execution_cursor: execution_cursors.get(symbol).copied(),
algo_request: Self::progressed_algorithm(algo_request, filled_qty, commission_state.get(&order_id).copied()),
value_budget: if algo_request.is_some() {self.remaining_algorithm_budget(value_budget,&report.fill_events[fill_start..])?} else {None},
reserved_cash: if algo_request.is_some() {Some(self.algorithm_cash_reservation(date,self.remaining_algorithm_budget(value_budget,&report.fill_events[fill_start..])?,remaining_qty,size_check_price,order_id,commission_state.get(&order_id).copied(),data.instruments().get(symbol),portfolio.cash())?)} else {None},
reason: reason.to_string(),
});
} else {
@@ -7015,7 +7235,7 @@ where
let status = if keep_open {
OrderStatus::PartiallyFilled
} else if filled_qty < requested_qty {
OrderStatus::Canceled
if self.algorithm_window_expired(algo_request, partial_fill_reason.as_deref().unwrap_or("")) {OrderStatus::Expired} else {OrderStatus::Canceled}
} else {
OrderStatus::Filled
};
@@ -7052,7 +7272,7 @@ where
status,
reason: order_reason,
});
if matches!(status, OrderStatus::Canceled | OrderStatus::Rejected) {
if matches!(status, OrderStatus::Canceled | OrderStatus::Rejected | OrderStatus::Expired) {
Self::emit_order_process_event(
report,
date,
@@ -7569,6 +7789,192 @@ where
})
}
fn normalized_algorithm(
&self,
date: NaiveDate,
quantity: u32,
order_id: u64,
commission: Option<f64>,
request: Option<&AlgoExecutionRequest>,
) -> Option<AlgoExecutionRequest> {
request
.copied()
.or_else(|| {
(self.matching_type == MatchingType::Vwap).then_some(AlgoExecutionRequest {
style: AlgoExecutionStyle::Vwap,
start_time: self.submission_time(),
end_time: None,
total_quantity: None,
filled_quantity: 0,
commission_remaining: None,
order_id: None,
})
})
.map(|mut request| {
request.total_quantity.get_or_insert(quantity);
request.order_id = Some(order_id);
request.commission_remaining = commission;
if request.start_time.is_none() {
request.start_time = self.execution_clock().or(self.submission_time());
}
if request.end_time.is_none() && request.style == AlgoExecutionStyle::Vwap {
request.end_time = Some(
self.post_close_execution_window(date)
.map(|(_, end)| end.time())
.unwrap_or_else(|| {
NaiveTime::from_hms_opt(15, 0, 0).expect("cash session close")
}),
);
}
request
})
}
fn cash_after_algorithm_reservations(
&self,
cash: f64,
except: Option<u64>,
) -> Result<f64, BacktestError> {
let reserved = FixedMoney::checked_sum_f64(
self.open_orders
.borrow()
.iter()
.filter(|order| except != Some(order.order_id))
.filter_map(|order| order.reserved_cash),
)
.and_then(|amount| amount.checked_add(self.runtime_unprocessed_algorithm_cash.get()))
.ok_or_else(|| BacktestError::Execution("algorithm reserved cash overflow".into()))?;
FixedMoney::from_f64(cash)
.and_then(|cash| cash.checked_sub(reserved))
.map(|available| available.max(FixedMoney::ZERO).to_f64())
.ok_or_else(|| BacktestError::Execution("algorithm available cash is invalid".into()))
}
#[allow(clippy::too_many_arguments)]
fn algorithm_cash_reservation(
&self,
date: NaiveDate,
budget: Option<f64>,
quantity: u32,
price: f64,
order_id: u64,
commission: Option<f64>,
instrument: Option<&Instrument>,
cash: f64,
) -> Result<f64, BacktestError> {
let available = self.cash_after_algorithm_reservations(cash, Some(order_id))?;
if let Some(budget) = budget.filter(|_| self.strict_value_budget) {
return Ok(budget.min(available));
}
let gross = budget.unwrap_or(price * f64::from(quantity));
if !gross.is_finite() || gross < 0. {
return Err(BacktestError::Execution(
"algorithm reservation requires a current price or explicit value budget".into(),
));
}
let mut state = commission
.map(|left| (order_id, left))
.into_iter()
.collect();
let cost = self.cost_model.calculate_with_order_state_for_instrument(
date,
OrderSide::Buy,
gross,
Some(order_id),
&mut state,
instrument,
);
FixedMoney::checked_sum_f64([gross, cost.total()])
.map(|amount| amount.to_f64().min(available))
.ok_or_else(|| BacktestError::Execution("algorithm cash reservation overflow".into()))
}
fn algorithm_still_working(
&self,
request: Option<&AlgoExecutionRequest>,
reason: Option<&str>,
) -> bool {
request.is_some_and(|request| {
self.runtime_intraday_end_time
.get()
.zip(request.end_time)
.is_some_and(|(clock, end)| clock < end)
}) && self
.runtime_time_in_force
.get()
.is_none_or(|tif| matches!(tif, OrderTimeInForce::Day | OrderTimeInForce::Gtc))
&& Self::limit_order_can_remain_open(reason)
}
fn unfilled_algorithm_status(
&self,
request: Option<&AlgoExecutionRequest>,
reason: &str,
) -> OrderStatus {
if self.algorithm_window_expired(request, reason) {
OrderStatus::Expired
} else {
zero_fill_status_for_reason(reason)
}
}
fn algorithm_window_expired(
&self,
request: Option<&AlgoExecutionRequest>,
reason: &str,
) -> bool {
request.is_some_and(|request| {
self.runtime_intraday_end_time
.get()
.zip(request.end_time)
.is_some_and(|(clock, end)| clock >= end)
}) && matches!(
reason,
"intraday quote liquidity exhausted"
| "no execution quotes after start"
| "no execution quotes at or before start"
)
}
fn progressed_algorithm(
request: Option<&AlgoExecutionRequest>,
filled: u32,
commission: Option<f64>,
) -> Option<AlgoExecutionRequest> {
request.copied().map(|mut request| {
request.filled_quantity = request.filled_quantity.saturating_add(filled);
request.commission_remaining = commission;
request
})
}
fn remaining_algorithm_budget(
&self,
budget: Option<f64>,
fills: &[FillEvent],
) -> Result<Option<f64>, BacktestError> {
let Some(budget) = budget else {
return Ok(None);
};
let spent = FixedMoney::checked_sum_f64(fills.iter().map(|fill| {
if self.strict_value_budget {
-fill.net_cash_flow
} else {
fill.gross_amount
}
}))
.ok_or_else(|| {
BacktestError::Execution("algorithm budget spent amount is invalid".into())
})?;
let remaining = FixedMoney::from_f64(budget)
.and_then(|budget| budget.checked_sub(spent))
.filter(|remaining| *remaining >= FixedMoney::ZERO)
.ok_or_else(|| {
BacktestError::Execution("algorithm spent more than its frozen value budget".into())
})?;
Ok(Some(remaining.to_f64()))
}
fn resolve_execution_fill(
&self,
date: NaiveDate,
@@ -7614,6 +8020,12 @@ where
{
Some(start_cursor.map_or(date.and_time(submitted), |cursor| cursor.max(date.and_time(submitted))))
} else { start_cursor };
let start_cursor = if algo_request.is_some() {
match (start_cursor, self.execution_clock().map(|time| date.and_time(time))) {
(Some(declared), Some(clock)) => Some(declared.max(clock)),
(start, _) => start,
}
} else { start_cursor };
let end_cursor = post_close_window.map(|window| {
runtime_end_time.map_or(window.1, |end| window.1.min(date.and_time(end)))
}).or_else(|| {
@@ -7630,10 +8042,17 @@ where
} else {
end_cursor
};
let end_cursor = if algo_request.is_some() {
match (end_cursor, runtime_end_time.map(|time| date.and_time(time))) {
(Some(declared), Some(clock)) => Some(declared.min(clock)),
(end, _) => end,
}
} else { end_cursor };
let quotes = data.execution_quotes_on(date, symbol);
let calibration = self.slippage_calibration(data, snapshot)?;
if let Some(fill) = self.select_execution_fill_with_ledger(
let previous_schedule = self.runtime_algo_schedule.replace(algo_request.copied());
let selected = self.select_execution_fill_with_ledger(
symbol,
snapshot,
quotes,
@@ -7652,7 +8071,9 @@ where
execution_ledger,
calibration.as_ref(),
data.instruments().get(symbol),
)? {
);
self.runtime_algo_schedule.set(previous_schedule);
if let Some(fill) = selected? {
return Ok(Some(fill));
}
@@ -7662,11 +8083,8 @@ where
|| runtime_end_time.is_some()
|| self.intraday_execution_start_time.is_some()
{
let next_cursor = algo_request
.and_then(|request| request.start_time)
.or(runtime_start_time)
.or(self.intraday_execution_start_time)
.map(|start_time| date.and_time(start_time) + Duration::seconds(1))
let next_cursor = start_cursor
.map(|time| time + Duration::seconds(1))
.unwrap_or_else(|| date.and_hms_opt(0, 0, 1).expect("valid midnight"));
return Ok(Some(ExecutionFill {
quantity: 0,
@@ -7778,16 +8196,24 @@ where
return Ok(None);
}
let algo_schedule = self.runtime_algo_schedule.get();
let mut preview_commission_state = BTreeMap::new();
let schedule_start = algo_schedule.and_then(|request| request.start_time)
.map(|time| snapshot.date.and_time(time)).or(start_cursor);
let schedule_end = algo_schedule.and_then(|request| request.end_time)
.map(|time| snapshot.date.and_time(time)).or(end_cursor);
let quote_quantity_limited =
self.quote_quantity_limited_for_window(matching_type, start_cursor, end_cursor);
self.quote_quantity_limited_for_window(matching_type, schedule_start, schedule_end);
let twap_schedule = (matching_type == MatchingType::Twap)
.then(|| TwapSchedule::new(start_cursor, end_cursor, requested_qty))
.then(|| TwapSchedule::new(schedule_start, schedule_end,
algo_schedule.and_then(|request|request.total_quantity).unwrap_or(requested_qty)))
.transpose()?;
let lot = round_lot.max(1);
let exact_time_order_quote = matching_type != MatchingType::MinuteLast
&& start_cursor.is_some()
&& end_cursor.is_some()
&& start_cursor == end_cursor;
&& start_cursor == end_cursor
&& !(algo_schedule.is_some() && schedule_start != schedule_end);
let use_decision_time_quote = !self.is_post_close_fixed_price(snapshot.date)
&& start_cursor.is_some()
&& (matching_type == MatchingType::MinuteLast || exact_time_order_quote);
@@ -7923,7 +8349,8 @@ where
}
let mut take_qty = if let Some(schedule) = &twap_schedule {
remaining_qty.min(available_qty).min(schedule.due_quantity(execution_at, filled_qty))
remaining_qty.min(available_qty).min(schedule.due_quantity(execution_at,
algo_schedule.map_or(0,|request|request.filled_quantity).saturating_add(filled_qty)))
} else {
remaining_qty.min(available_qty)
};
@@ -7984,10 +8411,16 @@ where
);
continue;
}
let candidate_cost = self
.cost_model
.calculate_for_instrument(snapshot.date, OrderSide::Buy, candidate_gross, instrument)
.total();
let candidate_cost = if let Some(request)=algo_schedule {
preview_commission_state.clear();
if let (Some(id),Some(remaining))=(request.order_id,request.commission_remaining) {
preview_commission_state.insert(id,remaining);
}
self.cost_model.calculate_with_order_state_for_instrument(snapshot.date,OrderSide::Buy,
candidate_gross,request.order_id,&mut preview_commission_state,instrument).total()
} else {
self.cost_model.calculate_for_instrument(snapshot.date,OrderSide::Buy,candidate_gross,instrument).total()
};
let candidate_cash =
FixedMoney::checked_sum_f64([candidate_gross, candidate_cost])
.expect("buy cash must be finite fixed-point money")
@@ -8252,6 +8685,8 @@ fn sell_reason(decision: &StrategyDecision, symbol: &str) -> &'static str {
#[cfg(test)]
mod tests {
mod algorithm_clock;
use std::collections::BTreeMap;
use chrono::NaiveTime;
@@ -8291,6 +8726,9 @@ mod tests {
time_in_force: OrderTimeInForce::Gtc,
commission_remaining: None,
execution_cursor: None,
algo_request: None,
value_budget: None,
reserved_cash: None,
reason: format!("order_{order_id}"),
}
}
@@ -0,0 +1,778 @@
use super::*;
fn time(minute: u32) -> NaiveTime {
NaiveTime::from_hms_opt(10, minute, 0).unwrap()
}
fn data(quotes: &[(u32, f64, u32)]) -> DataSet {
data_with_snapshot(quotes, limit_test_snapshot())
}
fn data_with_snapshot(quotes: &[(u32, f64, u32)], snapshot: DailyMarketSnapshot) -> DataSet {
DataSet::from_components_with_actions_and_quotes(
vec![limit_test_instrument()],
vec![snapshot],
vec![],
vec![limit_test_candidate(true, true)],
vec![limit_test_benchmark()],
vec![],
quotes
.iter()
.map(|&(minute, price, volume)| {
let mut quote = limit_test_quote(price, price, price);
quote.timestamp = quote.date.and_time(time(minute));
quote.volume_delta = u64::from(volume);
quote.amount_delta = price * f64::from(volume);
quote.bid1_volume = u64::from(volume / 100);
quote.ask1_volume = u64::from(volume / 100);
quote
})
.collect(),
)
.unwrap()
}
fn broker() -> BrokerSimulator<ChinaAShareCostModel, ChinaEquityRuleHooks> {
BrokerSimulator::new(
ChinaAShareCostModel::default()
.with_commission_rate(0.0003)
.with_minimum_commission(5.),
ChinaEquityRuleHooks,
)
.with_matching_type(MatchingType::MinuteLast)
.with_execution_price_field(PriceField::Last)
.with_intraday_execution_start_time(time(0))
.with_volume_limit(true)
.with_volume_percent(0.25)
.with_liquidity_limit(false)
.with_inactive_limit(false)
.with_strict_value_budget(true)
}
fn intent(style: AlgoOrderStyle, value: f64) -> StrategyDecision {
StrategyDecision {
order_intents: vec![OrderIntent::AlgoValue {
symbol: "000001.SZ".into(),
value,
style,
start_time: Some(time(0)),
end_time: Some(time(10)),
reason: "clock-algorithm".into(),
}],
..Default::default()
}
}
fn step(
broker: &BrokerSimulator<ChinaAShareCostModel, ChinaEquityRuleHooks>,
portfolio: &mut PortfolioState,
data: &DataSet,
minute: u32,
decision: &StrategyDecision,
) -> BrokerExecutionReport {
broker
.execute_between(
limit_test_snapshot().date,
portfolio,
data,
decision,
Some(time(minute)),
Some(time(minute)),
)
.unwrap()
}
#[test]
fn twap_clock_preserves_quantity_prices_fees_budget_and_parent_order() {
let data = data(&[
(0, 10., 4_000),
(2, 10.1, 4_000),
(5, 10.2, 4_000),
(10, 10.3, 4_000),
]);
let decision = intent(AlgoOrderStyle::Twap, 10_000.);
let mut synchronous_account = PortfolioState::new(20_000.);
let reference = broker()
.execute(
limit_test_snapshot().date,
&mut synchronous_account,
&data,
&decision,
)
.unwrap();
let broker = broker();
let mut account = PortfolioState::new(20_000.);
let mut fills = Vec::new();
let mut events = Vec::new();
let empty = StrategyDecision::default();
for minute in [0, 2, 5, 10] {
let batch = step(
&broker,
&mut account,
&data,
minute,
if minute == 0 { &decision } else { &empty },
);
assert!(
batch
.fill_events
.iter()
.all(|fill| fill.execution_timestamp.unwrap().time() <= time(minute))
);
fills.extend(batch.fill_events);
events.extend(batch.order_events);
}
let canonical = |rows: &[crate::events::FillEvent]| {
rows.iter()
.map(|fill| {
(
fill.quantity,
fill.price.to_bits(),
fill.commission.to_bits(),
fill.stamp_tax.to_bits(),
fill.transfer_fee.to_bits(),
fill.execution_timestamp,
fill.order_id,
)
})
.collect::<Vec<_>>()
};
assert_eq!(canonical(&fills), canonical(&reference.fill_events));
assert_eq!(account.cash(), synchronous_account.cash());
assert_eq!(fills.iter().map(|fill| fill.quantity).sum::<u32>(), 900);
assert_eq!(fills.iter().map(|fill| fill.commission).sum::<f64>(), 5.);
assert!(fills.iter().map(|fill| -fill.net_cash_flow).sum::<f64>() <= 10_000.);
assert!(events.iter().all(|event| event.order_id == Some(1)));
assert_eq!(events.last().unwrap().status, OrderStatus::Filled);
assert!(broker.open_order_views().is_empty());
}
#[test]
fn partial_algorithm_cancel_releases_reservation_and_never_executes_the_remainder() {
let data = data(&[
(0, 10., 4_000),
(2, 10., 4_000),
(5, 10., 4_000),
(10, 10., 4_000),
]);
let broker = broker();
let mut account = PortfolioState::new(20_000.);
step(
&broker,
&mut account,
&data,
0,
&intent(AlgoOrderStyle::Twap, 10_000.),
);
assert_eq!(broker.open_order_views()[0].reserved_cash, Some(10_000.));
let partial = step(
&broker,
&mut account,
&data,
2,
&StrategyDecision::default(),
);
assert_eq!(
partial
.fill_events
.iter()
.map(|fill| fill.quantity)
.sum::<u32>(),
100
);
let working = broker.open_order_views();
assert_eq!(working[0].order_id, 1);
assert_eq!(working[0].filled_quantity, 100);
assert_eq!(
working[0].reserved_cash,
Some(10_000. + partial.fill_events[0].net_cash_flow)
);
let cancel = step(
&broker,
&mut account,
&data,
3,
&StrategyDecision {
order_intents: vec![OrderIntent::CancelAll {
reason: "explicit-user-cancel".into(),
}],
..Default::default()
},
);
assert!(cancel.fill_events.is_empty());
assert_eq!(
cancel.order_events.last().unwrap().status,
OrderStatus::Canceled
);
assert_eq!(cancel.order_events.last().unwrap().filled_quantity, 100);
assert!(broker.open_order_views().is_empty());
assert!(
step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default()
)
.fill_events
.is_empty()
);
assert_eq!(account.position("000001.SZ").unwrap().quantity, 100);
}
#[test]
fn algorithm_expiry_without_a_quote_does_not_reuse_old_liquidity() {
let data = data(&[(0, 10., 4_000), (2, 10., 4_000)]);
let broker = broker();
let mut account = PortfolioState::new(20_000.);
step(
&broker,
&mut account,
&data,
0,
&intent(AlgoOrderStyle::Twap, 10_000.),
);
step(
&broker,
&mut account,
&data,
2,
&StrategyDecision::default(),
);
assert_eq!(
broker.next_day_order_expiry(limit_test_snapshot().date),
Some(time(10))
);
let terminal = step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default(),
);
assert!(terminal.fill_events.is_empty());
assert_eq!(
terminal.order_events.last().unwrap().status,
OrderStatus::Expired
);
assert_eq!(terminal.order_events.last().unwrap().filled_quantity, 100);
assert!(
terminal
.process_events
.iter()
.any(|event| event.detail.contains("Expired")),
"{:?}",
terminal.process_events
);
assert!(broker.open_order_views().is_empty());
}
#[test]
fn separate_buy_cannot_spend_the_working_algorithms_cash_budget() {
let data = data(&[
(0, 10., 4_000),
(1, 10., 4_000),
(2, 10., 4_000),
(10, 10., 4_000),
]);
let broker = broker();
let mut account = PortfolioState::new(11_000.);
step(
&broker,
&mut account,
&data,
0,
&intent(AlgoOrderStyle::Twap, 10_000.),
);
let other = step(
&broker,
&mut account,
&data,
1,
&StrategyDecision {
order_intents: vec![OrderIntent::Shares {
symbol: "000001.SZ".into(),
quantity: 1_000,
reason: "separate-buy".into(),
}],
..Default::default()
},
);
assert!(
other.fill_events.is_empty(),
"cash reserved for order 1 was spent: {:?}",
other.fill_events
);
let final_batch = step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default(),
);
assert!(
final_batch
.fill_events
.iter()
.all(|fill| fill.order_id == Some(1))
);
assert_eq!(account.position("000001.SZ").unwrap().quantity, 900);
assert!(account.cash() >= 1_000.);
}
#[test]
fn changing_the_later_daily_close_does_not_resize_an_algorithm_submitted_now() {
let quotes = [(0, 10., 4_000), (2, 10.1, 4_000), (10, 10.2, 4_000)];
let mut changed = limit_test_snapshot();
changed.close = 100.;
changed.last_price = 100.;
let run = |data: DataSet| {
let broker = broker();
let mut account = PortfolioState::new(20_000.);
let initial = step(
&broker,
&mut account,
&data,
0,
&intent(AlgoOrderStyle::Twap, 10_000.),
);
assert!(initial.fill_events.is_empty());
let quantity = broker.open_order_views()[0].requested_quantity;
let final_batch = step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default(),
);
(
quantity,
final_batch
.fill_events
.iter()
.map(|fill| {
(
fill.quantity,
fill.price.to_bits(),
fill.net_cash_flow.to_bits(),
)
})
.collect::<Vec<_>>(),
)
};
assert_eq!(
run(data(&quotes)),
run(data_with_snapshot(&quotes, changed))
);
}
#[test]
fn vwap_clock_preserves_cash_costs_and_does_not_spend_future_volume() {
let data = data(&[
(0, 10., 400),
(2, 10., 800),
(5, 10., 1_200),
(10, 10., 4_000),
]);
let decision = intent(AlgoOrderStyle::Vwap, 10_000.);
let mut synchronous_account = PortfolioState::new(20_000.);
let reference = broker()
.execute(
limit_test_snapshot().date,
&mut synchronous_account,
&data,
&decision,
)
.unwrap();
let broker = broker();
let mut account = PortfolioState::new(20_000.);
let empty = StrategyDecision::default();
let mut filled = 0;
let mut commission = 0.;
for (minute, expected) in [(0, 100), (2, 300), (5, 600), (10, 900)] {
let batch = step(
&broker,
&mut account,
&data,
minute,
if minute == 0 { &decision } else { &empty },
);
filled += batch
.fill_events
.iter()
.map(|fill| fill.quantity)
.sum::<u32>();
commission += batch
.fill_events
.iter()
.map(|fill| fill.commission)
.sum::<f64>();
assert_eq!(filled, expected);
assert!(batch.fill_events.iter().all(|fill| fill.order_id == Some(1)
&& fill.execution_timestamp.unwrap().time() <= time(minute)));
}
assert_eq!(account.cash(), synchronous_account.cash());
assert_eq!(
commission,
reference
.fill_events
.iter()
.map(|fill| fill.commission)
.sum::<f64>()
);
assert!(broker.open_order_views().is_empty());
}
#[test]
fn global_vwap_matching_keeps_the_same_working_order_between_clock_ticks() {
let data = data(&[(0, 10., 400), (2, 10., 400), (10, 10., 4_000)]);
let broker = broker().with_matching_type(MatchingType::Vwap);
let mut account = PortfolioState::new(20_000.);
let first = step(
&broker,
&mut account,
&data,
0,
&StrategyDecision {
order_intents: vec![OrderIntent::Shares {
symbol: "000001.SZ".into(),
quantity: 900,
reason: "configured-vwap".into(),
}],
..Default::default()
},
);
assert_eq!(
first
.fill_events
.iter()
.map(|fill| fill.quantity)
.sum::<u32>(),
100
);
assert_eq!(
broker.open_order_views().len(),
1,
"{:?}",
first.order_events
);
let second = step(
&broker,
&mut account,
&data,
2,
&StrategyDecision::default(),
);
assert_eq!(second.fill_events[0].quantity, 100);
assert_eq!(second.fill_events[0].order_id, Some(1));
let final_batch = step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default(),
);
assert_eq!(final_batch.fill_events[0].quantity, 700);
assert_eq!(final_batch.fill_events[0].order_id, Some(1));
assert!(broker.open_order_views().is_empty());
}
#[test]
fn algorithm_sell_honors_t_plus_one_and_keeps_original_quantity_after_partial_fills() {
let data = data(&[(0, 10., 400), (2, 10., 800), (10, 10., 4_000)]);
let date = limit_test_snapshot().date;
for acquired_today in [false, true] {
let broker = broker();
let mut account = PortfolioState::new(20_000.);
account.position_mut("000001.SZ").buy(
if acquired_today {
date
} else {
date.pred_opt().unwrap()
},
1_000,
10.,
);
let decision = intent(AlgoOrderStyle::Vwap, -10_000.);
let mut fills = Vec::new();
let mut events = Vec::new();
let empty = StrategyDecision::default();
for minute in [0, 2, 10] {
let batch = step(
&broker,
&mut account,
&data,
minute,
if minute == 0 { &decision } else { &empty },
);
fills.extend(batch.fill_events);
events.extend(batch.order_events);
}
assert_eq!(
fills.iter().map(|fill| fill.quantity).sum::<u32>(),
if acquired_today { 0 } else { 1_000 }
);
assert!(events.iter().all(|event| event.order_id == Some(1)));
if !acquired_today {
assert_eq!(events.last().unwrap().status, OrderStatus::Filled);
assert_eq!(events.last().unwrap().requested_quantity, 1_000);
assert_eq!(events.last().unwrap().filled_quantity, 1_000);
}
assert!(broker.open_order_views().is_empty());
}
}
#[test]
fn an_explicit_ioc_or_fok_does_not_become_a_persistent_algorithm() {
let data = data(&[(0, 10., 400), (2, 10., 4_000), (10, 10., 4_000)]);
for tif in [
OrderTimeInForce::Ioc,
OrderTimeInForce::Fok,
OrderTimeInForce::Day,
OrderTimeInForce::Gtc,
] {
let broker = broker();
let mut account = PortfolioState::new(20_000.);
let mut decision = intent(AlgoOrderStyle::Vwap, 10_000.);
if !decision.order_intents[0].supports_time_in_force(tif) {
decision.order_intents = decision
.order_intents
.into_iter()
.map(|intent| intent.with_time_in_force(tif))
.collect();
let error = broker
.execute_between(
limit_test_snapshot().date,
&mut account,
&data,
&decision,
Some(time(0)),
Some(time(0)),
)
.unwrap_err();
assert!(
error
.to_string()
.contains("is not supported for this order intent")
);
assert_eq!(account.cash(), 20_000.);
assert!(broker.open_order_views().is_empty());
continue;
}
decision.order_intents = decision
.order_intents
.into_iter()
.map(|intent| intent.with_time_in_force(tif))
.collect();
let first = step(&broker, &mut account, &data, 0, &decision);
let persists = matches!(tif, OrderTimeInForce::Day | OrderTimeInForce::Gtc);
assert_eq!(
!broker.open_order_views().is_empty(),
persists,
"{tif:?}: {:?}",
first.order_events
);
if !persists {
assert!(
step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default()
)
.fill_events
.is_empty()
);
}
}
}
#[test]
fn two_working_algorithms_reserve_only_real_cash_without_starving_the_first() {
let data = data(&[(0, 10., 40_000), (10, 10., 40_000)]);
let broker = broker();
let mut account = PortfolioState::new(15_000.);
let mut decision = intent(AlgoOrderStyle::Twap, 10_000.);
decision
.order_intents
.extend(intent(AlgoOrderStyle::Twap, 10_000.).order_intents);
step(&broker, &mut account, &data, 0, &decision);
assert_eq!(
broker
.open_order_views()
.iter()
.map(|order| order.reserved_cash.unwrap())
.collect::<Vec<_>>(),
vec![10_000., 5_000.]
);
let report = step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default(),
);
assert_eq!(
report
.fill_events
.iter()
.map(|fill| (fill.order_id, fill.quantity))
.collect::<Vec<_>>(),
vec![(Some(1), 900), (Some(2), 500)]
);
assert!(account.cash() >= 0.);
assert!(broker.open_order_views().is_empty());
}
#[test]
fn a_clock_slice_does_not_turn_window_twap_into_an_unlimited_instant_order() {
let data = data(&[(0, 10., 100), (2, 10., 100), (10, 10.1, 100)]);
let broker = broker()
.with_volume_limit(false)
.with_liquidity_limit(false);
let mut account = PortfolioState::new(20_000.);
step(
&broker,
&mut account,
&data,
0,
&intent(AlgoOrderStyle::Twap, 10_000.),
);
let first = step(
&broker,
&mut account,
&data,
2,
&StrategyDecision::default(),
);
let last = step(
&broker,
&mut account,
&data,
10,
&StrategyDecision::default(),
);
assert_eq!(
first
.fill_events
.iter()
.map(|fill| fill.quantity)
.sum::<u32>(),
100
);
assert_eq!(
last.fill_events
.iter()
.map(|fill| fill.quantity)
.sum::<u32>(),
100
);
assert_eq!(
last.order_events.last().unwrap().status,
OrderStatus::Expired
);
assert_eq!(last.order_events.last().unwrap().filled_quantity, 200);
assert!(broker.open_order_views().is_empty());
}
#[test]
fn non_matching_controls_amend_or_cancel_without_filling_a_crossing_quote() {
let data = data(&[(0, 10., 4_000), (2, 9.4, 4_000)]);
let broker = broker();
let mut account = PortfolioState::new(20_000.);
step(
&broker,
&mut account,
&data,
0,
&StrategyDecision {
order_intents: vec![
OrderIntent::LimitShares {
symbol: "000001.SZ".into(),
quantity: 100,
limit_price: 9.5,
reason: "resting".into(),
}
.with_time_in_force(OrderTimeInForce::Gtc),
],
..Default::default()
},
);
assert_eq!(broker.open_order_views().len(), 1);
let modify = broker
.execute_controls_without_matching(
limit_test_snapshot().date,
limit_test_snapshot().date,
&mut account,
&data,
&StrategyDecision {
order_intents: vec![OrderIntent::ModifyOrder {
order_id: 1,
new_total_quantity: Some(200),
new_limit_price: Some(9.3),
reason: "pre-open-amend".into(),
}],
..Default::default()
},
Some(time(2)),
)
.unwrap();
assert!(modify.fill_events.is_empty());
assert_eq!(broker.open_order_views()[0].limit_price, 9.3);
assert_eq!(broker.open_order_views()[0].requested_quantity, 200);
let cancel = broker
.execute_controls_without_matching(
limit_test_snapshot().date,
limit_test_snapshot().date,
&mut account,
&data,
&StrategyDecision {
order_intents: vec![OrderIntent::CancelAll {
reason: "pre-open-cancel".into(),
}],
..Default::default()
},
Some(time(2)),
)
.unwrap();
assert!(cancel.fill_events.is_empty());
assert_eq!(
cancel.order_events.last().unwrap().status,
OrderStatus::Canceled
);
assert_eq!(account.cash(), 20_000.);
assert!(broker.open_order_views().is_empty());
}
#[test]
fn control_only_phase_cannot_be_used_to_submit_an_order_or_leave_matching_disabled() {
let data = data(&[(0, 10., 4_000)]);
let broker = broker();
let mut account = PortfolioState::new(20_000.);
let submit = StrategyDecision {
order_intents: vec![OrderIntent::Shares {
symbol: "000001.SZ".into(),
quantity: 100,
reason: "normal-order".into(),
}],
..Default::default()
};
assert!(
broker
.execute_controls_without_matching(
limit_test_snapshot().date,
limit_test_snapshot().date,
&mut account,
&data,
&submit,
Some(time(0))
)
.is_err()
);
assert_eq!(account.cash(), 20_000.);
assert_eq!(
step(&broker, &mut account, &data, 0, &submit).fill_events[0].quantity,
100
);
}
File diff suppressed because it is too large Load Diff
+11
View File
@@ -28,6 +28,17 @@ impl FixedMoney {
self.0
}
pub fn to_decimal_string(self) -> String {
let magnitude = self.0.unsigned_abs();
let scale = MONEY_SCALE as u128;
let sign = if self.0 < 0 { "-" } else { "" };
let width = MONEY_SCALE.ilog10() as usize;
format!("{sign}{}.{:0width$}", magnitude / scale, magnitude % scale)
.trim_end_matches('0')
.trim_end_matches('.')
.to_string()
}
pub fn from_decimal_str(value: &str) -> Result<Self, String> {
let value = value.trim();
if value.is_empty() {
+1
View File
@@ -20,6 +20,7 @@ pub mod fixed_point;
pub mod futures;
pub mod instrument;
pub mod metrics;
pub mod manual_execution;
mod numeric_expr_vm;
pub mod platform_expr_strategy;
pub mod platform_runtime_schema;
+534
View File
@@ -0,0 +1,534 @@
//! Confirmed manual fills are external observations, not simulated broker fills.
//! The producer must bind these records to the runtime's durable order/audit facts.
use std::collections::BTreeSet;
use chrono::{DateTime, FixedOffset, NaiveDate, Timelike, Utc};
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::events::OrderSide;
use crate::{DataSet, FixedMoney, PortfolioState};
use rust_decimal::prelude::ToPrimitive;
pub const MANUAL_REPLAY_SCHEMA: &str = "fidc.observed-manual-executions/v1";
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub struct ManualExecutionReplay {
pub schema: String,
pub runtime_id: String,
pub account_id: String,
pub source_contract_sha256: String,
pub content_sha256: String,
pub observation_cutoff: DateTime<Utc>,
pub actions: Vec<ManualExecutionAction>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub struct ManualExecutionAction {
pub action_id: String,
pub source: ManualExecutionSource,
pub audit_event_ids: Vec<String>,
pub confirmed_at: DateTime<Utc>,
pub outcome: ManualActionOutcome,
pub orders: Vec<ManualExecutionOrder>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum ManualActionOutcome {
NoOrdersNeeded,
OrdersTerminal,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum ManualExecutionSource {
ManualSecurityTrade,
ManualPositionAction,
ManualRebalance,
StockPoolAllocation,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub struct ManualExecutionOrder {
pub order_id: String,
pub broker_order_id: Option<String>,
pub source_adapter: String,
pub symbol: String,
pub side: OrderSide,
pub quantity: u32,
pub submitted_at: DateTime<Utc>,
pub terminal_at: DateTime<Utc>,
pub terminal_status: ManualOrderTerminalStatus,
pub fills: Vec<ManualExecutionFill>,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum ManualOrderTerminalStatus {
Filled,
Cancelled,
Rejected,
Expired,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub struct ManualExecutionFill {
pub trade_id: String,
pub observation_event_id: String,
pub observation_sequence: u64,
pub trade_date: NaiveDate,
pub executed_at: DateTime<Utc>,
pub observed_at: DateTime<Utc>,
pub timestamp_precision: ManualTimestampPrecision,
pub quantity: u32,
#[serde(with = "rust_decimal::serde::str")]
pub price: Decimal,
#[serde(with = "rust_decimal::serde::str")]
pub commission: Decimal,
#[serde(with = "rust_decimal::serde::str")]
pub stamp_tax: Decimal,
#[serde(with = "rust_decimal::serde::str")]
pub transfer_fee: Decimal,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum ManualTimestampPrecision {
Second,
Millisecond,
Microsecond,
Nanosecond,
}
impl ManualTimestampPrecision {
fn nanoseconds(self) -> i64 {
match self {
Self::Second => 1_000_000_000,
Self::Millisecond => 1_000_000,
Self::Microsecond => 1_000,
Self::Nanosecond => 1,
}
}
}
impl ManualExecutionFill {
pub fn gross_amount(&self) -> Result<Decimal, String> {
self.price
.checked_mul(Decimal::from(self.quantity))
.ok_or_else(|| "manual fill gross amount overflow".into())
}
pub fn total_fees(&self) -> Result<Decimal, String> {
self.commission
.checked_add(self.stamp_tax)
.and_then(|sum| sum.checked_add(self.transfer_fee))
.ok_or_else(|| "manual fill fees overflow".into())
}
}
fn identifier(value: &str) -> Result<(), String> {
if value.is_empty()
|| value.trim() != value
|| value.len() > 256
|| value.chars().any(char::is_control)
{
return Err("manual execution identity is empty, untrimmed or invalid".into());
}
Ok(())
}
impl ManualExecutionReplay {
pub fn observations(&self) -> Result<Vec<ManualFillObservation<'_>>, String> {
self.validate()?;
let mut observations = Vec::new();
for action in &self.actions {
for order in &action.orders {
for fill in &order.fills {
observations.push(ManualFillObservation {
action,
order,
fill,
});
}
}
}
observations.sort_by_key(|entry| (entry.fill.observed_at, entry.fill.observation_sequence));
Ok(observations)
}
pub fn content_digest(&self) -> Result<String, String> {
let mut value = serde_json::to_value(self).map_err(|error| error.to_string())?;
value
.as_object_mut()
.ok_or("manual replay is not an object")?
.remove("contentSha256");
let bytes = serde_json::to_vec(&value).map_err(|error| error.to_string())?;
Ok(format!("{:x}", Sha256::digest(bytes)))
}
pub fn validate(&self) -> Result<(), String> {
if self.schema != MANUAL_REPLAY_SCHEMA {
return Err("unsupported manual replay schema".into());
}
identifier(&self.runtime_id)?;
identifier(&self.account_id)?;
if self.source_contract_sha256.len() != 64
|| !self
.source_contract_sha256
.bytes()
.all(|v| v.is_ascii_hexdigit())
{
return Err("manual replay source contract hash is invalid".into());
}
if self.content_digest()? != self.content_sha256 {
return Err("manual replay content digest mismatch".into());
}
if self.actions.len() > 100_000 {
return Err("manual replay action limit exceeded; trace was not truncated".into());
}
let shanghai = FixedOffset::east_opt(8 * 3600).unwrap();
let mut actions = BTreeSet::new();
let mut audits = BTreeSet::new();
let mut orders = BTreeSet::new();
let mut broker_orders = BTreeSet::new();
let mut trades = BTreeSet::new();
let mut observation_events = BTreeSet::new();
let mut observation_sequences = BTreeSet::new();
for action in &self.actions {
identifier(&action.action_id)?;
if !actions.insert(action.action_id.as_str())
|| action.confirmed_at > self.observation_cutoff
{
return Err("duplicate manual action or confirmation after cutoff".into());
}
if action.audit_event_ids.is_empty() {
return Err("manual action has no immutable audit binding".into());
}
if (action.outcome == ManualActionOutcome::NoOrdersNeeded) != action.orders.is_empty() {
return Err("manual action outcome does not prove its order coverage".into());
}
for id in &action.audit_event_ids {
identifier(id)?;
if !audits.insert(id.as_str()) {
return Err("manual audit event is bound more than once".into());
}
}
for order in &action.orders {
identifier(&order.order_id)?;
identifier(&order.source_adapter)?;
identifier(&order.symbol)?;
if let Some(id) = &order.broker_order_id {
identifier(id)?;
if !broker_orders.insert((
order.source_adapter.as_str(),
order.submitted_at.with_timezone(&shanghai).date_naive(),
id.as_str(),
)) {
return Err("manual local orders share one broker order identity".into());
}
}
if !order.fills.is_empty()
&& order.source_adapter != "paper"
&& order.broker_order_id.is_none()
{
return Err(
"manual broker fills require their original broker order identity".into(),
);
}
if !orders.insert(order.order_id.as_str())
|| order.quantity == 0
|| order.quantity > i32::MAX as u32
{
return Err("duplicate manual order or invalid quantity".into());
}
if order.submitted_at < action.confirmed_at
|| order.terminal_at < order.submitted_at
|| order.terminal_at > self.observation_cutoff
{
return Err(
"manual order confirmation/submission/terminal time is inconsistent".into(),
);
}
let mut filled = 0_u32;
for fill in &order.fills {
identifier(&fill.trade_id)?;
identifier(&fill.observation_event_id)?;
if fill.observation_sequence == 0
|| fill.observation_sequence > i64::MAX as u64
|| !observation_events.insert(fill.observation_event_id.as_str())
|| !observation_sequences.insert(fill.observation_sequence)
{
return Err(
"manual fill requires a unique durable observation event and sequence"
.into(),
);
}
if !trades.insert((fill.trade_date, fill.trade_id.as_str()))
|| fill.quantity == 0
{
return Err("duplicate manual trade or zero fill quantity".into());
}
if fill.executed_at.with_timezone(&shanghai).date_naive() != fill.trade_date
|| fill.observed_at > self.observation_cutoff
|| fill.observed_at < order.submitted_at
|| fill.observed_at < fill.executed_at
|| fill.executed_at > order.terminal_at
{
return Err("manual fill execution/observation time is inconsistent".into());
}
if i64::from(fill.executed_at.nanosecond())
% fill.timestamp_precision.nanoseconds()
!= 0
{
return Err(
"broker timestamp contains digits finer than its declared precision"
.into(),
);
}
let upper = fill
.executed_at
.checked_add_signed(chrono::Duration::nanoseconds(
fill.timestamp_precision.nanoseconds(),
))
.ok_or("manual execution timestamp overflow")?;
if fill.executed_at < order.submitted_at && order.submitted_at >= upper {
return Err("manual fill predates its submitted order".into());
}
if fill.price <= Decimal::ZERO
|| [fill.commission, fill.stamp_tax, fill.transfer_fee]
.iter()
.any(|fee| *fee < Decimal::ZERO)
{
return Err(
"manual fill requires a positive price and complete nonnegative fees"
.into(),
);
}
fill.gross_amount()?
.checked_add(fill.total_fees()?)
.ok_or("manual fill cash amount overflow")?;
filled = filled
.checked_add(fill.quantity)
.ok_or("manual cumulative fill quantity overflow")?;
}
if filled > order.quantity
|| (order.terminal_status == ManualOrderTerminalStatus::Filled
&& filled != order.quantity)
|| (order.terminal_status == ManualOrderTerminalStatus::Rejected && filled != 0)
|| (matches!(
order.terminal_status,
ManualOrderTerminalStatus::Cancelled | ManualOrderTerminalStatus::Expired
) && filled == order.quantity)
{
return Err("manual terminal status disagrees with cumulative fills".into());
}
}
}
Ok(())
}
}
#[derive(Debug, Clone, Copy)]
pub struct ManualFillObservation<'a> {
pub action: &'a ManualExecutionAction,
pub order: &'a ManualExecutionOrder,
pub fill: &'a ManualExecutionFill,
}
#[derive(Debug, Clone, PartialEq)]
pub struct AppliedManualFill {
pub gross: FixedMoney,
pub fees: FixedMoney,
pub cash_delta: FixedMoney,
pub quantity_after: u32,
}
/// One replay owns its immutable trace and progress. Advancing is atomic even
/// if a later receipt in the same step disagrees with the shadow account.
pub struct ManualReplayCursor {
replay: ManualExecutionReplay,
indices: Vec<(usize, usize, usize)>,
cursor: usize,
clock: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "camelCase")]
pub struct ManualReplayApplication {
pub action_id: String,
pub order_id: String,
pub trade_id: String,
pub observation_event_id: String,
pub observation_sequence: u64,
pub observed_at: DateTime<Utc>,
pub executed_at: DateTime<Utc>,
pub symbol: String,
pub side: OrderSide,
pub quantity: u32,
pub quantity_after: u32,
pub price: String,
pub commission: String,
pub stamp_tax: String,
pub transfer_fee: String,
pub source_gross_amount: String,
pub ledger_gross_amount: String,
pub ledger_fees: String,
pub cash_delta: String,
}
impl ManualReplayCursor {
pub fn new(replay: ManualExecutionReplay) -> Result<Self, String> {
replay.validate()?;
let mut indices = Vec::new();
for (a, action) in replay.actions.iter().enumerate() {
for (o, order) in action.orders.iter().enumerate() {
for f in 0..order.fills.len() {
indices.push((a, o, f));
}
}
}
indices.sort_by_key(|&(a, o, f)| {
let fill = &replay.actions[a].orders[o].fills[f];
(fill.observed_at, fill.observation_sequence)
});
Ok(Self {
replay,
indices,
cursor: 0,
clock: None,
})
}
pub fn next_observation_at(&self) -> Option<DateTime<Utc>> {
self.indices
.get(self.cursor)
.map(|&(a, o, f)| self.replay.actions[a].orders[o].fills[f].observed_at)
}
pub fn applied_count(&self) -> usize {
self.cursor
}
pub fn advance(
&mut self,
at: DateTime<Utc>,
portfolio: &mut PortfolioState,
data: &DataSet,
has_pending_orders: bool,
) -> Result<Vec<ManualReplayApplication>, String> {
if at > self.replay.observation_cutoff {
return Err("manual observation clock exceeds the frozen evidence cutoff".into());
}
if self.clock.is_some_and(|clock| at < clock) {
return Err("manual observation clock moved backwards".into());
}
let end = self.cursor
+ self.indices[self.cursor..]
.iter()
.take_while(|&&(a, o, f)| {
self.replay.actions[a].orders[o].fills[f].observed_at <= at
})
.count();
if end == self.cursor {
self.clock = Some(at);
return Ok(vec![]);
}
let mut next = portfolio.clone();
let mut applications = Vec::with_capacity(end - self.cursor);
for &(a, o, f) in &self.indices[self.cursor..end] {
let action = &self.replay.actions[a];
let order = &action.orders[o];
let fill = &order.fills[f];
let applied = ManualFillObservation {
action,
order,
fill,
}
.apply(&mut next, data, has_pending_orders)?;
applications.push(ManualReplayApplication {
action_id: action.action_id.clone(),
order_id: order.order_id.clone(),
trade_id: fill.trade_id.clone(),
observation_event_id: fill.observation_event_id.clone(),
observation_sequence: fill.observation_sequence,
observed_at: fill.observed_at,
executed_at: fill.executed_at,
symbol: order.symbol.clone(),
side: order.side,
quantity: fill.quantity,
quantity_after: applied.quantity_after,
price: fill.price.to_string(),
commission: fill.commission.to_string(),
stamp_tax: fill.stamp_tax.to_string(),
transfer_fee: fill.transfer_fee.to_string(),
source_gross_amount: fill.gross_amount()?.to_string(),
ledger_gross_amount: applied.gross.to_decimal_string(),
ledger_fees: applied.fees.to_decimal_string(),
cash_delta: applied.cash_delta.to_decimal_string(),
});
}
*portfolio = next;
self.cursor = end;
self.clock = Some(at);
Ok(applications)
}
}
impl ManualFillObservation<'_> {
pub(crate) fn apply(
&self,
portfolio: &mut PortfolioState,
data: &DataSet,
has_pending_orders: bool,
) -> Result<AppliedManualFill, String> {
if has_pending_orders {
return Err("manual observation conflicts with pending shadow orders".into());
}
let instrument = data
.instrument(&self.order.symbol)
.ok_or("manual observation instrument is absent from frozen source data")?;
if instrument
.dated_market_absence_reason(self.fill.trade_date)
.is_some()
{
return Err("manual execution contradicts the frozen instrument lifecycle".into());
}
let gross = FixedMoney::from_decimal_str(&self.fill.gross_amount()?.to_string())?;
let fees = FixedMoney::from_decimal_str(&self.fill.total_fees()?.to_string())?;
let price = self
.fill
.price
.to_f64()
.filter(|price| price.is_finite() && *price > 0.)
.ok_or("manual execution price cannot be represented for valuation")?;
// This is the real observed trade price, not a fabricated quote. The
// normal market clock remains responsible for subsequent marks.
let cash_delta = portfolio.apply_observed_manual_fill(
self.fill.trade_date,
&self.order.symbol,
self.order.side,
self.fill.quantity,
price,
price,
gross,
fees,
)?;
Ok(AppliedManualFill {
gross,
fees,
cash_delta,
quantity_after: portfolio
.position(&self.order.symbol)
.map_or(0, |position| position.quantity),
})
}
}
#[cfg(test)]
mod tests;
@@ -0,0 +1,416 @@
use super::*;
use serde_json::{Value, json};
fn sample() -> ManualExecutionReplay {
let mut input:ManualExecutionReplay=serde_json::from_value(json!({
"schema":MANUAL_REPLAY_SCHEMA,"runtimeId":"runtime-1","accountId":"account-1",
"sourceContractSha256":"a".repeat(64),"contentSha256":"", "observationCutoff":"2026-09-14T08:00:00Z",
"actions":[{"actionId":"action-1","source":"manual_security_trade","auditEventIds":["audit-1"],
"confirmedAt":"2026-09-14T01:30:00.500Z","outcome":"orders_terminal","orders":[{
"orderId":"order-1","brokerOrderId":"broker-1","sourceAdapter":"gt-api","symbol":"000001.SZ","side":"Buy","quantity":100,
"submittedAt":"2026-09-14T01:30:00.600Z","terminalAt":"2026-09-14T01:30:00.900Z","terminalStatus":"filled",
"fills":[{"tradeId":"trade-1","observationEventId":"received-1","observationSequence":1,"tradeDate":"2026-09-14","executedAt":"2026-09-14T01:30:00Z",
"observedAt":"2026-09-14T01:30:01Z","timestampPrecision":"second","quantity":100,
"price":"10.1234567891","commission":"0.1000001","stampTax":"0","transferFee":"0.02"}]
}]}]
})).unwrap();
reseal(&mut input);
input
}
fn reseal(input: &mut ManualExecutionReplay) {
input.content_sha256 = input.content_digest().unwrap();
}
fn semantic_result(input: &ManualExecutionReplay) -> Result<(), String> {
let mut input = input.clone();
reseal(&mut input);
input.validate()
}
#[test]
fn complete_exact_decimal_evidence_allows_later_observation_and_retains_source_digits() {
let input = sample();
input.validate().unwrap();
let fill = &input.actions[0].orders[0].fills[0];
assert_eq!(fill.gross_amount().unwrap().to_string(), "1012.3456789100");
assert_eq!(fill.total_fees().unwrap().to_string(), "0.1200001");
assert_eq!(
serde_json::to_value(&input).unwrap()["actions"][0]["orders"][0]["fills"][0]["price"],
"10.1234567891"
);
}
#[test]
fn all_required_money_and_binding_fields_reject_missing_or_wrong_values() {
let original = serde_json::to_value(sample()).unwrap();
for field in ["price", "commission", "stampTax", "transferFee"] {
let mut missing = original.clone();
missing["actions"][0]["orders"][0]["fills"][0]
.as_object_mut()
.unwrap()
.remove(field);
assert!(
serde_json::from_value::<ManualExecutionReplay>(missing).is_err(),
"{field}"
);
let mut numeric = original.clone();
numeric["actions"][0]["orders"][0]["fills"][0][field] = json!(1.1);
assert!(
serde_json::from_value::<ManualExecutionReplay>(numeric).is_err(),
"numeric {field}"
);
}
for mutate in [
("schema", json!("unknown")),
("sourceContractSha256", json!("broken")),
("accountId", json!(" ")),
] {
let mut value = original.clone();
value[mutate.0] = mutate.1;
assert!(
semantic_result(&serde_json::from_value::<ManualExecutionReplay>(value).unwrap())
.is_err()
);
}
}
#[test]
fn inconsistent_counts_terminals_audits_and_duplicate_facts_are_rejected() {
let original = sample();
let mut invalid = original.clone();
invalid.actions[0].orders[0].quantity = 200;
assert!(semantic_result(&invalid).is_err());
let mut invalid = original.clone();
invalid.actions[0].orders[0].terminal_status = ManualOrderTerminalStatus::Rejected;
assert!(semantic_result(&invalid).is_err());
let mut invalid = original.clone();
invalid.actions[0].audit_event_ids.clear();
assert!(semantic_result(&invalid).is_err());
let mut invalid = original.clone();
invalid.actions.push(invalid.actions[0].clone());
assert!(semantic_result(&invalid).is_err());
let mut invalid = original.clone();
let duplicate = invalid.actions[0].orders[0].fills[0].clone();
invalid.actions[0].orders[0].fills.push(duplicate);
assert!(semantic_result(&invalid).is_err());
let mut invalid = original.clone();
invalid.actions[0].orders[0].broker_order_id = None;
assert!(semantic_result(&invalid).is_err());
invalid.actions[0].orders[0].source_adapter = "paper".into();
reseal(&mut invalid);
invalid.validate().unwrap();
}
#[test]
fn source_time_precision_is_not_invented_and_submitted_time_must_fit_the_interval() {
let mut input = sample();
input.actions[0].orders[0].submitted_at = "2026-09-14T01:30:00.999999Z".parse().unwrap();
input.actions[0].orders[0].terminal_at = "2026-09-14T01:30:01.500Z".parse().unwrap();
input.actions[0].orders[0].fills[0].observed_at = "2026-09-14T01:30:02Z".parse().unwrap();
reseal(&mut input);
input.validate().unwrap();
input.actions[0].orders[0].submitted_at = "2026-09-14T01:30:01Z".parse().unwrap();
assert!(semantic_result(&input).is_err());
let mut input = sample();
input.actions[0].orders[0].fills[0].executed_at = "2026-09-14T01:30:00.800Z".parse().unwrap();
assert!(semantic_result(&input).is_err());
input.actions[0].orders[0].fills[0].timestamp_precision = ManualTimestampPrecision::Millisecond;
reseal(&mut input);
input.validate().unwrap();
input.actions[0].orders[0].fills[0].executed_at =
"2026-09-14T01:30:00.800001Z".parse().unwrap();
assert!(semantic_result(&input).is_err());
}
#[test]
fn confirmed_no_order_outcome_is_distinct_from_unconfirmed_or_unknown_work() {
let mut input = sample();
input.actions[0].orders.clear();
assert!(semantic_result(&input).is_err());
input.actions[0].outcome = ManualActionOutcome::NoOrdersNeeded;
reseal(&mut input);
input.validate().unwrap();
let mut value = serde_json::to_value(input).unwrap();
value["actions"][0]["outcome"] = json!("result_unknown");
assert!(serde_json::from_value::<ManualExecutionReplay>(value).is_err());
}
#[test]
fn raw_timezone_and_cutoff_are_required() {
let mut value = serde_json::to_value(sample()).unwrap();
value["actions"][0]["orders"][0]["fills"][0]["executedAt"] = json!("2026-09-14T09:30:00");
assert!(serde_json::from_value::<ManualExecutionReplay>(value).is_err());
let mut input = sample();
input.observation_cutoff = "2026-09-14T01:30:00.700Z".parse().unwrap();
assert!(semantic_result(&input).is_err());
let mut value = serde_json::to_value(sample()).unwrap();
value["actions"][0]["orders"][0]["fills"][0]["commission"] = Value::Null;
assert!(serde_json::from_value::<ManualExecutionReplay>(value).is_err());
}
#[test]
fn changing_any_external_price_or_identity_invalidates_the_frozen_trace() {
let input = sample();
let original = input.content_sha256.clone();
let mut changed = input.clone();
changed.actions[0].orders[0].fills[0].price += Decimal::ONE;
assert_ne!(changed.content_digest().unwrap(), original);
assert_eq!(
changed.validate().unwrap_err(),
"manual replay content digest mismatch"
);
let mut changed = input;
changed.account_id = "another-account".into();
assert_ne!(changed.content_digest().unwrap(), original);
assert!(changed.validate().is_err());
}
fn identity_data(listed: NaiveDate) -> DataSet {
DataSet::from_components(
vec![crate::Instrument {
symbol: "000001.SZ".into(),
name: "test".into(),
board: "SZ".into(),
round_lot: 100,
listed_at: Some(listed),
delisted_at: None,
status: "active".into(),
}],
vec![],
vec![],
vec![],
vec![crate::BenchmarkSnapshot {
date: listed,
benchmark: "000300.SH".into(),
open: 100.,
close: 100.,
prev_close: 100.,
volume: 0,
}],
)
.unwrap()
}
#[test]
fn confirmed_manual_fill_changes_cash_and_lots_but_not_external_cash_flow_units() {
let input = sample();
let observations = input.observations().unwrap();
let data = identity_data(NaiveDate::from_ymd_opt(2020, 1, 1).unwrap());
let mut account = PortfolioState::new(10_000.);
let applied = observations[0].apply(&mut account, &data, false).unwrap();
assert_eq!(
applied.gross,
FixedMoney::from_decimal_str("1012.345679").unwrap()
);
assert_eq!(applied.fees, FixedMoney::from_decimal_str("0.12").unwrap());
assert_eq!(account.cash(), 8987.534321);
assert_eq!(account.position("000001.SZ").unwrap().quantity, 100);
assert_eq!(
account
.position("000001.SZ")
.unwrap()
.sellable_qty(input.actions[0].orders[0].fills[0].trade_date),
0
);
assert_eq!(account.external_cash_flow_total(), 0.);
assert_eq!(account.starting_cash(), 10_000.);
}
#[test]
fn manual_mismatches_are_atomic_and_do_not_borrow_shares_cash_or_override_pending_orders() {
let data = identity_data(NaiveDate::from_ymd_opt(2020, 1, 1).unwrap());
let input = sample();
let observations = input.observations().unwrap();
let mut poor = PortfolioState::new(10.);
assert!(observations[0].apply(&mut poor, &data, false).is_err());
assert_eq!(poor.cash(), 10.);
assert!(poor.positions().is_empty());
let mut account = PortfolioState::new(10_000.);
assert!(observations[0].apply(&mut account, &data, true).is_err());
assert_eq!(account.cash(), 10_000.);
assert!(account.positions().is_empty());
observations[0].apply(&mut account, &data, false).unwrap();
let before = account.cash();
let mut sell = input.clone();
sell.actions[0].orders[0].side = OrderSide::Sell;
reseal(&mut sell);
assert!(
sell.observations().unwrap()[0]
.apply(&mut account, &data, false)
.unwrap_err()
.contains("T+1")
);
assert_eq!(account.cash(), before);
assert_eq!(account.position("000001.SZ").unwrap().quantity, 100);
let unlisted = identity_data(NaiveDate::from_ymd_opt(2027, 1, 1).unwrap());
assert!(
observations[0]
.apply(&mut account, &unlisted, false)
.unwrap_err()
.contains("lifecycle")
);
assert_eq!(account.cash(), before);
}
#[test]
fn the_next_day_manual_sale_keeps_the_actual_quantity_and_fee_contract() {
let data = identity_data(NaiveDate::from_ymd_opt(2020, 1, 1).unwrap());
let input = sample();
let mut account = PortfolioState::new(10_000.);
input.observations().unwrap()[0]
.apply(&mut account, &data, false)
.unwrap();
let mut sell = input.clone();
let order = &mut sell.actions[0].orders[0];
order.side = OrderSide::Sell;
order.submitted_at += chrono::Duration::days(1);
order.terminal_at += chrono::Duration::days(1);
order.fills[0].trade_date = order.fills[0].trade_date.succ_opt().unwrap();
order.fills[0].executed_at += chrono::Duration::days(1);
order.fills[0].observed_at += chrono::Duration::days(1);
sell.observation_cutoff += chrono::Duration::days(1);
reseal(&mut sell);
let applied = sell.observations().unwrap()[0]
.apply(&mut account, &data, false)
.unwrap();
assert_eq!(applied.quantity_after, 0);
assert_eq!(account.cash(), 9999.76);
assert_eq!(account.external_cash_flow_total(), 0.);
}
#[test]
fn observations_follow_durable_receipt_order_and_not_input_array_order() {
let mut input = sample();
let mut second = input.actions[0].orders[0].fills[0].clone();
second.trade_id = "trade-2".into();
second.observation_event_id = "received-2".into();
second.observation_sequence = 2;
input.actions[0].orders[0].quantity = 200;
input.actions[0].orders[0].fills.insert(0, second);
reseal(&mut input);
assert_eq!(
input
.observations()
.unwrap()
.iter()
.map(|row| row.fill.observation_sequence)
.collect::<Vec<_>>(),
vec![1, 2]
);
let mut invalid = input.clone();
invalid.actions[0].orders[0].fills[0].observation_sequence = 1;
assert!(
semantic_result(&invalid)
.unwrap_err()
.contains("observation")
);
let mut invalid = input;
invalid.actions[0].orders[0].fills[0].observation_event_id = "received-1".into();
assert!(
semantic_result(&invalid)
.unwrap_err()
.contains("observation")
);
}
#[test]
fn partial_cancel_is_valid_but_full_fill_cannot_be_reported_as_cancelled() {
let mut input = sample();
input.actions[0].orders[0].quantity = 200;
input.actions[0].orders[0].terminal_status = ManualOrderTerminalStatus::Cancelled;
semantic_result(&input).unwrap();
input.actions[0].orders[0].quantity = 100;
assert!(
semantic_result(&input)
.unwrap_err()
.contains("terminal status")
);
}
#[test]
fn cursor_waits_for_observation_and_never_reapplies_or_rewinds() {
let input = sample();
let at = input.actions[0].orders[0].fills[0].observed_at;
let mut replay = ManualReplayCursor::new(input).unwrap();
let data = identity_data(NaiveDate::from_ymd_opt(2020, 1, 1).unwrap());
let mut account = PortfolioState::new(10_000.);
assert_eq!(replay.next_observation_at(), Some(at));
assert!(
replay
.advance(
at - chrono::Duration::milliseconds(1),
&mut account,
&data,
false
)
.unwrap()
.is_empty()
);
assert_eq!(account.cash(), 10_000.);
let records = replay.advance(at, &mut account, &data, false).unwrap();
assert_eq!(records.len(), 1);
assert_eq!(records[0].cash_delta, "-1012.465679");
assert_eq!(replay.applied_count(), 1);
assert_eq!(replay.next_observation_at(), None);
let cash = account.cash();
assert!(
replay
.advance(at, &mut account, &data, false)
.unwrap()
.is_empty()
);
assert_eq!(account.cash(), cash);
assert!(
replay
.advance(
at - chrono::Duration::seconds(1),
&mut account,
&data,
false
)
.unwrap_err()
.contains("backwards")
);
}
#[test]
fn failed_multi_receipt_advance_keeps_both_progress_and_portfolio_unchanged() {
let mut input = sample();
let mut next = input.actions[0].orders[0].fills[0].clone();
next.trade_id = "trade-2".into();
next.observation_event_id = "received-2".into();
next.observation_sequence = 2;
input.actions[0].orders[0].quantity = 200;
input.actions[0].orders[0].fills.push(next);
reseal(&mut input);
let at = input.actions[0].orders[0].fills[0].observed_at;
let mut replay = ManualReplayCursor::new(input).unwrap();
let data = identity_data(NaiveDate::from_ymd_opt(2020, 1, 1).unwrap());
let mut account = PortfolioState::new(1_500.);
assert!(replay.advance(at, &mut account, &data, false).is_err());
assert_eq!(account.cash(), 1_500.);
assert!(account.positions().is_empty());
assert_eq!(replay.applied_count(), 0);
assert_eq!(replay.next_observation_at(), Some(at));
}
#[test]
fn fixed_money_decimal_text_preserves_micro_units_without_float_conversion() {
for text in [
"0",
"100",
"-100",
"0.000001",
"-0.000001",
"12345678901234567890123456.123456",
] {
assert_eq!(
FixedMoney::from_decimal_str(text)
.unwrap()
.to_decimal_string(),
text
);
}
let min = FixedMoney::from_raw(i128::MIN);
assert!(min.to_decimal_string().starts_with('-'));
}
@@ -36234,6 +36234,7 @@ mod tests {
avg_price: 0.0,
transaction_cost: 0.0,
limit_price: 10.2,
reserved_cash: None,
reason: "pending_limit_sell".to_string(),
}];
let subscriptions = BTreeSet::new();
@@ -36382,6 +36383,7 @@ mod tests {
avg_price: 0.0,
transaction_cost: 0.0,
limit_price: 9.9,
reserved_cash: None,
reason: "pending_limit_buy".to_string(),
},
OpenOrderView {
@@ -36396,6 +36398,7 @@ mod tests {
avg_price: 0.0,
transaction_cost: 0.0,
limit_price: 10.2,
reserved_cash: None,
reason: "pending_limit_sell".to_string(),
},
];
+129 -9
View File
@@ -138,18 +138,28 @@ impl Position {
if quantity == 0 {
return;
}
let gross_amount = fixed_money_or_panic(execution_price * quantity as f64, "position buy gross amount");
self.buy_with_fixed_gross(date,quantity,execution_price,mark_price,gross_amount);
}
fn buy_with_fixed_gross(
&mut self,
date: NaiveDate,
quantity: u32,
execution_price: f64,
mark_price: f64,
gross_amount: FixedMoney,
) {
let previous_quantity = self.quantity;
self.last_buy_date = Some(self.last_buy_date.map_or(date, |previous| previous.max(date)));
self.last_buy_date = Some(
self.last_buy_date
.map_or(date, |previous| previous.max(date)),
);
if previous_quantity == 0 {
self.opened_date = Some(date);
}
let previous_average_price = self.average_price;
let previous_average_cost = self.average_cost;
let gross_amount = fixed_money_or_panic(
execution_price * quantity as f64,
"position buy gross amount",
);
self.lots.push(PositionLot {
acquired_date: date,
quantity,
@@ -200,6 +210,20 @@ impl Position {
quantity: u32,
execution_price: f64,
mark_price: f64,
) -> Result<f64, String> {
if quantity > self.quantity {
return Err(format!("sell quantity {} exceeds current quantity {} for {}",quantity,self.quantity,self.symbol));
}
let total_proceeds = fixed_money(execution_price * quantity as f64,"position sell gross amount")?;
self.sell_with_fixed_gross(quantity,execution_price,mark_price,total_proceeds)
}
fn sell_with_fixed_gross(
&mut self,
quantity: u32,
execution_price: f64,
mark_price: f64,
total_proceeds: FixedMoney,
) -> Result<f64, String> {
if quantity > self.quantity {
return Err(format!(
@@ -208,10 +232,6 @@ impl Position {
));
}
let total_proceeds = fixed_money(
execution_price * quantity as f64,
"position sell gross amount",
)?;
let mut remaining = quantity;
let mut remaining_proceeds = total_proceeds;
let mut realized = FixedMoney::ZERO;
@@ -796,6 +816,106 @@ impl PortfolioState {
Ok(())
}
/// Apply one fully observed external fill atomically. Its money is already
/// quantized from the original decimal amounts, not from a float product.
pub(crate) fn apply_observed_manual_fill(
&mut self,
trade_date: NaiveDate,
symbol: &str,
side: crate::events::OrderSide,
quantity: u32,
price: f64,
mark_price: f64,
gross: FixedMoney,
fees: FixedMoney,
) -> Result<FixedMoney, String> {
use crate::events::OrderSide;
if symbol.trim().is_empty()
|| quantity == 0
|| quantity > i32::MAX as u32
|| !price.is_finite()
|| price <= 0.
|| !mark_price.is_finite()
|| mark_price <= 0.
|| gross <= FixedMoney::ZERO
|| fees < FixedMoney::ZERO
{
return Err("invalid observed manual fill".into());
}
let mut position = self
.positions
.get(symbol)
.cloned()
.unwrap_or_else(|| Position::new(symbol));
let delta = match side {
OrderSide::Buy => gross.checked_add(fees).and_then(FixedMoney::checked_neg),
OrderSide::Sell => gross.checked_sub(fees),
}
.ok_or("manual fill cash delta overflow")?;
let next_cash = self
.cash
.checked_add(delta)
.filter(|cash| *cash >= FixedMoney::ZERO)
.ok_or("manual fill disagrees with shadow available cash")?;
let next_cost = position
.day_trade_cost
.checked_add(fees)
.ok_or("manual trade cost overflow")?;
match side {
OrderSide::Buy => {
let total_quantity = position
.quantity
.checked_add(quantity)
.ok_or("manual position quantity overflow")?;
FixedMoney::from_f64(mark_price * f64::from(total_quantity))
.ok_or("manual marked position value overflow")?;
position
.day_buy_quantity
.checked_add(quantity)
.ok_or("manual daily buy quantity overflow")?;
position
.day_trade_quantity_delta
.checked_add(quantity as i32)
.ok_or("manual daily quantity delta overflow")?;
position
.day_buy_value
.checked_add(gross)
.ok_or("manual daily buy value overflow")?;
let total_basis = gross.checked_add(fees).ok_or("manual lot basis overflow")?;
position
.total_cost_basis()
.checked_add(total_basis)
.ok_or("manual aggregate position basis overflow")?;
position.buy_with_fixed_gross(trade_date, quantity, price, mark_price, gross);
position
.lots
.last_mut()
.ok_or("manual buy produced no lot")?
.cost_basis = total_basis;
position.average_cost += fees.to_f64() / f64::from(position.quantity);
}
OrderSide::Sell => {
if quantity > position.sellable_qty(trade_date) {
return Err("manual fill disagrees with shadow sellable holdings or T+1".into());
}
position
.day_sell_quantity
.checked_add(quantity)
.ok_or("manual daily sell quantity overflow")?;
position
.day_trade_quantity_delta
.checked_sub(quantity as i32)
.ok_or("manual daily quantity delta overflow")?;
position.sell_with_fixed_gross(quantity, price, mark_price, gross)?;
}
}
position.day_trade_cost = next_cost;
position.refresh_day_pnl();
self.positions.insert(symbol.to_string(), position);
self.cash = next_cash;
Ok(delta)
}
pub fn prune_flat_positions(&mut self) {
let mut sold_symbols = Vec::new();
self.positions.retain(|symbol, position| {
+74 -2
View File
@@ -102,6 +102,7 @@ pub struct OpenOrderView {
pub avg_price: f64,
pub transaction_cost: f64,
pub limit_price: f64,
pub reserved_cash: Option<f64>,
pub reason: String,
}
@@ -497,6 +498,7 @@ impl StrategyContext<'_> {
.iter()
.filter(|order| order.side == OrderSide::Buy)
.map(|order| {
if let Some(reserved) = order.reserved_cash { return reserved; }
let price = if order.limit_price.is_finite() {
order.limit_price.max(0.0)
} else {
@@ -988,6 +990,15 @@ pub struct StrategyDecision {
}
impl StrategyDecision {
pub(crate) fn is_portfolio_target_only(&self) -> bool {
(self.rebalance && self.order_intents.is_empty())
|| (self.order_intents.len() == 1
&& matches!(
self.order_intents[0].unwrapped(),
OrderIntent::StockPool { .. } | OrderIntent::TargetPortfolioSmart { .. }
))
}
pub fn potential_buy_symbols(&self, open_orders: &[OpenOrderView]) -> BTreeSet<String> {
let mut symbols = BTreeSet::new();
if self.rebalance {
@@ -1001,9 +1012,24 @@ impl StrategyDecision {
}
pub fn merge_from(&mut self, mut other: StrategyDecision) {
if self.is_portfolio_target_only() && other.is_portfolio_target_only() {
let mut previous = std::mem::replace(self, other);
previous
.diagnostics
.push("unsubmitted_portfolio_target_superseded".into());
self.notes.splice(0..0, previous.notes);
self.diagnostics.splice(0..0, previous.diagnostics);
return;
}
self.buy_denials.append(&mut other.buy_denials);
self.rebalance |= other.rebalance;
self.target_weights.append(&mut other.target_weights);
if other.rebalance {
// Rebalance targets are a complete portfolio, not an additive
// list. A newer unsent target replaces the earlier allocation.
self.rebalance = true;
self.target_weights = std::mem::take(&mut other.target_weights);
} else {
self.target_weights.append(&mut other.target_weights);
}
self.exit_symbols.append(&mut other.exit_symbols);
self.order_intents.append(&mut other.order_intents);
self.notes.append(&mut other.notes);
@@ -1023,6 +1049,52 @@ impl StrategyDecision {
}
}
#[cfg(test)]
mod decision_merge_tests {
use super::*;
#[test]
fn newer_complete_target_replaces_old_symbols_without_discarding_explicit_actions() {
let mut earlier = StrategyDecision {
rebalance: true,
target_weights: BTreeMap::from([("A".into(), 0.5), ("B".into(), 0.5)]),
exit_symbols: BTreeSet::from(["risk_exit".into()]),
order_intents: vec![OrderIntent::Shares {
symbol: "explicit".into(),
quantity: 100,
reason: "explicit action".into(),
}],
..Default::default()
};
earlier.merge_from(StrategyDecision {
rebalance: true,
target_weights: BTreeMap::from([("C".into(), 1.)]),
..Default::default()
});
assert_eq!(earlier.target_weights, BTreeMap::from([("C".into(), 1.)]));
assert!(earlier.rebalance);
assert!(earlier.exit_symbols.contains("risk_exit"));
assert_eq!(earlier.order_intents.len(), 1);
}
#[test]
fn explicit_empty_complete_target_replaces_old_allocation_but_empty_callback_does_not() {
let mut decision = StrategyDecision {
rebalance: true,
target_weights: BTreeMap::from([("A".into(), 1.)]),
..Default::default()
};
decision.merge_from(StrategyDecision::default());
assert_eq!(decision.target_weights.len(), 1);
decision.merge_from(StrategyDecision {
rebalance: true,
..Default::default()
});
assert!(decision.target_weights.is_empty());
assert!(decision.rebalance);
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AlgoOrderStyle {
Vwap,