Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 75 additions & 6 deletions crates/coven-cli/src/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12830,6 +12830,13 @@ fn pending_document_protected_targets(
/// Process one deterministic round-robin batch. The persistent filename cursor
/// prevents a prefix of human-only proposals from starving later due work.
pub(crate) fn process_due_threads_proposals(coven_home: &Path) -> Result<usize> {
#[cfg(feature = "threads-test-clock")]
if crate::threads_clock::fixture_mode_enabled(coven_home)? {
crate::daemon::append_daemon_recovery_log(
coven_home,
"threads_scheduler_checkpoint phase=pass-lock-wait",
);
}
let _pass_guard = threads_scheduler_pass_lock()
.lock()
.map_err(|_| anyhow::anyhow!("threads proposal scheduler lock is poisoned"))?;
Expand Down Expand Up @@ -15140,6 +15147,8 @@ struct DeterministicThreadsClockRequest {
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct DeterministicThreadsTickRequest {
capability: String,
#[serde(default)]
workers: Option<u8>,
}

#[cfg(test)]
Expand Down Expand Up @@ -17004,19 +17013,49 @@ fn deterministic_threads_tick_response(
);
}
};
let workers = request.workers.unwrap_or(1);
if !matches!(workers, 1 | 2) {
return api_error(
400,
"invalid_request",
"Deterministic Threads ticks support one or two recovery workers.",
None,
);
}
let _clock_guard = crate::threads_clock::fixture_control_lock()
.lock()
.map_err(|_| anyhow::anyhow!("deterministic Threads clock lock is poisoned"))?;
match crate::threads_clock::authorize_fixture(coven_home, &request.capability) {
Ok(snapshot) => json_response(
200,
&json!({
Ok(snapshot) => {
let worker_results = if workers == 2 {
std::thread::scope(|scope| -> Result<Vec<usize>> {
let first = scope.spawn(|| process_due_threads_proposals(coven_home));
let second = scope.spawn(|| process_due_threads_proposals(coven_home));
let first = first.join();
let second = second.join();
Ok(vec![
first.map_err(|_| {
anyhow::anyhow!("first Threads recovery worker panicked")
})??,
second.map_err(|_| {
anyhow::anyhow!("second Threads recovery worker panicked")
})??,
])
})?
} else {
vec![process_due_threads_proposals(coven_home)?]
};
let mut body = json!({
"ok": true,
"processed": process_due_threads_proposals(coven_home)?,
"processed": worker_results.iter().sum::<usize>(),
"now": snapshot.now.format(&time::format_description::well_known::Rfc3339)?,
"source": snapshot.source.as_str(),
}),
),
});
if workers == 2 {
body["workerResults"] = json!(worker_results);
}
json_response(200, &body)
}
Err(error) => map_deterministic_threads_clock_error(error),
}
}
Expand Down Expand Up @@ -39106,6 +39145,36 @@ tier = 0
Ok(())
}

#[cfg(feature = "threads-test-clock")]
#[test]
fn deterministic_threads_recovery_workers_are_bounded_and_authorized() -> Result<()> {
let temp = tempfile::tempdir()?;
let home = temp.path();
seed_deterministic_threads_clock(home, "fixture-cap", "2026-09-09T10:00:00Z")?;
for workers in [json!(0), json!(3), json!(256), json!(-1), json!("2")] {
let response = deterministic_threads_tick_response(
home,
Some(&json!({"capability": "fixture-cap", "workers": workers}).to_string()),
)?;
assert_eq!(response.status, 400, "{workers}: {}", response.body);
}
let denied = deterministic_threads_tick_response(
home,
Some(r#"{"capability":"wrong-cap","workers":2}"#),
)?;
assert_eq!(denied.status, 403, "{}", denied.body);
let inactive = deterministic_threads_tick_response(
tempfile::tempdir()?.path(),
Some(r#"{"capability":"fixture-cap","workers":2}"#),
)?;
assert_eq!(inactive.status, 404, "{}", inactive.body);
assert!(
!home.join("daemon-recovery.log").exists(),
"refused control started a recovery worker"
);
Ok(())
}

#[cfg(feature = "threads-test-clock")]
#[test]
fn deterministic_threads_clock_control_rejects_invalid_and_non_monotonic_inputs() -> Result<()>
Expand Down
174 changes: 174 additions & 0 deletions crates/coven-cli/tests/support/threads_corpus_closure_cases.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,174 @@
use super::*;

#[test]
fn retired_corpus_recovery_workers_contend_without_duplicate_apply() -> Result<()> {
let corpus = retired_ward_corpus()?;
let case = retired_review_case(&corpus)?;
let duration = case["approval"]["veto"]["duration_seconds"]
.as_i64()
.context("corpus veto duration")?;
run_clocked_journey(
"retired-corpus-recovery-worker-contention",
|home, workspace| seed_retired_review_case(home, workspace, case),
|fixture, capability| {
let staged = submit_retired_case(fixture, case)?;
let id = staged["proposalId"].as_str().context("proposal id")?;
tick_scheduler(fixture, capability)?;
fixture.restart_daemon()?;
advance_clock_to_offset(fixture, capability, duration + 1)?;
final_commit_cases::arm_final_commit_pause(fixture, capability)?;
let before = scheduler_entries(fixture)?;
thread::scope(|scope| -> Result<()> {
let first_home = fixture.coven_home.clone();
let first_payload = json!({"capability": capability, "workers": 2});
let first = scope.spawn(move || {
daemon_http_request(
&first_home,
"POST",
"/api/v1/internal/threads/test-clock/tick",
Some(&first_payload),
)
});
let paused = final_commit_cases::wait_for_final_commit_pause(fixture);
if let Err(error) = paused {
final_commit_cases::release_final_commit_pause(fixture, capability)?;
first.join().map_err(|_| anyhow::anyhow!("first recovery worker panicked"))??;
return Err(error);
}
let contended = (|| -> Result<()> {
let deadline = Instant::now() + Duration::from_secs(5);
while scheduler_entries(fixture)? < before + 2 {
anyhow::ensure!(
Instant::now() < deadline,
"second daemon recovery worker did not reach the held pass lock"
);
thread::sleep(Duration::from_millis(10));
}
assert_corpus_bytes(fixture, case, "before")?;
let terminals: i64 = fixture.store()?.query_row(
"SELECT COUNT(*) FROM ward_audit WHERE proposal_id=?1
AND event_type IN ('proposal_approved','proposal_rejected','proposal_vetoed')",
[id],
|row| row.get(0),
)?;
anyhow::ensure!(terminals == 0, "paused recovery already terminalized");
Ok(())
})();
final_commit_cases::release_final_commit_pause(fixture, capability)?;
let first = first.join().map_err(|_| anyhow::anyhow!("first recovery worker panicked"))??;
contended?;
let first_body: Value = serde_json::from_str(&first.1)?;
let mut workers: Vec<usize> = serde_json::from_value(first_body["workerResults"].clone())?;
workers.sort_unstable();
anyhow::ensure!(
first.0 == 200 && first_body["processed"] == 1 && workers == [0, 1],
"competing recovery results: {first:?}"
);
fs::write(
fixture.artifact_dir.join("recovery-workers.json"),
serde_json::to_vec_pretty(&json!({
"tick": first_body,
"both_workers_entered_before_final_commit_release": true,
}))?,
)?;
Ok(())
})?;
assert_corpus_bytes(fixture, case, "after")?;
assert_window_terminal(fixture, id, "proposal_approved", "applied", json!(true))?;
let intents: i64 = fixture.store()?.query_row(
"SELECT COUNT(*) FROM ward_audit WHERE proposal_id=?1 AND decision='proposal-apply-intent'",
[id],
|row| row.get(0),
)?;
anyhow::ensure!(intents == 1, "recovery created {intents} apply intents");
let committed_audit = ward_audit_jsonl(&fixture.coven_home.join("coven.sqlite3"))?;
fixture.restart_daemon()?;
tick_scheduler(fixture, capability)?;
assert_corpus_bytes(fixture, case, "after")?;
assert_window_terminal(fixture, id, "proposal_approved", "applied", json!(true))?;
anyhow::ensure!(
committed_audit == ward_audit_jsonl(&fixture.coven_home.join("coven.sqlite3"))?,
"restart replay appended evidence after the only terminal"
);
let remaining = fixture.request("GET", "/api/v1/threads/proposals", None)?;
anyhow::ensure!(remaining.body["proposals"] == json!([]), "recovery left pending authority");
Ok(())
},
)
}

fn scheduler_entries(fixture: &ThreadsFixture) -> Result<usize> {
Ok(fs::read_to_string(fixture.coven_home.join("daemon-recovery.log"))?
.lines()
.filter(|line| line.contains("threads_scheduler_checkpoint phase=pass-lock-wait"))
.count())
}

#[test]
fn unsupported_retired_identity_declaration_never_migrates_or_stages() -> Result<()> {
let corpus = retired_ward_corpus()?;
let mut case = retired_review_case(&corpus)?.clone();
let unsupported = corpus["unsupported_cases"]
.as_array()
.context("unsupported corpus cases")?
.iter()
.find(|case| case["id"] == "unknown-identity-fact")
.context("canonical unsupported declaration")?;
case["declarations"] = unsupported["declarations"].clone();
run_clocked_journey(
"retired-corpus-unsupported-identity-declaration",
|home, workspace| {
seed_retired_review_source(home, workspace, &case)?;
Ok(())
},
|fixture, _| {
let original = fs::read(fixture.workspace.join("ward.toml"))?;
let migration = migrate_retired_review_source(&fixture.coven_home)?;
let output = format!(
"{}\n{}",
String::from_utf8_lossy(&migration.stdout),
String::from_utf8_lossy(&migration.stderr),
);
anyhow::ensure!(
!migration.status.success()
&& output.contains(unsupported["expected_error_contains"].as_str().context("expected refusal")?),
"unsupported corpus declaration migrated: {output}"
);
fs::create_dir_all(&fixture.artifact_dir)?;
fs::write(
fixture.artifact_dir.join("migration-refusal.txt"),
fixture.sanitize_fixture_text(&output),
)?;
fs::write(
fixture.artifact_dir.join("corpus-unsupported.json"),
serde_json::to_vec_pretty(unsupported)?,
)?;
anyhow::ensure!(fs::read(fixture.workspace.join("ward.toml"))? == original);
anyhow::ensure!(!fixture.workspace.join("ward.toml.v01.bak").exists());
let edits: Vec<Value> = case["surfaces"].as_array().context("surfaces")?
.iter()
.map(|surface| json!({"target": surface["path"], "contents": surface["after"]}))
.collect();
let refused = fixture.request(
"POST",
"/api/v1/familiars/sage/edits",
Some(&json!({"edits": edits})),
)?;
anyhow::ensure!(
refused.status == 500 && refused.body["error"]["code"] == "ward_config_invalid",
"unsupported retired declaration reached staging: {refused:?}"
);
assert_corpus_bytes(fixture, &case, "before")?;
let listed = fixture.request("GET", "/api/v1/threads/proposals", None)?;
anyhow::ensure!(listed.body["proposals"] == json!([]));
let authority: i64 = fixture.store()?.query_row(
"SELECT COUNT(*) FROM ward_audit WHERE event_type IN
('proposal_submitted','proposal_window_opened','proposal_approved','apply_audit')",
[],
|row| row.get(0),
)?;
anyhow::ensure!(authority == 0, "unsupported migration gained proposal authority");
Ok(())
},
)
}
Loading