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