1use std::collections::{BTreeSet, HashMap};
9#[cfg(test)]
10use std::path::Path;
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::{Arc, Mutex, RwLock};
13use std::time::{Duration, Instant};
14
15use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
16use chrono::NaiveDateTime;
17#[cfg(test)]
18use data_preprocess::DataError;
19use data_preprocess::ParquetStore;
20#[cfg(test)]
21use data_preprocess::models::{BarQueryOpts, QueryOpts, Timeframe};
22use futures::StreamExt;
23use futures::stream::{self, BoxStream};
24use qs_backtest::BacktestResult;
25#[cfg(test)]
26use qs_backtest::data_feed::{DataFeed, MarketEvent, VecFeed, bars_to_feed, ticks_to_feed};
27use qs_backtest::data_feed::{EventBatchFeedError, KWayMergeError};
28use qs_backtest::evaluation::EvaluationOptions;
29use qs_backtest::profile::{ManagementProfile, ProfileRegistry, RawSignal};
30use qs_backtest::runner::{
31 BacktestConfig, BacktestRunner, FutureQuoteConfig, ReplayProgress, StreamingReplayError,
32};
33use qs_symbols::SymbolRegistry;
34use tokio::sync::watch;
35
36use crate::artifact_store::ArtifactStore;
37use crate::convert::{
38 account_currency_from_msg, config_from_msg, evaluation_options_from_msg_for_symbols,
39 future_config_from_msg, profile_from_msg, raw_signal_from_msg, result_to_msg,
40 validate_future_quote_scalars,
41};
42use crate::error::{BacktestServerError, Result};
43use crate::fx_loader::describe_future_stream;
44use crate::market_loader::{
45 CancellationCheck, MarketStreamDescription, MarketStreamError, describe_primary_market_stream,
46};
47use crate::replay_plan::{ReplayPlan, RequestedSymbolScope};
48use crate::rpc_types::*;
49
50pub struct ServerState {
52 pub symbol_registry: SymbolRegistry,
54 pub profile_registry: RwLock<ProfileRegistry>,
56 pub data_dir: String,
58 pub profiles_path: String,
60 pub start_time: Instant,
62 pub jobs: Mutex<HashMap<String, BacktestJob>>,
64 pub max_retained_jobs: usize,
66 pub artifact_store: ArtifactStore,
68}
69
70#[derive(Debug, Clone)]
72pub struct BacktestJob {
73 pub status: JobStatus,
75 pub submitted_at: Instant,
77 pub completed_at: Option<Instant>,
79 pub progress: BacktestProgress,
81 pub result: Option<BacktestResultMsg>,
83 pub artifact: Option<ResultArtifactRefMsg>,
85 pub inline_complete: bool,
87 pub artifact_consumed: bool,
89 pub error: Option<String>,
91 pub cancellation: JobCancellationToken,
93 pub worker_active: bool,
95 pub updates: watch::Sender<BacktestStatusResponse>,
97}
98
99#[derive(Debug, Clone, Default)]
101pub struct JobCancellationToken(Arc<AtomicBool>);
102
103impl JobCancellationToken {
104 pub fn cancel(&self) {
105 self.0.store(true, Ordering::Release);
106 }
107
108 pub fn is_cancelled(&self) -> bool {
109 self.0.load(Ordering::Acquire)
110 }
111}
112
113#[derive(Debug, Clone, PartialEq, Eq)]
115pub enum JobStatus {
116 Queued,
117 LoadingData,
118 Running,
119 Completed,
120 Failed,
121 Cancelled,
122}
123
124enum PreparedResult<T> {
125 Inline(T),
126 Artifact {
127 reference: ResultArtifactRefMsg,
128 summary: Option<T>,
129 },
130}
131
132fn prepare_result<T, F>(
133 state: &ServerState,
134 result: T,
135 delivery: Option<ResultDeliveryMsg>,
136 summarize: F,
137) -> std::result::Result<PreparedResult<T>, String>
138where
139 T: serde::Serialize,
140 F: FnOnce(&T) -> T,
141{
142 let bytes = serde_json::to_vec(&result)
143 .map_err(|error| format!("failed to serialize result JSON: {error}"))?;
144 let Some(delivery) = delivery else {
145 return Ok(PreparedResult::Inline(result));
146 };
147 let inline_limit = state.artifact_store.inline_limit_bytes();
148 match delivery {
149 ResultDeliveryMsg::Inline if bytes.len() > inline_limit => Err(format!(
150 "result JSON is {} bytes, exceeding the configured inline limit of {} bytes; use result_delivery 'auto' or 'artifact'",
151 bytes.len(),
152 inline_limit
153 )),
154 ResultDeliveryMsg::Inline => Ok(PreparedResult::Inline(result)),
155 ResultDeliveryMsg::Auto if bytes.len() <= inline_limit => {
156 Ok(PreparedResult::Inline(result))
157 }
158 ResultDeliveryMsg::Auto | ResultDeliveryMsg::Artifact => {
159 let reference = state
160 .artifact_store
161 .persist_json(&bytes)
162 .map_err(|error| format!("failed to persist result artifact: {error}"))?;
163 let summary = summarize(&result);
164 let summary = serde_json::to_vec(&summary)
165 .ok()
166 .filter(|bytes| bytes.len() <= inline_limit)
167 .map(|_| summary);
168 Ok(PreparedResult::Artifact { reference, summary })
169 }
170 }
171}
172
173fn compact_result_for_console(result: &BacktestResultMsg) -> BacktestResultMsg {
174 let mut summary = result.clone();
175 summary.equity_curve.clear();
176 summary.trade_log.truncate(30);
177 summary.positions.truncate(15);
178 if let Some(future) = summary.future.as_mut() {
179 future.recorded_fills = serde_json::Value::Null;
180 future.action_dispositions = serde_json::Value::Null;
181 future.close_events = serde_json::Value::Null;
182 future.completed_positions = serde_json::Value::Null;
183 future.open_positions = serde_json::Value::Null;
184 future.pending_orders = serde_json::Value::Null;
185 future.pending_order_lifecycle.clear();
186 future.mtm_equity_curve = serde_json::Value::Null;
187 }
188 summary
189}
190
191fn compact_profile_results(results: &[ProfileResult]) -> Vec<ProfileResult> {
192 results
193 .iter()
194 .cloned()
195 .map(|mut profile| {
196 profile.result = profile.result.as_ref().map(compact_result_for_console);
197 profile
198 })
199 .collect()
200}
201
202fn single_response_from_result(
203 state: &ServerState,
204 result: BacktestResultMsg,
205 start: Instant,
206 delivery: Option<ResultDeliveryMsg>,
207) -> RunBacktestResponse {
208 match prepare_result(state, result, delivery, compact_result_for_console) {
209 Ok(PreparedResult::Inline(result)) => RunBacktestResponse {
210 success: true,
211 error: None,
212 result: Some(result),
213 elapsed_ms: start.elapsed().as_millis() as u64,
214 artifact: None,
215 inline_complete: true,
216 },
217 Ok(PreparedResult::Artifact { reference, summary }) => RunBacktestResponse {
218 success: true,
219 error: None,
220 result: summary,
221 elapsed_ms: start.elapsed().as_millis() as u64,
222 artifact: Some(reference),
223 inline_complete: false,
224 },
225 Err(error) => RunBacktestResponse {
226 success: false,
227 error: Some(error),
228 result: None,
229 elapsed_ms: start.elapsed().as_millis() as u64,
230 artifact: None,
231 inline_complete: false,
232 },
233 }
234}
235
236fn multi_response_from_results(
237 state: &ServerState,
238 results: Vec<ProfileResult>,
239 start: Instant,
240 delivery: Option<ResultDeliveryMsg>,
241) -> RunBacktestMultiResponse {
242 let result_error = results
243 .iter()
244 .find(|result| !result.success)
245 .and_then(|result| result.error.clone());
246 let result_success = result_error.is_none();
247 match prepare_result(state, results, delivery, |results| {
248 compact_profile_results(results)
249 }) {
250 Ok(PreparedResult::Inline(results)) => RunBacktestMultiResponse {
251 success: result_success,
252 error: result_error,
253 results,
254 elapsed_ms: start.elapsed().as_millis() as u64,
255 artifact: None,
256 inline_complete: true,
257 },
258 Ok(PreparedResult::Artifact { reference, summary }) => RunBacktestMultiResponse {
259 success: result_success,
260 error: result_error,
261 results: summary.unwrap_or_default(),
262 elapsed_ms: start.elapsed().as_millis() as u64,
263 artifact: Some(reference),
264 inline_complete: false,
265 },
266 Err(error) => RunBacktestMultiResponse {
267 success: false,
268 error: Some(error),
269 results: Vec::new(),
270 elapsed_ms: start.elapsed().as_millis() as u64,
271 artifact: None,
272 inline_complete: false,
273 },
274 }
275}
276
277impl JobStatus {
278 pub fn as_str(&self) -> &'static str {
280 match self {
281 JobStatus::Queued => "Queued",
282 JobStatus::LoadingData => "LoadingData",
283 JobStatus::Running => "Running",
284 JobStatus::Completed => "Completed",
285 JobStatus::Failed => "Failed",
286 JobStatus::Cancelled => "Cancelled",
287 }
288 }
289
290 fn is_terminal(&self) -> bool {
291 matches!(self, Self::Completed | Self::Failed | Self::Cancelled)
292 }
293}
294
295fn job_status_response(job_id: &str, job: &BacktestJob) -> BacktestStatusResponse {
296 let elapsed_ms = job
297 .completed_at
298 .map(|completed| completed.duration_since(job.submitted_at).as_millis() as u64);
299 BacktestStatusResponse {
300 success: true,
301 job_id: job_id.to_owned(),
302 status: job.status.as_str().to_owned(),
303 error: job.error.clone(),
304 elapsed_ms,
305 progress: job.progress.clone(),
306 }
307}
308
309fn publish_job_status(job_id: &str, job: &BacktestJob) {
310 job.updates.send_replace(job_status_response(job_id, job));
311}
312
313pub fn subscribe_backtest_status(
315 state: &ServerState,
316 job_id: &str,
317) -> std::result::Result<(watch::Receiver<BacktestStatusResponse>, Instant), String> {
318 let jobs = state.jobs.lock().unwrap();
319 let job = jobs
320 .get(job_id)
321 .ok_or_else(|| format!("Job '{job_id}' not found"))?;
322 Ok((job.updates.subscribe(), job.submitted_at))
323}
324
325struct BacktestWatchState {
326 job_id: String,
327 updates: watch::Receiver<BacktestStatusResponse>,
328 submitted_at: Instant,
329 heartbeat: tokio::time::Interval,
330 emit_initial: bool,
331 finished: bool,
332}
333
334pub fn watch_backtest_stream(
336 state: Arc<ServerState>,
337 req: WatchBacktestRequest,
338 heartbeat_interval: Duration,
339) -> BoxStream<'static, std::result::Result<BacktestEvent, xrpc::RpcError>> {
340 let (updates, submitted_at) = match subscribe_backtest_status(&state, &req.job_id) {
341 Ok(subscription) => subscription,
342 Err(error) => {
343 return stream::once(async move { Err(xrpc::RpcError::ServerError(error)) }).boxed();
344 }
345 };
346
347 let period = heartbeat_interval.max(Duration::from_millis(1));
348 let start = tokio::time::Instant::now() + period;
349 let mut heartbeat = tokio::time::interval_at(start, period);
350 heartbeat.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
351 let initial = BacktestWatchState {
352 job_id: req.job_id,
353 updates,
354 submitted_at,
355 heartbeat,
356 emit_initial: true,
357 finished: false,
358 };
359
360 stream::unfold(initial, |mut state| async move {
361 if state.finished {
362 return None;
363 }
364
365 if state.emit_initial {
366 state.emit_initial = false;
367 let status = state.updates.borrow().clone();
368 state.finished = status.is_terminal();
369 return Some((Ok(BacktestEvent::Snapshot { status }), state));
370 }
371
372 tokio::select! {
373 changed = state.updates.changed() => {
374 match changed {
375 Ok(()) => {
376 let status = state.updates.borrow().clone();
377 state.finished = status.is_terminal();
378 Some((Ok(BacktestEvent::Snapshot { status }), state))
379 }
380 Err(_) => {
381 state.finished = true;
382 Some((Err(xrpc::RpcError::ServerError(format!(
383 "Backtest job '{}' update channel closed before a terminal snapshot",
384 state.job_id
385 ))), state))
386 }
387 }
388 }
389 _ = state.heartbeat.tick() => {
390 let event = BacktestEvent::Heartbeat {
391 job_id: state.job_id.clone(),
392 elapsed_ms: state.submitted_at.elapsed().as_millis() as u64,
393 };
394 Some((Ok(event), state))
395 }
396 }
397 })
398 .boxed()
399}
400
401fn progress_stage_rank(stage: &str) -> u8 {
402 match stage {
403 "queued" => 0,
404 "loading_data" => 1,
405 "replay" => 2,
406 "completed" | "failed" | "cancelled" => 3,
407 _ => 0,
408 }
409}
410
411fn merge_progress(current: &mut BacktestProgress, next: BacktestProgress) {
412 if progress_stage_rank(&next.stage) >= progress_stage_rank(¤t.stage) {
413 current.stage = next.stage;
414 }
415 current.processed_events = current.processed_events.max(next.processed_events);
416 current.total_events = current.total_events.max(next.total_events);
417 current.processed_signals = current.processed_signals.max(next.processed_signals);
418 current.total_signals = current.total_signals.max(next.total_signals);
419 current.processed_symbols = current.processed_symbols.max(next.processed_symbols);
420 current.total_symbols = current.total_symbols.max(next.total_symbols);
421}
422
423fn remove_oldest_terminal_job(jobs: &mut HashMap<String, BacktestJob>) -> Option<BacktestJob> {
424 let oldest = jobs
425 .iter()
426 .filter(|(_, job)| job.status.is_terminal() && !job.worker_active)
427 .min_by(|(left_id, left), (right_id, right)| {
428 left.completed_at
429 .cmp(&right.completed_at)
430 .then_with(|| left.submitted_at.cmp(&right.submitted_at))
431 .then_with(|| left_id.cmp(right_id))
432 })
433 .map(|(id, _)| id.clone())?;
434 jobs.remove(&oldest)
435}
436
437fn delete_job_artifacts(state: &ServerState, removed: &[BacktestJob]) {
438 for job in removed {
439 if !job.artifact_consumed
440 && let Some(artifact) = job.artifact.as_ref()
441 {
442 let _ = state.artifact_store.delete(&artifact.artifact_id);
443 }
444 }
445}
446
447pub fn cleanup_expired_jobs(state: &ServerState, retention: Duration) -> usize {
449 let removed = {
450 let mut jobs = state.jobs.lock().unwrap();
451 let expired = jobs
452 .iter()
453 .filter(|(_, job)| {
454 job.status.is_terminal()
455 && !job.worker_active
456 && job
457 .completed_at
458 .is_some_and(|completed| completed.elapsed() >= retention)
459 })
460 .map(|(id, _)| id.clone())
461 .collect::<Vec<_>>();
462 let mut removed = expired
463 .into_iter()
464 .filter_map(|id| jobs.remove(&id))
465 .collect::<Vec<_>>();
466 while jobs.len() > state.max_retained_jobs {
467 let Some(job) = remove_oldest_terminal_job(&mut jobs) else {
468 break;
469 };
470 removed.push(job);
471 }
472 removed
473 };
474 delete_job_artifacts(state, &removed);
475 removed.len()
476}
477
478pub fn cancel_active_jobs(state: &ServerState) -> usize {
480 let mut cancelled = 0;
481 let mut jobs = state.jobs.lock().unwrap();
482 for (job_id, job) in jobs.iter_mut().filter(|(_, job)| !job.status.is_terminal()) {
483 job.cancellation.cancel();
484 job.status = JobStatus::Cancelled;
485 job.completed_at = Some(Instant::now());
486 job.result = None;
487 job.artifact = None;
488 job.inline_complete = true;
489 job.artifact_consumed = false;
490 job.error = None;
491 job.progress.stage = "cancelled".into();
492 publish_job_status(job_id, job);
493 cancelled += 1;
494 }
495 cancelled
496}
497
498fn update_job_progress(state: &ServerState, job_id: &str, progress: BacktestProgress) {
499 let mut jobs = state.jobs.lock().unwrap();
500 if let Some(job) = jobs.get_mut(job_id)
501 && !job.status.is_terminal()
502 {
503 match progress.stage.as_str() {
504 "loading_data" => job.status = JobStatus::LoadingData,
505 "replay" => job.status = JobStatus::Running,
506 _ => {}
507 }
508 merge_progress(&mut job.progress, progress);
509 publish_job_status(job_id, job);
510 }
511}
512
513impl std::str::FromStr for JobStatus {
514 type Err = &'static str;
515
516 fn from_str(s: &str) -> std::result::Result<Self, Self::Err> {
517 match s {
518 "Queued" => Ok(Self::Queued),
519 "LoadingData" => Ok(Self::LoadingData),
520 "Running" => Ok(Self::Running),
521 "Completed" => Ok(Self::Completed),
522 "Failed" => Ok(Self::Failed),
523 "Cancelled" => Ok(Self::Cancelled),
524 _ => Err("invalid job status"),
525 }
526 }
527}
528
529pub fn handle_ping(state: &ServerState) -> PingResponse {
533 PingResponse {
534 status: "OK".into(),
535 uptime_secs: state.start_time.elapsed().as_secs(),
536 data_dir: state.data_dir.clone(),
537 }
538}
539
540pub fn handle_list_profiles(state: &ServerState) -> ListProfilesResponse {
544 let registry = state.profile_registry.read().unwrap();
545 let profiles = registry
546 .names()
547 .into_iter()
548 .filter_map(|name| {
549 let p = registry.get(name)?;
550 Some(ProfileInfo {
551 name: p.name.clone(),
552 use_targets: p.use_targets.clone(),
553 close_ratios: p.close_ratios.clone(),
554 stoploss_mode: format!("{:?}", p.stoploss_mode),
555 rules_count: p.rules.len(),
556 let_remainder_run: p.let_remainder_run,
557 })
558 })
559 .collect();
560 ListProfilesResponse { profiles }
561}
562
563pub fn handle_list_symbols(
567 state: &ServerState,
568 req: &ListSymbolsRequest,
569) -> std::result::Result<ListSymbolsResponse, String> {
570 let store = ParquetStore::open(&state.data_dir).map_err(|e| e.to_string())?;
571
572 let exchange_filter = req.exchange.as_deref();
573 let symbol_filter: Option<&str> = None;
574
575 let stat_rows = store
576 .stats(exchange_filter, symbol_filter)
577 .map_err(|e| e.to_string())?;
578
579 let symbols: Vec<SymbolAvailability> = stat_rows
580 .into_iter()
581 .filter(|row| {
582 if let Some(ref dt) = req.data_type {
584 let dt_lower = dt.to_lowercase();
585 if dt_lower == "tick" && row.data_type != "tick" {
586 return false;
587 }
588 if dt_lower == "bar" && row.data_type == "tick" {
589 return false;
590 }
591 }
592 true
593 })
594 .map(|row| {
595 let (data_type, timeframe) = parse_availability_data_type(&row.data_type);
598
599 SymbolAvailability {
600 exchange: row.exchange,
601 symbol: row.symbol,
602 data_type,
603 timeframe,
604 row_count: row.count,
605 earliest: row.ts_min.format("%Y-%m-%dT%H:%M:%S").to_string(),
606 latest: row.ts_max.format("%Y-%m-%dT%H:%M:%S").to_string(),
607 }
608 })
609 .collect();
610
611 Ok(ListSymbolsResponse { symbols })
612}
613
614fn parse_availability_data_type(raw: &str) -> (String, Option<String>) {
615 let trimmed = raw.trim();
616 let Some(rest) = trimmed.strip_prefix("bar") else {
617 return (trimmed.to_string(), None);
618 };
619 let timeframe = rest
620 .trim()
621 .strip_prefix('(')
622 .and_then(|value| value.strip_suffix(')'))
623 .map(str::trim)
624 .filter(|value| !value.is_empty())
625 .map(ToOwned::to_owned);
626 ("bar".into(), timeframe)
627}
628
629pub fn handle_run_backtest(state: &ServerState, req: &RunBacktestRequest) -> RunBacktestResponse {
633 let start = Instant::now();
634 if let Err(error) = validate_future_quote_scalars(&req.future) {
635 return RunBacktestResponse {
636 success: false,
637 error: Some(error.to_string()),
638 result: None,
639 elapsed_ms: start.elapsed().as_millis() as u64,
640 artifact: None,
641 inline_complete: true,
642 };
643 }
644 match execute_backtest_with_future(state, &req.request, &req.future, &req.evaluation) {
645 Ok(result) => single_response_from_result(
646 state,
647 result_to_msg(&result),
648 start,
649 Some(req.result_delivery),
650 ),
651 Err(error) => RunBacktestResponse {
652 success: false,
653 error: Some(error.to_string()),
654 result: None,
655 elapsed_ms: start.elapsed().as_millis() as u64,
656 artifact: None,
657 inline_complete: true,
658 },
659 }
660}
661
662fn execute_backtest_with_future(
663 state: &ServerState,
664 req: &BacktestRunSpec,
665 future: &FutureQuoteConfigMsg,
666 evaluation: &ProviderEvaluationOptionsMsg,
667) -> Result<BacktestResult> {
668 execute_backtest_with_future_controlled(state, req, future, evaluation, None, &mut |_| {})
669}
670
671fn ensure_not_cancelled(cancellation: Option<&JobCancellationToken>) -> Result<()> {
672 if cancellation.is_some_and(JobCancellationToken::is_cancelled) {
673 Err(BacktestServerError::Cancelled)
674 } else {
675 Ok(())
676 }
677}
678
679fn map_streaming_replay_error(
680 error: StreamingReplayError<MarketStreamError>,
681) -> BacktestServerError {
682 match error {
683 StreamingReplayError::Cancelled(_) => BacktestServerError::Cancelled,
684 StreamingReplayError::Feed(KWayMergeError::Source {
685 error: EventBatchFeedError::Source(error),
686 ..
687 }) => error,
688 StreamingReplayError::Feed(error) => BacktestServerError::MarketStream(error.to_string()),
689 }
690}
691
692fn execute_backtest_with_future_controlled(
693 state: &ServerState,
694 req: &BacktestRunSpec,
695 future: &FutureQuoteConfigMsg,
696 evaluation: &ProviderEvaluationOptionsMsg,
697 cancellation: Option<&JobCancellationToken>,
698 progress: &mut dyn FnMut(BacktestProgress),
699) -> Result<BacktestResult> {
700 ensure_not_cancelled(cancellation)?;
701 validate_request(req)?;
702 validate_future_quote_scalars(future)?;
703
704 if let Some(ref profile_msg) = req.profile_def {
706 let profile = profile_from_msg(profile_msg)?;
707 profile.validate().map_err(|error| {
708 BacktestServerError::InvalidRequest(format!("Invalid inline profile: {error}"))
709 })?;
710 }
711 let profile = resolve_profile(state, req)?;
712
713 let from = parse_optional_datetime(&req.from)?;
714 let to = parse_optional_datetime(&req.to)?;
715 let plan = build_replay_plan(
716 state,
717 &req.symbol,
718 &req.symbols,
719 req.all_symbols,
720 &req.raw_signals,
721 from,
722 to,
723 future.signal_latency_ms,
724 )?;
725 validate_replay_sizing(req, &plan)?;
726 let evaluation_options = evaluation_options_from_msg_for_symbols(
727 evaluation,
728 &state.symbol_registry,
729 plan.requested_symbols(),
730 )?;
731 let config = config_from_msg(&req.config, &state.symbol_registry, plan.active_symbols())?;
732 let account_currency = account_currency_from_msg(future)?;
733 let exchange = req.exchange.to_lowercase();
734 tracing::info!(
735 "run_backtest: requested_symbols={:?} active_symbols={:?} idle_symbols={:?} loading_start={:?} exchange={} data_type={}",
736 plan.requested_symbols(),
737 plan.active_symbols(),
738 plan.idle_explicit_symbols(),
739 plan.loading_start(),
740 exchange,
741 req.data_type
742 );
743 tracing::info!(
744 "run_backtest: {} signals after date filtering",
745 plan.retained_signals().len()
746 );
747 ensure_not_cancelled(cancellation)?;
748
749 let total_symbols = plan.active_symbols().len() as u64;
750 let total_signals = plan.retained_signals().len() as u64;
751 progress(BacktestProgress {
752 stage: "loading_data".into(),
753 total_signals,
754 total_symbols,
755 ..BacktestProgress::default()
756 });
757 tracing::info!(
758 "run_backtest: loading market data for {} active symbols...",
759 plan.active_symbols().len()
760 );
761 let mut result = {
762 let mut cancelled = || cancellation.is_some_and(JobCancellationToken::is_cancelled);
763 let primary = describe_primary_market_stream(
764 &state.data_dir,
765 &exchange,
766 plan.active_symbols(),
767 &req.data_type,
768 req.timeframe.as_deref(),
769 plan.loading_start(),
770 to,
771 &mut cancelled,
772 &mut |processed_symbols| {
773 progress(BacktestProgress {
774 stage: "loading_data".into(),
775 processed_symbols,
776 total_symbols,
777 total_signals,
778 ..BacktestProgress::default()
779 });
780 },
781 )?;
782 progress(BacktestProgress {
783 stage: "loading_conversion_data".into(),
784 processed_symbols: total_symbols,
785 total_symbols,
786 total_signals,
787 ..BacktestProgress::default()
788 });
789 let bundle = describe_future_stream(
790 &state.data_dir,
791 &exchange,
792 &state.symbol_registry,
793 &account_currency,
794 plan.active_symbols(),
795 &req.data_type,
796 plan.loading_start(),
797 primary,
798 &mut cancelled,
799 )?;
800 let primary_eod = bundle.description.primary_eod();
801 let future_config = future_config_from_msg(future, bundle.currency_plan)?;
802 let cancellation_token = cancellation.cloned();
803 let stream_cancellation: CancellationCheck = Arc::new(move || {
804 cancellation_token
805 .as_ref()
806 .is_some_and(JobCancellationToken::is_cancelled)
807 });
808 let mut feed = bundle.description.open(stream_cancellation)?;
809 ensure_not_cancelled(cancellation)?;
810 progress(BacktestProgress {
811 stage: "replay".into(),
812 total_signals,
813 processed_symbols: total_symbols,
814 total_symbols,
815 ..BacktestProgress::default()
816 });
817 tracing::info!("run_backtest: starting streaming FutureQuote engine...");
818 let mut runner = BacktestRunner::new_future(config, future_config);
819 runner = runner.with_evaluation_options(evaluation_options.clone());
820 runner
821 .run_raw_signals_future_streaming_controlled(
822 &mut feed,
823 primary_eod,
824 plan.retained_signals().to_vec(),
825 profile.as_ref(),
826 || cancellation.is_some_and(JobCancellationToken::is_cancelled),
827 |ReplayProgress {
828 processed_events,
829 total_events,
830 processed_signals,
831 total_signals,
832 }| {
833 progress(BacktestProgress {
834 stage: "replay".into(),
835 processed_events: processed_events as u64,
836 total_events: total_events as u64,
837 processed_signals: processed_signals as u64,
838 total_signals: total_signals as u64,
839 processed_symbols: total_symbols,
840 total_symbols,
841 });
842 },
843 )
844 .map_err(map_streaming_replay_error)?
845 };
846 ensure_not_cancelled(cancellation)?;
847
848 attach_future_reproducibility_metadata(
849 &mut result,
850 state,
851 req,
852 &plan,
853 future,
854 profile.as_ref(),
855 );
856 tracing::info!(
857 "run_backtest: done, {} trades, {} positions",
858 result.total_trades,
859 result.total_positions
860 );
861 Ok(result)
862}
863
864fn attach_future_reproducibility_metadata(
865 result: &mut BacktestResult,
866 state: &ServerState,
867 req: &BacktestRunSpec,
868 plan: &ReplayPlan,
869 future: &FutureQuoteConfigMsg,
870 profile: Option<&ManagementProfile>,
871) {
872 let Some(metadata) = result.execution_metadata.as_mut() else {
873 return;
874 };
875 let tags = &mut metadata.tags;
876 tags.insert("data.exchange".into(), req.exchange.to_lowercase());
877 tags.insert("data.type".into(), req.data_type.to_lowercase());
878 tags.insert(
879 "data.timeframe".into(),
880 req.timeframe.clone().unwrap_or_else(|| "none".into()),
881 );
882 tags.insert(
883 "data.requested_from".into(),
884 req.from.clone().unwrap_or_else(|| "unbounded".into()),
885 );
886 tags.insert(
887 "data.requested_to".into(),
888 req.to.clone().unwrap_or_else(|| "unbounded".into()),
889 );
890 tags.insert("data.symbols".into(), plan.active_symbols().join(","));
891 tags.insert(
892 "data.requested_symbols".into(),
893 plan.requested_symbols().join(","),
894 );
895 tags.insert(
896 "data.active_symbols".into(),
897 plan.active_symbols().join(","),
898 );
899 tags.insert(
900 "data.idle_symbols".into(),
901 plan.idle_explicit_symbols().join(","),
902 );
903 tags.insert("data.idle_run".into(), plan.is_idle().to_string());
904 tags.insert(
905 "data.loading_from".into(),
906 plan.loading_start()
907 .map(|timestamp| timestamp.format("%Y-%m-%dT%H:%M:%S%.f").to_string())
908 .unwrap_or_else(|| "none".into()),
909 );
910 tags.insert(
911 "execution.signal_latency_ms".into(),
912 future.signal_latency_ms.max(0).to_string(),
913 );
914 if req.data_type.eq_ignore_ascii_case("bar") {
915 tags.insert(
916 "data.bar_quote_convention".into(),
917 "close_only_zero_spread".into(),
918 );
919 tags.insert("data.intrabar_simulation".into(), "false".into());
920 }
921
922 match profile {
923 Some(profile) => {
924 tags.insert("profile.identity".into(), profile.name.clone());
925 tags.insert(
926 "profile.options".into(),
927 serde_json::to_string(profile).unwrap_or_else(|_| "unavailable".into()),
928 );
929 }
930 None => {
931 tags.insert("profile.identity".into(), "none".into());
932 tags.insert("profile.options".into(), "null".into());
933 }
934 }
935 match req.config.sizing.as_ref() {
936 Some(sizing) => {
937 let identity = match sizing {
938 SizingPolicyMsg::FixedLot { .. } => "fixed_lot",
939 SizingPolicyMsg::FixedRiskAmount { .. } => "fixed_risk_amount",
940 SizingPolicyMsg::BalanceRiskPercent { .. } => "balance_risk_percent",
941 };
942 tags.insert("sizing.identity".into(), identity.into());
943 tags.insert(
944 "sizing.options".into(),
945 serde_json::to_string(sizing).unwrap_or_else(|_| "unavailable".into()),
946 );
947 }
948 None => {
949 tags.insert("sizing.identity".into(), "signal_size".into());
950 tags.insert("sizing.options".into(), "null".into());
951 }
952 }
953
954 for symbol in plan.active_symbols() {
955 let Some(spec) = state.symbol_registry.spec(symbol) else {
956 continue;
957 };
958 let prefix = format!("symbol.{symbol}");
959 tags.insert(format!("{prefix}.canonical"), spec.canonical.clone());
960 tags.insert(
961 format!("{prefix}.pip_position"),
962 spec.pip_position.to_string(),
963 );
964 tags.insert(format!("{prefix}.digits"), spec.digits.to_string());
965 tags.insert(format!("{prefix}.category"), spec.category.clone());
966 tags.insert(
967 format!("{prefix}.lot_base_units"),
968 spec.lot_base_units.to_string(),
969 );
970 tags.insert(
971 format!("{prefix}.lot_step_units"),
972 spec.lot_step_units.to_string(),
973 );
974 tags.insert(
975 format!("{prefix}.lot_min_steps"),
976 spec.lot_min_steps.to_string(),
977 );
978 tags.insert(
979 format!("{prefix}.lot_max_steps"),
980 spec.lot_max_steps.to_string(),
981 );
982 }
983}
984
985pub fn handle_run_backtest_multi(
988 state: &ServerState,
989 req: &RunBacktestMultiRequest,
990) -> RunBacktestMultiResponse {
991 let start = Instant::now();
992 if req.request.profiles.is_empty() {
993 return RunBacktestMultiResponse {
994 success: false,
995 error: Some("At least one profile is required.".into()),
996 results: Vec::new(),
997 elapsed_ms: start.elapsed().as_millis() as u64,
998 artifact: None,
999 inline_complete: true,
1000 };
1001 }
1002 if let Err(error) = validate_future_quote_scalars(&req.future) {
1003 let error = error.to_string();
1004 return RunBacktestMultiResponse {
1005 success: false,
1006 error: Some(error.clone()),
1007 results: profile_error_results(&req.request, error),
1008 elapsed_ms: start.elapsed().as_millis() as u64,
1009 artifact: None,
1010 inline_complete: true,
1011 };
1012 }
1013
1014 let results =
1015 execute_backtest_multi_with_future(state, &req.request, &req.future, &req.evaluation);
1016 multi_response_from_results(state, results, start, Some(req.result_delivery))
1017}
1018
1019fn execute_backtest_multi_with_future(
1020 state: &ServerState,
1021 req: &BacktestMultiRunSpec,
1022 future: &FutureQuoteConfigMsg,
1023 evaluation: &ProviderEvaluationOptionsMsg,
1024) -> Vec<ProfileResult> {
1025 if let Err(error) = validate_future_quote_scalars(future) {
1026 return profile_error_results(req, error.to_string());
1027 }
1028
1029 let data_type = req.data_type.to_lowercase();
1031 if data_type != "tick" && data_type != "bar" {
1032 return req
1033 .profiles
1034 .iter()
1035 .map(|pr| ProfileResult {
1036 profile: profile_ref_name(pr),
1037 success: false,
1038 error: Some(format!(
1039 "Invalid data_type: '{}'. Must be 'tick' or 'bar'.",
1040 req.data_type
1041 )),
1042 result: None,
1043 })
1044 .collect();
1045 }
1046
1047 if data_type == "bar" && req.timeframe.is_none() {
1048 return req
1049 .profiles
1050 .iter()
1051 .map(|pr| ProfileResult {
1052 profile: profile_ref_name(pr),
1053 success: false,
1054 error: Some("timeframe is required when data_type is 'bar'.".into()),
1055 result: None,
1056 })
1057 .collect();
1058 }
1059
1060 if req.raw_signals.is_empty() {
1061 return req
1062 .profiles
1063 .iter()
1064 .map(|pr| ProfileResult {
1065 profile: profile_ref_name(pr),
1066 success: false,
1067 error: Some("At least one raw_signal is required.".into()),
1068 result: None,
1069 })
1070 .collect();
1071 }
1072
1073 let from = match parse_optional_datetime(&req.from) {
1074 Ok(value) => value,
1075 Err(error) => return profile_error_results(req, error.to_string()),
1076 };
1077 let to = match parse_optional_datetime(&req.to) {
1078 Ok(value) => value,
1079 Err(error) => return profile_error_results(req, error.to_string()),
1080 };
1081 let plan = match build_replay_plan(
1082 state,
1083 &req.symbol,
1084 &req.symbols,
1085 req.all_symbols,
1086 &req.raw_signals,
1087 from,
1088 to,
1089 future.signal_latency_ms,
1090 ) {
1091 Ok(plan) => plan,
1092 Err(error) => return profile_error_results(req, error.to_string()),
1093 };
1094 if let Err(error) = validate_replay_sizing_for_config(&req.config, &plan) {
1095 return profile_error_results(req, error.to_string());
1096 }
1097 let evaluation_options = match evaluation_options_from_msg_for_symbols(
1098 evaluation,
1099 &state.symbol_registry,
1100 plan.requested_symbols(),
1101 ) {
1102 Ok(options) => options,
1103 Err(error) => {
1104 return req
1105 .profiles
1106 .iter()
1107 .map(|profile| ProfileResult {
1108 profile: profile_ref_name(profile),
1109 success: false,
1110 error: Some(error.to_string()),
1111 result: None,
1112 })
1113 .collect();
1114 }
1115 };
1116 let config = match config_from_msg(&req.config, &state.symbol_registry, plan.active_symbols()) {
1117 Ok(config) => config,
1118 Err(error) => {
1119 return profile_error_results(req, error.to_string());
1120 }
1121 };
1122 let account_currency = match account_currency_from_msg(future) {
1123 Ok(account_currency) => account_currency,
1124 Err(error) => {
1125 return profile_error_results(req, error.to_string());
1126 }
1127 };
1128 let exchange = req.exchange.to_lowercase();
1129
1130 {
1131 let mut never_cancelled = || false;
1132 let primary = match describe_primary_market_stream(
1133 &state.data_dir,
1134 &exchange,
1135 plan.active_symbols(),
1136 &data_type,
1137 req.timeframe.as_deref(),
1138 plan.loading_start(),
1139 to,
1140 &mut never_cancelled,
1141 &mut |_| {},
1142 ) {
1143 Ok(primary) => primary,
1144 Err(error) => return profile_error_results(req, error.to_string()),
1145 };
1146 let bundle = match describe_future_stream(
1147 &state.data_dir,
1148 &exchange,
1149 &state.symbol_registry,
1150 &account_currency,
1151 plan.active_symbols(),
1152 &data_type,
1153 plan.loading_start(),
1154 primary,
1155 &mut never_cancelled,
1156 ) {
1157 Ok(bundle) => bundle,
1158 Err(error) => return profile_error_results(req, error.to_string()),
1159 };
1160 let future_config = match future_config_from_msg(future, bundle.currency_plan) {
1161 Ok(config) => config,
1162 Err(error) => return profile_error_results(req, error.to_string()),
1163 };
1164 let primary_eod = bundle.description.primary_eod();
1165 let metadata_request = single_request_from_multi(req);
1166
1167 req.profiles
1168 .iter()
1169 .map(|profile_ref| {
1170 let name = profile_ref_name(profile_ref);
1171 let profile = match profile_ref {
1172 ProfileRef::Named(profile_name) => state
1173 .profile_registry
1174 .read()
1175 .unwrap()
1176 .get(profile_name)
1177 .cloned()
1178 .ok_or_else(|| BacktestServerError::ProfileNotFound(profile_name.clone())),
1179 ProfileRef::Inline(message) => profile_from_msg(message).and_then(|profile| {
1180 profile.validate().map_err(|error| {
1181 BacktestServerError::InvalidRequest(format!(
1182 "Invalid inline profile: {error}"
1183 ))
1184 })?;
1185 Ok(profile)
1186 }),
1187 };
1188 let run_result = profile.and_then(|profile| {
1189 run_profile_streaming(
1190 &profile,
1191 plan.retained_signals(),
1192 &bundle.description,
1193 primary_eod,
1194 &config,
1195 &future_config,
1196 future,
1197 Some(&evaluation_options),
1198 state,
1199 &metadata_request,
1200 &plan,
1201 )
1202 });
1203 match run_result {
1204 Ok(result) => ProfileResult {
1205 profile: name,
1206 success: true,
1207 error: None,
1208 result: Some(result_to_msg(&result)),
1209 },
1210 Err(error) => ProfileResult {
1211 profile: name,
1212 success: false,
1213 error: Some(error.to_string()),
1214 result: None,
1215 },
1216 }
1217 })
1218 .collect()
1219 }
1220}
1221
1222fn single_request_from_multi(req: &BacktestMultiRunSpec) -> BacktestRunSpec {
1223 BacktestRunSpec {
1224 symbol: req.symbol.clone(),
1225 symbols: req.symbols.clone(),
1226 all_symbols: req.all_symbols,
1227 exchange: req.exchange.clone(),
1228 data_type: req.data_type.clone(),
1229 timeframe: req.timeframe.clone(),
1230 from: req.from.clone(),
1231 to: req.to.clone(),
1232 raw_signals: req.raw_signals.clone(),
1233 profile: None,
1234 profile_def: None,
1235 config: req.config.clone(),
1236 }
1237}
1238
1239fn profile_ref_name(pr: &ProfileRef) -> String {
1241 match pr {
1242 ProfileRef::Named(name) => name.clone(),
1243 ProfileRef::Inline(msg) => msg.name.clone(),
1244 }
1245}
1246
1247fn profile_error_results(req: &BacktestMultiRunSpec, error: String) -> Vec<ProfileResult> {
1248 req.profiles
1249 .iter()
1250 .map(|profile| ProfileResult {
1251 profile: profile_ref_name(profile),
1252 success: false,
1253 error: Some(error.clone()),
1254 result: None,
1255 })
1256 .collect()
1257}
1258
1259#[allow(clippy::too_many_arguments)]
1260fn run_profile_streaming(
1261 profile: &ManagementProfile,
1262 raw_signals: &[RawSignal],
1263 description: &MarketStreamDescription,
1264 primary_eod: Option<NaiveDateTime>,
1265 config: &BacktestConfig,
1266 future_config: &FutureQuoteConfig,
1267 future: &FutureQuoteConfigMsg,
1268 evaluation: Option<&EvaluationOptions>,
1269 state: &ServerState,
1270 metadata_request: &BacktestRunSpec,
1271 plan: &ReplayPlan,
1272) -> Result<BacktestResult> {
1273 let cancellation: CancellationCheck = Arc::new(|| false);
1274 let mut feed = description.open(cancellation)?;
1275 let mut runner = BacktestRunner::new_future(config.clone(), future_config.clone());
1276 if let Some(options) = evaluation {
1277 runner = runner.with_evaluation_options(options.clone());
1278 }
1279 let mut result = runner
1280 .run_raw_signals_future_streaming_controlled(
1281 &mut feed,
1282 primary_eod,
1283 raw_signals.to_vec(),
1284 Some(profile),
1285 || false,
1286 |_| {},
1287 )
1288 .map_err(map_streaming_replay_error)?;
1289 attach_future_reproducibility_metadata(
1290 &mut result,
1291 state,
1292 metadata_request,
1293 plan,
1294 future,
1295 Some(profile),
1296 );
1297 Ok(result)
1298}
1299
1300#[allow(clippy::too_many_arguments)]
1303fn build_replay_plan(
1304 state: &ServerState,
1305 symbol: &str,
1306 symbols: &[String],
1307 all_symbols: bool,
1308 raw_signal_msgs: &[RawSignalMsg],
1309 requested_from: Option<NaiveDateTime>,
1310 requested_to: Option<NaiveDateTime>,
1311 signal_latency_ms: i64,
1312) -> Result<ReplayPlan> {
1313 let scope =
1314 resolve_requested_symbol_scope(&state.symbol_registry, symbol, symbols, all_symbols)?;
1315 let raw_signals = raw_signal_msgs
1316 .iter()
1317 .map(|signal| raw_signal_from_msg(signal, scope.default_symbol(), &state.symbol_registry))
1318 .collect::<Result<Vec<_>>>()?;
1319 ReplayPlan::build(
1320 scope,
1321 raw_signals,
1322 requested_from,
1323 requested_to,
1324 signal_latency_ms,
1325 )
1326}
1327
1328fn resolve_requested_symbol_scope(
1329 registry: &SymbolRegistry,
1330 symbol: &str,
1331 symbols: &[String],
1332 all_symbols: bool,
1333) -> Result<RequestedSymbolScope> {
1334 if all_symbols {
1335 return Ok(RequestedSymbolScope::Inferred);
1336 }
1337
1338 let mut resolved = BTreeSet::new();
1339 if !symbols.is_empty() {
1340 for raw in symbols {
1341 let trimmed = raw.trim();
1342 if !trimmed.is_empty() {
1343 resolved.insert(normalize_symbol(registry, trimmed));
1344 }
1345 }
1346 if resolved.is_empty() {
1347 return Err(BacktestServerError::InvalidRequest(
1348 "symbols was provided but did not contain any non-empty symbol".into(),
1349 ));
1350 }
1351 } else {
1352 let trimmed = symbol.trim();
1353 if trimmed.is_empty() {
1354 return Err(BacktestServerError::InvalidRequest(
1355 "symbol is required unless symbols or all_symbols is provided".into(),
1356 ));
1357 }
1358 resolved.insert(normalize_symbol(registry, trimmed));
1359 }
1360
1361 Ok(RequestedSymbolScope::explicit(resolved))
1362}
1363
1364fn resolve_profile(
1366 state: &ServerState,
1367 req: &BacktestRunSpec,
1368) -> Result<Option<ManagementProfile>> {
1369 if let Some(ref profile_msg) = req.profile_def {
1370 let profile = profile_from_msg(profile_msg)?;
1371 profile.validate().map_err(|e| {
1372 BacktestServerError::InvalidRequest(format!("Invalid inline profile: {e}"))
1373 })?;
1374 Ok(Some(profile))
1375 } else if let Some(ref profile_name) = req.profile {
1376 let registry = state.profile_registry.read().unwrap();
1377 let profile = registry
1378 .get(profile_name)
1379 .ok_or_else(|| BacktestServerError::ProfileNotFound(profile_name.clone()))?;
1380 Ok(Some(profile.clone()))
1381 } else {
1382 Ok(None)
1383 }
1384}
1385
1386#[cfg(test)]
1390fn load_market_events(
1391 data_dir: &str,
1392 exchange: &str,
1393 symbol: &str,
1394 data_type: &str,
1395 timeframe: Option<&str>,
1396 from: Option<NaiveDateTime>,
1397 to: Option<NaiveDateTime>,
1398) -> Result<Vec<MarketEvent>> {
1399 load_market_events_controlled(
1400 data_dir, exchange, symbol, data_type, timeframe, from, to, None,
1401 )
1402}
1403
1404#[cfg(test)]
1405#[allow(clippy::too_many_arguments)]
1406fn load_market_events_controlled(
1407 data_dir: &str,
1408 exchange: &str,
1409 symbol: &str,
1410 data_type: &str,
1411 timeframe: Option<&str>,
1412 from: Option<NaiveDateTime>,
1413 to: Option<NaiveDateTime>,
1414 cancellation: Option<&JobCancellationToken>,
1415) -> Result<Vec<MarketEvent>> {
1416 ensure_not_cancelled(cancellation)?;
1417 let store = ParquetStore::open(data_dir)?;
1418 let dt = data_type.to_lowercase();
1419
1420 if dt == "tick" {
1421 let disk_exchange =
1422 resolve_partition_value(data_dir, "ticks", "exchange", exchange, "", cancellation)?;
1423 let disk_symbol = resolve_partition_value(
1424 data_dir,
1425 "ticks",
1426 "symbol",
1427 symbol,
1428 &format!("exchange={disk_exchange}"),
1429 cancellation,
1430 )?;
1431 let opts = QueryOpts {
1432 exchange: disk_exchange.clone(),
1433 symbol: disk_symbol.clone(),
1434 from,
1435 to,
1436 limit: 0,
1437 tail: false,
1438 descending: false,
1439 };
1440 let (ticks, _total) = store
1441 .query_ticks_cancellable(&opts, || {
1442 cancellation.is_some_and(JobCancellationToken::is_cancelled)
1443 })
1444 .map_err(map_data_cancellation)?;
1445 if ticks.is_empty() {
1446 return Err(BacktestServerError::NoDataFound {
1447 symbol: disk_symbol,
1448 exchange: disk_exchange,
1449 data_type: "tick".into(),
1450 });
1451 }
1452 let feed = ticks_to_feed(ticks);
1453 Ok(canonicalize_market_event_symbols(
1454 feed_to_events(feed),
1455 symbol,
1456 ))
1457 } else if dt == "bar" {
1458 let tf_str = timeframe.ok_or_else(|| {
1459 BacktestServerError::InvalidRequest("timeframe is required for bar data".into())
1460 })?;
1461 let tf = Timeframe::parse(tf_str).map_err(|_| {
1462 BacktestServerError::InvalidRequest(format!("Invalid timeframe: '{tf_str}'"))
1463 })?;
1464 let disk_exchange =
1465 resolve_partition_value(data_dir, "bars", "exchange", exchange, "", cancellation)?;
1466 let disk_symbol = resolve_partition_value(
1467 data_dir,
1468 "bars",
1469 "symbol",
1470 symbol,
1471 &format!("exchange={disk_exchange}"),
1472 cancellation,
1473 )?;
1474 let disk_timeframe = resolve_partition_value(
1475 data_dir,
1476 "bars",
1477 "timeframe",
1478 tf.as_str(),
1479 &format!("exchange={disk_exchange}/symbol={disk_symbol}"),
1480 cancellation,
1481 )?;
1482 let opts = BarQueryOpts {
1483 exchange: disk_exchange.clone(),
1484 symbol: disk_symbol.clone(),
1485 timeframe: disk_timeframe.clone(),
1486 from,
1487 to,
1488 limit: 0,
1489 tail: false,
1490 descending: false,
1491 };
1492 let (bars, _total) = store
1493 .query_bars_cancellable(&opts, || {
1494 cancellation.is_some_and(JobCancellationToken::is_cancelled)
1495 })
1496 .map_err(map_data_cancellation)?;
1497 if bars.is_empty() {
1498 return Err(BacktestServerError::NoDataFound {
1499 symbol: disk_symbol,
1500 exchange: disk_exchange,
1501 data_type: format!("bar({})", disk_timeframe),
1502 });
1503 }
1504 let feed = bars_to_feed(bars);
1505 Ok(canonicalize_market_event_symbols(
1506 feed_to_events(feed),
1507 symbol,
1508 ))
1509 } else {
1510 Err(BacktestServerError::InvalidRequest(format!(
1511 "Invalid data_type: '{}'. Must be 'tick' or 'bar'.",
1512 data_type
1513 )))
1514 }
1515}
1516
1517#[cfg(test)]
1518fn map_data_cancellation(error: DataError) -> BacktestServerError {
1519 match error {
1520 DataError::Cancelled => BacktestServerError::Cancelled,
1521 other => BacktestServerError::Database(other),
1522 }
1523}
1524
1525#[cfg(test)]
1527fn resolve_partition_value(
1528 data_dir: &str,
1529 data_subdir: &str,
1530 key: &str,
1531 requested: &str,
1532 parent: &str,
1533 cancellation: Option<&JobCancellationToken>,
1534) -> Result<String> {
1535 let dir = if parent.is_empty() {
1536 Path::new(data_dir).join(data_subdir)
1537 } else {
1538 Path::new(data_dir).join(data_subdir).join(parent)
1539 };
1540 let prefix = format!("{key}=");
1541 let mut case_insensitive_match = None;
1542
1543 ensure_not_cancelled(cancellation)?;
1544 if let Ok(entries) = std::fs::read_dir(dir) {
1545 for entry in entries.flatten() {
1546 ensure_not_cancelled(cancellation)?;
1547 let name = entry.file_name();
1548 let name = name.to_string_lossy();
1549 let Some(value) = name.strip_prefix(&prefix) else {
1550 continue;
1551 };
1552 if value == requested {
1553 return Ok(value.to_string());
1554 }
1555 if case_insensitive_match.is_none() && value.eq_ignore_ascii_case(requested) {
1556 case_insensitive_match = Some(value.to_string());
1557 }
1558 }
1559 }
1560
1561 ensure_not_cancelled(cancellation)?;
1562 Ok(case_insensitive_match.unwrap_or_else(|| requested.to_string()))
1563}
1564
1565#[cfg(test)]
1567fn canonicalize_market_event_symbols(
1568 mut events: Vec<MarketEvent>,
1569 canonical_symbol: &str,
1570) -> Vec<MarketEvent> {
1571 for event in &mut events {
1572 match event {
1573 MarketEvent::Tick { symbol, .. } | MarketEvent::Bar { symbol, .. } => {
1574 *symbol = canonical_symbol.to_string();
1575 }
1576 }
1577 }
1578 events
1579}
1580
1581#[cfg(test)]
1583fn feed_to_events(mut feed: VecFeed) -> Vec<MarketEvent> {
1584 let mut events = Vec::with_capacity(feed.total());
1585 while let Some(event) = feed.next_event() {
1586 events.push(event);
1587 }
1588 events
1589}
1590
1591fn validate_request(req: &BacktestRunSpec) -> Result<()> {
1595 let dt = req.data_type.to_lowercase();
1596 if dt != "tick" && dt != "bar" {
1597 return Err(BacktestServerError::InvalidRequest(format!(
1598 "Invalid data_type: '{}'. Must be 'tick' or 'bar'.",
1599 req.data_type
1600 )));
1601 }
1602 if dt == "bar" && req.timeframe.is_none() {
1603 return Err(BacktestServerError::InvalidRequest(
1604 "timeframe is required when data_type is 'bar'.".into(),
1605 ));
1606 }
1607 if req.raw_signals.is_empty() {
1608 return Err(BacktestServerError::InvalidRequest(
1609 "At least one raw_signal is required.".into(),
1610 ));
1611 }
1612 Ok(())
1613}
1614
1615fn validate_replay_sizing(req: &BacktestRunSpec, plan: &ReplayPlan) -> Result<()> {
1616 validate_replay_sizing_for_config(&req.config, plan)
1617}
1618
1619fn validate_replay_sizing_for_config(config: &BacktestConfigMsg, plan: &ReplayPlan) -> Result<()> {
1620 if !plan.is_idle() && config.sizing.is_none() {
1621 return Err(BacktestServerError::InvalidRequest(
1622 "Entry signals require an account sizing policy.".into(),
1623 ));
1624 }
1625 Ok(())
1626}
1627
1628fn normalize_symbol(registry: &SymbolRegistry, raw: &str) -> String {
1632 registry.normalize_or_passthrough(raw)
1633}
1634
1635fn parse_datetime(s: &str) -> Result<NaiveDateTime> {
1637 let formats = [
1639 "%Y-%m-%dT%H:%M:%S%.f",
1640 "%Y-%m-%dT%H:%M:%S",
1641 "%Y-%m-%d %H:%M:%S%.f",
1642 "%Y-%m-%d %H:%M:%S",
1643 "%Y-%m-%d",
1644 ];
1645 for fmt in &formats {
1646 if let Ok(dt) = NaiveDateTime::parse_from_str(s, fmt) {
1647 return Ok(dt);
1648 }
1649 }
1650 if let Ok(date) = chrono::NaiveDate::parse_from_str(s, "%Y-%m-%d") {
1652 return Ok(date.and_hms_opt(0, 0, 0).unwrap());
1653 }
1654 Err(BacktestServerError::InvalidRequest(format!(
1655 "Cannot parse datetime: '{s}'. Use ISO format (e.g. '2026-01-15T10:30:00' or '2026-01-15')."
1656 )))
1657}
1658
1659fn parse_optional_datetime(s: &Option<String>) -> Result<Option<NaiveDateTime>> {
1661 match s {
1662 Some(v) => parse_datetime(v).map(Some),
1663 None => Ok(None),
1664 }
1665}
1666
1667pub fn handle_add_profile(state: &ServerState, req: &AddProfileRequest) -> AddProfileResponse {
1671 let profile = match profile_from_msg(&req.profile) {
1672 Ok(p) => p,
1673 Err(e) => {
1674 let registry = state.profile_registry.read().unwrap();
1675 return AddProfileResponse {
1676 success: false,
1677 error: Some(e.to_string()),
1678 profile_count: registry.len(),
1679 };
1680 }
1681 };
1682
1683 let mut registry = state.profile_registry.write().unwrap();
1684 match registry.insert(profile, req.overwrite) {
1685 Ok(()) => AddProfileResponse {
1686 success: true,
1687 error: None,
1688 profile_count: registry.len(),
1689 },
1690 Err(e) => AddProfileResponse {
1691 success: false,
1692 error: Some(e.to_string()),
1693 profile_count: registry.len(),
1694 },
1695 }
1696}
1697
1698pub fn handle_remove_profile(
1700 state: &ServerState,
1701 req: &RemoveProfileRequest,
1702) -> RemoveProfileResponse {
1703 let mut registry = state.profile_registry.write().unwrap();
1704 let removed = registry.remove(&req.name);
1705 RemoveProfileResponse {
1706 success: removed,
1707 error: if removed {
1708 None
1709 } else {
1710 Some(format!("Profile '{}' not found", req.name))
1711 },
1712 profile_count: registry.len(),
1713 }
1714}
1715
1716pub fn handle_reload_profiles(state: &ServerState) -> ReloadProfilesResponse {
1718 if state.profiles_path.is_empty() {
1719 return ReloadProfilesResponse {
1720 success: false,
1721 error: Some("No profiles_path configured".into()),
1722 profile_count: state.profile_registry.read().unwrap().len(),
1723 loaded_from: String::new(),
1724 };
1725 }
1726
1727 match ProfileRegistry::load(&state.profiles_path) {
1728 Ok(new_registry) => {
1729 let count = new_registry.len();
1730 let mut registry = state.profile_registry.write().unwrap();
1731 *registry = new_registry;
1732 ReloadProfilesResponse {
1733 success: true,
1734 error: None,
1735 profile_count: count,
1736 loaded_from: state.profiles_path.clone(),
1737 }
1738 }
1739 Err(e) => ReloadProfilesResponse {
1740 success: false,
1741 error: Some(e.to_string()),
1742 profile_count: state.profile_registry.read().unwrap().len(),
1743 loaded_from: state.profiles_path.clone(),
1744 },
1745 }
1746}
1747
1748fn admit_backtest_job(state: &ServerState) -> SubmitBacktestResponse {
1752 let job_id = format!("job-{}", uuid_v4_simple());
1753 let initial_progress = BacktestProgress {
1754 stage: "queued".into(),
1755 ..BacktestProgress::default()
1756 };
1757 let (updates, _) = watch::channel(BacktestStatusResponse {
1758 success: true,
1759 job_id: job_id.clone(),
1760 status: JobStatus::Queued.as_str().into(),
1761 error: None,
1762 elapsed_ms: None,
1763 progress: initial_progress.clone(),
1764 });
1765 let job = BacktestJob {
1766 status: JobStatus::Queued,
1767 submitted_at: Instant::now(),
1768 completed_at: None,
1769 progress: initial_progress,
1770 result: None,
1771 artifact: None,
1772 inline_complete: false,
1773 artifact_consumed: false,
1774 error: None,
1775 cancellation: JobCancellationToken::default(),
1776 worker_active: true,
1777 updates,
1778 };
1779 if state.max_retained_jobs == 0 {
1780 return SubmitBacktestResponse {
1781 success: false,
1782 job_id: None,
1783 error: Some("Async job retention limit is zero".into()),
1784 };
1785 }
1786 let mut evicted = Vec::new();
1787 {
1788 let mut jobs = state.jobs.lock().unwrap();
1789 while jobs.len() >= state.max_retained_jobs {
1790 let Some(removed) = remove_oldest_terminal_job(&mut jobs) else {
1791 break;
1792 };
1793 evicted.push(removed);
1794 }
1795 if jobs.len() >= state.max_retained_jobs {
1796 drop(jobs);
1797 delete_job_artifacts(state, &evicted);
1798 return SubmitBacktestResponse {
1799 success: false,
1800 job_id: None,
1801 error: Some(format!(
1802 "Async job limit reached (max {})",
1803 state.max_retained_jobs
1804 )),
1805 };
1806 }
1807 jobs.insert(job_id.clone(), job);
1808 }
1809 delete_job_artifacts(state, &evicted);
1810 SubmitBacktestResponse {
1811 success: true,
1812 job_id: Some(job_id),
1813 error: None,
1814 }
1815}
1816
1817pub fn handle_submit_backtest(
1818 state: &ServerState,
1819 req: &SubmitBacktestRequest,
1820) -> SubmitBacktestResponse {
1821 let validation = (|| -> Result<()> {
1822 let request = &req.request.request;
1823 validate_future_quote_scalars(&req.request.future)?;
1824 validate_request(request)?;
1825 let from = parse_optional_datetime(&request.from)?;
1826 let to = parse_optional_datetime(&request.to)?;
1827 let plan = build_replay_plan(
1828 state,
1829 &request.symbol,
1830 &request.symbols,
1831 request.all_symbols,
1832 &request.raw_signals,
1833 from,
1834 to,
1835 req.request.future.signal_latency_ms,
1836 )?;
1837 validate_replay_sizing(request, &plan)?;
1838 account_currency_from_msg(&req.request.future)?;
1839 config_from_msg(
1840 &request.config,
1841 &state.symbol_registry,
1842 plan.active_symbols(),
1843 )?;
1844 evaluation_options_from_msg_for_symbols(
1845 &req.request.evaluation,
1846 &state.symbol_registry,
1847 plan.requested_symbols(),
1848 )?;
1849 Ok(())
1850 })();
1851 if let Err(error) = validation {
1852 return SubmitBacktestResponse {
1853 success: false,
1854 job_id: None,
1855 error: Some(error.to_string()),
1856 };
1857 }
1858 admit_backtest_job(state)
1859}
1860
1861pub fn handle_get_backtest_status(
1863 state: &ServerState,
1864 req: &GetBacktestStatusRequest,
1865) -> BacktestStatusResponse {
1866 let jobs = state.jobs.lock().unwrap();
1867 match jobs.get(&req.job_id) {
1868 Some(job) => job_status_response(&req.job_id, job),
1869 None => BacktestStatusResponse {
1870 success: false,
1871 job_id: req.job_id.clone(),
1872 status: "NotFound".into(),
1873 error: Some(format!("Job '{}' not found", req.job_id)),
1874 elapsed_ms: None,
1875 progress: BacktestProgress::default(),
1876 },
1877 }
1878}
1879
1880pub fn handle_get_backtest_result(
1882 state: &ServerState,
1883 req: &GetBacktestResultRequest,
1884) -> GetBacktestResultResponse {
1885 let jobs = state.jobs.lock().unwrap();
1886 match jobs.get(&req.job_id) {
1887 Some(job) if job.status == JobStatus::Completed && job.artifact_consumed => {
1888 GetBacktestResultResponse {
1889 success: false,
1890 job_id: req.job_id.clone(),
1891 result: job.result.clone(),
1892 error: Some("Job result artifact has already been consumed".into()),
1893 artifact: None,
1894 inline_complete: false,
1895 artifact_consumed: true,
1896 }
1897 }
1898 Some(job) if job.status == JobStatus::Completed => GetBacktestResultResponse {
1899 success: true,
1900 job_id: req.job_id.clone(),
1901 result: job.result.clone(),
1902 error: None,
1903 artifact: job.artifact.clone(),
1904 inline_complete: job.inline_complete,
1905 artifact_consumed: false,
1906 },
1907 Some(job) => GetBacktestResultResponse {
1908 success: false,
1909 job_id: req.job_id.clone(),
1910 result: None,
1911 error: Some(format!(
1912 "Job is not completed (status: {})",
1913 job.status.as_str()
1914 )),
1915 artifact: None,
1916 inline_complete: true,
1917 artifact_consumed: false,
1918 },
1919 None => GetBacktestResultResponse {
1920 success: false,
1921 job_id: req.job_id.clone(),
1922 result: None,
1923 error: Some(format!("Job '{}' not found", req.job_id)),
1924 artifact: None,
1925 inline_complete: true,
1926 artifact_consumed: false,
1927 },
1928 }
1929}
1930
1931pub fn handle_get_result_artifact_chunk(
1933 state: &ServerState,
1934 req: &GetResultArtifactChunkRequest,
1935) -> GetResultArtifactChunkResponse {
1936 match state
1937 .artifact_store
1938 .read_chunk(&req.artifact_id, req.offset)
1939 {
1940 Ok(chunk) => GetResultArtifactChunkResponse {
1941 success: true,
1942 artifact_id: req.artifact_id.clone(),
1943 offset: chunk.offset,
1944 data_base64: BASE64_STANDARD.encode(chunk.bytes),
1945 eof: chunk.eof,
1946 error: None,
1947 },
1948 Err(error) => GetResultArtifactChunkResponse {
1949 success: false,
1950 artifact_id: req.artifact_id.clone(),
1951 offset: req.offset,
1952 data_base64: String::new(),
1953 eof: false,
1954 error: Some(error.to_string()),
1955 },
1956 }
1957}
1958
1959pub fn handle_delete_result_artifact(
1961 state: &ServerState,
1962 req: &DeleteResultArtifactRequest,
1963) -> DeleteResultArtifactResponse {
1964 let mut jobs = state.jobs.lock().unwrap();
1965 match state.artifact_store.delete(&req.artifact_id) {
1966 Ok(deleted) => {
1967 for job in jobs.values_mut().filter(|job| {
1968 job.artifact
1969 .as_ref()
1970 .is_some_and(|artifact| artifact.artifact_id == req.artifact_id)
1971 }) {
1972 job.artifact = None;
1973 job.artifact_consumed = true;
1974 job.inline_complete = false;
1975 }
1976 if deleted {
1977 DeleteResultArtifactResponse {
1978 success: true,
1979 artifact_id: req.artifact_id.clone(),
1980 error: None,
1981 }
1982 } else {
1983 DeleteResultArtifactResponse {
1984 success: false,
1985 artifact_id: req.artifact_id.clone(),
1986 error: Some(format!("Artifact '{}' not found", req.artifact_id)),
1987 }
1988 }
1989 }
1990 Err(error) => DeleteResultArtifactResponse {
1991 success: false,
1992 artifact_id: req.artifact_id.clone(),
1993 error: Some(error.to_string()),
1994 },
1995 }
1996}
1997
1998pub fn handle_cancel_backtest(
2000 state: &ServerState,
2001 req: &CancelBacktestRequest,
2002) -> CancelBacktestResponse {
2003 let mut jobs = state.jobs.lock().unwrap();
2004 match jobs.get_mut(&req.job_id) {
2005 Some(job)
2006 if job.status == JobStatus::Queued
2007 || job.status == JobStatus::LoadingData
2008 || job.status == JobStatus::Running =>
2009 {
2010 job.cancellation.cancel();
2011 job.status = JobStatus::Cancelled;
2012 job.completed_at = Some(Instant::now());
2013 job.result = None;
2014 job.artifact = None;
2015 job.inline_complete = true;
2016 job.artifact_consumed = false;
2017 job.error = None;
2018 job.progress.stage = "cancelled".into();
2019 publish_job_status(&req.job_id, job);
2020 CancelBacktestResponse {
2021 success: true,
2022 job_id: req.job_id.clone(),
2023 error: None,
2024 }
2025 }
2026 Some(job) => CancelBacktestResponse {
2027 success: false,
2028 job_id: req.job_id.clone(),
2029 error: Some(format!(
2030 "Cannot cancel job in status: {}",
2031 job.status.as_str()
2032 )),
2033 },
2034 None => CancelBacktestResponse {
2035 success: false,
2036 job_id: req.job_id.clone(),
2037 error: Some(format!("Job '{}' not found", req.job_id)),
2038 },
2039 }
2040}
2041
2042pub fn run_job_and_store(state: Arc<ServerState>, job_id: String, req: RunBacktestRequest) {
2043 run_job_and_store_inner(
2044 state,
2045 job_id,
2046 req.request,
2047 req.future,
2048 req.evaluation,
2049 req.result_delivery,
2050 );
2051}
2052
2053fn run_job_and_store_inner(
2054 state: Arc<ServerState>,
2055 job_id: String,
2056 req: BacktestRunSpec,
2057 future: FutureQuoteConfigMsg,
2058 evaluation: ProviderEvaluationOptionsMsg,
2059 delivery: ResultDeliveryMsg,
2060) {
2061 let cancellation = {
2062 let mut jobs = state.jobs.lock().unwrap();
2063 let Some(job) = jobs.get_mut(&job_id) else {
2064 return;
2065 };
2066 if job.status == JobStatus::Cancelled || job.cancellation.is_cancelled() {
2067 job.worker_active = false;
2068 return;
2069 }
2070 job.status = JobStatus::LoadingData;
2071 job.progress.stage = "loading_data".into();
2072 publish_job_status(&job_id, job);
2073 job.cancellation.clone()
2074 };
2075
2076 let result = execute_backtest_with_future_controlled(
2077 &state,
2078 &req,
2079 &future,
2080 &evaluation,
2081 Some(&cancellation),
2082 &mut |progress| update_job_progress(&state, &job_id, progress),
2083 );
2084
2085 let execution_cancelled = matches!(&result, Err(BacktestServerError::Cancelled));
2086 let prepared = if execution_cancelled {
2087 None
2088 } else {
2089 Some(match result {
2090 Ok(backtest_result) => prepare_result(
2091 &state,
2092 result_to_msg(&backtest_result),
2093 Some(delivery),
2094 compact_result_for_console,
2095 ),
2096 Err(error) => Err(error.to_string()),
2097 })
2098 };
2099
2100 let mut jobs = state.jobs.lock().unwrap();
2101 let Some(job) = jobs.get_mut(&job_id) else {
2102 return;
2103 };
2104 job.worker_active = false;
2105 if cancellation.is_cancelled() || job.status == JobStatus::Cancelled || execution_cancelled {
2106 if let Some(Ok(PreparedResult::Artifact { reference, .. })) = prepared.as_ref() {
2107 let _ = state.artifact_store.delete(&reference.artifact_id);
2108 }
2109 job.cancellation.cancel();
2110 job.status = JobStatus::Cancelled;
2111 job.result = None;
2112 job.artifact = None;
2113 job.inline_complete = true;
2114 job.artifact_consumed = false;
2115 job.error = None;
2116 job.progress.stage = "cancelled".into();
2117 job.completed_at.get_or_insert_with(Instant::now);
2118 publish_job_status(&job_id, job);
2119 return;
2120 }
2121
2122 match prepared.expect("non-cancelled execution has a prepared result") {
2123 Ok(PreparedResult::Inline(result)) => {
2124 job.status = JobStatus::Completed;
2125 job.result = Some(result);
2126 job.artifact = None;
2127 job.inline_complete = true;
2128 job.artifact_consumed = false;
2129 job.error = None;
2130 job.progress.stage = "completed".into();
2131 job.completed_at = Some(Instant::now());
2132 publish_job_status(&job_id, job);
2133 }
2134 Ok(PreparedResult::Artifact { reference, summary }) => {
2135 job.status = JobStatus::Completed;
2136 job.result = summary;
2137 job.artifact = Some(reference);
2138 job.inline_complete = false;
2139 job.artifact_consumed = false;
2140 job.error = None;
2141 job.progress.stage = "completed".into();
2142 job.completed_at = Some(Instant::now());
2143 publish_job_status(&job_id, job);
2144 }
2145 Err(error) => {
2146 job.status = JobStatus::Failed;
2147 job.result = None;
2148 job.artifact = None;
2149 job.inline_complete = true;
2150 job.artifact_consumed = false;
2151 job.error = Some(error);
2152 job.progress.stage = "failed".into();
2153 job.completed_at = Some(Instant::now());
2154 publish_job_status(&job_id, job);
2155 }
2156 }
2157}
2158
2159fn uuid_v4_simple() -> String {
2161 use std::time::SystemTime;
2162 let now = SystemTime::now()
2163 .duration_since(SystemTime::UNIX_EPOCH)
2164 .unwrap_or_default();
2165 let nanos = now.as_nanos();
2166 format!("{:x}", nanos)
2167}
2168
2169#[cfg(test)]
2170mod tests {
2171 use super::*;
2172 #[allow(unused_imports)]
2173 use chrono::NaiveDate;
2174 #[allow(unused_imports)]
2175 use qs_core::types::{OrderType, Side};
2176 #[allow(unused_imports)]
2177 use std::sync::Arc;
2178
2179 fn test_request(request: BacktestRunSpec) -> RunBacktestRequest {
2180 RunBacktestRequest {
2181 request,
2182 future: FutureQuoteConfigMsg {
2183 account_currency: "USD".into(),
2184 ..FutureQuoteConfigMsg::default()
2185 },
2186 evaluation: ProviderEvaluationOptionsMsg::default(),
2187 result_delivery: ResultDeliveryMsg::Auto,
2188 }
2189 }
2190
2191 fn submit_for_test(state: &ServerState, request: &BacktestRunSpec) -> SubmitBacktestResponse {
2192 handle_submit_backtest(
2193 state,
2194 &SubmitBacktestRequest {
2195 request: test_request(request.clone()),
2196 },
2197 )
2198 }
2199
2200 fn run_job_for_test(state: Arc<ServerState>, job_id: String, request: BacktestRunSpec) {
2201 run_job_and_store(state, job_id, test_request(request));
2202 }
2203
2204 fn test_state() -> ServerState {
2205 let artifact_directory = std::env::temp_dir().join(format!(
2206 "qs_backtest_server_handler_artifacts_{}",
2207 std::process::id()
2208 ));
2209 ServerState {
2210 symbol_registry: SymbolRegistry::empty(),
2211 profile_registry: RwLock::new(ProfileRegistry::empty()),
2212 data_dir: "/tmp/test".into(),
2213 profiles_path: String::new(),
2214 start_time: Instant::now(),
2215 jobs: std::sync::Mutex::new(std::collections::HashMap::new()),
2216 max_retained_jobs: 1_000,
2217 artifact_store: ArtifactStore::new(
2218 artifact_directory,
2219 12 * 1024 * 1024,
2220 1024 * 1024,
2221 Duration::from_secs(3_600),
2222 1024 * 1024 * 1024,
2223 )
2224 .unwrap(),
2225 }
2226 }
2227
2228 #[test]
2229 fn parse_datetime_iso() {
2230 let dt = parse_datetime("2026-01-15T10:30:00").unwrap();
2231 assert_eq!(dt.to_string(), "2026-01-15 10:30:00");
2232 }
2233
2234 #[test]
2235 fn parse_datetime_space_separator() {
2236 let dt = parse_datetime("2026-01-15 10:30:00").unwrap();
2237 assert_eq!(dt.to_string(), "2026-01-15 10:30:00");
2238 }
2239
2240 #[test]
2241 fn parse_datetime_date_only() {
2242 let dt = parse_datetime("2026-01-15").unwrap();
2243 assert_eq!(dt.to_string(), "2026-01-15 00:00:00");
2244 }
2245
2246 #[test]
2247 fn parse_datetime_invalid() {
2248 assert!(parse_datetime("not-a-date").is_err());
2249 }
2250
2251 #[test]
2252 fn parse_optional_datetime_none() {
2253 assert!(parse_optional_datetime(&None).unwrap().is_none());
2254 }
2255
2256 #[test]
2257 fn parse_optional_datetime_some() {
2258 let result = parse_optional_datetime(&Some("2026-01-15".into())).unwrap();
2259 assert!(result.is_some());
2260 }
2261
2262 #[test]
2263 fn normalize_symbol_passthrough() {
2264 let reg = SymbolRegistry::empty();
2265 assert_eq!(normalize_symbol(®, "BTCUSD"), "btcusd");
2266 }
2267
2268 #[test]
2269 fn requested_symbol_scope_normalizes_explicit_symbols() {
2270 let scope = resolve_requested_symbol_scope(
2271 &SymbolRegistry::empty(),
2272 "",
2273 &["XAU/USD".into(), " GBPJPY ".into()],
2274 false,
2275 )
2276 .unwrap();
2277
2278 assert_eq!(
2279 scope,
2280 RequestedSymbolScope::explicit(["gbpjpy".into(), "xauusd".into()])
2281 );
2282 }
2283
2284 #[test]
2285 fn replay_plan_derives_inferred_symbols_from_retained_entries() {
2286 let state = test_state();
2287 let signals = vec![
2288 RawSignalMsg::Entry {
2289 ts: "2026-01-15T10:00:00".into(),
2290 symbol: "XAUUSD".into(),
2291 side: "Buy".into(),
2292 order_type: "Market".into(),
2293 price: Some(2000.0),
2294 risk: 1.0,
2295 stoploss: None,
2296 targets: vec![],
2297 group: None,
2298 trade_id: None,
2299 },
2300 RawSignalMsg::Entry {
2301 ts: "2026-01-15T10:01:00".into(),
2302 symbol: "GBP/JPY".into(),
2303 side: "Sell".into(),
2304 order_type: "Market".into(),
2305 price: Some(190.0),
2306 risk: 1.0,
2307 stoploss: None,
2308 targets: vec![],
2309 group: None,
2310 trade_id: None,
2311 },
2312 ];
2313 let plan = build_replay_plan(&state, "", &[], true, &signals, None, None, 0).unwrap();
2314
2315 assert_eq!(plan.active_symbols(), ["gbpjpy", "xauusd"]);
2316 }
2317
2318 #[test]
2319 fn replay_plan_preserves_explicit_single_symbol_default() {
2320 let state = test_state();
2321 let signals = vec![RawSignalMsg::Entry {
2322 ts: "2026-01-15T10:00:00".into(),
2323 symbol: "".into(),
2324 side: "Buy".into(),
2325 order_type: "Market".into(),
2326 price: Some(2000.0),
2327 risk: 1.0,
2328 stoploss: None,
2329 targets: vec![],
2330 group: None,
2331 trade_id: None,
2332 }];
2333 let replay =
2334 build_replay_plan(&state, "XAU/USD", &[], false, &signals, None, None, 0).unwrap();
2335
2336 assert_eq!(replay.active_symbols(), ["xauusd"]);
2337 }
2338
2339 #[test]
2340 fn replay_plan_rejects_missing_entry_symbol_in_explicit_multi_scope() {
2341 let state = test_state();
2342 let signals = vec![RawSignalMsg::Entry {
2343 ts: "2026-01-15T10:00:00".into(),
2344 symbol: "".into(),
2345 side: "Buy".into(),
2346 order_type: "Market".into(),
2347 price: Some(2000.0),
2348 risk: 1.0,
2349 stoploss: None,
2350 targets: vec![],
2351 group: None,
2352 trade_id: None,
2353 }];
2354 let error = build_replay_plan(
2355 &state,
2356 "",
2357 &["xauusd".into(), "gbpjpy".into()],
2358 false,
2359 &signals,
2360 None,
2361 None,
2362 0,
2363 )
2364 .unwrap_err();
2365
2366 assert!(error.to_string().contains("symbol is required"));
2367 }
2368
2369 #[test]
2370 fn validate_request_rejects_empty_signals() {
2371 let req = BacktestRunSpec {
2372 symbol: "eurusd".into(),
2373 symbols: Vec::new(),
2374 all_symbols: false,
2375 exchange: "ctrader".into(),
2376 data_type: "tick".into(),
2377 timeframe: None,
2378 from: None,
2379 to: None,
2380 raw_signals: vec![],
2381 profile: None,
2382 profile_def: None,
2383 config: BacktestConfigMsg {
2384 initial_balance: None,
2385 close_on_finish: None,
2386 fill_model: None,
2387 sizing: None,
2388 },
2389 };
2390 assert!(validate_request(&req).is_err());
2391 }
2392
2393 #[test]
2394 fn validate_request_rejects_invalid_data_type() {
2395 let req = BacktestRunSpec {
2396 symbol: "eurusd".into(),
2397 symbols: Vec::new(),
2398 all_symbols: false,
2399 exchange: "ctrader".into(),
2400 data_type: "invalid".into(),
2401 timeframe: None,
2402 from: None,
2403 to: None,
2404 raw_signals: vec![RawSignalMsg::Entry {
2405 ts: "2026-01-15T10:00:00".into(),
2406 symbol: "eurusd".into(),
2407 side: "Buy".into(),
2408 order_type: "Market".into(),
2409 price: None,
2410 risk: 1.0,
2411 stoploss: None,
2412 targets: vec![],
2413 group: None,
2414 trade_id: None,
2415 }],
2416 profile: None,
2417 profile_def: None,
2418 config: BacktestConfigMsg {
2419 initial_balance: None,
2420 close_on_finish: None,
2421 fill_model: None,
2422 sizing: None,
2423 },
2424 };
2425 assert!(validate_request(&req).is_err());
2426 }
2427
2428 #[test]
2429 fn validate_request_rejects_bar_without_timeframe() {
2430 let req = BacktestRunSpec {
2431 symbol: "eurusd".into(),
2432 symbols: Vec::new(),
2433 all_symbols: false,
2434 exchange: "ctrader".into(),
2435 data_type: "bar".into(),
2436 timeframe: None,
2437 from: None,
2438 to: None,
2439 raw_signals: vec![RawSignalMsg::Entry {
2440 ts: "2026-01-15T10:00:00".into(),
2441 symbol: "eurusd".into(),
2442 side: "Buy".into(),
2443 order_type: "Market".into(),
2444 price: None,
2445 risk: 1.0,
2446 stoploss: None,
2447 targets: vec![],
2448 group: None,
2449 trade_id: None,
2450 }],
2451 profile: None,
2452 profile_def: None,
2453 config: BacktestConfigMsg {
2454 initial_balance: None,
2455 close_on_finish: None,
2456 fill_model: None,
2457 sizing: None,
2458 },
2459 };
2460 assert!(validate_request(&req).is_err());
2461 }
2462
2463 #[test]
2464 fn validate_request_accepts_valid_tick() {
2465 let req = BacktestRunSpec {
2466 symbol: "eurusd".into(),
2467 symbols: Vec::new(),
2468 all_symbols: false,
2469 exchange: "ctrader".into(),
2470 data_type: "tick".into(),
2471 timeframe: None,
2472 from: None,
2473 to: None,
2474 raw_signals: vec![RawSignalMsg::Entry {
2475 ts: "2026-01-15T10:00:00".into(),
2476 symbol: "eurusd".into(),
2477 side: "Buy".into(),
2478 order_type: "Market".into(),
2479 price: None,
2480 risk: 1.0,
2481 stoploss: None,
2482 targets: vec![],
2483 group: None,
2484 trade_id: None,
2485 }],
2486 profile: None,
2487 profile_def: None,
2488 config: BacktestConfigMsg {
2489 initial_balance: None,
2490 close_on_finish: None,
2491 fill_model: None,
2492 sizing: Some(SizingPolicyMsg::FixedLot { lots: 0.01 }),
2493 },
2494 };
2495 assert!(validate_request(&req).is_ok());
2496 }
2497
2498 #[test]
2499 fn validate_request_rejects_entry_without_sizing() {
2500 let mut req = BacktestRunSpec {
2501 symbol: "eurusd".into(),
2502 symbols: Vec::new(),
2503 all_symbols: false,
2504 exchange: "ctrader".into(),
2505 data_type: "tick".into(),
2506 timeframe: None,
2507 from: None,
2508 to: None,
2509 raw_signals: vec![RawSignalMsg::Entry {
2510 ts: "2026-01-15T10:00:00".into(),
2511 symbol: "eurusd".into(),
2512 side: "Buy".into(),
2513 order_type: "Market".into(),
2514 price: None,
2515 risk: 1.0,
2516 stoploss: None,
2517 targets: vec![],
2518 group: None,
2519 trade_id: None,
2520 }],
2521 profile: None,
2522 profile_def: None,
2523 config: BacktestConfigMsg {
2524 initial_balance: None,
2525 close_on_finish: None,
2526 fill_model: None,
2527 sizing: Some(SizingPolicyMsg::FixedLot { lots: 0.01 }),
2528 },
2529 };
2530 req.config.sizing = None;
2531 let state = test_state();
2532 let replay = build_replay_plan(
2533 &state,
2534 &req.symbol,
2535 &req.symbols,
2536 req.all_symbols,
2537 &req.raw_signals,
2538 None,
2539 None,
2540 0,
2541 )
2542 .unwrap();
2543
2544 assert!(validate_request(&req).is_ok());
2545 assert!(validate_replay_sizing(&req, &replay).is_err());
2546 }
2547
2548 #[test]
2549 fn ping_returns_ok() {
2550 let state = test_state();
2551 let resp = handle_ping(&state);
2552 assert_eq!(resp.status, "OK");
2553 assert_eq!(resp.data_dir, "/tmp/test");
2554 }
2555
2556 #[allow(dead_code)]
2557 fn temp_data_dir(name: &str) -> std::path::PathBuf {
2558 let unique = std::time::SystemTime::now()
2559 .duration_since(std::time::UNIX_EPOCH)
2560 .unwrap()
2561 .as_nanos();
2562 std::env::temp_dir().join(format!(
2563 "qs_backtest_server_{name}_{}_{}",
2564 std::process::id(),
2565 unique
2566 ))
2567 }
2568
2569 #[test]
2570 fn resolve_partition_value_matches_case_insensitive_symbol() {
2571 let root = temp_data_dir("partition_symbol");
2572 let symbol_dir = root
2573 .join("ticks")
2574 .join("exchange=icmarkets")
2575 .join("symbol=AUDCAD");
2576 std::fs::create_dir_all(&symbol_dir).unwrap();
2577 let root_str = root.to_string_lossy().to_string();
2578
2579 let resolved = resolve_partition_value(
2580 &root_str,
2581 "ticks",
2582 "symbol",
2583 "audcad",
2584 "exchange=icmarkets",
2585 None,
2586 )
2587 .unwrap();
2588
2589 assert_eq!(resolved, "AUDCAD");
2590 std::fs::remove_dir_all(root).unwrap();
2591 }
2592
2593 #[test]
2594 fn load_market_events_resolves_uppercase_tick_partition() {
2595 let root = temp_data_dir("tick_partition");
2596 let store = ParquetStore::open(&root).unwrap();
2597 let ts = parse_datetime("2026-01-15T10:00:00").unwrap();
2598 let ticks = vec![data_preprocess::Tick {
2599 exchange: "icmarkets".into(),
2600 symbol: "AUDCAD".into(),
2601 ts,
2602 bid: Some(0.9000),
2603 ask: Some(0.9002),
2604 last: None,
2605 volume: None,
2606 flags: None,
2607 }];
2608 store.insert_ticks(&ticks).unwrap();
2609 let root_str = root.to_string_lossy().to_string();
2610
2611 let events =
2612 load_market_events(&root_str, "icmarkets", "audcad", "tick", None, None, None).unwrap();
2613
2614 assert_eq!(events.len(), 1);
2615 match &events[0] {
2616 MarketEvent::Tick {
2617 symbol, bid, ask, ..
2618 } => {
2619 assert_eq!(symbol, "audcad");
2620 assert_eq!(*bid, 0.9000);
2621 assert_eq!(*ask, 0.9002);
2622 }
2623 MarketEvent::Bar { .. } => panic!("expected tick event"),
2624 }
2625 std::fs::remove_dir_all(root).unwrap();
2626 }
2627
2628 #[test]
2629 fn load_market_events_resolves_uppercase_bar_partition() {
2630 let root = temp_data_dir("bar_partition");
2631 let store = ParquetStore::open(&root).unwrap();
2632 let ts = parse_datetime("2026-01-15T10:00:00").unwrap();
2633 let bars = vec![data_preprocess::Bar {
2634 exchange: "icmarkets".into(),
2635 symbol: "AUDCAD".into(),
2636 timeframe: Timeframe::M1,
2637 ts,
2638 open: 0.9000,
2639 high: 0.9010,
2640 low: 0.8990,
2641 close: 0.9005,
2642 tick_vol: 10,
2643 volume: 0,
2644 spread: 2,
2645 }];
2646 store.insert_bars(&bars).unwrap();
2647 let root_str = root.to_string_lossy().to_string();
2648
2649 let events = load_market_events(
2650 &root_str,
2651 "icmarkets",
2652 "audcad",
2653 "bar",
2654 Some("1m"),
2655 None,
2656 None,
2657 )
2658 .unwrap();
2659
2660 assert_eq!(events.len(), 1);
2661 match &events[0] {
2662 MarketEvent::Bar { symbol, close, .. } => {
2663 assert_eq!(symbol, "audcad");
2664 assert_eq!(*close, 0.9005);
2665 }
2666 MarketEvent::Tick { .. } => panic!("expected bar event"),
2667 }
2668 std::fs::remove_dir_all(root).unwrap();
2669 }
2670
2671 #[test]
2672 fn list_profiles_empty_registry() {
2673 let state = test_state();
2674 let resp = handle_list_profiles(&state);
2675 assert!(resp.profiles.is_empty());
2676 }
2677
2678 #[test]
2679 fn add_profile_success() {
2680 let state = test_state();
2681 let req = AddProfileRequest {
2682 profile: ManagementProfileMsg {
2683 name: "new_prof".into(),
2684 target_selection: None,
2685 use_targets: vec![1],
2686 close_ratios: vec![1.0],
2687 stoploss_mode: None,
2688 rules: vec![],
2689 group_override: None,
2690 let_remainder_run: false,
2691 },
2692 overwrite: false,
2693 };
2694 let resp = handle_add_profile(&state, &req);
2695 assert!(resp.success);
2696 assert!(resp.error.is_none());
2697 assert_eq!(resp.profile_count, 1);
2698 }
2699
2700 #[test]
2701 fn add_profile_duplicate_rejected() {
2702 let state = test_state();
2703 let req = AddProfileRequest {
2704 profile: ManagementProfileMsg {
2705 name: "dup".into(),
2706 target_selection: None,
2707 use_targets: vec![1],
2708 close_ratios: vec![1.0],
2709 stoploss_mode: None,
2710 rules: vec![],
2711 group_override: None,
2712 let_remainder_run: false,
2713 },
2714 overwrite: false,
2715 };
2716 let resp1 = handle_add_profile(&state, &req);
2717 assert!(resp1.success);
2718 let resp2 = handle_add_profile(&state, &req);
2719 assert!(!resp2.success);
2720 assert!(resp2.error.as_ref().unwrap().contains("Duplicate"));
2721 }
2722
2723 #[test]
2724 fn add_profile_overwrite_success() {
2725 let state = test_state();
2726 let req1 = AddProfileRequest {
2727 profile: ManagementProfileMsg {
2728 name: "ow".into(),
2729 target_selection: None,
2730 use_targets: vec![1],
2731 close_ratios: vec![1.0],
2732 stoploss_mode: None,
2733 rules: vec![],
2734 group_override: None,
2735 let_remainder_run: false,
2736 },
2737 overwrite: false,
2738 };
2739 handle_add_profile(&state, &req1);
2740 let req2 = AddProfileRequest {
2741 profile: ManagementProfileMsg {
2742 name: "ow".into(),
2743 target_selection: None,
2744 use_targets: vec![1, 2],
2745 close_ratios: vec![0.5, 0.5],
2746 stoploss_mode: None,
2747 rules: vec![],
2748 group_override: None,
2749 let_remainder_run: false,
2750 },
2751 overwrite: true,
2752 };
2753 let resp = handle_add_profile(&state, &req2);
2754 assert!(resp.success);
2755 assert_eq!(resp.profile_count, 1);
2756 }
2757
2758 #[test]
2759 fn add_profile_invalid_rejected() {
2760 let state = test_state();
2761 let req = AddProfileRequest {
2762 profile: ManagementProfileMsg {
2763 name: "bad".into(),
2764 target_selection: None,
2765 use_targets: vec![1, 2],
2766 close_ratios: vec![1.0], stoploss_mode: None,
2768 rules: vec![],
2769 group_override: None,
2770 let_remainder_run: false,
2771 },
2772 overwrite: false,
2773 };
2774 let resp = handle_add_profile(&state, &req);
2775 assert!(!resp.success);
2776 assert!(resp.error.is_some());
2777 assert_eq!(resp.profile_count, 0);
2778 }
2779
2780 #[test]
2781 fn remove_profile_success() {
2782 let state = test_state();
2783 let add_req = AddProfileRequest {
2785 profile: ManagementProfileMsg {
2786 name: "rm_me".into(),
2787 target_selection: None,
2788 use_targets: vec![1],
2789 close_ratios: vec![1.0],
2790 stoploss_mode: None,
2791 rules: vec![],
2792 group_override: None,
2793 let_remainder_run: false,
2794 },
2795 overwrite: false,
2796 };
2797 handle_add_profile(&state, &add_req);
2798 let resp = handle_remove_profile(
2799 &state,
2800 &RemoveProfileRequest {
2801 name: "rm_me".into(),
2802 },
2803 );
2804 assert!(resp.success);
2805 assert!(resp.error.is_none());
2806 assert_eq!(resp.profile_count, 0);
2807 }
2808
2809 #[test]
2810 fn remove_profile_not_found() {
2811 let state = test_state();
2812 let resp = handle_remove_profile(
2813 &state,
2814 &RemoveProfileRequest {
2815 name: "nope".into(),
2816 },
2817 );
2818 assert!(!resp.success);
2819 assert!(resp.error.as_ref().unwrap().contains("not found"));
2820 }
2821
2822 #[test]
2827 fn filter_signals_by_date_inclusive() {
2828 use qs_backtest::profile::RawSignal;
2829 let t = |d: u32| {
2830 NaiveDate::from_ymd_opt(2026, 3, d)
2831 .unwrap()
2832 .and_hms_opt(0, 0, 0)
2833 .unwrap()
2834 };
2835 let signals: Vec<RawSignal> = vec![
2836 RawSignal::Entry {
2837 ts: t(8),
2838 symbol: "X".into(),
2839 side: Side::Buy,
2840 order_type: OrderType::Market,
2841 price: None,
2842 risk_multiplier: 0.01,
2843 stoploss: None,
2844 targets: vec![],
2845 group: None,
2846 trade_id: None,
2847 },
2848 RawSignal::Entry {
2849 ts: t(9),
2850 symbol: "X".into(),
2851 side: Side::Sell,
2852 order_type: OrderType::Market,
2853 price: None,
2854 risk_multiplier: 0.01,
2855 stoploss: None,
2856 targets: vec![],
2857 group: None,
2858 trade_id: None,
2859 },
2860 RawSignal::Entry {
2861 ts: t(12),
2862 symbol: "X".into(),
2863 side: Side::Buy,
2864 order_type: OrderType::Market,
2865 price: None,
2866 risk_multiplier: 0.01,
2867 stoploss: None,
2868 targets: vec![],
2869 group: None,
2870 trade_id: None,
2871 },
2872 ];
2873 let replay = ReplayPlan::build(
2874 RequestedSymbolScope::explicit(["X".into()]),
2875 signals,
2876 Some(t(8)),
2877 Some(t(11)),
2878 0,
2879 )
2880 .unwrap();
2881 assert_eq!(replay.retained_signals().len(), 2);
2882 }
2883
2884 #[test]
2885 fn filter_signals_by_date_no_filter() {
2886 use qs_backtest::profile::RawSignal;
2887 let signals: Vec<RawSignal> = vec![RawSignal::Entry {
2888 ts: NaiveDate::from_ymd_opt(2026, 1, 1)
2889 .unwrap()
2890 .and_hms_opt(0, 0, 0)
2891 .unwrap(),
2892 symbol: "X".into(),
2893 side: Side::Buy,
2894 order_type: OrderType::Market,
2895 price: None,
2896 risk_multiplier: 0.01,
2897 stoploss: None,
2898 targets: vec![],
2899 group: None,
2900 trade_id: None,
2901 }];
2902 let replay = ReplayPlan::build(
2903 RequestedSymbolScope::explicit(["X".into()]),
2904 signals,
2905 None,
2906 None,
2907 0,
2908 )
2909 .unwrap();
2910 assert_eq!(replay.retained_signals().len(), 1);
2911 }
2912
2913 #[allow(dead_code)]
2918 fn job_test_state() -> ServerState {
2919 test_state()
2920 }
2921
2922 #[allow(dead_code)]
2923 fn valid_submit_request() -> BacktestRunSpec {
2924 BacktestRunSpec {
2925 symbol: "XAUUSD".into(),
2926 symbols: vec![],
2927 all_symbols: false,
2928 exchange: "icmarkets".into(),
2929 data_type: "tick".into(),
2930 timeframe: None,
2931 from: None,
2932 to: None,
2933 raw_signals: vec![RawSignalMsg::Entry {
2934 ts: "2026-01-01T00:00:00".into(),
2935 symbol: "xauusd".into(),
2936 side: "Buy".into(),
2937 order_type: "Market".into(),
2938 price: Some(5000.0),
2939 risk: 1.0,
2940 stoploss: Some(4990.0),
2941 targets: vec![],
2942 group: None,
2943 trade_id: None,
2944 }],
2945 profile: None,
2946 profile_def: None,
2947 config: BacktestConfigMsg {
2948 initial_balance: Some(10_000.0),
2949 close_on_finish: Some(true),
2950 fill_model: Some("BidAsk".into()),
2951 sizing: Some(SizingPolicyMsg::FixedLot { lots: 0.01 }),
2952 },
2953 }
2954 }
2955
2956 #[test]
2957 fn submit_and_cancel_job() {
2958 let state = job_test_state();
2959 let submit = submit_for_test(&state, &valid_submit_request());
2960 assert!(submit.success);
2961 let job_id = submit.job_id.unwrap();
2962
2963 let cancel = handle_cancel_backtest(
2964 &state,
2965 &CancelBacktestRequest {
2966 job_id: job_id.clone(),
2967 },
2968 );
2969 assert!(cancel.success);
2970
2971 let status = handle_get_backtest_status(&state, &GetBacktestStatusRequest { job_id });
2972 assert_eq!(status.status, "Cancelled");
2973 }
2974
2975 #[test]
2976 fn submit_invalid_request_rejected() {
2977 let state = job_test_state();
2978 let mut req = valid_submit_request();
2979 req.raw_signals = vec![];
2980 let submit = submit_for_test(&state, &req);
2981 assert!(!submit.success);
2982 assert!(submit.job_id.is_none());
2983 }
2984
2985 #[test]
2986 fn get_status_not_found() {
2987 let state = job_test_state();
2988 let status = handle_get_backtest_status(
2989 &state,
2990 &GetBacktestStatusRequest {
2991 job_id: "nonexistent".into(),
2992 },
2993 );
2994 assert!(!status.success);
2995 assert_eq!(status.status, "NotFound");
2996 }
2997
2998 #[tokio::test]
2999 async fn watch_stream_emits_initial_progress_terminal_and_end() {
3000 let state = Arc::new(job_test_state());
3001 let submit = submit_for_test(&state, &valid_submit_request());
3002 let job_id = submit.job_id.unwrap();
3003 let mut stream = watch_backtest_stream(
3004 state.clone(),
3005 WatchBacktestRequest {
3006 job_id: job_id.clone(),
3007 },
3008 Duration::from_secs(60),
3009 );
3010
3011 let initial = stream.next().await.unwrap().unwrap();
3012 assert!(matches!(
3013 initial,
3014 BacktestEvent::Snapshot { ref status }
3015 if status.job_id == job_id && status.status == "Queued"
3016 ));
3017
3018 update_job_progress(
3019 &state,
3020 &job_id,
3021 BacktestProgress {
3022 stage: "replay".into(),
3023 processed_events: 10,
3024 total_events: 100,
3025 processed_signals: 2,
3026 total_signals: 8,
3027 processed_symbols: 1,
3028 total_symbols: 2,
3029 },
3030 );
3031 let progress = stream.next().await.unwrap().unwrap();
3032 assert!(matches!(
3033 progress,
3034 BacktestEvent::Snapshot { ref status }
3035 if status.status == "Running"
3036 && status.progress.processed_events == 10
3037 && status.progress.total_events == 100
3038 ));
3039
3040 assert!(
3041 handle_cancel_backtest(
3042 &state,
3043 &CancelBacktestRequest {
3044 job_id: job_id.clone(),
3045 },
3046 )
3047 .success
3048 );
3049 let terminal = stream.next().await.unwrap().unwrap();
3050 assert!(matches!(
3051 terminal,
3052 BacktestEvent::Snapshot { ref status }
3053 if status.status == "Cancelled" && status.is_terminal()
3054 ));
3055 assert!(stream.next().await.is_none());
3056 }
3057
3058 #[tokio::test]
3059 async fn watch_stream_heartbeats_and_resubscribes_to_terminal_snapshot() {
3060 let state = Arc::new(job_test_state());
3061 let submit = submit_for_test(&state, &valid_submit_request());
3062 let job_id = submit.job_id.unwrap();
3063 let mut stream = watch_backtest_stream(
3064 state.clone(),
3065 WatchBacktestRequest {
3066 job_id: job_id.clone(),
3067 },
3068 Duration::from_millis(5),
3069 );
3070
3071 assert!(matches!(
3072 stream.next().await.unwrap().unwrap(),
3073 BacktestEvent::Snapshot { .. }
3074 ));
3075 let heartbeat = tokio::time::timeout(Duration::from_millis(100), stream.next())
3076 .await
3077 .unwrap()
3078 .unwrap()
3079 .unwrap();
3080 assert!(matches!(
3081 heartbeat,
3082 BacktestEvent::Heartbeat { job_id: ref id, .. } if id == &job_id
3083 ));
3084
3085 assert!(
3086 handle_cancel_backtest(
3087 &state,
3088 &CancelBacktestRequest {
3089 job_id: job_id.clone(),
3090 },
3091 )
3092 .success
3093 );
3094 drop(stream);
3095
3096 let mut resumed = watch_backtest_stream(
3097 state,
3098 WatchBacktestRequest {
3099 job_id: job_id.clone(),
3100 },
3101 Duration::from_secs(60),
3102 );
3103 assert!(matches!(
3104 resumed.next().await.unwrap().unwrap(),
3105 BacktestEvent::Snapshot { ref status }
3106 if status.job_id == job_id && status.status == "Cancelled"
3107 ));
3108 assert!(resumed.next().await.is_none());
3109 }
3110
3111 #[tokio::test]
3112 async fn watch_stream_reports_missing_job_as_stream_error() {
3113 let mut stream = watch_backtest_stream(
3114 Arc::new(job_test_state()),
3115 WatchBacktestRequest {
3116 job_id: "missing".into(),
3117 },
3118 Duration::from_secs(60),
3119 );
3120 assert!(matches!(
3121 stream.next().await,
3122 Some(Err(xrpc::RpcError::ServerError(error))) if error == "Job 'missing' not found"
3123 ));
3124 assert!(stream.next().await.is_none());
3125 }
3126
3127 #[tokio::test]
3128 async fn shutdown_cancellation_publishes_terminal_snapshot() {
3129 let state = Arc::new(job_test_state());
3130 let submit = submit_for_test(&state, &valid_submit_request());
3131 let job_id = submit.job_id.unwrap();
3132 let (mut updates, _) = subscribe_backtest_status(&state, &job_id).unwrap();
3133
3134 assert_eq!(cancel_active_jobs(&state), 1);
3135 updates.changed().await.unwrap();
3136 let status = updates.borrow().clone();
3137 assert_eq!(status.status, "Cancelled");
3138 assert!(status.is_terminal());
3139 }
3140
3141 #[test]
3142 fn cancel_nonexistent_job() {
3143 let state = job_test_state();
3144 let resp = handle_cancel_backtest(
3145 &state,
3146 &CancelBacktestRequest {
3147 job_id: "nope".into(),
3148 },
3149 );
3150 assert!(!resp.success);
3151 }
3152
3153 #[test]
3158 fn config_from_msg_populates_contract_sizes() {
3159 let toml = r#"
3160[[symbol]]
3161canonical = "xauusd"
3162aliases = ["gold"]
3163pip_position = 1
3164digits = 2
3165category = "metal"
3166base_currency = "XAU"
3167quote_currency = "USD"
3168pnl_currency = "USD"
3169lot_base_units = 100
3170lot_step_units = 1
3171lot_min_steps = 1
3172lot_max_steps = 0
3173
3174[[symbol]]
3175canonical = "gbpjpy"
3176aliases = []
3177pip_position = 2
3178digits = 3
3179category = "forex"
3180base_currency = "GBP"
3181quote_currency = "JPY"
3182pnl_currency = "JPY"
3183lot_base_units = 100000
3184lot_step_units = 1000
3185lot_min_steps = 1
3186lot_max_steps = 0
3187"#;
3188 let registry = SymbolRegistry::from_toml(toml).unwrap();
3189 let symbols = vec!["xauusd".to_string(), "gbpjpy".to_string()];
3190 let msg = BacktestConfigMsg {
3191 initial_balance: Some(10_000.0),
3192 close_on_finish: Some(true),
3193 fill_model: Some("BidAsk".into()),
3194 sizing: None,
3195 };
3196 let config = config_from_msg(&msg, ®istry, &symbols).unwrap();
3197 assert_eq!(config.contract_sizes.get("xauusd"), Some(&100.0));
3198 assert_eq!(config.contract_sizes.get("gbpjpy"), Some(&100_000.0));
3199 assert!(!config.symbol_specs.is_empty());
3200 }
3201
3202 #[tokio::test]
3203 async fn job_full_lifecycle_transitions_publish_terminal_snapshot() {
3204 use std::sync::Arc;
3205 let state = Arc::new(job_test_state());
3206
3207 let submit = submit_for_test(&state, &valid_submit_request());
3209 assert!(submit.success);
3210 let job_id = submit.job_id.unwrap();
3211
3212 let status = handle_get_backtest_status(
3214 &state,
3215 &GetBacktestStatusRequest {
3216 job_id: job_id.clone(),
3217 },
3218 );
3219 assert_eq!(status.status, "Queued");
3220 let mut stream = watch_backtest_stream(
3221 state.clone(),
3222 WatchBacktestRequest {
3223 job_id: job_id.clone(),
3224 },
3225 Duration::from_secs(60),
3226 );
3227 assert!(matches!(
3228 stream.next().await.unwrap().unwrap(),
3229 BacktestEvent::Snapshot { ref status } if status.status == "Queued"
3230 ));
3231
3232 let req = valid_submit_request();
3235 run_job_for_test(state.clone(), job_id.clone(), req);
3236
3237 let status = handle_get_backtest_status(
3239 &state,
3240 &GetBacktestStatusRequest {
3241 job_id: job_id.clone(),
3242 },
3243 );
3244 assert!(
3245 status.status == "Failed" || status.status == "Completed",
3246 "Expected Failed or Completed, got {}",
3247 status.status
3248 );
3249 let terminal = stream.next().await.unwrap().unwrap();
3250 assert!(matches!(
3251 terminal,
3252 BacktestEvent::Snapshot { ref status }
3253 if status.status == "Failed" || status.status == "Completed"
3254 ));
3255 assert!(stream.next().await.is_none());
3256
3257 let result_resp = handle_get_backtest_result(&state, &GetBacktestResultRequest { job_id });
3259 if status.status == "Failed" {
3260 assert!(!result_resp.success);
3261 }
3262 }
3263
3264 #[test]
3265 fn concurrent_jobs_independent() {
3266 let state = job_test_state();
3267
3268 let submit1 = submit_for_test(&state, &valid_submit_request());
3270 let submit2 = submit_for_test(&state, &valid_submit_request());
3271
3272 assert!(submit1.success);
3273 assert!(submit2.success);
3274
3275 let id1 = submit1.job_id.unwrap();
3276 let id2 = submit2.job_id.unwrap();
3277
3278 assert_ne!(id1, id2);
3280
3281 let s1 = handle_get_backtest_status(
3283 &state,
3284 &GetBacktestStatusRequest {
3285 job_id: id1.clone(),
3286 },
3287 );
3288 let s2 = handle_get_backtest_status(
3289 &state,
3290 &GetBacktestStatusRequest {
3291 job_id: id2.clone(),
3292 },
3293 );
3294 assert_eq!(s1.status, "Queued");
3295 assert_eq!(s2.status, "Queued");
3296
3297 handle_cancel_backtest(
3299 &state,
3300 &CancelBacktestRequest {
3301 job_id: id1.clone(),
3302 },
3303 );
3304
3305 let s1_after =
3306 handle_get_backtest_status(&state, &GetBacktestStatusRequest { job_id: id1 });
3307 let s2_after =
3308 handle_get_backtest_status(&state, &GetBacktestStatusRequest { job_id: id2 });
3309 assert_eq!(s1_after.status, "Cancelled");
3310 assert_eq!(s2_after.status, "Queued");
3311 }
3312
3313 #[test]
3314 fn job_cleanup_removes_expired() {
3315 let state = job_test_state();
3316
3317 let submit = submit_for_test(&state, &valid_submit_request());
3319 let job_id = submit.job_id.unwrap();
3320 handle_cancel_backtest(
3321 &state,
3322 &CancelBacktestRequest {
3323 job_id: job_id.clone(),
3324 },
3325 );
3326
3327 {
3329 let jobs = state.jobs.lock().unwrap();
3330 assert!(jobs.contains_key(&job_id));
3331 }
3332
3333 state
3335 .jobs
3336 .lock()
3337 .unwrap()
3338 .get_mut(&job_id)
3339 .unwrap()
3340 .worker_active = false;
3341
3342 assert_eq!(cleanup_expired_jobs(&state, Duration::ZERO), 1);
3344
3345 {
3347 let jobs = state.jobs.lock().unwrap();
3348 assert!(!jobs.contains_key(&job_id));
3349 }
3350 }
3351
3352 #[test]
3353 fn job_status_enum_roundtrip() {
3354 assert_eq!(JobStatus::Queued.as_str(), "Queued");
3355 assert_eq!(JobStatus::LoadingData.as_str(), "LoadingData");
3356 assert_eq!(JobStatus::Running.as_str(), "Running");
3357 assert_eq!(JobStatus::Completed.as_str(), "Completed");
3358 assert_eq!(JobStatus::Failed.as_str(), "Failed");
3359 assert_eq!(JobStatus::Cancelled.as_str(), "Cancelled");
3360
3361 assert_eq!("Queued".parse::<JobStatus>(), Ok(JobStatus::Queued));
3362 assert_eq!(
3363 "LoadingData".parse::<JobStatus>(),
3364 Ok(JobStatus::LoadingData)
3365 );
3366 assert_eq!("invalid".parse::<JobStatus>(), Err("invalid job status"));
3367 }
3368}