diff --git a/crates/coven-cli/src/api.rs b/crates/coven-cli/src/api.rs index 54e606d6d..f1a98d34d 100644 --- a/crates/coven-cli/src/api.rs +++ b/crates/coven-cli/src/api.rs @@ -7394,9 +7394,20 @@ fn proposal_decision_semantics( if decision != "reject" { return Err("proposal-expiry-decision-invalid"); } + let window_close = scheduled + .filter(|proposal| proposal.veto_deadline().is_some()) + .map(|_| coven_threads_core::ProposalWindowCloseAuditDetail { + reason: coven_threads_core::WindowCloseReason::RevalidationFailed, + replay_hash_matched: Some(false), + rationale: note.map(str::to_string), + }); return Ok(ProposalDecisionSemantics { rejection_event: coven_threads_core::AuditEventType::ProposalRejected, - rejection_decision: "expired", + rejection_decision: if window_close.is_some() { + coven_threads_core::WindowCloseReason::RevalidationFailed.tag() + } else { + "expired" + }, approval_path_label: scheduled .map(|proposal| { proposal @@ -7406,7 +7417,7 @@ fn proposal_decision_semantics( .to_string() }) .unwrap_or_else(|| "human_review".to_string()), - window_close: None, + window_close, }); } let Some(scheduled) = scheduled else { @@ -7498,6 +7509,21 @@ fn proposal_decision_semantics( }) } +fn scheduled_rejection_window_close( + scheduled: Option<&crate::proposal_scheduler::ScheduledProposal>, + reason: coven_threads_core::WindowCloseReason, + replay_hash_matched: Option, + rationale: Option<&str>, +) -> Option { + scheduled + .filter(|proposal| proposal.veto_deadline().is_some()) + .map(|_| coven_threads_core::ProposalWindowCloseAuditDetail { + reason, + replay_hash_matched, + rationale: rationale.map(str::to_string), + }) +} + fn revalidate_scheduled_materialized_before( workspace: &Path, scheduled: &crate::proposal_scheduler::ScheduledProposal, @@ -9393,12 +9419,29 @@ fn expire_threads_proposal(coven_home: &Path, proposal_id: &str) -> Result bool { + if response.status == 200 { + return true; + } + !path.exists() + && serde_json::from_str::(&response.body) + .ok() + .and_then(|body| body.get("terminal").and_then(Value::as_bool)) + == Some(true) +} + fn proposal_retention_expired( coven_home: &Path, conn: &rusqlite::Connection, document: &ProposalEnvelopeDocument, now: time::OffsetDateTime, ) -> bool { + if document + .scheduled() + .is_some_and(|proposal| proposal.veto_deadline().is_some()) + { + return false; + } document .pending() .staged_at @@ -9747,10 +9790,49 @@ fn decide_threads_proposal_inner( } }; drop(raw); + if document.pending().id.0 != proposal_uuid { + audit_reservation.release_if_unneeded()?; + return json_response(409, &json!({ "blocked": true, "why": "proposal-corrupt" })); + } let Some(review_kind) = PendingReviewKind::from_stored_value(document.review_kind.as_deref()) else { + let opened_window = proposal_window_context(&conn, proposal_id)?; if document.decision_state.is_some() { + return quarantine_proposal_recovery_claim( + coven_home, + &mut claim, + audit_reservation, + proposal_id, + "recovery-corrupt", + "proposal-recovery-corrupt", + "applied recovery state is corrupt", + ); + } else if let Some(opened) = opened_window.as_ref() { + let approval_path_label = document + .scheduled() + .map(|proposal| proposal.classification().approval_path.display_label()) + .unwrap_or("familiar_review"); claim.preserve(); + append_open_window_revalidation_failure( + &conn, + proposal_id, + document.pending(), + opened, + approval_path_label, + None, + decision_now, + )?; + audit_reservation.finish()?; + claim.consume()?; + return json_response( + 409, + &json!({ + "blocked": true, + "why": "proposal-corrupt", + "proposalId": proposal_id, + "terminal": true, + }), + ); } else { claim.restore_pending(&document)?; } @@ -9839,13 +9921,9 @@ fn decide_threads_proposal_inner( &json!({ "blocked": true, "why": "proposal-revision-required" }), ); } - } + }; let scheduled = document.scheduled(); let pending = document.pending(); - if pending.id.0 != proposal_uuid { - audit_reservation.release_if_unneeded()?; - return json_response(409, &json!({ "blocked": true, "why": "proposal-corrupt" })); - } if decision == "approve" { if let Err(error) = ward::validate_staged_edit_budget(&pending.edits) { if let Some(limit) = ward::ward_edit_budget_failure(&error) { @@ -9881,21 +9959,148 @@ fn decide_threads_proposal_inner( } else { note }; - let Some(familiar_id) = human_familiar_id_for_weave(coven_home, pending.familiar_id)? else { - audit_reservation.release_if_unneeded()?; - return json_response( - 409, - &json!({ "blocked": true, "why": "proposal-familiar-missing" }), - ); + let approval_path_label = scheduled + .map(|proposal| proposal.classification().approval_path.display_label()) + .unwrap_or("human_review"); + let mut opened_window = proposal_window_context(&conn, proposal_id)?; + let familiar_id = match human_familiar_id_for_weave(coven_home, pending.familiar_id) { + Ok(Some(familiar_id)) => familiar_id, + Ok(None) | Err(_) if applying_state.is_some() => { + return quarantine_proposal_recovery_claim( + coven_home, + &mut claim, + audit_reservation, + proposal_id, + "familiar-unavailable", + "proposal-recovery-familiar-unavailable", + "interrupted apply has no readable live familiar authority", + ); + } + Ok(None) | Err(_) if opened_window.is_some() && applying_state.is_none() => { + claim.preserve(); + append_open_window_revalidation_failure( + &conn, + proposal_id, + pending, + opened_window + .as_ref() + .expect("guard established opened window context"), + approval_path_label, + note.as_deref(), + decision_now, + )?; + audit_reservation.finish()?; + claim.consume()?; + return json_response( + 409, + &json!({ + "blocked": true, + "why": "proposal-familiar-missing", + "proposalId": proposal_id, + "terminal": true, + }), + ); + } + Ok(None) => { + audit_reservation.release_if_unneeded()?; + return json_response( + 409, + &json!({ "blocked": true, "why": "proposal-familiar-missing" }), + ); + } + Err(error) => return Err(error), }; let workspace = crate::cockpit_sources::familiar_workspace(coven_home, &familiar_id); - let Some(config) = ward::WardConfig::load(&workspace)? else { - audit_reservation.release_if_unneeded()?; + let config = match ward::WardConfig::load(&workspace) { + Ok(Some(config)) => config, + Ok(None) | Err(_) if applying_state.is_some() => { + return quarantine_proposal_recovery_claim( + coven_home, + &mut claim, + audit_reservation, + proposal_id, + "ward-unavailable", + "proposal-recovery-ward-unavailable", + "interrupted apply has no readable live Ward authority", + ); + } + Ok(None) | Err(_) if opened_window.is_some() && applying_state.is_none() => { + claim.preserve(); + append_open_window_revalidation_failure( + &conn, + proposal_id, + pending, + opened_window + .as_ref() + .expect("guard established opened window context"), + approval_path_label, + note.as_deref(), + decision_now, + )?; + audit_reservation.finish()?; + claim.consume()?; + return json_response( + 409, + &json!({ + "blocked": true, + "why": "ward-not-configured", + "proposalId": proposal_id, + "terminal": true, + }), + ); + } + Ok(None) => { + audit_reservation.release_if_unneeded()?; + return json_response( + 409, + &json!({ "blocked": true, "why": "ward-not-configured" }), + ); + } + Err(error) => return Err(error), + }; + if scheduled.is_some_and(|proposal| proposal.veto_deadline().is_some()) { + if !proposal_window_opened(&conn, proposal_id)? { + ensure_proposal_window_opened_audit_with_conn( + coven_home, + &conn, + scheduled.expect("guard established scheduled proposal"), + decision_now, + )?; + opened_window = proposal_window_context(&conn, proposal_id)?; + } + } else if applying_state.is_some() && opened_window.is_some() { + return quarantine_proposal_recovery_claim( + coven_home, + &mut claim, + audit_reservation, + proposal_id, + "window-state-inconsistent", + "proposal-recovery-window-state-inconsistent", + "applied recovery state conflicts with an opened veto window", + ); + } else if let Some(opened) = opened_window.as_ref() { + claim.preserve(); + append_open_window_revalidation_failure( + &conn, + proposal_id, + pending, + opened, + approval_path_label, + note.as_deref(), + decision_now, + )?; + audit_reservation.finish()?; + claim.consume()?; return json_response( 409, - &json!({ "blocked": true, "why": "ward-not-configured" }), + &json!({ + "blocked": true, + "why": "proposal-window-state-inconsistent", + "proposalId": proposal_id, + "terminal": true, + }), ); - }; + } let authorization = authorization_from_writer(&pending.writer); let edits = match staged_edits_to_ward_edits(pending) { Ok(edits) => edits, @@ -9905,7 +10110,46 @@ fn decide_threads_proposal_inner( } }; let targets: Vec = edits.iter().map(|edit| edit.target.clone()).collect(); - let ward = ward::Ward::new(&workspace, config.clone())?; + let ward = match ward::Ward::new(&workspace, config.clone()) { + Ok(ward) => ward, + Err(error) if applying_state.is_some() => { + return quarantine_proposal_recovery_claim( + coven_home, + &mut claim, + audit_reservation, + proposal_id, + "ward-invalid", + "proposal-recovery-ward-invalid", + &format!("interrupted apply cannot reconstruct live Ward authority: {error:#}"), + ); + } + Err(_error) if opened_window.is_some() && applying_state.is_none() => { + claim.preserve(); + append_open_window_revalidation_failure( + &conn, + proposal_id, + pending, + opened_window + .as_ref() + .expect("guard established opened window context"), + approval_path_label, + note.as_deref(), + decision_now, + )?; + audit_reservation.finish()?; + claim.consume()?; + return json_response( + 409, + &json!({ + "blocked": true, + "why": "proposal-revalidation-failed", + "proposalId": proposal_id, + "terminal": true, + }), + ); + } + Err(error) => return Err(error), + }; let adjudication = ward.evaluate(&ward::Proposal { targets: targets.clone(), authorization: authorization.clone(), @@ -9929,7 +10173,7 @@ fn decide_threads_proposal_inner( }), ); } - let state = crate::threads_gate::build_weave_state_for_writer_at( + let state_result = crate::threads_gate::build_weave_state_for_writer_at( &conn, &familiar_id, &workspace, @@ -9938,6 +10182,12 @@ fn decide_threads_proposal_inner( false, decision_now, Some(&pending.writer), + ); + let weave_hash = revalidation_failure_weave_hash( + coven_home, + proposal_id, + state_result, + opened_window.as_ref(), )?; let approval_path_label = scheduled .map(|proposal| { @@ -9969,7 +10219,7 @@ fn decide_threads_proposal_inner( event_type: coven_threads_core::AuditEventType::ProposalRejected, proposal_id, familiar_id: &familiar_id, - weave_hash: state.weave.weave_hash(), + weave_hash: &weave_hash, approver: Some(&pending.writer), files_touched: &targets, decision: "protected-target-not-proposable", @@ -10157,7 +10407,7 @@ fn decide_threads_proposal_inner( ) })); if coherence_revalidation_failed { - let state = crate::threads_gate::build_weave_state_at( + let state_result = crate::threads_gate::build_weave_state_at( &conn, &familiar_id, &workspace, @@ -10165,22 +10415,59 @@ fn decide_threads_proposal_inner( &[], false, decision_now, - )?; - append_proposal_refusal_audit( - &conn, + ); + let weave_hash = revalidation_failure_weave_hash( + coven_home, proposal_id, - &familiar_id, - state.weave.weave_hash(), - &pending.writer, - &targets, - pending.channel, - decision_now, + state_result, + opened_window.as_ref().filter(|_| applying_state.is_none()), )?; - audit_reservation.release_if_unneeded()?; - if applying_state.is_none() { - claim.restore_pending(&document)?; - } else { + if applying_state.is_none() + && scheduled.is_some_and(|proposal| proposal.veto_deadline().is_some()) + { claim.preserve(); + let window_close = scheduled_rejection_window_close( + scheduled, + coven_threads_core::WindowCloseReason::RevalidationFailed, + Some(false), + note.as_deref(), + ); + append_proposal_decision_audit( + &conn, + ProposalDecisionAudit { + event_type: coven_threads_core::AuditEventType::ProposalRejected, + proposal_id, + familiar_id: &familiar_id, + weave_hash: &weave_hash, + approver: Some(&pending.writer), + files_touched: &targets, + decision: coven_threads_core::WindowCloseReason::RevalidationFailed.tag(), + approval_rationale: note.as_deref(), + approval_path_label: &decision_semantics.approval_path_label, + window_close: window_close.as_ref(), + channel: pending.channel, + }, + decision_now, + )?; + audit_reservation.finish()?; + claim.consume()?; + } else { + append_proposal_refusal_audit( + &conn, + proposal_id, + &familiar_id, + &weave_hash, + &pending.writer, + &targets, + pending.channel, + decision_now, + )?; + audit_reservation.release_if_unneeded()?; + if applying_state.is_none() { + claim.restore_pending(&document)?; + } else { + claim.preserve(); + } } return json_response( 409, @@ -10189,6 +10476,8 @@ fn decide_threads_proposal_inner( "why": "proposal-revalidation-failed", "proposalId": proposal_id, "reviewKind": review_kind.as_str(), + "terminal": applying_state.is_none() + && scheduled.is_some_and(|proposal| proposal.veto_deadline().is_some()), }), ); } @@ -10197,7 +10486,7 @@ fn decide_threads_proposal_inner( && adjudication.is_blocked() && !coherence_rejection { - let state = crate::threads_gate::build_weave_state_for_writer_at( + let state_result = crate::threads_gate::build_weave_state_for_writer_at( &conn, &familiar_id, &workspace, @@ -10206,21 +10495,37 @@ fn decide_threads_proposal_inner( false, decision_now, Some(&pending.writer), + ); + let weave_hash = revalidation_failure_weave_hash( + coven_home, + proposal_id, + state_result, + opened_window.as_ref(), )?; claim.preserve(); + let window_close = scheduled_rejection_window_close( + scheduled, + coven_threads_core::WindowCloseReason::RevalidationFailed, + Some(false), + note.as_deref(), + ); append_proposal_decision_audit( &conn, ProposalDecisionAudit { event_type: coven_threads_core::AuditEventType::ProposalRejected, proposal_id, familiar_id: &familiar_id, - weave_hash: state.weave.weave_hash(), + weave_hash: &weave_hash, approver: Some(&pending.writer), files_touched: &targets, - decision: "proposal-live-adjudication-failed", + decision: window_close + .as_ref() + .map_or("proposal-live-adjudication-failed", |close| { + close.reason.tag() + }), approval_rationale: note.as_deref(), approval_path_label: &decision_semantics.approval_path_label, - window_close: None, + window_close: window_close.as_ref(), channel: pending.channel, }, decision_now, @@ -10233,6 +10538,7 @@ fn decide_threads_proposal_inner( "blocked": true, "why": "proposal-live-adjudication-failed", "proposalId": proposal_id, + "terminal": true, }), ); } @@ -10275,7 +10581,7 @@ fn decide_threads_proposal_inner( .map(|d| d.resolved.clone()) .collect() }; - let state = if review_kind == PendingReviewKind::Coherence { + let state_result = if review_kind == PendingReviewKind::Coherence { crate::threads_gate::build_weave_state_at( &conn, &familiar_id, @@ -10284,7 +10590,7 @@ fn decide_threads_proposal_inner( &[], false, decision_now, - )? + ) } else if scheduled.is_some() { crate::threads_gate::build_weave_state_for_writer_at( &conn, @@ -10295,7 +10601,7 @@ fn decide_threads_proposal_inner( false, decision_now, Some(&pending.writer), - )? + ) } else { crate::threads_gate::build_weave_state_at( &conn, @@ -10305,7 +10611,53 @@ fn decide_threads_proposal_inner( &gated_targets, false, decision_now, - )? + ) + }; + let state = match state_result { + Ok(state) => state, + Err(_error) + if applying_state.is_none() + && scheduled.is_some_and(|proposal| proposal.veto_deadline().is_some()) => + { + let weave_hash = proposal_window_opened_weave_hash(&conn, proposal_id)? + .context("opened proposal window is missing its committed weave hash")?; + let window_close = scheduled_rejection_window_close( + scheduled, + coven_threads_core::WindowCloseReason::RevalidationFailed, + Some(false), + note.as_deref(), + ); + claim.preserve(); + append_proposal_decision_audit( + &conn, + ProposalDecisionAudit { + event_type: coven_threads_core::AuditEventType::ProposalRejected, + proposal_id, + familiar_id: &familiar_id, + weave_hash: &weave_hash, + approver: Some(&pending.writer), + files_touched: &targets, + decision: coven_threads_core::WindowCloseReason::RevalidationFailed.tag(), + approval_rationale: note.as_deref(), + approval_path_label: &decision_semantics.approval_path_label, + window_close: window_close.as_ref(), + channel: pending.channel, + }, + decision_now, + )?; + audit_reservation.finish()?; + claim.consume()?; + return json_response( + 409, + &json!({ + "blocked": true, + "why": "proposal-evidence-replay-failed", + "proposalId": proposal_id, + "terminal": true, + }), + ); + } + Err(error) => return Err(error), }; if decision == "approve" && review_kind == PendingReviewKind::Authority @@ -10399,6 +10751,17 @@ fn decide_threads_proposal_inner( }; if let Some(reason) = rejection { claim.preserve(); + let close_reason = if reason == "proposal-evidence-diverged" { + coven_threads_core::WindowCloseReason::EvidenceDiverged + } else { + coven_threads_core::WindowCloseReason::RevalidationFailed + }; + let window_close = scheduled_rejection_window_close( + Some(scheduled), + close_reason, + Some(false), + note.as_deref(), + ); append_proposal_decision_audit( &conn, ProposalDecisionAudit { @@ -10408,10 +10771,12 @@ fn decide_threads_proposal_inner( weave_hash: state.weave.weave_hash(), approver: Some(&pending.writer), files_touched: &targets, - decision: reason, + decision: window_close + .as_ref() + .map_or(reason, |close| close.reason.tag()), approval_rationale: note.as_deref(), approval_path_label: &decision_semantics.approval_path_label, - window_close: None, + window_close: window_close.as_ref(), channel: pending.channel, }, decision_now, @@ -10424,6 +10789,7 @@ fn decide_threads_proposal_inner( "blocked": true, "why": reason, "proposalId": proposal_id, + "terminal": true, }), ); } @@ -11023,17 +11389,50 @@ fn quarantine_scheduler_candidate(coven_home: &Path, path: &Path, error: &anyhow ); } -fn read_pending_proposal_document(path: &Path) -> Result { - let raw = read_pending_proposal_file(path) - .with_context(|| format!("reading scheduled proposal {}", path.display()))?; - ProposalEnvelopeDocument::parse_preflighted(&raw) - .with_context(|| format!("parsing scheduled proposal {}", path.display())) -} - -fn parse_scheduler_authority_document( - path: &Path, - mut document: ProposalEnvelopeDocument, -) -> Result<( +fn quarantine_proposal_recovery_claim( + coven_home: &Path, + claim: &mut PendingDecisionClaim, + audit_reservation: store::WardAuditReservation<'_>, + proposal_id: &str, + quarantine_reason: &str, + why: &str, + summary: &str, +) -> Result { + claim.preserve(); + audit_reservation.release_if_unneeded()?; + let quarantine = crate::proposal_store::quarantine(coven_home, &claim.path, quarantine_reason)? + .context("proposal recovery claim disappeared before quarantine")?; + crate::daemon::append_daemon_recovery_log( + coven_home, + &format!( + "threads proposal {proposal_id}: {summary}; quarantined at {} for manual recovery", + quarantine.display() + ), + ); + json_response( + 409, + &json!({ + "blocked": true, + "why": why, + "proposalId": proposal_id, + "terminal": false, + "manualRecoveryRequired": true, + "quarantined": true, + }), + ) +} + +fn read_pending_proposal_document(path: &Path) -> Result { + let raw = read_pending_proposal_file(path) + .with_context(|| format!("reading scheduled proposal {}", path.display()))?; + ProposalEnvelopeDocument::parse_preflighted(&raw) + .with_context(|| format!("parsing scheduled proposal {}", path.display())) +} + +fn parse_scheduler_authority_document( + path: &Path, + mut document: ProposalEnvelopeDocument, +) -> Result<( Uuid, time::OffsetDateTime, Option, @@ -11115,8 +11514,41 @@ pub(crate) fn process_due_threads_proposals(coven_home: &Path) -> Result continue; } }; + let (opened_window, retention_expired) = { + let conn = store::open_store(&store_path(coven_home))?; + ( + proposal_window_opened(&conn, &document.pending().id.0.to_string())?, + proposal_retention_expired(coven_home, &conn, &document, now), + ) + }; + let mut revalidation_required = opened_window + && document + .scheduled() + .is_none_or(|proposal| proposal.veto_deadline().is_none()); + if revalidation_required { + crate::daemon::append_daemon_recovery_log( + coven_home, + &format!( + "threads scheduler: proposal {} has opened-window history without a \ + veto-window approval path; recovering through the terminal decision boundary", + path.display() + ), + ); + } let protected_targets = match pending_document_protected_targets(coven_home, &document) { Ok(targets) => targets, + Err(error) if opened_window => { + crate::daemon::append_daemon_recovery_log( + coven_home, + &format!( + "threads scheduler: live classification failed for opened proposal {}: \ + {error:#}; recovering through the terminal decision boundary", + path.display() + ), + ); + revalidation_required = true; + Vec::new() + } Err(error) => { crate::daemon::append_daemon_recovery_log( coven_home, @@ -11129,7 +11561,7 @@ pub(crate) fn process_due_threads_proposals(coven_home: &Path) -> Result continue; } }; - if !protected_targets.is_empty() { + if !protected_targets.is_empty() || revalidation_required { if let Some(proposal) = document.scheduled() { if let Err(error) = ensure_proposal_window_opened_audit(coven_home, proposal, now) { crate::daemon::append_daemon_recovery_log( @@ -11157,31 +11589,21 @@ pub(crate) fn process_due_threads_proposals(coven_home: &Path) -> Result decision, body.as_deref(), ) { - Ok(response) if response.status == 200 => completed += 1, - Ok(response) => { - let body: Value = serde_json::from_str(&response.body)?; - if response.status == 409 - && body["terminal"] == true - && body["why"] == "protected-target-not-proposable" - { - completed += 1; - } else { - crate::daemon::append_daemon_recovery_log( - coven_home, - &format!( - "threads scheduler: protected proposal rejection failed for {}: \ - HTTP {} {}; retained for retry", - path.display(), - response.status, - response.body - ), - ); - } - } + Ok(response) if scheduler_response_completed(&response, &path) => completed += 1, + Ok(response) => crate::daemon::append_daemon_recovery_log( + coven_home, + &format!( + "threads scheduler: proposal refusal failed for {}: \ + HTTP {} {}; retained for retry", + path.display(), + response.status, + response.body + ), + ), Err(error) => crate::daemon::append_daemon_recovery_log( coven_home, &format!( - "threads scheduler: protected proposal rejection failed for {}: \ + "threads scheduler: proposal refusal failed for {}: \ {error:#}; retained for retry", path.display() ), @@ -11189,7 +11611,7 @@ pub(crate) fn process_due_threads_proposals(coven_home: &Path) -> Result } continue; } - let (proposal_id, staged_at, proposal, request) = + let (proposal_id, _staged_at, proposal, request) = match parse_scheduler_authority_document(&path, document) { Ok(parsed) => parsed, Err(error) => { @@ -11208,16 +11630,11 @@ pub(crate) fn process_due_threads_proposals(coven_home: &Path) -> Result &request.decision, body.as_deref(), )?; - return Ok(response.status == 200); + return Ok(scheduler_response_completed(&response, &path)); } - let expired = staged_at - .checked_add(time::Duration::days( - crate::proposal_store::PENDING_PROPOSAL_RETENTION_DAYS, - )) - .is_some_and(|deadline| now >= deadline); - if expired { + if retention_expired { let response = expire_threads_proposal(coven_home, &proposal_id.to_string())?; - if response.status == 200 { + if scheduler_response_completed(&response, &path) { return Ok(true); } let reason = anyhow::anyhow!( @@ -11256,7 +11673,23 @@ pub(crate) fn process_due_threads_proposals(coven_home: &Path) -> Result "approve", None, )?; - Ok(response.status == 200) + if scheduler_response_completed(&response, &path) { + return Ok(true); + } + let reason = serde_json::from_str::(&response.body) + .ok() + .and_then(|body| body.get("why").and_then(Value::as_str).map(str::to_string)) + .unwrap_or_else(|| "unknown".to_string()); + crate::daemon::append_daemon_recovery_log( + coven_home, + &format!( + "threads scheduler: due proposal {} did not terminalize: HTTP {} {reason}; \ + retained for retry", + path.display(), + response.status + ), + ); + Ok(false) })(); match result { Ok(true) => completed += 1, @@ -11312,13 +11745,23 @@ fn recover_proposal_claim_document( decision, body.as_deref(), )?; - Ok(response.status == 200) + Ok(scheduler_response_completed(&response, claim_path)) } fn ensure_proposal_window_opened_audit( coven_home: &Path, proposal: &crate::proposal_scheduler::ScheduledProposal, now: time::OffsetDateTime, +) -> Result<()> { + let conn = store::open_store(&store_path(coven_home))?; + ensure_proposal_window_opened_audit_with_conn(coven_home, &conn, proposal, now) +} + +fn ensure_proposal_window_opened_audit_with_conn( + coven_home: &Path, + conn: &rusqlite::Connection, + proposal: &crate::proposal_scheduler::ScheduledProposal, + now: time::OffsetDateTime, ) -> Result<()> { let (Some(deadline), Some(earliest_close)) = (proposal.veto_deadline(), proposal.earliest_close()) @@ -11326,6 +11769,9 @@ fn ensure_proposal_window_opened_audit( return Ok(()); }; let pending = proposal.pending(); + if proposal_window_opened(conn, &pending.id.0.to_string())? { + return Ok(()); + } let Some(familiar_id) = human_familiar_id_for_weave(coven_home, pending.familiar_id)? else { anyhow::bail!("scheduled proposal familiar is missing"); }; @@ -11337,9 +11783,8 @@ fn ensure_proposal_window_opened_audit( .iter() .map(|edit| edit.surface.as_str().to_string()) .collect(); - let conn = store::open_store(&store_path(coven_home))?; let state = crate::threads_gate::build_weave_state_for_writer_at( - &conn, + conn, &familiar_id, &workspace, &config, @@ -12203,6 +12648,11 @@ struct ProposalTerminalAudit { files_touched: Vec, } +struct ProposalWindowContext { + familiar_id: String, + weave_hash: Vec, +} + fn terminal_decision_label(event_type: &str) -> &'static str { match event_type { "proposal_approved" => "approved", @@ -12240,6 +12690,127 @@ fn proposal_terminal_event( .transpose() } +fn proposal_window_opened(conn: &rusqlite::Connection, proposal_id: &str) -> Result { + conn.query_row( + "SELECT EXISTS ( + SELECT 1 + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type = 'proposal_window_opened' + )", + [proposal_id], + |row| row.get(0), + ) + .context("checking proposal window audit state") +} + +fn proposal_window_context( + conn: &rusqlite::Connection, + proposal_id: &str, +) -> Result> { + use rusqlite::OptionalExtension; + + conn.query_row( + "SELECT familiar_id, ward_hash + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type = 'proposal_window_opened' + ORDER BY id DESC + LIMIT 1", + [proposal_id], + |row| { + Ok(ProposalWindowContext { + familiar_id: row.get(0)?, + weave_hash: row.get(1)?, + }) + }, + ) + .optional() + .context("loading proposal window audit context") +} + +fn proposal_window_opened_weave_hash( + conn: &rusqlite::Connection, + proposal_id: &str, +) -> Result>> { + use rusqlite::OptionalExtension; + + conn.query_row( + "SELECT ward_hash + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type = 'proposal_window_opened' + ORDER BY id DESC + LIMIT 1", + [proposal_id], + |row| row.get(0), + ) + .optional() + .context("loading proposal window weave hash") +} + +fn revalidation_failure_weave_hash( + coven_home: &Path, + proposal_id: &str, + state: Result, + opened: Option<&ProposalWindowContext>, +) -> Result> { + match state { + Ok(state) => Ok(state.weave.weave_hash().to_vec()), + Err(error) => { + let Some(opened) = opened else { + return Err(error); + }; + crate::daemon::append_daemon_recovery_log( + coven_home, + &format!( + "threads proposal {proposal_id}: rejection replay failed: {error:#}; \ + closing against committed window evidence with replay_hash_matched=false" + ), + ); + Ok(opened.weave_hash.clone()) + } + } +} + +fn append_open_window_revalidation_failure( + conn: &rusqlite::Connection, + proposal_id: &str, + pending: &coven_threads_core::PendingProposal, + opened: &ProposalWindowContext, + approval_path_label: &str, + rationale: Option<&str>, + now: time::OffsetDateTime, +) -> Result<()> { + let files_touched: Vec = pending + .edits + .iter() + .map(|edit| edit.surface.as_str().to_string()) + .collect(); + let window_close = coven_threads_core::ProposalWindowCloseAuditDetail { + reason: coven_threads_core::WindowCloseReason::RevalidationFailed, + replay_hash_matched: Some(false), + rationale: rationale.map(str::to_string), + }; + append_proposal_decision_audit( + conn, + ProposalDecisionAudit { + event_type: coven_threads_core::AuditEventType::ProposalRejected, + proposal_id, + familiar_id: &opened.familiar_id, + weave_hash: &opened.weave_hash, + approver: Some(&pending.writer), + files_touched: &files_touched, + decision: coven_threads_core::WindowCloseReason::RevalidationFailed.tag(), + approval_rationale: rationale, + approval_path_label, + window_close: Some(&window_close), + channel: pending.channel, + }, + now, + ) +} + fn find_any_pending_decision_claim( coven_home: &Path, proposal_id: &str, @@ -12505,6 +13076,27 @@ fn append_proposal_decision_audit_with_probe_summary( probe_summary: Option<&crate::ward_probes::ProbeSummary>, now: time::OffsetDateTime, ) -> Result<()> { + let terminal = matches!( + audit.event_type, + coven_threads_core::AuditEventType::ProposalApproved + | coven_threads_core::AuditEventType::ProposalRejected + | coven_threads_core::AuditEventType::ProposalVetoed + ); + if terminal { + anyhow::ensure!( + proposal_terminal_event(conn, audit.proposal_id)?.is_none(), + "proposal already has a terminal audit event" + ); + let opened = proposal_window_opened(conn, audit.proposal_id)?; + anyhow::ensure!( + opened == audit.window_close.is_some(), + if opened { + "opened proposal window requires typed terminal close evidence" + } else { + "human proposal path must not fabricate window close evidence" + } + ); + } let files_touched = serde_json::to_string(audit.files_touched)?; let detail = match audit.event_type { coven_threads_core::AuditEventType::ProposalApproved => { @@ -31585,7 +32177,9 @@ tier = 0 assert!(!pending.exists()); let conn = store::open_store(&home.join("coven.sqlite3"))?; let (event, detail): (String, String) = conn.query_row( - "SELECT event_type, detail FROM ward_audit WHERE proposal_id = ?1", + "SELECT event_type, detail + FROM ward_audit + WHERE proposal_id = ?1 AND event_type = 'proposal_vetoed'", [&proposal_id], |row| Ok((row.get(0)?, row.get(1)?)), )?; @@ -31655,7 +32249,7 @@ tier = 0 } #[test] - fn threads_scheduled_deadline_replay_refuses_diverged_before_image() -> Result<()> { + fn threads_scheduled_human_path_rejects_existing_window_state() -> Result<()> { let temp = tempfile::tempdir()?; let home = temp.path(); let (pending, proposal_id) = stage_scheduled_reviewed_edit( @@ -31663,8 +32257,30 @@ tier = 0 coven_threads_core::ApprovalPath::HumanApproval, time::OffsetDateTime::now_utc(), )?; - let target = home.join("familiars/sage/reviewed/skill.md"); - std::fs::write(&target, "concurrent")?; + let now = time::OffsetDateTime::now_utc(); + let detail = coven_threads_core::ProposalWindowAuditDetail { + approval_path_label: "familiar_review".to_string(), + deadline: now + time::Duration::minutes(5), + earliest_close: now + time::Duration::minutes(1), + evidence_replay_hash_hex: "00".repeat(32), + affected_regions: Vec::new(), + }; + let conn = store::open_store(&home.join("coven.sqlite3"))?; + conn.execute( + "INSERT INTO ward_audit ( + event_type, proposal_id, familiar_id, ward_hash, decision, detail, + files_touched, submitted_at, decided_at + ) VALUES ( + 'proposal_window_opened', ?1, 'sage', zeroblob(32), 'window-opened', ?2, + '[]', ?3, ?3 + )", + rusqlite::params![ + proposal_id, + serde_json::to_string(&detail)?, + now.format(&time::format_description::well_known::Rfc3339)?, + ], + )?; + drop(conn); let decision_body = scheduled_decision_body(home, &proposal_id, None)?; let response = handle_request_with_body( @@ -31677,111 +32293,981 @@ tier = 0 assert_eq!(response.status, 409, "got {}", response.body); let body: Value = serde_json::from_str(&response.body)?; - assert_eq!(body["why"], "proposal-evidence-diverged"); - assert_eq!(std::fs::read_to_string(target)?, "concurrent"); - assert!(!pending.exists(), "failed deadline replay is terminal"); - let conn = store::open_store(&home.join("coven.sqlite3"))?; - let event: String = conn.query_row( - "SELECT event_type FROM ward_audit - WHERE proposal_id = ?1 AND event_type = 'proposal_rejected'", - [&proposal_id], - |row| row.get(0), - )?; - assert_eq!(event, "proposal_rejected"); - Ok(()) - } - - #[test] - fn threads_scheduled_rejects_live_promotion_to_protected_tier() -> Result<()> { - let temp = tempfile::tempdir()?; - let home = temp.path(); - let veto = coven_threads_core::VetoWindow::new( - std::time::Duration::from_secs(300), - std::time::Duration::from_secs(60), - ); - let (pending, proposal_id) = stage_scheduled_reviewed_edit( - home, - coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, - time::OffsetDateTime::now_utc(), - )?; - assert_eq!(process_due_threads_proposals(home)?, 0); - let ward_path = home.join("familiars/sage/ward.toml"); - let ward = std::fs::read_to_string(&ward_path)? - .replace( - "protected_surface = [\"SOUL.md\"]", - "protected_surface = [\"SOUL.md\", \"reviewed/skill.md\"]", - ) - .replace( - "path = \"reviewed/\"\ntier = 1", - "path = \"reviewed/skill.md\"\ntier = 0", - ); - std::fs::write(&ward_path, ward)?; - - assert_eq!(process_due_threads_proposals(home)?, 1); + assert_eq!(body["why"], "proposal-window-state-inconsistent"); + assert_eq!(body["terminal"], true); assert!(!pending.exists()); assert_eq!( std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, "before" ); let conn = store::open_store(&home.join("coven.sqlite3"))?; - let (decision, detail): (String, String) = conn.query_row( - "SELECT decision, detail FROM ward_audit + let detail: String = conn.query_row( + "SELECT detail + FROM ward_audit WHERE proposal_id = ?1 AND event_type = 'proposal_rejected'", [&proposal_id], - |row| Ok((row.get(0)?, row.get(1)?)), + |row| row.get(0), )?; - assert_eq!(decision, "protected-target-not-proposable"); let close: coven_threads_core::ProposalWindowCloseAuditDetail = serde_json::from_str(&detail)?; assert_eq!( close.reason, coven_threads_core::WindowCloseReason::RevalidationFailed ); - assert_eq!(close.replay_hash_matched, Some(false)); - assert_eq!( - close.rationale.as_deref(), - Some("protected-target-not-proposable") - ); - - let retry = handle_request_with_body( - "POST", - &format!("/api/v1/threads/proposals/{proposal_id}/approve"), - home, - None, - Some("{}"), - )?; - assert_eq!(retry.status, 409, "got {}", retry.body); - let terminal_count: i64 = conn.query_row( - "SELECT COUNT(*) FROM ward_audit - WHERE proposal_id = ?1 - AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", - [&proposal_id], - |row| row.get(0), - )?; - assert_eq!(terminal_count, 1); Ok(()) } #[test] - fn threads_scheduler_first_pass_protected_rejection_records_window_close_detail() -> Result<()> + fn threads_scheduled_window_corruption_checks_internal_id_before_terminal_close() -> Result<()> { let temp = tempfile::tempdir()?; let home = temp.path(); - #[cfg(feature = "threads-test-clock")] - let staged_at = - seed_deterministic_threads_clock(home, "fixture-cap", "2026-09-09T10:00:00Z")?; - #[cfg(not(feature = "threads-test-clock"))] - let staged_at = time::OffsetDateTime::now_utc(); - let veto = coven_threads_core::VetoWindow::new( - std::time::Duration::from_secs(300), - std::time::Duration::from_secs(60), - ); let (pending, proposal_id) = stage_scheduled_reviewed_edit( home, - coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, - staged_at, + coven_threads_core::ApprovalPath::HumanApproval, + time::OffsetDateTime::now_utc(), )?; - let ward_path = home.join("familiars/sage/ward.toml"); + let now = time::OffsetDateTime::now_utc(); + let detail = coven_threads_core::ProposalWindowAuditDetail { + approval_path_label: "familiar_review".to_string(), + deadline: now + time::Duration::minutes(5), + earliest_close: now + time::Duration::minutes(1), + evidence_replay_hash_hex: "00".repeat(32), + affected_regions: Vec::new(), + }; + let conn = store::open_store(&home.join("coven.sqlite3"))?; + conn.execute( + "INSERT INTO ward_audit ( + event_type, proposal_id, familiar_id, ward_hash, decision, detail, + files_touched, submitted_at, decided_at + ) VALUES ( + 'proposal_window_opened', ?1, 'sage', zeroblob(32), 'window-opened', ?2, + '[]', ?3, ?3 + )", + rusqlite::params![ + proposal_id, + serde_json::to_string(&detail)?, + now.format(&time::format_description::well_known::Rfc3339)?, + ], + )?; + drop(conn); + + let mismatched_id = Uuid::new_v4().to_string(); + let mut staged: Value = serde_json::from_str( + &String::from_utf8(std::fs::read(&pending)?)?.replace(&proposal_id, &mismatched_id), + )?; + staged["reviewKind"] = json!("automatic"); + std::fs::write(&pending, serde_json::to_vec_pretty(&staged)?)?; + + let response = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some("{}"), + )?; + + assert_eq!(response.status, 409, "got {}", response.body); + let body: Value = serde_json::from_str(&response.body)?; + assert_eq!(body["why"], "proposal-corrupt"); + assert!( + body.get("terminal").is_none(), + "mismatched proposal ids must fail before any terminal close path" + ); + assert!(pending.exists(), "corrupt proposal must remain inspectable"); + assert!( + find_pending_decision_claim(home, &proposal_id, "approve").is_none(), + "corrupt proposal must not leave a hidden claim marker" + ); + let preserved: Value = serde_json::from_slice(&std::fs::read(&pending)?)?; + assert_eq!(preserved["reviewKind"], "automatic"); + assert_eq!(preserved["pending"]["id"], mismatched_id); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 0); + Ok(()) + } + + #[test] + fn threads_scheduled_human_window_recovery_quarantines_applied_claim() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let (_, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::HumanApprovalWithRationale, + time::OffsetDateTime::now_utc(), + )?; + let decision_body = + scheduled_decision_body(home, &proposal_id, Some("principal reviewed"))?; + set_proposal_decision_failpoint(Some(( + ProposalDecisionFailpoint::ApplyBeforeAudit, + proposal_id.clone(), + ))); + + let interrupted = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + ); + assert!( + interrupted.is_err(), + "failpoint must interrupt the decision" + ); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "after" + ); + let claim = find_pending_decision_claim(home, &proposal_id, "approve") + .expect("interrupted approval leaves a durable claim"); + + let now = time::OffsetDateTime::now_utc(); + let detail = coven_threads_core::ProposalWindowAuditDetail { + approval_path_label: "familiar_review".to_string(), + deadline: now + time::Duration::minutes(5), + earliest_close: now + time::Duration::minutes(1), + evidence_replay_hash_hex: "00".repeat(32), + affected_regions: Vec::new(), + }; + let conn = store::open_store(&home.join("coven.sqlite3"))?; + conn.execute( + "INSERT INTO ward_audit ( + event_type, proposal_id, familiar_id, ward_hash, decision, detail, + files_touched, submitted_at, decided_at + ) VALUES ( + 'proposal_window_opened', ?1, 'sage', zeroblob(32), 'window-opened', ?2, + '[]', ?3, ?3 + )", + rusqlite::params![ + proposal_id, + serde_json::to_string(&detail)?, + now.format(&time::format_description::well_known::Rfc3339)?, + ], + )?; + drop(conn); + + let response = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + None, + )?; + + assert_eq!(response.status, 409, "got {}", response.body); + let body: Value = serde_json::from_str(&response.body)?; + assert_eq!(body["why"], "proposal-recovery-window-state-inconsistent"); + assert_eq!(body["terminal"], false); + assert_eq!(body["manualRecoveryRequired"], true); + assert_eq!(body["quarantined"], true); + assert!( + !claim.exists(), + "active recovery claim must stop hot-looping" + ); + let quarantine_entries = std::fs::read_dir(home.join("pending/quarantine"))? + .collect::>>()?; + assert_eq!( + quarantine_entries.len(), + 1, + "the recovery artifact must remain available for repair" + ); + let quarantined: Value = + serde_json::from_slice(&std::fs::read(quarantine_entries[0].path())?)?; + assert!( + quarantined.get("decisionState").is_some(), + "quarantine must retain the applied recovery evidence" + ); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "after" + ); + assert_eq!(process_due_threads_proposals(home)?, 0); + let recovery_log = std::fs::read_to_string(crate::daemon::daemon_recovery_log_path(home))?; + assert!(recovery_log.contains("quarantined")); + assert!(recovery_log.contains("manual recovery")); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 0, "recovery must not invent an outcome"); + Ok(()) + } + + #[test] + fn threads_scheduler_quarantines_corrupt_applied_claim_for_manual_recovery() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let (_, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::HumanApprovalWithRationale, + time::OffsetDateTime::now_utc(), + )?; + let decision_body = + scheduled_decision_body(home, &proposal_id, Some("principal reviewed"))?; + set_proposal_decision_failpoint(Some(( + ProposalDecisionFailpoint::ApplyBeforeAudit, + proposal_id.clone(), + ))); + + let interrupted = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + ); + assert!( + interrupted.is_err(), + "failpoint must interrupt the decision" + ); + let claim = find_pending_decision_claim(home, &proposal_id, "approve") + .expect("interrupted approval leaves a durable claim"); + let store_path = home.join("coven.sqlite3"); + let conn = store::open_store(&store_path)?; + let reservation_bytes = proposal_decision_audit_reservation_bytes(&conn, "approve")?; + store::WardAuditReservation::acquire( + &conn, + &store_path, + format!("proposal:{proposal_id}:approve"), + "proposal-approval", + reservation_bytes, + )? + .preserve()?; + drop(conn); + + let mut corrupted: Value = serde_json::from_slice(&std::fs::read(&claim)?)?; + corrupted["reviewKind"] = json!("automatic"); + std::fs::write(&claim, serde_json::to_vec_pretty(&corrupted)?)?; + + assert_eq!(process_due_threads_proposals(home)?, 0); + assert_eq!(process_due_threads_proposals(home)?, 0); + assert!( + !claim.exists(), + "corrupt recovery claim must stop retrying in-place" + ); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "after" + ); + let quarantine_entries = std::fs::read_dir(home.join("pending/quarantine"))? + .collect::>>()?; + assert_eq!( + quarantine_entries.len(), + 1, + "manual recovery must keep exactly one quarantined artifact" + ); + let quarantined: Value = + serde_json::from_slice(&std::fs::read(quarantine_entries[0].path())?)?; + assert!( + quarantined.get("decisionState").is_some(), + "quarantine must retain the applied recovery state" + ); + assert!( + quarantined.get("decisionRequest").is_some(), + "quarantine must retain the applied intent" + ); + let conn = store::open_store(&store_path)?; + let reservations: i64 = conn.query_row( + "SELECT COUNT(*) + FROM coven_ward_audit_reservations + WHERE token = ?1", + [format!("proposal:{proposal_id}:approve")], + |row| row.get(0), + )?; + assert_eq!( + reservations, 1, + "manual recovery quarantine must preserve the reservation evidence" + ); + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 0, "recovery must not invent an outcome"); + let recovery_log = std::fs::read_to_string(crate::daemon::daemon_recovery_log_path(home))?; + assert!(recovery_log.contains("corrupt")); + assert!(recovery_log.contains("manual recovery")); + Ok(()) + } + + #[test] + fn threads_decision_closes_new_window_when_live_ward_construction_fails() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { + veto: coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ), + }, + crate::threads_clock::now(home)? - time::Duration::minutes(10), + )?; + let decision_body = scheduled_decision_body(home, &proposal_id, None)?; + let ward_path = home.join("familiars/sage/ward.toml"); + let config = + std::fs::read_to_string(&ward_path)?.replace("path = \"reviewed/\"", "path = \"[\""); + std::fs::write(ward_path, config)?; + let response = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + )?; + assert_eq!(response.status, 409, "{}", response.body); + assert_eq!( + serde_json::from_str::(&response.body)?["terminal"], + true + ); + assert!(!pending.exists()); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "before" + ); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + assert!(proposal_window_opened(&conn, &proposal_id)?); + let detail: String = conn.query_row( + "SELECT detail FROM ward_audit WHERE proposal_id = ?1 + AND event_type = 'proposal_rejected'", + [&proposal_id], + |row| row.get(0), + )?; + let close: coven_threads_core::ProposalWindowCloseAuditDetail = + serde_json::from_str(&detail)?; + assert_eq!( + close.reason, + coven_threads_core::WindowCloseReason::RevalidationFailed + ); + assert_eq!(close.replay_hash_matched, Some(false)); + assert_eq!(process_due_threads_proposals(home)?, 0); + Ok(()) + } + + #[test] + fn threads_scheduler_quarantines_applied_claim_with_unavailable_authority() -> Result<()> { + for failure in [ + "missing-familiar", + "invalid-familiar", + "missing-ward", + "invalid-ward", + "invalid-ward-glob", + ] { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let (_, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::HumanApproval, + crate::threads_clock::now(home)?, + )?; + let body = scheduled_decision_body(home, &proposal_id, None)?; + set_proposal_decision_failpoint(Some(( + ProposalDecisionFailpoint::ApplyBeforeAudit, + proposal_id.clone(), + ))); + assert!(handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&body), + ) + .is_err()); + let claim = find_pending_decision_claim(home, &proposal_id, "approve") + .context("durable interrupted claim")?; + let original_claim = std::fs::read(&claim)?; + let ward_path = home.join("familiars/sage/ward.toml"); + match failure { + "missing-familiar" => std::fs::write(home.join("familiars.toml"), "")?, + "invalid-familiar" => std::fs::write(home.join("familiars.toml"), "invalid = [")?, + "missing-ward" => std::fs::remove_file(&ward_path)?, + "invalid-ward" => std::fs::write(&ward_path, "invalid = [")?, + "invalid-ward-glob" => { + let config = std::fs::read_to_string(&ward_path)? + .replace("path = \"reviewed/\"", "path = \"[\""); + std::fs::write(&ward_path, config)?; + } + _ => unreachable!(), + } + assert_eq!(process_due_threads_proposals(home)?, 0, "{failure}"); + assert!( + !claim.exists(), + "{failure}: unrecoverable claim must stop hot-looping" + ); + assert_eq!(process_due_threads_proposals(home)?, 0, "{failure}"); + let quarantined = std::fs::read_dir(home.join("pending/quarantine"))? + .collect::>>()?; + assert_eq!(quarantined.len(), 1, "{failure}"); + assert_eq!( + std::fs::read(quarantined[0].path())?, + original_claim, + "{failure}" + ); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "after", + "{failure}" + ); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + assert!( + proposal_terminal_event(&conn, &proposal_id)?.is_none(), + "{failure}" + ); + assert!( + load_proposal_apply_intent(&conn, &proposal_id)?.is_some(), + "{failure}" + ); + } + Ok(()) + } + + #[test] + fn threads_scheduled_deadline_replay_refuses_diverged_before_image() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let veto = coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, + time::OffsetDateTime::now_utc() - time::Duration::minutes(10), + )?; + let document = read_pending_proposal_document(&pending)?; + let (_, _, scheduled, _) = parse_scheduler_authority_document(&pending, document)?; + ensure_proposal_window_opened_audit( + home, + scheduled + .as_ref() + .context("fixture carries scheduled proposal")?, + crate::threads_clock::now(home)?, + )?; + let target = home.join("familiars/sage/reviewed/skill.md"); + std::fs::write(&target, "concurrent")?; + let decision_body = scheduled_decision_body(home, &proposal_id, None)?; + + let response = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + )?; + + assert_eq!(response.status, 409, "got {}", response.body); + let body: Value = serde_json::from_str(&response.body)?; + assert_eq!(body["why"], "proposal-evidence-diverged"); + assert_eq!(std::fs::read_to_string(target)?, "concurrent"); + assert!(!pending.exists(), "failed deadline replay is terminal"); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let (event, detail): (String, String) = conn.query_row( + "SELECT event_type, detail FROM ward_audit + WHERE proposal_id = ?1 AND event_type = 'proposal_rejected'", + [&proposal_id], + |row| Ok((row.get(0)?, row.get(1)?)), + )?; + assert_eq!(event, "proposal_rejected"); + let close: coven_threads_core::ProposalWindowCloseAuditDetail = + serde_json::from_str(&detail)?; + assert_eq!( + close.reason, + coven_threads_core::WindowCloseReason::EvidenceDiverged + ); + assert_eq!(close.replay_hash_matched, Some(false)); + let retry = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + )?; + assert_eq!(retry.status, 409, "got {}", retry.body); + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 1); + Ok(()) + } + + #[test] + fn threads_scheduled_deadline_replay_failure_closes_window() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let veto = coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, + time::OffsetDateTime::now_utc() - time::Duration::minutes(10), + )?; + let document = read_pending_proposal_document(&pending)?; + let (_, _, scheduled, _) = parse_scheduler_authority_document(&pending, document)?; + ensure_proposal_window_opened_audit( + home, + scheduled + .as_ref() + .context("fixture carries scheduled proposal")?, + crate::threads_clock::now(home)?, + )?; + let target = home.join("familiars/sage/reviewed/skill.md"); + std::fs::remove_file(&target)?; + std::fs::create_dir(&target)?; + let decision_body = scheduled_decision_body(home, &proposal_id, None)?; + + let response = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + )?; + + assert_eq!(response.status, 409, "got {}", response.body); + let body: Value = serde_json::from_str(&response.body)?; + assert_eq!(body["why"], "proposal-evidence-replay-failed"); + assert_eq!(body["terminal"], true); + assert!(!pending.exists()); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let detail: String = conn.query_row( + "SELECT detail + FROM ward_audit + WHERE proposal_id = ?1 AND event_type = 'proposal_rejected'", + [&proposal_id], + |row| row.get(0), + )?; + let close: coven_threads_core::ProposalWindowCloseAuditDetail = + serde_json::from_str(&detail)?; + assert_eq!( + close.reason, + coven_threads_core::WindowCloseReason::RevalidationFailed + ); + assert_eq!(close.replay_hash_matched, Some(false)); + Ok(()) + } + + #[test] + fn threads_scheduled_missing_ward_closes_existing_window() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let veto = coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, + time::OffsetDateTime::now_utc() - time::Duration::minutes(10), + )?; + let document = read_pending_proposal_document(&pending)?; + let (_, _, scheduled, _) = parse_scheduler_authority_document(&pending, document)?; + ensure_proposal_window_opened_audit( + home, + scheduled + .as_ref() + .context("fixture carries scheduled proposal")?, + crate::threads_clock::now(home)?, + )?; + std::fs::remove_file(home.join("familiars/sage/ward.toml"))?; + let decision_body = scheduled_decision_body(home, &proposal_id, None)?; + + let response = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + )?; + + assert_eq!(response.status, 409, "got {}", response.body); + let body: Value = serde_json::from_str(&response.body)?; + assert_eq!(body["why"], "ward-not-configured"); + assert_eq!(body["terminal"], true); + assert!(!pending.exists()); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) + FROM ward_audit + WHERE proposal_id = ?1 AND event_type = 'proposal_rejected'", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 1); + Ok(()) + } + + #[test] + fn threads_scheduler_closes_existing_window_when_live_authority_is_unavailable() -> Result<()> { + for failure in [ + "missing-ward", + "invalid-ward", + "invalid-ward-glob", + "missing-familiar", + "invalid-familiar", + "invalid-protected-baseline", + ] { + let temp = tempfile::tempdir()?; + let home = temp.path(); + #[cfg(feature = "threads-test-clock")] + seed_deterministic_threads_clock(home, "fixture-cap", "2026-09-09T10:00:00Z")?; + let now = crate::threads_clock::now(home)?; + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { + veto: coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ), + }, + now - time::Duration::minutes(10), + )?; + let document = read_pending_proposal_document(&pending)?; + ensure_proposal_window_opened_audit( + home, + document.scheduled().context("scheduled fixture")?, + now, + )?; + let ward_path = home.join("familiars/sage/ward.toml"); + match failure { + "missing-ward" => std::fs::remove_file(&ward_path)?, + "invalid-ward" => std::fs::write(&ward_path, "invalid = [")?, + "invalid-ward-glob" => { + let config = std::fs::read_to_string(&ward_path)? + .replace("path = \"reviewed/\"", "path = \"[\""); + std::fs::write(&ward_path, config)?; + } + "missing-familiar" => std::fs::write(home.join("familiars.toml"), "")?, + "invalid-familiar" => std::fs::write(home.join("familiars.toml"), "invalid = [")?, + "invalid-protected-baseline" => { + std::fs::remove_file(home.join("familiars/sage/SOUL.md"))?; + std::fs::create_dir(home.join("familiars/sage/SOUL.md"))?; + } + _ => unreachable!(), + } + + assert_eq!( + process_due_threads_proposals(home)?, + 1, + "{failure}: {}", + std::fs::read_to_string(crate::daemon::daemon_recovery_log_path(home)) + .unwrap_or_default() + ); + assert!(!pending.exists(), "{failure}"); + assert_eq!(process_due_threads_proposals(home)?, 0, "{failure}"); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "before", + "{failure}" + ); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let (event, detail, decided_at): (String, String, String) = conn.query_row( + "SELECT event_type, detail, decided_at FROM ward_audit + WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + assert_eq!(event, "proposal_rejected", "{failure}"); + let close: coven_threads_core::ProposalWindowCloseAuditDetail = + serde_json::from_str(&detail)?; + assert_eq!( + close.reason, + coven_threads_core::WindowCloseReason::RevalidationFailed, + "{failure}" + ); + assert_eq!(close.replay_hash_matched, Some(false), "{failure}"); + let decided_at = time::OffsetDateTime::parse( + &decided_at, + &time::format_description::well_known::Rfc3339, + )?; + #[cfg(feature = "threads-test-clock")] + assert_eq!(decided_at, now, "{failure}"); + #[cfg(not(feature = "threads-test-clock"))] + assert!(decided_at >= now, "{failure}"); + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) FROM ward_audit WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 1, "{failure}"); + let reservations: i64 = conn.query_row( + "SELECT COUNT(*) FROM coven_ward_audit_reservations", + [], + |row| row.get(0), + )?; + assert_eq!(reservations, 0, "{failure}"); + } + Ok(()) + } + + #[test] + fn threads_scheduler_closes_rejected_window_when_replay_is_unavailable() -> Result<()> { + assert_rejected_window_with_unavailable_replay("protected-promotion")?; + #[cfg(unix)] + assert_rejected_window_with_unavailable_replay("blocked-symlink")?; + Ok(()) + } + + fn assert_rejected_window_with_unavailable_replay(failure: &str) -> Result<()> { + let temp = tempfile::tempdir()?; + #[cfg(unix)] + let outside = tempfile::tempdir()?; + let home = temp.path(); + let now = crate::threads_clock::now(home)?; + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { + veto: coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ), + }, + now - time::Duration::minutes(10), + )?; + let document = read_pending_proposal_document(&pending)?; + ensure_proposal_window_opened_audit( + home, + document.scheduled().context("scheduled fixture")?, + now, + )?; + let workspace = home.join("familiars/sage"); + let original = match failure { + "protected-promotion" => { + let ward_path = workspace.join("ward.toml"); + let config = std::fs::read_to_string(&ward_path)? + .replace( + "protected_surface = [\"SOUL.md\"]", + "protected_surface = [\"SOUL.md\", \"reviewed/skill.md\"]", + ) + .replace( + "path = \"reviewed/\"\ntier = 1", + "path = \"reviewed/skill.md\"\ntier = 0", + ); + std::fs::write(&ward_path, config)?; + workspace.join("reviewed/skill.md") + } + #[cfg(unix)] + "blocked-symlink" => { + std::fs::rename( + workspace.join("reviewed"), + workspace.join("reviewed-before"), + )?; + std::os::unix::fs::symlink(outside.path(), workspace.join("reviewed"))?; + workspace.join("reviewed-before/skill.md") + } + _ => unreachable!(), + }; + let config = ward::WardConfig::load(&workspace)?.context("valid changed Ward")?; + let ward = ward::Ward::new(&workspace, config)?; + let adjudication = ward.evaluate(&ward::Proposal { + targets: vec!["reviewed/skill.md".to_string()], + authorization: authorization_from_writer(&document.pending().writer), + }); + match failure { + "protected-promotion" => { + assert!(!protected_proposal_targets(&adjudication, &ward)?.is_empty()); + } + #[cfg(unix)] + "blocked-symlink" => { + assert!(adjudication.is_blocked()); + assert!(protected_proposal_targets(&adjudication, &ward)?.is_empty()); + } + _ => unreachable!(), + } + std::fs::remove_file(workspace.join("SOUL.md"))?; + std::fs::create_dir(workspace.join("SOUL.md"))?; + + assert_eq!( + process_due_threads_proposals(home)?, + 1, + "{failure}: {}", + std::fs::read_to_string(crate::daemon::daemon_recovery_log_path(home)) + .unwrap_or_default() + ); + assert!(!pending.exists(), "{failure}"); + assert_eq!(process_due_threads_proposals(home)?, 0, "{failure}"); + assert_eq!(std::fs::read_to_string(original)?, "before", "{failure}"); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let detail: String = conn.query_row( + "SELECT detail FROM ward_audit WHERE proposal_id = ?1 + AND event_type = 'proposal_rejected'", + [&proposal_id], + |row| row.get(0), + )?; + let close: coven_threads_core::ProposalWindowCloseAuditDetail = + serde_json::from_str(&detail)?; + assert_eq!( + close.reason, + coven_threads_core::WindowCloseReason::RevalidationFailed, + "{failure}" + ); + assert_eq!(close.replay_hash_matched, Some(false), "{failure}"); + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) FROM ward_audit WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 1, "{failure}"); + let reservations: i64 = conn.query_row( + "SELECT COUNT(*) FROM coven_ward_audit_reservations", + [], + |row| row.get(0), + )?; + assert_eq!(reservations, 0, "{failure}"); + Ok(()) + } + + #[test] + fn threads_scheduled_human_replay_rejection_has_no_window_close() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::HumanApproval, + time::OffsetDateTime::now_utc(), + )?; + std::fs::write(home.join("familiars/sage/reviewed/skill.md"), "concurrent")?; + let decision_body = scheduled_decision_body(home, &proposal_id, None)?; + + let response = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some(&decision_body), + )?; + + assert_eq!(response.status, 409, "got {}", response.body); + assert!(!pending.exists()); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let detail: Option = conn.query_row( + "SELECT detail + FROM ward_audit + WHERE proposal_id = ?1 AND event_type = 'proposal_rejected'", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(detail, None); + Ok(()) + } + + #[test] + fn threads_scheduled_rejects_live_promotion_to_protected_tier() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let veto = coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, + time::OffsetDateTime::now_utc(), + )?; + assert_eq!(process_due_threads_proposals(home)?, 0); + let ward_path = home.join("familiars/sage/ward.toml"); + let ward = std::fs::read_to_string(&ward_path)? + .replace( + "protected_surface = [\"SOUL.md\"]", + "protected_surface = [\"SOUL.md\", \"reviewed/skill.md\"]", + ) + .replace( + "path = \"reviewed/\"\ntier = 1", + "path = \"reviewed/skill.md\"\ntier = 0", + ); + std::fs::write(&ward_path, ward)?; + + assert_eq!(process_due_threads_proposals(home)?, 1); + assert!(!pending.exists()); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "before" + ); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let (decision, detail): (String, String) = conn.query_row( + "SELECT decision, detail FROM ward_audit + WHERE proposal_id = ?1 AND event_type = 'proposal_rejected'", + [&proposal_id], + |row| Ok((row.get(0)?, row.get(1)?)), + )?; + assert_eq!(decision, "protected-target-not-proposable"); + let close: coven_threads_core::ProposalWindowCloseAuditDetail = + serde_json::from_str(&detail)?; + assert_eq!( + close.reason, + coven_threads_core::WindowCloseReason::RevalidationFailed + ); + assert_eq!(close.replay_hash_matched, Some(false)); + assert_eq!( + close.rationale.as_deref(), + Some("protected-target-not-proposable") + ); + + let retry = handle_request_with_body( + "POST", + &format!("/api/v1/threads/proposals/{proposal_id}/approve"), + home, + None, + Some("{}"), + )?; + assert_eq!(retry.status, 409, "got {}", retry.body); + let terminal_count: i64 = conn.query_row( + "SELECT COUNT(*) FROM ward_audit + WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| row.get(0), + )?; + assert_eq!(terminal_count, 1); + Ok(()) + } + + #[test] + fn threads_scheduler_first_pass_protected_rejection_records_window_close_detail() -> Result<()> + { + let temp = tempfile::tempdir()?; + let home = temp.path(); + #[cfg(feature = "threads-test-clock")] + let staged_at = + seed_deterministic_threads_clock(home, "fixture-cap", "2026-09-09T10:00:00Z")?; + #[cfg(not(feature = "threads-test-clock"))] + let staged_at = time::OffsetDateTime::now_utc(); + let veto = coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, + staged_at, + )?; + let ward_path = home.join("familiars/sage/ward.toml"); let ward = std::fs::read_to_string(&ward_path)? .replace( "protected_surface = [\"SOUL.md\"]", @@ -32329,6 +33815,109 @@ tier = 0 Ok(()) } + #[test] + fn threads_scheduler_revalidates_old_window_instead_of_expiring_it() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let veto = coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ); + let (_, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, + time::OffsetDateTime::now_utc() - time::Duration::days(31), + )?; + + assert_eq!(process_due_threads_proposals(home)?, 1); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "after" + ); + + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let terminal: (String, String, String) = conn.query_row( + "SELECT event_type, decision, detail + FROM ward_audit + WHERE proposal_id = ?1 + AND event_type IN ('proposal_approved', 'proposal_rejected', 'proposal_vetoed')", + [&proposal_id], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + assert_eq!(terminal.0, "proposal_approved"); + assert_eq!(terminal.1, "approved"); + let detail: coven_threads_core::ProposalApprovalAuditDetail = + serde_json::from_str(&terminal.2)?; + let close = detail + .window_close + .expect("old delayed proposal must close through replay"); + assert_eq!(close.reason, coven_threads_core::WindowCloseReason::Applied); + assert_eq!(close.replay_hash_matched, Some(true)); + Ok(()) + } + + #[test] + fn threads_scheduler_terminalizes_unknown_review_kind_after_opening_window() -> Result<()> { + let temp = tempfile::tempdir()?; + let home = temp.path(); + let veto = coven_threads_core::VetoWindow::new( + std::time::Duration::from_secs(300), + std::time::Duration::from_secs(60), + ); + let (pending, proposal_id) = stage_scheduled_reviewed_edit( + home, + coven_threads_core::ApprovalPath::FamiliarCoherence { veto }, + time::OffsetDateTime::now_utc() - time::Duration::minutes(10), + )?; + let mut staged: Value = serde_json::from_slice(&std::fs::read(&pending)?)?; + staged["reviewKind"] = json!("automatic"); + std::fs::write(&pending, serde_json::to_vec_pretty(&staged)?)?; + + assert_eq!(process_due_threads_proposals(home)?, 1); + assert!( + !pending.exists(), + "terminal corrupt proposal must be consumed" + ); + assert_eq!( + std::fs::read_to_string(home.join("familiars/sage/reviewed/skill.md"))?, + "before" + ); + let conn = store::open_store(&home.join("coven.sqlite3"))?; + let rows: Vec<(String, Option)> = { + let mut statement = conn.prepare( + "SELECT event_type, detail + FROM ward_audit + WHERE proposal_id = ?1 + ORDER BY id", + )?; + let rows = statement + .query_map([&proposal_id], |row| Ok((row.get(0)?, row.get(1)?)))? + .collect::>()?; + rows + }; + assert_eq!(rows.len(), 2); + assert_eq!(rows[0].0, "proposal_window_opened"); + assert_eq!(rows[1].0, "proposal_rejected"); + let close: coven_threads_core::ProposalWindowCloseAuditDetail = + serde_json::from_str(rows[1].1.as_deref().context("terminal close detail")?)?; + assert_eq!( + close.reason, + coven_threads_core::WindowCloseReason::RevalidationFailed + ); + assert_eq!(close.replay_hash_matched, Some(false)); + let recovery_log = std::fs::read_to_string(crate::daemon::daemon_recovery_log_path(home)) + .unwrap_or_default(); + assert!( + !recovery_log.contains("did not terminalize"), + "terminal 409 closes must not be logged as retriable: {recovery_log}" + ); + assert!( + !recovery_log.contains("retained for retry"), + "terminal 409 closes must not be logged as retained: {recovery_log}" + ); + Ok(()) + } + #[test] fn threads_scheduler_batches_advance_without_starving_later_due_proposals() -> Result<()> { let temp = tempfile::tempdir()?; @@ -33831,7 +35420,7 @@ tier = 0 } #[test] - fn threads_approve_recovery_preserves_claim_if_ward_disappears() -> Result<()> { + fn threads_approve_recovery_quarantines_claim_if_ward_disappears() -> Result<()> { let temp = tempfile::tempdir()?; let home = temp.path(); let (_, proposal_id) = stage_pending_protected_edit(home)?; @@ -33851,6 +35440,7 @@ tier = 0 assert!(interrupted.is_err()); let claim = find_pending_decision_claim(home, &proposal_id, "approve") .expect("interrupted approval leaves a recovery claim"); + let claim_bytes = std::fs::read(&claim)?; std::fs::remove_file(workspace.join("ward.toml"))?; let retry = handle_request_with_body( @@ -33863,11 +35453,17 @@ tier = 0 assert_eq!(retry.status, 409, "got {}", retry.body); let body: Value = serde_json::from_str(&retry.body)?; - assert_eq!(body["why"], "ward-not-configured"); + assert_eq!(body["why"], "proposal-recovery-ward-unavailable"); + assert_eq!(body["terminal"], false); + assert_eq!(body["manualRecoveryRequired"], true); assert!( - claim.exists(), - "recovery claim must not downgrade to pending" + !claim.exists(), + "unverifiable claim must leave the retry queue" ); + let quarantined = std::fs::read_dir(home.join("pending/quarantine"))? + .collect::>>()?; + assert_eq!(quarantined.len(), 1); + assert_eq!(std::fs::read(quarantined[0].path())?, claim_bytes); let conn = store::open_store(&home.join("coven.sqlite3"))?; let terminal_count: i64 = conn.query_row( "SELECT COUNT(*) @@ -33885,7 +35481,7 @@ tier = 0 } #[test] - fn threads_approve_recovery_preserves_claim_if_familiar_disappears() -> Result<()> { + fn threads_approve_recovery_quarantines_claim_if_familiar_disappears() -> Result<()> { let temp = tempfile::tempdir()?; let home = temp.path(); let (_, proposal_id) = stage_pending_protected_edit(home)?; @@ -33905,6 +35501,7 @@ tier = 0 assert!(interrupted.is_err()); let claim = find_pending_decision_claim(home, &proposal_id, "approve") .expect("interrupted approval leaves a recovery claim"); + let claim_bytes = std::fs::read(&claim)?; std::fs::remove_file(home.join("familiars.toml"))?; let retry = handle_request_with_body( @@ -33917,11 +35514,17 @@ tier = 0 assert_eq!(retry.status, 409, "got {}", retry.body); let body: Value = serde_json::from_str(&retry.body)?; - assert_eq!(body["why"], "proposal-familiar-missing"); + assert_eq!(body["why"], "proposal-recovery-familiar-unavailable"); + assert_eq!(body["terminal"], false); + assert_eq!(body["manualRecoveryRequired"], true); assert!( - claim.exists(), - "missing familiar must leave the recovery claim intact" + !claim.exists(), + "unverifiable claim must leave the retry queue" ); + let quarantined = std::fs::read_dir(home.join("pending/quarantine"))? + .collect::>>()?; + assert_eq!(quarantined.len(), 1); + assert_eq!(std::fs::read(quarantined[0].path())?, claim_bytes); let conn = store::open_store(&home.join("coven.sqlite3"))?; let terminal_count: i64 = conn.query_row( "SELECT COUNT(*) @@ -34267,6 +35870,18 @@ tier = 0 assert!(first.is_err(), "failpoint must interrupt pending cleanup"); let claim = find_pending_decision_claim(home, &proposal_id, "reject") .expect("committed rejection leaves its claim until recovery"); + let store_path = home.join("coven.sqlite3"); + let conn = store::open_store(&store_path)?; + let reservation_bytes = proposal_decision_audit_reservation_bytes(&conn, "approve")?; + store::WardAuditReservation::acquire( + &conn, + &store_path, + format!("proposal:{proposal_id}:approve"), + "injected-stale-terminal-reservation", + reservation_bytes, + )? + .preserve()?; + drop(conn); let retry = handle_request_with_body( "POST", @@ -34287,6 +35902,17 @@ tier = 0 |row| row.get(0), )?; assert_eq!(rejected_count, 1); + let reservations: i64 = conn.query_row( + "SELECT COUNT(*) + FROM coven_ward_audit_reservations + WHERE token IN (?1, ?2)", + rusqlite::params![ + format!("proposal:{proposal_id}:approve"), + format!("proposal:{proposal_id}:reject"), + ], + |row| row.get(0), + )?; + assert_eq!(reservations, 0); Ok(()) } diff --git a/crates/coven-cli/tests/fixtures/threads_daemon.rs b/crates/coven-cli/tests/fixtures/threads_daemon.rs index 55afb9a54..1be4a5783 100644 --- a/crates/coven-cli/tests/fixtures/threads_daemon.rs +++ b/crates/coven-cli/tests/fixtures/threads_daemon.rs @@ -19,6 +19,17 @@ pub const PRINCIPAL_FINGERPRINT: &str = "fpr-e2e-synthetic"; const LIFECYCLE_TIMEOUT: Duration = Duration::from_secs(15); pub fn run_journey(journey: impl FnOnce(&mut ThreadsFixture) -> Result<()>) -> Result<()> { + // These are authority journeys, not concurrent cold-start throughput tests. + // Keep native Windows process admission isolated, as in + // windows_daemon_lifecycle; transport-only tests remain concurrent. + #[cfg(windows)] + let _admission = { + static ADMISSION: std::sync::OnceLock> = std::sync::OnceLock::new(); + ADMISSION + .get_or_init(|| std::sync::Mutex::new(())) + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + }; let mut fixture = ThreadsFixture::start()?; let result = journey(&mut fixture); match (result, fixture.stop_daemon()) { @@ -43,6 +54,8 @@ pub struct ThreadsFixture { pub workspace: PathBuf, daemon_pid: Option, stopped: bool, + #[cfg(windows)] + owned_daemon: Option, } impl ThreadsFixture { @@ -92,6 +105,8 @@ forbidden = ["(?i)ignore previous"] workspace, daemon_pid: None, stopped: true, + #[cfg(windows)] + owned_daemon: None, }; fixture.start_daemon()?; Ok(fixture) @@ -99,13 +114,128 @@ forbidden = ["(?i)ignore previous"] pub fn start_daemon(&mut self) -> Result<()> { anyhow::ensure!(self.stopped, "fixture daemon must be stopped before start"); - self.stopped = false; - self.daemon_command("start")?; - self.wait_for_health()?; + #[cfg(windows)] + { + self.start_owned_windows_daemon()?; + } + #[cfg(unix)] + { + self.stopped = false; + self.daemon_command("start")?; + self.wait_for_health()?; + } self.daemon_pid = Some(self.read_daemon_pid()?); Ok(()) } + #[cfg(windows)] + fn start_owned_windows_daemon(&mut self) -> Result<()> { + let stdout = fs::OpenOptions::new() + .create(true) + .append(true) + .open(self.coven_home.join("fixture-serve.stdout.log"))?; + let stderr = fs::OpenOptions::new() + .create(true) + .append(true) + .open(self.coven_home.join("fixture-serve.stderr.log"))?; + // Exercise the real server, not the separate two-second launcher SLA. + let child = self + .daemon_command_builder("serve") + .stdin(std::process::Stdio::null()) + .stdout(stdout) + .stderr(stderr) + .spawn()?; + let pid = child.id(); + self.owned_daemon = Some(child); + self.daemon_pid = Some(pid); + self.stopped = false; + self.wait_for_owned_windows_health(pid).with_context(|| { + format!( + "foreground fixture daemon {pid} failed readiness; stdout: {:?}; stderr: {:?}", + fs::read_to_string(self.coven_home.join("fixture-serve.stdout.log")), + fs::read_to_string(self.coven_home.join("fixture-serve.stderr.log")), + ) + }) + } + + #[cfg(windows)] + fn wait_for_owned_windows_health(&mut self, pid: u32) -> Result<()> { + let started = Instant::now(); + let deadline = started + LIFECYCLE_TIMEOUT; + let pipe = coven_client::owner_only_windows_pipe_name(&self.coven_home)?; + let mut last_pending = None; + loop { + let child = self + .owned_daemon + .as_mut() + .context("owned fixture process")?; + if let Some(status) = child.try_wait()? { + anyhow::bail!("fixture daemon exited before readiness: {status}"); + } + anyhow::ensure!( + Instant::now() < deadline, + "fixture daemon did not become healthy after {:?}: {last_pending:?}", + started.elapsed(), + ); + let probe_deadline = deadline.min(Instant::now() + Duration::from_millis(250)); + match coven_client::probe_windows_daemon_health_with_identity_until( + &pipe, + probe_deadline, + ) { + Ok(Some(probe)) => { + let body: Value = serde_json::from_slice(&probe.body)?; + anyhow::ensure!( + probe.status == 200 && body["ok"] == true, + "fixture daemon returned invalid health: HTTP {} {body}", + probe.status, + ); + anyhow::ensure!( + probe.server_pid == pid + && body["daemon"]["pid"] == pid + && body["daemon"]["socket"] == pipe, + "fixture health did not identify its owned process and pipe: {body}", + ); + anyhow::ensure!( + Instant::now() < deadline, + "fixture daemon readiness exceeded its deadline after {:?}", + started.elapsed(), + ); + return Ok(()); + } + Ok(None) => {} + Err(error) if is_pending_windows_startup_error(&error) => { + last_pending = Some(error.to_string()); + } + Err(error) => return Err(error.into()), + } + thread::sleep( + Duration::from_millis(25).min(deadline.saturating_duration_since(Instant::now())), + ); + } + } + + #[cfg(windows)] + fn reap_owned_daemon(&mut self, terminate: bool) -> Result<()> { + let Some(child) = self.owned_daemon.as_mut() else { + return Ok(()); + }; + if child.try_wait()?.is_none() && terminate { + child.kill()?; + } + let started = Instant::now(); + while child.try_wait()?.is_none() { + anyhow::ensure!( + started.elapsed() < LIFECYCLE_TIMEOUT, + "owned fixture daemon {} was not reaped after {:?}", + child.id(), + started.elapsed(), + ); + thread::sleep(Duration::from_millis(25)); + } + self.owned_daemon = None; + Ok(()) + } + pub fn request( &self, method: &'static str, @@ -171,7 +301,7 @@ forbidden = ["(?i)ignore previous"] .map_err(Into::into) } - fn daemon_command(&self, operation: &str) -> Result { + fn daemon_command_builder(&self, operation: &str) -> Command { let mut command = Command::new(env!("CARGO_BIN_EXE_coven")); command.env_clear(); // Keep only OS launch necessities, not developer credentials or test controls. @@ -180,7 +310,7 @@ forbidden = ["(?i)ignore previous"] command.env(key, value); } } - let output = command + command .args(["daemon", operation]) .current_dir(&self.workspace) .env("COVEN_HOME", &self.coven_home) @@ -191,13 +321,22 @@ forbidden = ["(?i)ignore previous"] .env("XDG_CACHE_HOME", self._temp.path().join("cache")) .env("TMPDIR", self._temp.path()) .env("TEMP", self._temp.path()) - .env("TMP", self._temp.path()) - .output()?; + .env("TMP", self._temp.path()); + command + } + + fn daemon_command(&self, operation: &str) -> Result { + let started = Instant::now(); + let output = self.daemon_command_builder(operation).output()?; anyhow::ensure!( output.status.success(), - "daemon {operation} failed\nstdout:\n{}\nstderr:\n{}", + "daemon {operation} failed after {:?}\nstdout:\n{}\nstderr:\n{}\n\ + fixture status: {:?}\nfixture recovery log: {:?}", + started.elapsed(), String::from_utf8_lossy(&output.stdout), - String::from_utf8_lossy(&output.stderr) + String::from_utf8_lossy(&output.stderr), + fs::read_to_string(self.coven_home.join("daemon.json")), + fs::read_to_string(self.coven_home.join("daemon-recovery.log")), ); Ok(output) } @@ -217,6 +356,7 @@ forbidden = ["(?i)ignore previous"] .or_else(|| self.read_daemon_pid().ok().filter(|pid| pid_is_alive(*pid))) } + #[cfg(unix)] fn wait_for_health(&self) -> Result<()> { let started = Instant::now(); loop { @@ -242,6 +382,8 @@ forbidden = ["(?i)ignore previous"] } let pid = self.current_daemon_pid(); self.daemon_command("stop")?; + #[cfg(windows)] + self.reap_owned_daemon(false)?; let started = Instant::now(); while pid.is_some_and(pid_is_alive) || self.coven_home.join("daemon.json").exists() @@ -263,8 +405,16 @@ forbidden = ["(?i)ignore previous"] let before = self .current_daemon_pid() .context("restart requires a live daemon")?; - self.daemon_command("restart")?; - self.wait_for_health()?; + #[cfg(windows)] + { + self.stop_daemon()?; + self.start_daemon()?; + } + #[cfg(unix)] + { + self.daemon_command("restart")?; + self.wait_for_health()?; + } let after = self.read_daemon_pid()?; self.daemon_pid = Some(after); anyhow::ensure!( @@ -288,6 +438,13 @@ impl Drop for ThreadsFixture { if let Err(error) = self.stop_daemon() { eprintln!("threads fixture graceful shutdown failed: {error:#}"); } + #[cfg(windows)] + if self.owned_daemon.is_some() { + if let Err(error) = self.reap_owned_daemon(true) { + eprintln!("failed terminating owned fixture daemon: {error:#}"); + } + return; + } if let Some(pid) = pid.filter(|pid| pid_is_alive(*pid)) { eprintln!("threads fixture fallback terminating its daemon pid {pid}"); if let Err(error) = self.terminate_daemon_process(pid) { @@ -297,6 +454,17 @@ impl Drop for ThreadsFixture { } } +#[cfg(any(windows, test))] +fn is_pending_windows_startup_error(error: &coven_client::ClientError) -> bool { + match error { + coven_client::ClientError::Io { source, .. } => source.kind() == io::ErrorKind::TimedOut, + coven_client::ClientError::InvalidHttpResponse(message) => { + message == coven_client::EMPTY_RESPONSE_TIMEOUT_MESSAGE + } + _ => false, + } +} + #[cfg(unix)] fn pid_is_alive(pid: u32) -> bool { Command::new("kill") @@ -532,6 +700,29 @@ fn windows_pipe_connected(stream: &WindowsPipe) -> io::Result<()> { #[cfg(test)] mod transport_tests { use super::*; + + #[test] + fn startup_wait_retries_only_pending_transport_not_identity_or_protocol_errors() { + use coven_client::ClientError; + assert!(is_pending_windows_startup_error(&ClientError::Io { + operation: coven_client::WINDOWS_CONNECT_OPERATION, + source: io::ErrorKind::TimedOut.into(), + })); + assert!(is_pending_windows_startup_error( + &ClientError::InvalidHttpResponse(coven_client::EMPTY_RESPONSE_TIMEOUT_MESSAGE.into()) + )); + for error in [ + ClientError::DaemonInstanceChanged, + ClientError::Discovery("wrong owner".into()), + ClientError::InvalidHttpResponse("partial response timed out".into()), + ClientError::Io { + operation: coven_client::WINDOWS_CONNECT_OPERATION, + source: io::ErrorKind::PermissionDenied.into(), + }, + ] { + assert!(!is_pending_windows_startup_error(&error), "{error}"); + } + } use std::collections::VecDeque; struct ScriptedStream(VecDeque<&'static [u8]>); diff --git a/crates/coven-cli/tests/threads_terminal_recovery.rs b/crates/coven-cli/tests/threads_terminal_recovery.rs new file mode 100644 index 000000000..e5bf91122 --- /dev/null +++ b/crates/coven-cli/tests/threads_terminal_recovery.rs @@ -0,0 +1,404 @@ +//! Bounded real-daemon regressions for #932 / #886. +//! Public Tier-1 coherence intake is real; opened-window rows below are explicitly +//! synthetic, inconsistent LEGACY recovery history, not supported scheduled +//! intake/publication. Supersession, the reviewed-pin/human gate, and full #886 +//! closure remain deferred. + +use std::fs; +use std::path::PathBuf; +use std::time::{Duration, Instant}; + +use anyhow::{Context, Result}; +use coven_threads_core::{ + ProposalApprovalAuditDetail, ProposalWindowAuditDetail, ProposalWindowCloseAuditDetail, + WindowCloseReason, +}; +use rusqlite::{params, Connection, OpenFlags}; +use serde_json::{json, Value}; +use time::{format_description::well_known::Rfc3339, OffsetDateTime}; + +#[path = "fixtures/threads_daemon.rs"] +mod threads_daemon; +use threads_daemon::{run_journey, ThreadsFixture, FAMILIAR_ID}; + +const EDITS: &str = "/api/v1/familiars/sage/edits"; +const PROPOSALS: &str = "/api/v1/threads/proposals"; +const TARGET: &str = "reviewed/skill.md"; +const BEFORE: &[u8] = b"before"; +const AFTER: &str = "after"; +// Hang guard, not a promptness claim: several scheduler intervals tolerate +// loaded native CI runners. Only the durable terminal row establishes success. +const RECOVERY_TIMEOUT: Duration = Duration::from_secs(120); + +#[test] +fn public_coherence_approval_has_no_fabricated_window_close() -> Result<()> { + run_journey(|fixture| { + let proposal = stage_public_coherence_edit(fixture)?; + let submitted = audit_rows(fixture, &proposal)?; + fixture.restart_daemon()?; + assert_eq!(audit_rows(fixture, &proposal)?, submitted); + assert!( + proposal.pending.is_file(), + "human review must remain pending" + ); + assert_eq!(fs::read(fixture.workspace.join(TARGET))?, BEFORE); + let approved = fixture.request( + "POST", + &proposal.approval_path(), + Some(&proposal.decision_body()), + )?; + assert_eq!(approved.status, 200, "{approved:?}"); + assert_eq!(approved.body["decision"], "approved", "{approved:?}"); + assert_eq!(approved.body["reviewKind"], "coherence", "{approved:?}"); + assert_eq!(fs::read(fixture.workspace.join(TARGET))?, AFTER.as_bytes()); + assert_consumed(fixture, &proposal)?; + + let rows = audit_rows(fixture, &proposal)?; + let terminal = only_terminal(&rows); + assert_eq!(terminal.event, "proposal_approved", "{rows:?}"); + let detail: ProposalApprovalAuditDetail = + serde_json::from_str(terminal.detail.as_deref().context("approval detail")?)?; + assert!(detail.window_close.is_none(), "{detail:?}"); + assert!( + rows.iter().all(|row| row.event != "proposal_window_opened"), + "ordinary human approval invented an opened window: {rows:?}" + ); + + fixture.restart_daemon()?; + let replay = fixture.request( + "POST", + &proposal.approval_path(), + Some(&proposal.decision_body()), + )?; + assert_eq!(replay.status, 200, "{replay:?}"); + assert_eq!(replay.body["idempotent"], true, "{replay:?}"); + assert_eq!(audit_rows(fixture, &proposal)?, rows); + assert_eq!(fs::read(fixture.workspace.join(TARGET))?, AFTER.as_bytes()); + assert_consumed(fixture, &proposal) + }) +} + +#[test] +fn synthetic_legacy_opened_history_is_rejected_once_across_restart() -> Result<()> { + assert_legacy_recovery(None) +} + +#[test] +fn synthetic_legacy_opened_history_with_missing_ward_is_rejected_once() -> Result<()> { + assert_legacy_recovery(Some(MissingPrerequisite::Ward)) +} + +#[test] +fn synthetic_legacy_opened_history_with_missing_familiar_is_rejected_once() -> Result<()> { + assert_legacy_recovery(Some(MissingPrerequisite::Familiar)) +} + +enum MissingPrerequisite { + Ward, + Familiar, +} + +fn assert_legacy_recovery(missing: Option) -> Result<()> { + run_journey(|fixture| { + let proposal = stage_public_coherence_edit(fixture)?; + let opened = seed_synthetic_legacy_opened_history(fixture, &proposal)?; + match missing { + Some(MissingPrerequisite::Ward) => { + fs::remove_file(fixture.workspace.join("ward.toml"))?; + } + Some(MissingPrerequisite::Familiar) => { + fs::write(fixture.coven_home.join("familiars.toml"), "")?; + } + None => {} + } + fixture.start_daemon()?; + wait_for_startup_terminal(fixture, &proposal)?; + let rows = assert_revalidation_rejected(fixture, &proposal, &opened)?; + + for after_restart in [false, true] { + if after_restart { + fixture.restart_daemon()?; + } + let replay = fixture.request( + "POST", + &proposal.approval_path(), + Some(&proposal.decision_body()), + )?; + assert_eq!(replay.status, 409, "{replay:?}"); + assert_eq!(replay.body["why"], "proposal-already-decided", "{replay:?}"); + assert_eq!(replay.body["eventType"], "proposal_rejected", "{replay:?}"); + let same_decision = fixture.request( + "POST", + &format!("{PROPOSALS}/{}/reject", proposal.id), + Some(&proposal.decision_body()), + )?; + assert_eq!(same_decision.status, 200, "{same_decision:?}"); + assert_eq!(same_decision.body["idempotent"], true, "{same_decision:?}"); + assert_eq!( + assert_revalidation_rejected(fixture, &proposal, &opened)?, + rows + ); + } + Ok(()) + }) +} + +fn wait_for_startup_terminal(fixture: &ThreadsFixture, proposal: &PublicProposal) -> Result<()> { + let started = Instant::now(); + loop { + let rows = audit_rows(fixture, proposal)?; + if rows.iter().any(AuditRow::is_terminal) { + return Ok(()); + } + if started.elapsed() >= RECOVERY_TIMEOUT { + anyhow::bail!( + "startup recovery did not terminalize legacy opened history without a decision POST \ + after {:?}; pending exists: {}; target bytes: {:?}; audit rows: {rows:?}; \ + recovery log: {:?}", + started.elapsed(), + proposal.pending.exists(), + fs::read(fixture.workspace.join(TARGET))?, + fs::read_to_string(fixture.coven_home.join("daemon-recovery.log")), + ); + } + // Polling backoff only; elapsed time never establishes recovery success. + std::thread::sleep(Duration::from_millis(50)); + } +} + +struct PublicProposal { + id: String, + pending: PathBuf, + revision: String, +} + +impl PublicProposal { + fn approval_path(&self) -> String { + format!("{PROPOSALS}/{}/approve", self.id) + } + + fn decision_body(&self) -> Value { + json!({"expectedRevision": self.revision, "note": "Synthetic human coherence review"}) + } +} + +fn stage_public_coherence_edit(fixture: &ThreadsFixture) -> Result { + fs::create_dir_all(fixture.workspace.join("reviewed"))?; + fs::write(fixture.workspace.join(TARGET), BEFORE)?; + let staged = fixture.request( + "POST", + EDITS, + Some(&json!({"edits": [{"target": TARGET, "contents": AFTER}]})), + )?; + assert_eq!(staged.status, 202, "{staged:?}"); + assert_eq!(staged.body["disposition"], "staged", "{staged:?}"); + assert_eq!(staged.body["reviewKind"], "coherence", "{staged:?}"); + let id = staged.body["proposalId"] + .as_str() + .context("proposal id")? + .to_owned(); + let pending = PathBuf::from( + staged.body["pendingPath"] + .as_str() + .context("pending path")?, + ); + assert!( + pending.is_file(), + "public intake did not persist {pending:?}" + ); + assert_eq!(fs::read(fixture.workspace.join(TARGET))?, BEFORE); + let detail = fixture.request("GET", &format!("{PROPOSALS}/{id}"), None)?; + assert_eq!(detail.status, 200, "{detail:?}"); + assert_eq!( + detail.body["proposal"]["reviewKind"], "coherence", + "{detail:?}" + ); + let revision = detail.body["proposal"]["proposalRevision"] + .as_str() + .context("public proposal revision")? + .to_owned(); + let proposal = PublicProposal { + id, + pending, + revision, + }; + let rows = audit_rows(fixture, &proposal)?; + assert_eq!( + rows.iter() + .filter(|row| row.event == "proposal_submitted") + .count(), + 1, + "public intake must submit the proposal: {rows:?}" + ); + assert!( + rows.iter() + .all(|row| row.event != "proposal_window_opened" && !row.is_terminal()), + "coherence intake unexpectedly opened or decided a window: {rows:?}" + ); + Ok(proposal) +} + +fn seed_synthetic_legacy_opened_history( + fixture: &mut ThreadsFixture, + proposal: &PublicProposal, +) -> Result { + fixture.stop_daemon()?; + let pending_before = fs::read(&proposal.pending)?; + let document: Value = serde_json::from_slice(&pending_before)?; + assert!(document.get("decisionRequest").is_none(), "{document}"); + assert!(document.get("decisionState").is_none(), "{document}"); + let now = OffsetDateTime::now_utc(); + let detail = serde_json::to_string(&ProposalWindowAuditDetail { + approval_path_label: "familiar_review".to_owned(), + deadline: now + time::Duration::minutes(5), + earliest_close: now + time::Duration::minutes(1), + evidence_replay_hash_hex: "00".repeat(32), + affected_regions: Vec::new(), + })?; + // The ONLY writable store access: append fixture-only legacy history to a + // stopped daemon's existing database. Keep its real pending document and + // submission intact; never seed approved/applying authority or alter triggers. + let conn = Connection::open_with_flags( + fixture.coven_home.join("coven.sqlite3"), + OpenFlags::SQLITE_OPEN_READ_WRITE, + )?; + let inserted = conn.execute( + "INSERT INTO ward_audit ( + event_type, proposal_id, familiar_id, ward_hash, decision, detail, + files_touched, submitted_at, decided_at + ) + SELECT 'proposal_window_opened', proposal_id, familiar_id, ward_hash, + 'window-opened', ?2, files_touched, submitted_at, ?3 + FROM ward_audit + WHERE proposal_id = ?1 AND event_type = 'proposal_submitted'", + params![proposal.id, detail, now.format(&Rfc3339)?], + )?; + assert_eq!( + inserted, 1, + "legacy fixture must append exactly one opened row" + ); + drop(conn); + assert_eq!(fs::read(&proposal.pending)?, pending_before); + Ok(detail) +} + +#[derive(Debug, PartialEq, Eq)] +struct AuditRow { + id: i64, + event: String, + decision: String, + detail: Option, + decided_at: String, +} + +impl AuditRow { + fn is_terminal(&self) -> bool { + matches!( + self.event.as_str(), + "proposal_approved" | "proposal_rejected" | "proposal_vetoed" + ) + } +} + +fn audit_rows(fixture: &ThreadsFixture, proposal: &PublicProposal) -> Result> { + let conn = fixture.store()?; + let mut statement = conn.prepare( + "SELECT id, event_type, decision, detail, decided_at FROM ward_audit + WHERE proposal_id = ?1 AND familiar_id = ?2 ORDER BY id", + )?; + let rows = statement.query_map(params![proposal.id, FAMILIAR_ID], |row| { + Ok(AuditRow { + id: row.get(0)?, + event: row.get(1)?, + decision: row.get(2)?, + detail: row.get(3)?, + decided_at: row.get(4)?, + }) + })?; + rows.collect::>>() + .map_err(Into::into) +} + +fn only_terminal(rows: &[AuditRow]) -> &AuditRow { + let terminals: Vec<_> = rows.iter().filter(|row| row.is_terminal()).collect(); + assert_eq!( + terminals.len(), + 1, + "expected one durable terminal row: {rows:?}" + ); + terminals[0] +} + +fn assert_revalidation_rejected( + fixture: &ThreadsFixture, + proposal: &PublicProposal, + opened_detail: &str, +) -> Result> { + let rows = audit_rows(fixture, proposal)?; + let terminal = only_terminal(&rows); + assert_eq!(terminal.event, "proposal_rejected", "{rows:?}"); + let raw_detail = terminal + .detail + .as_deref() + .context("typed terminal detail")?; + let close: ProposalWindowCloseAuditDetail = serde_json::from_str(raw_detail)?; + assert_eq!( + close.reason, + WindowCloseReason::RevalidationFailed, + "{close:?}" + ); + assert_eq!(close.replay_hash_matched, Some(false), "{close:?}"); + assert_eq!( + serde_json::from_str::(raw_detail)?["reason"], + "revalidation_failed" + ); + let opened: Vec<_> = rows + .iter() + .filter(|row| row.event == "proposal_window_opened") + .collect(); + assert_eq!( + opened.len(), + 1, + "recovery invented another opened row: {rows:?}" + ); + assert_eq!(opened[0].detail.as_deref(), Some(opened_detail)); + assert!( + rows.iter() + .all(|row| row.event != "apply_audit" && row.decision != "proposal-apply-intent"), + "rejected history gained write authority: {rows:?}" + ); + assert_eq!(fs::read(fixture.workspace.join(TARGET))?, BEFORE); + assert_eq!(fs::read(fixture.workspace.join("SOUL.md"))?, b"# Sage\n"); + assert_consumed(fixture, proposal)?; + Ok(rows) +} + +fn assert_consumed(fixture: &ThreadsFixture, proposal: &PublicProposal) -> Result<()> { + assert!( + !proposal.pending.exists(), + "terminal proposal remains pending" + ); + let pending_dir = fixture.coven_home.join("pending"); + if pending_dir.exists() { + let entries = fs::read_dir(&pending_dir)? + .map(|entry| entry.map(|entry| entry.file_name())) + .collect::>>()?; + // Persistent directory bookkeeping is not a pending proposal artifact. + assert!( + entries + .iter() + .all(|name| name == ".quota.lock" || name == ".scheduler-cursor"), + "terminal proposal left pending artifacts: {entries:?}" + ); + } + let listed = fixture.request("GET", PROPOSALS, None)?; + assert_eq!(listed.status, 200, "{listed:?}"); + assert_eq!(listed.body["proposals"], json!([]), "{listed:?}"); + let reservations: i64 = fixture.store()?.query_row( + "SELECT COUNT(*) FROM coven_ward_audit_reservations", + [], + |row| row.get(0), + )?; + assert_eq!(reservations, 0, "terminal proposal retained audit capacity"); + Ok(()) +} diff --git a/docs/design/threads-terminal-recovery.md b/docs/design/threads-terminal-recovery.md new file mode 100644 index 000000000..42269fb1f --- /dev/null +++ b/docs/design/threads-terminal-recovery.md @@ -0,0 +1,86 @@ +--- +summary: "Keep opened Threads veto windows tied to one typed terminal outcome." +read_when: + - Changing Threads proposal decisions or startup recovery +title: "Threads terminal recovery" +description: "Typed close evidence, replay failures, and the boundary between rejection and quarantine." +source_adjacent_reason: "Documents the terminal audit guard and proposal recovery policy." +--- + +# Threads terminal recovery + +An opened Threads veto window requires one typed terminal outcome, so you can +distinguish an applied change, a veto, and a failed revalidation in the audit. + +```sh +cargo test --locked -p coven-cli --test threads_terminal_recovery +``` + +This target runs the actual daemon. Its historical-window fixtures seed audit +state only while the daemon is stopped. They exercise recovery, not a supported +scheduled-proposal publication route. + +On Windows, the fixture owns a foreground `coven daemon serve` child and waits +for authenticated health under a bounded startup guard. It uses the real +`daemon stop` command and replaces the process before checking restart recovery. +This covers server authority and durable recovery, not the separate two-second +CLI launcher or atomic `daemon restart` contract. Dedicated lifecycle tests +retain those contracts. Unix journeys continue to use the normal launcher. + +## Terminal evidence + +The central decision audit writer rejects duplicate terminal rows. It also +requires an opened-window row if and only if the decision carries typed +`window_close` evidence. + +| Terminal event | Typed close reason | +| --- | --- | +| `proposal_approved` | `applied` | +| `proposal_vetoed` | `vetoed` | +| `proposal_rejected` | `evidence_diverged`, `revalidation_failed`, or `superseded` | + +The `superseded` reason is part of the contract. This checkpoint does not add a +production supersession or scheduled-publication path. + +Deadline expiry starts revalidation. It does not authorize a write or replace +an opened window with `proposal_expired`. Human approval without a window +remains valid and does not fabricate close evidence. + +## Recover without granting authority + +An unapplied proposal with an existing window fails closed when its familiar, +Ward configuration, or authoritative replay becomes unavailable. Its rejection +uses `revalidation_failed` and `replay_hash_matched = false`. If live weave +construction fails, recovery retains the committed window hash as audit +context. That hash is not evidence that current bytes matched. + +Live promotion to a protected target still rejects before mutation. An +unavailable replay cannot turn that rejection into an indefinitely pending +window. Recovery preserves a previously recorded decision request rather than +switching verbs and conflicting with an interrupted approval. + +Inconsistent human-labelled pending history with an existing window must be +rejected, never normalized into approval. Once a terminal row is durable, +repeated recovery consumes leftover pending state and reservations without +appending another terminal row or applying edits again. + +## Quarantine is not a close + +An interrupted apply with unverifiable or inconsistent committed intent is +different from an unapplied proposal. Recovery quarantines that evidence rather +than asserting that no write occurred. Quarantine does **not** satisfy the typed +terminal invariant and requires separate resolution. + +If the familiar or Ward authority disappears, becomes unreadable, or cannot be +constructed, an interrupted apply leaves the automatic retry queue. Its claim +bytes and recorded apply intent remain available for manual recovery. The +response reports `terminal = false` and `manualRecoveryRequired = true`; simply +restoring the configuration does not automatically resume a quarantined apply. + +Do not delete quarantined intent or classify it as a completed rejection to +make the terminal count balance. This checkpoint also does not repair arbitrary +historical database corruption or infer missing close evidence for past writes. + +For deterministic scheduler deadlines, use the opt-in +[Threads test clock](threads-test-clock.md). Neither the clock nor seeded +historical recovery fixtures grant protected-write authority. diff --git a/docs/design/threads-test-clock.md b/docs/design/threads-test-clock.md index 3f92a92b4..ec173ee43 100644 --- a/docs/design/threads-test-clock.md +++ b/docs/design/threads-test-clock.md @@ -71,3 +71,7 @@ earliest close, and beyond the deadline. Submit proposals through supported intake, and inspect filesystem, pending, and `ward_audit` effects after each step. The fixture token controls time only; it never supplies proposal, approval, or protected-write authority. + +See [Threads terminal recovery](threads-terminal-recovery.md) for typed close +evidence and the distinction between recovering pending work and quarantining +an unverifiable interrupted apply.