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