fix(runner): preserve OpenCode terminal results

This commit is contained in:
Dotta 2026-09-03 01:49:00 -05:00
parent 3d8720c16f
commit 3be5f48053
1 changed files with 339 additions and 33 deletions

View File

@ -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<Vec<NormalizedProviderEvent>, 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<NormalizedProviderEvent> {
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<Normaliz
"runTerminalState": if succeeded { "succeeded" } else if cancelled { "cancelled" } else { "failed" },
"reportedWorkDisposition": disposition,
});
vec![
NormalizedProviderEvent {
let mut events = Vec::new();
if state.active_provider_result_fingerprint.is_none() {
events.push(NormalizedProviderEvent {
event_type: "run.result.proposed".to_owned(),
priority: EventPriority::P0,
payload: result,
},
NormalizedProviderEvent {
event_type: "run.terminal".to_owned(),
priority: EventPriority::P0,
payload: terminal,
},
]
});
}
events.push(NormalizedProviderEvent {
event_type: "run.terminal".to_owned(),
priority: EventPriority::P0,
payload: terminal,
});
events
}
fn relabel_provider_event(
@ -396,6 +527,10 @@ struct CodexProviderState {
receipt_limit_interrupt_attempts: u8,
#[serde(default)]
receipt_limit_interrupt_deadline_unix_ms: Option<u64>,
#[serde(default)]
active_provider_result_fingerprint: Option<String>,
#[serde(default)]
active_provider_result_disposition: Option<String>,
last_agent_message: Option<String>,
#[serde(default)]
pending_events: VecDeque<PolledEvent>,
@ -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, &params)
.into_iter()
.map(|event| relabel_provider_event(event, &provider_name))
.collect::<Vec<_>>();
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, &params)?;
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", &params).unwrap();
let replay_events =
normalize_provider_notification(&mut state, "paperclip/runResult", &params).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",
&params(result),
)
.unwrap_err()
.to_string()
.contains("contract revision"));
let mut malformed = opencode_result_state();
assert!(normalize_provider_notification(
&mut malformed,
"paperclip/runResult",
&params(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",
&params(result),
)
.unwrap_err()
.to_string()
.contains("bound completion criteria"));
let mut conflicting = opencode_result_state();
normalize_provider_notification(
&mut conflicting,
"paperclip/runResult",
&params(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",
&params(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",
&params(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(),