Skip to main content

backtest_server/
handlers.rs

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