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