diff --git a/.github/scripts/test_ci_changes.py b/.github/scripts/test_ci_changes.py index 4ea596f2ab..97cc7410f9 100644 --- a/.github/scripts/test_ci_changes.py +++ b/.github/scripts/test_ci_changes.py @@ -949,12 +949,66 @@ def test_breg_tutorial_runs_native_example_recovery(self) -> None: breg_job = workflow.split("\n breg-tutorial:\n", 1)[1].split( "\n breg-evidence-composition:\n", 1 )[0] - self.assertIn( + test = ( "dev::examples::tests::" - "native_create_and_apply_recover_after_process_exit_without_duplicate_revisions", - breg_job, + "native_create_recovers_after_process_exit_without_duplicate_records" ) + self.assertIn(test, breg_job) self.assertIn("-- --ignored --exact", breg_job) + # A renamed test would otherwise leave the step running zero tests. + self.assertIn( + "grep -q 'test result: ok\\. 1 passed' " + '"${RUNNER_TEMP}/breg-example-recovery.log" || ' + '{ echo "::error::expected exactly one passing test for ' + f'{test.rsplit("::", 1)[1]}"; exit 1; }}', + breg_job, + ) + source = ( + Path("crates/registry-bregctl/src/dev/examples.rs").read_text() + ) + self.assertIn(f"fn {test.rsplit('::', 1)[1]}()", source) + + def test_casework_postgres_runs_task_approval_and_local_session_exactly( + self, + ) -> None: + workflow = Path(".github/workflows/ci.yml").read_text() + casework_job = workflow.split("\n casework-postgres:\n", 1)[1].split( + "\n scheduling-contracts:\n", 1 + )[0] + cases = ( + ( + "task_grants::native_exchange_tests::" + "approved_casework_tasks_reach_evidence_breg_and_scheduling_through_stock_thunderid", + "casework-task-approval.log", + ), + ( + "task_grants::local_session_tests::" + "source_backed_dev_approves_exchanges_and_revokes_on_stock_issuer", + "casework-local-session.log", + ), + ) + source_dir = Path("crates/registry-casework/src/task_grants") + for test, log in cases: + with self.subTest(test=test): + self.assertIn(test, casework_job) + # A renamed test would otherwise leave the step running zero + # tests, and an inexact filter could silently match more than + # one, so the step names each test exactly and fails unless + # it reports one pass. + self.assertRegex( + casework_job, + rf"-- --ignored --exact \\\s*\n\s*{re.escape(test)} \\", + ) + self.assertIn( + "grep -q 'test result: ok\\. 1 passed' " + f'"${{RUNNER_TEMP}}/{log}" || ' + '{ echo "::error::expected exactly one passing test for ' + f'{test.rsplit("::", 1)[1]}"; exit 1; }}', + casework_job, + ) + module, name = test.rsplit("::", 2)[1:] + module_source = (source_dir / f"{module}.rs").read_text() + self.assertIn(f"fn {name}()", module_source) def test_breg_tutorial_inputs_cover_every_replayed_tutorial(self) -> None: # Each page's tutorial_test frontmatter is the source of truth for diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8c549cef8d..934eccf501 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -408,6 +408,22 @@ jobs: - name: Cargo deny run: cargo deny check + - name: Standalone Cargo lockfiles + # Fuzz, SDK and editor workspaces keep their own Cargo.lock, and most + # build root crates by path, so a root change can leave one stale. + # Resolving with --locked fails on a stale lock and reads only the + # registry index, not crate sources. + shell: bash + run: | + set -euo pipefail + mapfile -t locks < <(git ls-files '*/Cargo.lock') + (( ${#locks[@]} > 0 )) + for lock in "${locks[@]}"; do + echo "Checking ${lock}" + cargo update --workspace --locked \ + --manifest-path "$(dirname "${lock}")/Cargo.toml" + done + rust-quality: name: Rust format and clippy needs: @@ -752,6 +768,8 @@ jobs: POSTGRES_CONTAINER_ID: ${{ job.services.postgres.id }} run: docker exec "$POSTGRES_CONTAINER_ID" createdb -U casework casework_task_approval - name: Casework task approval through stock issuer to Evidence, BReg, and Scheduling + # An inexact name filter that selects nothing still exits 0, so the + # step names each test exactly and fails unless it reports one pass. shell: bash env: CASEWORK_ASSIGNMENT_TEST_DATABASE_URL: postgresql://casework:casework_test@localhost:${{ job.services.postgres.ports['5432'] }}/casework_task_approval @@ -763,13 +781,17 @@ jobs: SCHEDULING_AUTH_PROBE_BIN="${GITHUB_WORKSPACE}/target/debug/examples/scheduling-auth-probe" \ cargo test --locked -p registry-casework \ --features postgres-test --lib \ - approved_casework_tasks_reach_evidence_breg_and_scheduling_through_stock_thunderid \ - -- --ignored + -- --ignored --exact \ + task_grants::native_exchange_tests::approved_casework_tasks_reach_evidence_breg_and_scheduling_through_stock_thunderid \ + | tee "${RUNNER_TEMP}/casework-task-approval.log" + grep -q 'test result: ok\. 1 passed' "${RUNNER_TEMP}/casework-task-approval.log" || { echo "::error::expected exactly one passing test for approved_casework_tasks_reach_evidence_breg_and_scheduling_through_stock_thunderid"; exit 1; } cargo build --locked -p registry-casework -p registry-caseworkctl cargo test --locked -p registry-casework \ --features postgres-test --lib \ - source_backed_dev_approves_exchanges_and_revokes_on_stock_issuer \ - -- --ignored + -- --ignored --exact \ + task_grants::local_session_tests::source_backed_dev_approves_exchanges_and_revokes_on_stock_issuer \ + | tee "${RUNNER_TEMP}/casework-local-session.log" + grep -q 'test result: ok\. 1 passed' "${RUNNER_TEMP}/casework-local-session.log" || { echo "::error::expected exactly one passing test for source_backed_dev_approves_exchanges_and_revokes_on_stock_issuer"; exit 1; } scheduling-contracts: name: Scheduling product contracts @@ -1496,12 +1518,16 @@ jobs: --bin-dir target/debug --verification - name: Verify retained example recovery after process exit + # An inexact name filter that selects nothing still exits 0, so the + # step names the test exactly and fails unless it reports one pass. shell: bash run: | set -euo pipefail cargo test --locked -p registry-bregctl --lib \ - dev::examples::tests::native_create_and_apply_recover_after_process_exit_without_duplicate_revisions \ - -- --ignored --exact + dev::examples::tests::native_create_recovers_after_process_exit_without_duplicate_records \ + -- --ignored --exact \ + | tee "${RUNNER_TEMP}/breg-example-recovery.log" + grep -q 'test result: ok\. 1 passed' "${RUNNER_TEMP}/breg-example-recovery.log" || { echo "::error::expected exactly one passing test for native_create_recovers_after_process_exit_without_duplicate_records"; exit 1; } - name: Execute the Base Registry Engine tutorials from a reader directory shell: bash diff --git a/crates/registry-bregctl/src/dev/examples.rs b/crates/registry-bregctl/src/dev/examples.rs index b935108aa7..4f983d9543 100644 --- a/crates/registry-bregctl/src/dev/examples.rs +++ b/crates/registry-bregctl/src/dev/examples.rs @@ -1807,7 +1807,7 @@ mod tests { /// An optional structured field exercises untyped nested reference refusal. fn prepare_native_fixture(project: &Path) { let path = project.join("registry.yaml"); - let mut source: Value = serde_json::from_slice(&fs::read(&path).unwrap()).unwrap(); + let mut source: Value = serde_norway::from_slice(&fs::read(&path).unwrap()).unwrap(); let entity = source["entities"] .as_array_mut() .unwrap() @@ -1861,6 +1861,20 @@ mod tests { grant[fields].as_array_mut().unwrap().push(json!("notes")); } fs::write(path, serde_json::to_vec_pretty(&source).unwrap()).unwrap(); + // The starter's change requests name the `casework` review authority, + // which a start activates. No scenario here submits a request, so the + // binding needs only a declared client and an unused loopback endpoint. + let path = project.join("dev-clients.yaml"); + let mut clients: Value = serde_norway::from_slice(&fs::read(&path).unwrap()).unwrap(); + clients["clients"].as_array_mut().unwrap().push(json!({ + "id":"casework-producer", "accessProfiles":[], + "scopes":["casework:reviews:request"], "claims":{} + })); + clients["reviewAuthorities"] = json!({"casework":{ + "endpoint":"http://127.0.0.1:9/", "profile":"integration-requester", + "producerId":"registry-breg", "recoveryDays":7, "client":"casework-producer" + }}); + fs::write(path, serde_norway::to_string(&clients).unwrap()).unwrap(); } fn preflight_refuses_before_writes(project: &Path, first: &Value) { diff --git a/crates/registry-caseworkctl/src/dev/tests.rs b/crates/registry-caseworkctl/src/dev/tests.rs index ffc29a63f7..8f902b6ef0 100644 --- a/crates/registry-caseworkctl/src/dev/tests.rs +++ b/crates/registry-caseworkctl/src/dev/tests.rs @@ -43,18 +43,41 @@ fn standalone(root: &Path) -> PathBuf { project } -fn wait_for_process_exit(pid: rustix::process::Pid, grace: Duration) -> bool { - let deadline = Instant::now() + grace; - while rustix::process::test_kill_process(pid).is_ok() && Instant::now() < deadline { +/// Only a hung child reaches this: each wait ends as soon as its event +/// happens, so a starved runner makes a test slower, never wrong. Matches +/// the outer bound the language-server and evidencectl process tests use. +const TEST_EVENT_BOUND: Duration = Duration::from_secs(120); + +/// Polls `condition` until it holds or [`TEST_EVENT_BOUND`] passes. +fn eventually(mut condition: impl FnMut() -> bool) -> bool { + let deadline = Instant::now() + TEST_EVENT_BOUND; + while !condition() { + if Instant::now() >= deadline { + return false; + } thread::sleep(Duration::from_millis(5)); } - rustix::process::test_kill_process(pid).is_err() + true +} + +fn wait_for_process_exit(pid: rustix::process::Pid) -> bool { + eventually(|| rustix::process::test_kill_process(pid).is_err()) } fn file_has_bytes(path: &Path) -> bool { fs::metadata(path).is_ok_and(|metadata| metadata.len() > 0) } +/// The PID a child wrote to `path`, or `None` if it wrote nothing within +/// [`TEST_EVENT_BOUND`]. +fn announced_pid(path: &Path) -> Option { + if !eventually(|| file_has_bytes(path)) { + return None; + } + let pid = fs::read_to_string(path).unwrap().parse::().unwrap(); + Some(rustix::process::Pid::from_raw(pid).unwrap()) +} + #[test] fn init_clients_bind_the_standalone_template() { let root = tempfile::tempdir().unwrap(); @@ -1348,7 +1371,8 @@ fn failed_start_waits_for_the_supervisor_lock_to_be_released() { elapsed >= Duration::from_millis(150), "elapsed: {elapsed:?}" ); - assert!(elapsed < Duration::from_secs(1), "elapsed: {elapsed:?}"); + // A supervisor that outlived its release grace would be killed instead. + assert!(elapsed < SUPERVISOR_RELEASE_GRACE, "elapsed: {elapsed:?}"); } #[test] @@ -1397,7 +1421,9 @@ fn database_readiness_commands_stop_at_the_aggregate_deadline() { ); assert!(refusal.contains("timed out"), "{refusal}"); - assert!(started.elapsed() < Duration::from_secs(1)); + // Ignoring the aggregate deadline would run to the child's own + // CHILD_DEADLINE instead. + assert!(started.elapsed() < CHILD_DEADLINE / 2); } #[test] @@ -1409,12 +1435,11 @@ fn interrupted_native_prerequisite_is_killed_and_reaped() { let signal = Arc::clone(&terminate); let marker_for_signal = marker.clone(); let interrupter = thread::spawn(move || { - let deadline = Instant::now() + Duration::from_secs(1); - while !marker_for_signal.exists() { - assert!(Instant::now() < deadline, "prerequisite did not start"); - thread::sleep(Duration::from_millis(5)); - } + let pid = announced_pid(&marker_for_signal); + // A missed PID still interrupts, so the prerequisite is reaped + // before the assertion below reports it. signal.store(true, Ordering::Relaxed); + pid }); let started = Instant::now(); @@ -1433,12 +1458,15 @@ fn interrupted_native_prerequisite_is_killed_and_reaped() { ) .unwrap_err() ); - interrupter.join().unwrap(); - let pid = fs::read_to_string(&marker).unwrap().parse::().unwrap(); - let pid = rustix::process::Pid::from_raw(pid).unwrap(); + let pid = interrupter + .join() + .unwrap() + .expect("prerequisite did not start"); assert!(refusal.contains("interrupted"), "{refusal}"); - assert!(started.elapsed() < Duration::from_secs(1)); + // The prerequisite never exits by itself, so only the interruption ends + // it long before its own CHILD_DEADLINE. + assert!(started.elapsed() < CHILD_DEADLINE / 2); assert!(rustix::process::test_kill_process(pid).is_err()); } @@ -1447,8 +1475,9 @@ fn native_pump_setup_failures_reap_the_child_and_join_started_pumps() { let root = tempfile::tempdir().unwrap(); private::directory(&root.path().join("logs")).unwrap(); for fail_on in [1, 2] { + // The child outlasts every wait below, so only cleanup can end it. let child = Command::new("/bin/sleep") - .arg("5") + .arg("180") .stdout(Stdio::piped()) .stderr(Stdio::piped()) .spawn() @@ -1490,7 +1519,7 @@ fn native_pump_setup_failures_reap_the_child_and_join_started_pumps() { let reader = if fail_on == 1 { "output" } else { "diagnostic" }; assert!(refusal.contains(reader), "{refusal}"); - assert!(started.elapsed() < Duration::from_secs(1)); + assert!(started.elapsed() < TEST_EVENT_BOUND); assert_eq!(joined.load(Ordering::Relaxed), fail_on == 2); assert!(rustix::process::test_kill_process(pid).is_err()); } @@ -1550,7 +1579,8 @@ fn failed_native_stdin_write_reaps_the_child_and_joins_pumps() { refusal.contains("write native prerequisite input"), "{refusal}" ); - assert!(started.elapsed() < Duration::from_secs(1)); + // The child loops forever, so waiting on it would reach CHILD_DEADLINE. + assert!(started.elapsed() < CHILD_DEADLINE / 2); assert!(stdout_joined.load(Ordering::Relaxed)); assert!(stderr_joined.load(Ordering::Relaxed)); assert!(rustix::process::test_kill_process(pid).is_err()); @@ -1561,21 +1591,16 @@ fn active_http_prerequisite_stops_promptly_when_interrupted() { let listener = TcpListener::bind(("127.0.0.1", 0)).unwrap(); let address = listener.local_addr().unwrap(); let terminate = Arc::new(AtomicBool::new(false)); - let release = Arc::new(AtomicBool::new(false)); + let (release, released) = mpsc::channel::<()>(); let signal = Arc::clone(&terminate); - let release_server = Arc::clone(&release); let server = thread::spawn(move || { let (mut stream, _) = listener.accept().unwrap(); - stream - .set_read_timeout(Some(Duration::from_secs(1))) - .unwrap(); + stream.set_read_timeout(Some(TEST_EVENT_BOUND)).unwrap(); let mut request = [0u8; 1]; assert_eq!(stream.read(&mut request).unwrap(), 1); signal.store(true, Ordering::Relaxed); - let deadline = Instant::now() + Duration::from_secs(2); - while !release_server.load(Ordering::Relaxed) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } + // Hold the connection open, never responding, until the client is done. + let _ = released.recv_timeout(TEST_EVENT_BOUND); }); let started = Instant::now(); @@ -1588,26 +1613,23 @@ fn active_http_prerequisite_stops_promptly_when_interrupted() { &terminate, ); let elapsed = started.elapsed(); - release.store(true, Ordering::Relaxed); + drop(release); server.join().unwrap(); let refusal = format!("{:#}", result.unwrap_err()); assert!(refusal.contains("interrupted"), "{refusal}"); - // The budget only has to show the terminate flag beat HTTP_TIMEOUT (10s, - // dev/mod.rs), so 5s still leaves a 2x margin. The window being measured - // is not just cancellation: http_with_timeout builds a fresh - // current-thread runtime and reqwest client inside it, and that client - // build loads the system root certificate store on first use in the - // process. Cancellation itself is bounded by the 50ms sleep poll loop in - // http_with_timeout, so a 1s budget was really measuring process - // warm-up, which is the volatile term here. - assert!(elapsed < Duration::from_secs(5), "elapsed: {elapsed:?}"); + // The request timeout reports the prerequisite unavailable, not + // interrupted, so the refusal above already shows the terminate flag won. + // The window also covers building a fresh runtime and reqwest client, + // which loads the system root store and slows under load, so the only + // bound on it is the one the flag had to beat. + assert!(elapsed < HTTP_TIMEOUT, "elapsed: {elapsed:?}"); } #[test] fn service_http_readiness_stops_at_the_phase_deadline() { let child = Command::new("/bin/sleep") - .arg("5") + .arg("180") .process_group(0) .spawn() .unwrap(); @@ -1632,7 +1654,8 @@ fn service_http_readiness_stops_at_the_phase_deadline() { let _ = service.stop_with_grace(Duration::from_millis(50), Duration::from_millis(10)); assert!(refusal.contains("readiness timed out"), "{refusal}"); - assert!(elapsed < Duration::from_secs(1), "elapsed: {elapsed:?}"); + // A probe handed the full request timeout would sleep through it. + assert!(elapsed < HTTP_TIMEOUT, "elapsed: {elapsed:?}"); assert!(!request_timeouts.is_empty()); assert!(request_timeouts[0] <= deadline.duration_since(started)); } @@ -1898,12 +1921,14 @@ fn service_guard_process_helper() { let interruption = StartInterruption::install().unwrap(); let mut command = Command::new(binary); command.args(arguments).stdin(Stdio::null()); + // A graceful service's shutdown must fit inside this grace even on a + // starved runner; a stubborn service spends all of it before forced KILL. let forced = guard_service_command( command, std::io::stdin(), Arc::clone(&interruption.requested), - Duration::from_millis(500), - Duration::from_millis(500), + Duration::from_secs(5), + Duration::from_secs(5), ) .unwrap(); if forced { @@ -1944,12 +1969,8 @@ fn nonzero_outer_guard_helper() { .arg(&service_pid_file) .spawn() .unwrap(); - let deadline = Instant::now() + Duration::from_secs(5); - while !file_has_bytes(&service_pid_file) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } assert!( - file_has_bytes(&service_pid_file), + eventually(|| file_has_bytes(&service_pid_file)), "service did not become ready" ); std::process::exit(23); @@ -1965,6 +1986,7 @@ fn guarded_service_stops_after_its_supervisor_is_killed() { ) .unwrap(); fs::set_permissions(&binary, fs::Permissions::from_mode(0o700)).unwrap(); + let service_pid_file = root.path().join("service.pid"); let mut supervisor = Command::new(std::env::current_exe().unwrap()) .args([ "--exact", @@ -1978,62 +2000,43 @@ fn guarded_service_stops_after_its_supervisor_is_killed() { .stderr(Stdio::null()) .spawn() .unwrap(); - let service_pid_file = root.path().join("service.pid"); - let deadline = Instant::now() + Duration::from_secs(5); - while !file_has_bytes(&service_pid_file) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(10)); - } - if !file_has_bytes(&service_pid_file) { + let Some(service_pid) = announced_pid(&service_pid_file) else { supervisor.kill().unwrap(); supervisor.wait().unwrap(); panic!("guarded service did not start"); - } - let service_pid = fs::read_to_string(&service_pid_file) - .unwrap() - .parse::() - .unwrap(); - let service_pid = rustix::process::Pid::from_raw(service_pid).unwrap(); - let started = Instant::now(); + }; // SIGKILL skips every supervisor destructor. The kernel still closes the // supervisor's liveness writer, which must stop the exact guarded child. + // The service loops forever by itself, so only that path can end it. supervisor.kill().unwrap(); supervisor.wait().unwrap(); - let deadline = Instant::now() + Duration::from_secs(5); - while rustix::process::test_kill_process(service_pid).is_ok() && Instant::now() < deadline { - thread::sleep(Duration::from_millis(10)); - } - let stopped = rustix::process::test_kill_process(service_pid).is_err(); + let stopped = wait_for_process_exit(service_pid); if !stopped { rustix::process::kill_process(service_pid, rustix::process::Signal::KILL).unwrap(); } assert!(stopped, "guarded service survived supervisor death"); - assert!(started.elapsed() < Duration::from_secs(5)); } #[test] fn service_guard_owns_a_stubborn_child_during_startup_interruption() { let root = tempfile::tempdir().unwrap(); let binary = root.path().join("stubborn.sh"); - let service_pid_file = root.path().join("service.pid"); fs::write( &binary, b"#!/bin/sh\ntrap '' TERM\nprintf '%s' \"$$\" > \"$1\"\nwhile :; do sleep 0.02; done\n", ) .unwrap(); fs::set_permissions(&binary, fs::Permissions::from_mode(0o700)).unwrap(); + let service_pid_file = root.path().join("service.pid"); + let marker = service_pid_file.clone(); let (reader, writer) = UnixStream::pair().unwrap(); let terminate = Arc::new(AtomicBool::new(false)); let request = Arc::clone(&terminate); - let marker = service_pid_file.clone(); let requester = thread::spawn(move || { - let deadline = Instant::now() + Duration::from_secs(5); - while !marker.exists() && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } - let ready = marker.exists(); - if ready { + let ready = announced_pid(&marker); + if ready.is_some() { request.store(true, Ordering::Relaxed); } // EOF is also a cleanup request if readiness failed, so the assertion @@ -2052,19 +2055,13 @@ fn service_guard_owns_a_stubborn_child_during_startup_interruption() { Duration::from_millis(50), ) .unwrap(); - let ready = requester.join().unwrap(); - assert!( - ready, - "stubborn service did not reach its startup handshake" - ); - let service_pid = fs::read_to_string(&service_pid_file) + let service_pid = requester + .join() .unwrap() - .parse::() - .unwrap(); - let service_pid = rustix::process::Pid::from_raw(service_pid).unwrap(); + .expect("stubborn service did not reach its startup handshake"); assert!(forced); - assert!(wait_for_process_exit(service_pid, Duration::from_secs(2))); + assert!(wait_for_process_exit(service_pid)); } #[test] @@ -2072,14 +2069,16 @@ fn service_guard_does_not_force_kill_after_a_fast_term_exit() { let (reader, writer) = UnixStream::pair().unwrap(); drop(writer); let mut command = Command::new("/bin/sleep"); - command.arg("5").stdin(Stdio::null()); + command.arg("180").stdin(Stdio::null()); + // `sleep` exits at once on TERM. A grace far longer than any scheduling + // delay leaves forced KILL as the result only of ignoring that exit. let forced = guard_service_command( command, reader, Arc::new(AtomicBool::new(false)), - Duration::from_millis(100), - Duration::from_millis(100), + TEST_EVENT_BOUND, + TEST_EVENT_BOUND, ) .unwrap(); @@ -2108,11 +2107,7 @@ fn established_service_keeps_its_graceful_shutdown_window() { "guarded", ) .unwrap(); - let deadline = Instant::now() + Duration::from_secs(5); - while !file_has_bytes(&service_pid_file) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } - if !file_has_bytes(&service_pid_file) { + if announced_pid(&service_pid_file).is_none() { let _ = service.stop(); panic!("service did not reach its startup handshake"); } @@ -2122,7 +2117,9 @@ fn established_service_keeps_its_graceful_shutdown_window() { assert!(graceful.exists(), "guardian truncated graceful shutdown"); assert!(started.elapsed() >= Duration::from_millis(150)); - assert!(started.elapsed() < Duration::from_secs(1)); + // Returning early means stop followed the service's own exit rather than + // spending its whole signal grace. + assert!(started.elapsed() < SERVICE_SIGNAL_GRACE); } #[test] @@ -2130,33 +2127,24 @@ fn established_stubborn_service_reports_forced_shutdown() { let root = tempfile::tempdir().unwrap(); private::directory(&root.path().join("logs")).unwrap(); let binary = root.path().join("stubborn.sh"); - let service_pid_file = root.path().join("service.pid"); fs::write( &binary, b"#!/bin/sh\ntrap '' TERM\nprintf '%s' \"$$\" > \"$1\"\nwhile :; do sleep 0.02; done\n", ) .unwrap(); fs::set_permissions(&binary, fs::Permissions::from_mode(0o700)).unwrap(); + let service_pid_file = root.path().join("service.pid"); let mut service = service(&binary, &[], &service_pid_file, &[], root.path(), "guarded").unwrap(); - let deadline = Instant::now() + Duration::from_secs(5); - while !file_has_bytes(&service_pid_file) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } - if !file_has_bytes(&service_pid_file) { + let Some(service_pid) = announced_pid(&service_pid_file) else { let _ = service.stop(); panic!("stubborn service did not reach its startup handshake"); - } - let service_pid = fs::read_to_string(&service_pid_file) - .unwrap() - .parse::() - .unwrap(); - let service_pid = rustix::process::Pid::from_raw(service_pid).unwrap(); + }; let refusal = format!("{:#}", service.stop().unwrap_err()); assert!(refusal.contains("required forced shutdown"), "{refusal}"); - assert!(wait_for_process_exit(service_pid, Duration::from_secs(2))); + assert!(wait_for_process_exit(service_pid)); } #[test] @@ -2164,30 +2152,16 @@ fn killed_guard_leaves_the_supervisor_to_clean_its_exact_service_group() { let root = tempfile::tempdir().unwrap(); private::directory(&root.path().join("logs")).unwrap(); let binary = root.path().join("stubborn.sh"); - let service_pid_file = root.path().join("service.pid"); fs::write( &binary, b"#!/bin/sh\ntrap '' TERM\nprintf '%s' \"$$\" > \"$1\"\nwhile :; do sleep 0.02; done\n", ) .unwrap(); fs::set_permissions(&binary, fs::Permissions::from_mode(0o700)).unwrap(); + let service_pid_file = root.path().join("service.pid"); let mut service = service(&binary, &[], &service_pid_file, &[], root.path(), "guarded").unwrap(); - let deadline = Instant::now() + Duration::from_secs(5); - while !file_has_bytes(&service_pid_file) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } - assert!( - file_has_bytes(&service_pid_file), - "guarded service did not start" - ); - let service_pid = rustix::process::Pid::from_raw( - fs::read_to_string(&service_pid_file) - .unwrap() - .parse::() - .unwrap(), - ) - .unwrap(); + let service_pid = announced_pid(&service_pid_file).expect("guarded service did not start"); assert_eq!( rustix::process::getpgid(Some(service_pid)).unwrap(), service.guard_pgid, @@ -2195,12 +2169,8 @@ fn killed_guard_leaves_the_supervisor_to_clean_its_exact_service_group() { ); rustix::process::kill_process(service.guard_pid, rustix::process::Signal::KILL).unwrap(); - let deadline = Instant::now() + Duration::from_secs(2); - while service.guard_exit().unwrap().is_none() && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } assert!( - service.guard_exit().unwrap().is_some(), + eventually(|| service.guard_exit().unwrap().is_some()), "killed guard did not become waitable" ); assert!( @@ -2216,7 +2186,7 @@ fn killed_guard_leaves_the_supervisor_to_clean_its_exact_service_group() { ); assert!(refusal.contains("guard exited abnormally"), "{refusal}"); - assert!(wait_for_process_exit(service_pid, Duration::from_secs(2))); + assert!(wait_for_process_exit(service_pid)); } #[test] @@ -2241,33 +2211,22 @@ fn nonzero_guard_exit_is_detected_without_waiting_for_pump_eof() { .env("CASEWORKCTL_TEST_NONZERO_GUARD_PID", &service_pid_file) .env("CASEWORKCTL_TEST_NONZERO_GUARD_SERVICE", &service_binary); let mut service = service_with_guard_command(guard, root.path(), "guarded").unwrap(); - // The helper itself permits five seconds for the service PID file. Give - // the parent that complete startup budget plus scheduling margin when the - // full test suite is running concurrently. - let deadline = Instant::now() + Duration::from_secs(6); - while (!file_has_bytes(&service_pid_file) || service.guard_exit().unwrap().is_none()) - && Instant::now() < deadline - { - thread::sleep(Duration::from_millis(5)); - } - assert!( - file_has_bytes(&service_pid_file), - "guard did not create its service" - ); + // The helper exits only after its service has written its PID. assert!( - service.guard_exit().unwrap().is_some(), + eventually(|| service.guard_exit().unwrap().is_some()), "nonzero guard did not become waitable" ); let service_pid = rustix::process::Pid::from_raw( fs::read_to_string(&service_pid_file) - .unwrap() + .expect("guard did not create its service") .parse::() .unwrap(), ) .unwrap(); assert!(rustix::process::test_kill_process(service_pid).is_ok()); - let started = Instant::now(); + // The service ignores HUP and TERM and holds the pump pipes open, so + // waiting for pump EOF before group cleanup would never return. let refusal = format!( "{:#}", service @@ -2276,8 +2235,7 @@ fn nonzero_guard_exit_is_detected_without_waiting_for_pump_eof() { ); assert!(refusal.contains("guard exited abnormally"), "{refusal}"); - assert!(started.elapsed() < Duration::from_secs(1)); - assert!(wait_for_process_exit(service_pid, Duration::from_secs(2))); + assert!(wait_for_process_exit(service_pid)); } #[test] @@ -2288,30 +2246,20 @@ fn live_guard_timeout_kills_the_pinned_group_before_reaping() { let service_pid_file = root.path().join("service.pid"); // Publish the descendant PID from the guard that created it. This proves // the group member exists without depending on when that child is scheduled. + // Both sleeps outlast every wait below, so only the group KILL ends them. fs::write( &guard_binary, - b"#!/bin/sh\ntrap '' TERM\n/bin/sleep 60 &\nprintf '%s' \"$!\" > \"$1\"\nexec /bin/sleep 60\n", + b"#!/bin/sh\ntrap '' TERM\n/bin/sleep 180 &\nprintf '%s' \"$!\" > \"$1\"\nexec /bin/sleep 180\n", ) .unwrap(); fs::set_permissions(&guard_binary, fs::Permissions::from_mode(0o700)).unwrap(); let mut guard = Command::new(&guard_binary); guard.arg(&service_pid_file); let mut service = service_with_guard_command(guard, root.path(), "guarded").unwrap(); - let deadline = Instant::now() + Duration::from_secs(5); - while !file_has_bytes(&service_pid_file) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } - if !file_has_bytes(&service_pid_file) { + let Some(service_pid) = announced_pid(&service_pid_file) else { let _ = service.stop_with_grace(Duration::from_millis(75), Duration::from_millis(10)); panic!("guarded service did not start"); - } - let service_pid = rustix::process::Pid::from_raw( - fs::read_to_string(&service_pid_file) - .unwrap() - .parse::() - .unwrap(), - ) - .unwrap(); + }; assert!(rustix::process::test_kill_process(service_pid).is_ok()); assert_eq!( rustix::process::getpgid(Some(service_pid)).unwrap(), @@ -2329,14 +2277,15 @@ fn live_guard_timeout_kills_the_pinned_group_before_reaping() { assert!(refusal.contains("required forced shutdown"), "{refusal}"); assert!(started.elapsed() >= Duration::from_millis(70)); - assert!(started.elapsed() < Duration::from_secs(1)); - assert!(wait_for_process_exit(service_pid, Duration::from_secs(2))); + // Reaping before the KILL would block on the guard's own three-minute sleep. + assert!(started.elapsed() < TEST_EVENT_BOUND); + assert!(wait_for_process_exit(service_pid)); } #[test] fn post_kill_wait_is_bounded_when_a_guard_does_not_become_waitable() { let mut guard = Command::new("/bin/sleep") - .arg("5") + .arg("180") .process_group(0) .spawn() .unwrap(); @@ -2345,7 +2294,7 @@ fn post_kill_wait_is_bounded_when_a_guard_does_not_become_waitable() { // Model a kernel reporting successful group KILL without making the guard // waitable. Cleanup must return at its own bound instead of entering a - // blocking Child::wait. + // blocking Child::wait, which would last the guard's three-minute sleep. let refusal = format!( "{:#}", kill_guard_group_and_reap_with(&mut guard, guard_pid, Duration::from_millis(40), |_pgid| { @@ -2356,7 +2305,7 @@ fn post_kill_wait_is_bounded_when_a_guard_does_not_become_waitable() { assert!(refusal.contains("bounded cleanup wait"), "{refusal}"); assert!(started.elapsed() >= Duration::from_millis(35)); - assert!(started.elapsed() < Duration::from_secs(1)); + assert!(started.elapsed() < TEST_EVENT_BOUND); assert!(guard_exit(guard_pid).unwrap().is_none()); rustix::process::kill_process(guard_pid, rustix::process::Signal::KILL).unwrap(); guard.wait().unwrap(); @@ -2365,16 +2314,18 @@ fn post_kill_wait_is_bounded_when_a_guard_does_not_become_waitable() { #[test] fn failed_group_kill_never_enters_a_blocking_guard_wait() { let mut guard = Command::new("/bin/sleep") - .arg("5") + .arg("180") .process_group(0) .spawn() .unwrap(); let guard_pid = rustix::process::Pid::from_raw(guard.id() as i32).unwrap(); let started = Instant::now(); + // A failed KILL must return at once. Waiting for the guard instead would + // spend the whole post-kill grace, or the guard's three-minute sleep. let refusal = format!( "{:#}", - kill_guard_group_and_reap_with(&mut guard, guard_pid, Duration::from_secs(1), |_pgid| Err( + kill_guard_group_and_reap_with(&mut guard, guard_pid, TEST_EVENT_BOUND, |_pgid| Err( anyhow::anyhow!("injected group KILL failure") ),) .unwrap_err() @@ -2382,7 +2333,7 @@ fn failed_group_kill_never_enters_a_blocking_guard_wait() { assert!(refusal.contains("cannot KILL"), "{refusal}"); assert!(refusal.contains("injected group KILL failure"), "{refusal}"); - assert!(started.elapsed() < Duration::from_millis(100)); + assert!(started.elapsed() < TEST_EVENT_BOUND / 2); assert!(guard_exit(guard_pid).unwrap().is_none()); rustix::process::kill_process(guard_pid, rustix::process::Signal::KILL).unwrap(); guard.wait().unwrap(); @@ -2428,7 +2379,8 @@ fn service_pump_setup_failures_reap_the_child_and_join_started_pumps() { let reader = if fail_on == 1 { "output" } else { "diagnostic" }; assert!(refusal.contains(reader), "{refusal}"); - assert!(started.elapsed() < Duration::from_secs(1)); + // The child waits on stdin forever, so only cleanup can end it. + assert!(started.elapsed() < TEST_EVENT_BOUND); assert_eq!(joined.load(Ordering::Relaxed), fail_on == 2); assert!(rustix::process::test_kill_process(pid).is_err()); } @@ -2440,13 +2392,13 @@ fn guardian_pump_setup_failure_reaps_a_stubborn_owned_service() { let logs = root.path().join("logs"); private::directory(&logs).unwrap(); let binary = root.path().join("stubborn.sh"); - let service_pid_file = root.path().join("service.pid"); fs::write( &binary, b"#!/bin/sh\ntrap '' TERM\nprintf '%s' \"$$\" > \"$1\"\nwhile :; do sleep 0.02; done\n", ) .unwrap(); fs::set_permissions(&binary, fs::Permissions::from_mode(0o700)).unwrap(); + let service_pid_file = root.path().join("service.pid"); let mut guardian = service_guard_command(&binary, &[service_pid_file.as_os_str().to_owned()]).unwrap(); let child = guardian @@ -2457,34 +2409,22 @@ fn guardian_pump_setup_failure_reaps_a_stubborn_owned_service() { .spawn() .unwrap(); let journal = RetainedJournal::open(&logs.join("casework.log")).unwrap(); - let started = Instant::now(); - let mut ready = false; + let mut service_pid = None; let refusal = format!( "{:#}", service_with_pump_spawner(child, journal, |_stream, _task| { - let deadline = Instant::now() + Duration::from_secs(5); - while !file_has_bytes(&service_pid_file) && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } - ready = file_has_bytes(&service_pid_file); + service_pid = announced_pid(&service_pid_file); Err(std::io::Error::other("injected pump spawn failure")) }) .err() .expect("injected guardian pump spawn must fail") ); - assert!( - ready, - "stubborn service did not reach its startup handshake" - ); - let service_pid = fs::read_to_string(&service_pid_file) - .unwrap() - .parse::() - .unwrap(); - let service_pid = rustix::process::Pid::from_raw(service_pid).unwrap(); + let service_pid = service_pid.expect("stubborn service did not reach its startup handshake"); assert!(refusal.contains("output reader"), "{refusal}"); - assert!(started.elapsed() < Duration::from_secs(5)); + // The service ignores TERM and loops forever, so it is gone only because + // setup cleanup killed its group. assert!(rustix::process::test_kill_process(service_pid).is_err()); } @@ -2617,11 +2557,10 @@ fn service_cleanup_joins_every_log_pump() { let mut children = Children { casework: Some(Service::from_guard(child, vec![pump]).unwrap()), }; - let deadline = Instant::now() + Duration::from_secs(2); - while !children.exited().unwrap() && Instant::now() < deadline { - thread::sleep(Duration::from_millis(5)); - } - assert!(children.exited().unwrap(), "guard did not become waitable"); + assert!( + eventually(|| children.exited().unwrap()), + "guard did not become waitable" + ); children.stop().unwrap(); assert!(joined.load(Ordering::Relaxed)); diff --git a/crates/registry-evidencectl/tests/production_build.rs b/crates/registry-evidencectl/tests/production_build.rs index 733f877221..4e9057b673 100644 --- a/crates/registry-evidencectl/tests/production_build.rs +++ b/crates/registry-evidencectl/tests/production_build.rs @@ -869,6 +869,10 @@ fn identical_inputs_produce_identical_bundle_bytes_revision_and_stable_report_sh #[test] fn termination_cancels_evidence_and_removes_only_current_build_staging() { + // Only a hung build reaches this; a runner starved by parallel builds just + // makes each wait below longer. + const BUILD_EVENT_DEADLINE: Duration = Duration::from_secs(120); + for signal in [rustix::process::Signal::INT, rustix::process::Signal::TERM] { let fixture = Fixture::new(); let ready = fixture.root.join("blocked-evidence-ready"); @@ -885,7 +889,7 @@ fn termination_cancels_evidence_and_removes_only_current_build_staging() { .spawn() .expect("blocking evidencectl build starts"); - let ready_deadline = Instant::now() + Duration::from_secs(10); + let ready_deadline = Instant::now() + BUILD_EVENT_DEADLINE; while !ready.exists() { assert!( child.try_wait().expect("poll blocking build").is_none(), @@ -904,7 +908,7 @@ fn termination_cancels_evidence_and_removes_only_current_build_staging() { .expect("evidencectl PID is positive"); rustix::process::kill_process(pid, signal).expect("send signal to evidencectl build"); - let exit_deadline = Instant::now() + Duration::from_secs(10); + let exit_deadline = Instant::now() + BUILD_EVENT_DEADLINE; loop { if child.try_wait().expect("poll interrupted build").is_some() { break; diff --git a/crates/registry-language-server/tests/evidence_cards_protocol.rs b/crates/registry-language-server/tests/evidence_cards_protocol.rs index 0ca24b1bad..7368bd1267 100644 --- a/crates/registry-language-server/tests/evidence_cards_protocol.rs +++ b/crates/registry-language-server/tests/evidence_cards_protocol.rs @@ -16,13 +16,13 @@ use std::{ fs, io::{BufRead, BufReader, Read, Write}, path::Path, - process::{Child, ChildStdin, Command, Stdio}, + process::{Child, ChildStdin, Command, ExitStatus, Stdio}, sync::{ mpsc::{self, Receiver, RecvTimeoutError}, Arc, Mutex, }, thread::{self, JoinHandle}, - time::Duration, + time::{Duration, Instant}, }; use serde_json::{json, Value}; @@ -626,9 +626,9 @@ async fn a_request_immediately_after_a_change_never_serves_the_edit_it_replaced( // duplicate the small amount of raw framing `tests/protocol.rs` already hand-rolls, because that // module's helpers are private to their own binary and there is no third crate for both to share. -/// How long a test waits for one more message before concluding the server has stopped producing -/// them. -const MESSAGE_TIMEOUT: Duration = Duration::from_secs(10); +/// How long a test waits for the response to one request before concluding the server is hung. +/// Far above any healthy response time, since a busy runner can starve the server for seconds. +const RESPONSE_DEADLINE: Duration = Duration::from_secs(120); fn send(stdin: &mut ChildStdin, message: Value) { let body = serde_json::to_vec(&message).unwrap(); @@ -658,8 +658,8 @@ fn read_one_message(reader: &mut R) -> Option { } /// Reads framed LSP messages off a reader on a background thread and forwards each one to the -/// returned channel, so a caller waiting for a message that never arrives times out on -/// `Receiver::recv_timeout` instead of blocking on the pipe. The join handle lets a caller wait for +/// returned channel, so a caller waiting for a response that never arrives fails at +/// [`RESPONSE_DEADLINE`] instead of blocking on the pipe. The join handle lets a caller wait for /// the thread to observe end-of-stream, which matters for the stdout-hygiene test below: it has to /// know every byte the process ever wrote has been accounted for before it inspects them. fn spawn_message_reader(reader: R) -> (Receiver, JoinHandle<()>) { @@ -675,20 +675,22 @@ fn spawn_message_reader(reader: R) -> (Receiver (receiver, handle) } -fn receive(messages: &Receiver) -> Value { - match messages.recv_timeout(MESSAGE_TIMEOUT) { - Ok(message) => message, - Err(RecvTimeoutError::Timeout) => { - panic!("language server sent nothing for {MESSAGE_TIMEOUT:?}") - } - Err(RecvTimeoutError::Disconnected) => panic!("language server closed stdout"), - } -} - +/// Waits for the response to request `id`, setting aside any notifications that arrive first. fn receive_response(messages: &Receiver, id: i64) -> Value { + let deadline = Instant::now() + RESPONSE_DEADLINE; let mut others = Vec::new(); loop { - let message = receive(messages); + let remaining = deadline.saturating_duration_since(Instant::now()); + let message = match messages.recv_timeout(remaining) { + Ok(message) => message, + Err(RecvTimeoutError::Timeout) => panic!( + "language server did not return response {id} within {RESPONSE_DEADLINE:?}; \ + received instead: {others:?}" + ), + Err(RecvTimeoutError::Disconnected) => panic!( + "language server closed stdout before response {id}; received instead: {others:?}" + ), + }; if message.get("id").and_then(Value::as_i64) == Some(id) { return message; } @@ -725,10 +727,25 @@ fn shut_down(stdin: &mut ChildStdin, stdout: &Receiver, mut child: Child, stdin, json!({"jsonrpc": "2.0", "method": "exit", "params": null}), ); - assert!(child - .wait() - .expect("the server process is waitable") - .success()); + assert!(wait_for_exit(&mut child).success()); +} + +/// Waits for the process to exit within [`RESPONSE_DEADLINE`], killing it and panicking with a +/// clear message otherwise. A server that answers `shutdown` but never acts on `exit` would +/// otherwise hang the test suite on an unbounded [`Child::wait`]. +fn wait_for_exit(child: &mut Child) -> ExitStatus { + let deadline = Instant::now() + RESPONSE_DEADLINE; + loop { + if let Some(status) = child.try_wait().expect("the server process is waitable") { + return status; + } + if Instant::now() >= deadline { + let _ = child.kill(); + let _ = child.wait(); + panic!("language server did not exit within {RESPONSE_DEADLINE:?} after `exit`"); + } + thread::sleep(Duration::from_millis(50)); + } } /// The handshake declares exactly what completion and hover implement, and a client that asks the diff --git a/release/scripts/check-gates-inventory.py b/release/scripts/check-gates-inventory.py index 5641ebf26e..ee3bb4fbc5 100644 --- a/release/scripts/check-gates-inventory.py +++ b/release/scripts/check-gates-inventory.py @@ -120,6 +120,11 @@ "ci-result:\n name: CI result", ), ("Cargo deny", "run: cargo deny check"), + ( + "Standalone Cargo lockfiles", + "cargo update --workspace --locked \\\n" + ' --manifest-path "$(dirname "${lock}")/Cargo.toml"', + ), ( "Platform path filter", "platform: ${{ steps.filter.outputs.platform }}",