From 87968c5081d0686738eb8ee6927422241cdc18f3 Mon Sep 17 00:00:00 2001 From: Sonic Shih Date: Tue, 4 Aug 2026 00:40:21 +0800 Subject: [PATCH 1/4] feat(alpha): gate MCTS on immutable CEX baselines (#600) --- rust_hft/alpha-harness/app/src/mission.rs | 255 +++++++++++++++++++--- 1 file changed, 220 insertions(+), 35 deletions(-) diff --git a/rust_hft/alpha-harness/app/src/mission.rs b/rust_hft/alpha-harness/app/src/mission.rs index b35876027..04046edb8 100644 --- a/rust_hft/alpha-harness/app/src/mission.rs +++ b/rust_hft/alpha-harness/app/src/mission.rs @@ -3,13 +3,12 @@ use crate::{ data_mission, }; use alpha_domain::{ - canonical_json_hash, CexGpPolicyV1, CexResearchContentRefV1, MissionStatus, ResearchMission, + canonical_json_hash, CexBaselineArtifactV1, CexBaselineGateV1, CexBaselineModelKindV1, + CexFactorBankRevisionV2, CexGpPolicyV1, CexResearchContentRefV1, CexResearchMissionArtifactV1, + MissionStatus, ResearchMission, }; use alpha_engine::{ - engines::{ - CexMctsSearchIdentityV1, GeneticProgrammingEngine, MctsEngine, OfflineRlEngine, - OfflineTrace, - }, + engines::{GeneticProgrammingEngine, OfflineRlEngine, OfflineTrace}, evaluation::prepare_dataset, formula_evaluator::FormulaEvaluator, learning::{close_learning_loop, FailureCritic, LearningConfig}, @@ -24,6 +23,11 @@ pub(crate) const BAYESIAN_WINDOW_SEARCH_LIVE_CAPABILITY_ERROR: &str = "Bayesian window search is research-only and cannot produce live-executable formulas"; pub(crate) const OFFLINE_RL_LIVE_CAPABILITY_ERROR: &str = "Offline RL search is research-only and cannot produce live-executable formulas"; +const CEX_FACTOR_BANK_REGISTRY_KIND: &str = "cex_factor_bank"; +const CEX_BASELINE_RIDGE_REGISTRY_KIND: &str = "cex_baseline_ridge"; +const CEX_BASELINE_CART_REGISTRY_KIND: &str = "cex_baseline_cart"; +const CEX_BASELINE_GATE_REGISTRY_KIND: &str = "cex_baseline_gate"; +const CEX_MCTS_CANDIDATE_SPACE_ID: &str = "live-factor-ast-add-secondary-v1"; #[derive(Debug, Clone, serde::Serialize)] #[serde(rename_all = "snake_case")] @@ -90,6 +94,7 @@ fn execute_mission_inner( )?; let labels = manifest.evaluation_label_spec()?; let protocol = args.dataset.validation.evaluation_protocol(&labels)?; + let evaluation_protocol_hash = protocol.content_hash()?; let dataset = prepare_dataset(rows, &protocol)?; let research_context = dataset.engine_context(); let research_dataset_sha256 = canonical_json_hash(&research_context.rows())?; @@ -105,18 +110,20 @@ fn execute_mission_inner( id: format!("cex-walk-forward-partition-{walk_forward_partition_sha256}"), content_sha256: walk_forward_partition_sha256, }; + if matches!(args.engine, EngineChoice::Mcts) { + validate_mcts_baseline_gate( + &store, + &mission, + &research_dataset, + &walk_forward_partition, + &evaluation_protocol_hash, + )?; + bail!( + "MCTS candidate space {CEX_MCTS_CANDIDATE_SPACE_ID} is not authorized for CEX Factor Bank missions; Factor-Bank subset adapter (#601) is required" + ); + } let evaluator = FormulaEvaluator::for_mission(&mission).map_err(anyhow::Error::msg)?; - let evaluator_config_hash = - canonical_json_hash(&evaluator.config_evidence().map_err(anyhow::Error::msg)?)?; - let evaluation_protocol_hash = protocol.content_hash()?; - let proposal_engine = build_engine( - args, - &dataset, - &mission, - &evaluation_protocol_hash, - &evaluator_config_hash, - governed_gp, - )?; + let proposal_engine = build_engine(args, &dataset, &mission, governed_gp)?; let mut kernel = AutoResearchKernel::new(&mut store, proposal_engine, evaluator); let outcome = kernel.run( &args.mission_id, @@ -142,6 +149,155 @@ fn execute_mission_inner( }) } +fn validate_mcts_baseline_gate( + store: &AlphaStore, + mission: &ResearchMission, + research_dataset: &CexResearchContentRefV1, + walk_forward_partition: &CexResearchContentRefV1, + evaluation_protocol_hash: &str, +) -> anyhow::Result<()> { + let gate_id = mission + .baseline_artifact_id + .as_deref() + .context("MCTS requires a baseline gate identity")?; + let gate_revision = store + .get_registry_revision(gate_id) + .with_context(|| format!("MCTS baseline gate registry revision {gate_id} is missing"))?; + if gate_revision.registry_kind != CEX_BASELINE_GATE_REGISTRY_KIND + || gate_revision.revision_id != gate_id + { + bail!("MCTS baseline gate registry kind or identity is invalid"); + } + let gate: CexBaselineGateV1 = serde_json::from_value(gate_revision.payload.clone()) + .context("MCTS baseline gate payload is invalid")?; + gate.validate().map_err(anyhow::Error::msg)?; + + let factor_bank_revision = store + .get_registry_revision(&gate.factor_bank_revision_id) + .with_context(|| { + format!( + "MCTS Factor Bank registry revision {} is missing", + gate.factor_bank_revision_id + ) + })?; + if factor_bank_revision.registry_kind != CEX_FACTOR_BANK_REGISTRY_KIND + || factor_bank_revision.revision_id != gate.factor_bank_revision_id + || factor_bank_revision.parent_revision_id.as_deref() + != Some(mission.dataset_manifest_id.as_str()) + { + bail!("MCTS Factor Bank registry binding drifted"); + } + if gate_revision.parent_revision_id.as_deref() != Some(gate.factor_bank_revision_id.as_str()) { + bail!("MCTS baseline gate parent revision drifted"); + } + let factor_bank: CexFactorBankRevisionV2 = + serde_json::from_value(factor_bank_revision.payload.clone()) + .context("MCTS Factor Bank payload is invalid")?; + factor_bank.validate().map_err(anyhow::Error::msg)?; + if factor_bank.research_dataset != *research_dataset + || factor_bank.walk_forward_partition != *walk_forward_partition + { + bail!("MCTS Factor Bank research identity drifted"); + } + let control_revision = store + .get_registry_revision(&gate.mission_id) + .with_context(|| { + format!( + "MCTS source research mission registry revision {} is missing", + gate.mission_id + ) + })?; + if control_revision.registry_kind != "cex_research_mission" + || control_revision.revision_id != gate.mission_id + || control_revision.parent_revision_id.as_deref() + != Some(mission.dataset_manifest_id.as_str()) + { + bail!("MCTS source research mission registry binding drifted"); + } + let control_mission: CexResearchMissionArtifactV1 = + serde_json::from_value(control_revision.payload.clone()) + .context("MCTS source research mission payload is invalid")?; + control_mission.validate().map_err(anyhow::Error::msg)?; + if control_mission.semantic_id().map_err(anyhow::Error::msg)? != gate.mission_id { + bail!("MCTS source research mission identity drifted"); + } + + let ridge = match gate.ridge_artifact_id.as_deref() { + Some(ridge_id) => Some(read_mcts_baseline_artifact( + store, + ridge_id, + CEX_BASELINE_RIDGE_REGISTRY_KIND, + &factor_bank, + &control_mission, + research_dataset, + walk_forward_partition, + evaluation_protocol_hash, + CexBaselineModelKindV1::Ridge, + )?), + None => None, + }; + let cart = match gate.cart_artifact_id.as_deref() { + Some(cart_id) => Some(read_mcts_baseline_artifact( + store, + cart_id, + CEX_BASELINE_CART_REGISTRY_KIND, + &factor_bank, + &control_mission, + research_dataset, + walk_forward_partition, + evaluation_protocol_hash, + CexBaselineModelKindV1::ShallowCart, + )?), + None => None, + }; + gate.validate_binding(&factor_bank, ridge.as_ref(), cart.as_ref()) + .map_err(anyhow::Error::msg)?; + if !gate.passed { + bail!("MCTS baseline gate did not pass"); + } + Ok(()) +} + +#[allow(clippy::too_many_arguments)] +fn read_mcts_baseline_artifact( + store: &AlphaStore, + artifact_id: &str, + registry_kind: &str, + factor_bank: &CexFactorBankRevisionV2, + control_mission: &CexResearchMissionArtifactV1, + research_dataset: &CexResearchContentRefV1, + walk_forward_partition: &CexResearchContentRefV1, + evaluation_protocol_hash: &str, + model_kind: CexBaselineModelKindV1, +) -> anyhow::Result { + let revision = store.get_registry_revision(artifact_id).with_context(|| { + format!("MCTS baseline artifact registry revision {artifact_id} is missing") + })?; + if revision.registry_kind != registry_kind + || revision.revision_id != artifact_id + || revision.parent_revision_id.as_deref() != Some(factor_bank.revision_id.as_str()) + { + bail!("MCTS baseline artifact registry binding drifted"); + } + let artifact: CexBaselineArtifactV1 = serde_json::from_value(revision.payload.clone()) + .context("MCTS baseline artifact payload is invalid")?; + artifact.validate().map_err(anyhow::Error::msg)?; + artifact + .validate_binding(control_mission, factor_bank) + .map_err(anyhow::Error::msg)?; + if artifact.model_kind != model_kind + || artifact.mission_id != control_mission.semantic_id().map_err(anyhow::Error::msg)? + || artifact.factor_bank_revision_id != factor_bank.revision_id + || artifact.research_dataset != *research_dataset + || artifact.walk_forward_partition != *walk_forward_partition + || artifact.evaluation_policy.content_sha256 != evaluation_protocol_hash + || artifact.evaluation.protocol_binding()?.1 != evaluation_protocol_hash + { + bail!("MCTS baseline artifact identity drifted"); + } + Ok(artifact) +} + pub fn mission_status(args: MissionStatusArgs) -> anyhow::Result<()> { let store = AlphaStore::open(&args.db)?; let lineage = store.mission_lineage(&args.mission_id)?; @@ -191,8 +347,6 @@ fn build_engine( args: &RunMissionArgs, dataset: &alpha_engine::evaluation::PreparedDataset, mission: &ResearchMission, - evaluation_protocol_hash: &str, - evaluator_config_hash: &str, governed_gp: Option<(&CexGpPolicyV1, &str)>, ) -> anyhow::Result> { validate_live_mission_args(args)?; @@ -215,7 +369,6 @@ fn build_engine( let fields = fields.into_iter().collect::>(); validate_live_feature_fields(&fields)?; let primary = fields[0].clone(); - let secondary = fields.get(1).cloned().unwrap_or_else(|| primary.clone()); let engine: Box = match args.engine { EngineChoice::Gp => match governed_gp { Some((policy, candidate_namespace)) => { @@ -238,22 +391,8 @@ fn build_engine( .map_err(anyhow::Error::msg)?, ), }, - EngineChoice::Mcts => Box::new( - MctsEngine::new_live_bound( - args.seed, - primary.clone(), - secondary, - 1.414, - 5, - CexMctsSearchIdentityV1::for_mission( - mission, - evaluation_protocol_hash.to_string(), - evaluator_config_hash.to_string(), - args.max_new_iterations, - ) - .map_err(anyhow::Error::msg)?, - ) - .map_err(anyhow::Error::msg)?, + EngineChoice::Mcts => bail!( + "MCTS candidate space {CEX_MCTS_CANDIDATE_SPACE_ID} is not authorized for CEX Factor Bank missions; Factor-Bank subset adapter (#601) is required" ), EngineChoice::Bayesian => bail!(BAYESIAN_WINDOW_SEARCH_LIVE_CAPABILITY_ERROR), EngineChoice::OfflineRl => { @@ -337,6 +476,9 @@ fn live_event_domain_name(domain: LiveEventDomain) -> &'static str { mod tests { use super::*; use crate::cli::{DatasetArgs, ValidationArgs}; + use alpha_domain::{MissionCompletionPolicy, SearchBudget, ValidatorMode}; + use chrono::Utc; + use hft_research_manifest::ManifestId; use std::path::PathBuf; #[test] @@ -351,6 +493,49 @@ mod tests { ); } + #[test] + fn mcts_shared_fan_in_requires_an_immutable_baseline_gate() { + let store = AlphaStore::open_in_memory().unwrap(); + let mission = ResearchMission { + mission_id: "consumer-mission".to_string(), + objective: "test objective".to_string(), + hypothesis_scope: "test hypothesis".to_string(), + mutable_scope: vec!["factor_ast".to_string()], + dataset_manifest_id: ManifestId::new("dataset-1").unwrap(), + baseline_artifact_id: None, + validation_mode: ValidatorMode::MissionValidator, + validator_spec: serde_json::json!({}), + search_budget: SearchBudget { + max_candidates: 1, + max_expansions: 1, + max_tokens: 0, + max_seconds: 1, + }, + completion_policy: MissionCompletionPolicy::default(), + prompt_snapshot_id: None, + search_policy_snapshot_id: "policy-1".to_string(), + status: MissionStatus::Pending, + terminal_reason: None, + created_at: Utc::now(), + updated_at: Utc::now(), + }; + let reference = |id: &str| CexResearchContentRefV1 { + id: id.to_string(), + content_sha256: "a".repeat(64), + }; + + let error = validate_mcts_baseline_gate( + &store, + &mission, + &reference("cex-research-dataset-1"), + &reference("cex-walk-forward-partition-1"), + &"b".repeat(64), + ) + .unwrap_err(); + + assert!(error.to_string().contains("baseline gate identity")); + } + #[test] fn mission_preflight_rejects_live_capability_before_opening_store() { let cases = [ From da92019357412b3dd02901d689438820c844843f Mon Sep 17 00:00:00 2001 From: Sonic Shih Date: Tue, 4 Aug 2026 05:33:24 +0800 Subject: [PATCH 2/4] fix(alpha): replay baseline evidence before MCTS (#600) --- rust_hft/alpha-harness/app/src/mission.rs | 33 ++++++++++++++++++----- 1 file changed, 27 insertions(+), 6 deletions(-) diff --git a/rust_hft/alpha-harness/app/src/mission.rs b/rust_hft/alpha-harness/app/src/mission.rs index 04046edb8..7df381aeb 100644 --- a/rust_hft/alpha-harness/app/src/mission.rs +++ b/rust_hft/alpha-harness/app/src/mission.rs @@ -4,12 +4,13 @@ use crate::{ }; use alpha_domain::{ canonical_json_hash, CexBaselineArtifactV1, CexBaselineGateV1, CexBaselineModelKindV1, - CexFactorBankRevisionV2, CexGpPolicyV1, CexResearchContentRefV1, CexResearchMissionArtifactV1, - MissionStatus, ResearchMission, + CexBaselinePolicyV1, CexFactorBankRevisionV2, CexGpPolicyV1, CexResearchContentRefV1, + CexResearchMissionArtifactV1, MissionStatus, ResearchMission, }; use alpha_engine::{ + baselines::verify_cex_baseline_artifact, engines::{GeneticProgrammingEngine, OfflineRlEngine, OfflineTrace}, - evaluation::prepare_dataset, + evaluation::{prepare_dataset, EngineContext}, formula_evaluator::FormulaEvaluator, learning::{close_learning_loop, FailureCritic, LearningConfig}, llm::{LlmConfig, LlmProposalEngine, OpenAiCompatibleClient}, @@ -114,6 +115,7 @@ fn execute_mission_inner( validate_mcts_baseline_gate( &store, &mission, + Some(&research_context), &research_dataset, &walk_forward_partition, &evaluation_protocol_hash, @@ -152,6 +154,7 @@ fn execute_mission_inner( fn validate_mcts_baseline_gate( store: &AlphaStore, mission: &ResearchMission, + research_context: Option<&EngineContext<'_>>, research_dataset: &CexResearchContentRefV1, walk_forward_partition: &CexResearchContentRefV1, evaluation_protocol_hash: &str, @@ -221,6 +224,9 @@ fn validate_mcts_baseline_gate( if control_mission.semantic_id().map_err(anyhow::Error::msg)? != gate.mission_id { bail!("MCTS source research mission identity drifted"); } + let baseline_policy = + CexBaselinePolicyV1::controlled_v1(control_mission.spec.policies.baseline.id.clone())?; + baseline_policy.validate_binding(&control_mission.spec.policies.baseline)?; let ridge = match gate.ridge_artifact_id.as_deref() { Some(ridge_id) => Some(read_mcts_baseline_artifact( @@ -229,6 +235,8 @@ fn validate_mcts_baseline_gate( CEX_BASELINE_RIDGE_REGISTRY_KIND, &factor_bank, &control_mission, + &baseline_policy, + research_context.context("MCTS baseline artifact verification context is missing")?, research_dataset, walk_forward_partition, evaluation_protocol_hash, @@ -243,6 +251,8 @@ fn validate_mcts_baseline_gate( CEX_BASELINE_CART_REGISTRY_KIND, &factor_bank, &control_mission, + &baseline_policy, + research_context.context("MCTS baseline artifact verification context is missing")?, research_dataset, walk_forward_partition, evaluation_protocol_hash, @@ -250,8 +260,14 @@ fn validate_mcts_baseline_gate( )?), None => None, }; - gate.validate_binding(&factor_bank, ridge.as_ref(), cart.as_ref()) - .map_err(anyhow::Error::msg)?; + gate.validate_binding( + &control_mission, + &baseline_policy, + &factor_bank, + ridge.as_ref(), + cart.as_ref(), + ) + .map_err(anyhow::Error::msg)?; if !gate.passed { bail!("MCTS baseline gate did not pass"); } @@ -265,6 +281,8 @@ fn read_mcts_baseline_artifact( registry_kind: &str, factor_bank: &CexFactorBankRevisionV2, control_mission: &CexResearchMissionArtifactV1, + baseline_policy: &CexBaselinePolicyV1, + research_context: &EngineContext<'_>, research_dataset: &CexResearchContentRefV1, walk_forward_partition: &CexResearchContentRefV1, evaluation_protocol_hash: &str, @@ -283,7 +301,7 @@ fn read_mcts_baseline_artifact( .context("MCTS baseline artifact payload is invalid")?; artifact.validate().map_err(anyhow::Error::msg)?; artifact - .validate_binding(control_mission, factor_bank) + .validate_binding(control_mission, baseline_policy, factor_bank) .map_err(anyhow::Error::msg)?; if artifact.model_kind != model_kind || artifact.mission_id != control_mission.semantic_id().map_err(anyhow::Error::msg)? @@ -295,6 +313,8 @@ fn read_mcts_baseline_artifact( { bail!("MCTS baseline artifact identity drifted"); } + verify_cex_baseline_artifact(research_context, factor_bank, &artifact) + .map_err(anyhow::Error::msg)?; Ok(artifact) } @@ -527,6 +547,7 @@ mod tests { let error = validate_mcts_baseline_gate( &store, &mission, + None, &reference("cex-research-dataset-1"), &reference("cex-walk-forward-partition-1"), &"b".repeat(64), From d863d422e28abc0d02ba8e4424b925565871c9a1 Mon Sep 17 00:00:00 2001 From: Sonic Shih Date: Tue, 4 Aug 2026 07:44:36 +0800 Subject: [PATCH 3/4] test(alpha): cover CEX MCTS baseline fan-in gate --- .../alpha-harness/app/src/mission_runner.rs | 164 ++++++++++++++++++ 1 file changed, 164 insertions(+) diff --git a/rust_hft/alpha-harness/app/src/mission_runner.rs b/rust_hft/alpha-harness/app/src/mission_runner.rs index 019393e6a..3f0121805 100644 --- a/rust_hft/alpha-harness/app/src/mission_runner.rs +++ b/rust_hft/alpha-harness/app/src/mission_runner.rs @@ -1489,6 +1489,170 @@ mod tests { std::fs::remove_dir_all(fixture.root).unwrap(); } + #[test] + fn mission_execute_then_mcts_refuses_after_replaying_passing_baselines() { + let mut fixture = fixture("mcts-baseline-gate"); + fixture.mission.spec.feature_fields = vec!["book_imbalance".to_string()]; + rewrite_features(&mut fixture, |row| { + let direction = row.label.signum(); + row.features.insert("book_imbalance".to_string(), direction); + row.label = direction * 0.001; + }); + + execute(fixture.args.clone()).unwrap(); + + let results = fixture.args.work_dir.join("results"); + let db = results.join("alpha.duckdb"); + let producer_id = fixture.mission.semantic_id().unwrap(); + let mut store = AlphaStore::open(&db).unwrap(); + let producer = store.get_mission(&producer_id).unwrap(); + let gate: serde_json::Value = + serde_json::from_slice(&std::fs::read(results.join("baseline-gate.json")).unwrap()) + .unwrap(); + assert_eq!(gate["passed"], true); + let gate_id = gate["gate_id"].as_str().unwrap().to_string(); + let gate_revision = store.get_registry_revision(&gate_id).unwrap(); + assert_eq!(gate_revision.registry_kind, "cex_baseline_gate"); + assert_eq!( + gate_revision.parent_revision_id, + gate["factor_bank_revision_id"].as_str().map(str::to_owned) + ); + + let mut consumer = producer; + consumer.mission_id = format!("{producer_id}-mcts-consumer"); + consumer.baseline_artifact_id = Some(gate_id); + consumer.status = MissionStatus::Pending; + consumer.terminal_reason = None; + consumer.created_at = Utc::now(); + consumer.updated_at = consumer.created_at; + let consumer_id = consumer.mission_id.clone(); + store.create_mission(&consumer).unwrap(); + drop(store); + + let args = RunMissionArgs { + db, + mission_id: consumer_id, + engine: EngineChoice::Mcts, + seed: fixture.mission.spec.search.seed, + feature_fields: fixture.mission.spec.feature_fields.clone(), + offline_trace: None, + max_new_iterations: Some(1), + dataset: DatasetArgs { + dataset_manifest: results.join("cex-replay-dataset-manifest.json"), + validation: ValidationArgs::from_protocol( + &fixture.mission.spec.evaluation_protocol, + ), + }, + }; + let error = mission::execute_mission(&args, false).unwrap_err(); + + assert!(error + .to_string() + .contains("Factor-Bank subset adapter (#601) is required")); + let store = AlphaStore::open(&args.db).unwrap(); + assert_eq!( + store.get_mission(&args.mission_id).unwrap().status, + MissionStatus::Pending + ); + assert!(store + .mission_lineage(&args.mission_id) + .unwrap() + .iterations + .is_empty()); + std::fs::remove_dir_all(fixture.root).unwrap(); + } + + #[test] + fn mcts_rejects_a_tampered_published_gate_before_transition() { + let mut fixture = fixture("mcts-tampered-gate"); + fixture.mission.spec.feature_fields = vec!["book_imbalance".to_string()]; + rewrite_features(&mut fixture, |row| { + let direction = row.label.signum(); + row.features.insert("book_imbalance".to_string(), direction); + row.label = direction * 0.001; + }); + execute(fixture.args.clone()).unwrap(); + + let results = fixture.args.work_dir.join("results"); + let db = results.join("alpha.duckdb"); + let producer_id = fixture.mission.semantic_id().unwrap(); + let mut store = AlphaStore::open(&db).unwrap(); + let producer = store.get_mission(&producer_id).unwrap(); + let gate_revision = store + .get_registry_revision( + serde_json::from_slice::( + &std::fs::read(results.join("baseline-gate.json")).unwrap(), + ) + .unwrap()["gate_id"] + .as_str() + .unwrap(), + ) + .unwrap(); + let mut tampered_gate: alpha_domain::CexBaselineGateV1 = + serde_json::from_value(gate_revision.payload.clone()).unwrap(); + tampered_gate.policy_hash = "0".repeat(64); + tampered_gate.gate_id.clear(); + tampered_gate.gate_id = format!( + "cex-baseline-gate-{}", + canonical_json_hash(&tampered_gate).unwrap() + ); + tampered_gate.validate().unwrap(); + let tampered_gate_id = tampered_gate.gate_id.clone(); + store + .put_registry_revision(&alpha_store::RegistryRevision { + revision_id: tampered_gate_id.clone(), + registry_kind: gate_revision.registry_kind, + asset_id: gate_revision.asset_id, + parent_revision_id: gate_revision.parent_revision_id, + payload: serde_json::to_value(&tampered_gate).unwrap(), + created_at: Utc::now(), + }) + .unwrap(); + + let mut consumer = producer; + consumer.mission_id = format!("{producer_id}-mcts-tampered"); + consumer.baseline_artifact_id = Some(tampered_gate_id); + consumer.status = MissionStatus::Pending; + consumer.terminal_reason = None; + consumer.created_at = Utc::now(); + consumer.updated_at = consumer.created_at; + let consumer_id = consumer.mission_id.clone(); + store.create_mission(&consumer).unwrap(); + drop(store); + + let args = RunMissionArgs { + db, + mission_id: consumer_id, + engine: EngineChoice::Mcts, + seed: fixture.mission.spec.search.seed, + feature_fields: fixture.mission.spec.feature_fields.clone(), + offline_trace: None, + max_new_iterations: Some(1), + dataset: DatasetArgs { + dataset_manifest: results.join("cex-replay-dataset-manifest.json"), + validation: ValidationArgs::from_protocol( + &fixture.mission.spec.evaluation_protocol, + ), + }, + }; + let error = mission::execute_mission(&args, false).unwrap_err(); + + assert!(error + .to_string() + .contains("baseline gate producer binding drifted")); + let store = AlphaStore::open(&args.db).unwrap(); + assert_eq!( + store.get_mission(&args.mission_id).unwrap().status, + MissionStatus::Pending + ); + assert!(store + .mission_lineage(&args.mission_id) + .unwrap() + .iterations + .is_empty()); + std::fs::remove_dir_all(fixture.root).unwrap(); + } + #[test] fn execute_rejects_mixed_hypothesis_targets_before_side_effects() { let mut fixture = fixture("mixed-hypothesis-targets"); From 71f749d765223f7ee3c65a67fb9e84422afe5eb6 Mon Sep 17 00:00:00 2001 From: Sonic Shih Date: Tue, 4 Aug 2026 08:06:14 +0800 Subject: [PATCH 4/4] fix(alpha): close CEX MCTS gate review gaps --- .../alpha-harness/app/src/mission_runner.rs | 81 ++++++++ .../alpha-harness/engine/src/engines/mcts.rs | 196 +----------------- .../alpha-harness/engine/src/engines/mod.rs | 2 +- 3 files changed, 93 insertions(+), 186 deletions(-) diff --git a/rust_hft/alpha-harness/app/src/mission_runner.rs b/rust_hft/alpha-harness/app/src/mission_runner.rs index 3f0121805..f96bbb798 100644 --- a/rust_hft/alpha-harness/app/src/mission_runner.rs +++ b/rust_hft/alpha-harness/app/src/mission_runner.rs @@ -1562,6 +1562,87 @@ mod tests { std::fs::remove_dir_all(fixture.root).unwrap(); } + #[test] + fn mcts_rejects_a_domain_valid_failed_gate_before_transition() { + let fixture = fixture("mcts-failed-empty-gate"); + execute(fixture.args.clone()).unwrap(); + + let results = fixture.args.work_dir.join("results"); + let db = results.join("alpha.duckdb"); + let producer_id = fixture.mission.semantic_id().unwrap(); + let mut store = AlphaStore::open(&db).unwrap(); + let producer = store.get_mission(&producer_id).unwrap(); + let factor_bank: CexFactorBankRevisionV2 = + serde_json::from_slice(&std::fs::read(results.join("factor-bank.json")).unwrap()) + .unwrap(); + assert!(factor_bank.entries.is_empty()); + let baseline_policy = + CexBaselinePolicyV1::controlled_v1(fixture.mission.spec.policies.baseline.id.clone()) + .unwrap(); + let gate: alpha_domain::CexBaselineGateV1 = + serde_json::from_slice(&std::fs::read(results.join("baseline-gate.json")).unwrap()) + .unwrap(); + assert!(!gate.passed); + assert_eq!( + gate.failure_codes, + vec![alpha_domain::CexBaselineFailureCodeV1::EmptyFactorBank] + ); + gate.validate().unwrap(); + assert_eq!( + gate, + alpha_domain::CexBaselineGateV1::empty_factor_bank( + &producer_id, + &baseline_policy, + &factor_bank + ) + .unwrap() + ); + let gate_revision = store.get_registry_revision(&gate.gate_id).unwrap(); + assert_eq!(gate_revision.registry_kind, "cex_baseline_gate"); + assert_eq!(gate_revision.payload, serde_json::to_value(&gate).unwrap()); + + let mut consumer = producer; + consumer.mission_id = format!("{producer_id}-mcts-failed"); + consumer.baseline_artifact_id = Some(gate.gate_id.clone()); + consumer.status = MissionStatus::Pending; + consumer.terminal_reason = None; + consumer.created_at = Utc::now(); + consumer.updated_at = consumer.created_at; + let consumer_id = consumer.mission_id.clone(); + store.create_mission(&consumer).unwrap(); + drop(store); + + let args = RunMissionArgs { + db, + mission_id: consumer_id, + engine: EngineChoice::Mcts, + seed: fixture.mission.spec.search.seed, + feature_fields: fixture.mission.spec.feature_fields.clone(), + offline_trace: None, + max_new_iterations: Some(1), + dataset: DatasetArgs { + dataset_manifest: results.join("cex-replay-dataset-manifest.json"), + validation: ValidationArgs::from_protocol( + &fixture.mission.spec.evaluation_protocol, + ), + }, + }; + let error = mission::execute_mission(&args, false).unwrap_err(); + + assert_eq!(error.to_string(), "MCTS baseline gate did not pass"); + let store = AlphaStore::open(&args.db).unwrap(); + assert_eq!( + store.get_mission(&args.mission_id).unwrap().status, + MissionStatus::Pending + ); + assert!(store + .mission_lineage(&args.mission_id) + .unwrap() + .iterations + .is_empty()); + std::fs::remove_dir_all(fixture.root).unwrap(); + } + #[test] fn mcts_rejects_a_tampered_published_gate_before_transition() { let mut fixture = fixture("mcts-tampered-gate"); diff --git a/rust_hft/alpha-harness/engine/src/engines/mcts.rs b/rust_hft/alpha-harness/engine/src/engines/mcts.rs index c8331141f..897b034e1 100644 --- a/rust_hft/alpha-harness/engine/src/engines/mcts.rs +++ b/rust_hft/alpha-harness/engine/src/engines/mcts.rs @@ -3,21 +3,13 @@ use crate::{ evaluation::ProposalContext, CandidateEvaluation, EngineProposal, HistoricalObservation, ProposalEngine, ProposalEngineCheckpoint, RemainingBudget, }; -use alpha_domain::{ - CandidateArtifact, EngineKind, MissionCompletionPolicy, ResearchMission, SearchBudget, - WALK_FORWARD_EVALUATOR_VERSION, -}; +use alpha_domain::{CandidateArtifact, EngineKind}; use hft_factor_dsl::{validate_live_formula, FactorAst, FactorOperator, FactorTerminal}; use hft_search_kernel::{backpropagate, select_expandable, validate_tree, UctNode, UctStats}; use serde::{Deserialize, Serialize}; use std::collections::{BTreeMap, BTreeSet}; -pub const MCTS_CHECKPOINT_VERSION: u32 = 3; -const CEX_MCTS_SEARCH_IDENTITY_VERSION: &str = "cex-mcts-search-identity-v1"; -const CEX_MCTS_CANDIDATE_SPACE_ID: &str = "live-factor-ast-add-secondary-v1"; -const CEX_MCTS_STOPPING_RULE_ID: &str = "min-kept-search-budget-or-max-new-iterations-pause-v1"; -const CEX_MCTS_SELECTION_RULE_ID: &str = "one-canonical-walk-forward-candidate-v1"; -const CEX_MCTS_UCT_TIE_BREAK_RULE_ID: &str = "last-child-on-total-cmp-tie-v1"; +pub const MCTS_CHECKPOINT_VERSION: u32 = 4; fn expansion_actions(live_only: bool) -> Vec { if live_only { @@ -70,106 +62,25 @@ impl UctNode for Node { #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] -struct MctsConfigV3 { +struct MctsConfigV4 { seed: u64, root_ast: FactorAst, secondary_field: String, exploration: f64, max_depth: usize, live_only: bool, - #[serde(default, skip_serializing_if = "Option::is_none")] - research_identity: Option, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] -struct MctsCheckpointV3 { - config: MctsConfigV3, +struct MctsCheckpointV4 { + config: MctsConfigV4, rng: DeterministicRng, nodes: Vec, candidates: BTreeMap, seen: BTreeSet, } -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct CexMctsSearchIdentityV1 { - pub schema_version: String, - pub mission_id: String, - pub dataset_manifest_id: String, - pub search_policy_snapshot_id: String, - pub search_budget: SearchBudget, - pub completion_policy: MissionCompletionPolicy, - pub max_new_iterations: Option, - pub evaluator_version: String, - pub evaluator_config_hash: String, - pub evaluation_protocol_hash: String, - pub candidate_space_id: String, - pub stopping_rule_id: String, - pub selection_rule_id: String, - pub uct_tie_break_rule_id: String, -} - -impl CexMctsSearchIdentityV1 { - pub fn for_mission( - mission: &ResearchMission, - evaluation_protocol_hash: String, - evaluator_config_hash: String, - max_new_iterations: Option, - ) -> Result { - mission.validate().map_err(|error| error.to_string())?; - let identity = Self { - schema_version: CEX_MCTS_SEARCH_IDENTITY_VERSION.to_string(), - mission_id: mission.mission_id.clone(), - dataset_manifest_id: mission.dataset_manifest_id.as_str().to_string(), - search_policy_snapshot_id: mission.search_policy_snapshot_id.clone(), - search_budget: mission.search_budget.clone(), - completion_policy: mission.completion_policy.clone(), - max_new_iterations, - evaluator_version: WALK_FORWARD_EVALUATOR_VERSION.to_string(), - evaluator_config_hash, - evaluation_protocol_hash, - candidate_space_id: CEX_MCTS_CANDIDATE_SPACE_ID.to_string(), - stopping_rule_id: CEX_MCTS_STOPPING_RULE_ID.to_string(), - selection_rule_id: CEX_MCTS_SELECTION_RULE_ID.to_string(), - uct_tie_break_rule_id: CEX_MCTS_UCT_TIE_BREAK_RULE_ID.to_string(), - }; - identity.validate()?; - Ok(identity) - } - - pub fn validate(&self) -> Result<(), String> { - if self.schema_version != CEX_MCTS_SEARCH_IDENTITY_VERSION - || self.mission_id.trim().is_empty() - || self.dataset_manifest_id.trim().is_empty() - || self.search_policy_snapshot_id.trim().is_empty() - || self.evaluator_version != WALK_FORWARD_EVALUATOR_VERSION - || self.candidate_space_id != CEX_MCTS_CANDIDATE_SPACE_ID - || self.stopping_rule_id != CEX_MCTS_STOPPING_RULE_ID - || self.selection_rule_id != CEX_MCTS_SELECTION_RULE_ID - || self.uct_tie_break_rule_id != CEX_MCTS_UCT_TIE_BREAK_RULE_ID - || self.max_new_iterations == Some(0) - || !is_sha256(&self.evaluator_config_hash) - || !is_sha256(&self.evaluation_protocol_hash) - { - return Err("invalid CEX MCTS research identity".to_string()); - } - self.search_budget - .validate() - .map_err(|error| error.to_string())?; - self.completion_policy - .validate() - .map_err(|error| error.to_string()) - } -} - -fn is_sha256(value: &str) -> bool { - value.len() == 64 - && value - .bytes() - .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) -} - #[derive(Debug, Clone, PartialEq)] pub struct MctsNodeSnapshot { pub node_id: usize, @@ -187,7 +98,6 @@ pub struct MctsEngine { max_depth: usize, secondary_field: String, live_only: bool, - research_identity: Option, nodes: Vec, candidates: BTreeMap, seen: BTreeSet, @@ -240,20 +150,6 @@ impl MctsEngine { ) } - pub fn new_live_bound( - seed: u64, - root_field: impl Into, - secondary_field: impl Into, - exploration: f64, - max_depth: usize, - research_identity: CexMctsSearchIdentityV1, - ) -> Result { - research_identity.validate()?; - let mut engine = Self::new_live(seed, root_field, secondary_field, exploration, max_depth)?; - engine.research_identity = Some(research_identity); - Ok(engine) - } - fn new_with_mode( seed: u64, root_field: impl Into, @@ -278,7 +174,6 @@ impl MctsEngine { max_depth, secondary_field, live_only, - research_identity: None, nodes: vec![Node { ast: FactorAst::Terminal(FactorTerminal::Field(root_field)), parent: None, @@ -309,8 +204,8 @@ impl MctsEngine { .collect() } - fn config(&self) -> Result { - Ok(MctsConfigV3 { + fn config(&self) -> Result { + Ok(MctsConfigV4 { seed: self.seed, root_ast: self .nodes @@ -322,7 +217,6 @@ impl MctsEngine { exploration: self.exploration, max_depth: self.max_depth, live_only: self.live_only, - research_identity: self.research_identity.clone(), }) } @@ -440,7 +334,7 @@ impl ProposalEngine for MctsEngine { } fn checkpoint(&self) -> Result { - let state = MctsCheckpointV3 { + let state = MctsCheckpointV4 { config: self.config()?, rng: self.rng.clone(), nodes: self.nodes.clone(), @@ -464,14 +358,14 @@ impl ProposalEngine for MctsEngine { if checkpoint.kind != EngineKind::Mcts || checkpoint.version != MCTS_CHECKPOINT_VERSION { return Err("MCTS checkpoint kind or version mismatch".to_string()); } - let state: MctsCheckpointV3 = serde_json::from_value(checkpoint.state.clone()) + let state: MctsCheckpointV4 = serde_json::from_value(checkpoint.state.clone()) .map_err(|error| format!("invalid MCTS checkpoint state: {error}"))?; state.validate()?; if state.config != self.config()? { return Err("MCTS checkpoint configuration mismatch".to_string()); } - let MctsCheckpointV3 { + let MctsCheckpointV4 { config, rng, nodes, @@ -484,7 +378,6 @@ impl ProposalEngine for MctsEngine { self.max_depth = config.max_depth; self.secondary_field = config.secondary_field; self.live_only = config.live_only; - self.research_identity = config.research_identity; self.nodes = nodes; self.candidates = candidates; self.seen = seen; @@ -492,7 +385,7 @@ impl ProposalEngine for MctsEngine { } } -impl MctsCheckpointV3 { +impl MctsCheckpointV4 { fn validate(&self) -> Result<(), String> { if self.config.secondary_field.trim().is_empty() || !self.config.exploration.is_finite() @@ -501,9 +394,6 @@ impl MctsCheckpointV3 { { return Err("invalid MCTS checkpoint configuration".to_string()); } - if let Some(identity) = &self.config.research_identity { - identity.validate()?; - } self.config .root_ast .validate() @@ -911,70 +801,6 @@ mod tests { assert_eq!(run(&dataset(false)), run(&dataset(true))); } - #[test] - fn bound_checkpoint_rejects_snapshot_or_search_identity_drift() { - fn identity(dataset: &str, max_candidates: usize) -> CexMctsSearchIdentityV1 { - let now = chrono::Utc::now(); - let mission = alpha_domain::ResearchMission { - mission_id: "mission-1".to_string(), - objective: "test".to_string(), - hypothesis_scope: "test".to_string(), - mutable_scope: vec!["factor_ast".to_string()], - dataset_manifest_id: hft_research_manifest::ManifestId::new(dataset).unwrap(), - baseline_artifact_id: None, - validation_mode: alpha_domain::ValidatorMode::MissionValidator, - validator_spec: serde_json::json!({}), - search_budget: alpha_domain::SearchBudget { - max_candidates, - max_expansions: 32, - max_tokens: 0, - max_seconds: 60, - }, - completion_policy: alpha_domain::MissionCompletionPolicy::default(), - prompt_snapshot_id: None, - search_policy_snapshot_id: "mcts-lob-pit-v1".to_string(), - status: alpha_domain::MissionStatus::Pending, - terminal_reason: None, - created_at: now, - updated_at: now, - }; - CexMctsSearchIdentityV1::for_mission(&mission, "a".repeat(64), "b".repeat(64), Some(1)) - .unwrap() - } - - let engine = - MctsEngine::new_live_bound(3, "best_bid", "best_ask", 1.4, 3, identity("dataset-a", 4)) - .unwrap(); - let checkpoint = engine.checkpoint().unwrap(); - assert_eq!(checkpoint.version, 3); - - let mut restored = - MctsEngine::new_live_bound(3, "best_bid", "best_ask", 1.4, 3, identity("dataset-a", 4)) - .unwrap(); - restored.restore_checkpoint(&checkpoint, &[]).unwrap(); - assert_eq!(restored.checkpoint().unwrap(), checkpoint); - - for mismatched in [identity("dataset-b", 4), identity("dataset-a", 5)] { - let mut restored = - MctsEngine::new_live_bound(3, "best_bid", "best_ask", 1.4, 3, mismatched).unwrap(); - assert!(restored.restore_checkpoint(&checkpoint, &[]).is_err()); - } - - let mut protocol_drift = identity("dataset-a", 4); - protocol_drift.evaluation_protocol_hash = "c".repeat(64); - let mut restored = - MctsEngine::new_live_bound(3, "best_bid", "best_ask", 1.4, 3, protocol_drift).unwrap(); - assert!(restored.restore_checkpoint(&checkpoint, &[]).is_err()); - - let mut forged_tie_rule = checkpoint.clone(); - forged_tie_rule.state["config"]["research_identity"]["uct_tie_break_rule_id"] = - serde_json::json!("first-child-on-tie"); - let mut restored = - MctsEngine::new_live_bound(3, "best_bid", "best_ask", 1.4, 3, identity("dataset-a", 4)) - .unwrap(); - assert!(restored.restore_checkpoint(&forged_tie_rule, &[]).is_err()); - } - #[test] fn abandon_removes_pending_candidate() { let mut engine = MctsEngine::new(3, "oi", "imbalance", 1.4, 3).unwrap(); diff --git a/rust_hft/alpha-harness/engine/src/engines/mod.rs b/rust_hft/alpha-harness/engine/src/engines/mod.rs index 873762edd..c672dd0a9 100644 --- a/rust_hft/alpha-harness/engine/src/engines/mod.rs +++ b/rust_hft/alpha-harness/engine/src/engines/mod.rs @@ -8,7 +8,7 @@ use hft_search_kernel::DeterministicRng; pub(crate) use bayesian::solve; pub use bayesian::BayesianOptimizerEngine; pub use gp::GeneticProgrammingEngine; -pub use mcts::{CexMctsSearchIdentityV1, MctsEngine, MctsNodeSnapshot, MCTS_CHECKPOINT_VERSION}; +pub use mcts::{MctsEngine, MctsNodeSnapshot, MCTS_CHECKPOINT_VERSION}; pub use offline_rl::{OfflineRlEngine, OfflineTrace}; #[cfg(test)]