Skip to main content

backtest_server/
handlers.rs

1//! RPC handler implementations for the backtest server.
2//!
3//! Each handler receives a request message, processes it against the shared
4//! server state (symbol registry, profile registry, Parquet store), and
5//! returns a response message. Errors are captured in the response rather
6//! than crashing the server.
7
8mod 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
61/// Shared state accessible by all client handlers.
62pub struct ServerState {
63    /// Symbol registry for compatibility normalization and current currency metadata.
64    pub symbol_registry: SymbolRegistry,
65    /// Immutable instrument catalog and stored-series identity policy.
66    pub instrument_domain: InstrumentDomain,
67    /// Management profile registry (TOML-loaded + dynamically added).
68    pub profile_registry: RwLock<ProfileRegistry>,
69    /// Root directory for Parquet market data.
70    pub data_dir: String,
71    /// Path to the profiles TOML file (for reload).
72    pub profiles_path: String,
73    /// Server start time for uptime reporting.
74    pub start_time: Instant,
75    /// Async backtest job storage (Issue 2).
76    pub jobs: Mutex<HashMap<String, BacktestJob>>,
77    /// Hard bound for queued, running, and retained terminal jobs.
78    pub max_retained_jobs: usize,
79    /// Filesystem storage for complete large result JSON payloads.
80    pub artifact_store: ArtifactStore,
81    /// Limits for configured strategy runs and the gate that runs one search at a time.
82    pub strategies: StrategyServiceState,
83}
84
85/// Configured strategy limits plus the gate that lets only one parameter search hold its market data in memory at a time.
86#[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/// Validated input a retained job runs once its worker starts.
109#[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/// What a retained job produces, which decides the endpoint that returns its output.
118#[derive(Debug, Clone, Copy, PartialEq, Eq)]
119pub enum JobKind {
120    /// A backtest result, returned by `get_backtest_result`.
121    Backtest,
122    /// A search output, returned by `get_search_result`.
123    Search,
124}
125
126/// Internal representation of an async backtest job.
127#[derive(Debug, Clone)]
128pub struct BacktestJob {
129    /// Current job status.
130    pub status: JobStatus,
131    /// When the job was submitted.
132    pub submitted_at: Instant,
133    /// When the job completed (if finished).
134    pub completed_at: Option<Instant>,
135    /// Structured loading and replay progress.
136    pub progress: BacktestProgress,
137    /// Complete inline result or compact console summary.
138    pub result: Option<BacktestResultMsg>,
139    /// Complete result artifact when the full inline object was released.
140    pub artifact: Option<ResultArtifactRefMsg>,
141    /// Separately labeled resumable checkpoint for a cancelled search.
142    pub checkpoint_artifact: Option<ResultArtifactRefMsg>,
143    /// True when the complete result is present in `result`.
144    pub inline_complete: bool,
145    /// True after the job artifact has been deleted following delivery.
146    pub artifact_consumed: bool,
147    /// Error message (when failed).
148    pub error: Option<String>,
149    /// Cooperative cancellation shared with the blocking worker.
150    pub cancellation: JobCancellationToken,
151    /// True while a blocking worker slot is reserved or still accessing this job.
152    pub worker_active: bool,
153    /// Coalesced current status published to server-streaming subscribers.
154    pub updates: watch::Sender<BacktestStatusResponse>,
155    /// Accepted request and immutable profile snapshot until the worker starts.
156    pub accepted: Option<AcceptedJobInput>,
157    /// Which result endpoint serves this job.
158    pub kind: JobKind,
159    /// Summary of a completed search; its complete output is the job artifact.
160    pub search: Option<SearchSummaryMsg>,
161}
162
163/// Lightweight per-job cancellation token without an additional runtime dependency.
164#[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/// Typed job status for the async backtest API.
178#[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    /// Convert to string for wire transport.
343    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
377/// Price point size per active symbol, used to express stored bar spreads in price units.
378fn 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
392/// Subscribe to the current and future coalesced snapshots of a retained job.
393pub 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
413/// Build the server stream for a retained backtest job.
414pub 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(&current.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
529/// Remove terminal jobs older than `retention` and enforce the configured bound.
530pub 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
560/// Cooperatively cancel every active job during server shutdown.
561pub 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
611// ── Ping ────────────────────────────────────────────────────────────────────
612
613/// Handle `ping` — returns server status and uptime.
614pub 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
622// ── List Profiles ───────────────────────────────────────────────────────────
623
624/// Handle `list_profiles` — returns all loaded management profiles.
625pub 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
646// ── List Symbols ────────────────────────────────────────────────────────────
647
648/// Handle `list_symbols` — returns available data from the Parquet store.
649pub 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            // Apply data_type filter if requested.
666            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            // The store reports bars as `bar (1m)`; tolerate the historical
679            // no-space spelling too so existing data remains discoverable.
680            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
712// ── Run Backtest ────────────────────────────────────────────────────────────
713
714/// Handle the canonical FutureQuoteV1 execution endpoint.
715pub 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
1052// ── Run Backtest Multi ──────────────────────────────────────────────────────
1053
1054pub 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    // Early validation common to all profiles.
1097    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(&registry, &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(&registry, 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
1269/// Extract the profile name from a `ProfileRef`.
1270fn 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// ── Signal Conversion ───────────────────────────────────────────────────────
1333
1334#[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
1427/// Resolve a run default profile and entry-class routes from the request fields every run shape shares.
1428fn 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(&registry, &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// ── Data Loading ────────────────────────────────────────────────────────────
1465
1466/// Load raw market events from Parquet (shared between single and multi runs).
1467#[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/// Resolve a Hive partition value by scanning child directories case-insensitively.
1604#[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/// Rewrite loaded market events to the canonical symbol used by signals.
1644#[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/// Drain a VecFeed into its underlying Vec<MarketEvent>.
1660#[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
1669// ── Validation ──────────────────────────────────────────────────────────────
1670
1671/// Validate the common fields of a run_backtest request.
1672fn 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
1706// ── Parsing Helpers ─────────────────────────────────────────────────────────
1707
1708/// Normalize a symbol name via the registry, falling back to passthrough.
1709fn normalize_symbol(registry: &SymbolRegistry, raw: &str) -> String {
1710    registry.normalize_or_passthrough(raw)
1711}
1712
1713/// Parse an ISO datetime string into NaiveDateTime.
1714pub(crate) fn parse_datetime(s: &str) -> Result<NaiveDateTime> {
1715    // Try multiple common formats.
1716    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    // Try date-only (appends midnight).
1729    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
1737/// Parse an optional datetime string.
1738fn 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
1745// ── Phase 2: Profile Management Handlers ────────────────────────────────────
1746
1747/// Handle `add_profile` — add or overwrite a management profile at runtime.
1748pub 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
1776/// Handle `remove_profile` — remove a management profile by name.
1777pub 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
1794/// Handle `reload_profiles` — reload profiles from the configured TOML file.
1795pub 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
1826// ── Async Job API Handlers (Issue 2) ─────────────────────────────────────────
1827
1828/// Admit a validated request to the bounded async job store.
1829fn 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
1950/// Handle `get_backtest_status` — poll the status of a submitted job.
1951pub 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
1969/// Handle `get_backtest_result` — fetch the result of a completed job.
1970pub 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
2029/// Return one base64-encoded raw result artifact chunk.
2030pub 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
2057/// Delete a complete result artifact.
2058pub 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
2103/// Handle `cancel_backtest` — cancel a submitted job.
2104pub 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
2203/// Run one retained job that produces a backtest result, then store it inline or as an artifact exactly as every result job does.
2204fn 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
2268/// Move a queued job to loading and return its cancellation token, or `None` when it was cancelled before its worker started.
2269fn 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
2307/// Generate a simple unique ID without external dependencies.
2308fn 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(&reg, "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], // mismatch
2948                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        // Add a profile first.
2967        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    //
3008    // Issue 1: filter_signals_by_date tests
3009    //
3010
3011    #[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    //
3103    // Issue 2: Async job tests
3104    //
3105
3106    #[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    //
3365    // Issue 3: Contract size wiring test
3366    //
3367
3368    #[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, &registry, &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        // Submit a valid job.
3420        let submit = submit_for_test(&state, &valid_submit_request());
3421        assert!(submit.success);
3422        let job_id = submit.job_id.unwrap();
3423
3424        // Status should be Queued.
3425        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        // Run the job via the canonical worker (will fail since no real data,
3445        // but should transition through LoadingData -> Failed).
3446        let req = valid_submit_request();
3447        run_job_for_test(state.clone(), job_id.clone(), req);
3448
3449        // Status should be Failed (no market data at /tmp/test).
3450        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        // Fetch result should fail since job is not Completed.
3470        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        // Submit two jobs.
3481        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        // IDs must be different.
3491        assert_ne!(id1, id2);
3492
3493        // Both should be Queued.
3494        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        // Cancel one, the other should still be Queued.
3510        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        // Submit and cancel a job.
3530        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        // Job should exist.
3540        {
3541            let jobs = state.jobs.lock().unwrap();
3542            assert!(jobs.contains_key(&job_id));
3543        }
3544
3545        // Simulate the spawned worker acknowledging the pre-start cancellation.
3546        state
3547            .jobs
3548            .lock()
3549            .unwrap()
3550            .get_mut(&job_id)
3551            .unwrap()
3552            .worker_active = false;
3553
3554        // Cleanup with max_age=0 removes completed/cancelled jobs.
3555        assert_eq!(cleanup_expired_jobs(&state, Duration::ZERO), 1);
3556
3557        // Job should be gone.
3558        {
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}