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