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