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