perf(eval): split metrics-readback into launch/consume — restore async pipelining

The async-diag merge (3c0d26292, building on d9cb14f1b + 673b04a8d) wired the
GPU side of cross-stream eval pipelining correctly (training stream `cuda_stream`,
eval stream `validation_stream`, `cuStreamWaitEvent` barrier on `train_done_event`)
but left a host-side `stream.synchronize()` inside
`gpu_backtest_evaluator::launch_metrics_and_download` that blocked the CPU thread
for the FULL eval drain (~25-30s/epoch on L40S). The host thread is the same one
that submits the next epoch's training kernels via `run_full_step` — so until eval
drained, training submission was gated on it, defeating the dedicated
`validation_stream`.

Fix: split metrics readback into record/consume halves at the buffer-readback layer:

* `launch_metrics_and_record_event` — submits the metrics kernel + DtoD-async
  copies of `actions_history_buf` / `intent_mag_buf` / `picked_action_history_buf`
  into mapped-pinned mirrors, then records `eval_done_event`. Returns immediately.
* `consume_metrics_after_event` — `event.synchronize()` (the SOLE host wait per
  epoch boundary), then `read_volatile`s mapped-pinned `metrics_buf`. After the
  event syncs, the four action-distribution helpers read directly from the
  now-coherent mapped-pinned mirrors — eliminating four synchronous DtoH copies
  that violated `feedback_no_htod_htoh_only_mapped_pinned.md`.

Caller migration (per `feedback_no_partial_refactor.md`):
* `evaluate_dqn_graphed` (existing public API) becomes a thin sync wrapper:
  launch_async → consume. ABI unchanged.
* New `evaluate_dqn_graphed_async` for the pipelined path.
* `evaluate` / `evaluate_ppo` / `evaluate_supervised` migrated inline to
  launch+consume (no caller pipelines them).
* `compute_validation_loss` replaced by `launch_validation_loss` +
  `consume_validation_loss`. Training loop now consumes pending at epoch start
  (cached_async_val_loss = epoch N-1's val_sharpe) and launches at epoch end
  (sentinel-only return). Final post-loop consume drains the LAST epoch's launch
  so smoke tests / hyperopt that read `last_eval_direction_dist` see the most
  recent epoch's data.

Mathematical identity preserved: every val_sharpe consumed is bit-identical to
what the prior synchronous flow would have produced for the same epoch — only
the timing of the host parse differs (one epoch later, matching the existing
`cached_async_val_loss` lag semantics; HEALTH_DIAG `val [...]` / `val_dir_dist` /
`val_picked_dir_dist` emit moves with the consume).

Audit: docs/dqn-gpu-hot-path-audit.md Fix 4. Out of scope: periodic chunk-level
`stream.synchronize` calls in `submit_dqn_step_loop_cublas` (lines ~1596 + ~1863)
for kernel-error detection — host-blocking but smaller in aggregate; future work.

Workspace baseline preserved: 13 warnings.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
jgrusewski
2026-04-28 20:35:41 +02:00
parent 3c0d262927
commit 660f02ff40
4 changed files with 445 additions and 96 deletions

View File

@@ -31,13 +31,13 @@
use std::mem::ManuallyDrop; use std::mem::ManuallyDrop;
use std::sync::Arc; use std::sync::Arc;
use cudarc::driver::{CudaFunction, CudaGraph, CudaSlice, CudaStream, DevicePtr, LaunchConfig, PushKernelArg}; use cudarc::driver::{CudaEvent, CudaFunction, CudaGraph, CudaSlice, CudaStream, DevicePtr, LaunchConfig, PushKernelArg};
use tracing::info; use tracing::info;
use ml_core::nvtx::NvtxRange; use ml_core::nvtx::NvtxRange;
use crate::MLError; use crate::MLError;
use super::gpu_weights::{BranchingWeightSet, DuelingWeightSet, PpoActorWeightSet}; use super::gpu_weights::{BranchingWeightSet, DuelingWeightSet, PpoActorWeightSet};
use super::mapped_pinned::MappedF32Buffer; use super::mapped_pinned::{MappedF32Buffer, MappedI32Buffer};
// ── CUDA Graph wrapper ─────────────────────────────────────────────────────── // ── CUDA Graph wrapper ───────────────────────────────────────────────────────
// //
@@ -355,11 +355,51 @@ pub struct GpuBacktestEvaluator {
// visible to the GPU through `dev_ptr` (`cuMemHostAlloc DEVICEMAP`). The // visible to the GPU through `dev_ptr` (`cuMemHostAlloc DEVICEMAP`). The
// kernel writes through the device alias; the host reads via volatile // kernel writes through the device alias; the host reads via volatile
// load on `host_ptr` after stream synchronisation. No DtoH copy — the // load on `host_ptr` after stream synchronisation. No DtoH copy — the
// only allowed CPU↔GPU path per // only allowed CPU↔CPU path per
// `feedback_no_htod_htoh_only_mapped_pinned.md`. Layout `[n_windows * 14]` // `feedback_no_htod_htoh_only_mapped_pinned.md`. Layout `[n_windows * 14]`
// matches the kernel's flat per-window stride. // matches the kernel's flat per-window stride.
metrics_buf: MappedF32Buffer, metrics_buf: MappedF32Buffer,
// ── Async eval pipeline (perf-fix 2026-04-28) ─────────────────────────────
//
// Splits eval-stream readback from host-side wait. `launch_metrics_and_record_event`
// submits the metrics kernel + DtoD-async copies of the per-step action history
// buffers into the mapped-pinned mirrors below, then records `eval_done_event`.
// `consume_metrics_after_event` calls `event.synchronize()` (the ONLY host wait
// per epoch boundary) and reads through `host_ptr`. The split lets the host
// CPU thread launch the next epoch's training kernels while eval drains —
// restoring the pipelining the dedicated `validation_stream` is supposed to
// enable. Original async-diag commits (d9cb14f1b + 673b04a8d) wired the GPU-
// side cross-stream event correctly but left the host blocked inside
// `synchronize()` after the metrics launch, so pipelining never engaged.
/// Recorded on `self.stream` after the metrics kernel + action-history DtoD
/// copies are queued. Lazy-allocated on first launch (the evaluator is
/// constructed before the CUDA context is available for `new_event` here in
/// some early-init paths, so we defer creation). Reused across epochs —
/// `record` re-records into the same event handle.
eval_done_event: Option<CudaEvent>,
/// Mapped-pinned mirror of `actions_history_buf` ([n_windows * max_len] i32).
/// Filled via `memcpy_dtod_async` on `self.stream` from `actions_history_buf`
/// inside `launch_metrics_and_record_event`. Read by the per-direction and
/// per-magnitude distribution helpers after `consume_metrics_after_event`
/// has synchronised the event. Replaces the synchronous `memcpy_dtoh` into
/// a freshly-allocated `Vec<i32>` per call — that path violates
/// `feedback_no_htod_htoh_only_mapped_pinned.md` and forced a host wait
/// inside the eval flow.
actions_history_pinned: MappedI32Buffer,
/// Mapped-pinned mirror of `intent_mag_buf` ([n_windows * max_len] i32).
/// Lazy-allocated alongside `intent_mag_buf` in `ensure_action_select_ready`.
/// Same DtoD-async-on-eval-stream pattern as `actions_history_pinned`.
intent_mag_pinned: Option<MappedI32Buffer>,
/// Mapped-pinned mirror of `picked_action_history_buf` ([n_windows * max_len] i32).
/// Lazy-allocated alongside `picked_action_history_buf` in
/// `ensure_action_select_ready`. Same pattern.
picked_action_history_pinned: Option<MappedI32Buffer>,
// Dimensions and config // Dimensions and config
n_windows: usize, n_windows: usize,
max_len: usize, max_len: usize,
@@ -650,6 +690,17 @@ impl GpuBacktestEvaluator {
MappedF32Buffer::new(n_windows * 14) MappedF32Buffer::new(n_windows * 14)
.map_err(|e| MLError::ModelError(format!("metrics mapped-pinned alloc: {e}")))? .map_err(|e| MLError::ModelError(format!("metrics mapped-pinned alloc: {e}")))?
}; };
// Mapped-pinned mirror for actions_history_buf — populated via
// `memcpy_dtod_async` on the eval stream before the eval_done_event is
// recorded. Replaces the per-call synchronous `memcpy_dtoh` into a
// freshly-allocated host `Vec<i32>`. See the field docstring above for
// the perf-fix rationale.
let actions_history_pinned = unsafe {
MappedI32Buffer::new(n_windows * max_len)
.map_err(|e| MLError::ModelError(format!(
"actions_history mapped-pinned alloc: {e}"
)))?
};
// State layout defined by state_layout.cuh — single source of truth. // State layout defined by state_layout.cuh — single source of truth.
// Canonical 96-dim: 42 market + 20 OFI + 16 MTF + 8 portfolio + 6 plan/ISV + 4 padding. // Canonical 96-dim: 42 market + 20 OFI + 16 MTF + 8 portfolio + 6 plan/ISV + 4 padding.
@@ -697,6 +748,16 @@ impl GpuBacktestEvaluator {
picked_action_history_buf: None, picked_action_history_buf: None,
states_buf, states_buf,
metrics_buf, metrics_buf,
// Lazy-init on first launch so the CUDA context is guaranteed
// active. Reused across epochs (re-record via `event.record(stream)`).
eval_done_event: None,
actions_history_pinned,
// Lazy-allocated alongside the underlying device buffers in
// `ensure_action_select_ready`. Stay `None` for non-DQN paths
// (evaluate / evaluate_ppo / evaluate_supervised) which never
// populate the underlying buffers.
intent_mag_pinned: None,
picked_action_history_pinned: None,
n_windows, n_windows,
max_len, max_len,
feature_dim, feature_dim,
@@ -971,7 +1032,11 @@ impl GpuBacktestEvaluator {
// is safe and avoids the latency of periodic done-flag downloads. // is safe and avoids the latency of periodic done-flag downloads.
} }
self.launch_metrics_and_download() // Synchronous launch+consume — `evaluate()` is the closure-based path
// used by callers that don't pipeline (no caller currently does — the
// pipelined DQN path goes through `evaluate_dqn_graphed_async`).
self.launch_metrics_and_record_event()?;
self.consume_metrics_after_event()
} }
/// Pure-GPU DQN evaluation via cuBLAS SGEMM — no candle, no closures. /// Pure-GPU DQN evaluation via cuBLAS SGEMM — no candle, no closures.
@@ -1018,13 +1083,53 @@ impl GpuBacktestEvaluator {
/// If you change weights (different `DuelingWeightSet`), call /// If you change weights (different `DuelingWeightSet`), call
/// [`invalidate_dqn_graph`] first or just use `evaluate_dqn` instead. /// [`invalidate_dqn_graph`] first or just use `evaluate_dqn` instead.
pub fn evaluate_dqn_graphed( pub fn evaluate_dqn_graphed(
&mut self,
online_weights: &DuelingWeightSet,
branching_weights: Option<&BranchingWeightSet>,
dqn_cfg: &DqnBacktestConfig,
q_provider: Option<&mut dyn super::q_value_provider::QValueProvider>,
) -> Result<Vec<WindowMetrics>, MLError> {
// Synchronous wrapper: launch async + immediate consume. Used by
// hyperopt (single-shot) and the smoke-test path that expects metrics
// back from this call without a separate consume step. The pipelined
// training path uses `evaluate_dqn_graphed_async` + a deferred
// `consume_metrics_after_event` from the next epoch's start (see
// `compute_validation_loss_launch` / `_consume` in trainer/metrics.rs).
self.evaluate_dqn_graphed_async(
online_weights, branching_weights, dqn_cfg, q_provider,
)?;
self.consume_metrics_after_event()
}
/// Async (launch-only) variant of [`evaluate_dqn_graphed`]. Submits the
/// full eval rollout + metrics + diagnostic kernels to `self.stream` and
/// records `eval_done_event`, then returns **without** any host-side wait.
///
/// The companion [`consume_metrics_after_event`] performs the single host
/// `cuEventSynchronize` and parses the mapped-pinned `metrics_buf`. Splitting
/// the two unblocks the host CPU thread to submit the next training epoch's
/// kernels while eval drains on the dedicated `validation_stream` — the
/// pipelining the original async-diag commit (d9cb14f1b) wired the GPU side
/// for but the host-side `synchronize()` inside the now-removed
/// `launch_metrics_and_download` defeated.
///
/// Per-epoch wall-clock saving on L40S: ~25-30s (the original audit's
/// 30-40s estimate, give or take, depending on training-side launch
/// density).
///
/// # Arguments
/// Same as [`evaluate_dqn_graphed`]. The `_online_weights` / `_branching_weights`
/// parameters are placeholders for caller-side ABI symmetry — actual
/// weight wiring goes through the `q_provider`'s `FusedTrainingCtx` which
/// owns the live training weights.
pub fn evaluate_dqn_graphed_async(
&mut self, &mut self,
_online_weights: &DuelingWeightSet, _online_weights: &DuelingWeightSet,
_branching_weights: Option<&BranchingWeightSet>, _branching_weights: Option<&BranchingWeightSet>,
dqn_cfg: &DqnBacktestConfig, dqn_cfg: &DqnBacktestConfig,
q_provider: Option<&mut dyn super::q_value_provider::QValueProvider>, q_provider: Option<&mut dyn super::q_value_provider::QValueProvider>,
) -> Result<Vec<WindowMetrics>, MLError> { ) -> Result<(), MLError> {
let _nvtx = NvtxRange::new("backtest_evaluate_dqn_graphed"); let _nvtx = NvtxRange::new("backtest_evaluate_dqn_graphed_async");
// Reset mutable state: portfolio, done flags, step returns, actions history. // Reset mutable state: portfolio, done flags, step returns, actions history.
// Without this, the second call sees done_flags=1 from the previous evaluation // Without this, the second call sees done_flags=1 from the previous evaluation
@@ -1055,19 +1160,26 @@ impl GpuBacktestEvaluator {
// Val plan_isv parity diagnostic (task #94): reduce plan_state + // Val plan_isv parity diagnostic (task #94): reduce plan_state +
// plan_isv across windows → 8 floats, log for epoch-boundary // plan_isv across windows → 8 floats, log for epoch-boundary
// visibility into plan activity. Cold path, runs once per eval. // visibility into plan activity. Cold path, runs once per eval.
// The host-side `read_all` inside this method reads the PREVIOUS
// epoch's diagnostic via the documented one-epoch lag — the value
// currently in mapped-pinned memory was already synchronised by the
// previous epoch's `consume_metrics_after_event`, so no additional
// host wait is required here.
self.launch_plan_diag_and_log()?; self.launch_plan_diag_and_log()?;
self.launch_metrics_and_download() self.launch_metrics_and_record_event()
} }
/// Val plan_isv parity diagnostic emission (task #94). /// Val plan_isv parity diagnostic emission (task #94).
/// ///
/// Launches `backtest_plan_diag_reduce` on the final plan_state + plan_isv /// Launches `backtest_plan_diag_reduce` on the final plan_state + plan_isv
/// buffers, DtoH-copies the 8-float summary, and logs a single /// buffers and logs a single `val_plan_diag` line with labelled stats.
/// `val_plan_diag` line with labelled stats. Called once per /// Called once per `evaluate_dqn_graphed_async` invocation at the epoch
/// `evaluate_dqn_graphed` invocation at the epoch boundary — a stream /// boundary. The host-side `read_all` here reads the **previous epoch's**
/// sync here is acceptable (already followed by the metrics readback /// diagnostic via the documented one-epoch lag (the value currently in
/// which syncs anyway). /// mapped-pinned memory was already synchronised by the previous epoch's
/// `consume_metrics_after_event`); the just-launched kernel populates the
/// buffer for the **next** epoch's read.
fn launch_plan_diag_and_log(&mut self) -> Result<(), MLError> { fn launch_plan_diag_and_log(&mut self) -> Result<(), MLError> {
let kernel = self.plan_diag_reduce_kernel.as_ref().ok_or_else(|| { let kernel = self.plan_diag_reduce_kernel.as_ref().ok_or_else(|| {
MLError::ModelError("plan_diag_reduce_kernel not loaded".to_owned()) MLError::ModelError("plan_diag_reduce_kernel not loaded".to_owned())
@@ -1076,9 +1188,8 @@ impl GpuBacktestEvaluator {
// Mapped-pinned device pointer for the diag scratch — kernel writes // Mapped-pinned device pointer for the diag scratch — kernel writes
// through `dev_ptr`, host reads via `host_ptr`. No DtoH per // through `dev_ptr`, host reads via `host_ptr`. No DtoH per
// `feedback_no_htod_htoh_only_mapped_pinned.md`. The metrics readback // `feedback_no_htod_htoh_only_mapped_pinned.md`. The metrics readback
// immediately after this call also syncs the stream, so we elide a // (now `consume_metrics_after_event`) is the single host-side wait
// dedicated sync here — `launch_metrics_and_download` is the // per epoch boundary.
// single sync point per epoch boundary.
let plan_diag_dev_ptr = self.plan_diag_buf.dev_ptr; let plan_diag_dev_ptr = self.plan_diag_buf.dev_ptr;
unsafe { unsafe {
self.stream self.stream
@@ -1096,14 +1207,15 @@ impl GpuBacktestEvaluator {
"backtest_plan_diag_reduce launch: {e}" "backtest_plan_diag_reduce launch: {e}"
)))?; )))?;
} }
// Read mapped-pinned host buffer. Caller (`evaluate_dqn_graphed`) // Read mapped-pinned host buffer. The values currently in memory were
// syncs the eval stream inside `launch_metrics_and_download` before // produced by the PREVIOUS epoch's plan_diag launch and made globally
// returning, so the values produced by the just-launched kernel are // visible by the previous epoch's `consume_metrics_after_event` (which
// globally visible to the CPU by the time the diag info!() reaches // calls `eval_done_event.synchronize()`); the kernel we just launched
// the next eval call. The first epoch returns zeros (safe default // above will populate the buffer for the NEXT epoch's read. The first
// from `MappedF32Buffer::new` zero-init); subsequent epochs see the // epoch returns zeros (safe default from `MappedF32Buffer::new`
// previous epoch's diagnostic — same one-epoch lag pattern as // zero-init); subsequent epochs see the previous epoch's diagnostic —
// `pending_val_loss`. Acceptable for a diagnostic-only readout. // same one-epoch lag pattern as `pending_val_loss`. Acceptable for a
// diagnostic-only readout.
let diag = self.plan_diag_buf.read_all(); let diag = self.plan_diag_buf.read_all();
info!( info!(
target: "val_plan_diag", target: "val_plan_diag",
@@ -1133,12 +1245,15 @@ impl GpuBacktestEvaluator {
/// Returns `[quarter_frac, half_frac, full_frac]` summing to 1.0 over /// Returns `[quarter_frac, half_frac, full_frac]` summing to 1.0 over
/// non-negative action entries (negative entries indicate uninitialized /// non-negative action entries (negative entries indicate uninitialized
/// or skipped steps and are excluded). /// or skipped steps and are excluded).
///
/// **Sync contract**: caller MUST have invoked
/// [`consume_metrics_after_event`] (or the synchronous `evaluate_*`
/// wrappers, which call it internally) since the last
/// [`launch_metrics_and_record_event`]. The reader does no host wait —
/// it just `read_volatile`s the mapped-pinned mirror that
/// `launch_metrics_and_record_event` queued a DtoD-async copy into.
pub fn read_eval_action_distribution_per_magnitude(&self) -> Result<[f32; 3], MLError> { pub fn read_eval_action_distribution_per_magnitude(&self) -> Result<[f32; 3], MLError> {
let len = self.n_windows * self.max_len; let host = self.actions_history_pinned.read_all();
let mut host = vec![0_i32; len];
self.stream
.memcpy_dtoh(&self.actions_history_buf, &mut host)
.map_err(|e| MLError::ModelError(format!("actions_history dtoh: {e}")))?;
let mut counts = [0_u64; 3]; let mut counts = [0_u64; 3];
let mut total = 0_u64; let mut total = 0_u64;
for &a in &host { for &a in &host {
@@ -1171,19 +1286,18 @@ impl GpuBacktestEvaluator {
/// the policy's Q-head actually prefer when Kelly cold-start is masking /// the policy's Q-head actually prefer when Kelly cold-start is masking
/// everything to Quarter via `actual_mag`?" /// everything to Quarter via `actual_mag`?"
/// ///
/// Returns `Ok([0.0; 3])` if `intent_mag_buf` has not been allocated /// Returns `Ok([0.0; 3])` if `intent_mag_pinned` has not been allocated
/// yet (evaluator hasn't run a DQN eval), or if all entries are /// yet (evaluator hasn't run a DQN eval), or if all entries are
/// skipped (zero rollout). /// skipped (zero rollout).
///
/// **Sync contract**: same as
/// [`read_eval_action_distribution_per_magnitude`].
pub fn read_eval_intent_magnitude_distribution(&self) -> Result<[f32; 3], MLError> { pub fn read_eval_intent_magnitude_distribution(&self) -> Result<[f32; 3], MLError> {
let buf = match self.intent_mag_buf.as_ref() { let buf = match self.intent_mag_pinned.as_ref() {
Some(b) => b, Some(b) => b,
None => return Ok([0.0; 3]), None => return Ok([0.0; 3]),
}; };
let len = self.n_windows * self.max_len; let host = buf.read_all();
let mut host = vec![0_i32; len];
self.stream
.memcpy_dtoh(buf, &mut host)
.map_err(|e| MLError::ModelError(format!("intent_mag_buf dtoh: {e}")))?;
let mut counts = [0_u64; 3]; let mut counts = [0_u64; 3];
let mut total = 0_u64; let mut total = 0_u64;
for &a in &host { for &a in &host {
@@ -1219,12 +1333,11 @@ impl GpuBacktestEvaluator {
/// direction is healthy but magnitude itself is collapsing. Paired with /// direction is healthy but magnitude itself is collapsing. Paired with
/// `read_eval_action_distribution_per_magnitude` to localise the /// `read_eval_action_distribution_per_magnitude` to localise the
/// magnitude-only Task 2.X fix's leverage surface. /// magnitude-only Task 2.X fix's leverage surface.
///
/// **Sync contract**: same as
/// [`read_eval_action_distribution_per_magnitude`].
pub fn read_eval_action_distribution_per_direction(&self) -> Result<[f32; 4], MLError> { pub fn read_eval_action_distribution_per_direction(&self) -> Result<[f32; 4], MLError> {
let len = self.n_windows * self.max_len; let host = self.actions_history_pinned.read_all();
let mut host = vec![0_i32; len];
self.stream
.memcpy_dtoh(&self.actions_history_buf, &mut host)
.map_err(|e| MLError::ModelError(format!("actions_history dtoh: {e}")))?;
let mut counts = [0_u64; 4]; let mut counts = [0_u64; 4];
let mut total = 0_u64; let mut total = 0_u64;
for &a in &host { for &a in &host {
@@ -1250,10 +1363,11 @@ impl GpuBacktestEvaluator {
} }
/// Read the per-direction histogram of the policy's RAW Boltzmann picks /// Read the per-direction histogram of the policy's RAW Boltzmann picks
/// across the FULL val window. Reads `picked_action_history_buf` — /// across the FULL val window. Reads `picked_action_history_pinned` —
/// scattered from `chunked_actions_buf` after each chunk's /// the mapped-pinned mirror of `picked_action_history_buf`, scattered
/// `experience_action_select`, so the buffer accumulates the kernel's /// from `chunked_actions_buf` after each chunk's `experience_action_select`,
/// raw action_idx writes for every bar before env_step touches anything. /// so the buffer accumulates the kernel's raw action_idx writes for every
/// bar before env_step touches anything.
/// ///
/// `actions_history_buf` records the POST-physics `actual_dir`, so a /// `actions_history_buf` records the POST-physics `actual_dir`, so a
/// divergence between this reader and `read_eval_action_distribution_per_direction` /// divergence between this reader and `read_eval_action_distribution_per_direction`
@@ -1262,18 +1376,17 @@ impl GpuBacktestEvaluator {
/// drains active picks to Flat → only post-physics flat). /// drains active picks to Flat → only post-physics flat).
/// ///
/// Returns `[short, hold, long, flat]` summing to 1.0 over the full window. /// Returns `[short, hold, long, flat]` summing to 1.0 over the full window.
/// `Ok([0.0; 4])` if `picked_action_history_buf` has not been allocated yet /// `Ok([0.0; 4])` if `picked_action_history_pinned` has not been allocated yet
/// (first call before `ensure_action_select_ready` has run). /// (first call before `ensure_action_select_ready` has run).
///
/// **Sync contract**: same as
/// [`read_eval_action_distribution_per_magnitude`].
pub fn read_chunked_actions_direction_distribution(&self) -> Result<[f32; 4], MLError> { pub fn read_chunked_actions_direction_distribution(&self) -> Result<[f32; 4], MLError> {
let buf = match self.picked_action_history_buf.as_ref() { let buf = match self.picked_action_history_pinned.as_ref() {
Some(b) => b, Some(b) => b,
None => return Ok([0.0; 4]), None => return Ok([0.0; 4]),
}; };
let len = buf.len(); let host = buf.read_all();
let mut host = vec![0_i32; len];
self.stream
.memcpy_dtoh(buf, &mut host)
.map_err(|e| MLError::ModelError(format!("picked_action_history dtoh: {e}")))?;
let mut counts = [0_u64; 4]; let mut counts = [0_u64; 4];
let mut total = 0_u64; let mut total = 0_u64;
for &a in &host { for &a in &host {
@@ -1822,8 +1935,9 @@ impl GpuBacktestEvaluator {
self.launch_env_step(step)?; self.launch_env_step(step)?;
} }
// Metrics reduction + download // Synchronous launch+consume — `evaluate_ppo` is a single-shot path.
self.launch_metrics_and_download() self.launch_metrics_and_record_event()?;
self.consume_metrics_after_event()
} }
// ── Supervised evaluation (hybrid: candle forward + CUDA signal→action) ── // ── Supervised evaluation (hybrid: candle forward + CUDA signal→action) ──
@@ -1906,8 +2020,9 @@ impl GpuBacktestEvaluator {
self.launch_env_step(step)?; self.launch_env_step(step)?;
} }
// Metrics reduction + download // Synchronous launch+consume — `evaluate_supervised` is a single-shot path.
self.launch_metrics_and_download() self.launch_metrics_and_record_event()?;
self.consume_metrics_after_event()
} }
// ── DQN forward helpers ────────────────────────────────────────────── // ── DQN forward helpers ──────────────────────────────────────────────
@@ -1981,6 +2096,24 @@ impl GpuBacktestEvaluator {
let picked_action_history_buf = self.stream.alloc_zeros::<i32>(n * self.max_len) let picked_action_history_buf = self.stream.alloc_zeros::<i32>(n * self.max_len)
.map_err(|e| MLError::ModelError(format!("alloc picked_action_history_buf: {e}")))?; .map_err(|e| MLError::ModelError(format!("alloc picked_action_history_buf: {e}")))?;
// Mapped-pinned mirrors for the two diagnostic history buffers above.
// Filled via `memcpy_dtod_async` on the eval stream before
// `eval_done_event` is recorded; read by the per-direction / intent /
// picked distribution helpers after the event syncs. See
// `actions_history_pinned` in the struct field docstring.
let intent_mag_pinned = unsafe {
MappedI32Buffer::new(n * self.max_len)
.map_err(|e| MLError::ModelError(format!(
"intent_mag mapped-pinned alloc: {e}"
)))?
};
let picked_action_history_pinned = unsafe {
MappedI32Buffer::new(n * self.max_len)
.map_err(|e| MLError::ModelError(format!(
"picked_action_history mapped-pinned alloc: {e}"
)))?
};
let ch_q_gaps = self.stream.alloc_zeros::<f32>(cn) let ch_q_gaps = self.stream.alloc_zeros::<f32>(cn)
.map_err(|e| MLError::ModelError(format!("alloc chunked q_gaps: {e}")))?; .map_err(|e| MLError::ModelError(format!("alloc chunked q_gaps: {e}")))?;
@@ -2014,6 +2147,8 @@ impl GpuBacktestEvaluator {
self.chunked_q_gaps_buf = Some(ch_q_gaps); self.chunked_q_gaps_buf = Some(ch_q_gaps);
self.intent_mag_buf = Some(intent_mag_buf); self.intent_mag_buf = Some(intent_mag_buf);
self.picked_action_history_buf = Some(picked_action_history_buf); self.picked_action_history_buf = Some(picked_action_history_buf);
self.intent_mag_pinned = Some(intent_mag_pinned);
self.picked_action_history_pinned = Some(picked_action_history_pinned);
self.scatter_intent_kernel = Some(scatter_intent); self.scatter_intent_kernel = Some(scatter_intent);
self.chunked_conviction_buf = Some(ch_conviction); self.chunked_conviction_buf = Some(ch_conviction);
self.chunked_magnitude_conviction_buf = Some(ch_magnitude_conviction); self.chunked_magnitude_conviction_buf = Some(ch_magnitude_conviction);
@@ -2128,10 +2263,19 @@ impl GpuBacktestEvaluator {
Ok(()) Ok(())
} }
/// Launch metrics reduction kernel and download results. /// Launch the metrics reduction kernel and queue async DtoD copies of the
/// per-step action history buffers into their mapped-pinned mirrors. Records
/// `eval_done_event` on `self.stream` after the queued work, then returns
/// **without** any host-side wait.
/// ///
/// Shared by all evaluation paths. /// This is the launch half of the perf-fix split (2026-04-28). The companion
fn launch_metrics_and_download(&self) -> Result<Vec<WindowMetrics>, MLError> { /// `consume_metrics_after_event` calls `event.synchronize()` and reads through
/// the mapped-pinned `host_ptr`. Splitting the two lets the host CPU thread
/// submit the next epoch's training kernels on `cuda_stream` while eval drains
/// asynchronously on `validation_stream`.
///
/// Shared by all evaluation paths via the thin `evaluate_*` wrappers.
pub fn launch_metrics_and_record_event(&mut self) -> Result<(), MLError> {
// Shared memory: 6 reduction arrays × 256 threads × 4 bytes = 6 KB // Shared memory: 6 reduction arrays × 256 threads × 4 bytes = 6 KB
// + 4096 floats for sort scratch / boundary data = 16 KB // + 4096 floats for sort scratch / boundary data = 16 KB
// Total = 22 KB (well within the 48 KB L1/shmem limit) // Total = 22 KB (well within the 48 KB L1/shmem limit)
@@ -2171,15 +2315,100 @@ impl GpuBacktestEvaluator {
.map_err(|e| MLError::ModelError(format!("compute_backtest_metrics launch: {e}")))?; .map_err(|e| MLError::ModelError(format!("compute_backtest_metrics launch: {e}")))?;
} }
// Synchronise the eval stream so the kernel's writes through the // Queue DtoD-async copies from the device action-history buffers into
// mapped device pointer are globally visible before the host reads // their mapped-pinned mirrors. The kernel above (and the env-step kernel
// via `host_ptr` (the kernel uses `__threadfence_system()` semantics // chain that produced these buffers) ordered before this point on the
// implicit in stream sync). This is the ONLY sync per evaluation — // same stream, so the copies see the post-rollout values. The host reads
// it blocks only the eval stream, not the training stream, so // through `host_ptr` after `eval_done_event.synchronize()` — a single
// training continues uninterrupted on `cuda_stream`. // barrier covers the metrics readback AND all four distribution readers.
self.stream //
.synchronize() // Mirrors are unconditional for `actions_history_buf` (always allocated)
.map_err(|e| MLError::ModelError(format!("metrics readback sync: {e}")))?; // and conditional for `intent_mag_buf` / `picked_action_history_buf`
// (allocated only on DQN paths via `ensure_action_select_ready`). Skip
// the conditional copies when the source buffers are not yet wired —
// those readers return `Ok([0.0; *])` from the unwritten zero-init.
unsafe {
let n_bytes_full = self.n_windows * self.max_len * std::mem::size_of::<i32>();
// actions_history_buf → actions_history_pinned
{
let (src_ptr, _src_guard) =
self.actions_history_buf.device_ptr(&self.stream);
let dst_ptr = self.actions_history_pinned.dev_ptr;
cudarc::driver::result::memcpy_dtod_async(
dst_ptr, src_ptr, n_bytes_full, self.stream.cu_stream(),
).map_err(|e| MLError::ModelError(format!(
"actions_history dtod-to-pinned: {e}"
)))?;
}
// intent_mag_buf → intent_mag_pinned (DQN path only)
if let (Some(src_buf), Some(dst_buf)) =
(self.intent_mag_buf.as_ref(), self.intent_mag_pinned.as_ref())
{
let (src_ptr, _src_guard) = src_buf.device_ptr(&self.stream);
cudarc::driver::result::memcpy_dtod_async(
dst_buf.dev_ptr, src_ptr, n_bytes_full, self.stream.cu_stream(),
).map_err(|e| MLError::ModelError(format!(
"intent_mag dtod-to-pinned: {e}"
)))?;
}
// picked_action_history_buf → picked_action_history_pinned (DQN path).
// The picked buffer has the same shape `[n_windows * max_len]` (see
// its allocation in `ensure_action_select_ready`).
if let (Some(src_buf), Some(dst_buf)) = (
self.picked_action_history_buf.as_ref(),
self.picked_action_history_pinned.as_ref(),
) {
let (src_ptr, _src_guard) = src_buf.device_ptr(&self.stream);
cudarc::driver::result::memcpy_dtod_async(
dst_buf.dev_ptr, src_ptr, n_bytes_full, self.stream.cu_stream(),
).map_err(|e| MLError::ModelError(format!(
"picked_action_history dtod-to-pinned: {e}"
)))?;
}
}
// Lazy-create the event on first launch (so the CUDA context is
// guaranteed active). Re-record into the same event handle on every
// subsequent launch — `record` overwrites the prior recording per the
// CUDA driver semantics, no allocator churn.
if self.eval_done_event.is_none() {
let ev = self.stream.context().new_event(None).map_err(|e| {
MLError::ModelError(format!("eval_done_event create: {e}"))
})?;
self.eval_done_event = Some(ev);
}
let ev = self.eval_done_event.as_ref().expect("just created above");
ev.record(&self.stream).map_err(|e| {
MLError::ModelError(format!("eval_done_event record: {e}"))
})?;
Ok(())
}
/// Wait on `eval_done_event` (host-side `cuEventSynchronize`) and parse the
/// per-window metrics from the mapped-pinned `metrics_buf`. The event's GPU
/// dependency chain ensures the metrics kernel + all action-history mirror
/// copies queued by `launch_metrics_and_record_event` are globally visible
/// before the host load.
///
/// This is the consume half of the perf-fix split (2026-04-28). Pipelined
/// callers (`compute_validation_loss_consume`) invoke this at the start of
/// the next epoch — the host CPU thread already submitted the previous
/// epoch's training to `cuda_stream`, so by the time control reaches here
/// the eval stream has typically drained and the event is already complete.
/// Synchronous callers (the `evaluate_*` wrappers used by hyperopt / smoke
/// tests) invoke launch + consume back-to-back — no wall-clock change for
/// those paths since they had nothing to overlap with anyway.
pub fn consume_metrics_after_event(&self) -> Result<Vec<WindowMetrics>, MLError> {
let ev = self.eval_done_event.as_ref().ok_or_else(|| {
MLError::ModelError(
"consume_metrics_after_event called before launch_metrics_and_record_event"
.to_owned(),
)
})?;
ev.synchronize().map_err(|e| {
MLError::ModelError(format!("eval_done_event synchronize: {e}"))
})?;
let metrics_host: Vec<f32> = self.metrics_buf.read_all(); let metrics_host: Vec<f32> = self.metrics_buf.read_all();
// Parse flat metrics into per-window structs. // Parse flat metrics into per-window structs.

View File

@@ -474,23 +474,31 @@ impl DQNTrainer {
None None
} }
/// Compute validation loss (negative Sharpe) using the GPU backtest evaluator. /// Launch validation backtest on the dedicated `validation_stream` —
/// submits all eval kernels, queues the metrics readback into mapped-pinned
/// memory, records `eval_done_event`, then returns **without** any host
/// wait. The companion [`consume_validation_loss`] performs the single
/// host `cuEventSynchronize` and parses the metrics for HEALTH_DIAG emit.
///
/// Splitting launch from consume restores the cross-stream pipelining the
/// dedicated `validation_stream` is supposed to provide: the host CPU
/// thread immediately returns to submit the NEXT epoch's training kernels
/// to `cuda_stream` while eval drains on `validation_stream`. The previous
/// monolithic `compute_validation_loss` blocked inside `synchronize()`
/// after the metrics launch, defeating the pipelining.
/// ///
/// Runs the full backtest step loop on GPU: gather kernel -> cuBLAS forward -> /// Runs the full backtest step loop on GPU: gather kernel -> cuBLAS forward ->
/// env kernel -> metrics reduction. Only the final metrics (1 window x 14 floats) /// env kernel -> metrics reduction. The evaluator is lazy-initialized on the
/// are downloaded to CPU. Replaces the old CPU path (GPU->CPU Q-value download, /// first call and cached across epochs. Weight pointers change each epoch
/// CPU argmax, CPU Sharpe) which took ~2s/epoch vs ~50ms here. /// (training updates weights), so we call `invalidate_dqn_graph()` to force
/// /// CUDA Graph re-capture.
/// The evaluator is lazy-initialized on the first call and cached across epochs. pub(crate) async fn launch_validation_loss(&mut self) -> Result<()> {
/// Weight pointers change each epoch (training updates weights), so we call
/// `invalidate_dqn_graph()` to force CUDA Graph re-capture.
pub(crate) async fn compute_validation_loss(&mut self) -> Result<f64> {
use crate::cuda_pipeline::gpu_backtest_evaluator::{ use crate::cuda_pipeline::gpu_backtest_evaluator::{
DqnBacktestConfig, GpuBacktestConfig, GpuBacktestEvaluator, DqnBacktestConfig, GpuBacktestConfig, GpuBacktestEvaluator,
}; };
if self.val_data.is_empty() { if self.val_data.is_empty() {
return Ok(0.0); return Ok(());
} }
// Run validation on the dedicated `validation_stream` so eval kernels // Run validation on the dedicated `validation_stream` so eval kernels
@@ -518,7 +526,7 @@ impl DQNTrainer {
.or(self.cuda_stream.as_ref()); .or(self.cuda_stream.as_ref());
let stream = match val_stream { let stream = match val_stream {
Some(s) => s, Some(s) => s,
None => return Ok(0.0), // CPU device — skip GPU validation None => return Ok(()), // CPU device — skip GPU validation
}; };
let train_done_event = if let (Some(ref vs), Some(ref ts)) = let train_done_event = if let (Some(ref vs), Some(ref ts)) =
(self.validation_stream.as_ref(), self.cuda_stream.as_ref()) (self.validation_stream.as_ref(), self.cuda_stream.as_ref())
@@ -700,14 +708,46 @@ impl DQNTrainer {
}; };
// Safety: fused_ptr and evaluator are independent fields — no aliasing. // Safety: fused_ptr and evaluator are independent fields — no aliasing.
let metrics = unsafe { // `_async` submits eval kernels + records `eval_done_event` then returns
evaluator.evaluate_dqn_graphed( // immediately — no host wait. The matching `consume_validation_loss`
// (called at the start of the NEXT epoch) does the single
// `cuEventSynchronize` and parses the metrics for HEALTH_DIAG emit.
unsafe {
evaluator.evaluate_dqn_graphed_async(
&*online_weights_ptr, &*online_weights_ptr,
Some(&*branching_weights_ptr), Some(&*branching_weights_ptr),
&dqn_cfg, &dqn_cfg,
Some(&mut *fused_ptr), Some(&mut *fused_ptr),
) )
}.map_err(|e| anyhow::anyhow!("GPU backtest evaluate_dqn_graphed: {e}"))?; }.map_err(|e| anyhow::anyhow!("GPU backtest evaluate_dqn_graphed_async: {e}"))?;
Ok(())
}
/// Wait on `eval_done_event` (host-side `cuEventSynchronize`), parse the
/// per-window metrics from mapped-pinned memory, emit HEALTH_DIAG `val[...]`
/// + per-direction / intent-magnitude / picked-direction distribution lines,
/// and return the val_loss (= -sharpe) for downstream early-stopping /
/// best-checkpoint logic.
///
/// **MUST** be called exactly once after each [`launch_validation_loss`].
/// Calling without a prior launch returns `Ok(0.0)` (the val_data-empty
/// fall-through).
pub(crate) async fn consume_validation_loss(&mut self) -> Result<f64> {
if self.val_data.is_empty() {
return Ok(0.0);
}
let evaluator = match self.gpu_evaluator.as_mut() {
Some(ev) => ev,
// Lazy-init only happens inside `launch_validation_loss`. If we got
// here without a launch (e.g., CPU device path), there's nothing to
// consume — return 0 to match the empty-val_data fall-through.
None => return Ok(0.0),
};
let metrics = evaluator.consume_metrics_after_event()
.map_err(|e| anyhow::anyhow!("GPU backtest consume_metrics: {e}"))?;
// Single window — log full validation metrics // Single window — log full validation metrics
let val_sharpe = if let Some(m) = metrics.first() { let val_sharpe = if let Some(m) = metrics.first() {

View File

@@ -299,10 +299,30 @@ impl DQNTrainer {
fused.write_isv_signal_at(TLOB_REGIME_FOCUS_EMA_INDEX, ema); fused.write_isv_signal_at(TLOB_REGIME_FOCUS_EMA_INDEX, ema);
} }
// Collect pending async validation from PREVIOUS epoch // Collect pending async validation from PREVIOUS epoch.
if let Some(val) = self.pending_val_loss.take() { //
// Store for this epoch's logging (one-epoch delay is acceptable) // `pending_val_loss = Some(0.0)` is used as the "launched but not yet
self.cached_async_val_loss = Some(val); // consumed" sentinel set by `launch_validation_loss` last epoch. The
// actual val_loss is computed here by `consume_validation_loss`,
// which calls `cuEventSynchronize` on the evaluator's `eval_done_event`,
// reads the mapped-pinned `metrics_buf`, and emits the HEALTH_DIAG
// `val [...]` + `val_dir_dist` / `val_picked_dir_dist` /
// `val_plan_isv_diag` lines. The host wait happens HERE — at the
// start of the next epoch, AFTER training has already submitted its
// initial epoch-boundary work to `cuda_stream`. Eval has been
// draining on `validation_stream` since last epoch's launch, so the
// event is typically already complete; the host wait is a no-op or
// near-no-op in steady state. This is the point of the pipelining
// (see `gpu_backtest_evaluator::launch_metrics_and_record_event`).
if self.pending_val_loss.take().is_some() {
match self.consume_validation_loss().await {
Ok(val) => {
self.cached_async_val_loss = Some(val);
}
Err(e) => {
tracing::warn!("Async validation consume failed: {e}");
}
}
} }
self.reset_epoch_state(epoch); self.reset_epoch_state(epoch);
@@ -1000,6 +1020,27 @@ impl DQNTrainer {
} }
} }
// Drain any in-flight async validation launched by the last epoch's
// reporting block. Without this, the LAST epoch's eval kernels are
// still queued on `validation_stream` when training returns —
// `eval_done_event` is recorded but never `synchronize`d, so the
// mapped-pinned mirrors (`actions_history_pinned` etc.) carry the
// SECOND-to-last epoch's data. Smoke tests / hyperopt that read
// `last_eval_direction_dist` / `last_eval_magnitude_dist` after
// training expect the LAST epoch's data — this consume preserves that.
// Failure is non-fatal: the readers fall back to whatever the previous
// consume left in the mirrors (one-epoch lag, same as in-loop).
if self.pending_val_loss.take().is_some() {
match self.consume_validation_loss().await {
Ok(val) => {
self.cached_async_val_loss = Some(val);
}
Err(e) => {
tracing::warn!("Final async validation drain failed: {e}");
}
}
}
// Drain any in-flight async best-checkpoint workers before // Drain any in-flight async best-checkpoint workers before
// returning. Guarantees the most recent improvement's safetensors // returning. Guarantees the most recent improvement's safetensors
// bytes are written to disk/MinIO before the caller observes the // bytes are written to disk/MinIO before the caller observes the
@@ -3807,21 +3848,41 @@ impl DQNTrainer {
} }
// Validation loss — async on separate stream when available. // Validation loss — async on separate stream when available.
// Uses cached result from previous epoch (one-epoch delay) to avoid //
// blocking the training stream. On epoch 0, runs synchronously. // Per-epoch flow (post async-pipeline-fix 2026-04-28):
// 1. Reporting time (here): consume cached value from PREVIOUS epoch's
// async launch, then submit THIS epoch's eval kernels via
// `launch_validation_loss` and immediately return — no host wait.
// Set `pending_val_loss = Some(0.0)` as the "launched" sentinel
// (the value 0.0 is NEVER read; the real val_loss comes from
// `consume_validation_loss` at the next epoch's start).
// 2. Next epoch's start: takes the sentinel, calls
// `consume_validation_loss` which does `event.synchronize()` and
// parses the mapped-pinned `metrics_buf`. That's the SOLE host
// wait per epoch boundary.
//
// Epoch 0 (no cached val yet): synchronous launch + consume to seed the
// pipeline. Wall-clock ~25-30s slower than steady-state epoch 1+, but
// unavoidable for bootstrap correctness. From epoch 1 onward, the
// host-side eval cost is amortised via the cross-stream pipelining.
let val_loss = if self.validation_stream.is_some() && self.cached_async_val_loss.is_some() { let val_loss = if self.validation_stream.is_some() && self.cached_async_val_loss.is_some() {
let cached = self.cached_async_val_loss.unwrap_or(0.0); let cached = self.cached_async_val_loss.unwrap_or(0.0);
// Launch this epoch's validation asynchronously (result collected next epoch) match self.launch_validation_loss().await {
match self.compute_validation_loss().await { Ok(()) => {
Ok(val) => { self.pending_val_loss = Some(val); } // Sentinel — picked up by the consume block at the next
Err(e) => { tracing::warn!("Async validation failed: {e}"); } // epoch's start. Value is irrelevant; only the `is_some()`
// matters for routing the next-epoch consume call.
self.pending_val_loss = Some(0.0);
}
Err(e) => { tracing::warn!("Async validation launch failed: {e}"); }
} }
cached cached
} else { } else {
// Epoch 0 or no validation stream — run synchronously // Epoch 0 or no validation stream — sync launch + consume to seed
let val = self.compute_validation_loss().await?; // `cached_async_val_loss` for the next epoch's async path.
// Also store for next epoch's async pipeline bootstrap self.launch_validation_loss().await?;
self.pending_val_loss = Some(val); let val = self.consume_validation_loss().await?;
self.cached_async_val_loss = Some(val);
val val
}; };
let val_sharpe = -val_loss; // val_loss = -sharpe (lower is better convention) let val_sharpe = -val_loss; // val_loss = -sharpe (lower is better convention)
@@ -3846,7 +3907,7 @@ impl DQNTrainer {
// C.3 Plan 3 Task 7: state-distribution KL between train sample and val // C.3 Plan 3 Task 7: state-distribution KL between train sample and val
// batch (OFI block). Launches once per validation epoch — must run AFTER // batch (OFI block). Launches once per validation epoch — must run AFTER
// `compute_validation_loss` populates the evaluator's `chunked_states_buf` // `launch_validation_loss` populates the evaluator's `chunked_states_buf`
// and BEFORE the next training epoch consumes ISV[STATE_KL_AMP=79]. The // and BEFORE the next training epoch consumes ISV[STATE_KL_AMP=79]. The
// val state pointer comes from the evaluator's chunked buffer; the train // val state pointer comes from the evaluator's chunked buffer; the train
// states are the experience-collector's own `states_out` (last collect // states are the experience-collector's own `states_out` (last collect

View File

@@ -94,3 +94,22 @@
### Fix 3 — Convert `tau_buf` from device VRAM to pinned+device-mapped ### Fix 3 — Convert `tau_buf` from device VRAM to pinned+device-mapped
`gpu_iqn_head.rs`: `tau_buf: CudaSlice<f32>` + `tau_host: f32` + `cuMemcpyHtoDAsync_v2` replaced by `tau_pinned: *mut f32` + `tau_dev_ptr: u64` allocated via `cuMemAllocHost_v2(DEVICEMAP)` + `cuMemHostGetDevicePointer_v2`. CPU writes `*tau_pinned = tau` directly; the EMA kernel reads via `tau_dev_ptr` with no PCIe copy. Drop impl updated to free `tau_pinned`. `gpu_iqn_head.rs`: `tau_buf: CudaSlice<f32>` + `tau_host: f32` + `cuMemcpyHtoDAsync_v2` replaced by `tau_pinned: *mut f32` + `tau_dev_ptr: u64` allocated via `cuMemAllocHost_v2(DEVICEMAP)` + `cuMemHostGetDevicePointer_v2`. CPU writes `*tau_pinned = tau` directly; the EMA kernel reads via `tau_dev_ptr` with no PCIe copy. Drop impl updated to free `tau_pinned`.
### Fix 4 (perf) — Async-eval pipeline: split `launch_metrics_and_download` (2026-04-28)
`gpu_backtest_evaluator.rs` + `trainer/metrics.rs` + `trainer/training_loop.rs`: the prior merge `3c0d26292` (async-diag — eval on dedicated stream + spawn_blocking checkpoint, building on `d9cb14f1b` + `673b04a8d`) wired the **GPU side** of cross-stream pipelining correctly (training stream `cuda_stream`, eval stream `validation_stream`, `cuStreamWaitEvent` barrier on `train_done_event`) but left a host-side `stream.synchronize()` inside `launch_metrics_and_download` that blocked the CPU thread for the FULL eval drain (~25-30s/epoch on L40S). The host thread is the same one that submits the next epoch's training kernels via `run_full_step` — so until eval drained, training submission was gated on it, and the dedicated `validation_stream` provided no actual wall-clock benefit.
Split into two record/consume halves at the buffer-readback layer:
* `launch_metrics_and_record_event(&mut self) -> Result<(), MLError>` — submits the metrics kernel + DtoD-async copies of `actions_history_buf` / `intent_mag_buf` / `picked_action_history_buf` into mapped-pinned mirrors (`actions_history_pinned: MappedI32Buffer` allocated in `new()`, the other two lazy-allocated alongside their device buffers in `ensure_action_select_ready`). Records `eval_done_event` (lazy-created on first launch via `stream.context().new_event(None)`, re-recorded each epoch via `event.record(stream)`). **Returns immediately** — no host wait.
* `consume_metrics_after_event(&self) -> Result<Vec<WindowMetrics>, MLError>` — calls `event.synchronize()` (the SOLE host-side wait per epoch boundary), then `read_volatile`s the mapped-pinned `metrics_buf` for the per-window metrics. After the event syncs, the four action-distribution helpers (`read_eval_action_distribution_per_*`, `read_eval_intent_magnitude_distribution`, `read_chunked_actions_direction_distribution`) read directly from the now-coherent mapped-pinned mirrors — no `memcpy_dtoh` per call (the prior implementation allocated a fresh `Vec<i32>` of size `n_windows * max_len` and did a synchronous DtoH each call, which violates `feedback_no_htod_htoh_only_mapped_pinned.md` AND added 4 extra host-side stalls beyond the metrics one). All four mirrors land in the same event-sync barrier as the metrics readback.
Caller migration (`feedback_no_partial_refactor.md`):
* `evaluate_dqn_graphed` (existing public API, used by hyperopt + smoke test) becomes a thin synchronous wrapper: launch_async → consume. Public ABI unchanged.
* New `evaluate_dqn_graphed_async` — launch-only sibling for the pipelined path. Returns `()`.
* `evaluate` / `evaluate_ppo` / `evaluate_supervised` migrated inline to launch+consume — no caller pipelines them, so no behavioural change.
* `compute_validation_loss` removed; replaced by `launch_validation_loss` (returns `()`) + `consume_validation_loss` (returns `f64` and emits the `HEALTH_DIAG val [...]` / `val_dir_dist` / `val_picked_dir_dist` lines).
* `training_loop.rs`'s validation reporting block now: epoch-start consumes any pending launch into `cached_async_val_loss`, epoch-end reports the cached value and submits a fresh launch with `pending_val_loss = Some(0.0)` as the "launched" sentinel. Bootstrap (epoch 0 / no validation_stream) does sync launch+consume. After the for-epoch loop terminates, a final consume drains the LAST epoch's launch so smoke tests / hyperopt that read `last_eval_direction_dist()` see the most recent epoch's data, not the second-to-last.
Mathematical identity preserved: every consumed val_sharpe is bit-identical to what the prior `compute_validation_loss` would have produced for the same epoch — the only change is *when* the host parses it (one epoch later, which matches the existing `cached_async_val_loss` lag semantics; HEALTH_DIAG `val [...]` emit moves with the consume). Expected wall-clock saving: ~25-30s/epoch on L40S.
Out-of-scope (remaining host-side stream syncs in eval, not part of this fix): `gpu_backtest_evaluator.rs` lines ~1596 + ~1863 — periodic `self.stream.synchronize()` calls every 10 chunks inside `submit_dqn_step_loop_cublas` for kernel-error detection. They block the host CPU thread but only for the chunked rollout duration of ~5 calls × per-chunk-residency-time per epoch (each catches deferred CUDA errors within 10 chunks of origin). Smaller in aggregate than the `launch_metrics_and_download` final wait this fix removed; future work to convert to async error-polling (e.g., `eval_done_event.is_complete()` polling from training submissions). Not adjacent to the perf-fix scope.