From 3be5f480536c82d8881ef850aaa1525bcf0fa0ca Mon Sep 17 00:00:00 2001 From: Dotta Date: Thu, 3 Sep 2026 01:49:00 -0500 Subject: [PATCH] fix(runner): preserve OpenCode terminal results --- .../runner-core/src/provider_backend.rs | 372 ++++++++++++++++-- 1 file changed, 339 insertions(+), 33 deletions(-) diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs b/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs index 60fb12b77b..93cfb1d0e9 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs @@ -236,13 +236,142 @@ fn semantic_result_event( } } +fn validate_opencode_run_result( + state: &CodexProviderState, + params: &Value, +) -> Result<(Value, String, String), DurableRunnerError> { + if state.config.provider != "opencode" { + return Err(DurableRunnerError::invalid( + "paperclip/runResult is reserved for the verified OpenCode provider", + )); + } + if state.lifecycle != "turn_active" || state.active_provider_turn_id.is_none() { + return Err(DurableRunnerError::invalid( + "OpenCode emitted paperclip/runResult without an active provider turn", + )); + } + if params.get("threadId").and_then(Value::as_str) != state.thread_id.as_deref() + || params.get("turnId").and_then(Value::as_str) != state.active_provider_turn_id.as_deref() + { + return Err(DurableRunnerError::invalid( + "OpenCode paperclip/runResult is not bound to the active provider turn", + )); + } + let result = params.get("result").cloned().ok_or_else(|| { + DurableRunnerError::invalid("OpenCode paperclip/runResult omitted its result") + })?; + let schema: Value = serde_json::from_str(include_str!( + "../../../../protocol/schemas/result.schema.json" + )) + .map_err(|_| DurableRunnerError::invalid("embedded Paperclip result schema is invalid"))?; + let validator = jsonschema::validator_for(&schema).map_err(|_| { + DurableRunnerError::invalid("embedded Paperclip result schema cannot compile") + })?; + if !validator.is_valid(&result) { + return Err(DurableRunnerError::invalid( + "OpenCode paperclip/runResult failed the Paperclip result schema", + )); + } + let contract = state.completion_contract.as_ref().ok_or_else(|| { + DurableRunnerError::invalid("OpenCode paperclip/runResult has no bound completion contract") + })?; + let claim = result.get("completionClaim").ok_or_else(|| { + DurableRunnerError::invalid("OpenCode paperclip/runResult omitted its completion claim") + })?; + if claim.get("contractRevision").and_then(Value::as_str) != Some(contract.revision.as_str()) { + return Err(DurableRunnerError::invalid( + "OpenCode paperclip/runResult changed its completion contract revision", + )); + } + let criteria = claim + .get("criteria") + .and_then(Value::as_array) + .ok_or_else(|| { + DurableRunnerError::invalid( + "OpenCode paperclip/runResult omitted its completion criteria", + ) + })?; + let mut reported_criterion_ids = HashSet::new(); + for criterion in criteria { + let criterion_id = criterion + .get("criterionId") + .and_then(Value::as_str) + .ok_or_else(|| { + DurableRunnerError::invalid( + "OpenCode paperclip/runResult has an invalid completion criterion", + ) + })?; + if !reported_criterion_ids.insert(criterion_id) { + return Err(DurableRunnerError::invalid( + "OpenCode paperclip/runResult repeated a completion criterion", + )); + } + } + if reported_criterion_ids.len() != contract.criterion_ids.len() + || !contract + .criterion_ids + .iter() + .all(|criterion_id| reported_criterion_ids.contains(criterion_id.as_str())) + { + return Err(DurableRunnerError::invalid( + "OpenCode paperclip/runResult changed its bound completion criteria", + )); + } + let disposition = result + .get("reportedWorkDisposition") + .and_then(Value::as_str) + .expect("the validated result schema requires a disposition") + .to_owned(); + let fingerprint = semantic_value_digest(&result); + Ok((result, fingerprint, disposition)) +} + +fn normalize_provider_notification( + state: &mut CodexProviderState, + method: &str, + params: &Value, +) -> Result, DurableRunnerError> { + if method != "paperclip/runResult" { + return Ok(normalize_codex_notification(method, params) + .into_iter() + .map(|event| relabel_provider_event(event, &state.config.provider)) + .collect()); + } + let (result, fingerprint, disposition) = validate_opencode_run_result(state, params)?; + match ( + state.active_provider_result_fingerprint.as_deref(), + state.active_provider_result_disposition.as_deref(), + ) { + (None, None) => { + state.active_provider_result_fingerprint = Some(fingerprint); + state.active_provider_result_disposition = Some(disposition); + Ok(vec![NormalizedProviderEvent { + event_type: "run.result.proposed".to_owned(), + priority: EventPriority::P0, + payload: result, + }]) + } + (Some(existing_fingerprint), Some(existing_disposition)) + if existing_fingerprint == fingerprint && existing_disposition == disposition => + { + Ok(Vec::new()) + } + _ => Err(DurableRunnerError::invalid( + "OpenCode emitted conflicting paperclip/runResult notifications for one turn", + )), + } +} + fn terminal_events(state: &CodexProviderState, event_type: &str) -> Vec { let Some(contract) = state.completion_contract.as_ref() else { return Vec::new(); }; let succeeded = event_type == "turn.completed"; let cancelled = matches!(event_type, "turn.cancelled" | "turn.interrupted"); - let disposition = if succeeded { "done" } else { "needs_review" }; + let disposition = state + .active_provider_result_disposition + .as_deref() + .unwrap_or(if succeeded { "done" } else { "needs_review" }); let provider = state.config.provider.as_str(); let provider_name = if provider == "opencode" { "OpenCode" @@ -308,18 +437,20 @@ fn terminal_events(state: &CodexProviderState, event_type: &str) -> Vec, + #[serde(default)] + active_provider_result_fingerprint: Option, + #[serde(default)] + active_provider_result_disposition: Option, last_agent_message: Option, #[serde(default)] pending_events: VecDeque, @@ -466,6 +601,8 @@ impl CodexProviderState { receipt_limit_interrupt_accepted: false, receipt_limit_interrupt_attempts: 0, receipt_limit_interrupt_deadline_unix_ms: None, + active_provider_result_fingerprint: None, + active_provider_result_disposition: None, last_agent_message: None, pending_events: VecDeque::new(), queued_events: VecDeque::new(), @@ -566,6 +703,25 @@ impl CodexProviderState { || self .receipt_limit_interrupt_deadline_unix_ms .is_some_and(|deadline| deadline == 0 || !self.receipt_limit_interrupt_pending) + || self.active_provider_result_fingerprint.is_some() + != self.active_provider_result_disposition.is_some() + || self + .active_provider_result_fingerprint + .as_ref() + .is_some_and(|fingerprint| { + fingerprint.len() != 71 + || !fingerprint.starts_with("sha256:") + || !fingerprint[7..] + .chars() + .all(|character| character.is_ascii_hexdigit()) + }) + || self + .active_provider_result_disposition + .as_deref() + .is_some_and(|disposition| { + !matches!(disposition, "done" | "blocked" | "needs_review" | "yielded") + || self.config.provider != "opencode" + }) || (matches!( self.lifecycle.as_str(), "prepared" | "session_open" | "closed" @@ -776,6 +932,8 @@ impl CodexProviderState { self.completed_turn_process_generation = None; self.completed_provider_turn_id = None; self.ambiguous_turn_start_pending = false; + self.active_provider_result_fingerprint = None; + self.active_provider_result_disposition = None; self.last_agent_message = None; } self.lifecycle = if self.active_provider_turn_id.is_some() { @@ -1056,6 +1214,8 @@ impl CodexCommandExecutor { state.receipt_limit_interrupt_accepted = false; state.receipt_limit_interrupt_attempts = 0; state.receipt_limit_interrupt_deadline_unix_ms = None; + state.active_provider_result_fingerprint = None; + state.active_provider_result_disposition = None; state.last_agent_message = None; state.lifecycle = "closed".to_owned(); let _ = state.push_terminal_event(NormalizedProviderEvent { @@ -1108,6 +1268,8 @@ impl CodexCommandExecutor { state.receipt_limit_interrupt_accepted = false; state.receipt_limit_interrupt_attempts = 0; state.receipt_limit_interrupt_deadline_unix_ms = None; + state.active_provider_result_fingerprint = None; + state.active_provider_result_disposition = None; state.last_agent_message = None; state.lifecycle = "closed".to_owned(); // Closing the provider is the safety boundary. Preserve that @@ -1510,6 +1672,8 @@ impl CodexCommandExecutor { next_state.receipt_limit_interrupt_accepted = false; next_state.receipt_limit_interrupt_attempts = 0; next_state.receipt_limit_interrupt_deadline_unix_ms = None; + next_state.active_provider_result_fingerprint = None; + next_state.active_provider_result_disposition = None; next_state.last_agent_message = None; if let Some(provider) = self.provider.as_mut() { provider.shutdown().map_err(|error| { @@ -1630,6 +1794,8 @@ impl CodexCommandExecutor { state.receipt_limit_interrupt_accepted = false; state.receipt_limit_interrupt_attempts = 0; state.receipt_limit_interrupt_deadline_unix_ms = None; + state.active_provider_result_fingerprint = None; + state.active_provider_result_disposition = None; state.last_agent_message = None; state.lifecycle = "closed".to_owned(); let provider_label = state.config.provider.clone(); @@ -1854,6 +2020,8 @@ impl CodexCommandExecutor { state.completed_turn_authoritative = false; state.completed_turn_process_generation = None; state.completed_provider_turn_id = None; + state.active_provider_result_fingerprint = None; + state.active_provider_result_disposition = None; state.last_agent_message = None; } self.save_state()?; @@ -1890,6 +2058,8 @@ impl CodexCommandExecutor { state.receipt_limit_interrupt_accepted = false; state.receipt_limit_interrupt_attempts = 0; state.receipt_limit_interrupt_deadline_unix_ms = None; + state.active_provider_result_fingerprint = None; + state.active_provider_result_disposition = None; state.last_agent_message = None; state.lifecycle = "turn_active".to_owned(); let provider_label = state.config.provider.clone(); @@ -2400,6 +2570,8 @@ impl CodexCommandExecutor { state.receipt_limit_interrupt_accepted = false; state.receipt_limit_interrupt_attempts = 0; state.receipt_limit_interrupt_deadline_unix_ms = None; + state.active_provider_result_fingerprint = None; + state.active_provider_result_disposition = None; state.lifecycle = "closed".to_owned(); let thread_id = state.thread_id.clone(); let provider_name = state.config.provider.clone(); @@ -2498,29 +2670,7 @@ impl CodexCommandExecutor { } else { None }; - let provider_name = self - .state - .as_ref() - .map(|state| state.config.provider.clone()) - .unwrap_or_else(|| "codex".to_owned()); - let normalized = normalize_codex_notification(&method, ¶ms) - .into_iter() - .map(|event| relabel_provider_event(event, &provider_name)) - .collect::>(); - let normalized_event_count = normalized.len(); - let terminal_event_type = normalized - .iter() - .find(|event| event.event_type.starts_with("turn.")) - .map(|event| event.event_type.clone()) - .filter(|event_type| { - matches!( - event_type.as_str(), - "turn.completed" - | "turn.failed" - | "turn.cancelled" - | "turn.interrupted" - ) - }); + let terminal_event_type = normalized_terminal_type.map(str::to_owned); let identity = self.event_identity.clone(); let state = self .state @@ -2544,6 +2694,8 @@ impl CodexCommandExecutor { })?; state.reconcile_active_provider_turn(Some(provider_turn_id)); } + let normalized = normalize_provider_notification(state, &method, ¶ms)?; + let normalized_event_count = normalized.len(); if terminal_event_type.is_some() { state.settle_active_provider_turn_identity()?; let settled = state @@ -2829,6 +2981,158 @@ impl CommandExecutor for CodexCommandExecutor { mod tests { use super::*; + fn opencode_result_state() -> CodexProviderState { + let mut state = CodexProviderState::new( + CodexProviderConfig { + provider: "opencode".to_owned(), + driver: "opencode_server".to_owned(), + provider_version: "1.18.17".to_owned(), + command: PathBuf::from("node"), + args: Vec::new(), + cwd: std::env::current_dir() + .unwrap() + .to_string_lossy() + .into_owned(), + model: Some("openrouter/model".to_owned()), + provider_session_id: None, + instructions: String::new(), + approval_policy: "never".to_owned(), + }, + Some(CompletionContractBinding { + revision: "revision-1".to_owned(), + criterion_ids: vec!["criterion-1".to_owned()], + }), + ProviderToolBridge::default(), + ); + state.thread_id = Some("thread-1".to_owned()); + state.active_provider_turn_id = Some("turn-1".to_owned()); + state.lifecycle = "turn_active".to_owned(); + state + } + + fn valid_opencode_result() -> Value { + json!({ + "schema": "paperclip.run_result.v1", + "reportedWorkDisposition": "done", + "summary": "Finished the requested work.", + "completionClaim": { + "contractRevision": "revision-1", + "objectiveSatisfied": true, + "criteria": [{ + "criterionId": "criterion-1", + "status": "satisfied", + "evidenceRefs": ["provider:opencode:agent-message"], + }], + "remainingWork": [], + }, + "evidence": [{"ref": "provider:opencode:agent-message"}], + "verification": [], + "attentionRequests": [], + "artifacts": [], + }) + } + + #[test] + fn preserves_one_verified_opencode_result_before_its_terminal() { + let mut state = opencode_result_state(); + let params = json!({ + "threadId": "thread-1", + "turnId": "turn-1", + "itemId": "semantic-result", + "result": valid_opencode_result(), + }); + + let result_events = + normalize_provider_notification(&mut state, "paperclip/runResult", ¶ms).unwrap(); + let replay_events = + normalize_provider_notification(&mut state, "paperclip/runResult", ¶ms).unwrap(); + let terminal = terminal_events(&state, "turn.completed"); + + assert_eq!(result_events.len(), 1); + assert_eq!(result_events[0].event_type, "run.result.proposed"); + assert_eq!(result_events[0].priority, EventPriority::P0); + assert!(replay_events.is_empty()); + assert_eq!(terminal.len(), 1); + assert_eq!(terminal[0].event_type, "run.terminal"); + assert_eq!(terminal[0].payload["reportedWorkDisposition"], "done"); + assert!(state.validate().is_ok()); + } + + #[test] + fn rejects_unbound_conflicting_or_spoofed_opencode_results() { + let params = |result: Value| { + json!({ + "threadId": "thread-1", + "turnId": "turn-1", + "itemId": "semantic-result", + "result": result, + }) + }; + + let mut wrong_revision = opencode_result_state(); + let mut result = valid_opencode_result(); + result["completionClaim"]["contractRevision"] = json!("revision-2"); + assert!(normalize_provider_notification( + &mut wrong_revision, + "paperclip/runResult", + ¶ms(result), + ) + .unwrap_err() + .to_string() + .contains("contract revision")); + + let mut malformed = opencode_result_state(); + assert!(normalize_provider_notification( + &mut malformed, + "paperclip/runResult", + ¶ms(json!({"schema": "paperclip.run_result.v1"})), + ) + .unwrap_err() + .to_string() + .contains("failed the Paperclip result schema")); + + let mut wrong_criteria = opencode_result_state(); + let mut result = valid_opencode_result(); + result["completionClaim"]["criteria"][0]["criterionId"] = json!("criterion-2"); + assert!(normalize_provider_notification( + &mut wrong_criteria, + "paperclip/runResult", + ¶ms(result), + ) + .unwrap_err() + .to_string() + .contains("bound completion criteria")); + + let mut conflicting = opencode_result_state(); + normalize_provider_notification( + &mut conflicting, + "paperclip/runResult", + ¶ms(valid_opencode_result()), + ) + .unwrap(); + let mut result = valid_opencode_result(); + result["summary"] = json!("A conflicting second result."); + assert!(normalize_provider_notification( + &mut conflicting, + "paperclip/runResult", + ¶ms(result), + ) + .unwrap_err() + .to_string() + .contains("conflicting")); + + let mut spoofed = opencode_result_state(); + spoofed.config.provider = "codex".to_owned(); + assert!(normalize_provider_notification( + &mut spoofed, + "paperclip/runResult", + ¶ms(valid_opencode_result()), + ) + .unwrap_err() + .to_string() + .contains("reserved for the verified OpenCode provider")); + } + #[test] fn opencode_terminal_fallback_uses_its_actual_provider_identity() { let mut state = CodexProviderState::new( @@ -2907,6 +3211,8 @@ mod tests { receipt_limit_interrupt_accepted: false, receipt_limit_interrupt_attempts: 0, receipt_limit_interrupt_deadline_unix_ms: None, + active_provider_result_fingerprint: None, + active_provider_result_disposition: None, last_agent_message: None, pending_events: VecDeque::new(), queued_events: VecDeque::new(),