diff --git a/PROJECT_STATUS.md b/PROJECT_STATUS.md index 477708d..c26311a 100644 --- a/PROJECT_STATUS.md +++ b/PROJECT_STATUS.md @@ -29,7 +29,7 @@ one. | Multi-need and structured MCP | Sequential and steer delivery, bounded ledger, structured JSON tools, cancellation, shared resolver | App Server simulator and one structured MCP cache-hit observation | **Implemented; offline validated; live calibration** | | Claim-level reuse | Validator-extracted claims, claim proofs, mixed planning, and bounded authoritative location, runtime-flow, and focused-test claims | Deterministic freshness, mutation, negative, projection, economics, and performance cases | **Implemented; offline validated** | | Verified changes | Isolated patch preparation, independent verifier, one repair, explicit journaled apply | Simulator and focused persistence, isolation, drift, and recovery tests | **Implemented; offline validated** | -| Codex role-profile control plane | Canonical Codex role definitions, bounded policies, immutable revisions, state-digest CAS, SQLite V14 persistence, audit records, explicit WorkerProfile projection, bounded digest-bound HTTP API, and local editor | Focused deterministic Rust and frontend tests; configuration-only boundary (no worker/session binding or lifecycle execution) | **Implemented; offline validated** | +| Codex role-profile control plane | Canonical Codex role definitions, bounded policies, immutable revisions, state-digest CAS, SQLite persistence, audit records, explicit WorkerProfile projection, bounded digest-bound HTTP API, local editor, and frozen session/worker/cache provenance | Focused deterministic Rust and frontend tests; no parent-owned lifecycle execution | **Implemented; offline validated** | | Codex development lifecycle orchestration | Evidence, patch, test, verification, approval, and apply primitives exist; the configurable parent-owned role lifecycle is not integrated | Component-level offline evidence only | **Pending** | | Other-host subagent configuration | Configuration-only interoperability is planned for Claude Code and Cursor, followed by OpenCode and Antigravity | Not available | **Pending** | | Multi-host orchestration | Execution remains Codex-only; non-Codex execution follows configuration interoperability, a host contract, and conformance evidence | Not available | **Pending** | @@ -61,8 +61,9 @@ provider-backed claim-authority observation exists. | Embedded React control plane | **Implemented; frontend and local end-to-end validation** | | Needs, proofs, claims, changes, runs, models, cache, settings, approvals | **Implemented; development interface** | | Canonical named Codex role-profile domain and revision store | **Implemented; offline validated** | -| Named role-profile HTTP/editor | **Implemented; offline validated; configuration-only** | -| Role-profile session binding and lifecycle integration | **Pending; Codex-first** | +| Named role-profile HTTP/editor | **Implemented; offline validated; configuration mutations only** | +| Role-profile session, worker, cache, attempt, and audit provenance | **Implemented; offline validated; Codex-first** | +| Parent-owned role-profile lifecycle integration | **Pending; Codex-first** | | Non-Codex subagent configuration | **Pending; configuration only before execution** | | Non-Codex execution and orchestration | **Pending; later milestone** | | Stable public API or configuration compatibility | **Pending** | @@ -125,9 +126,10 @@ validation. provider-backed evidence. - Verified changes have no provider-backed patcher or verifier observation. - Canonical role-profile definitions, revision persistence, bounded HTTP/editor - flows, and request-time preflight are implemented and offline validated. - They do not provide session or worker binding, a lifecycle executor, or - automatic profile activation; activation is an explicit configuration change. + flows, request-time preflight, and frozen session/worker/cache/attempt/audit + provenance are implemented and offline validated. They do not provide a + parent-owned lifecycle executor or automatic profile activation; activation + is an explicit configuration change. - The verifier handles a deterministic serial set of up to four distinct associated certified test plans; exact duplicates collapse to one execution, while over-cap and unavailable plans fail closed. This behavior is offline diff --git a/crates/needle-app/src/artifact_cache_main_replay.rs b/crates/needle-app/src/artifact_cache_main_replay.rs index a70b91e..e1133ed 100644 --- a/crates/needle-app/src/artifact_cache_main_replay.rs +++ b/crates/needle-app/src/artifact_cache_main_replay.rs @@ -1,6 +1,6 @@ use super::{ AppError, HookConfig, absolute_run_path, canonical_child_path, option_value, - repository_status_clean, required_value, resolve_codex, + provision_experiment_role_profile, repository_status_clean, required_value, resolve_codex, }; use needle_bench::{ArtifactCacheReplayReport, run_artifact_cache_replay}; use needle_core::{ @@ -47,6 +47,7 @@ impl WorkerExecutor for ForbiddenWorker { discarded_facts: 0, worker_session_id: None, session_cleanup_success: Some(true), + role_profile_provenance: None, })) } } @@ -102,6 +103,16 @@ pub(super) fn run(arguments: &[String]) -> Result<(), AppError> { let store = RuntimeStore::new(artifact_root.join("needle.sqlite3")); let profile = HookConfig::default().profile().map_err(|error| AppError::Experiment(error.to_string()))?; + let role_profile_id = provision_experiment_role_profile( + &store, + "artifact-cache-replay.explorer", + profile.definition_digest, + "recorded-r35-fixture", + "low", + "default", + 1, + false, + )?; let instructions = profile.rendered_context_owned(); let main_config = WorkerConfig { executable: simulator.display().to_string(), @@ -110,6 +121,7 @@ pub(super) fn run(arguments: &[String]) -> Result<(), AppError> { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start(MainSessionConfig { codex: &main_config, @@ -126,11 +138,12 @@ pub(super) fn run(arguments: &[String]) -> Result<(), AppError> { .map_err(AppError::Experiment)?; let session_id = session.thread_id().to_owned(); store - .record_session_start( + .record_session_start_profiled( &session_id, profile.definition_digest, Some("simulated-main-r35-cache"), source_repository.to_str(), + &role_profile_id, ) .map_err(|error| AppError::Experiment(error.to_string()))?; diff --git a/crates/needle-app/src/main.rs b/crates/needle-app/src/main.rs index ee35180..6b99125 100644 --- a/crates/needle-app/src/main.rs +++ b/crates/needle-app/src/main.rs @@ -7,7 +7,10 @@ use needle_bench::{ evaluate_pilot_pair, parse_codex_jsonl, parse_jsonl, parse_task_fixture, redact_jsonl, }; use needle_core::{ - Digest, EvidenceFailurePolicy, FORMAT_REVISION, NeedKey, NeedRequest, TestPlan, WorkerConfig, + CodexHost, CodexRole, CommandPolicy, Digest, EvidenceFailurePolicy, FORMAT_REVISION, + FallbackPolicy, FilesystemPolicy, NeedKey, NeedRequest, NetworkPolicy, ReasoningLevel, + RepairPolicy, RoleProfileBudget, RoleProfileDefinition, RoleProfileDefinitionInput, + RoleProfileId, ServiceTier, TestPlan, TestPolicy, ToolPolicy, WorkerConfig, }; use needle_platform_codex::{ CodexWorker, CompactInput, HookConfig, SessionEndInput, SessionStartInput, StopInput, @@ -508,15 +511,33 @@ fn run_hook(arguments: Vec) -> Result<(), AppError> { if let Some(session_id) = parsed.session_id.as_deref() { let store = hook_runtime_store()?; let profile_digest = config.profile()?.definition_digest; - if let Err(error) = store.initialize().and_then(|_| { - store.record_session_start( - session_id, - profile_digest, - parsed.model.as_deref(), - parsed.cwd.as_deref(), - ) - }) { - eprintln!("needle: cannot record product session ({error}); fail-open"); + let selector = env::var("NEEDLE_ROLE_PROFILE_ID").ok(); + let result = match selector { + Some(value) => match RoleProfileId::new(value) { + Ok(profile_id) => store.initialize().and_then(|_| { + store.record_session_start_profiled( + session_id, + profile_digest, + parsed.model.as_deref(), + parsed.cwd.as_deref(), + &profile_id, + ) + }), + Err(error) => Err(needle_runtime::StoreError::RoleProfileValidation( + error.to_string(), + )), + }, + None => { + eprintln!( + "needle: NEEDLE_ROLE_PROFILE_ID is missing; session provenance is unknown" + ); + Ok(()) + } + }; + if let Err(error) = result { + eprintln!( + "needle: cannot record profiled product session ({error}); fail-open" + ); } } serde_json::to_value(output)? @@ -675,6 +696,10 @@ fn run_mcp(arguments: Vec) -> Result<(), AppError> { .or_else(|| env::var("NEEDLE_MCP_MAIN_MODEL").ok()) .unwrap_or_else(|| "unknown".to_owned()); validate_model_value(&main_model, "main model")?; + let role_profile = required_value(&arguments, "--role-profile").and_then(|value| { + RoleProfileId::new(value) + .map_err(|error| AppError::Usage(format!("invalid role profile: {error}"))) + })?; return mcp::serve(mcp::ProductMcpConfig { data_directory, repository_root, @@ -682,12 +707,13 @@ fn run_mcp(arguments: Vec) -> Result<(), AppError> { cache_only: arguments.iter().any(|argument| argument == "--cache-only"), calibration_reuse: env::var("NEEDLE_INTERNAL_CALIBRATION_REUSE").as_deref() == Ok("partial-tests-live"), + role_profile_id: role_profile, }) .map_err(AppError::Runtime); } if arguments.first().map(String::as_str) != Some("serve-benchmark") || arguments.len() != 1 { return Err(AppError::Usage( - "mcp serve [--data-dir ] [--repository ] [--main-model ] [--cache-only] | mcp serve-benchmark" + "mcp serve --role-profile [--data-dir ] [--repository ] [--main-model ] [--cache-only] | mcp serve-benchmark" .to_owned(), )); } @@ -699,7 +725,7 @@ fn validate_mcp_serve_arguments(arguments: &[String]) -> Result<(), AppError> { while index < arguments.len() { match arguments[index].as_str() { "--cache-only" => index += 1, - "--data-dir" | "--repository" | "--main-model" => { + "--data-dir" | "--repository" | "--main-model" | "--role-profile" => { if arguments.get(index + 1).is_none() { return Err(AppError::Usage(format!("{} requires a value", arguments[index]))); } @@ -808,6 +834,7 @@ fn transport_preflight_run(arguments: &[String]) -> Result<(), AppError> { service_tier: Some(service_tier), timeout_seconds: 30, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let report = CodexWorker::new(&data_root) .preflight_transport(&config, &repository) @@ -3058,6 +3085,77 @@ fn validate_model_value(value: &str, label: &str) -> Result<(), AppError> { validate_slug(value, label) } +#[allow(clippy::too_many_arguments)] +fn provision_experiment_role_profile( + store: &RuntimeStore, + profile_id: &str, + prompt_profile_digest: Digest, + model: &str, + reasoning: &str, + service_tier: &str, + timeout_seconds: u64, + repair_once: bool, +) -> Result { + let reasoning = match reasoning { + "low" => ReasoningLevel::Low, + "medium" => ReasoningLevel::Medium, + "high" => ReasoningLevel::High, + "xhigh" => ReasoningLevel::Xhigh, + value => { + return Err(AppError::Experiment(format!( + "experiment role profile cannot represent reasoning `{value}`" + ))); + } + }; + let service_tier = match service_tier { + "default" => ServiceTier::Default, + "priority" => ServiceTier::Priority, + value => { + return Err(AppError::Experiment(format!( + "experiment role profile cannot represent service tier `{value}`" + ))); + } + }; + let profile_id = + RoleProfileId::new(profile_id).map_err(|error| AppError::Experiment(error.to_string()))?; + let definition = RoleProfileDefinition::new(RoleProfileDefinitionInput { + profile_id: profile_id.clone(), + role: CodexRole::Explorer, + host: CodexHost::Codex, + model: model.to_owned(), + reasoning, + service_tier, + timeout_seconds, + budget: RoleProfileBudget { + max_turns: 8, + max_output_tokens: 2000, + max_cost_microusd: 1_000_000_000, + }, + prompt_profile_digest, + output_contract_digest: Digest::blake3(needle_core::ARTIFACT_RESULT_SCHEMA_ID), + tool_policy: ToolPolicy::ReadOnly, + command_policy: CommandPolicy::ReadOnly, + filesystem_policy: FilesystemPolicy::ReadOnlyCheckout, + network_policy: NetworkPolicy::Denied, + test_policy: TestPolicy::Disabled, + repair_policy: if repair_once { RepairPolicy::Once } else { RepairPolicy::None }, + fallback_policy: FallbackPolicy::Disabled, + concurrency: 1, + route_assignments: Vec::new(), + }) + .map_err(|error| AppError::Experiment(error.to_string()))?; + let revision = store + .create_role_profile(definition) + .map_err(|error| AppError::Experiment(error.to_string()))?; + let state = store + .role_profile_state(&profile_id) + .map_err(|error| AppError::Experiment(error.to_string()))?; + store + .activate_role_profile(&profile_id, revision.revision, state.state_digest) + .map_err(|error| AppError::Experiment(error.to_string()))?; + Ok(profile_id) +} + fn parse_experiment_arm(value: &str) -> Result { match value { "P0" => Ok(ExperimentArm::P0), diff --git a/crates/needle-app/src/mcp.rs b/crates/needle-app/src/mcp.rs index bffcf48..b9f4e1a 100644 --- a/crates/needle-app/src/mcp.rs +++ b/crates/needle-app/src/mcp.rs @@ -9,7 +9,8 @@ use change_schema::{ use needle_core::{ CacheResolution, CanonicalHasher, Digest, MultiNeedPolicy, Need, NeedCoordination, NeedDelivery, NeedStep, NeedStepRelation, NeedStepState, PredicateKind, - ReuseSufficiencyCertificateId, VerificationStatus, classify_need_step, + ReuseSufficiencyCertificateId, RoleProfileId, RoleProfileProvenance, VerificationStatus, + classify_need_step, }; use needle_platform_codex::{CodexPatchWorker, CodexVerifier, PatchContextItem}; use needle_runtime::{ @@ -87,6 +88,7 @@ pub(crate) struct ProductMcpConfig { pub(crate) main_model: String, pub(crate) cache_only: bool, pub(crate) calibration_reuse: bool, + pub(crate) role_profile_id: RoleProfileId, } pub(crate) fn serve(config: ProductMcpConfig) -> Result<(), String> { @@ -445,6 +447,7 @@ struct ProductMcpServer { repository_lineage: Digest, main_model: String, session_id: String, + role_profile_provenance: RoleProfileProvenance, next_turn: u64, protocol: Option, initialized: bool, @@ -511,7 +514,7 @@ impl ProductMcpServer { let profile_digest = Digest::blake3(b"needle.mcp-json-profile/1"); resolver .store() - .record_session_start_for_transport( + .record_session_start_for_transport_profiled( &session_id, profile_digest, Some(&config.main_model), @@ -519,14 +522,18 @@ impl ProductMcpServer { "mcp", transport_definition_digest, None, + &config.role_profile_id, ) .map_err(|error| error.to_string())?; - let policy = resolver + let session = resolver .store() .session(&session_id) .map_err(|error| error.to_string())? - .ok_or_else(|| "MCP session was not persisted".to_owned())? - .multi_need_policy; + .ok_or_else(|| "MCP session was not persisted".to_owned())?; + let policy = session.multi_need_policy; + let role_profile_provenance = session + .role_profile_provenance + .ok_or_else(|| "MCP session role-profile provenance was not persisted".to_owned())?; Ok(Self { resolver, cancellation, @@ -534,6 +541,7 @@ impl ProductMcpServer { repository_lineage: snapshot.repository_id, main_model: config.main_model, session_id, + role_profile_provenance, next_turn: 1, protocol: None, initialized: false, @@ -839,6 +847,7 @@ impl ProductMcpServer { &mapped.request.route, relation, &outcome, + &self.role_profile_provenance, ) { return tool_error(id, &format!("cannot persist MCP observation: {error}")); } @@ -865,17 +874,23 @@ impl ProductMcpServer { Ok(settings) => settings, Err(error) => return tool_error(id, &format!("cannot load worker settings: {error}")), }; + let worker_config = match self + .resolver + .store() + .resolve_session_worker_config(&self.session_id, settings.codex_executable.clone()) + { + Ok(config) => config, + Err(error) => { + return tool_error(id, &format!("cannot resolve role-profile config: {error}")); + } + }; let patcher = CodexPatchWorker::new(self.resolver.data_directory()) .with_cancellation(Arc::clone(&self.cancellation)); - let outcome = match patcher.prepare( - &settings.worker_config(), - &self.repository_root, - &request, - &context, - ) { - Ok(outcome) => outcome, - Err(error) => return tool_error(id, &error), - }; + let outcome = + match patcher.prepare(&worker_config, &self.repository_root, &request, &context) { + Ok(outcome) => outcome, + Err(error) => return tool_error(id, &error), + }; let request_digest = outcome.request_digest; let response = McpPrepareChangeResponse::from_outcome( outcome, @@ -1002,7 +1017,16 @@ impl ProductMcpServer { }; let verifier = CodexVerifier::new(self.resolver.data_directory()) .with_cancellation(Arc::clone(&self.cancellation)); - let worker_config = settings.worker_config(); + let worker_config = match self + .resolver + .store() + .resolve_session_worker_config(&self.session_id, settings.codex_executable.clone()) + { + Ok(config) => config, + Err(error) => { + return tool_error(id, &format!("cannot resolve role-profile config: {error}")); + } + }; let first = match verifier.verify(&worker_config, &self.repository_root, &request.change_id) { Ok(outcome) => outcome, @@ -1026,9 +1050,10 @@ impl ProductMcpServer { &request.change_id, ) { Ok(outcome) => outcome, - Err(error) => match verifier.record_inconclusive( + Err(error) => match verifier.record_inconclusive_with_provenance( &request.change_id, &format!("verification after the one-shot repair failed: {error}"), + worker_config.role_profile_provenance.as_ref(), ) { Ok(outcome) => outcome, Err(record_error) => return tool_error(id, &record_error), @@ -1036,9 +1061,10 @@ impl ProductMcpServer { }; } Err(error) => { - outcome = match verifier.record_inconclusive( + outcome = match verifier.record_inconclusive_with_provenance( &request.change_id, &format!("one-shot repair failed: {error}"), + worker_config.role_profile_provenance.as_ref(), ) { Ok(outcome) => outcome, Err(record_error) => return tool_error(id, &record_error), @@ -1096,6 +1122,10 @@ impl ProductMcpServer { worker_avoided: false, main_discovery_tainted: false, }; + let audit = serde_json::to_string(&json!({ + "role_profile_provenance": self.role_profile_provenance, + })) + .expect("bounded role-profile provenance serialization"); let persisted = self .resolver .store() @@ -1110,14 +1140,14 @@ impl ProductMcpServer { self.resolver.store().append_need_step_event( step.id, NeedStepState::Resolving, - "{}", + &audit, ) }) .and_then(|_| { self.resolver.store().append_need_step_event( step.id, NeedStepState::Cancelled, - "{}", + &audit, ) }); if let Err(error) = persisted { @@ -1537,6 +1567,7 @@ fn record_observation( route: &str, relation: NeedStepRelation, outcome: &ResolveOutcome, + role_profile_provenance: &RoleProfileProvenance, ) -> Result<(), String> { let path = data_directory.join(OBSERVATION_FILE); let mut file = OpenOptions::new() @@ -1547,7 +1578,7 @@ fn record_observation( serde_json::to_writer( &mut file, &json!({ - "schema": "needle.mcp-observation/3", + "schema": "needle.mcp-observation/4", "transport": "mcp", "request_format": "json", "turn_id": turn_id, @@ -1561,6 +1592,7 @@ fn record_observation( "calibration": outcome.calibration, "result_digest": outcome.result_digest, "semantic_artifact_ids": outcome.semantic_artifact_ids, + "role_profile_provenance": role_profile_provenance, }), ) .map_err(|error| error.to_string())?; @@ -1588,7 +1620,12 @@ fn rpc_error(id: Value, code: i64, message: &str) -> Value { #[cfg(test)] mod tests { use super::*; - use needle_core::{EvidenceFailurePolicy, MultiNeedPolicy}; + use needle_core::{ + CodexHost, CodexRole, CommandPolicy, EvidenceFailurePolicy, FallbackPolicy, + FilesystemPolicy, MultiNeedPolicy, NetworkPolicy, RepairPolicy, RoleProfileBudget, + RoleProfileDefinition, RoleProfileDefinitionInput, RoleProfileId, ServiceTier, TestPolicy, + ToolPolicy, + }; use needle_runtime::{RuntimeSettings, RuntimeStore}; use std::process::Command; @@ -1993,12 +2030,48 @@ mod tests { multi_need_policy: MultiNeedPolicy::default(), }) .unwrap(); + let profile_id = RoleProfileId::new("explorer.default").unwrap(); + let definition = RoleProfileDefinition::new(RoleProfileDefinitionInput { + profile_id: profile_id.clone(), + role: CodexRole::Explorer, + host: CodexHost::Codex, + model: "worker".to_owned(), + reasoning: needle_core::ReasoningLevel::Medium, + service_tier: ServiceTier::Default, + timeout_seconds: 5, + budget: RoleProfileBudget { + max_turns: 2, + max_output_tokens: 1200, + max_cost_microusd: 1000, + }, + prompt_profile_digest: Digest::blake3(b"prompt"), + output_contract_digest: Digest::blake3(b"output"), + tool_policy: ToolPolicy::ReadOnly, + command_policy: CommandPolicy::ReadOnly, + filesystem_policy: FilesystemPolicy::ReadOnlyCheckout, + network_policy: NetworkPolicy::Denied, + test_policy: TestPolicy::Disabled, + repair_policy: RepairPolicy::None, + fallback_policy: FallbackPolicy::Native, + concurrency: 1, + route_assignments: Vec::new(), + }) + .unwrap(); + store.create_role_profile(definition).unwrap(); + store + .activate_role_profile( + &profile_id, + 1, + store.role_profile_state(&profile_id).unwrap().state_digest, + ) + .unwrap(); let server = ProductMcpServer::new(ProductMcpConfig { data_directory, repository_root: repository, main_model: "main".to_owned(), cache_only: true, calibration_reuse: false, + role_profile_id: profile_id, }) .unwrap(); (root, server) diff --git a/crates/needle-app/src/minimal_live_pilot.rs b/crates/needle-app/src/minimal_live_pilot.rs index 0fa3fcf..a99ba8e 100644 --- a/crates/needle-app/src/minimal_live_pilot.rs +++ b/crates/needle-app/src/minimal_live_pilot.rs @@ -2,8 +2,8 @@ use super::{ AppError, HookConfig, absolute_run_path, canonical_child_path, clone_local_checkout, ensure_cache_pilot_hook_binary, ensure_codex_authenticated, ensure_dedicated_codex_home, ensure_product_pilot_hook_isolation, option_value, price_usage_observation_optional, - repository_status_clean, required_value, resolve_codex, validate_model_value, - validate_reasoning, validate_service_tier, validate_slug, + provision_experiment_role_profile, repository_status_clean, required_value, resolve_codex, + validate_model_value, validate_reasoning, validate_service_tier, validate_slug, }; use needle_bench::{ BenchmarkRoute, CachePilotResolveOutcome, ECONOMIC_EQUIVALENT_HIT_PROMPT, FinalArm, @@ -12,8 +12,8 @@ use needle_bench::{ }; use needle_core::{ CacheResolution, CapabilityMode, Digest, EvidenceFailurePolicy, ModelPolicy, MultiNeedPolicy, - NeedStep, PredicateKind, ReuseUnit, SelectedPlan, SemanticWorkerArtifact, WorkerConfig, - WorkerProfile, + NeedStep, PredicateKind, ReuseUnit, RoleProfileId, SelectedPlan, SemanticWorkerArtifact, + WorkerConfig, WorkerProfile, }; use needle_platform_codex::{CodexWorker, TransportPreflightReport}; use needle_runtime::{ @@ -735,6 +735,16 @@ fn execute(context: Execution<'_>) -> Result<(), AppError> { let profile = HookConfig::default().profile().map_err(|error| AppError::Experiment(error.to_string()))?; + let role_profile_id = provision_experiment_role_profile( + &store, + "minimal-live-pilot.explorer", + profile.definition_digest, + context.worker_model, + context.worker_reasoning, + context.service_tier, + context.timeout.as_secs().min(600), + false, + )?; let main_instructions = protocol::pilot_main_instructions(&profile.rendered_context_owned()); let miss_quality_spec = quality_spec(context.protocol)?; let hit_quality_spec = coverage_hit_quality_spec(context.protocol)?; @@ -822,6 +832,7 @@ fn execute(context: Execution<'_>) -> Result<(), AppError> { source_snapshot_digest: initial_snapshot.source_digest, repository_id: initial_snapshot.repository_id, prompt_profile_digest: profile.definition_digest, + role_profile_id: &role_profile_id, main_instructions: &main_instructions, prompt: publication_prompt, declared_test_plan: &declared_test_plan, @@ -963,6 +974,7 @@ fn execute(context: Execution<'_>) -> Result<(), AppError> { source_snapshot_digest: initial_snapshot.source_digest, repository_id: initial_snapshot.repository_id, prompt_profile_digest: profile.definition_digest, + role_profile_id: &role_profile_id, main_instructions: &main_instructions, prompt: if context.economic { if context.trace_reuse { @@ -1081,6 +1093,7 @@ struct Observe<'a> { source_snapshot_digest: Digest, repository_id: Digest, prompt_profile_digest: Digest, + role_profile_id: &'a RoleProfileId, main_instructions: &'a str, prompt: &'a str, declared_test_plan: &'a needle_core::TestPlan, @@ -1560,6 +1573,7 @@ fn run_declared_test_transport_preflight( service_tier: Some(service_tier.to_owned()), timeout_seconds: 30, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; CodexWorker::with_codex_home(data_root, codex_home) .preflight_transport_for_test_plan( diff --git a/crates/needle-app/src/minimal_live_pilot/direct_main.rs b/crates/needle-app/src/minimal_live_pilot/direct_main.rs index 549b835..7aca2a6 100644 --- a/crates/needle-app/src/minimal_live_pilot/direct_main.rs +++ b/crates/needle-app/src/minimal_live_pilot/direct_main.rs @@ -34,6 +34,7 @@ pub(super) fn observe(context: DirectObservation<'_>) -> Result Result<(), AppError> { service_tier: Some(service_tier.clone()), timeout_seconds, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let temporary = TemporaryRunRoot::create(&artifact_root)?; let request = capture_worker_request(temporary.path(), &source_repository, &protocol, &config)?; @@ -519,6 +521,7 @@ impl WorkerExecutor for CaptureWorker { discarded_facts: 0, worker_session_id: None, session_cleanup_success: None, + role_profile_provenance: _config.role_profile_provenance.clone(), })) } } @@ -543,14 +546,25 @@ fn capture_worker_request( .map_err(|error| AppError::Experiment(error.to_string()))?; let profile = HookConfig::default().profile().map_err(|error| AppError::Experiment(error.to_string()))?; + let role_profile_id = provision_experiment_role_profile( + &store, + "worker-diagnostic.explorer", + profile.definition_digest, + &config.model, + &config.reasoning, + config.service_tier.as_deref().unwrap_or("default"), + config.timeout_seconds, + false, + )?; let session = "worker-diagnostic-capture"; let turn = "worker-diagnostic-turn"; store - .record_session_start( + .record_session_start_profiled( session, profile.definition_digest, Some("diagnostic-main"), repository.to_str(), + &role_profile_id, ) .map_err(|error| AppError::Experiment(error.to_string()))?; store @@ -751,6 +765,7 @@ mod tests { discarded_facts: 0, worker_session_id: Some("thread".to_owned()), session_cleanup_success: Some(true), + role_profile_provenance: None, }; let metrics = attempt_metrics(None, Some(&failure)).unwrap(); assert_eq!(metrics.logical_worker_spawns, 1); diff --git a/crates/needle-bench/src/minimal_pilot.rs b/crates/needle-bench/src/minimal_pilot.rs index 4c8e010..064073d 100644 --- a/crates/needle-bench/src/minimal_pilot.rs +++ b/crates/needle-bench/src/minimal_pilot.rs @@ -3,10 +3,12 @@ use crate::{ preflight_frozen_corpus, }; use needle_core::{ - CacheResolution, CapabilityMode, Claim, CommandExecutionEvidence, Digest, - EvidenceFailurePolicy, EvidenceReference, NeedIr, NeedResult, PredicateKind, ReuseUnit, - SemanticArtifactResult, SemanticInterrupt, SemanticWorkerArtifact, TestPlan, Uncertainty, - WorkerConfig, WorkerFailure, WorkerOutcome, WorkerRequest, + CacheResolution, CapabilityMode, Claim, CodexHost, CodexRole, CommandExecutionEvidence, + CommandPolicy, Digest, EvidenceFailurePolicy, EvidenceReference, FallbackPolicy, + FilesystemPolicy, NeedIr, NeedResult, NetworkPolicy, PredicateKind, RepairPolicy, ReuseUnit, + RoleProfileBudget, RoleProfileDefinition, RoleProfileDefinitionInput, RoleProfileId, + SemanticArtifactResult, SemanticInterrupt, SemanticWorkerArtifact, ServiceTier, TestPlan, + TestPolicy, ToolPolicy, Uncertainty, WorkerConfig, WorkerFailure, WorkerOutcome, WorkerRequest, }; use needle_runtime::{ ResolveOutcome, ResolveRequest, RouteCostObservation, RuntimeEngine, RuntimeError, @@ -152,6 +154,7 @@ impl WorkerExecutor for DeterministicPilotWorker { discarded_facts: 0, worker_session_id: None, session_cleanup_success: Some(true), + role_profile_provenance: None, }) })?; } @@ -172,6 +175,7 @@ impl WorkerExecutor for DeterministicPilotWorker { discarded_facts: 0, worker_session_id: None, session_cleanup_success: Some(true), + role_profile_provenance: None, }) })?; retain_location_artifact(&mut semantic_result); @@ -204,6 +208,7 @@ impl WorkerExecutor for DeterministicPilotWorker { discarded_facts: 0, worker_session_id: None, session_cleanup_success: Some(true), + role_profile_provenance: None, }) })?; let evidence_id = "implementation-location".to_owned(); @@ -260,6 +265,7 @@ impl WorkerExecutor for DeterministicPilotWorker { discarded_facts: 0, worker_session_id: None, session_cleanup_success: Some(true), + role_profile_provenance: config.role_profile_provenance.clone(), }) } } @@ -340,6 +346,7 @@ fn run_minimal_pilot_dry_run_with_hit_prompt( trusted_test_execution: false, multi_need_policy: needle_core::MultiNeedPolicy::default(), })?; + let role_profile_id = deterministic_role_profile(&store)?; store.mark_utility_gate_passed()?; let spawns = Arc::new(AtomicU32::new(0)); let engine = RuntimeEngine::new( @@ -362,6 +369,7 @@ fn run_minimal_pilot_dry_run_with_hit_prompt( &source_repository, Some(declared_test_plan.clone()), true, + &role_profile_id, )?; let miss_interrupt_digest = semantic_interrupt_digest(&miss_request)?; let miss = engine.resolve(&miss_request)?; @@ -420,6 +428,7 @@ fn run_minimal_pilot_dry_run_with_hit_prompt( &source_repository, (!reworded_hit).then_some(declared_test_plan), !reworded_hit, + &role_profile_id, )?; let hit_interrupt_digest = semantic_interrupt_digest(&hit_request)?; let semantic_interrupt_digest_matches = miss_interrupt_digest == hit_interrupt_digest; @@ -551,6 +560,44 @@ fn semantic_interrupt_digest(request: &ResolveRequest) -> Result Result { + let profile_id = RoleProfileId::new("benchmark.explorer") + .map_err(|error| MinimalPilotDryRunError::Invalid(error.to_string()))?; + let definition = RoleProfileDefinition::new(RoleProfileDefinitionInput { + profile_id: profile_id.clone(), + role: CodexRole::Explorer, + host: CodexHost::Codex, + model: "deterministic-offline-fixture".to_owned(), + reasoning: needle_core::ReasoningLevel::Low, + service_tier: ServiceTier::Default, + timeout_seconds: 1, + budget: RoleProfileBudget { + max_turns: 2, + max_output_tokens: 1200, + max_cost_microusd: 1000, + }, + prompt_profile_digest: Digest::blake3(b"minimal-pilot-v04-locate-profile"), + output_contract_digest: Digest::blake3(needle_core::ARTIFACT_RESULT_SCHEMA_ID), + tool_policy: ToolPolicy::ReadOnly, + command_policy: CommandPolicy::ReadOnly, + filesystem_policy: FilesystemPolicy::ReadOnlyCheckout, + network_policy: NetworkPolicy::Denied, + test_policy: TestPolicy::Disabled, + repair_policy: RepairPolicy::None, + fallback_policy: FallbackPolicy::Native, + concurrency: 1, + route_assignments: Vec::new(), + }) + .map_err(|error| MinimalPilotDryRunError::Invalid(error.to_string()))?; + store.create_role_profile(definition)?; + let state = store.role_profile_state(&profile_id)?; + store.activate_role_profile(&profile_id, 1, state.state_digest)?; + Ok(profile_id) +} + +#[allow(clippy::too_many_arguments)] fn resolve_request( store: &RuntimeStore, task_id: &str, @@ -559,15 +606,17 @@ fn resolve_request( source_repository: &Path, declared_test_plan: Option, require_focused_tests: bool, + role_profile_id: &RoleProfileId, ) -> Result { let session = format!("minimal-pilot-{task_id}-{arm}"); let turn = "turn"; let profile_digest = Digest::blake3(b"minimal-pilot-v04-locate-profile"); - store.record_session_start( + store.record_session_start_profiled( &session, profile_digest, Some("frontier"), source_repository.to_str(), + role_profile_id, )?; store.record_user_prompt(&session, Some(turn), prompt, source_repository.to_str())?; let focused_tests = if require_focused_tests { diff --git a/crates/needle-bench/src/shadow_replay.rs b/crates/needle-bench/src/shadow_replay.rs index 2b9fe1b..60484f9 100644 --- a/crates/needle-bench/src/shadow_replay.rs +++ b/crates/needle-bench/src/shadow_replay.rs @@ -288,7 +288,7 @@ mod tests { Digest, EvidenceFailurePolicy, NeedCacheEntry, NeedCacheIdentity, NeedKey, NeedResult, WorkerArtifactResult, WorkerObservationTrace, WorkerOutcome, }; - use needle_runtime::RuntimeSettings; + use needle_runtime::{RuntimeSettings, StoreError}; use std::collections::BTreeMap; #[test] @@ -325,6 +325,7 @@ mod tests { normalized_request_digest: Digest::blake3(b"request"), worker_configuration_digest: Digest::blake3(b"worker"), output_schema_digest: Digest::blake3(b"schema"), + role_profile_provenance: None, }; let result = NeedResult { complete: true, @@ -365,11 +366,12 @@ mod tests { discarded_facts: 0, worker_session_id: None, session_cleanup_success: Some(true), + role_profile_provenance: None, }, created_unix_ms: 1, hit_count: 0, }; - store.publish(&entry).unwrap(); + assert!(matches!(store.publish(&entry), Err(StoreError::ArtifactIdentity(_)))); fs::write(root.join("report.json"), "{}").unwrap(); let report = run_shadow_replay( diff --git a/crates/needle-core/src/domain.rs b/crates/needle-core/src/domain.rs index bc52fa4..2fc6b80 100644 --- a/crates/needle-core/src/domain.rs +++ b/crates/needle-core/src/domain.rs @@ -1,6 +1,7 @@ use crate::{ ArtifactKind, Digest, HARD_RESULT_BYTES, HARD_RESULT_TOKENS, NeedFragment, NeedKey, - SemanticArtifactResult, TestPlan, WorkerArtifactResult, normalize_line_endings, + RoleProfileProvenance, SemanticArtifactResult, TestPlan, WorkerArtifactResult, + normalize_line_endings, }; use serde::{Deserialize, Serialize}; @@ -106,18 +107,29 @@ pub struct WorkerConfig { pub timeout_seconds: u64, #[serde(default)] pub evidence_failure_policy: EvidenceFailurePolicy, + #[serde(default)] + pub role_profile_provenance: Option, } impl WorkerConfig { pub fn digest(&self) -> Digest { Digest::blake3(format!( - "needle-worker-config\n{}\n{}\n{}\n{}\n{}\n{}\n", + "needle-worker-config\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n", self.executable, self.model, self.reasoning, self.service_tier.as_deref().unwrap_or_default(), self.timeout_seconds, self.evidence_failure_policy.as_str(), + self.role_profile_provenance + .as_ref() + .map(|provenance| { + format!( + "{}@{}#{}", + provenance.profile_id, provenance.revision, provenance.definition_digest + ) + }) + .unwrap_or_default(), )) } } @@ -186,6 +198,8 @@ pub struct WorkerOutcome { pub worker_session_id: Option, #[serde(default)] pub session_cleanup_success: Option, + #[serde(default)] + pub role_profile_provenance: Option, } const fn one_u32() -> u32 { @@ -208,6 +222,8 @@ pub struct WorkerFailure { pub discarded_facts: u32, pub worker_session_id: Option, pub session_cleanup_success: Option, + #[serde(default)] + pub role_profile_provenance: Option, } #[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] @@ -338,12 +354,14 @@ pub struct NeedCacheIdentity { pub normalized_request_digest: Digest, pub worker_configuration_digest: Digest, pub output_schema_digest: Digest, + #[serde(default)] + pub role_profile_provenance: Option, } impl NeedCacheIdentity { pub fn digest(&self) -> Digest { Digest::blake3(format!( - "needle-cache-identity\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n", + "needle-cache-identity\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n", self.repository_id, self.source_snapshot_digest, self.prompt_profile_digest, @@ -353,12 +371,13 @@ impl NeedCacheIdentity { self.normalized_request_digest, self.worker_configuration_digest, self.output_schema_digest, + provenance_identity(self.role_profile_provenance.as_ref()), )) } pub fn logical_digest(&self) -> Digest { Digest::blake3(format!( - "needle-cache-logical\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n", + "needle-cache-logical\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n{}\n", self.repository_id, self.prompt_profile_digest, self.route_definition_digest, @@ -367,10 +386,17 @@ impl NeedCacheIdentity { self.normalized_request_digest, self.worker_configuration_digest, self.output_schema_digest, + provenance_identity(self.role_profile_provenance.as_ref()), )) } } +fn provenance_identity(provenance: Option<&RoleProfileProvenance>) -> String { + provenance + .map(|value| format!("{}@{}#{}", value.profile_id, value.revision, value.definition_digest)) + .unwrap_or_else(|| "unknown".to_owned()) +} + #[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] #[serde(deny_unknown_fields)] pub struct NeedCacheEntry { diff --git a/crates/needle-core/src/role_profile.rs b/crates/needle-core/src/role_profile.rs index 2ee1f5b..4fe61f6 100644 --- a/crates/needle-core/src/role_profile.rs +++ b/crates/needle-core/src/role_profile.rs @@ -1,5 +1,6 @@ use crate::{ - CanonicalHasher, Digest, HARD_MAX_NEEDS_PER_TASK, HARD_RESULT_TOKENS, NeedKey, WorkerProfile, + CanonicalHasher, Digest, EvidenceFailurePolicy, HARD_MAX_NEEDS_PER_TASK, HARD_RESULT_TOKENS, + NeedKey, WorkerConfig, WorkerProfile, }; use serde::{Deserialize, Deserializer, Serialize, Serializer}; use std::fmt; @@ -615,12 +616,88 @@ pub struct RoleProfileRevision { pub activated_unix_ms: Option, } +/// The bounded, immutable identity of the role-profile revision used by a +/// session. This intentionally carries no executable, model prompt, policy, +/// or other host-local configuration. +#[derive(Clone, Debug, Eq, PartialEq, Hash, Ord, PartialOrd, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct RoleProfileProvenance { + pub profile_id: RoleProfileId, + pub revision: u64, + pub definition_digest: Digest, +} + +/// Compatibility spelling used by persistence and runtime APIs. +pub type RoleProfileBinding = RoleProfileProvenance; + +impl RoleProfileProvenance { + pub fn new( + profile_id: RoleProfileId, + revision: u64, + definition_digest: Digest, + ) -> Result { + let provenance = Self { profile_id, revision, definition_digest }; + provenance.validate()?; + Ok(provenance) + } + + pub fn from_revision( + revision: &RoleProfileRevision, + ) -> Result { + revision.validate()?; + Self::new( + revision.profile_id.clone(), + revision.revision, + revision.definition.definition_digest, + ) + } + + pub fn validate(&self) -> Result<(), RoleProfileValidationError> { + if self.revision == 0 { + return Err(RoleProfileValidationError::Revision); + } + // RoleProfileId is private-field validated by its constructor and + // Deserialize implementation. Calling new here also protects values + // assembled by direct struct literals in persistence code. + RoleProfileId::new(self.profile_id.as_str().to_owned()) + .map_err(|_| RoleProfileValidationError::IdentityMismatch)?; + Ok(()) + } +} + impl RoleProfileRevision { pub fn to_worker_profile(&self) -> Result { self.validate()?; self.definition.to_worker_profile() } + /// Deterministically project this frozen historical revision to the + /// existing Codex worker boundary. The executable is host-local input and + /// therefore deliberately supplied by the parent rather than persisted in + /// the role profile. + pub fn to_worker_config( + &self, + executable: impl Into, + ) -> Result { + self.validate()?; + let provenance = RoleProfileProvenance::from_revision(self)?; + Ok(WorkerConfig { + executable: executable.into(), + model: self.definition.model.clone(), + reasoning: self.definition.reasoning.as_str().to_owned(), + service_tier: match self.definition.service_tier { + ServiceTier::Default => None, + ServiceTier::Priority => Some("priority".to_owned()), + }, + timeout_seconds: self.definition.timeout_seconds, + evidence_failure_policy: match self.definition.repair_policy { + RepairPolicy::None => EvidenceFailurePolicy::DiscardInvalidFact, + RepairPolicy::Once => EvidenceFailurePolicy::RepairOnce, + }, + role_profile_provenance: Some(provenance), + }) + } + pub fn validate(&self) -> Result<(), RoleProfileValidationError> { if self.revision == 0 { return Err(RoleProfileValidationError::Revision); diff --git a/crates/needle-platform-codex/src/patcher.rs b/crates/needle-platform-codex/src/patcher.rs index fada30d..4e7cb24 100644 --- a/crates/needle-platform-codex/src/patcher.rs +++ b/crates/needle-platform-codex/src/patcher.rs @@ -121,12 +121,13 @@ impl CodexPatchWorker { let change_id = unique_change_id(request_digest); let store = RuntimeStore::new(self.data_directory.join("needle.sqlite3")); if let Err(error) = store.initialize().and_then(|()| { - store.record_change_request( + store.record_change_request_with_provenance( &change_id, sandbox.snapshot().repository_id, sandbox.snapshot().source_digest, request_digest, request, + config.role_profile_provenance.as_ref(), ) }) { return cleanup_sandbox_after_error(sandbox, error.to_string()); @@ -342,20 +343,21 @@ impl CodexPatchWorker { ), discrepancies, }; - if let Err(error) = store.record_prepared_change( + if let Err(error) = store.record_prepared_change_with_provenance( sandbox.snapshot().repository_id, request_digest, request, &patch, &turn.response, &blobs, + config.role_profile_provenance.as_ref(), ) { return cleanup_sandbox_after_error( sandbox, format!("cannot persist patch before cleanup: {error}"), ); } - if let Err(error) = store.record_patch_attempt( + if let Err(error) = store.record_patch_attempt_with_provenance( &change_id, patch_id, &json!({ @@ -373,6 +375,7 @@ impl CodexPatchWorker { }), None, current_unix_ms(), + config.role_profile_provenance.as_ref(), ) { let reason = format!("cannot persist patch attempt accounting: {error}"); return cleanup_sandbox_after_error(sandbox, reason); diff --git a/crates/needle-platform-codex/src/verifier.rs b/crates/needle-platform-codex/src/verifier.rs index 3695902..184eb85 100644 --- a/crates/needle-platform-codex/src/verifier.rs +++ b/crates/needle-platform-codex/src/verifier.rs @@ -291,7 +291,13 @@ impl CodexVerifier { "output_tokens": output_tokens }); store - .record_verification_artifact(&artifact, &attempt, &usage, None) + .record_verification_artifact_with_provenance( + &artifact, + &attempt, + &usage, + None, + config.role_profile_provenance.as_ref(), + ) .map_err(|error| error.to_string())?; Ok(VerifyChangeOutcome { artifact, @@ -307,6 +313,15 @@ impl CodexVerifier { &self, change_id: &ChangeId, reason: &str, + ) -> Result { + self.record_inconclusive_with_provenance(change_id, reason, None) + } + + pub fn record_inconclusive_with_provenance( + &self, + change_id: &ChangeId, + reason: &str, + role_profile_provenance: Option<&needle_core::RoleProfileProvenance>, ) -> Result { let store = RuntimeStore::new(self.data_directory.join("needle.sqlite3")); store.initialize().map_err(|error| error.to_string())?; @@ -356,7 +371,7 @@ impl CodexVerifier { created_unix_ms, }; store - .record_verification_artifact( + .record_verification_artifact_with_provenance( &artifact, &json!({"phase": "repair_orchestration", "verifier_started": false}), &json!({ @@ -365,6 +380,7 @@ impl CodexVerifier { "output_tokens": 0 }), None, + role_profile_provenance, ) .map_err(|error| error.to_string())?; Ok(VerifyChangeOutcome { diff --git a/crates/needle-platform-codex/src/worker.rs b/crates/needle-platform-codex/src/worker.rs index 287782a..98274f9 100644 --- a/crates/needle-platform-codex/src/worker.rs +++ b/crates/needle-platform-codex/src/worker.rs @@ -507,6 +507,7 @@ impl CodexWorker { discarded_facts, worker_session_id: session_id, session_cleanup_success: cleanup_success, + role_profile_provenance: config.role_profile_provenance.clone(), }) } Err((code, diagnostic)) => Err(Box::new(WorkerFailure { @@ -522,6 +523,7 @@ impl CodexWorker { discarded_facts, worker_session_id: session_id, session_cleanup_success: cleanup_success, + role_profile_provenance: config.role_profile_provenance.clone(), })), } } diff --git a/crates/needle-platform-codex/tests/main_interrupt.rs b/crates/needle-platform-codex/tests/main_interrupt.rs index 9f35fbd..a03e9eb 100644 --- a/crates/needle-platform-codex/tests/main_interrupt.rs +++ b/crates/needle-platform-codex/tests/main_interrupt.rs @@ -1,7 +1,10 @@ use needle_core::{ - ApprovalDecision, ApprovalDecisionSource, CommandClassification, Digest, EvidenceFailurePolicy, - NeedCoordination, NeedDelivery, NeedStepRelation, PredicateKind, TestPlan, WorkerConfig, - built_in_route_contracts, classify_need_step, compile_need, + ApprovalDecision, ApprovalDecisionSource, CodexHost, CodexRole, CommandClassification, + CommandPolicy, Digest, EvidenceFailurePolicy, FallbackPolicy, FilesystemPolicy, + NeedCoordination, NeedDelivery, NeedStepRelation, NetworkPolicy, PredicateKind, RepairPolicy, + RoleProfileBudget, RoleProfileDefinition, RoleProfileDefinitionInput, RoleProfileId, + ServiceTier, TestPlan, TestPolicy, ToolPolicy, WorkerConfig, built_in_route_contracts, + classify_need_step, compile_need, }; use needle_platform_codex::{ CodexMainSession, CodexWorker, HookConfig, MainNeedRelation, MainSessionConfig, MainTurnResult, @@ -17,6 +20,40 @@ use std::time::{Duration, Instant}; const SIMULATOR: &str = env!("CARGO_BIN_EXE_needle-sim-codex"); static NEXT_ID: AtomicU64 = AtomicU64::new(1); +fn active_role_profile(store: &RuntimeStore, prompt_profile_digest: Digest) -> RoleProfileId { + let profile_id = RoleProfileId::new("main-interrupt.explorer").unwrap(); + let definition = RoleProfileDefinition::new(RoleProfileDefinitionInput { + profile_id: profile_id.clone(), + role: CodexRole::Explorer, + host: CodexHost::Codex, + model: "simulated-worker".to_owned(), + reasoning: needle_core::ReasoningLevel::Medium, + service_tier: ServiceTier::Default, + timeout_seconds: 10, + budget: RoleProfileBudget { + max_turns: 2, + max_output_tokens: 1200, + max_cost_microusd: 1000, + }, + prompt_profile_digest, + output_contract_digest: Digest::blake3(needle_core::ARTIFACT_RESULT_SCHEMA_ID), + tool_policy: ToolPolicy::ReadOnly, + command_policy: CommandPolicy::ReadOnly, + filesystem_policy: FilesystemPolicy::ReadOnlyCheckout, + network_policy: NetworkPolicy::Denied, + test_policy: TestPolicy::Disabled, + repair_policy: RepairPolicy::None, + fallback_policy: FallbackPolicy::Native, + concurrency: 1, + route_assignments: Vec::new(), + }) + .unwrap(); + store.create_role_profile(definition).unwrap(); + let state = store.role_profile_state(&profile_id).unwrap(); + store.activate_role_profile(&profile_id, 1, state.state_digest).unwrap(); + profile_id +} + #[test] fn direct_main_completes_without_semantic_interrupt() { let root = temporary_root(); @@ -38,6 +75,7 @@ fn direct_main_completes_without_semantic_interrupt() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start_pilot(MainSessionConfig { codex: &config, @@ -88,6 +126,7 @@ fn pilot_main_auto_approves_one_bounded_repository_read() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start_pilot(MainSessionConfig { codex: &config, @@ -137,6 +176,7 @@ fn pilot_main_fails_fast_on_r84_style_script_and_preserves_usage() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start_pilot(MainSessionConfig { codex: &config, @@ -191,6 +231,7 @@ fn semantic_message_interrupts_before_tools_and_continues_same_thread() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let source_digest = Digest::blake3(b"main-interrupt-source"); let mut session = CodexMainSession::start(MainSessionConfig { @@ -255,6 +296,7 @@ fn nested_continuation_need_is_preserved_and_classified_before_cleanup() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start(MainSessionConfig { codex: &config, @@ -321,6 +363,7 @@ fn sequential_turns_accept_two_needs_then_return_a_final_response() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start(MainSessionConfig { codex: &config, @@ -548,6 +591,7 @@ fn r35_cache_main_simulation_returns_the_frontier_answer_without_tools() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start(MainSessionConfig { codex: &config, @@ -615,6 +659,7 @@ fn malformed_subject_is_interrupted_and_preserved_before_validation() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut session = CodexMainSession::start(MainSessionConfig { codex: &config, @@ -681,6 +726,7 @@ fn supervised_main_resolves_with_worker_then_continues_without_discovery() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let source_digest = Digest::blake3(b"main-resolve-source"); let repository_id = Digest::blake3(b"main-resolve-repository"); @@ -698,12 +744,14 @@ fn supervised_main_resolves_with_worker_then_continues_without_discovery() { }) .unwrap(); let session_id = session.thread_id().to_owned(); + let role_profile_id = active_role_profile(&store, profile.definition_digest); store - .record_session_start( + .record_session_start_profiled( &session_id, profile.definition_digest, Some("simulated-main"), repository.to_str(), + &role_profile_id, ) .unwrap(); @@ -774,6 +822,7 @@ fn main_session(scenario: &str) -> (PathBuf, PathBuf, CodexMainSession) { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let session = CodexMainSession::start(MainSessionConfig { codex: &config, diff --git a/crates/needle-platform-codex/tests/offline_n1.rs b/crates/needle-platform-codex/tests/offline_n1.rs index fea1937..422a28a 100644 --- a/crates/needle-platform-codex/tests/offline_n1.rs +++ b/crates/needle-platform-codex/tests/offline_n1.rs @@ -1,4 +1,9 @@ -use needle_core::{Digest, EvidenceFailurePolicy, NeedRequest, TestPlan, WorkerConfig}; +use needle_core::{ + CodexHost, CodexRole, CommandPolicy, Digest, EvidenceFailurePolicy, FallbackPolicy, + FilesystemPolicy, NeedRequest, NetworkPolicy, RepairPolicy, RoleProfileBudget, + RoleProfileDefinition, RoleProfileDefinitionInput, RoleProfileId, ServiceTier, TestPlan, + TestPolicy, ToolPolicy, WorkerConfig, +}; use needle_platform_codex::{CodexWorker, HookConfig, StopInput, handle_stop_with_resolver}; use needle_runtime::{ResolveRequest, RuntimeEngine, RuntimeSettings, RuntimeStore}; use std::fs; @@ -103,6 +108,7 @@ fn transport_preflight_reports_optional_runner_unavailable_without_a_model_turn( service_tier: None, timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }, &repository, "trace.state-flow", @@ -165,6 +171,40 @@ struct Simulation { sandboxes_cleaned: bool, } +fn active_role_profile(store: &RuntimeStore, prompt_profile_digest: Digest) -> RoleProfileId { + let profile_id = RoleProfileId::new("offline.explorer").unwrap(); + let definition = RoleProfileDefinition::new(RoleProfileDefinitionInput { + profile_id: profile_id.clone(), + role: CodexRole::Explorer, + host: CodexHost::Codex, + model: "simulated-worker".to_owned(), + reasoning: needle_core::ReasoningLevel::Medium, + service_tier: ServiceTier::Default, + timeout_seconds: 10, + budget: RoleProfileBudget { + max_turns: 2, + max_output_tokens: 1200, + max_cost_microusd: 1000, + }, + prompt_profile_digest, + output_contract_digest: Digest::blake3(needle_core::ARTIFACT_RESULT_SCHEMA_ID), + tool_policy: ToolPolicy::ReadOnly, + command_policy: CommandPolicy::ReadOnly, + filesystem_policy: FilesystemPolicy::ReadOnlyCheckout, + network_policy: NetworkPolicy::Denied, + test_policy: TestPolicy::Disabled, + repair_policy: RepairPolicy::Once, + fallback_policy: FallbackPolicy::Native, + concurrency: 1, + route_assignments: Vec::new(), + }) + .unwrap(); + store.create_role_profile(definition).unwrap(); + let state = store.role_profile_state(&profile_id).unwrap(); + store.activate_role_profile(&profile_id, 1, state.state_digest).unwrap(); + profile_id +} + impl Simulation { fn run(scenario: &str, repetition: u32) -> Self { let root = temporary_root(scenario, repetition); @@ -188,12 +228,14 @@ impl Simulation { let session_id = format!("offline-session-{repetition}"); let turn_id = format!("offline-turn-{repetition}"); let profile = HookConfig::default().profile().unwrap(); + let role_profile_id = active_role_profile(&store, profile.definition_digest); store - .record_session_start( + .record_session_start_profiled( &session_id, profile.definition_digest, Some("simulated-main"), repository.to_str(), + &role_profile_id, ) .unwrap(); store @@ -242,11 +284,12 @@ impl Simulation { let cache_session = format!("offline-cache-session-{repetition}"); let cache_turn = format!("offline-cache-turn-{repetition}"); store - .record_session_start( + .record_session_start_profiled( &cache_session, profile.definition_digest, Some("simulated-main"), repository.to_str(), + &role_profile_id, ) .unwrap(); store @@ -275,11 +318,12 @@ impl Simulation { let hit_session = format!("offline-hit-session-{repetition}"); let hit_turn = format!("offline-hit-turn-{repetition}"); store - .record_session_start( + .record_session_start_profiled( &hit_session, profile.definition_digest, Some("simulated-main"), repository.to_str(), + &role_profile_id, ) .unwrap(); store @@ -308,6 +352,12 @@ impl Simulation { false }; let worker = store.latest_worker_run().unwrap().expect("worker run"); + let worker_provenance = worker + .role_profile_provenance + .as_ref() + .expect("profiled sessions must retain worker provenance"); + assert_eq!(&worker_provenance.profile_id, &role_profile_id); + assert_eq!(worker_provenance.revision, 1); let worker_run_count = store.worker_run_count().unwrap(); let command_evidence_count = store.command_evidence_count().unwrap(); let no_pending_approvals = store.pending_approvals().unwrap().is_empty(); diff --git a/crates/needle-platform-codex/tests/patcher_offline.rs b/crates/needle-platform-codex/tests/patcher_offline.rs index 1d6c652..298f5a3 100644 --- a/crates/needle-platform-codex/tests/patcher_offline.rs +++ b/crates/needle-platform-codex/tests/patcher_offline.rs @@ -52,6 +52,7 @@ fn patcher_changes_only_disposable_checkout_and_persists_filesystem_patch() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let store = RuntimeStore::new(data.join("needle.sqlite3")); store @@ -192,6 +193,7 @@ fn repairable_patch_gets_exactly_one_revision_and_independent_reverification() { service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let store = RuntimeStore::new(data.join("needle.sqlite3")); store @@ -280,6 +282,7 @@ fn run_certified_plan_scenario(name: &str, plan_count: usize, duplicate: bool, s service_tier: Some("default".to_owned()), timeout_seconds: 10, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; store .initialize_defaults(&RuntimeSettings { diff --git a/crates/needle-runtime/src/orchestrator.rs b/crates/needle-runtime/src/orchestrator.rs index b14b343..0ac02fd 100644 --- a/crates/needle-runtime/src/orchestrator.rs +++ b/crates/needle-runtime/src/orchestrator.rs @@ -12,9 +12,9 @@ use needle_core::{ ARTIFACT_RESULT_SCHEMA_ID, Artifact, ArtifactContract, ArtifactId, ArtifactKind, ArtifactRequest, BehaviorStep, BehaviorTrace, CacheLookup, CacheResolution, CacheScope, CodeLocation, Dependency, DependencyManifest, Digest, EvidenceBrief, FrontierItem, - FrontierView, ModelPolicy, Need, NeedCacheEntry, NeedCacheIdentity, NeedFragment, NeedIr, - NeedRequest, Obligation, PredicateKind, RouteContract, RoutePlan, SemanticWorkerArtifact, - TestPlan, ValidationRecord, WorkerArtifact, WorkerArtifactResult, WorkerConfig, WorkerFailure, + FrontierView, Need, NeedCacheEntry, NeedCacheIdentity, NeedFragment, NeedIr, NeedRequest, + Obligation, PredicateKind, RouteContract, RoutePlan, SemanticWorkerArtifact, TestPlan, + ValidationRecord, WorkerArtifact, WorkerArtifactResult, WorkerConfig, WorkerFailure, WorkerOutcome, WorkerRequest, built_in_route_contracts, built_in_route_plans, compile_need, need_fragment, }; @@ -53,6 +53,8 @@ pub enum RuntimeError { Worker(Box), #[error("session context is unavailable")] MissingSession, + #[error("session role-profile provenance is unknown")] + RoleProfileProvenanceUnknown, #[error("no route matches this request")] NoRoute, #[error("worker result changed the source snapshot")] @@ -177,6 +179,9 @@ impl RuntimeEngine { ) -> Result { let session = self.store.session(&request.session_id)?.ok_or(RuntimeError::MissingSession)?; + if session.role_profile_provenance.is_none() { + return Err(RuntimeError::RoleProfileProvenanceUnknown); + } let root_task = session.root_task.ok_or(RuntimeError::MissingSession)?; if session.turn_id.as_deref().is_some_and(|turn| turn != request.turn_id) { return Err(RuntimeError::MissingSession); @@ -261,7 +266,10 @@ impl RuntimeEngine { } let preset = self.store.preset(&route.preset_id)?.ok_or(RuntimeError::NoRoute)?; let settings = self.store.settings()?; - let worker_config = settings.worker_config(); + let worker_config = self.store.resolve_session_worker_config( + &request.session_id, + settings.codex_executable.clone(), + )?; let trusted_test_execution = settings.trusted_test_execution; let artifact_request = evidence_brief_request(&request.need, snapshot.repository_id, snapshot.source_digest); @@ -275,6 +283,7 @@ impl RuntimeEngine { normalized_request_digest: request.need.digest(), worker_configuration_digest: worker_config.digest(), output_schema_digest: Digest::blake3(ARTIFACT_RESULT_SCHEMA_ID), + role_profile_provenance: worker_config.role_profile_provenance.clone(), }; let mut worker_request = WorkerRequest { root_task, @@ -580,15 +589,13 @@ impl RuntimeEngine { return Err(RuntimeError::LeaseExpired); } } - let outcome = - match self.generate_with_model_policy(&worker_config, &worker_request, identity_digest) - { - Ok(outcome) => outcome, - Err(error) => { - self.store.release_lease(identity_digest, &owner)?; - return Err(error); - } - }; + let outcome = match self.generate(&worker_config, &worker_request, identity_digest) { + Ok(outcome) => outcome, + Err(error) => { + self.store.release_lease(identity_digest, &owner)?; + return Err(error); + } + }; let entry = NeedCacheEntry { identity, result: outcome.result.clone(), @@ -685,79 +692,6 @@ impl RuntimeEngine { Ok(outcome) } - fn generate_with_model_policy( - &self, - base: &WorkerConfig, - request: &WorkerRequest, - identity_digest: Digest, - ) -> Result { - let policy = self.store.model_policy()?; - let (profiles, repair_once, native_fallback) = match policy { - ModelPolicy::FixedOrder { profiles, repair_once, native_fallback } => { - (profiles, repair_once, native_fallback) - } - ModelPolicy::CheapestValidatedFirst { promoted_profiles, native_fallback } => { - let promoted = self.store.promoted_profile_digests(request.need_key.as_str())?; - if promoted_profiles - .iter() - .any(|profile| !promoted.contains(&profile.definition_digest)) - { - return Err(RuntimeError::ModelPolicy( - "CheapestValidatedFirst includes an unpromoted route/profile pair" - .to_owned(), - )); - } - (promoted_profiles, false, native_fallback) - } - }; - if profiles.is_empty() { - return Err(RuntimeError::ModelPolicy("the fixed model order is empty".to_owned())); - } - let mut last_error = None; - for profile in profiles { - if profile.platform != "codex" { - return Err(RuntimeError::ModelPolicy( - "only Codex worker profiles are supported".to_owned(), - )); - } - let mut config = base.clone(); - config.model = profile.model; - config.reasoning = profile.reasoning; - config.service_tier = profile.service_tier; - if repair_once { - config.evidence_failure_policy = needle_core::EvidenceFailurePolicy::RepairOnce; - } - let attempt_identity = Digest::blake3(format!( - "needle-ladder-attempt\n{identity_digest}\n{}\n", - profile.definition_digest - )); - if self.store.negative_attempt(attempt_identity)?.is_some() { - continue; - } - match self.generate(&config, request, attempt_identity) { - Ok(outcome) => return Ok(outcome), - Err(error) => { - if let Some((code, diagnostic)) = cacheable_negative_failure(&error) { - self.store.record_negative_attempt( - attempt_identity, - code, - diagnostic, - now_ms().saturating_add(300_000), - )?; - } - last_error = Some(error); - } - } - } - if native_fallback { - Err(RuntimeError::NativeFallback) - } else { - Err(last_error.unwrap_or_else(|| { - RuntimeError::ModelPolicy("model ladder produced no attempt".to_owned()) - })) - } - } - fn wait_for_result( &self, context: WaitForResult<'_>, @@ -834,15 +768,6 @@ impl RuntimeEngine { } } -fn cacheable_negative_failure(error: &RuntimeError) -> Option<(&str, &str)> { - let RuntimeError::Worker(failure) = error else { - return None; - }; - ["test_evidence_invalid", "no_valid_evidence", "evidence_binding_failed", "semantic_validation"] - .contains(&failure.code.as_str()) - .then_some((failure.code.as_str(), failure.diagnostic.as_str())) -} - fn evidence_brief_request( need: &NeedRequest, repository_id: Digest, @@ -2335,10 +2260,12 @@ mod tests { use super::*; use crate::{RuntimeSettings, validate_semantic_artifact}; use needle_core::{ - Claim, ClaimKind, ClaimPayload, CommandExecutionEvidence, EvidenceFailurePolicy, - EvidenceReference, FlowStepRole, LocationRole, NeedResult, SemanticArtifactResult, - SemanticFlowStep, SemanticLocation, Uncertainty, WorkerFailure, WorkerOutcome, - WorkerProfile, + Claim, ClaimKind, ClaimPayload, CodexHost, CodexRole, CommandExecutionEvidence, + CommandPolicy, EvidenceFailurePolicy, EvidenceReference, FallbackPolicy, FilesystemPolicy, + FlowStepRole, LocationRole, ModelPolicy, NeedResult, NetworkPolicy, RepairPolicy, + RoleProfileBudget, RoleProfileDefinition, RoleProfileDefinitionInput, RoleProfileId, + SemanticArtifactResult, SemanticFlowStep, SemanticLocation, ServiceTier, TestPolicy, + ToolPolicy, Uncertainty, WorkerFailure, WorkerOutcome, WorkerProfile, }; use std::fs; use std::process::Command; @@ -2481,6 +2408,7 @@ mod tests { discarded_facts: 0, worker_session_id: None, session_cleanup_success: None, + role_profile_provenance: None, })); } FakeWorker { spawns: Arc::new(AtomicUsize::new(0)), delay: Duration::ZERO } @@ -2537,6 +2465,7 @@ mod tests { discarded_facts: 0, worker_session_id: None, session_cleanup_success: None, + role_profile_provenance: config.role_profile_provenance.clone(), }) } } @@ -2749,6 +2678,7 @@ mod tests { root: PathBuf, store: RuntimeStore, prompt_digest: Digest, + profile_id: RoleProfileId, } impl TestContext { @@ -2790,13 +2720,54 @@ mod tests { multi_need_policy: needle_core::MultiNeedPolicy::default(), }) .unwrap(); + let profile_id = RoleProfileId::new("test.profile").unwrap(); + let definition = RoleProfileDefinition::new(RoleProfileDefinitionInput { + profile_id: profile_id.clone(), + role: CodexRole::Explorer, + host: CodexHost::Codex, + model: "worker".to_owned(), + reasoning: needle_core::ReasoningLevel::Medium, + service_tier: ServiceTier::Default, + timeout_seconds: 5, + budget: RoleProfileBudget { + max_turns: 2, + max_output_tokens: 1200, + max_cost_microusd: 1000, + }, + prompt_profile_digest: Digest::blake3(b"prompt"), + output_contract_digest: Digest::blake3(b"output"), + tool_policy: ToolPolicy::ReadOnly, + command_policy: CommandPolicy::ReadOnly, + filesystem_policy: FilesystemPolicy::ReadOnlyCheckout, + network_policy: NetworkPolicy::Denied, + test_policy: TestPolicy::Disabled, + repair_policy: RepairPolicy::None, + fallback_policy: FallbackPolicy::Native, + concurrency: 1, + route_assignments: Vec::new(), + }) + .unwrap(); + store.create_role_profile(definition).unwrap(); + store + .activate_role_profile( + &profile_id, + 1, + store.role_profile_state(&profile_id).unwrap().state_digest, + ) + .unwrap(); store.mark_utility_gate_passed().unwrap(); - Self { root, store, prompt_digest: Digest::blake3("profile") } + Self { root, store, prompt_digest: Digest::blake3("profile"), profile_id } } fn request(&self, session: &str) -> ResolveRequest { self.store - .record_session_start(session, self.prompt_digest, Some("main"), self.root.to_str()) + .record_session_start_profiled( + session, + self.prompt_digest, + Some("main"), + self.root.to_str(), + &self.profile_id, + ) .unwrap(); self.store .record_user_prompt(session, Some("turn"), "Trace the answer.", self.root.to_str()) @@ -2971,11 +2942,12 @@ mod tests { let semantic_request = |session: &str, body: &str| { context .store - .record_session_start( + .record_session_start_profiled( session, context.prompt_digest, Some("main"), context.root.to_str(), + &context.profile_id, ) .unwrap(); context @@ -3295,11 +3267,12 @@ mod tests { .unwrap(); context .store - .record_session_start( + .record_session_start_profiled( "claim-authority-hit", context.prompt_digest, Some("main"), context.root.to_str(), + &context.profile_id, ) .unwrap(); context @@ -3357,11 +3330,12 @@ mod tests { .compatibility_request(); context .store - .record_session_start( + .record_session_start_profiled( "claim-authority-partial", context.prompt_digest, Some("main"), context.root.to_str(), + &context.profile_id, ) .unwrap(); context @@ -3478,11 +3452,12 @@ mod tests { let semantic_request = |session: &str, marker: &str| { context .store - .record_session_start( + .record_session_start_profiled( session, context.prompt_digest, Some("main"), context.root.to_str(), + &context.profile_id, ) .unwrap(); context @@ -3636,11 +3611,12 @@ mod tests { let semantic_request = |session: &str, marker: &str| { context .store - .record_session_start( + .record_session_start_profiled( session, context.prompt_digest, Some("main"), context.root.to_str(), + &context.profile_id, ) .unwrap(); context @@ -3886,7 +3862,7 @@ mod tests { } #[test] - fn fixed_model_order_escalates_after_one_worker_repairs_and_fails() { + fn active_session_ignores_mutable_model_policy_ladder() { let context = TestContext::create(); context .store @@ -3904,7 +3880,7 @@ mod tests { RuntimeEngine::new(context.store.clone(), LadderWorker { calls: calls.clone() }); let outcome = engine.resolve(&context.request("ladder")).unwrap(); assert_eq!(outcome.status, "generated"); - assert_eq!(*calls.lock().unwrap(), vec!["cheap", "strong"]); + assert_eq!(*calls.lock().unwrap(), vec!["worker"]); let _ = fs::remove_dir_all(context.root.parent().unwrap()); } @@ -4123,6 +4099,7 @@ mod tests { discarded_facts: 0, worker_session_id: Some("worker-r43-offline".to_owned()), session_cleanup_success: Some(true), + role_profile_provenance: None, }; let reused = BTreeMap::new(); let materialized = materialize_worker_artifact( diff --git a/crates/needle-runtime/src/store.rs b/crates/needle-runtime/src/store.rs index 47de655..1c275be 100644 --- a/crates/needle-runtime/src/store.rs +++ b/crates/needle-runtime/src/store.rs @@ -5,10 +5,10 @@ use needle_core::{ CommandClassification, CommandExecutionEvidence, Digest, EvidenceFailurePolicy, MainTurnOutcome, ModelPolicy, MultiNeedPolicy, Need, NeedCacheEntry, NeedCacheIdentity, NeedFragment, NeedIr, NeedStep, NeedStepRelation, NeedStepState, Preset, - ReuseSufficiencyCertificate, Route, SelectedPlan, SemanticInterrupt, TestPlan, WorkerConfig, - WorkerFailure, WorkerOutcome, WorkerProfile, built_in_capability_classes, - built_in_claim_capability_classes, built_in_predicate_contracts, built_in_route_contracts, - built_in_route_plans, + ReuseSufficiencyCertificate, RoleProfileId, RoleProfileProvenance, RoleProfileRevision, Route, + SelectedPlan, SemanticInterrupt, TestPlan, WorkerConfig, WorkerFailure, WorkerOutcome, + WorkerProfile, built_in_capability_classes, built_in_claim_capability_classes, + built_in_predicate_contracts, built_in_route_contracts, built_in_route_plans, }; use rusqlite::{Connection, OptionalExtension, params}; use serde::{Deserialize, Serialize}; @@ -791,6 +791,99 @@ BEGIN END; "#; +const MIGRATION_V15: &str = r#" +ALTER TABLE sessions ADD COLUMN role_profile_id TEXT; +ALTER TABLE sessions ADD COLUMN role_profile_revision INTEGER; +ALTER TABLE sessions ADD COLUMN role_profile_definition_digest TEXT; +ALTER TABLE worker_runs ADD COLUMN role_profile_id TEXT; +ALTER TABLE worker_runs ADD COLUMN role_profile_revision INTEGER; +ALTER TABLE worker_runs ADD COLUMN role_profile_definition_digest TEXT; +ALTER TABLE change_attempts ADD COLUMN role_profile_id TEXT; +ALTER TABLE change_attempts ADD COLUMN role_profile_revision INTEGER; +ALTER TABLE change_attempts ADD COLUMN role_profile_definition_digest TEXT; +ALTER TABLE change_requests ADD COLUMN role_profile_id TEXT; +ALTER TABLE change_requests ADD COLUMN role_profile_revision INTEGER; +ALTER TABLE change_requests ADD COLUMN role_profile_definition_digest TEXT; +CREATE TRIGGER sessions_role_profile_provenance_all_or_none_insert +BEFORE INSERT ON sessions +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'session role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER sessions_role_profile_provenance_all_or_none_update +BEFORE UPDATE OF role_profile_id, role_profile_revision, role_profile_definition_digest +ON sessions +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'session role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER sessions_role_profile_provenance_immutable +BEFORE UPDATE OF role_profile_id, role_profile_revision, role_profile_definition_digest +ON sessions +WHEN NEW.role_profile_id IS NOT OLD.role_profile_id + OR NEW.role_profile_revision IS NOT OLD.role_profile_revision + OR NEW.role_profile_definition_digest IS NOT OLD.role_profile_definition_digest +BEGIN + SELECT RAISE(ABORT, 'session role-profile provenance is immutable'); +END; +CREATE TRIGGER worker_runs_role_profile_provenance_all_or_none_insert +BEFORE INSERT ON worker_runs +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'worker-run role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER worker_runs_role_profile_provenance_all_or_none_update +BEFORE UPDATE OF role_profile_id, role_profile_revision, role_profile_definition_digest +ON worker_runs +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'worker-run role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER change_attempts_role_profile_provenance_all_or_none_insert +BEFORE INSERT ON change_attempts +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'change-attempt role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER change_attempts_role_profile_provenance_all_or_none_update +BEFORE UPDATE OF role_profile_id, role_profile_revision, role_profile_definition_digest +ON change_attempts +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'change-attempt role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER change_requests_role_profile_provenance_all_or_none_insert +BEFORE INSERT ON change_requests +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'change-request role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER change_requests_role_profile_provenance_all_or_none_update +BEFORE UPDATE OF role_profile_id, role_profile_revision, role_profile_definition_digest +ON change_requests +WHEN (NEW.role_profile_id IS NULL) != (NEW.role_profile_revision IS NULL) + OR (NEW.role_profile_id IS NULL) != (NEW.role_profile_definition_digest IS NULL) +BEGIN + SELECT RAISE(ABORT, 'change-request role-profile provenance must be all NULL or all set'); +END; +CREATE TRIGGER change_requests_role_profile_provenance_immutable +BEFORE UPDATE OF role_profile_id, role_profile_revision, role_profile_definition_digest +ON change_requests +WHEN NEW.role_profile_id IS NOT OLD.role_profile_id + OR NEW.role_profile_revision IS NOT OLD.role_profile_revision + OR NEW.role_profile_definition_digest IS NOT OLD.role_profile_definition_digest +BEGIN + SELECT RAISE(ABORT, 'change-request role-profile provenance is immutable'); +END; +"#; + #[derive(Debug, Error)] pub enum StoreError { #[error("database operation failed: {0}")] @@ -912,6 +1005,7 @@ pub struct SessionRecord { pub route_set: Vec, pub multi_need_policy: MultiNeedPolicy, pub multi_need_policy_digest: Digest, + pub role_profile_provenance: Option, } #[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] @@ -973,6 +1067,7 @@ impl RuntimeSettings { service_tier: None, timeout_seconds: self.worker_timeout_seconds, evidence_failure_policy: self.evidence_failure_policy, + role_profile_provenance: None, } } } @@ -1010,6 +1105,7 @@ pub struct WorkerRunRecord { pub repair_performed: bool, pub worker_session_id: Option, pub session_cleanup_success: Option, + pub role_profile_provenance: Option, } #[derive(Clone, Debug, Eq, PartialEq, Serialize)] @@ -1154,6 +1250,7 @@ impl RuntimeStore { apply_migration(&mut connection, 12, MIGRATION_V12)?; apply_migration(&mut connection, 13, MIGRATION_V13)?; apply_migration(&mut connection, 14, MIGRATION_V14)?; + apply_migration(&mut connection, 15, MIGRATION_V15)?; connection.execute( "INSERT OR IGNORE INTO settings(key, value) VALUES('utility_gate_passed', '0')", [], @@ -2051,6 +2148,158 @@ impl RuntimeStore { Ok(()) } + /// Starts a production session bound to the currently active revision of + /// an explicitly selected role profile. The active lookup and immutable + /// session insert share one immediate transaction. + pub fn record_session_start_profiled( + &self, + session_id: &str, + prompt_profile_digest: Digest, + model: Option<&str>, + cwd: Option<&str>, + profile_id: &RoleProfileId, + ) -> Result<(), StoreError> { + self.record_session_start_for_transport_profiled( + session_id, + prompt_profile_digest, + model, + cwd, + "hook", + needle_core::need_grammar_definition_digest(), + Some(needle_core::need_grammar_definition_digest()), + profile_id, + ) + } + + #[allow(clippy::too_many_arguments)] + pub fn record_session_start_for_transport_profiled( + &self, + session_id: &str, + prompt_profile_digest: Digest, + model: Option<&str>, + cwd: Option<&str>, + transport: &str, + transport_definition_digest: Digest, + need_grammar_digest: Option, + profile_id: &RoleProfileId, + ) -> Result<(), StoreError> { + self.initialize()?; + let route_set = self.routes()?; + let route_set_digest = route_set_digest(&route_set); + let multi_need_policy = self.multi_need_policy()?; + let multi_need_policy_digest = multi_need_policy.digest(); + let mut connection = self.connection()?; + let transaction = + connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?; + + let active_revision: Option = transaction + .query_row( + "SELECT active_revision FROM role_profile_state WHERE profile_id=?1", + [profile_id.as_str()], + |row| row.get(0), + ) + .optional()? + .flatten(); + let Some(revision) = active_revision else { + return Err(StoreError::RoleProfileConflict(format!( + "profile {profile_id} has no active revision" + ))); + }; + let (definition_json, row_digest, created_unix_ms, activated_unix_ms): ( + String, + String, + u64, + Option, + ) = transaction.query_row( + "SELECT definition_json, definition_digest, created_unix_ms, activated_unix_ms + FROM role_profile_revisions WHERE profile_id=?1 AND revision=?2", + params![profile_id.as_str(), revision], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)), + )?; + let definition: needle_core::RoleProfileDefinition = + serde_json::from_str(&definition_json)?; + definition + .validate() + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string()))?; + let parsed_row_digest = Digest::parse(&row_digest) + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string()))?; + if definition.profile_id != *profile_id + || definition.definition_digest != parsed_row_digest + || activated_unix_ms.is_none() + { + return Err(StoreError::RoleProfileCorruption( + "active role-profile revision identity or activation metadata is invalid" + .to_owned(), + )); + } + let revision_record = RoleProfileRevision { + profile_id: profile_id.clone(), + revision, + definition, + state: needle_core::RoleProfileState::Active, + created_unix_ms, + activated_unix_ms, + }; + let provenance = RoleProfileProvenance::from_revision(&revision_record) + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string()))?; + transaction.execute( + "INSERT INTO sessions( + session_id, prompt_profile_digest, route_set_digest, model, cwd, updated_unix_ms, + route_set_json, need_grammar_digest, multi_need_policy_json, + multi_need_policy_digest, transport, transport_definition_digest, + semantic_definition_digest, role_profile_id, role_profile_revision, + role_profile_definition_digest + ) + VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16) + ON CONFLICT(session_id) DO NOTHING", + params![ + session_id, + prompt_profile_digest.to_string(), + route_set_digest.to_string(), + model, + cwd, + now_ms(), + serde_json::to_string(&route_set)?, + need_grammar_digest.map(|digest| digest.to_string()), + serde_json::to_string(&multi_need_policy)?, + multi_need_policy_digest.to_string(), + transport, + transport_definition_digest.to_string(), + needle_core::need_ir_definition_digest().to_string(), + provenance.profile_id.as_str(), + provenance.revision, + provenance.definition_digest.to_string(), + ], + )?; + let existing: (Option, Option, Option) = transaction.query_row( + "SELECT role_profile_id, role_profile_revision, role_profile_definition_digest + FROM sessions WHERE session_id=?1", + [session_id], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + let existing = parse_role_profile_provenance(existing)?; + if existing.as_ref() != Some(&provenance) { + return Err(StoreError::RoleProfileConflict(format!( + "session {session_id} is already bound to a different role-profile revision" + ))); + } + transaction.commit()?; + Ok(()) + } + + /// Explicitly named legacy insertion helper retained for migration and + /// backward-compatibility tests. Production entry points must use one of + /// the profiled methods above. + pub fn record_legacy_session_start( + &self, + session_id: &str, + prompt_profile_digest: Digest, + model: Option<&str>, + cwd: Option<&str>, + ) -> Result<(), StoreError> { + self.record_session_start(session_id, prompt_profile_digest, model, cwd) + } + pub fn record_user_prompt( &self, session_id: &str, @@ -2086,7 +2335,9 @@ impl RuntimeStore { "SELECT session_id, turn_id, root_task, prompt_profile_digest, route_set_digest, model, cwd, route_set_json, need_grammar_digest, multi_need_policy_json, multi_need_policy_digest, transport, - transport_definition_digest, semantic_definition_digest + transport_definition_digest, semantic_definition_digest, + role_profile_id, role_profile_revision, + role_profile_definition_digest FROM sessions WHERE session_id=?1", [session_id], |row| { @@ -2105,6 +2356,9 @@ impl RuntimeStore { row.get::<_, Option>(11)?, row.get::<_, Option>(12)?, row.get::<_, Option>(13)?, + row.get::<_, Option>(14)?, + row.get::<_, Option>(15)?, + row.get::<_, Option>(16)?, )) }, ) @@ -2126,11 +2380,24 @@ impl RuntimeStore { transport, transport_definition_digest, semantic_definition_digest, + role_profile_id, + role_profile_revision, + role_profile_definition_digest, )| { let route_set = routes .map(|value| serde_json::from_str(&value)) .transpose()? .unwrap_or_default(); + let role_profile_provenance = parse_role_profile_provenance(( + role_profile_id, + role_profile_revision, + role_profile_definition_digest, + ))?; + if let Some(provenance) = &role_profile_provenance { + provenance.validate().map_err(|error| { + StoreError::RoleProfileCorruption(error.to_string()) + })?; + } Ok(SessionRecord { session_id, turn_id, @@ -2169,12 +2436,70 @@ impl RuntimeStore { .transpose() .map_err(|error| StoreError::Digest(error.to_string()))? .unwrap_or_else(|| MultiNeedPolicy::default().digest()), + role_profile_provenance, }) }, ) .transpose() } + /// Resolve the exact historical revision frozen on a session. Active + /// pointers are intentionally never consulted here. + pub fn resolve_session_worker_config( + &self, + session_id: &str, + executable: impl Into, + ) -> Result { + let session = self + .session(session_id)? + .ok_or_else(|| StoreError::RoleProfileNotFound(format!("session {session_id}")))?; + let provenance = session.role_profile_provenance.ok_or_else(|| { + StoreError::RoleProfileConflict( + "session has unknown role-profile provenance".to_owned(), + ) + })?; + let revision = + self.read_role_profile_revision(&provenance.profile_id, provenance.revision)?; + let actual = RoleProfileProvenance::from_revision(&revision) + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string()))?; + if actual != provenance { + return Err(StoreError::RoleProfileCorruption( + "session role-profile provenance does not match historical revision".to_owned(), + )); + } + revision + .to_worker_config(executable) + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string())) + } + + pub fn worker_config_for_session( + &self, + session_id: &str, + executable: impl Into, + ) -> Result { + self.resolve_session_worker_config(session_id, executable) + } + + /// Checks a bounded provenance value against immutable historical storage. + /// Current activation state is irrelevant; the exact revision and digest + /// must still exist and validate. + pub fn role_profile_provenance_is_historical( + &self, + provenance: &RoleProfileProvenance, + ) -> Result { + provenance + .validate() + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string()))?; + match self.read_role_profile_revision_by_digest( + &provenance.profile_id, + provenance.definition_digest, + ) { + Ok(revision) => Ok(revision.revision == provenance.revision), + Err(StoreError::RoleProfileNotFound(_)) => Ok(false), + Err(error) => Err(error), + } + } + pub fn record_need_step( &self, session_id: &str, @@ -2619,6 +2944,12 @@ impl RuntimeStore { } pub fn cache_lookup(&self, identity: &NeedCacheIdentity) -> Result { + let Some(requested_provenance) = identity.role_profile_provenance.as_ref() else { + return Ok(CacheLookup::Bypass("role-profile-provenance-unknown".to_owned())); + }; + if !self.role_profile_provenance_is_historical(requested_provenance)? { + return Ok(CacheLookup::Bypass("role-profile-provenance-invalid".to_owned())); + } let connection = self.connection()?; let digest = identity.digest().to_string(); let json: Option = connection @@ -2629,16 +2960,24 @@ impl RuntimeStore { ) .optional()?; if let Some(json) = json { - connection.execute( - "UPDATE cache_entries SET hit_count=hit_count+1 WHERE identity_digest=?1", - [&digest], - )?; let mut entry: NeedCacheEntry = serde_json::from_str(&json)?; if entry.identity.digest().to_string() != digest { return Err(StoreError::Digest( "cache entry identity does not match its primary key".to_owned(), )); } + if entry.identity.role_profile_provenance.as_ref() != Some(requested_provenance) + || entry.worker_outcome.role_profile_provenance.as_ref() + != Some(requested_provenance) + { + return Err(StoreError::ArtifactIdentity( + "cache entry role-profile provenance is inconsistent".to_owned(), + )); + } + connection.execute( + "UPDATE cache_entries SET hit_count=hit_count+1 WHERE identity_digest=?1", + [&digest], + )?; entry.hit_count = entry.hit_count.saturating_add(1); return Ok(CacheLookup::Hit(Box::new(entry))); } @@ -2654,6 +2993,21 @@ impl RuntimeStore { } pub fn publish(&self, entry: &NeedCacheEntry) -> Result<(), StoreError> { + let Some(provenance) = entry.identity.role_profile_provenance.as_ref() else { + return Err(StoreError::ArtifactIdentity( + "cannot publish cache entry without role-profile provenance".to_owned(), + )); + }; + if entry.worker_outcome.role_profile_provenance.as_ref() != Some(provenance) { + return Err(StoreError::ArtifactIdentity( + "cache identity and worker outcome role-profile provenance differ".to_owned(), + )); + } + if !self.role_profile_provenance_is_historical(provenance)? { + return Err(StoreError::ArtifactIdentity( + "cache entry references an unknown role-profile revision".to_owned(), + )); + } let connection = self.connection()?; connection.execute( "INSERT OR REPLACE INTO cache_entries(identity_digest, logical_digest, source_digest, entry_json, created_unix_ms, hit_count) @@ -3852,13 +4206,33 @@ impl RuntimeStore { pub fn record_worker_run(&self, entry: &NeedCacheEntry) -> Result<(), StoreError> { let outcome = &entry.worker_outcome; + let provenance = match ( + entry.identity.role_profile_provenance.as_ref(), + outcome.role_profile_provenance.as_ref(), + ) { + (Some(identity), Some(outcome)) if identity == outcome => { + if !self.role_profile_provenance_is_historical(identity)? { + return Err(StoreError::ArtifactIdentity( + "worker run references an unknown role-profile revision".to_owned(), + )); + } + Some(identity) + } + (None, None) => None, + _ => { + return Err(StoreError::ArtifactIdentity( + "worker run identity and outcome role-profile provenance differ".to_owned(), + )); + } + }; let connection = self.connection()?; connection.execute( "INSERT INTO worker_runs(identity_digest, model, reasoning, status, duration_ms, input_tokens, cached_input_tokens, output_tokens, result_digest, created_unix_ms, failure_code, failure_diagnostic, discarded_facts, logical_worker_spawns, - worker_turns, repair_performed, worker_session_id, session_cleanup_success) - VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL, NULL, ?11, ?12, ?13, ?14, ?15, ?16)", + worker_turns, repair_performed, worker_session_id, session_cleanup_success, + role_profile_id, role_profile_revision, role_profile_definition_digest) + VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL, NULL, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19)", params![ entry.identity.digest().to_string(), outcome.worker_model, @@ -3876,6 +4250,11 @@ impl RuntimeStore { outcome.repair_performed, outcome.worker_session_id, outcome.session_cleanup_success, + provenance.map(|value| value.profile_id.as_str()), + provenance.map(|value| value.revision), + provenance + .as_ref() + .map(|value| value.definition_digest.to_string()), ], )?; Ok(()) @@ -3887,13 +4266,33 @@ impl RuntimeStore { config: &WorkerConfig, failure: &WorkerFailure, ) -> Result<(), StoreError> { + let provenance = match ( + config.role_profile_provenance.as_ref(), + failure.role_profile_provenance.as_ref(), + ) { + (Some(config), Some(failure)) if config == failure => { + if !self.role_profile_provenance_is_historical(config)? { + return Err(StoreError::ArtifactIdentity( + "worker failure references an unknown role-profile revision".to_owned(), + )); + } + Some(config) + } + (None, None) => None, + _ => { + return Err(StoreError::ArtifactIdentity( + "worker config and failure role-profile provenance differ".to_owned(), + )); + } + }; let connection = self.connection()?; connection.execute( "INSERT INTO worker_runs(identity_digest, model, reasoning, status, duration_ms, input_tokens, cached_input_tokens, output_tokens, result_digest, created_unix_ms, failure_code, failure_diagnostic, discarded_facts, logical_worker_spawns, - worker_turns, repair_performed, worker_session_id, session_cleanup_success) - VALUES(?1, ?2, ?3, 'failed', ?4, ?5, ?6, ?7, NULL, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)", + worker_turns, repair_performed, worker_session_id, session_cleanup_success, + role_profile_id, role_profile_revision, role_profile_definition_digest) + VALUES(?1, ?2, ?3, 'failed', ?4, ?5, ?6, ?7, NULL, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19)", params![ identity.to_string(), config.model, @@ -3911,6 +4310,9 @@ impl RuntimeStore { failure.repair_performed, failure.worker_session_id, failure.session_cleanup_success, + provenance.map(|value| value.profile_id.as_str()), + provenance.map(|value| value.revision), + provenance.map(|value| value.definition_digest.to_string()), ], )?; Ok(()) @@ -3940,6 +4342,7 @@ impl RuntimeStore { discarded_facts: outcome.discarded_facts, worker_session_id: outcome.worker_session_id.clone(), session_cleanup_success: outcome.session_cleanup_success, + role_profile_provenance: outcome.role_profile_provenance.clone(), }, ) } @@ -4051,7 +4454,8 @@ impl RuntimeStore { .query_row( "SELECT input_tokens, cached_input_tokens, output_tokens, result_digest, failure_code, failure_diagnostic, discarded_facts, logical_worker_spawns, - worker_turns, repair_performed, worker_session_id, session_cleanup_success + worker_turns, repair_performed, worker_session_id, session_cleanup_success, + role_profile_id, role_profile_revision, role_profile_definition_digest FROM worker_runs ORDER BY id DESC LIMIT 1", [], |row| { @@ -4068,6 +4472,9 @@ impl RuntimeStore { row.get::<_, bool>(9)?, row.get::<_, Option>(10)?, row.get::<_, Option>(11)?, + row.get::<_, Option>(12)?, + row.get::<_, Option>(13)?, + row.get::<_, Option>(14)?, )) }, ) @@ -4087,7 +4494,15 @@ impl RuntimeStore { repair_performed, worker_session_id, session_cleanup_success, + role_profile_id, + role_profile_revision, + role_profile_definition_digest, )| { + let role_profile_provenance = parse_role_profile_provenance(( + role_profile_id, + role_profile_revision, + role_profile_definition_digest, + ))?; Ok(WorkerRunRecord { input_tokens, cached_input_tokens, @@ -4101,6 +4516,7 @@ impl RuntimeStore { repair_performed, worker_session_id, session_cleanup_success, + role_profile_provenance, }) }, ) @@ -4115,7 +4531,8 @@ impl RuntimeStore { let mut statement = connection.prepare( "SELECT input_tokens, cached_input_tokens, output_tokens, result_digest, failure_code, failure_diagnostic, discarded_facts, logical_worker_spawns, - worker_turns, repair_performed, worker_session_id, session_cleanup_success + worker_turns, repair_performed, worker_session_id, session_cleanup_success, + role_profile_id, role_profile_revision, role_profile_definition_digest FROM worker_runs ORDER BY id LIMIT -1 OFFSET ?1", )?; let rows = statement.query_map([previous_count], |row| { @@ -4132,6 +4549,9 @@ impl RuntimeStore { row.get::<_, bool>(9)?, row.get::<_, Option>(10)?, row.get::<_, Option>(11)?, + row.get::<_, Option>(12)?, + row.get::<_, Option>(13)?, + row.get::<_, Option>(14)?, )) })?; rows.map(|row| { @@ -4148,7 +4568,15 @@ impl RuntimeStore { repair_performed, worker_session_id, session_cleanup_success, + role_profile_id, + role_profile_revision, + role_profile_definition_digest, ) = row?; + let role_profile_provenance = parse_role_profile_provenance(( + role_profile_id, + role_profile_revision, + role_profile_definition_digest, + ))?; Ok(WorkerRunRecord { input_tokens, cached_input_tokens, @@ -4162,6 +4590,7 @@ impl RuntimeStore { repair_performed, worker_session_id, session_cleanup_success, + role_profile_provenance, }) }) .collect() @@ -4339,6 +4768,27 @@ fn parse_digest(value: &str) -> Result { Digest::parse(value).map_err(|error| StoreError::Digest(error.to_string())) } +fn parse_role_profile_provenance( + value: (Option, Option, Option), +) -> Result, StoreError> { + let (profile_id, revision, digest) = value; + match (profile_id, revision, digest) { + (None, None, None) => Ok(None), + (Some(profile_id), Some(revision), Some(digest)) => { + let profile_id = RoleProfileId::new(profile_id) + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string()))?; + let digest = Digest::parse(&digest) + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string()))?; + RoleProfileProvenance::new(profile_id, revision, digest) + .map(Some) + .map_err(|error| StoreError::RoleProfileCorruption(error.to_string())) + } + _ => Err(StoreError::RoleProfileCorruption( + "role-profile provenance columns are partially populated".to_owned(), + )), + } +} + fn apply_migration( connection: &mut Connection, version: u32, @@ -4727,7 +5177,7 @@ mod tests { .unwrap() .collect::, _>>() .unwrap(); - assert_eq!(versions, vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14]); + assert_eq!(versions, vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15]); let columns = connection .prepare("PRAGMA table_info(worker_runs)") .unwrap() @@ -5043,6 +5493,7 @@ mod tests { service_tier: None, timeout_seconds: 180, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; let mut repaired = base.clone(); repaired.evidence_failure_policy = EvidenceFailurePolicy::RepairOnce; @@ -5092,6 +5543,7 @@ mod tests { service_tier: None, timeout_seconds: 180, evidence_failure_policy: EvidenceFailurePolicy::DiscardInvalidFact, + role_profile_provenance: None, }; store .record_worker_failure( @@ -5110,6 +5562,7 @@ mod tests { discarded_facts: 3, worker_session_id: Some("session-1".to_owned()), session_cleanup_success: Some(false), + role_profile_provenance: None, }, ) .unwrap(); @@ -5130,6 +5583,7 @@ mod tests { discarded_facts: 3, worker_session_id: Some("session-2".to_owned()), session_cleanup_success: None, + role_profile_provenance: None, }, ) .unwrap(); diff --git a/crates/needle-runtime/src/store/changes.rs b/crates/needle-runtime/src/store/changes.rs index 7b25b74..5f21d91 100644 --- a/crates/needle-runtime/src/store/changes.rs +++ b/crates/needle-runtime/src/store/changes.rs @@ -4,7 +4,74 @@ use needle_core::{ PatchArtifact, PatchId, VerificationArtifact, VerificationStatus, }; -type PreparedChangeRow = (String, String, String, String, String, String, String, bool, u64); +type PreparedChangeRow = ( + String, + String, + String, + String, + String, + String, + String, + bool, + u64, + Option, + Option, + Option, +); + +type ExistingChangeRequestRow = + (String, String, String, String, Option, Option, Option); + +fn with_role_profile_provenance( + value: &serde_json::Value, + provenance: Option<&RoleProfileProvenance>, +) -> serde_json::Value { + let Some(provenance) = provenance else { + return value.clone(); + }; + let bounded = serde_json::json!({ + "profile_id": provenance.profile_id, + "revision": provenance.revision, + "definition_digest": provenance.definition_digest, + }); + if let Some(object) = value.as_object() { + let mut object = object.clone(); + object.insert("role_profile_provenance".to_owned(), bounded); + serde_json::Value::Object(object) + } else { + serde_json::json!({ + "attempt": value, + "role_profile_provenance": bounded, + }) + } +} + +fn change_role_profile_provenance( + connection: &Connection, + change_id: &ChangeId, +) -> Result, StoreError> { + let columns: (Option, Option, Option) = connection.query_row( + "SELECT role_profile_id, role_profile_revision, role_profile_definition_digest + FROM change_requests WHERE change_id=?1", + [change_id.to_string()], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + parse_role_profile_provenance(columns) +} + +fn require_change_role_profile_provenance( + connection: &Connection, + change_id: &ChangeId, + expected: Option<&RoleProfileProvenance>, +) -> Result, StoreError> { + let stored = change_role_profile_provenance(connection, change_id)?; + if stored.as_ref() != expected { + return Err(StoreError::ChangeConflict(format!( + "{change_id}: role-profile provenance changed" + ))); + } + Ok(stored) +} #[derive(Clone, Debug, Eq, PartialEq)] pub struct PatchFileBlob { @@ -25,6 +92,7 @@ pub struct PreparedChangeRecord { pub declared_output: serde_json::Value, pub repair_attempted: bool, pub created_unix_ms: u64, + pub role_profile_provenance: Option, } #[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)] @@ -36,6 +104,7 @@ pub struct ChangeAttemptRecord { pub usage: serde_json::Value, pub cost_microusd: Option, pub created_unix_ms: u64, + pub role_profile_provenance: Option, } impl RuntimeStore { @@ -47,24 +116,74 @@ impl RuntimeStore { request_digest: Digest, request: &ChangeRequest, ) -> Result<(), StoreError> { + self.record_change_request_with_provenance( + change_id, + repository_id, + source_snapshot, + request_digest, + request, + None, + ) + } + + pub fn record_change_request_with_provenance( + &self, + change_id: &ChangeId, + repository_id: Digest, + source_snapshot: Digest, + request_digest: Digest, + request: &ChangeRequest, + role_profile_provenance: Option<&RoleProfileProvenance>, + ) -> Result<(), StoreError> { + if let Some(provenance) = role_profile_provenance + && !self.role_profile_provenance_is_historical(provenance)? + { + return Err(StoreError::ChangeConflict(format!( + "{change_id}: unknown role-profile revision" + ))); + } let request_json = serde_json::to_string(request)?; let now = now_ms(); let mut connection = self.connection()?; let transaction = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?; - let existing: Option<(String, String, String, String)> = transaction + let existing: Option = transaction .query_row( - "SELECT request_digest, repository_id, source_snapshot_digest, request_json + "SELECT request_digest, repository_id, source_snapshot_digest, request_json, + role_profile_id, role_profile_revision, + role_profile_definition_digest FROM change_requests WHERE change_id=?1", [change_id.to_string()], - |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)), + |row| { + Ok(( + row.get(0)?, + row.get(1)?, + row.get(2)?, + row.get(3)?, + row.get(4)?, + row.get(5)?, + row.get(6)?, + )) + }, ) .optional()?; - if let Some((stored_request, stored_repository, stored_snapshot, stored_json)) = existing { + if let Some(( + stored_request, + stored_repository, + stored_snapshot, + stored_json, + profile_id, + profile_revision, + profile_digest, + )) = existing + { + let stored_provenance = + parse_role_profile_provenance((profile_id, profile_revision, profile_digest))?; if stored_request == request_digest.to_string() && stored_repository == repository_id.to_string() && stored_snapshot == source_snapshot.to_string() && stored_json == request_json + && stored_provenance.as_ref() == role_profile_provenance { transaction.commit()?; return Ok(()); @@ -75,21 +194,28 @@ impl RuntimeStore { "INSERT INTO change_requests( change_id, request_digest, repository_id, source_snapshot_digest, state, request_json, latest_patch_revision, repair_attempted, - created_unix_ms, updated_unix_ms - ) VALUES(?1, ?2, ?3, ?4, 'requested', ?5, 0, 0, ?6, ?6)", + created_unix_ms, updated_unix_ms, role_profile_id, + role_profile_revision, role_profile_definition_digest + ) VALUES(?1, ?2, ?3, ?4, 'requested', ?5, 0, 0, ?6, ?6, ?7, ?8, ?9)", params![ change_id.to_string(), request_digest.to_string(), repository_id.to_string(), source_snapshot.to_string(), request_json, - now + now, + role_profile_provenance.map(|value| value.profile_id.as_str()), + role_profile_provenance.map(|value| value.revision), + role_profile_provenance.map(|value| value.definition_digest.to_string()), ], )?; - let event_json = serde_json::to_string(&serde_json::json!({ - "request_digest": request_digest, - "state": "requested" - }))?; + let event_json = serde_json::to_string(&with_role_profile_provenance( + &serde_json::json!({ + "request_digest": request_digest, + "state": "requested" + }), + role_profile_provenance, + ))?; transaction.execute( "INSERT INTO change_events( change_id, event_type, payload_digest, payload_json, created_unix_ms @@ -115,6 +241,36 @@ impl RuntimeStore { declared_output: &serde_json::Value, file_blobs: &[PatchFileBlob], ) -> Result<(), StoreError> { + self.record_prepared_change_with_provenance( + repository_id, + request_digest, + request, + patch, + declared_output, + file_blobs, + None, + ) + } + + #[allow(clippy::too_many_arguments)] + pub fn record_prepared_change_with_provenance( + &self, + repository_id: Digest, + request_digest: Digest, + request: &ChangeRequest, + patch: &PatchArtifact, + declared_output: &serde_json::Value, + file_blobs: &[PatchFileBlob], + role_profile_provenance: Option<&RoleProfileProvenance>, + ) -> Result<(), StoreError> { + if let Some(provenance) = role_profile_provenance + && !self.role_profile_provenance_is_historical(provenance)? + { + return Err(StoreError::ChangeConflict(format!( + "{}: unknown role-profile revision", + patch.change_id + ))); + } if patch.id != PatchArtifact::compute_id(patch.source_snapshot, &patch.files) { return Err(StoreError::PatchArtifact( "patch id does not match the filesystem manifest".to_owned(), @@ -147,11 +303,14 @@ impl RuntimeStore { let manifest_json = serde_json::to_string(&patch.files)?; let declared_output_json = serde_json::to_string(declared_output)?; let discrepancies_json = serde_json::to_string(&patch.discrepancies)?; - let event_payload = serde_json::json!({ - "patch_id": patch.id, - "revision": patch.revision, - "state": "prepared" - }); + let event_payload = with_role_profile_provenance( + &serde_json::json!({ + "patch_id": patch.id, + "revision": patch.revision, + "state": "prepared" + }), + role_profile_provenance, + ); let event_json = serde_json::to_string(&event_payload)?; let event_digest = Digest::blake3(event_json.as_bytes()); let now = now_ms(); @@ -166,6 +325,13 @@ impl RuntimeStore { |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?, row.get(4)?)), ) .optional()?; + if existing.is_some() { + require_change_role_profile_provenance( + &transaction, + &patch.change_id, + role_profile_provenance, + )?; + } match &existing { None if patch.revision != 1 => { return Err(StoreError::PatchArtifact( @@ -190,8 +356,9 @@ impl RuntimeStore { transaction.execute( "INSERT INTO change_requests( change_id, request_digest, repository_id, source_snapshot_digest, - state, request_json, latest_patch_revision, created_unix_ms, updated_unix_ms - ) VALUES(?1, ?2, ?3, ?4, 'prepared', ?5, ?6, ?7, ?7) + state, request_json, latest_patch_revision, created_unix_ms, updated_unix_ms, + role_profile_id, role_profile_revision, role_profile_definition_digest + ) VALUES(?1, ?2, ?3, ?4, 'prepared', ?5, ?6, ?7, ?7, ?8, ?9, ?10) ON CONFLICT(change_id) DO NOTHING", params![ patch.change_id.to_string(), @@ -200,7 +367,10 @@ impl RuntimeStore { patch.source_snapshot.to_string(), request_json, patch.revision, - now + now, + role_profile_provenance.map(|value| value.profile_id.as_str()), + role_profile_provenance.map(|value| value.revision), + role_profile_provenance.map(|value| value.definition_digest.to_string()), ], )?; transaction.execute( @@ -268,7 +438,9 @@ impl RuntimeStore { .query_row( "SELECT c.request_json, c.request_digest, c.repository_id, c.source_snapshot_digest, c.state, p.artifact_json, - p.declared_output_json, c.repair_attempted, c.created_unix_ms + p.declared_output_json, c.repair_attempted, c.created_unix_ms, + c.role_profile_id, c.role_profile_revision, + c.role_profile_definition_digest FROM change_requests c JOIN patch_artifacts p ON p.change_id=c.change_id AND p.revision=c.latest_patch_revision @@ -285,6 +457,9 @@ impl RuntimeStore { row.get(6)?, row.get(7)?, row.get(8)?, + row.get(9)?, + row.get(10)?, + row.get(11)?, )) }, ) @@ -299,6 +474,9 @@ impl RuntimeStore { output, repair_attempted, created, + role_profile_id, + role_profile_revision, + role_profile_definition_digest, )) = row else { return Ok(None); @@ -315,6 +493,11 @@ impl RuntimeStore { declared_output: serde_json::from_str(&output)?, repair_attempted, created_unix_ms: created, + role_profile_provenance: parse_role_profile_provenance(( + role_profile_id, + role_profile_revision, + role_profile_definition_digest, + ))?, })) } @@ -346,12 +529,21 @@ impl RuntimeStore { change_id: &ChangeId, reason: &str, ) -> Result<(), StoreError> { - let payload = serde_json::json!({"reason": reason}); - let payload_json = serde_json::to_string(&payload)?; - let payload_digest = Digest::blake3(payload_json.as_bytes()); let now = now_ms(); let mut connection = self.connection()?; let transaction = connection.transaction()?; + let provenance = change_role_profile_provenance(&transaction, change_id)?; + // Provider and cleanup diagnostics can contain paths or prompt text. + // Persist only a stable classification and digest in the audit event. + let payload = with_role_profile_provenance( + &serde_json::json!({ + "reason": "change preparation failed", + "reason_digest": Digest::blake3(reason.as_bytes()), + }), + provenance.as_ref(), + ); + let payload_json = serde_json::to_string(&payload)?; + let payload_digest = Digest::blake3(payload_json.as_bytes()); let changed = transaction.execute( "UPDATE change_requests SET state='failed', updated_unix_ms=?2 WHERE change_id=?1", params![change_id.to_string(), now], @@ -404,6 +596,7 @@ impl RuntimeStore { && artifact.is_canonical() }) .ok_or_else(|| StoreError::ChangeConflict(change_id.to_string()))?; + let provenance = change_role_profile_provenance(&transaction, change_id)?; let changed = transaction.execute( "UPDATE change_requests SET state='repairing', repair_attempted=1, updated_unix_ms=?2 @@ -414,11 +607,14 @@ impl RuntimeStore { if changed != 1 { return Err(StoreError::ChangeConflict(change_id.to_string())); } - let payload_json = serde_json::to_string(&serde_json::json!({ - "patch_id": patch_id, - "verification_id": verification.id, - "state": "repairing" - }))?; + let payload_json = serde_json::to_string(&with_role_profile_provenance( + &serde_json::json!({ + "patch_id": patch_id, + "verification_id": verification.id, + "state": "repairing" + }), + provenance.as_ref(), + ))?; transaction.execute( "INSERT INTO change_events( change_id, event_type, payload_digest, payload_json, created_unix_ms @@ -454,6 +650,23 @@ impl RuntimeStore { attempt: &serde_json::Value, usage: &serde_json::Value, cost_microusd: Option, + ) -> Result<(), StoreError> { + self.record_verification_artifact_with_provenance( + artifact, + attempt, + usage, + cost_microusd, + None, + ) + } + + pub fn record_verification_artifact_with_provenance( + &self, + artifact: &VerificationArtifact, + attempt: &serde_json::Value, + usage: &serde_json::Value, + cost_microusd: Option, + role_profile_provenance: Option<&RoleProfileProvenance>, ) -> Result<(), StoreError> { if !artifact.is_canonical() || artifact.verdict == VerificationStatus::NotRequested { return Err(StoreError::PatchArtifact( @@ -461,17 +674,11 @@ impl RuntimeStore { )); } let artifact_json = serde_json::to_string(artifact)?; - let attempt_json = serde_json::to_string(attempt)?; + let attempt_json = + serde_json::to_string(&with_role_profile_provenance(attempt, role_profile_provenance))?; let usage_json = serde_json::to_string(usage)?; let verdict = serde_json::to_value(artifact.verdict)?.as_str().unwrap_or("inconclusive").to_owned(); - let event_payload = serde_json::json!({ - "verification_id": artifact.id, - "patch_id": artifact.patch_id, - "verdict": verdict - }); - let event_json = serde_json::to_string(&event_payload)?; - let event_digest = Digest::blake3(event_json.as_bytes()); let mut connection = self.connection()?; let transaction = connection.transaction()?; let patch_exists: u64 = transaction.query_row( @@ -488,6 +695,21 @@ impl RuntimeStore { "verification references an unknown patch".to_owned(), )); } + let stored_provenance = require_change_role_profile_provenance( + &transaction, + &artifact.change_id, + role_profile_provenance, + )?; + let event_payload = with_role_profile_provenance( + &serde_json::json!({ + "verification_id": artifact.id, + "patch_id": artifact.patch_id, + "verdict": verdict + }), + stored_provenance.as_ref(), + ); + let event_json = serde_json::to_string(&event_payload)?; + let event_digest = Digest::blake3(event_json.as_bytes()); transaction.execute( "INSERT INTO verification_artifacts( verification_id, change_id, patch_id, verdict, artifact_json, created_unix_ms @@ -504,15 +726,19 @@ impl RuntimeStore { transaction.execute( "INSERT INTO change_attempts( change_id, patch_id, role, attempt_json, usage_json, cost_microusd, - created_unix_ms - ) VALUES(?1, ?2, 'verifier', ?3, ?4, ?5, ?6)", + created_unix_ms, role_profile_id, role_profile_revision, + role_profile_definition_digest + ) VALUES(?1, ?2, 'verifier', ?3, ?4, ?5, ?6, ?7, ?8, ?9)", params![ artifact.change_id.to_string(), artifact.patch_id.to_string(), attempt_json, usage_json, cost_microusd, - artifact.created_unix_ms + artifact.created_unix_ms, + role_profile_provenance.map(|value| value.profile_id.as_str()), + role_profile_provenance.map(|value| value.revision), + role_profile_provenance.map(|value| value.definition_digest.to_string()), ], )?; transaction.execute( @@ -542,6 +768,28 @@ impl RuntimeStore { usage: &serde_json::Value, cost_microusd: Option, created_unix_ms: u64, + ) -> Result<(), StoreError> { + self.record_patch_attempt_with_provenance( + change_id, + patch_id, + attempt, + usage, + cost_microusd, + created_unix_ms, + None, + ) + } + + #[allow(clippy::too_many_arguments)] + pub fn record_patch_attempt_with_provenance( + &self, + change_id: &ChangeId, + patch_id: PatchId, + attempt: &serde_json::Value, + usage: &serde_json::Value, + cost_microusd: Option, + created_unix_ms: u64, + role_profile_provenance: Option<&RoleProfileProvenance>, ) -> Result<(), StoreError> { let connection = self.connection()?; let patch_exists: u64 = connection.query_row( @@ -554,18 +802,26 @@ impl RuntimeStore { "patch attempt references an unknown patch".to_owned(), )); } + require_change_role_profile_provenance(&connection, change_id, role_profile_provenance)?; connection.execute( "INSERT INTO change_attempts( change_id, patch_id, role, attempt_json, usage_json, cost_microusd, - created_unix_ms - ) VALUES(?1, ?2, 'patcher', ?3, ?4, ?5, ?6)", + created_unix_ms, role_profile_id, role_profile_revision, + role_profile_definition_digest + ) VALUES(?1, ?2, 'patcher', ?3, ?4, ?5, ?6, ?7, ?8, ?9)", params![ change_id.to_string(), patch_id.to_string(), - serde_json::to_string(attempt)?, + serde_json::to_string(&with_role_profile_provenance( + attempt, + role_profile_provenance, + ))?, serde_json::to_string(usage)?, cost_microusd, - created_unix_ms + created_unix_ms, + role_profile_provenance.map(|value| value.profile_id.as_str()), + role_profile_provenance.map(|value| value.revision), + role_profile_provenance.map(|value| value.definition_digest.to_string()), ], )?; Ok(()) @@ -577,7 +833,8 @@ impl RuntimeStore { ) -> Result, StoreError> { let connection = self.connection()?; let mut statement = connection.prepare( - "SELECT role, patch_id, attempt_json, usage_json, cost_microusd, created_unix_ms + "SELECT role, patch_id, attempt_json, usage_json, cost_microusd, created_unix_ms, + role_profile_id, role_profile_revision, role_profile_definition_digest FROM change_attempts WHERE change_id=?1 ORDER BY attempt_id", )?; let rows = statement.query_map([change_id.to_string()], |row| { @@ -588,10 +845,23 @@ impl RuntimeStore { row.get::<_, String>(3)?, row.get::<_, Option>(4)?, row.get::<_, u64>(5)?, + row.get::<_, Option>(6)?, + row.get::<_, Option>(7)?, + row.get::<_, Option>(8)?, )) })?; rows.map(|row| { - let (role, patch_id, attempt, usage, cost_microusd, created_unix_ms) = row?; + let ( + role, + patch_id, + attempt, + usage, + cost_microusd, + created_unix_ms, + role_profile_id, + role_profile_revision, + role_profile_definition_digest, + ) = row?; Ok(ChangeAttemptRecord { role, patch_id: PatchId( @@ -601,6 +871,11 @@ impl RuntimeStore { usage: serde_json::from_str(&usage)?, cost_microusd, created_unix_ms, + role_profile_provenance: parse_role_profile_provenance(( + role_profile_id, + role_profile_revision, + role_profile_definition_digest, + ))?, }) }) .collect() @@ -655,7 +930,9 @@ impl RuntimeStore { .query_row( "SELECT c.request_json, c.request_digest, c.repository_id, c.source_snapshot_digest, c.state, p.artifact_json, - p.declared_output_json, c.repair_attempted, c.created_unix_ms + p.declared_output_json, c.repair_attempted, c.created_unix_ms, + c.role_profile_id, c.role_profile_revision, + c.role_profile_definition_digest FROM change_requests c JOIN patch_artifacts p ON p.change_id=c.change_id AND p.revision=c.latest_patch_revision @@ -672,6 +949,9 @@ impl RuntimeStore { row.get(6)?, row.get(7)?, row.get(8)?, + row.get(9)?, + row.get(10)?, + row.get(11)?, )) }, ) @@ -686,6 +966,9 @@ impl RuntimeStore { output, repair_attempted, created, + role_profile_id, + role_profile_revision, + role_profile_definition_digest, )) = current_row else { return Err(StoreError::ChangeConflict(record.change_id.to_string())); @@ -702,6 +985,11 @@ impl RuntimeStore { declared_output: serde_json::from_str(&output)?, repair_attempted, created_unix_ms: created, + role_profile_provenance: parse_role_profile_provenance(( + role_profile_id, + role_profile_revision, + role_profile_definition_digest, + ))?, }; let verification_json: Option = transaction .query_row( @@ -758,11 +1046,15 @@ impl RuntimeStore { "UPDATE change_requests SET state='applying', updated_unix_ms=?2 WHERE change_id=?1", params![record.change_id.to_string(), record.created_unix_ms], )?; - let event_json = serde_json::to_string(&serde_json::json!({ - "apply_id": record.id, - "patch_id": record.patch_id, - "state": "applying" - }))?; + let provenance = change_role_profile_provenance(&transaction, &record.change_id)?; + let event_json = serde_json::to_string(&with_role_profile_provenance( + &serde_json::json!({ + "apply_id": record.id, + "patch_id": record.patch_id, + "state": "applying" + }), + provenance.as_ref(), + ))?; transaction.execute( "INSERT INTO change_events( change_id, event_type, payload_digest, payload_json, created_unix_ms @@ -795,6 +1087,9 @@ impl RuntimeStore { [apply_id.to_string()], |row| row.get(0), )?; + let parsed_change_id = ChangeId::parse(&change_id) + .map_err(|error| StoreError::PatchArtifact(error.to_owned()))?; + let provenance = change_role_profile_provenance(&transaction, &parsed_change_id)?; transaction.execute( "UPDATE change_applies SET status=?2, post_snapshot_digest=?3, completed_unix_ms=?4 @@ -816,11 +1111,14 @@ impl RuntimeStore { "UPDATE change_requests SET state=?2, updated_unix_ms=?3 WHERE change_id=?1", params![change_id, change_state, completed_unix_ms], )?; - let event_json = serde_json::to_string(&serde_json::json!({ - "apply_id": apply_id, - "state": apply_status_name(status), - "post_snapshot": post_snapshot - }))?; + let event_json = serde_json::to_string(&with_role_profile_provenance( + &serde_json::json!({ + "apply_id": apply_id, + "state": apply_status_name(status), + "post_snapshot": post_snapshot + }), + provenance.as_ref(), + ))?; transaction.execute( "INSERT INTO change_events( change_id, event_type, payload_digest, payload_json, created_unix_ms @@ -974,14 +1272,51 @@ fn decode_apply_record(row: &rusqlite::Row<'_>) -> rusqlite::Result RoleProfileProvenance { + let definition = RoleProfileDefinition::new(RoleProfileDefinitionInput { + profile_id: RoleProfileId::new("change.implementer").unwrap(), + role: CodexRole::Implementer, + host: CodexHost::Codex, + model: "offline-change-worker".to_owned(), + reasoning: needle_core::ReasoningLevel::Medium, + service_tier: ServiceTier::Default, + timeout_seconds: 30, + budget: RoleProfileBudget { + max_turns: 2, + max_output_tokens: 1200, + max_cost_microusd: 1000, + }, + prompt_profile_digest: Digest::blake3(b"change-prompt"), + output_contract_digest: Digest::blake3(b"change-output"), + tool_policy: ToolPolicy::IsolatedWrite, + command_policy: CommandPolicy::Denied, + filesystem_policy: FilesystemPolicy::DisposableCheckout, + network_policy: NetworkPolicy::Denied, + test_policy: TestPolicy::Disabled, + repair_policy: RepairPolicy::Once, + fallback_policy: FallbackPolicy::Native, + concurrency: 1, + route_assignments: Vec::new(), + }) + .unwrap(); + let revision = store.create_role_profile(definition).unwrap(); + let state = store.role_profile_state(&revision.profile_id).unwrap(); + store + .activate_role_profile(&revision.profile_id, revision.revision, state.state_digest) + .unwrap(); + RoleProfileProvenance::from_revision(&revision).unwrap() + } + #[test] fn failed_request_is_audited_before_any_patch_exists() { let suffix = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos(); @@ -989,6 +1324,7 @@ mod tests { .join(format!("needle-change-request-audit-{}-{suffix}.sqlite3", std::process::id())); let store = RuntimeStore::new(&path); store.initialize().unwrap(); + let provenance = active_profile_provenance(&store); let request = ChangeRequest { task: "Update the fixture.".to_owned(), acceptance_criteria: vec!["The fixture changes.".to_owned()], @@ -1004,15 +1340,17 @@ mod tests { let request_digest = request.digest(source); let change_id = ChangeId::from_digest(Digest::blake3(b"failed-change")); store - .record_change_request( + .record_change_request_with_provenance( &change_id, Digest::blake3(b"repository"), source, request_digest, &request, + Some(&provenance), ) .unwrap(); - store.record_change_failure(&change_id, "worker failed").unwrap(); + let sensitive_reason = "worker failed at C:\\private\\repo with raw prompt text"; + store.record_change_failure(&change_id, sensitive_reason).unwrap(); let connection = rusqlite::Connection::open(&path).unwrap(); let (state, revision): (String, u32) = connection @@ -1032,6 +1370,27 @@ mod tests { ) .unwrap(); assert_eq!(events, 2); + let payloads = connection + .prepare("SELECT payload_json FROM change_events WHERE change_id=?1 ORDER BY event_id") + .unwrap() + .query_map([change_id.to_string()], |row| row.get::<_, String>(0)) + .unwrap() + .collect::, _>>() + .unwrap(); + for payload in &payloads { + let payload: serde_json::Value = serde_json::from_str(payload).unwrap(); + assert_eq!( + payload.get("role_profile_provenance"), + Some(&serde_json::json!({ + "profile_id": provenance.profile_id, + "revision": provenance.revision, + "definition_digest": provenance.definition_digest, + })) + ); + } + assert!(!payloads[1].contains(sensitive_reason)); + assert!(!payloads[1].contains("C:\\\\private")); + assert!(payloads[1].contains("reason_digest")); drop(connection); drop(store); fs::remove_file(path).unwrap(); diff --git a/crates/needle-runtime/src/store/role_profiles/tests.rs b/crates/needle-runtime/src/store/role_profiles/tests.rs index 2c47e50..06c59e3 100644 --- a/crates/needle-runtime/src/store/role_profiles/tests.rs +++ b/crates/needle-runtime/src/store/role_profiles/tests.rs @@ -1,8 +1,9 @@ use super::*; use needle_core::{ - CodexHost, CodexRole, CommandPolicy, FallbackPolicy, FilesystemPolicy, NetworkPolicy, - RepairPolicy, RoleProfileBudget, RoleProfileDefinitionInput, RoleProfileId, ServiceTier, - TestPolicy, ToolPolicy, + CacheLookup, CodexHost, CodexRole, CommandPolicy, FallbackPolicy, FilesystemPolicy, + NeedCacheEntry, NeedCacheIdentity, NeedKey, NeedResult, NetworkPolicy, RepairPolicy, + RoleProfileBudget, RoleProfileDefinitionInput, RoleProfileId, RoleProfileProvenance, + ServiceTier, TestPolicy, ToolPolicy, WorkerOutcome, }; use std::path::PathBuf; use std::time::{SystemTime, UNIX_EPOCH}; @@ -14,14 +15,23 @@ fn temporary_store() -> (PathBuf, RuntimeStore) { } fn definition(id: &str) -> RoleProfileDefinition { + definition_with_execution(id, "gpt-5", 120, RepairPolicy::None) +} + +fn definition_with_execution( + id: &str, + model: &str, + timeout_seconds: u64, + repair_policy: RepairPolicy, +) -> RoleProfileDefinition { RoleProfileDefinition::new(RoleProfileDefinitionInput { profile_id: RoleProfileId::new(id).unwrap(), role: CodexRole::Explorer, host: CodexHost::Codex, - model: "gpt-5".to_owned(), + model: model.to_owned(), reasoning: needle_core::ReasoningLevel::Medium, service_tier: ServiceTier::Default, - timeout_seconds: 120, + timeout_seconds, budget: RoleProfileBudget { max_turns: 2, max_output_tokens: 1200, @@ -34,7 +44,7 @@ fn definition(id: &str) -> RoleProfileDefinition { filesystem_policy: FilesystemPolicy::ReadOnlyCheckout, network_policy: NetworkPolicy::Denied, test_policy: TestPolicy::Disabled, - repair_policy: RepairPolicy::None, + repair_policy, fallback_policy: FallbackPolicy::Native, concurrency: 1, route_assignments: vec![], @@ -80,7 +90,7 @@ fn migration_and_revision_lifecycle_are_atomic_and_immutable() { .unwrap() .collect::>() .unwrap(); - assert_eq!(versions, (1..=14).collect::>()); + assert_eq!(versions, (1..=15).collect::>()); for name in ["role_profiles", "role_profile_revisions", "role_profile_state", "role_profile_audit"] { @@ -216,6 +226,82 @@ fn migration_and_revision_lifecycle_are_atomic_and_immutable() { let _ = std::fs::remove_file(path); } +#[test] +fn session_binding_is_idempotent_conflict_checked_and_historical() { + let (path, store) = temporary_store(); + let first = store.create_role_profile(definition("explorer.session")).unwrap(); + let profile_id = first.profile_id.clone(); + let state = store.role_profile_state(&profile_id).unwrap(); + store.activate_role_profile(&profile_id, 1, state.state_digest).unwrap(); + + let prompt = Digest::blake3(b"session-prompt"); + store + .record_session_start_profiled("session-a", prompt, Some("main"), None, &profile_id) + .unwrap(); + store + .record_session_start_profiled("session-a", prompt, Some("main"), None, &profile_id) + .unwrap(); + let frozen = store.worker_config_for_session("session-a", "codex-a").unwrap(); + assert_eq!(frozen.model, "gpt-5"); + assert_eq!(frozen.timeout_seconds, 120); + assert_eq!( + frozen.evidence_failure_policy, + needle_core::EvidenceFailurePolicy::DiscardInvalidFact + ); + + let state = store.role_profile_state(&profile_id).unwrap(); + let second = store + .revise_role_profile( + &profile_id, + state.state_digest, + definition_with_execution("explorer.session", "gpt-5-mini", 240, RepairPolicy::Once), + ) + .unwrap(); + let state = store.role_profile_state(&profile_id).unwrap(); + store.activate_role_profile(&profile_id, second.revision, state.state_digest).unwrap(); + + assert!(matches!( + store.record_session_start_profiled("session-a", prompt, Some("main"), None, &profile_id,), + Err(StoreError::RoleProfileConflict(_)) + )); + let still_frozen = store.worker_config_for_session("session-a", "codex-a").unwrap(); + assert_eq!(still_frozen.model, "gpt-5"); + assert_eq!(still_frozen.timeout_seconds, 120); + + store + .record_session_start_profiled("session-b", prompt, Some("main"), None, &profile_id) + .unwrap(); + let current = store.worker_config_for_session("session-b", "codex-a").unwrap(); + assert_eq!(current.model, "gpt-5-mini"); + assert_eq!(current.timeout_seconds, 240); + assert_eq!(current.evidence_failure_policy, needle_core::EvidenceFailurePolicy::RepairOnce); + + let connection = Connection::open(&path).unwrap(); + assert!( + connection + .execute( + "UPDATE sessions + SET role_profile_revision=?2, role_profile_definition_digest=?3 + WHERE session_id=?1", + rusqlite::params![ + "session-a", + second.revision, + second.definition.definition_digest.to_string(), + ], + ) + .is_err() + ); + drop(connection); + + store.record_legacy_session_start("legacy", prompt, None, None).unwrap(); + assert!(matches!( + store.record_session_start_profiled("legacy", prompt, None, None, &profile_id), + Err(StoreError::RoleProfileConflict(_)) + )); + assert!(store.session("legacy").unwrap().unwrap().role_profile_provenance.is_none()); + let _ = std::fs::remove_file(path); +} + #[test] fn bounded_revision_listing_reads_only_the_latest_ordered_window() { let (path, store) = temporary_store(); @@ -249,6 +335,106 @@ fn bounded_revision_listing_reads_only_the_latest_ordered_window() { let _ = std::fs::remove_file(path); } +#[test] +fn cache_identity_separates_revisions_and_rejects_unknown_or_mismatched_provenance() { + let (path, store) = temporary_store(); + let first = store.create_role_profile(definition("explorer.cache")).unwrap(); + let profile_id = first.profile_id.clone(); + let state = store.role_profile_state(&profile_id).unwrap(); + let second = store + .revise_role_profile( + &profile_id, + state.state_digest, + definition_with_execution("explorer.cache", "gpt-5-mini", 180, RepairPolicy::Once), + ) + .unwrap(); + let first_provenance = RoleProfileProvenance::from_revision(&first).unwrap(); + let second_provenance = RoleProfileProvenance::from_revision(&second).unwrap(); + let identity = |provenance: Option| NeedCacheIdentity { + repository_id: Digest::blake3(b"repository"), + source_snapshot_digest: Digest::blake3(b"source"), + prompt_profile_digest: Digest::blake3(b"prompt"), + route_definition_digest: Digest::blake3(b"route"), + preset_definition_digest: Digest::blake3(b"preset"), + need_key: NeedKey::new("trace.state-flow").unwrap(), + normalized_request_digest: Digest::blake3(b"request"), + worker_configuration_digest: Digest::blake3(b"worker"), + output_schema_digest: Digest::blake3(b"schema"), + role_profile_provenance: provenance, + }; + let first_identity = identity(Some(first_provenance.clone())); + let second_identity = identity(Some(second_provenance.clone())); + assert_ne!(first_identity.digest(), second_identity.digest()); + assert_ne!(first_identity.logical_digest(), second_identity.logical_digest()); + + let unknown = identity(Some( + RoleProfileProvenance::new( + RoleProfileId::new("ghost").unwrap(), + 1, + Digest::blake3(b"ghost"), + ) + .unwrap(), + )); + assert!(matches!( + store.cache_lookup(&unknown).unwrap(), + CacheLookup::Bypass(reason) if reason == "role-profile-provenance-invalid" + )); + assert!(matches!( + store.cache_lookup(&identity(None)).unwrap(), + CacheLookup::Bypass(reason) if reason == "role-profile-provenance-unknown" + )); + + let result = NeedResult { + complete: true, + summary: "bounded".to_owned(), + claims: Vec::new(), + evidence: Vec::new(), + suggested_reads: Vec::new(), + suggested_commands: Vec::new(), + uncertainty: Vec::new(), + }; + let outcome = |provenance: RoleProfileProvenance| WorkerOutcome { + result: result.clone(), + artifact_result: None, + semantic_artifact_result: None, + worker_model: "gpt-5".to_owned(), + worker_reasoning: "medium".to_owned(), + codex_version: "test".to_owned(), + input_tokens: Some(1), + cached_input_tokens: Some(0), + output_tokens: Some(1), + duration_ms: 1, + process_status: "success".to_owned(), + logical_worker_spawns: 1, + worker_turns: 1, + repair_performed: false, + discarded_facts: 0, + worker_session_id: None, + session_cleanup_success: Some(true), + role_profile_provenance: Some(provenance), + }; + let mismatched = NeedCacheEntry { + identity: first_identity.clone(), + result: result.clone(), + worker_outcome: outcome(second_provenance), + created_unix_ms: 1, + hit_count: 0, + }; + assert!(matches!(store.publish(&mismatched), Err(StoreError::ArtifactIdentity(_)))); + + let matching = NeedCacheEntry { + identity: first_identity.clone(), + result: result.clone(), + worker_outcome: outcome(first_provenance), + created_unix_ms: 1, + hit_count: 0, + }; + store.publish(&matching).unwrap(); + assert!(matches!(store.cache_lookup(&first_identity).unwrap(), CacheLookup::Hit(_))); + assert!(matches!(store.cache_lookup(&second_identity).unwrap(), CacheLookup::Miss)); + let _ = std::fs::remove_file(path); +} + #[test] fn revising_an_active_profile_preserves_prior_active_audit_pointer() { let (path, store) = temporary_store(); @@ -364,7 +550,7 @@ fn inconsistent_state_pointers_fail_closed_without_historical_fallback() { } #[test] -fn v14_upgrades_a_valid_v13_database_and_rejects_checksum_drift() { +fn v15_upgrades_a_valid_v14_database_without_attributing_legacy_rows() { let (path, store) = temporary_store(); let connection = Connection::open(&path).unwrap(); let migrations = [ @@ -381,6 +567,7 @@ fn v14_upgrades_a_valid_v13_database_and_rejects_checksum_drift() { (11, super::super::MIGRATION_V11), (12, super::super::MIGRATION_V12), (13, super::super::MIGRATION_V13), + (14, super::super::MIGRATION_V14), ]; for (version, migration) in migrations { connection.execute_batch(migration).unwrap(); @@ -391,15 +578,36 @@ fn v14_upgrades_a_valid_v13_database_and_rejects_checksum_drift() { ) .unwrap(); } + connection + .execute( + "INSERT INTO sessions( + session_id, prompt_profile_digest, route_set_digest, updated_unix_ms + ) VALUES('legacy-session', ?1, ?2, 0)", + rusqlite::params![ + Digest::blake3(b"prompt").to_string(), + Digest::blake3(b"routes").to_string() + ], + ) + .unwrap(); drop(connection); store.initialize().unwrap(); let connection = Connection::open(&path).unwrap(); let version: u32 = connection .query_row("SELECT MAX(version) FROM schema_migrations", [], |row| row.get(0)) .unwrap(); - assert_eq!(version, 14); + assert_eq!(version, 15); + let legacy: (Option, Option, Option) = connection + .query_row( + "SELECT role_profile_id, role_profile_revision, + role_profile_definition_digest + FROM sessions WHERE session_id='legacy-session'", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + ) + .unwrap(); + assert_eq!(legacy, (None, None, None)); connection - .execute("UPDATE schema_migrations SET checksum='b3:invalid' WHERE version=14", []) + .execute("UPDATE schema_migrations SET checksum='b3:invalid' WHERE version=15", []) .unwrap(); drop(connection); let drifted = RuntimeStore::new(&path); diff --git a/docs/CONFIGURATION.md b/docs/CONFIGURATION.md index 0a853f4..76e865c 100644 --- a/docs/CONFIGURATION.md +++ b/docs/CONFIGURATION.md @@ -44,8 +44,8 @@ Credentials are neither imported nor exported. Worker execution and orchestration currently use Codex only. The web control plane can edit Codex model policy, runtime bounds, and canonical named role -profiles; role-profile changes remain configuration-only and do not bind a -worker or session. +profiles. The HTTP/editor changes configuration only; the product hook and MCP +entry points bind an explicitly selected active revision to each new session. The planned sequence is Codex-first role configuration and lifecycle orchestration, followed by configuration-only interoperability for Claude Code @@ -90,13 +90,22 @@ projects to the existing `WorkerProfile` representation with `None`, while distinct from the role-profile definition/revision digest. Projection is explicit and does not read or modify `ModelPolicy`. -The domain and SQLite migration remain an offline configuration boundary. The -local HTTP/editor exposes bounded, authenticated routes for list/detail/history, -audit, request-time preflight, draft CAS, and explicit activation/deactivation. -Preflight is recomputed at request time and is never persisted. No endpoint -launches a worker, binds a session, executes a lifecycle, or configures a -non-Codex host. Role profiles never carry credentials, host paths, raw -transcripts, or network access. +Role-profile definitions and revisions remain local SQLite state. New +production sessions bind an explicitly selected active revision; the binding +stores only profile ID, revision, and definition digest. Historical sessions +reload that exact revision even after a later activation. Legacy rows remain +unknown and cannot be reused for profile-dependent cache or worker execution. +The local HTTP/editor exposes bounded, authenticated routes for +list/detail/history, audit, request-time preflight, draft CAS, and explicit +activation/deactivation. Preflight is recomputed at request time and is never +persisted. These endpoints do not launch workers, bind sessions, execute a +lifecycle, or configure a non-Codex host. Role profiles never carry +credentials, host paths, raw transcripts, or network access. + +The product hook selects a profile with `NEEDLE_ROLE_PROFILE_ID`. If that +variable is missing or invalid, the hook remains fail-open but records no +session row, so the runtime cannot silently attribute the session to a current +active revision. MCP uses the required `--role-profile ` selector. ## Export and import diff --git a/docs/DEVELOPER_SETUP.md b/docs/DEVELOPER_SETUP.md index 3948f65..a166db7 100644 --- a/docs/DEVELOPER_SETUP.md +++ b/docs/DEVELOPER_SETUP.md @@ -101,11 +101,14 @@ cargo run --locked -p needle-app -- mcp serve \ --data-dir \ --repository \ --main-model \ + --role-profile \ --cache-only ``` -`--cache-only` disables worker fallback and change tools. Omit it only in a +The `--cache-only` option disables worker fallback and change tools. Omit it only in a deliberately configured development profile. See [MCP transport](MCP_TRANSPORT.md). +The role-profile selector is mandatory and freezes the selected active revision +for the MCP session. ## Validate plugins diff --git a/docs/MCP_TRANSPORT.md b/docs/MCP_TRANSPORT.md index 1306b39..02b139e 100644 --- a/docs/MCP_TRANSPORT.md +++ b/docs/MCP_TRANSPORT.md @@ -14,11 +14,13 @@ needle mcp serve \ --data-dir \ --repository \ --main-model \ + --role-profile \ --cache-only ``` The connection freezes repository, profile, enabled routes, semantic-definition -digest, main model, and worker policy. `--cache-only` disables worker fallback +digest, main model, selected role-profile revision, and worker policy. The +role-profile selector is mandatory. `--cache-only` disables worker fallback and change tools. The stdio server accepts bounded JSON-RPC lines and requires `initialize`