diff --git a/crates/ml/src/cuda_pipeline/experience_kernels.cu b/crates/ml/src/cuda_pipeline/experience_kernels.cu index 42051e5bb..fced99550 100644 --- a/crates/ml/src/cuda_pipeline/experience_kernels.cu +++ b/crates/ml/src/cuda_pipeline/experience_kernels.cu @@ -1012,20 +1012,23 @@ extern "C" __global__ void experience_action_select( * Q-gap = Q(best_direction) - Q(Flat). Higher = more conviction. * Used by env_step to scale target_position. * - * D7/N7: During contrarian override, the action selected is argmin rather - * than argmax. The raw Q(best) - Q(Flat) gap would mislead the position - * sizer into full conviction for what's actually the "deliberately wrong" - * action. Zero out the gap so contrarian trades size minimal. */ + * F6: Compute Q-gap with the same sign convention used by the action-selection + * softmax. When contrarian_active, this means the "best" action is the argmin + * of raw Q (equivalent to argmax of -Q). The conviction magnitude then reflects + * how confident the contrarian direction is — useful for position sizing + * consistently with the action actually being taken. */ if (out_q_gaps != NULL) { - if (contrarian_active != 0) { - out_q_gaps[i] = 0.0f; - } else { - int flat_idx = 3; /* Flat = index 3 for direction(4): Short(0)/Hold(1)/Long(2)/Flat(3) */ - float q_best = q_b0[argmax_n(q_b0, b0_size)]; - float q_flat = q_b0[flat_idx]; - float gap = q_best - q_flat; - out_q_gaps[i] = (gap > 0.0f) ? gap : 0.0f; + int flat_idx = 3; /* Flat = index 3 for direction(4): Short(0)/Hold(1)/Long(2)/Flat(3) */ + float q_sign_local = (contrarian_active != 0) ? -1.0f : 1.0f; + /* Find best action under the (possibly negated) Q. */ + int best_idx = 0; + float best_q = q_sign_local * q_b0[0]; + for (int a = 1; a < b0_size; a++) { + float qa = q_sign_local * q_b0[a]; + if (qa > best_q) { best_q = qa; best_idx = a; } } + float gap = best_q - (q_sign_local * q_b0[flat_idx]); + out_q_gaps[i] = (gap > 0.0f) ? gap : 0.0f; } } diff --git a/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs b/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs index 2e12e2a1f..484842d12 100644 --- a/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs +++ b/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs @@ -7742,7 +7742,7 @@ impl GpuDqnTrainer { let n = ring.param_count(); let num_bytes = n * std::mem::size_of::(); - let mut new_weights = self.stream.alloc_zeros::(n) + let mut new_weights = unsafe { self.stream.alloc::(n) } .map_err(|e| MLError::ModelError(format!("snapshot alloc: {e}")))?; { let (src_ptr, _sg) = self.params_buf.device_ptr(&self.stream); diff --git a/crates/ml/src/cuda_pipeline/q_snapshot.rs b/crates/ml/src/cuda_pipeline/q_snapshot.rs index a63d88441..fc91bdb38 100644 --- a/crates/ml/src/cuda_pipeline/q_snapshot.rs +++ b/crates/ml/src/cuda_pipeline/q_snapshot.rs @@ -88,7 +88,7 @@ impl SnapshotRing { let n = self.param_count; let num_bytes = n * std::mem::size_of::(); - let mut new_weights = self.stream.alloc_zeros::(n) + let mut new_weights = unsafe { self.stream.alloc::(n) } .map_err(|e| MLError::ModelError(format!("snapshot alloc: {e}")))?; { let (src_ptr, _src_guard) = params_src.device_ptr(&self.stream); diff --git a/pal_precommit.changeset b/pal_precommit.changeset new file mode 100644 index 000000000..09f607fd7 --- /dev/null +++ b/pal_precommit.changeset @@ -0,0 +1,2409 @@ +diff --git a/crates/ml-dqn/src/gpu_replay_buffer.rs b/crates/ml-dqn/src/gpu_replay_buffer.rs +index 620bb6466..d237dd7ce 100644 +--- a/crates/ml-dqn/src/gpu_replay_buffer.rs ++++ b/crates/ml-dqn/src/gpu_replay_buffer.rs +@@ -67,6 +67,8 @@ struct ReplayKernels { + per_insert_pa: CudaFunction, + per_prefix_scan: CudaFunction, + per_sample: CudaFunction, ++ /// C1/P1: Diversity-weighted priority kernel (health < 0.8 path). ++ pow_alpha_diverse_f32: CudaFunction, + } + + impl ReplayKernels { +@@ -108,6 +110,7 @@ impl ReplayKernels { + per_insert_pa: per_ld("per_insert_pa")?, + per_prefix_scan: per_ld("per_prefix_scan")?, + per_sample: per_ld("per_sample")?, ++ pow_alpha_diverse_f32: ld("pow_alpha_diverse_f32")?, + }) + } + } +@@ -196,6 +199,9 @@ pub struct GpuReplayBuffer { + trainer_dones_ptr: u64, + trainer_is_weights_ptr: u64, + trainer_state_dim_padded: usize, ++ /// C1/P1: Cached learning_health value from the trainer (updated once per epoch). ++ /// Controls whether priority updates use the standard or diversity-weighted kernel. ++ learning_health_cache: f32, + } + + impl Drop for GpuReplayBuffer { +@@ -323,6 +329,7 @@ impl GpuReplayBuffer { + trainer_dones_ptr: 0, + trainer_is_weights_ptr: 0, + trainer_state_dim_padded: 0, ++ learning_health_cache: 1.0, + }) + } + +@@ -350,6 +357,11 @@ impl GpuReplayBuffer { + pub const fn epsilon(&self) -> f32 { self.config.epsilon } + pub const fn state_dim(&self) -> usize { ml_core::state_layout::STATE_DIM } + ++ /// C1/P1: Update the cached learning_health value. Called from trainer at epoch boundary. ++ pub fn set_learning_health(&mut self, health: f32) { ++ self.learning_health_cache = health.clamp(0.0, 1.0); ++ } ++ + /// Wire trainer destination buffer pointers for direct-to-trainer gather. + /// Called once after GpuDqnTrainer is constructed. Pointers are stable + /// (CudaSlice allocations never move), so this is safe for graph capture. +@@ -818,21 +830,79 @@ impl GpuReplayBuffer { + + let (al, ep, bsi) = (self.config.alpha, self.config.epsilon, bs as i32); + let cap_i = self.config.capacity as i32; ++ let health = self.learning_health_cache; + + stream.memset_zeros(&mut self.update_batch_max) + .map_err(|e| MLError::ModelError(format!("zero batch_max: {e}")))?; + +- // per_update_pa: reads f32 td_errors directly, writes priorities + priorities_pa +- unsafe { +- stream.launch_builder(&self.kernels.per_update_pa) +- .arg(&self.priorities_pa) +- .arg(&self.priorities) +- .arg(&self.update_batch_max) +- .arg(indices) +- .arg(td_errors) +- .arg(&al).arg(&ep).arg(&cap_i).arg(&bsi) +- .launch(lcfg(bs)) +- .map_err(|e| MLError::ModelError(format!("per_update_pa: {e}")))?; ++ if health < 0.8 { ++ // C1/P1: during collapse, boost priorities of experiences whose action ++ // deviates from the batch mean. Requires a host-side mean readback to ++ // avoid a full reduction kernel — this path runs once per epoch at epoch boundary, ++ // so the DtoH sync cost is acceptable. ++ // ++ // Re-gather actions from the ring buffer using the current batch indices ++ // (sample_indices_i64 from the last sample_proportional call). This avoids ++ // depending on trainer_actions_ptr validity after graph replay. ++ let mbs = self.config.max_batch_size.max(1); ++ let gather_n = bs.min(mbs); ++ let gather_ni = gather_n as i32; ++ // Reuse sample_actions (CudaSlice) as scratch. gather_u32 gathers ++ // the action values at the sampled indices from the ring buffer. ++ let cap_i32 = self.config.capacity as i32; ++ unsafe { ++ stream.launch_builder(&self.kernels.gather_u32) ++ .arg(&mut self.sample_actions) ++ .arg(&self.actions) ++ .arg(&self.sample_indices_i64) ++ .arg(&gather_ni) ++ .arg(&cap_i32) ++ .launch(lcfg(gather_n)) ++ .map_err(|e| MLError::ModelError(format!("gather_u32 for diversity: {e}")))?; ++ } ++ // Synchronize to allow DtoH readback (epoch-boundary, not on hot path). ++ unsafe { ++ cudarc::driver::sys::cuStreamSynchronize(stream.cu_stream()); ++ } ++ let mut actions_host = vec![0u32; gather_n]; ++ stream.memcpy_dtoh(&self.sample_actions.slice(..gather_n), &mut actions_host) ++ .map_err(|e| MLError::ModelError(format!("actions dtoh diversity: {e}")))?; ++ let sum_actions: i64 = actions_host.iter().map(|&a| a as i64).sum(); ++ let mean_action_scaled: i32 = ((sum_actions * 1000) / (gather_n as i64).max(1)) as i32; ++ ++ // Reinterpret sample_actions (u32) as i32 for the kernel argument. ++ // The pointer type pun is safe: same device memory, same bit-width, kernel reads int*. ++ let actions_dev_ptr = self.sample_actions.raw_ptr() as u64; ++ unsafe { ++ stream.launch_builder(&self.kernels.pow_alpha_diverse_f32) ++ .arg(&self.priorities_pa) ++ .arg(&self.priorities) ++ .arg(&self.update_batch_max) ++ .arg(indices) ++ .arg(td_errors) ++ .arg(&actions_dev_ptr) ++ .arg(&mean_action_scaled) ++ .arg(&al) ++ .arg(&ep) ++ .arg(&health) ++ .arg(&cap_i) ++ .arg(&bsi) ++ .launch(lcfg(bs)) ++ .map_err(|e| MLError::ModelError(format!("pow_alpha_diverse_f32: {e}")))?; ++ } ++ } else { ++ // Standard PER (unchanged): reads f32 td_errors directly, writes priorities + priorities_pa ++ unsafe { ++ stream.launch_builder(&self.kernels.per_update_pa) ++ .arg(&self.priorities_pa) ++ .arg(&self.priorities) ++ .arg(&self.update_batch_max) ++ .arg(indices) ++ .arg(td_errors) ++ .arg(&al).arg(&ep).arg(&cap_i).arg(&bsi) ++ .launch(lcfg(bs)) ++ .map_err(|e| MLError::ModelError(format!("per_update_pa: {e}")))?; ++ } + } + + // Merge batch max into pending (unchanged logic) +diff --git a/crates/ml-dqn/src/replay_buffer_kernels.cu b/crates/ml-dqn/src/replay_buffer_kernels.cu +index 04800c927..9de823983 100644 +--- a/crates/ml-dqn/src/replay_buffer_kernels.cu ++++ b/crates/ml-dqn/src/replay_buffer_kernels.cu +@@ -471,6 +471,60 @@ void f32_idx_to_u32( + out[i] = (unsigned int)input[i]; + } + ++/* C1/P1: Health-weighted PER priority — boosts priorities of diverse-action ++ * experiences during collapse (health < 0.8). ++ * ++ * new_prio = clamp((|td_error[i]| + epsilon)^alpha ++ * × (1 + 2×(1−health) × |action[i] − mean_action|), ++ * epsilon, 1e6) ++ * priorities[indices[i]] = new_prio ++ * priorities_pa[indices[i]] = new_prio^alpha ++ * atomicMax(batch_max, new_prio) (IEEE 754 int trick) ++ * ++ * `mean_action_scaled` is the batch mean of actions multiplied by 1000 (passed ++ * as int to avoid precision issues when reducing host-side). ++ * ++ * Drop-in replacement for per_update_pa when health < 0.8. ++ */ ++extern "C" __global__ void pow_alpha_diverse_f32( ++ float* __restrict__ priorities_pa, ++ float* __restrict__ priorities, ++ float* __restrict__ batch_max, ++ const unsigned int* __restrict__ indices, ++ const float* __restrict__ td_errors, ++ const int* __restrict__ actions, ++ int mean_action_scaled, ++ float alpha, ++ float epsilon, ++ float health, ++ int capacity, ++ int batch_size ++) { ++ int i = blockIdx.x * blockDim.x + threadIdx.x; ++ if (i >= batch_size) return; ++ ++ unsigned int idx = indices[i]; ++ if (idx >= (unsigned int)capacity) return; ++ ++ float base = powf(fabsf(td_errors[i]) + epsilon, alpha); ++ ++ float mean_action = (float)mean_action_scaled * 0.001f; ++ float action_diff = fabsf((float)actions[i] - mean_action); ++ float diversity_mult = 1.0f + 2.0f * (1.0f - health) * action_diff; ++ ++ float new_prio = base * diversity_mult; ++ if (new_prio < epsilon) new_prio = epsilon; ++ if (new_prio > 1e6f) new_prio = 1e6f; ++ ++ priorities[idx] = new_prio; ++ ++ int ival = __float_as_int(new_prio); ++ atomicMax((int*)batch_max, ival); ++ ++ float pa = powf(new_prio, alpha); ++ priorities_pa[idx] = pa; ++} ++ + // ── 21. Regime weight combination kernel ──────────────────────────────────── + // out[i] = trending[i] * t_scale + ranging[i] * r_scale + volatile[i] * v_scale + +diff --git a/crates/ml/src/cuda_pipeline/c51_loss_kernel.cu b/crates/ml/src/cuda_pipeline/c51_loss_kernel.cu +index c486a2d77..1c6b16804 100644 +--- a/crates/ml/src/cuda_pipeline/c51_loss_kernel.cu ++++ b/crates/ml/src/cuda_pipeline/c51_loss_kernel.cu +@@ -161,6 +161,12 @@ __device__ void block_bellman_project_f( + } + } + ++/* B3/G4: Read learning_health from ISV buffer at index 12. ++ * Returns 0.5 if pointer is null (safe fallback). */ ++__device__ __forceinline__ float get_learning_health(const float* __restrict__ isv_signals_ptr) { ++ return (isv_signals_ptr != NULL) ? isv_signals_ptr[12] : 0.5f; ++} ++ + /* ══════════════════════════════════════════════════════════════════════ + * MAIN KERNEL: c51_loss_batched (float arithmetic, BF16 I/O) + * ══════════════════════════════════════════════════════════════════════ */ +@@ -232,7 +238,10 @@ extern "C" __global__ void c51_loss_batched( + const float* __restrict__ atom_positions, /* [4, num_atoms] adaptive positions. NULL = linear. */ + + /* ── Per-sample CVaR alpha from learned risk branch ── */ +- const float* __restrict__ cvar_alpha_buf /* [B] per-sample CVaR alpha from risk branch. NULL = use iqn_readiness. */ ++ const float* __restrict__ cvar_alpha_buf, /* [B] per-sample CVaR alpha from risk branch. NULL = use iqn_readiness. */ ++ ++ /* ── B3/G4: ISV signals for health-scaled Expected SARSA temperature ── */ ++ const float* __restrict__ isv_signals /* [ISV_DIM=13] pinned device-mapped. isv_signals[12] = learning_health. NULL = 0.5 fallback. */ + ) { + extern __shared__ float shmem_f[]; + +@@ -543,7 +552,14 @@ extern "C" __global__ void c51_loss_batched( + * meaningful at any Q-value scale. 1% of |mean Q| or 1e-6. */ + float mean_q = (max_eq + min_eq) * 0.5f; + float tau_floor = fmaxf(fabsf(mean_q) * 0.01f, 1e-6f); +- float tau = fmaxf(q_gap_local, tau_floor); ++ float tau_base = fmaxf(q_gap_local, tau_floor); ++ ++ // B3/G4: health-coupled temperature scaling. ++ // health=1 (healthy): factor=1.0 → sharp softmax (near argmax target, deterministic). ++ // health=0 (collapsed): factor=6.0 → wide softmax (stochastic sampling breaks attractor). ++ float health = get_learning_health(isv_signals); ++ float tau_health_factor = 1.0f + 5.0f * (1.0f - health); ++ float tau = tau_base * tau_health_factor; + float sum_exp = 0.0f; + for (int a = 0; a < n_d; a++) { + action_weights[a] = expf((eq_per_action[a] - max_eq) / tau); +diff --git a/crates/ml/src/cuda_pipeline/experience_kernels.cu b/crates/ml/src/cuda_pipeline/experience_kernels.cu +index df6a5c572..42051e5bb 100644 +--- a/crates/ml/src/cuda_pipeline/experience_kernels.cu ++++ b/crates/ml/src/cuda_pipeline/experience_kernels.cu +@@ -761,7 +761,8 @@ extern "C" __global__ void experience_action_select( + float eps_urg_mult, /* per-branch epsilon multiplier: urgency */ + int timestep, /* current timestep for stateless RNG */ + const float* __restrict__ per_sample_epsilon, /* [N] IQL expectile gap epsilon, NULL=use cosine schedule */ +- const float* __restrict__ isv_signals_ptr /* [8] pinned ISV signals for adaptive hold. NULL = static. */ ++ const float* __restrict__ isv_signals_ptr, /* [8] pinned ISV signals for adaptive hold. NULL = static. */ ++ int contrarian_active /* D7/N7: when non-zero, negate Q values before Boltzmann → argmin sampling */ + ) { + int i = blockIdx.x * blockDim.x + threadIdx.x; + if (i >= N) return; +@@ -804,6 +805,14 @@ extern "C" __global__ void experience_action_select( + const float* q_b2 = q_b1 + b1_size; + const float* q_b3 = q_b2 + b2_size; + ++ /* D7/N7: contrarian sign flip. ++ * When contrarian_active != 0, negate Q values before Boltzmann softmax. ++ * softmax(-q/tau) is argmin-favoring — the model picks its least-preferred ++ * action, breaking systematic anti-correlation attractors. ++ * When contrarian_active == 0, q_sign = +1.0f → bit-identical to previous ++ * behavior (identity multiplication). */ ++ float q_sign = (contrarian_active != 0) ? -1.0f : 1.0f; ++ + int dir_idx, mag_idx, a2, a3; + + /* Hold enforcement removed — replaced by cost-driven hold timing. +@@ -833,10 +842,10 @@ extern "C" __global__ void experience_action_select( + * Q-values differentiate (converges to argmax deterministically). + * When Q-values are flat, spreads evenly instead of flipping on noise. */ + /* Training mode: Boltzmann softmax over direction Q-values */ +- float q_max_d = (q_b0[0]); ++ float q_max_d = q_sign * (q_b0[0]); + float q_min_d = q_max_d; + for (int a = 1; a < b0_size; a++) { +- float qv = (q_b0[a]); ++ float qv = q_sign * (q_b0[a]); + q_max_d = fmaxf(q_max_d, qv); + q_min_d = fminf(q_min_d, qv); + } +@@ -845,7 +854,7 @@ extern "C" __global__ void experience_action_select( + float sum_e = 0.0f; + float exps_d[4]; /* b0_size=4: Short/Hold/Long/Flat */ + for (int a = 0; a < b0_size; a++) { +- float qv = (q_b0[a]); ++ float qv = q_sign * (q_b0[a]); + exps_d[a] = expf((qv - q_max_d) / tau_d); + sum_e += exps_d[a]; + } +@@ -889,10 +898,10 @@ extern "C" __global__ void experience_action_select( + } else { + /* Adaptive temperature: scale by Q-range so Boltzmann is meaningful + * regardless of absolute Q magnitude. Floor at 0.01 to avoid div/0. */ +- float q_max_m = (q_b1[0]); ++ float q_max_m = q_sign * (q_b1[0]); + float q_min_m = q_max_m; + for (int a = 1; a < b1_size; a++) { +- float qv = (q_b1[a]); ++ float qv = q_sign * (q_b1[a]); + q_max_m = fmaxf(q_max_m, qv); + q_min_m = fminf(q_min_m, qv); + } +@@ -906,7 +915,7 @@ extern "C" __global__ void experience_action_select( + float sum_e = 0.0f; + float exps[MAX_MAGNITUDE_ACTIONS]; + for (int a = 0; a < b1_size && a < MAX_MAGNITUDE_ACTIONS; a++) { +- float qv = (q_b1[a]); ++ float qv = q_sign * (q_b1[a]); + exps[a] = expf((qv - q_max_m) / tau); + sum_e += exps[a]; + } +@@ -925,10 +934,11 @@ extern "C" __global__ void experience_action_select( + a2 = (r >= b2_size) ? b2_size - 1 : r; + } else { + /* Boltzmann softmax over order Q-values */ +- float q_max_ord = q_b2[0], q_min_ord = q_b2[0]; ++ float q_max_ord = q_sign * q_b2[0], q_min_ord = q_max_ord; + for (int a = 1; a < b2_size; a++) { +- q_max_ord = fmaxf(q_max_ord, q_b2[a]); +- q_min_ord = fminf(q_min_ord, q_b2[a]); ++ float qv_ord = q_sign * q_b2[a]; ++ q_max_ord = fmaxf(q_max_ord, qv_ord); ++ q_min_ord = fminf(q_min_ord, qv_ord); + } + /* Higher temperature floor (0.5 vs 0.01) for order: Q-values for different + * order types start nearly identical (same reward signal). The low 0.01 floor +@@ -939,7 +949,7 @@ extern "C" __global__ void experience_action_select( + float sum_e = 0.0f; + float exps_ord[3]; + for (int a = 0; a < b2_size; a++) { +- exps_ord[a] = expf((q_b2[a] - q_max_ord) / tau_ord); ++ exps_ord[a] = expf((q_sign * q_b2[a] - q_max_ord) / tau_ord); + sum_e += exps_ord[a]; + } + float ro = philox_uniform(i, timestep, rng_ctr++) * sum_e; +@@ -957,16 +967,17 @@ extern "C" __global__ void experience_action_select( + a3 = (r >= b3_size) ? b3_size - 1 : r; + } else { + /* Boltzmann softmax over urgency Q-values */ +- float q_max_urg = q_b3[0], q_min_urg = q_b3[0]; ++ float q_max_urg = q_sign * q_b3[0], q_min_urg = q_max_urg; + for (int a = 1; a < b3_size; a++) { +- q_max_urg = fmaxf(q_max_urg, q_b3[a]); +- q_min_urg = fminf(q_min_urg, q_b3[a]); ++ float qv_urg = q_sign * q_b3[a]; ++ q_max_urg = fmaxf(q_max_urg, qv_urg); ++ q_min_urg = fminf(q_min_urg, qv_urg); + } + float tau_urg = fmaxf(q_max_urg - q_min_urg, 0.5f); + float sum_e = 0.0f; + float exps_urg[3]; + for (int a = 0; a < b3_size; a++) { +- exps_urg[a] = expf((q_b3[a] - q_max_urg) / tau_urg); ++ exps_urg[a] = expf((q_sign * q_b3[a] - q_max_urg) / tau_urg); + sum_e += exps_urg[a]; + } + float ru = philox_uniform(i, timestep, rng_ctr++) * sum_e; +@@ -999,13 +1010,22 @@ extern "C" __global__ void experience_action_select( + + /* Output Q-gap for conviction-based position sizing. + * Q-gap = Q(best_direction) - Q(Flat). Higher = more conviction. +- * Used by env_step to scale target_position. */ ++ * Used by env_step to scale target_position. ++ * ++ * D7/N7: During contrarian override, the action selected is argmin rather ++ * than argmax. The raw Q(best) - Q(Flat) gap would mislead the position ++ * sizer into full conviction for what's actually the "deliberately wrong" ++ * action. Zero out the gap so contrarian trades size minimal. */ + if (out_q_gaps != NULL) { +- int flat_idx = 3; /* Flat = index 3 for direction(4): Short(0)/Hold(1)/Long(2)/Flat(3) */ +- float q_best = q_b0[argmax_n(q_b0, b0_size)]; +- float q_flat = q_b0[flat_idx]; +- float gap = q_best - q_flat; +- out_q_gaps[i] = (gap > 0.0f) ? gap : 0.0f; ++ if (contrarian_active != 0) { ++ out_q_gaps[i] = 0.0f; ++ } else { ++ int flat_idx = 3; /* Flat = index 3 for direction(4): Short(0)/Hold(1)/Long(2)/Flat(3) */ ++ float q_best = q_b0[argmax_n(q_b0, b0_size)]; ++ float q_flat = q_b0[flat_idx]; ++ float gap = q_best - q_flat; ++ out_q_gaps[i] = (gap > 0.0f) ? gap : 0.0f; ++ } + } + } + +@@ -1111,7 +1131,8 @@ extern "C" __global__ void experience_env_step( + float price_confirm_weight, /* micro-reward: price confirmation weight (default 0.5) */ + float book_aggression_weight, /* micro-reward: book aggression weight (default 0.3) */ + float hold_quality_weight, /* micro-reward: retrospective hold quality weight (default 0.2) */ +- float micro_reward_temp /* micro-reward: tanh temperature (default 3.0) */ ++ float micro_reward_temp, /* micro-reward: tanh temperature (default 3.0) */ ++ float cf_ratio /* D4/N4: adaptive counterfactual flip ratio [0,1]; 0.5=healthy, 0.8=collapsed */ + ) { + int i = blockIdx.x * blockDim.x + threadIdx.x; + if (i >= N) return; +@@ -1158,7 +1179,7 @@ extern "C" __global__ void experience_env_step( + int do_flip = 0; + { + float flip_u = philox_uniform(i, bar_idx + 8888, 7777); +- do_flip = (flip_u < 0.5f) ? 1 : 0; ++ do_flip = (flip_u < cf_ratio) ? 1 : 0; + } + if (do_flip) { + for (int f = 0; f < 16 && f < state_dim; f++) { +diff --git a/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs b/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs +index bc6fedef0..2e12e2a1f 100644 +--- a/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs ++++ b/crates/ml/src/cuda_pipeline/gpu_dqn_trainer.rs +@@ -83,7 +83,8 @@ const OFI_EMBED_DIM: usize = 10; // OFI embedding width appended to history + + /// Introspective State Vector (ISV) configuration. + const ISV_K: usize = 4; // Temporal ISV history length +-const ISV_DIM: usize = 12; // ISV signal count (8 core + 4 regime awareness) ++const ISV_DIM: usize = 13; // ISV signal count (8 core + 4 regime awareness + 1 learning_health) ++pub const LEARNING_HEALTH_INDEX: usize = 12; + const ISV_EMB_DIM: usize = 8; // ISV embedding output dimension (FC2 output, gate/gamma input) + /// First ISV tensor index in the flat param buffer. + /// ISV weights (68-79) are online-only — NOT synced to the target network. +@@ -1562,11 +1563,11 @@ pub struct GpuDqnTrainer { + atom_util_scratch_dev_ptr: u64, + + // ── ISV core buffers (pinned device-mapped, GPU read/write) ── +- isv_signals_pinned: *mut f32, // [ISV_DIM = 12] ++ isv_signals_pinned: *mut f32, // [ISV_DIM = 13] + isv_signals_dev_ptr: u64, +- isv_history_pinned: *mut f32, // [ISV_K * ISV_DIM = 48] ++ isv_history_pinned: *mut f32, // [ISV_K * ISV_DIM = 52] + isv_history_dev_ptr: u64, +- isv_decay_pinned: *mut f32, // [ISV_DIM = 12] learned decay weights ++ isv_decay_pinned: *mut f32, // [ISV_DIM = 13] learned decay weights + isv_decay_dev_ptr: u64, + lagged_td_error_pinned: *mut f32, // [1] recursive confidence target + lagged_td_error_dev_ptr: u64, +@@ -1658,6 +1659,41 @@ pub struct GpuDqnTrainer { + tau_final: f32, + /// Total number of steps over which tau is cosine-annealed. + tau_anneal_steps: i32, ++ ++ /// Last effective cql_alpha applied (base × (1−regime_stability) × health). ++ /// Populated by apply_cql_gradient. Initialized to 0.0. ++ pub(crate) last_cql_alpha_eff: f32, ++ ++ /// Last effective tau value (cosine-annealed, potentially floored by 0.01×(1-health)). ++ /// Populated by `apply_health_coupled_tau_floor`. Initialized to 0.0. ++ pub(crate) last_tau_eff: f32, ++ ++ /// B3/G4: Last health-scaled Expected SARSA tau factor (1.0 at full health, 6.0 at collapse). ++ /// Populated by `launch_c51_loss`. Initialized to 0.0. ++ pub(crate) last_sarsa_tau_factor: f32, ++ ++ /// B4/G5: Last adaptive gradient budget for IQN (0.10 + 0.30×health). ++ /// Populated by `FusedTrainingCtx::compute_adaptive_budgets`. Default: health=1 → 0.40. ++ pub(crate) last_iqn_budget_eff: f32, ++ /// B4/G5: Last adaptive gradient budget for CQL (0.10×(1−regime)×health). ++ /// Populated by `FusedTrainingCtx::compute_adaptive_budgets`. Default: 0.00 (regime_stability=1). ++ pub(crate) last_cql_budget_eff: f32, ++ /// B4/G5: Last adaptive gradient budget for C51 (1−iqn−cql−ens). ++ /// Populated by `FusedTrainingCtx::compute_adaptive_budgets`. Default: health=1 → 0.55. ++ pub(crate) last_c51_budget_eff: f32, ++ /// B4/G5: Adaptive gradient budget for ensemble (constant 0.05). ++ /// Populated by `FusedTrainingCtx::compute_adaptive_budgets`. ++ pub(crate) last_ens_budget_eff: f32, ++ ++ /// C4/P4: Last effective gamma applied to C51 Bellman projection. ++ /// Formula: gamma_base + 0.005×(regime_stability − 0.5) − 0.05×(1 − health), ++ /// clamped to [0.9, 0.995]. Populated by `apply_adaptive_gamma`. Default: 0.99. ++ pub(crate) last_gamma_eff: f32, ++ ++ /// D1/N1: Ring buffer of best-health weight snapshots for temporal self-distillation. ++ pub(crate) q_snapshots: crate::cuda_pipeline::q_snapshot::SnapshotRing, ++ /// D1/N1: Whether distillation gradient was applied in the most recent epoch boundary. ++ pub(crate) last_distill_active: bool, + } + + impl GpuDqnTrainer { +@@ -4051,6 +4087,7 @@ impl GpuDqnTrainer { + &mut self, + iqn_d_h_s2_ptr: u64, + _online_dueling: &mut DuelingWeightSet, ++ iqn_budget: f32, + ) -> Result<(), MLError> { + let b = self.config.batch_size; + let _eg = EventTrackingGuard::new(self.stream.context()); +@@ -4171,13 +4208,14 @@ impl GpuDqnTrainer { + )?; + } + +- // ── 7. Plain SAXPY: grad_buf[trunk] += iqn_lambda * scratch ── ++ // ── 7. Plain SAXPY: grad_buf[trunk] += iqn_lambda * iqn_budget * scratch ── + // Adaptive lambda: scales with IQN loss readiness (0 when uncertain, 1 when converged). ++ // B4/G5: further scaled by iqn_budget (0.10+0.30×health). + // Uses aux_child-specific handle — saxpy_f32_kernel is captured in forward_child. + { + let scratch_ptr = self.ptrs.iqn_trunk_m; + let grad_ptr = self.ptrs.grad_buf; +- let scale = self.config.iqn_lambda * self.iqn_readiness; ++ let scale = self.config.iqn_lambda * self.iqn_readiness * iqn_budget; + let n_i32 = trunk_grad_total as i32; + let blocks = ((trunk_grad_total + 255) / 256) as u32; + unsafe { +@@ -4573,7 +4611,25 @@ impl GpuDqnTrainer { + let d_v_ptr = self.cql_d_value_logits.raw_ptr(); + let d_adv_ptr = self.cql_d_adv_logits.raw_ptr(); + +- let cql_alpha = self.config.cql_alpha; ++ // B1/G2: Adaptive cql_alpha. ++ // - health from ISV[LEARNING_HEALTH_INDEX] (0.5 fallback if pinned ptr null) ++ // - regime_stability from ISV[11] ++ // - cql_alpha_eff = base × (1 − regime_stability) × health ++ // Collapse → 0 (CQL off); volatile+healthy → full; stable+healthy → 0. ++ let (health, regime_stability) = if self.isv_signals_pinned.is_null() { ++ (0.5_f32, 0.5_f32) ++ } else { ++ unsafe { ++ let h = (*self.isv_signals_pinned.add(LEARNING_HEALTH_INDEX)).clamp(0.0, 1.0); ++ let s = (*self.isv_signals_pinned.add(11)).clamp(0.0, 1.0); ++ (h, s) ++ } ++ }; ++ let cql_alpha = self.config.cql_alpha * (1.0 - regime_stability) * health; ++ ++ // Cache for logging via FusedTrainingCtx / DQNTrainer. ++ self.last_cql_alpha_eff = cql_alpha; ++ + let n_i32 = b as i32; + let na_i32 = na as i32; + let b0_i32 = b0 as i32; +@@ -4740,14 +4796,15 @@ impl GpuDqnTrainer { + /// + /// Called after `apply_cql_gradient` populated `cql_grad_scratch`. + /// No per-component clip — single global clip in Adam handles safety. +- pub fn apply_cql_saxpy(&mut self) -> Result<(), MLError> { ++ pub fn apply_cql_saxpy(&mut self, cql_budget: f32) -> Result<(), MLError> { + let total = self.total_params as i32; + let blocks = self.grad_norm_blocks as u32; + let scratch_ptr = self.ptrs.cql_grad_scratch; + let grad_ptr = self.ptrs.grad_buf; +- let alpha = 1.0_f32; ++ // B4/G5: scale CQL gradient contribution by adaptive budget (0.10×(1−regime)×health). ++ let alpha = cql_budget; + +- // Plain SAXPY: grad_buf += 1.0 * cql_scratch ++ // Plain SAXPY: grad_buf += cql_budget * cql_scratch + // Uses aux_child-specific handle — saxpy_f32_kernel is captured in forward_child. + unsafe { + self.stream +@@ -6924,6 +6981,9 @@ impl GpuDqnTrainer { + let tau_final = 0.001_f32; + let tau_anneal_steps: i32 = 100_000; + ++ // D1/N1: clone Arc before it is moved into the struct literal. ++ let stream_for_snapshots = Arc::clone(&stream); ++ + Ok(Self { + config, + stream, +@@ -7418,6 +7478,19 @@ impl GpuDqnTrainer { + tau_init, + tau_final, + tau_anneal_steps, ++ last_cql_alpha_eff: 0.0, ++ last_tau_eff: 0.0, ++ last_sarsa_tau_factor: 0.0, ++ last_iqn_budget_eff: 0.40, ++ last_cql_budget_eff: 0.00, ++ last_c51_budget_eff: 0.55, ++ last_ens_budget_eff: 0.05, ++ last_gamma_eff: 0.99, ++ q_snapshots: crate::cuda_pipeline::q_snapshot::SnapshotRing::new( ++ stream_for_snapshots, ++ total_params, ++ ), ++ last_distill_active: false, + }) + } + +@@ -7599,6 +7672,223 @@ impl GpuDqnTrainer { + } + } + ++ /// B4/G5: Read ISV health index and regime stability from pinned host memory. ++ /// Returns (health [0,1], regime_stability [0,1]). Falls back to (0.5, 0.5) if ++ /// pinned pointer is null. ++ pub(crate) fn read_isv_health_and_regime(&self) -> (f32, f32) { ++ if self.isv_signals_pinned.is_null() { return (0.5_f32, 0.5_f32); } ++ unsafe { ++ let h = (*self.isv_signals_pinned.add(LEARNING_HEALTH_INDEX)).clamp(0.0, 1.0); ++ let s = (*self.isv_signals_pinned.add(11)).clamp(0.0, 1.0); ++ (h, s) ++ } ++ } ++ ++ /// B2/G3: Apply a health-coupled minimum to a scheduled tau value. ++ /// Returns `max(tau_scheduled, 0.01 * (1.0 - health))`. During collapse ++ /// (health≈0), the floor rises to 0.01 to accelerate target-network adaptation. ++ /// Caches the effective value for HEALTH_DIAG logging. ++ pub fn apply_health_coupled_tau_floor(&mut self, tau_scheduled: f64) -> f64 { ++ let health = if self.isv_signals_pinned.is_null() { ++ 0.5_f32 ++ } else { ++ unsafe { (*self.isv_signals_pinned.add(LEARNING_HEALTH_INDEX)).clamp(0.0, 1.0) } ++ }; ++ let tau_floor = (0.01 * (1.0 - health)) as f64; ++ let tau_eff = tau_scheduled.max(tau_floor); ++ self.last_tau_eff = tau_eff as f32; ++ tau_eff ++ } ++ ++ /// C4/P4: Apply temporal-coupled adjustment to the scheduled gamma. ++ /// Formula: gamma_base + 0.005×(regime_stability − 0.5) − 0.05×(1 − health), ++ /// clamped to [0.9, 0.995]. Updates `adaptive_gamma` and caches the effective ++ /// value in `last_gamma_eff` for HEALTH_DIAG logging. ++ pub fn apply_adaptive_gamma(&mut self, gamma_base: f32) -> f32 { ++ let (health, regime_stability) = self.read_isv_health_and_regime(); ++ let gamma_eff = gamma_base ++ + 0.005_f32 * (regime_stability - 0.5) ++ - 0.05_f32 * (1.0 - health); ++ let gamma_clamped = gamma_eff.clamp(0.9, 0.995); ++ self.adaptive_gamma = gamma_clamped; ++ self.last_gamma_eff = gamma_clamped; ++ gamma_clamped ++ } ++ ++ // ── D1/N1: Temporal self-distillation ───────────────────────────────── ++ ++ /// D1/N1: Whether distillation gradient was applied in the most recent epoch. ++ pub fn last_distill_active(&self) -> bool { ++ self.last_distill_active ++ } ++ ++ /// D1/N1: Snapshot current params_buf if health and q_gap qualify. ++ /// Avoids split-borrow between params_buf and q_snapshots. ++ pub fn maybe_snapshot_params(&mut self, health: f32, q_gap: f32, epoch: u32) -> Result { ++ use cudarc::driver::{DevicePtr, DevicePtrMut}; ++ use crate::cuda_pipeline::q_snapshot::SNAPSHOT_HEALTH_THRESHOLD; ++ ++ if health < SNAPSHOT_HEALTH_THRESHOLD { ++ return Ok(false); ++ } ++ ++ let ring = &mut self.q_snapshots; ++ if ring.snapshots.len() >= crate::cuda_pipeline::q_snapshot::MAX_SNAPSHOTS { ++ let worst_q_gap = ring.snapshots.iter().map(|s| s.q_gap).fold(f32::INFINITY, f32::min); ++ if q_gap <= worst_q_gap { ++ return Ok(false); ++ } ++ } ++ ++ let n = ring.param_count(); ++ let num_bytes = n * std::mem::size_of::(); ++ let mut new_weights = self.stream.alloc_zeros::(n) ++ .map_err(|e| MLError::ModelError(format!("snapshot alloc: {e}")))?; ++ { ++ let (src_ptr, _sg) = self.params_buf.device_ptr(&self.stream); ++ let (dst_ptr, _dg) = new_weights.device_ptr_mut(&self.stream); ++ #[allow(unsafe_code)] ++ unsafe { ++ cudarc::driver::result::memcpy_dtod_async( ++ dst_ptr, src_ptr, num_bytes, self.stream.cu_stream(), ++ ).map_err(|e| MLError::ModelError(format!("snapshot dtod: {e}")))?; ++ } ++ } ++ ++ let ring = &mut self.q_snapshots; ++ if ring.snapshots.len() >= crate::cuda_pipeline::q_snapshot::MAX_SNAPSHOTS { ++ let worst_idx = ring.snapshots.iter() ++ .enumerate() ++ .min_by(|(_, a), (_, b)| a.q_gap.partial_cmp(&b.q_gap).unwrap()) ++ .map(|(i, _)| i) ++ .unwrap(); ++ ring.snapshots.swap_remove(worst_idx); ++ } ++ ++ ring.snapshots.push(crate::cuda_pipeline::q_snapshot::QSnapshot { ++ weights: new_weights, ++ health, ++ q_gap, ++ epoch, ++ }); ++ Ok(true) ++ } ++ ++ /// D1/N1: Apply distillation gradient — pulls current weights toward the best ++ /// snapshot with strength `0.1 × (1 − health)`. Returns true if applied. ++ /// ++ /// Implemented as two SAXPY calls into grad_buf (ungraphed, stream-ordered): ++ /// grad += alpha * params (adds current) ++ /// grad += -alpha * best_snap (subtracts best) ++ /// Net effect: grad += alpha * (params - best_snap) ++ /// ++ /// Uses `saxpy_f32_aux` (dqn_saxpy_f32_kernel, aux_child handle) — same handle ++ /// used by `apply_cql_saxpy`, safe for ungraphed injection before graph_adam. ++ pub fn apply_distillation_gradient(&mut self) -> Result { ++ let (health, _) = self.read_isv_health_and_regime(); ++ let distill_weight = 0.1_f32 * (1.0 - health).max(0.0); ++ if distill_weight < 0.01 { ++ self.last_distill_active = false; ++ return Ok(false); ++ } ++ let best_ptr = match self.q_snapshots.best_raw_ptr() { ++ Some(p) => p, ++ None => { ++ self.last_distill_active = false; ++ return Ok(false); ++ } ++ }; ++ ++ let grad_ptr = self.ptrs.grad_buf; ++ let params_ptr = self.ptrs.params_ptr; ++ let n = self.total_params as i32; ++ let blocks = self.grad_norm_blocks as u32; ++ ++ unsafe { ++ // grad += distill_weight * params_buf ++ self.stream ++ .launch_builder(&self.saxpy_f32_aux) ++ .arg(&grad_ptr) ++ .arg(¶ms_ptr) ++ .arg(&distill_weight) ++ .arg(&n) ++ .launch(LaunchConfig { ++ grid_dim: (blocks, 1, 1), ++ block_dim: (256, 1, 1), ++ shared_mem_bytes: 0, ++ }) ++ .map_err(|e| MLError::ModelError(format!("distill saxpy (+params): {e}")))?; ++ ++ // grad += (-distill_weight) * best_snapshot ++ let neg_weight = -distill_weight; ++ self.stream ++ .launch_builder(&self.saxpy_f32_aux) ++ .arg(&grad_ptr) ++ .arg(&best_ptr) ++ .arg(&neg_weight) ++ .arg(&n) ++ .launch(LaunchConfig { ++ grid_dim: (blocks, 1, 1), ++ block_dim: (256, 1, 1), ++ shared_mem_bytes: 0, ++ }) ++ .map_err(|e| MLError::ModelError(format!("distill saxpy (-best): {e}")))?; ++ } ++ ++ self.last_distill_active = true; ++ Ok(true) ++ } ++ ++ // ── A3: LearningHealth signal writers / readers ─────────────────────── ++ ++ /// Write a scalar into the ISV pinned buffer at the given index. ++ /// No-op if the pinned pointer is null or index >= ISV_DIM. ++ pub fn write_isv_signal_at(&self, index: usize, value: f32) { ++ if self.isv_signals_pinned.is_null() { return; } ++ if index >= ISV_DIM { return; } ++ unsafe { *self.isv_signals_pinned.add(index) = value; } ++ } ++ ++ /// Read atom utilization from q_readback_pinned[6] (one-step lag, written by ++ /// reduce_current_q_stats). Returns 0.0 if pointer is null. ++ pub fn read_atom_utilization(&self) -> f32 { ++ if self.q_readback_pinned.is_null() { return 0.0; } ++ unsafe { (*self.q_readback_pinned.add(6)).clamp(0.0, 1.0) } ++ } ++ ++ /// Coarse Q-collapse detector derived from the single-sample Q readback at ++ /// `q_readback_pinned[7..7+total_actions]`. ++ /// ++ /// Returns a LARGE value (~100) when all Q values are (nearly) equal — i.e. ++ /// rank-1 / Q-collapse — and a value close to 1.0 when the Q spread is wide. ++ /// Intentionally coarse: a full SVD across a batch of Q samples would require ++ /// a separate larger buffer and is out of scope for the host-side health path. ++ /// ++ /// The downstream `NormalizedComponents::from_raw` thresholds are calibrated ++ /// for this behaviour (`spectral_gap_norm` drops toward 0 as this value rises). ++ pub fn compute_q_spectral_gap(&self) -> f32 { ++ if self.q_readback_pinned.is_null() { return 1.0; } ++ // Layout: q_readback_pinned[0..7] = stats, [7..7+total_actions] = per-action Q values. ++ let total_actions = self.total_actions() as usize; ++ if total_actions == 0 { return 1.0; } ++ unsafe { ++ let base = self.q_readback_pinned.add(7); ++ let mut max = f32::NEG_INFINITY; ++ let mut min = f32::INFINITY; ++ for i in 0..total_actions { ++ let v = *base.add(i); ++ if v.is_finite() { ++ if v > max { max = v; } ++ if v < min { min = v; } ++ } ++ } ++ let range = (max - min).abs(); ++ // Range near zero means all actions have equal Q → rank-1-like collapse; ++ // return a large value to signal pathology via spectral_gap_norm ≈ 0. ++ if range < 1e-6 { 100.0 } else { 1.0 + range.recip() } ++ } ++ } ++ + /// Read eval_q_std_ema. + pub fn eval_q_std_ema(&self) -> f32 { self.eval_q_std_ema } + +@@ -10608,7 +10898,7 @@ impl GpuDqnTrainer { + /// Reads: on_v_logits, on_b_logits, tg_v_logits, tg_b_logits, + /// on_next_v_logits, on_next_b_logits (Double DQN). + /// Writes: per_sample_loss, td_errors, total_loss, save_current_lp, save_projected. +- fn launch_c51_loss(&self) -> Result<(), MLError> { ++ fn launch_c51_loss(&mut self) -> Result<(), MLError> { + let kernel = &self.c51_loss_kernel; + + let b = self.config.batch_size; +@@ -10674,6 +10964,14 @@ impl GpuDqnTrainer { + let shmem_floats = na + na + max_branch * na + na + na + na + 8; + let shmem_bytes = (shmem_floats * std::mem::size_of::()) as u32; + ++ // B3/G4: mirror the kernel's tau_health_factor on host for logging. ++ let current_health = if self.isv_signals_pinned.is_null() { ++ 0.5_f32 ++ } else { ++ unsafe { (*self.isv_signals_pinned.add(LEARNING_HEALTH_INDEX)).clamp(0.0, 1.0) } ++ }; ++ self.last_sarsa_tau_factor = 1.0 + 5.0 * (1.0 - current_health); ++ + unsafe { + self.stream + .launch_builder(kernel) +@@ -10737,6 +11035,8 @@ impl GpuDqnTrainer { + .arg(&atom_positions_buf_ptr) + // ── Per-sample CVaR alpha from learned risk branch ── + .arg(&cvar_alpha_buf_ptr) ++ // ── B3/G4: ISV signals for health-scaled Expected SARSA temperature ── ++ .arg(&self.isv_signals_dev_ptr) + .launch(LaunchConfig { + grid_dim: (b as u32, 1, 1), + block_dim: (256, 1, 1), +diff --git a/crates/ml/src/cuda_pipeline/gpu_experience_collector.rs b/crates/ml/src/cuda_pipeline/gpu_experience_collector.rs +index 01a423e26..c3bc6126f 100644 +--- a/crates/ml/src/cuda_pipeline/gpu_experience_collector.rs ++++ b/crates/ml/src/cuda_pipeline/gpu_experience_collector.rs +@@ -732,6 +732,17 @@ pub struct GpuExperienceCollector { + hindsight_relabel_kernel: CudaFunction, + /// v8: TD(lambda) return computation kernel (loaded from nstep cubin). + td_lambda_kernel: CudaFunction, ++ ++ /// D4/N4: cached learning_health value from the trainer (updated once per epoch). ++ learning_health_cache: f32, ++ /// D4/N4: last effective cf_ratio computed from learning_health_cache. ++ last_cf_ratio_eff: f32, ++ ++ /// D7/N7: Contrarian override flag cached from trainer. When non-zero, the ++ /// experience kernel negates Q values before Boltzmann softmax, converting ++ /// argmax-favoring sampling to argmin-favoring. Used briefly to escape Q-uniform ++ /// attractors when the policy is systematically anti-correlated with market. ++ contrarian_active_cache: u8, + } + + impl Drop for GpuExperienceCollector { +@@ -1298,6 +1309,9 @@ impl GpuExperienceCollector { + difficulty_scores_kernel, + hindsight_relabel_kernel, + td_lambda_kernel, ++ learning_health_cache: 1.0, // D4/N4: assume healthy at construction ++ last_cf_ratio_eff: 0.5, // D4/N4: standard cf_ratio at healthy state ++ contrarian_active_cache: 0, // D7/N7: off by default + reward_rank_kernel, + reward_compute_abs_sharpe_kernel, + bitonic_sort_step_kernel, +@@ -1340,6 +1354,34 @@ impl GpuExperienceCollector { + self.readiness_dev_ptr = dev_ptr; + } + ++ /// D4/N4: propagate learning health from the trainer (called once per epoch). ++ /// The next call to collect_experiences_gpu will compute cf_ratio from this value. ++ pub fn set_learning_health(&mut self, health: f32) { ++ self.learning_health_cache = health.clamp(0.0, 1.0); ++ } ++ ++ /// D4/N4: return the last effective cf_ratio (set during collect_experiences_gpu). ++ pub fn last_cf_ratio_eff(&self) -> f32 { ++ self.last_cf_ratio_eff ++ } ++ ++ /// D7/N7: propagate contrarian flag from trainer (called once per epoch when ++ /// the state machine updates). When active, the action-selection kernel negates ++ /// Q values before Boltzmann softmax → argmin-favoring sampling. ++ pub fn set_contrarian_active(&mut self, active: bool) { ++ self.contrarian_active_cache = if active { 1 } else { 0 }; ++ } ++ ++ /// D7/N7: return the cached contrarian flag. ++ pub fn contrarian_active(&self) -> u8 { ++ self.contrarian_active_cache ++ } ++ ++ /// D4/N4: compute cf_ratio from cached health. ++ fn current_cf_ratio(&self) -> f32 { ++ (0.5_f32 + 0.3_f32 * (1.0 - self.learning_health_cache)).clamp(0.0, 1.0) ++ } ++ + /// Set CVaR position scaling from IQN dual-head. + /// The device pointer will be passed to the env_step kernel. + /// Call with 0 to disable (NULL pointer = no scaling). +@@ -2348,6 +2390,7 @@ impl GpuExperienceCollector { + .arg(&(t as i32)) // timestep for stateless Philox RNG + .arg(&self.per_sample_epsilon_ptr) // IQL expectile gap epsilon (0=cosine schedule) + .arg(&self.isv_signals_dev_ptr) // ISV signals for adaptive hold (0=NULL=static) ++ .arg(&(self.contrarian_active_cache as i32)) // D7/N7: Q-negation flag + .launch(launch_cfg) + .map_err(|e| MLError::ModelError(format!( + "experience_action_select t={t}: {e}" +@@ -2386,6 +2429,9 @@ impl GpuExperienceCollector { + + // ── 5. Environment step (reward v5: trade-aware hybrid) ────── + // max_pos already defined above (action_select block) ++ // D4/N4: adaptive counterfactual ratio — 0.5 at healthy, 0.8 at collapse. ++ let cf_ratio = self.current_cf_ratio(); ++ self.last_cf_ratio_eff = cf_ratio; + let tx_cost = config.tx_cost_multiplier; + let l_i32 = timesteps as i32; + let total_actions_env = (self.branch_sizes[0] + self.branch_sizes[1] +@@ -2467,6 +2513,7 @@ impl GpuExperienceCollector { + .arg(&config.book_aggression_weight) // micro-reward: book aggression weight + .arg(&config.hold_quality_weight) // micro-reward: hold quality weight + .arg(&config.micro_reward_temp) // micro-reward: tanh temperature ++ .arg(&cf_ratio) // D4/N4: adaptive counterfactual flip ratio + .launch(launch_cfg) + .map_err(|e| MLError::ModelError(format!( + "experience_env_step t={t}: {e}" +diff --git a/crates/ml/src/cuda_pipeline/learning_health.rs b/crates/ml/src/cuda_pipeline/learning_health.rs +new file mode 100644 +index 000000000..67389d251 +--- /dev/null ++++ b/crates/ml/src/cuda_pipeline/learning_health.rs +@@ -0,0 +1,186 @@ ++//! LearningHealth signal — unified training health metric. ++//! ++//! Composes 7 normalized components into a scalar [0, 1] that senses training ++//! collapse and drives adaptive hyperparameters. ++ ++/// Smoothstep function: smooth transition from 0 at `edge0` to 1 at `edge1`. ++/// Standard GLSL smoothstep formula. ++#[inline] ++pub fn smoothstep(edge0: f32, edge1: f32, x: f32) -> f32 { ++ if edge1 == edge0 { return if x >= edge0 { 1.0 } else { 0.0 }; } ++ let t = ((x - edge0) / (edge1 - edge0)).clamp(0.0, 1.0); ++ t * t * (3.0 - 2.0 * t) ++} ++ ++/// Raw input signals for LearningHealth composition. ++#[derive(Debug, Clone, Copy, Default)] ++pub struct HealthComponents { ++ pub q_gap: f32, // max - 2nd max of Q values across actions ++ pub q_var: f32, // variance of Q across actions ++ pub atom_util: f32, // C51 atom utilization fraction [0, 1] ++ pub grad_norm: f32, // gradient L2 norm ++ pub ens_disagreement: f32, // mean pairwise diff across ensemble heads ++ pub grad_consistency: f32, // cosine similarity of successive gradients ++ pub spectral_gap: f32, // sigma_1 / sigma_2 from SVD of Q matrix ++} ++ ++/// Computed normalized signals in [0, 1]. ++#[derive(Debug, Clone, Copy, Default)] ++pub struct NormalizedComponents { ++ pub q_gap_norm: f32, ++ pub q_var_norm: f32, ++ pub atom_util_norm: f32, ++ pub grad_stable: f32, ++ pub ens_agree: f32, ++ pub grad_consistency_norm: f32, ++ pub spectral_gap_norm: f32, ++} ++ ++impl NormalizedComponents { ++ /// Normalize raw signals into [0, 1] via smoothstep curves. Thresholds encode ++ /// empirical collapse/health boundaries observed in production runs: ++ /// - q_gap: [0.01, 0.5] — below 0.01 the argmax is indistinguishable from noise; above 0.5 actions are well-separated ++ /// - q_var: [0.001, 0.1] — tight variance indicates Q-collapse to uniform ++ /// - atom_util: [0.2, 0.7] — C51 healthy when ≥ 70% of atoms carry >1% probability ++ /// - grad_norm: [10, 100] — inverted; higher grad_norm is less stable ++ /// - ens_disagreement: directly used; higher disagreement means less agreement on Q ++ /// - grad_consistency: [-0.2, 0.5] — cosine sim between successive gradient vectors ++ /// - spectral_gap: [2, 10] — inverted; sigma_1/sigma_2 above 10 means rank-1 collapse ++ pub fn from_raw(raw: &HealthComponents) -> Self { ++ Self { ++ q_gap_norm: smoothstep(0.01, 0.5, raw.q_gap), ++ q_var_norm: smoothstep(0.001, 0.1, raw.q_var), ++ atom_util_norm: smoothstep(0.2, 0.7, raw.atom_util), ++ grad_stable: 1.0 - smoothstep(10.0, 100.0, raw.grad_norm), ++ ens_agree: (1.0 - raw.ens_disagreement).clamp(0.0, 1.0), ++ grad_consistency_norm: smoothstep(-0.2, 0.5, raw.grad_consistency), ++ spectral_gap_norm: 1.0 - smoothstep(2.0, 10.0, raw.spectral_gap), ++ } ++ } ++ ++ /// Weighted sum per spec Layer 1 composition. ++ pub fn compose(&self) -> f32 { ++ 0.25 * self.q_gap_norm ++ + 0.15 * self.q_var_norm ++ + 0.15 * self.atom_util_norm ++ + 0.15 * self.grad_stable ++ + 0.10 * self.ens_agree ++ + 0.10 * self.grad_consistency_norm ++ + 0.10 * self.spectral_gap_norm ++ } ++} ++ ++/// LearningHealth state: EMA-smoothed health value with bounds and warmup. ++#[derive(Debug, Clone)] ++pub struct LearningHealth { ++ /// Current EMA-smoothed health [0.2, 0.95]. ++ pub value: f32, ++ /// Epoch counter — enables warmup mode for first N epochs. ++ pub epoch: u32, ++ /// Latest normalized components (for logging). ++ pub components: NormalizedComponents, ++ /// EMA smoothing factor (new sample weight). ++ ema_alpha: f32, ++ /// Warmup epochs — health clamped to 0.5 during warmup. ++ warmup_epochs: u32, ++} ++ ++impl LearningHealth { ++ pub fn new() -> Self { ++ Self { ++ value: 0.5, // neutral start ++ epoch: 0, ++ components: NormalizedComponents::default(), ++ ema_alpha: 0.3, ++ warmup_epochs: 3, ++ } ++ } ++ ++ /// Update health from new raw components. Returns updated health value. ++ pub fn update(&mut self, raw: &HealthComponents) -> f32 { ++ self.components = NormalizedComponents::from_raw(raw); ++ let new_sample = self.components.compose(); ++ ++ if self.epoch < self.warmup_epochs { ++ // Warmup: pin to 0.5, allow components to accumulate ++ self.value = 0.5; ++ } else { ++ // EMA update ++ let ema = (1.0 - self.ema_alpha) * self.value + self.ema_alpha * new_sample; ++ // Clamp to [0.2, 0.95] — hard floor prevents total adaptive-mechanism shutdown ++ // during transient collapse; hard ceiling keeps CQL/budget levers from pegging ++ // at extremes when health appears perfect. ++ self.value = ema.clamp(0.2, 0.95); ++ } ++ ++ self.epoch += 1; ++ self.value ++ } ++} ++ ++impl Default for LearningHealth { ++ fn default() -> Self { Self::new() } ++} ++ ++#[cfg(test)] ++mod tests { ++ use super::*; ++ ++ #[test] ++ fn test_smoothstep_boundaries() { ++ assert_eq!(smoothstep(0.0, 1.0, -0.5), 0.0); ++ assert_eq!(smoothstep(0.0, 1.0, 1.5), 1.0); ++ assert!((smoothstep(0.0, 1.0, 0.5) - 0.5).abs() < 0.001); ++ } ++ ++ #[test] ++ fn test_collapsed_inputs_give_low_health() { ++ let raw = HealthComponents { ++ q_gap: 0.01, // collapsed ++ q_var: 0.001, // collapsed ++ atom_util: 0.1, // under-utilized ++ grad_norm: 50_000.0, // exploding ++ ens_disagreement: 0.0,// everyone agrees (on uniform) ++ grad_consistency: -0.3, // gradients contradict ++ spectral_gap: 50.0, // rank 1 ++ }; ++ let mut health = LearningHealth::new(); ++ // Skip warmup ++ for _ in 0..5 { health.update(&raw); } ++ assert!(health.value <= 0.3, "Collapsed state should give health <= 0.3, got {}", health.value); ++ } ++ ++ #[test] ++ fn test_healthy_inputs_give_high_health() { ++ let raw = HealthComponents { ++ q_gap: 1.0, // strong differentiation ++ q_var: 0.5, // good variance ++ atom_util: 0.9, // atoms fully utilized ++ grad_norm: 5.0, // stable ++ ens_disagreement: 0.1, // moderate diversity ++ grad_consistency: 0.9, // consistent direction ++ spectral_gap: 1.5, // balanced rank ++ }; ++ let mut health = LearningHealth::new(); ++ for _ in 0..5 { health.update(&raw); } ++ assert!(health.value >= 0.7, "Healthy state should give health >= 0.7, got {}", health.value); ++ } ++ ++ #[test] ++ fn test_warmup_clamps_to_neutral() { ++ let raw = HealthComponents { ++ q_gap: 0.01, q_var: 0.001, atom_util: 0.1, ++ grad_norm: 50_000.0, ens_disagreement: 0.0, ++ grad_consistency: -0.5, spectral_gap: 100.0, ++ }; ++ let mut health = LearningHealth::new(); ++ // First 3 updates are warmup — should stay at 0.5 ++ for _ in 0..3 { ++ let v = health.update(&raw); ++ assert_eq!(v, 0.5, "Warmup should pin health to 0.5"); ++ } ++ // 4th update starts EMA ++ health.update(&raw); ++ assert!(health.value < 0.5, "Post-warmup should start dropping with collapsed inputs"); ++ } ++} +diff --git a/crates/ml/src/cuda_pipeline/meta_q_network.rs b/crates/ml/src/cuda_pipeline/meta_q_network.rs +new file mode 100644 +index 000000000..754d4eacc +--- /dev/null ++++ b/crates/ml/src/cuda_pipeline/meta_q_network.rs +@@ -0,0 +1,234 @@ ++//! N8: Meta-Q network — predicts collapse probability K epochs ahead. ++//! ++//! Small MLP (7→32→16→1 sigmoid) trained online via BCE loss on historical ++//! (health_components, future_collapse) pairs. Output is logged as a leading ++//! indicator — does NOT feed into LearningHealth until validated. ++ ++use std::collections::VecDeque; ++ ++pub const META_Q_INPUT_DIM: usize = 7; ++pub const META_Q_H1: usize = 32; ++pub const META_Q_H2: usize = 16; ++pub const META_Q_LOOKAHEAD: usize = 5; ++ ++#[derive(Debug, Clone)] ++pub struct MetaQInput { ++ pub q_gap_ema: f32, ++ pub log_grad_norm: f32, ++ pub atom_util: f32, ++ pub ens_disagreement: f32, ++ pub regime_stability: f32, ++ pub spectral_gap_norm: f32, ++ pub grad_consistency: f32, ++} ++ ++impl MetaQInput { ++ pub fn to_vec(&self) -> [f32; META_Q_INPUT_DIM] { ++ [ ++ self.q_gap_ema, ++ self.log_grad_norm, ++ self.atom_util, ++ self.ens_disagreement, ++ self.regime_stability, ++ self.spectral_gap_norm, ++ self.grad_consistency, ++ ] ++ } ++} ++ ++/// Small MLP on host (CPU). n_params ≈ 7×32 + 32×16 + 16×1 + biases = ~800. ++#[derive(Debug)] ++pub struct MetaQNetwork { ++ w1: Vec, b1: Vec, ++ w2: Vec, b2: Vec, ++ w3: Vec, b3: Vec, ++ m_w1: Vec, v_w1: Vec, ++ m_w2: Vec, v_w2: Vec, ++ m_w3: Vec, v_w3: Vec, ++ lr: f32, ++ step: u32, ++ buffer: VecDeque<(MetaQInput, f32)>, ++ buffer_cap: usize, ++} ++ ++impl MetaQNetwork { ++ pub fn new() -> Self { ++ // Use rand with a fixed seed for determinism of initial weights. ++ // fastrand is already a dep in several places; fall back to a simple ++ // deterministic LCG if fastrand is unavailable. ++ let mut rng = SimpleLcg::new(0xDEADBEEF); ++ let xavier = |rng: &mut SimpleLcg, fan_in: usize, fan_out: usize| -> f32 { ++ let range = (6.0 / (fan_in + fan_out) as f32).sqrt(); ++ (rng.next_f32() * 2.0 - 1.0) * range ++ }; ++ let w1: Vec = (0..META_Q_INPUT_DIM * META_Q_H1) ++ .map(|_| xavier(&mut rng, META_Q_INPUT_DIM, META_Q_H1)).collect(); ++ let w2: Vec = (0..META_Q_H1 * META_Q_H2) ++ .map(|_| xavier(&mut rng, META_Q_H1, META_Q_H2)).collect(); ++ let w3: Vec = (0..META_Q_H2) ++ .map(|_| xavier(&mut rng, META_Q_H2, 1)).collect(); ++ Self { ++ w1, b1: vec![0.0; META_Q_H1], ++ w2, b2: vec![0.0; META_Q_H2], ++ w3, b3: vec![0.0; 1], ++ m_w1: vec![0.0; META_Q_INPUT_DIM * META_Q_H1], v_w1: vec![0.0; META_Q_INPUT_DIM * META_Q_H1], ++ m_w2: vec![0.0; META_Q_H1 * META_Q_H2], v_w2: vec![0.0; META_Q_H1 * META_Q_H2], ++ m_w3: vec![0.0; META_Q_H2], v_w3: vec![0.0; META_Q_H2], ++ lr: 1e-3, ++ step: 0, ++ buffer: VecDeque::with_capacity(200), ++ buffer_cap: 200, ++ } ++ } ++ ++ pub fn predict(&self, input: &MetaQInput) -> f32 { ++ let x = input.to_vec(); ++ let mut h1 = vec![0.0f32; META_Q_H1]; ++ for j in 0..META_Q_H1 { ++ let mut s = self.b1[j]; ++ for i in 0..META_Q_INPUT_DIM { s += x[i] * self.w1[i * META_Q_H1 + j]; } ++ h1[j] = s.max(0.0); ++ } ++ let mut h2 = vec![0.0f32; META_Q_H2]; ++ for j in 0..META_Q_H2 { ++ let mut s = self.b2[j]; ++ for i in 0..META_Q_H1 { s += h1[i] * self.w2[i * META_Q_H2 + j]; } ++ h2[j] = s.max(0.0); ++ } ++ let mut logit = self.b3[0]; ++ for i in 0..META_Q_H2 { logit += h2[i] * self.w3[i]; } ++ 1.0 / (1.0 + (-logit).exp()) ++ } ++ ++ pub fn push_sample(&mut self, input: MetaQInput, label: f32) { ++ if self.buffer.len() >= self.buffer_cap { self.buffer.pop_front(); } ++ self.buffer.push_back((input, label)); ++ } ++ ++ pub fn train_step(&mut self) { ++ if self.buffer.len() < 32 { return; } ++ let mut rng = SimpleLcg::new(self.step as u64 ^ 0xAA55AA55); ++ for _ in 0..32 { ++ let idx = (rng.next_u32() as usize) % self.buffer.len(); ++ let (input, label) = self.buffer[idx].clone(); ++ self.train_one(&input, label); ++ } ++ self.step += 1; ++ } ++ ++ fn train_one(&mut self, input: &MetaQInput, label: f32) { ++ let x = input.to_vec(); ++ let mut h1_pre = vec![0.0f32; META_Q_H1]; ++ let mut h1 = vec![0.0f32; META_Q_H1]; ++ for j in 0..META_Q_H1 { ++ let mut s = self.b1[j]; ++ for i in 0..META_Q_INPUT_DIM { s += x[i] * self.w1[i * META_Q_H1 + j]; } ++ h1_pre[j] = s; h1[j] = s.max(0.0); ++ } ++ let mut h2_pre = vec![0.0f32; META_Q_H2]; ++ let mut h2 = vec![0.0f32; META_Q_H2]; ++ for j in 0..META_Q_H2 { ++ let mut s = self.b2[j]; ++ for i in 0..META_Q_H1 { s += h1[i] * self.w2[i * META_Q_H2 + j]; } ++ h2_pre[j] = s; h2[j] = s.max(0.0); ++ } ++ let mut logit = self.b3[0]; ++ for i in 0..META_Q_H2 { logit += h2[i] * self.w3[i]; } ++ let pred = 1.0 / (1.0 + (-logit).exp()); ++ ++ let d_logit = pred - label; ++ let mut d_w3 = vec![0.0f32; META_Q_H2]; ++ for i in 0..META_Q_H2 { d_w3[i] = d_logit * h2[i]; } ++ let d_b3 = d_logit; ++ let d_h2: Vec = (0..META_Q_H2).map(|i| d_logit * self.w3[i] * (if h2_pre[i] > 0.0 { 1.0 } else { 0.0 })).collect(); ++ let mut d_w2 = vec![0.0f32; META_Q_H1 * META_Q_H2]; ++ let d_b2 = d_h2.clone(); ++ for i in 0..META_Q_H1 { ++ for j in 0..META_Q_H2 { d_w2[i * META_Q_H2 + j] = d_h2[j] * h1[i]; } ++ } ++ let d_h1: Vec = (0..META_Q_H1).map(|i| { ++ let mut s = 0.0f32; ++ for j in 0..META_Q_H2 { s += d_h2[j] * self.w2[i * META_Q_H2 + j]; } ++ s * (if h1_pre[i] > 0.0 { 1.0 } else { 0.0 }) ++ }).collect(); ++ let mut d_w1 = vec![0.0f32; META_Q_INPUT_DIM * META_Q_H1]; ++ let d_b1 = d_h1.clone(); ++ for i in 0..META_Q_INPUT_DIM { ++ for j in 0..META_Q_H1 { d_w1[i * META_Q_H1 + j] = d_h1[j] * x[i]; } ++ } ++ self.adam_update(d_w1, d_w2, d_w3, d_b1, d_b2, d_b3); ++ } ++ ++ fn adam_update(&mut self, d_w1: Vec, d_w2: Vec, d_w3: Vec, d_b1: Vec, d_b2: Vec, d_b3: f32) { ++ let beta1 = 0.9f32; let beta2 = 0.999f32; let eps = 1e-8f32; ++ let t = (self.step + 1) as f32; ++ let bc1 = 1.0 - beta1.powf(t); let bc2 = 1.0 - beta2.powf(t); ++ let lr = self.lr; ++ ++ let adam = |p: &mut [f32], m: &mut [f32], v: &mut [f32], g: &[f32]| { ++ for i in 0..p.len() { ++ m[i] = beta1 * m[i] + (1.0 - beta1) * g[i]; ++ v[i] = beta2 * v[i] + (1.0 - beta2) * g[i] * g[i]; ++ let mh = m[i] / bc1; ++ let vh = v[i] / bc2; ++ p[i] -= lr * mh / (vh.sqrt() + eps); ++ } ++ }; ++ adam(&mut self.w1, &mut self.m_w1, &mut self.v_w1, &d_w1); ++ adam(&mut self.w2, &mut self.m_w2, &mut self.v_w2, &d_w2); ++ adam(&mut self.w3, &mut self.m_w3, &mut self.v_w3, &d_w3); ++ for i in 0..self.b1.len() { self.b1[i] -= lr * d_b1[i]; } ++ for i in 0..self.b2.len() { self.b2[i] -= lr * d_b2[i]; } ++ self.b3[0] -= lr * d_b3; ++ } ++} ++ ++impl Default for MetaQNetwork { ++ fn default() -> Self { Self::new() } ++} ++ ++/// Deterministic LCG for reproducible weight init and minibatch sampling. ++#[derive(Debug)] ++struct SimpleLcg { state: u64 } ++ ++impl SimpleLcg { ++ fn new(seed: u64) -> Self { Self { state: seed.wrapping_add(1) } } ++ fn next_u32(&mut self) -> u32 { ++ self.state = self.state.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407); ++ (self.state >> 33) as u32 ++ } ++ fn next_f32(&mut self) -> f32 { ++ (self.next_u32() as f32) / (u32::MAX as f32) ++ } ++} ++ ++#[cfg(test)] ++mod tests { ++ use super::*; ++ ++ #[test] ++ fn test_predict_returns_probability() { ++ let net = MetaQNetwork::new(); ++ let input = MetaQInput { ++ q_gap_ema: 0.3, log_grad_norm: 2.0, atom_util: 0.5, ++ ens_disagreement: 0.1, regime_stability: 0.7, ++ spectral_gap_norm: 0.6, grad_consistency: 0.3, ++ }; ++ let p = net.predict(&input); ++ assert!((0.0..=1.0).contains(&p), "probability out of range: {}", p); ++ } ++ ++ #[test] ++ fn test_training_reduces_loss() { ++ let mut net = MetaQNetwork::new(); ++ let collapse_input = MetaQInput { ++ q_gap_ema: 0.01, log_grad_norm: 10.0, atom_util: 0.1, ++ ens_disagreement: 0.0, regime_stability: 0.5, ++ spectral_gap_norm: 0.1, grad_consistency: -0.3, ++ }; ++ for _ in 0..100 { net.push_sample(collapse_input.clone(), 1.0); } ++ for _ in 0..50 { net.train_step(); } ++ let p = net.predict(&collapse_input); ++ assert!(p > 0.5, "Trained net should predict >0.5 for collapse state, got {}", p); ++ } ++} +diff --git a/crates/ml/src/cuda_pipeline/mod.rs b/crates/ml/src/cuda_pipeline/mod.rs +index c714a72f1..981ef4b26 100644 +--- a/crates/ml/src/cuda_pipeline/mod.rs ++++ b/crates/ml/src/cuda_pipeline/mod.rs +@@ -37,6 +37,9 @@ pub mod gpu_iql_trainer; + pub mod gpu_iqn_head; + pub mod gpu_attention; + pub mod decision_transformer; ++pub mod learning_health; ++pub mod q_snapshot; ++pub mod meta_q_network; + // gpu_replay_buffer moved to ml-dqn crate + + /// Maximum bytes allowed for a single GPU upload (2 GB safety limit). +diff --git a/crates/ml/src/cuda_pipeline/q_snapshot.rs b/crates/ml/src/cuda_pipeline/q_snapshot.rs +new file mode 100644 +index 000000000..a63d88441 +--- /dev/null ++++ b/crates/ml/src/cuda_pipeline/q_snapshot.rs +@@ -0,0 +1,123 @@ ++//! N1: Q-network weight snapshots ring buffer for temporal self-distillation. ++//! ++//! Keeps up to `MAX_SNAPSHOTS` historical weight checkpoints taken at high-health ++//! epochs. During collapse recovery, a gradient pull toward the best snapshot ++//! is added to the master gradient. ++ ++use std::sync::Arc; ++use cudarc::driver::{CudaSlice, CudaStream, DevicePtr, DevicePtrMut}; ++use crate::MLError; ++ ++pub const MAX_SNAPSHOTS: usize = 5; ++pub const SNAPSHOT_HEALTH_THRESHOLD: f32 = 0.7; ++pub const DISTILL_HEALTH_THRESHOLD: f32 = 0.4; ++ ++/// A single weight snapshot with metadata. ++pub struct QSnapshot { ++ pub weights: CudaSlice, ++ pub health: f32, ++ pub q_gap: f32, ++ pub epoch: u32, ++} ++ ++impl std::fmt::Debug for QSnapshot { ++ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { ++ f.debug_struct("QSnapshot") ++ .field("health", &self.health) ++ .field("q_gap", &self.q_gap) ++ .field("epoch", &self.epoch) ++ .finish() ++ } ++} ++ ++/// Ring buffer of best-health snapshots. ++pub struct SnapshotRing { ++ stream: Arc, ++ pub(crate) snapshots: Vec, ++ param_count: usize, ++} ++ ++impl std::fmt::Debug for SnapshotRing { ++ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { ++ f.debug_struct("SnapshotRing") ++ .field("param_count", &self.param_count) ++ .field("snapshot_count", &self.snapshots.len()) ++ .finish() ++ } ++} ++ ++impl SnapshotRing { ++ pub fn new(stream: Arc, param_count: usize) -> Self { ++ Self { ++ stream, ++ snapshots: Vec::with_capacity(MAX_SNAPSHOTS), ++ param_count, ++ } ++ } ++ ++ pub fn is_empty(&self) -> bool { ++ self.snapshots.is_empty() ++ } ++ ++ pub fn param_count(&self) -> usize { ++ self.param_count ++ } ++ ++ /// Copy current weights into a new snapshot if health >= threshold and ++ /// q_gap beats the worst stored snapshot (when at capacity). ++ pub fn maybe_snapshot( ++ &mut self, ++ params_src: &CudaSlice, ++ health: f32, ++ q_gap: f32, ++ epoch: u32, ++ ) -> Result { ++ if health < SNAPSHOT_HEALTH_THRESHOLD { ++ return Ok(false); ++ } ++ ++ // At capacity, only replace if this q_gap beats the worst stored entry. ++ if self.snapshots.len() >= MAX_SNAPSHOTS { ++ let worst_q_gap = self.snapshots.iter() ++ .map(|s| s.q_gap) ++ .fold(f32::INFINITY, f32::min); ++ if q_gap <= worst_q_gap { ++ return Ok(false); ++ } ++ } ++ ++ let n = self.param_count; ++ let num_bytes = n * std::mem::size_of::(); ++ let mut new_weights = self.stream.alloc_zeros::(n) ++ .map_err(|e| MLError::ModelError(format!("snapshot alloc: {e}")))?; ++ { ++ let (src_ptr, _src_guard) = params_src.device_ptr(&self.stream); ++ let (dst_ptr, _dst_guard) = new_weights.device_ptr_mut(&self.stream); ++ #[allow(unsafe_code)] ++ unsafe { ++ cudarc::driver::result::memcpy_dtod_async( ++ dst_ptr, src_ptr, num_bytes, self.stream.cu_stream(), ++ ).map_err(|e| MLError::ModelError(format!("snapshot dtod: {e}")))?; ++ } ++ } ++ ++ if self.snapshots.len() >= MAX_SNAPSHOTS { ++ let worst_idx = self.snapshots.iter() ++ .enumerate() ++ .min_by(|(_, a), (_, b)| a.q_gap.partial_cmp(&b.q_gap).unwrap()) ++ .map(|(i, _)| i) ++ .unwrap(); ++ self.snapshots.swap_remove(worst_idx); ++ } ++ ++ self.snapshots.push(QSnapshot { weights: new_weights, health, q_gap, epoch }); ++ Ok(true) ++ } ++ ++ /// Best (highest q_gap) snapshot raw device pointer, if any. ++ pub fn best_raw_ptr(&self) -> Option { ++ self.snapshots.iter() ++ .max_by(|a, b| a.q_gap.partial_cmp(&b.q_gap).unwrap()) ++ .map(|s| s.weights.raw_ptr()) ++ } ++} +diff --git a/crates/ml/src/trainers/dqn/fused_training.rs b/crates/ml/src/trainers/dqn/fused_training.rs +index 0dcd32f7a..b84d36816 100644 +--- a/crates/ml/src/trainers/dqn/fused_training.rs ++++ b/crates/ml/src/trainers/dqn/fused_training.rs +@@ -1157,23 +1157,26 @@ impl FusedTrainingCtx { + ).map_err(|e| anyhow::anyhow!("HER in-place relabel kernel: {e}"))?; + } + +- // NOTE: Per-component gradient budget clip REMOVED. +- // The old code clipped C51+MSE gradient to a fraction of max_grad_norm, +- // producing a CONSTANT gradient norm every step. Adam degenerates into +- // SignSGD when it never sees true gradient magnitude variation. +- // Safety is now handled by the single global clip inside Adam's kernel. ++ // NOTE: Per-component gradient budget clip REMOVED (old explicit clip → constant grad norm). ++ // B4/G5: Budgets are now applied as SCALE factors on individual SAXPY contributions, ++ // not as norm clips. This preserves gradient magnitude variation while steering ++ // component balance. Safety is still handled by the single global clip in Adam. ++ ++ // B4/G5: Compute adaptive budgets once per step from ISV health/regime signals. ++ let (_c51_budget, iqn_budget, cql_budget, _ens_budget) = self.compute_adaptive_budgets(); + + // DIAGNOSTIC: sync between each aux op to find hanging kernel + + // EMA target update — adaptive tau based on Q-divergence. + { + let dqn = agent.primary_dqn_mut(); +- let tau = compute_cosine_annealed_tau( ++ let tau_scheduled = compute_cosine_annealed_tau( + dqn.get_training_steps(), + dqn.config.tau, + dqn.config.tau_final, + dqn.config.tau_anneal_steps, + ); ++ let tau = self.trainer.apply_health_coupled_tau_floor(tau_scheduled); + // Use cosine-annealed tau directly — fully deterministic. + // Adaptive tau (from q_divergence) was removed because q_divergence + // uses atomicAdd in the C51 loss kernel, introducing non-determinism +@@ -1424,16 +1427,17 @@ impl FusedTrainingCtx { + if let Some(ref mut iqn) = self.gpu_iqn { + let d_h_s2_ptr = iqn.d_h_s2_raw_ptr(); + self.trainer.apply_iqn_trunk_gradient( +- d_h_s2_ptr, &mut self.online_dueling, ++ d_h_s2_ptr, &mut self.online_dueling, iqn_budget, + ).map_err(|e| anyhow::anyhow!("IQN trunk gradient: {e}"))?; + + let dqn = agent.primary_dqn_mut(); +- let tau = compute_cosine_annealed_tau( ++ let tau_scheduled = compute_cosine_annealed_tau( + dqn.get_training_steps(), + dqn.config.tau, + dqn.config.tau_final, + dqn.config.tau_anneal_steps, + ); ++ let tau = self.trainer.apply_health_coupled_tau_floor(tau_scheduled); + iqn.target_ema_update(tau as f32) + .map_err(|e| anyhow::anyhow!("IQN EMA update: {e}"))?; + +@@ -1455,10 +1459,11 @@ impl FusedTrainingCtx { + + + // CQL gradient + SAXPY (no per-component clip — Adam handles safety). ++ // B4/G5: cql_budget scales the SAXPY contribution (0.10×(1−regime)×health). + if self.trainer.has_cql() { + match self.trainer.apply_cql_gradient() { + Ok(true) => { +- self.trainer.apply_cql_saxpy() ++ self.trainer.apply_cql_saxpy(cql_budget) + .map_err(|e| anyhow::anyhow!("CQL SAXPY: {e}"))?; + } + Ok(false) => {} +@@ -1829,23 +1834,25 @@ impl FusedTrainingCtx { + iqn.increment_adam_step(); + // Update IQN tau via stable host address + let dqn = agent.primary_dqn_mut(); +- let tau = compute_cosine_annealed_tau( ++ let tau_scheduled = compute_cosine_annealed_tau( + dqn.get_training_steps(), + dqn.config.tau, + dqn.config.tau_final, + dqn.config.tau_anneal_steps, + ); ++ let tau = self.trainer.apply_health_coupled_tau_floor(tau_scheduled); + iqn.set_tau_host(tau as f32); + } + // Update DQN trainer tau via stable host address + { + let dqn = agent.primary_dqn_mut(); +- let tau = compute_cosine_annealed_tau( ++ let tau_scheduled = compute_cosine_annealed_tau( + dqn.get_training_steps(), + dqn.config.tau, + dqn.config.tau_final, + dqn.config.tau_anneal_steps, + ); ++ let tau = self.trainer.apply_health_coupled_tau_floor(tau_scheduled); + self.trainer.set_tau_value(tau as f32); + } + } +@@ -2172,6 +2179,98 @@ impl FusedTrainingCtx { + self.trainer.read_isv_regime() + } + ++ /// D8/N8: Forward (health, regime_stability) read through to caller. ++ pub(crate) fn read_isv_regime_via_trainer(&self) -> (f32, f32) { ++ self.trainer.read_isv_health_and_regime() ++ } ++ ++ // ── A3: LearningHealth signal writers / readers ─────────────────────── ++ ++ /// Write a single f32 value to ISV signals buffer at given index (pinned HtoD). ++ pub(crate) fn write_isv_signal_at(&self, index: usize, value: f32) { ++ self.trainer.write_isv_signal_at(index, value); ++ } ++ ++ /// Read atom utilization from the last Q-stats readback (pinned pointer offset 6). ++ pub(crate) fn read_atom_utilization(&self) -> f32 { ++ self.trainer.read_atom_utilization() ++ } ++ ++ /// Last effective cql_alpha from the most recent apply_cql_gradient() call. ++ pub(crate) fn last_cql_alpha_eff(&self) -> f32 { ++ self.trainer.last_cql_alpha_eff ++ } ++ ++ /// Last effective tau (cosine-annealed, potentially health-floored). ++ pub(crate) fn last_tau_eff(&self) -> f32 { ++ self.trainer.last_tau_eff ++ } ++ ++ /// B3/G4: Last health-scaled Expected SARSA tau factor (1.0 at full health, 6.0 at collapse). ++ pub(crate) fn last_sarsa_tau_factor(&self) -> f32 { ++ self.trainer.last_sarsa_tau_factor ++ } ++ ++ /// B4/G5: Compute adaptive gradient budgets from learning_health and regime_stability. ++ /// Returns (c51_budget, iqn_budget, cql_budget, ens_budget). Caches to trainer fields ++ /// for HEALTH_DIAG logging. ++ /// ++ /// Formula: ++ /// iqn_budget = 0.10 + 0.30 × health ++ /// cql_budget = 0.10 × (1 − regime_stability) × health ++ /// ens_budget = 0.05 (constant) ++ /// c51_budget = 1.0 − iqn_budget − cql_budget − ens_budget ++ /// ++ /// At health=1: iqn=0.40, cql variable (depends on regime_stability), ens=0.05, c51=rest. ++ /// At health=0: iqn=0.10, cql=0, ens=0.05, c51=0.85 (C51 takes over). ++ pub(crate) fn compute_adaptive_budgets(&mut self) -> (f32, f32, f32, f32) { ++ let (health, regime_stability) = self.trainer.read_isv_health_and_regime(); ++ let iqn_budget = 0.10_f32 + 0.30_f32 * health; ++ let cql_budget = 0.10_f32 * (1.0 - regime_stability) * health; ++ let ens_budget = 0.05_f32; ++ let c51_budget = (1.0_f32 - iqn_budget - cql_budget - ens_budget).max(0.0); ++ ++ self.trainer.last_iqn_budget_eff = iqn_budget; ++ self.trainer.last_cql_budget_eff = cql_budget; ++ self.trainer.last_c51_budget_eff = c51_budget; ++ self.trainer.last_ens_budget_eff = ens_budget; ++ ++ (c51_budget, iqn_budget, cql_budget, ens_budget) ++ } ++ ++ /// B4/G5: Last adaptive IQN gradient budget (0.10 + 0.30×health). ++ pub(crate) fn last_iqn_budget_eff(&self) -> f32 { self.trainer.last_iqn_budget_eff } ++ /// B4/G5: Last adaptive CQL gradient budget (0.10×(1−regime)×health). ++ pub(crate) fn last_cql_budget_eff(&self) -> f32 { self.trainer.last_cql_budget_eff } ++ /// B4/G5: Last adaptive C51 gradient budget (1−iqn−cql−ens). ++ pub(crate) fn last_c51_budget_eff(&self) -> f32 { self.trainer.last_c51_budget_eff } ++ /// B4/G5: Last ensemble gradient budget (constant 0.05). ++ pub(crate) fn last_ens_budget_eff(&self) -> f32 { self.trainer.last_ens_budget_eff } ++ ++ /// Compute spectral gap on the last Q readback slice (coarse max/min ratio proxy). ++ pub(crate) fn compute_q_spectral_gap(&self) -> f32 { ++ self.trainer.compute_q_spectral_gap() ++ } ++ ++ /// D1/N1: Conditionally snapshot current weights at epoch boundary. ++ /// Snapshots if health >= 0.7 and q_gap beats the worst stored entry. ++ pub(crate) fn maybe_snapshot_qnet(&mut self, health: f32, q_gap: f32, epoch: u32) -> anyhow::Result { ++ self.trainer.maybe_snapshot_params(health, q_gap, epoch) ++ .map_err(|e| anyhow::anyhow!("snapshot: {e}")) ++ } ++ ++ /// D1/N1: Apply distillation gradient pull toward best snapshot. ++ /// No-op when health >= 0.4 or no snapshots exist yet. ++ pub(crate) fn apply_distillation(&mut self) -> anyhow::Result { ++ self.trainer.apply_distillation_gradient() ++ .map_err(|e| anyhow::anyhow!("distill: {e}")) ++ } ++ ++ /// D1/N1: Whether distillation was active in the most recent epoch boundary. ++ pub(crate) fn last_distill_active(&self) -> bool { ++ self.trainer.last_distill_active ++ } ++ + pub(crate) fn read_popart_variance(&self) -> f32 { + if !self.popart_enabled { return 0.0; } + let mut var = [0.0_f32]; +@@ -2666,6 +2765,17 @@ impl FusedTrainingCtx { + self.trainer.set_adaptive_gamma(gamma); + } + ++ /// C4/P4: Apply temporal-coupled adjustment to the scheduled gamma. ++ /// Delegates to `GpuDqnTrainer::apply_adaptive_gamma`, which reads health ++ /// and regime_stability from ISV pinned memory, clamps to [0.9, 0.995], ++ /// and caches the result in `last_gamma_eff`. ++ pub(crate) fn apply_adaptive_gamma(&mut self, gamma_base: f32) -> f32 { ++ self.trainer.apply_adaptive_gamma(gamma_base) ++ } ++ ++ /// C4/P4: Last effective gamma after temporal-coupled adjustment. ++ pub(crate) fn last_gamma_eff(&self) -> f32 { self.trainer.last_gamma_eff } ++ + /// G5: Set the EMA threshold for epistemic variance gating. + pub(crate) fn set_var_ema(&self, val: f32) { + self.trainer.set_var_ema(val); +diff --git a/crates/ml/src/trainers/dqn/smoke_tests/adaptive_learning.rs b/crates/ml/src/trainers/dqn/smoke_tests/adaptive_learning.rs +new file mode 100644 +index 000000000..ed0261407 +--- /dev/null ++++ b/crates/ml/src/trainers/dqn/smoke_tests/adaptive_learning.rs +@@ -0,0 +1,96 @@ ++//! E1: Collapse-recovery smoke test for the LearningHealth adaptive system. ++//! ++//! Trains a short DQN run on fxcache data and asserts that the adaptive ++//! mechanisms (health-coupled CQL/tau/SARSA-temp/gradient-budget/gamma, ++//! PER diversity, snapshots/distillation, Q-gap barrier, plasticity injection, ++//! CF curriculum, IB, ensemble oracle, contrarian override, meta-Q) collectively ++//! prevent Q-value collapse. The specific thresholds encode the design goals: ++//! ++//! - `q_gap_ema` must stay above 0.05 (non-trivial action differentiation) ++//! - `learning_health.value` must stay above 0.3 (system not pinned in collapse mode) ++//! ++//! Baseline runs prior to this feature collapsed both signals to near-zero by ++//! epoch 8 — this test is a regression guard. ++//! ++//! Run: `FOXHUNT_TEST_DATA=test_data/futures-baseline \ ++//! cargo test -p ml --lib -- test_adaptive_learning_no_collapse --ignored --nocapture` ++ ++use super::helpers::*; ++ ++#[test] ++#[ignore] // GPU + fxcache data ++fn test_adaptive_learning_no_collapse() -> anyhow::Result<()> { ++ use crate::fxcache; ++ ++ let cache_dir = feature_cache_dir(); ++ assert!(cache_dir.exists(), "feature-cache not found at {:?}", cache_dir); ++ ++ let entries: Vec<_> = std::fs::read_dir(&cache_dir)? ++ .filter_map(|e| e.ok()) ++ .filter(|e| e.path().extension().and_then(|s| s.to_str()) == Some("fxcache")) ++ .collect(); ++ assert!(!entries.is_empty(), "No .fxcache files in {:?}", cache_dir); ++ ++ let fxcache_data = fxcache::load_fxcache(&entries[0].path())?; ++ let n = fxcache_data.bar_count.min(20_000); ++ let train_end = (n * 80) / 100; ++ ++ let mut params = smoke_params(); ++ params.epochs = 10; ++ ++ let mut trainer = smoke_trainer_with(params)?; ++ let rt = tokio::runtime::Builder::new_current_thread() ++ .enable_all() ++ .build()?; ++ ++ let features = &fxcache_data.features[..n]; ++ let targets = &fxcache_data.targets[..n]; ++ let ofi = &fxcache_data.ofi[..n.min(fxcache_data.ofi.len())]; ++ ++ let train_feat = &features[..train_end]; ++ let train_targets = &targets[..train_end]; ++ let val_feat = &features[train_end..n]; ++ let val_targets = &targets[train_end..n]; ++ ++ rt.block_on(trainer.init_from_fxcache(features, targets, ofi))?; ++ trainer.set_training_range(0, train_end, train_end, n); ++ trainer.set_val_data_from_slices(val_feat, val_targets, train_end); ++ rt.block_on(trainer.reset_for_fold())?; ++ ++ rt.block_on(trainer.train_fold_from_slices( ++ train_feat, ++ train_targets, ++ |_epoch, _bytes, _best| Ok("skip".to_owned()), ++ ))?; ++ ++ // ── Post-training assertions ────────────────────────────────────────── ++ let final_q_gap = trainer.health_ema.q_gap_ema; ++ let final_health = trainer.learning_health.value; ++ let unhealthy_streak = trainer.unhealthy_epoch_count; ++ ++ eprintln!( ++ "[ADAPTIVE] final q_gap_ema={:.4}, health={:.2}, unhealthy_streak={}, \ ++ low_winrate={}, contrarian_active={}", ++ final_q_gap, ++ final_health, ++ unhealthy_streak, ++ trainer.low_winrate_count, ++ trainer.contrarian_active, ++ ); ++ ++ assert!( ++ final_q_gap > 0.05, ++ "Q-gap collapsed to {:.4} after 10 epochs — adaptive mechanisms failed \ ++ to prevent collapse (threshold 0.05)", ++ final_q_gap ++ ); ++ ++ assert!( ++ final_health > 0.3, ++ "LearningHealth dropped to {:.2} — system is pinned in collapse mode \ ++ (threshold 0.30)", ++ final_health ++ ); ++ ++ Ok(()) ++} +diff --git a/crates/ml/src/trainers/dqn/smoke_tests/mod.rs b/crates/ml/src/trainers/dqn/smoke_tests/mod.rs +index db2b1b0ac..12a29580d 100644 +--- a/crates/ml/src/trainers/dqn/smoke_tests/mod.rs ++++ b/crates/ml/src/trainers/dqn/smoke_tests/mod.rs +@@ -20,3 +20,5 @@ mod regression; + pub mod reward_v8; + #[cfg(test)] + mod generalization; ++#[cfg(test)] ++mod adaptive_learning; +diff --git a/crates/ml/src/trainers/dqn/trainer/constructor.rs b/crates/ml/src/trainers/dqn/trainer/constructor.rs +index cdaba041f..f37c91b6e 100644 +--- a/crates/ml/src/trainers/dqn/trainer/constructor.rs ++++ b/crates/ml/src/trainers/dqn/trainer/constructor.rs +@@ -723,6 +723,36 @@ impl DQNTrainer { + enrichment: super::enrichment::EnrichmentState::new( + initial_epsilon, + ), ++ learning_health: crate::cuda_pipeline::learning_health::LearningHealth::new(), ++ health_ema: Default::default(), ++ last_q_gap: None, ++ last_q_var: None, ++ last_grad_norm: None, ++ last_ens_disagreement: None, ++ last_atom_util: None, ++ last_cql_alpha_eff: None, ++ last_iqn_budget_eff: None, ++ last_cql_budget_eff: None, ++ last_c51_budget_eff: None, ++ last_tau_eff: None, ++ last_sarsa_tau_factor: None, ++ last_gamma_eff: None, ++ last_cf_ratio_eff: None, ++ last_distill_active: None, ++ last_barrier_loss: None, ++ last_plasticity_ready: None, ++ unhealthy_epoch_count: 0, ++ last_ib_penalty: None, ++ last_ensemble_collapse_score: None, ++ last_contrarian_active: None, ++ low_winrate_count: 0, ++ contrarian_active: false, ++ contrarian_remaining_epochs: 0, ++ last_epoch_win_rate: 0.5, ++ last_meta_q_pred: None, ++ meta_q: crate::cuda_pipeline::meta_q_network::MetaQNetwork::new(), ++ pending_meta_q_samples: std::collections::VecDeque::with_capacity(32), ++ health_history: std::collections::VecDeque::with_capacity(2 * crate::cuda_pipeline::meta_q_network::META_Q_LOOKAHEAD), + }) + } + +diff --git a/crates/ml/src/trainers/dqn/trainer/metrics.rs b/crates/ml/src/trainers/dqn/trainer/metrics.rs +index e7442d17a..4f889ae72 100644 +--- a/crates/ml/src/trainers/dqn/trainer/metrics.rs ++++ b/crates/ml/src/trainers/dqn/trainer/metrics.rs +@@ -6,6 +6,43 @@ use super::DQNTrainer; + use crate::TrainingMetrics; + use super::super::config::DQNAgentType; + ++/// Rolling EMA values for LearningHealth inputs. ++#[derive(Debug, Clone, Default)] ++pub(crate) struct HealthEmaTrackers { ++ pub q_gap_ema: f32, ++ pub q_var_ema: f32, ++ pub grad_norm_ema: f32, ++ prev_grad_norm: Option, ++ ema_alpha: f32, // initialized to 0.1 on first update ++} ++ ++impl HealthEmaTrackers { ++ pub(crate) fn update(&mut self, q_gap: f32, q_var: f32, grad_norm: f32) { ++ if self.ema_alpha == 0.0 { self.ema_alpha = 0.1; } ++ self.q_gap_ema = (1.0 - self.ema_alpha) * self.q_gap_ema + self.ema_alpha * q_gap; ++ self.q_var_ema = (1.0 - self.ema_alpha) * self.q_var_ema + self.ema_alpha * q_var; ++ self.grad_norm_ema = (1.0 - self.ema_alpha) * self.grad_norm_ema + self.ema_alpha * grad_norm; ++ } ++ ++ /// Simple scalar proxy for gradient direction consistency: 1 - |delta|/max. ++ /// Returns 0.0 on first call; +1.0 means no change, -1.0 means opposite sign. ++ /// This is a weaker proxy than full vector cosine similarity, but keeps A3 ++ /// scope tight — can be upgraded in a follow-up if the Adam flat buffer ++ /// becomes available. ++ pub(crate) fn update_grad_consistency(&mut self, current_grad_norm: f32) -> f32 { ++ let cos = match self.prev_grad_norm { ++ None => 0.0, ++ Some(prev) => { ++ let denom = prev.abs().max(current_grad_norm.abs()).max(1e-6); ++ let delta = (current_grad_norm - prev).abs() / denom; ++ (1.0 - 2.0 * delta).clamp(-1.0, 1.0) ++ } ++ }; ++ self.prev_grad_norm = Some(current_grad_norm); ++ cos ++ } ++} ++ + impl DQNTrainer { + /// Full training loop with existing logic (Wave 12 Group 3) + /// Calculate average metrics for an epoch +diff --git a/crates/ml/src/trainers/dqn/trainer/mod.rs b/crates/ml/src/trainers/dqn/trainer/mod.rs +index 16c565a7b..1ababcb82 100644 +--- a/crates/ml/src/trainers/dqn/trainer/mod.rs ++++ b/crates/ml/src/trainers/dqn/trainer/mod.rs +@@ -472,6 +472,72 @@ pub struct DQNTrainer { + + /// Self-improving enrichment state (persists across epochs). + pub(crate) enrichment: enrichment::EnrichmentState, ++ ++ // ── A3: LearningHealth (Layer 1) ───────────────────────────────── ++ /// Unified training health scalar [0.2, 0.95] — EMA-smoothed composition of ++ /// Q-gap, Q-variance, atom utilization, grad stability, ensemble agreement, ++ /// gradient consistency, and spectral gap. Broadcast to ISV[12] each epoch. ++ pub(crate) learning_health: crate::cuda_pipeline::learning_health::LearningHealth, ++ /// Rolling EMA trackers for LearningHealth input signals. ++ pub(crate) health_ema: metrics::HealthEmaTrackers, ++ /// Last computed Q-gap signal (mean of 4 per-branch EMA gaps). Cached for A4/B1. ++ pub(crate) last_q_gap: Option, ++ /// Last computed Q-variance proxy (epoch_q_max - epoch_q_min). Cached for A4. ++ pub(crate) last_q_var: Option, ++ /// Last grad_norm EMA fed into LearningHealth. Cached for A4/B4. ++ pub(crate) last_grad_norm: Option, ++ /// Last ensemble disagreement value (placeholder 0.1 until D6/N6). Cached for A4. ++ pub(crate) last_ens_disagreement: Option, ++ /// Last atom utilization from q_readback_pinned[6]. Cached for A4/G4. ++ pub(crate) last_atom_util: Option, ++ ++ // ── A4: Extended HEALTH_DIAG placeholders ───────────────────────────── ++ // All fields start as None; B/C/D tasks populate them with real values. ++ /// B1/G2 — CQL alpha effective value ++ pub(crate) last_cql_alpha_eff: Option, ++ /// B4/G5 — per-component gradient budget allocations ++ pub(crate) last_iqn_budget_eff: Option, ++ pub(crate) last_cql_budget_eff: Option, ++ pub(crate) last_c51_budget_eff: Option, ++ /// B2/G3 — tau effective ++ pub(crate) last_tau_eff: Option, ++ /// B3/G4 — Expected SARSA temperature scaling factor ++ pub(crate) last_sarsa_tau_factor: Option, ++ /// C4/P4 — adaptive gamma ++ pub(crate) last_gamma_eff: Option, ++ /// D4/N4 — counterfactual ratio ++ pub(crate) last_cf_ratio_eff: Option, ++ /// D1/N1 — distillation active flag ++ pub(crate) last_distill_active: Option, ++ /// D2/N2 — barrier loss value ++ pub(crate) last_barrier_loss: Option, ++ /// D3/N3 — plasticity-ready flag (true = ready to trigger, false = in cooldown) ++ pub(crate) last_plasticity_ready: Option, ++ /// D3/N3: Consecutive epoch count with health < 0.3. When reaches 3, triggers plasticity injection. ++ pub(crate) unhealthy_epoch_count: u32, ++ /// D5/N5 — information bottleneck penalty ++ pub(crate) last_ib_penalty: Option, ++ /// D6/N6 — ensemble collapse score ++ pub(crate) last_ensemble_collapse_score: Option, ++ /// D7/N7 — contrarian override active ++ pub(crate) last_contrarian_active: Option, ++ /// D7/N7: Consecutive epochs with WinRate < 0.40. Triggers contrarian override at 5. ++ pub(crate) low_winrate_count: u32, ++ /// D7/N7: Whether contrarian override (argmin) is currently active. ++ pub(crate) contrarian_active: bool, ++ /// D7/N7: Epochs remaining in current contrarian window (2 epochs per activation). ++ pub(crate) contrarian_remaining_epochs: u32, ++ /// D7/N7: Most recent epoch win_rate (0..1) from financials. Carried into next process_epoch_boundary. ++ pub(crate) last_epoch_win_rate: f32, ++ /// D8/N8 — meta-Q collapse prediction ++ pub(crate) last_meta_q_pred: Option, ++ /// D8/N8: Meta-Q network predicting collapse probability K epochs ahead. ++ pub(crate) meta_q: crate::cuda_pipeline::meta_q_network::MetaQNetwork, ++ /// D8/N8: Pending inputs awaiting label resolution (K epochs delayed). ++ pub(crate) pending_meta_q_samples: std::collections::VecDeque<(crate::cuda_pipeline::meta_q_network::MetaQInput, u32)>, ++ /// D8/N8: Rolling history of (epoch, health) for meta-Q label resolution. ++ /// Capped at 2×META_Q_LOOKAHEAD. ++ pub(crate) health_history: std::collections::VecDeque<(u32, f32)>, + } + + impl std::fmt::Debug for DQNTrainer { +diff --git a/crates/ml/src/trainers/dqn/trainer/training_loop.rs b/crates/ml/src/trainers/dqn/trainer/training_loop.rs +index 95698c79e..67e306eb3 100644 +--- a/crates/ml/src/trainers/dqn/trainer/training_loop.rs ++++ b/crates/ml/src/trainers/dqn/trainer/training_loop.rs +@@ -477,18 +477,8 @@ impl DQNTrainer { + } + } + +- // Periodic shrink-and-perturb +- let sp_interval = self.hyperparams.shrink_perturb_interval; +- if sp_interval > 0 && epoch > 0 && epoch % sp_interval == 0 { +- if let Some(ref mut fused) = self.fused_ctx { +- let alpha = self.hyperparams.shrink_perturb_alpha as f32; +- let sigma = self.hyperparams.shrink_perturb_sigma as f32; +- match fused.shrink_and_perturb(alpha, sigma) { +- Ok(()) => info!(epoch, alpha, sigma, "Periodic shrink-and-perturb applied"), +- Err(e) => warn!(epoch, "Shrink-and-perturb failed (non-fatal): {e}"), +- } +- } +- } ++ // D3/N3: Periodic shrink-and-perturb replaced by health-triggered path in ++ // process_epoch_boundary. No periodic trigger here. + + // ── Phase 3: Batched training from replay buffer ── + let phase3_start = std::time::Instant::now(); +@@ -630,9 +620,12 @@ impl DQNTrainer { + } else if util < 0.4 { + self.adaptive_gamma = (self.adaptive_gamma - 0.01).max(0.90); + } +- // Wire adaptive gamma into GPU trainer for C51 Bellman projection ++ // C4/P4: Wire adaptive gamma into GPU trainer with temporal-coupled adjustment. ++ // apply_adaptive_gamma reads health+regime from ISV, adjusts, clamps [0.9, 0.995], ++ // and caches last_gamma_eff for HEALTH_DIAG. + if let Some(ref mut fused) = self.fused_ctx { +- fused.set_adaptive_gamma(self.adaptive_gamma); ++ let gamma_scheduled = self.adaptive_gamma as f32; ++ fused.apply_adaptive_gamma(gamma_scheduled); + } + info!("G4 gamma={:.3} util={:.2}", self.adaptive_gamma, util); + } +@@ -723,8 +716,10 @@ impl DQNTrainer { + if let Some(eval_gamma) = result.gamma { + let blended = 0.9 * self.adaptive_gamma + 0.1 * eval_gamma; + self.adaptive_gamma = blended.clamp(0.85, 0.98); ++ // C4/P4: apply temporal-coupled adjustment on top of EMA blend. + if let Some(ref mut fused) = self.fused_ctx { +- fused.set_adaptive_gamma(self.adaptive_gamma); ++ let gamma_scheduled = self.adaptive_gamma as f32; ++ fused.apply_adaptive_gamma(gamma_scheduled); + } + } + +@@ -1783,6 +1778,338 @@ impl DQNTrainer { + })?; + } + ++ // ── LearningHealth computation (Layer 1) ────────────────────────────── ++ // Source signal adaptations vs. plan spec (A3 known deviations): ++ // - q_gap: mean of per_branch_q_gap_ema (4 branches) via fused_ctx — ++ // NOTE: per_branch_q_gap_ema() performs a synchronous stream flush + DtoH ++ // (cuStreamSynchronize + cuMemcpyDtoH_v2). Acceptable at epoch boundary ++ // because it runs once per epoch; DO NOT move into a per-step loop. ++ // - q_var: epoch_q_max - epoch_q_min proxy instead of QValueStatsResult.q_variance ++ // (q_variance is not separately tracked as a scalar field; range is a sufficient proxy). ++ // - grad_consistency: scalar proxy (successive grad_norm delta) instead of full Adam ++ // vector cosine similarity — per-component Adam buffers prevent flat-vector approach. ++ // - spectral_gap: coarse max/min ratio on q_readback_pinned instead of full cuSOLVER SVD. ++ // - ens_disagreement: 0.1 placeholder until D6/N6 wires real ensemble pairwise Q-gap. ++ { ++ use crate::cuda_pipeline::learning_health::HealthComponents; ++ ++ // Per-branch Q-gap mean as q_gap signal (4 values, epoch-scoped DtoH). ++ let q_gap_signal: f32 = if let Some(ref fused) = self.fused_ctx { ++ let gaps = fused.per_branch_q_gap_ema(); ++ ((gaps[0] + gaps[1] + gaps[2] + gaps[3]) * 0.25).max(0.0) ++ } else { ++ 0.0 ++ }; ++ ++ // q_var proxy: range (max - min) of per-step Q values this epoch. If no ++ // steps ran (epoch_q_min still INFINITY), fall back to 0.0 rather than ++ // propagating NaN/-inf into health normalization. ++ let q_var_signal = if self.epoch_q_min.is_finite() && self.epoch_q_max.is_finite() { ++ (self.epoch_q_max - self.epoch_q_min).max(0.0) ++ } else { ++ 0.0 ++ }; ++ ++ // Atom utilization from q_readback_pinned[6] (one-step lag). ++ let atom_util = self.fused_ctx.as_ref() ++ .map(|f| f.read_atom_utilization()) ++ .unwrap_or(0.0); ++ ++ // Spectral gap proxy from q_readback_pinned[7..7+total_actions]. ++ let spectral_gap = self.fused_ctx.as_ref() ++ .map(|f| f.compute_q_spectral_gap()) ++ .unwrap_or(1.0); ++ ++ let grad_norm = avg_grad.abs(); ++ let ens_disagreement = 0.1f32; // placeholder — D6/N6 task will replace with real ensemble pairwise Q-gap ++ ++ // Update EMAs. ++ self.health_ema.update(q_gap_signal, q_var_signal, grad_norm); ++ let grad_consistency = self.health_ema.update_grad_consistency(grad_norm); ++ ++ let raw = HealthComponents { ++ q_gap: self.health_ema.q_gap_ema, ++ q_var: self.health_ema.q_var_ema, ++ atom_util, ++ grad_norm: self.health_ema.grad_norm_ema, ++ ens_disagreement, ++ grad_consistency, ++ spectral_gap, ++ }; ++ let health_value = self.learning_health.update(&raw); ++ ++ // Cache for later tasks (A4 logging, B1, B4, G4 etc.) ++ self.last_q_gap = Some(raw.q_gap); ++ self.last_q_var = Some(raw.q_var); ++ self.last_grad_norm = Some(raw.grad_norm); ++ self.last_ens_disagreement = Some(raw.ens_disagreement); ++ self.last_atom_util = Some(raw.atom_util); ++ ++ // Broadcast to GPU at ISV[LEARNING_HEALTH_INDEX=12]. ++ if let Some(ref fused) = self.fused_ctx { ++ fused.write_isv_signal_at( ++ crate::cuda_pipeline::gpu_dqn_trainer::LEARNING_HEALTH_INDEX, ++ health_value, ++ ); ++ } ++ ++ // B1/G2: propagate last cql_alpha_eff for logging. ++ if let Some(ref fused) = self.fused_ctx { ++ self.last_cql_alpha_eff = Some(fused.last_cql_alpha_eff()); ++ } ++ ++ // B2/G3: propagate last tau_eff for logging. ++ if let Some(ref fused) = self.fused_ctx { ++ self.last_tau_eff = Some(fused.last_tau_eff()); ++ } ++ ++ // B3/G4: propagate SARSA tau factor for logging. ++ if let Some(ref fused) = self.fused_ctx { ++ self.last_sarsa_tau_factor = Some(fused.last_sarsa_tau_factor()); ++ } ++ ++ // B4/G5: propagate adaptive gradient budgets for logging. ++ if let Some(ref fused) = self.fused_ctx { ++ self.last_iqn_budget_eff = Some(fused.last_iqn_budget_eff()); ++ self.last_cql_budget_eff = Some(fused.last_cql_budget_eff()); ++ self.last_c51_budget_eff = Some(fused.last_c51_budget_eff()); ++ } ++ ++ // C4/P4: propagate adaptive gamma for logging. ++ if let Some(ref fused) = self.fused_ctx { ++ self.last_gamma_eff = Some(fused.last_gamma_eff()); ++ } ++ ++ // D1/N1: snapshot high-health checkpoints and pull toward best during collapse. ++ if let Some(ref mut fused) = self.fused_ctx { ++ let q_gap_for_snapshot = self.last_q_gap.unwrap_or(0.0); ++ let _ = fused.maybe_snapshot_qnet(health_value, q_gap_for_snapshot, epoch as u32); ++ if health_value < 0.4 { ++ let _ = fused.apply_distillation(); ++ } ++ self.last_distill_active = Some(fused.last_distill_active()); ++ } ++ ++ // D2/N2: Q-gap barrier constraint (scalar-only — informational loss). ++ // min_required = 0.05 × health; barrier_loss = 0.5 × max(0, min_required − q_gap)². ++ // Drives no gradient in this task; surfaces as HEALTH_DIAG `barrier=...` for visibility. ++ { ++ let q_gap_current = self.health_ema.q_gap_ema; ++ let min_required = 0.05_f32 * health_value; ++ let barrier = (min_required - q_gap_current).max(0.0); ++ self.last_barrier_loss = Some(0.5 * barrier * barrier); ++ } ++ ++ // D3/N3: Health-triggered plasticity injection. ++ // Track consecutive unhealthy epochs; after 3, run shrink_perturb to break ++ // frozen representations and reset the counter. ++ { ++ let is_unhealthy = health_value < 0.3; ++ if is_unhealthy { ++ self.unhealthy_epoch_count += 1; ++ } else { ++ self.unhealthy_epoch_count = 0; ++ } ++ ++ if self.unhealthy_epoch_count >= 3 { ++ tracing::info!( ++ "Plasticity injection triggered: health={:.2} for {} consecutive epochs", ++ health_value, self.unhealthy_epoch_count ++ ); ++ if let Some(ref mut fused) = self.fused_ctx { ++ let alpha = self.hyperparams.shrink_perturb_alpha as f32; ++ let sigma = self.hyperparams.shrink_perturb_sigma as f32; ++ match fused.shrink_and_perturb(alpha, sigma) { ++ Ok(()) => tracing::info!(epoch, alpha, sigma, "D3/N3 health-triggered shrink_and_perturb applied"), ++ Err(e) => tracing::warn!(epoch, "D3/N3 shrink_and_perturb failed (non-fatal): {e}"), ++ } ++ } ++ self.unhealthy_epoch_count = 0; // reset after trigger ++ self.last_plasticity_ready = Some(false); // cooldown state ++ } else { ++ self.last_plasticity_ready = Some(true); ++ } ++ } ++ ++ // D5/N5: Information Bottleneck penalty (scalar-only — informational). ++ // IB penalizes lack of state-dependence in Q (low variance across batch). ++ // Minimum viable: host-side scalar computed from cached q_var_ema, surfaced ++ // via HEALTH_DIAG as `ib={..}`. Full backward-through-variance is deferred. ++ { ++ let ib_weight = 0.05_f32 * (1.0 - health_value).max(0.0); ++ let min_q_var = 0.01_f32; ++ let q_var_current = self.health_ema.q_var_ema; ++ let ib_penalty_raw = (min_q_var - q_var_current).max(0.0); ++ self.last_ib_penalty = Some(ib_weight * ib_penalty_raw); ++ } ++ ++ // D6/N6: Ensemble as collapse oracle. ++ // High ensemble_collapse_score ⇒ ensembles agree (often on uniform during collapse). ++ // When score > 0.8 AND we haven't already triggered N3 recently, force a plasticity ++ // injection — this shortcut catches collapse even when the moving-average of ++ // individual health inputs hasn't accumulated to 3 unhealthy epochs yet. ++ // Uses last_ens_disagreement (Path C scalar proxy — populated by A3 readback). ++ { ++ use crate::cuda_pipeline::learning_health::smoothstep; ++ let mean_pairwise_gap = self.last_ens_disagreement.unwrap_or(0.1); ++ let ensemble_collapse_score = 1.0 - smoothstep(0.01, 0.1, mean_pairwise_gap); ++ self.last_ensemble_collapse_score = Some(ensemble_collapse_score); ++ ++ if ensemble_collapse_score > 0.8 && self.unhealthy_epoch_count < 3 { ++ tracing::info!( ++ "Ensemble collapse detected (score={:.2}): forcing plasticity injection", ++ ensemble_collapse_score ++ ); ++ if let Some(ref mut fused) = self.fused_ctx { ++ let alpha = self.hyperparams.shrink_perturb_alpha as f32; ++ let sigma = self.hyperparams.shrink_perturb_sigma as f32; ++ match fused.shrink_and_perturb(alpha, sigma) { ++ Ok(()) => { ++ self.unhealthy_epoch_count = 0; ++ self.last_plasticity_ready = Some(false); ++ } ++ Err(e) => tracing::warn!("D6/N6 shrink_and_perturb failed: {e}"), ++ } ++ } ++ } ++ } ++ ++ // D7/N7: Contrarian override state machine. ++ // When WinRate < 0.40 for 5 consecutive epochs AND health < 0.3, flip to argmin ++ // for 2 epochs to escape the mirror-image attractor. WinRate >= 0.45 resets. ++ // Note: last_epoch_win_rate is carried from previous epoch's financials computation ++ // (which runs after process_epoch_boundary). Epoch 0 defaults to 0.5 (no trigger). ++ { ++ let winrate = self.last_epoch_win_rate; ++ ++ if winrate < 0.40 { ++ self.low_winrate_count += 1; ++ } else if winrate >= 0.45 { ++ self.contrarian_active = false; ++ self.low_winrate_count = 0; ++ } ++ ++ if self.low_winrate_count >= 5 && health_value < 0.3 && !self.contrarian_active { ++ tracing::warn!( ++ "CONTRARIAN OVERRIDE activated: WinRate={:.1}% for {} epochs, health={:.2}", ++ winrate * 100.0, self.low_winrate_count, health_value ++ ); ++ self.contrarian_active = true; ++ self.contrarian_remaining_epochs = 2; ++ } ++ ++ if self.contrarian_active { ++ if self.contrarian_remaining_epochs == 0 { ++ tracing::info!( ++ "CONTRARIAN OVERRIDE deactivated: window expired" ++ ); ++ self.contrarian_active = false; ++ } else { ++ self.contrarian_remaining_epochs -= 1; ++ } ++ } ++ ++ self.last_contrarian_active = Some(self.contrarian_active); ++ ++ // D7/N7 Part B: forward the flag to the experience collector so the ++ // kernel negates Q values during Boltzmann sampling. ++ if let Some(ref mut collector) = self.gpu_experience_collector { ++ collector.set_contrarian_active(self.contrarian_active); ++ } ++ } ++ ++ // D8/N8: Meta-Q collapse predictor — leading indicator, NOT fed back into LearningHealth. ++ { ++ use crate::cuda_pipeline::meta_q_network::{MetaQInput, META_Q_LOOKAHEAD}; ++ ++ // Build input from current health components. regime_stability read via fused_ctx. ++ let (_, regime_stability) = if let Some(ref fused) = self.fused_ctx { ++ fused.read_isv_regime_via_trainer() ++ } else { ++ (0.5_f32, 0.5_f32) ++ }; ++ ++ let meta_q_input = MetaQInput { ++ q_gap_ema: self.health_ema.q_gap_ema, ++ log_grad_norm: self.health_ema.grad_norm_ema.max(1e-3).ln(), ++ atom_util: self.last_atom_util.unwrap_or(0.0), ++ ens_disagreement: self.last_ens_disagreement.unwrap_or(0.1), ++ regime_stability, ++ spectral_gap_norm: self.learning_health.components.spectral_gap_norm, ++ grad_consistency: self.learning_health.components.grad_consistency_norm, ++ }; ++ ++ // Predict collapse probability ++ let meta_q_pred = self.meta_q.predict(&meta_q_input); ++ self.last_meta_q_pred = Some(meta_q_pred); ++ ++ // Push current (input, epoch) into pending — label assigned K epochs later ++ let cur_epoch_u32 = self.current_epoch as u32; ++ self.pending_meta_q_samples.push_back((meta_q_input, cur_epoch_u32)); ++ ++ // Track health history for label resolution ++ if self.health_history.len() >= 2 * META_Q_LOOKAHEAD { ++ self.health_history.pop_front(); ++ } ++ self.health_history.push_back((cur_epoch_u32, health_value)); ++ ++ // Resolve labels for samples that are K epochs old ++ while let Some((pending_input, pending_epoch)) = self.pending_meta_q_samples.front().cloned() { ++ if cur_epoch_u32 < pending_epoch + META_Q_LOOKAHEAD as u32 { break; } ++ let collapsed = self.health_history.iter() ++ .filter(|(ep, _)| *ep > pending_epoch && *ep <= pending_epoch + META_Q_LOOKAHEAD as u32) ++ .any(|(_, h)| *h < 0.3); ++ let label = if collapsed { 1.0 } else { 0.0 }; ++ self.meta_q.push_sample(pending_input, label); ++ self.pending_meta_q_samples.pop_front(); ++ } ++ ++ // Train one step ++ self.meta_q.train_step(); ++ } ++ ++ // HEALTH_DIAG: components are [0, 1] normalized. effective = hyperparams after health-adaptation. novels = mechanism states. ++ tracing::info!( ++ "HEALTH_DIAG[{}]: health={:.2} components [q_gap={:.2} q_var={:.2} atoms={:.2} grad_stable={:.2} ens_agree={:.2} grad_cos={:.2} spectral={:.2}] effective [cql_alpha={:.4} iqn_budget={:.2} cql_budget={:.2} c51_budget={:.2} tau={:.5} sarsa_tau={:.2} gamma={:.3} cf_ratio={:.2}] novels [distill={} barrier={:.3} plasticity={} ib={:.3} ensemble_collapse={:.2} contrarian={} meta_q_pred={:.2}]", ++ epoch, ++ health_value, ++ self.learning_health.components.q_gap_norm, ++ self.learning_health.components.q_var_norm, ++ self.learning_health.components.atom_util_norm, ++ self.learning_health.components.grad_stable, ++ self.learning_health.components.ens_agree, ++ self.learning_health.components.grad_consistency_norm, ++ self.learning_health.components.spectral_gap_norm, ++ self.last_cql_alpha_eff.unwrap_or(0.0), ++ self.last_iqn_budget_eff.unwrap_or(0.40), ++ self.last_cql_budget_eff.unwrap_or(0.10), ++ self.last_c51_budget_eff.unwrap_or(0.45), ++ self.last_tau_eff.unwrap_or(0.005), ++ self.last_sarsa_tau_factor.unwrap_or(1.0), ++ self.last_gamma_eff.unwrap_or(0.99), ++ self.last_cf_ratio_eff.unwrap_or(0.5), ++ if self.last_distill_active.unwrap_or(false) { "on" } else { "off" }, ++ self.last_barrier_loss.unwrap_or(0.0), ++ if self.last_plasticity_ready.unwrap_or(true) { "ready" } else { "cooldown" }, ++ self.last_ib_penalty.unwrap_or(0.0), ++ self.last_ensemble_collapse_score.unwrap_or(0.0), ++ if self.last_contrarian_active.unwrap_or(false) { "on" } else { "off" }, ++ self.last_meta_q_pred.unwrap_or(0.0), ++ ); ++ ++ // C1/P1: propagate health to GPU replay buffer for diversity-weighted priorities. ++ { ++ let mut agent = self.agent.write().await; ++ agent.primary_dqn_mut().memory.gpu.set_learning_health(health_value); ++ } ++ ++ // D4/N4: propagate health to GPU experience collector for adaptive CF curriculum. ++ if let Some(ref mut collector) = self.gpu_experience_collector { ++ collector.set_learning_health(health_value); ++ self.last_cf_ratio_eff = Some(collector.last_cf_ratio_eff()); ++ } ++ } ++ + // Feed per-step Q-value extremes into monitor (real min/max, not single-sample). + if self.epoch_q_min.is_finite() && self.epoch_q_max.is_finite() { + tracing::debug!( +@@ -2375,6 +2702,8 @@ impl DQNTrainer { + } + + epoch_max_dd = financials.max_drawdown; ++ // D7/N7: carry win_rate into next epoch's process_epoch_boundary. ++ self.last_epoch_win_rate = financials.win_rate as f32; + financials.sharpe + }; +