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