refactor(async): prefer tokio sleeps and lint std::thread

Add a Clippy disallowed-methods guardrail for std::thread sleep/spawn
and convert the CLI polling paths to tokio::time::sleep so they no
longer block Tokio workers. Keep the intentional OS-thread sites with
narrow #[expect(...)] annotations that explain why std::thread is
required there.
This commit is contained in:
Bryan Helmkamp 2026-04-12 12:36:19 -04:00
parent 2bf35ab184
commit 708c37aed1
22 changed files with 119 additions and 24 deletions

View file

@ -108,6 +108,7 @@ use_self = "warn"
# Project-specific lints
wildcard_imports = "warn"
absolute_paths = "warn"
disallowed_methods = "warn"
[profile.release]
lto = "thin"

View file

@ -1,2 +1,7 @@
absolute-paths-max-segments = 2
absolute-paths-allowed-crates = ["std", "core", "alloc"]
disallowed-methods = [
{ path = "std::thread::sleep", reason = "Prefer tokio::time::sleep on Tokio paths; document intentional blocking sleeps with #[expect(clippy::disallowed_methods, reason = \"...\")]", replacement = "tokio::time::sleep" },
{ path = "std::thread::spawn", reason = "Prefer Tokio task APIs on async paths; document intentional dedicated OS threads with #[expect(clippy::disallowed_methods, reason = \"...\")]" },
{ path = "std::thread::Builder::spawn", reason = "Prefer Tokio task APIs on async paths; document intentional dedicated OS threads with #[expect(clippy::disallowed_methods, reason = \"...\")]" },
]

View file

@ -114,6 +114,10 @@ enum WorkerControlStreamEvent {
Eof,
}
#[expect(
clippy::disallowed_methods,
reason = "Worker control reads blocking stdin on a dedicated OS thread and forwards lines into Tokio."
)]
fn spawn_worker_control_stream(
interviewer: Arc<ControlInterviewer>,
cancel_token: Arc<AtomicBool>,

View file

@ -6,6 +6,7 @@ use fabro_util::printer::Printer;
use fabro_util::terminal::Styles;
use fabro_workflow::records::Conclusion;
use fabro_workflow::run_status::RunStatus;
use tokio::time;
use tracing::info;
use crate::args::{GlobalArgs, WaitArgs};
@ -65,9 +66,9 @@ pub(crate) async fn run(
run_id
);
}
std::thread::sleep(interval.min(dl - now));
time::sleep(interval.min(dl - now)).await;
} else {
std::thread::sleep(interval);
time::sleep(interval).await;
}
};

View file

@ -55,7 +55,7 @@ pub(crate) async fn dispatch(
}) => {
let settings = user_config::load_settings_with_storage_dir(storage_dir.as_deref())?;
let storage_dir = user_config::storage_dir(&settings)?;
stop::execute(&storage_dir, Duration::from_secs(timeout), printer);
stop::execute(&storage_dir, Duration::from_secs(timeout), printer).await;
Ok(())
}
ServerCommand::Status(ServerStatusArgs { storage_dir, json }) => {

View file

@ -1,5 +1,4 @@
use std::path::{Path, PathBuf};
use std::thread;
use std::time::Duration;
use anyhow::{Result, bail};
@ -12,6 +11,7 @@ use fabro_server::serve;
use fabro_server::serve::{DEFAULT_TCP_PORT, ServeArgs};
use fabro_util::printer::Printer;
use fabro_util::terminal::Styles;
use tokio::time;
use super::record;
@ -28,11 +28,11 @@ pub(crate) async fn execute(
if foreground {
Box::pin(execute_foreground(bind, serve_args, storage_dir, styles)).await
} else {
execute_daemon(&bind, &serve_args, &storage_dir, true, printer)
execute_daemon(&bind, &serve_args, &storage_dir, true, printer).await
}
}
pub(crate) fn ensure_server_running_for_storage(
pub(crate) async fn ensure_server_running_for_storage(
storage_dir: &Path,
config_path: &Path,
) -> Result<Bind> {
@ -41,20 +41,20 @@ pub(crate) fn ensure_server_running_for_storage(
}
let bind = Bind::Unix(default_socket_path());
ensure_server_running_with_bind(bind, config_path, storage_dir)
ensure_server_running_with_bind(bind, config_path, storage_dir).await
}
pub(crate) fn ensure_server_running_on_socket(
pub(crate) async fn ensure_server_running_on_socket(
socket_path: &Path,
config_path: &Path,
storage_dir: &Path,
) -> Result<()> {
let bind = Bind::Unix(socket_path.to_path_buf());
let _ = ensure_server_running_with_bind(bind, config_path, storage_dir)?;
let _ = ensure_server_running_with_bind(bind, config_path, storage_dir).await?;
Ok(())
}
fn ensure_server_running_with_bind(
async fn ensure_server_running_with_bind(
bind: Bind,
config_path: &Path,
storage_dir: &Path,
@ -93,7 +93,9 @@ fn ensure_server_running_with_bind(
storage_dir,
false,
Printer::Silent,
) {
)
.await
{
Ok(()) => Ok(bind),
Err(err) => {
if let Some(existing) = record::active_server_record(storage_dir) {
@ -122,7 +124,7 @@ async fn execute_foreground(
storage_dir: PathBuf,
styles: &'static Styles,
) -> Result<()> {
let lock_file = acquire_lock(&storage_dir)?;
let lock_file = acquire_lock(&storage_dir).await?;
let _lock_file = lock_file; // keep alive for the duration
if let Some(existing) = record::active_server_record(&storage_dir) {
@ -171,14 +173,14 @@ async fn execute_foreground(
// Daemon mode
// ---------------------------------------------------------------------------
fn execute_daemon(
async fn execute_daemon(
bind: &BindRequest,
serve_args: &ServeArgs,
storage_dir: &Path,
announce: bool,
printer: Printer,
) -> Result<()> {
let lock_file = acquire_lock(storage_dir)?;
let lock_file = acquire_lock(storage_dir).await?;
let _lock_file = lock_file; // keep alive until function returns
if let Some(existing) = record::active_server_record(storage_dir) {
@ -288,7 +290,7 @@ fn execute_daemon(
bail!("Server exited during startup with status {status}");
}
thread::sleep(poll_interval);
time::sleep(poll_interval).await;
elapsed += poll_interval;
}
@ -306,7 +308,7 @@ fn execute_daemon(
// Helpers
// ---------------------------------------------------------------------------
fn acquire_lock(storage_dir: &Path) -> Result<std::fs::File> {
async fn acquire_lock(storage_dir: &Path) -> Result<std::fs::File> {
let lock_path = Storage::new(storage_dir).server_state().lock_path();
if let Some(parent) = lock_path.parent() {
std::fs::create_dir_all(parent)?;
@ -325,7 +327,7 @@ fn acquire_lock(storage_dir: &Path) -> Result<std::fs::File> {
if elapsed >= timeout {
bail!("timed out waiting for server lock");
}
thread::sleep(poll_interval);
time::sleep(poll_interval).await;
elapsed += poll_interval;
}

View file

@ -1,13 +1,13 @@
use std::path::Path;
use std::thread;
use std::time::Duration;
use fabro_server::bind::Bind;
use fabro_util::printer::Printer;
use tokio::time;
use super::record;
pub(crate) fn execute(storage_dir: &Path, timeout: Duration, printer: Printer) {
pub(crate) async fn execute(storage_dir: &Path, timeout: Duration, printer: Printer) {
let Some(active) = record::active_server_record_details(storage_dir) else {
fabro_util::printerr!(printer, "Server is not running");
std::process::exit(1);
@ -22,13 +22,13 @@ pub(crate) fn execute(storage_dir: &Path, timeout: Duration, printer: Printer) {
if !fabro_proc::process_alive(record.pid) {
break;
}
thread::sleep(poll_interval);
time::sleep(poll_interval).await;
elapsed += poll_interval;
}
if fabro_proc::process_alive(record.pid) {
fabro_proc::sigkill(record.pid);
thread::sleep(Duration::from_millis(100));
time::sleep(Duration::from_millis(100)).await;
}
record::remove_server_record(&active.record_path);

View file

@ -60,7 +60,7 @@ pub(crate) async fn run_uninstall(
return Ok(());
}
execute_uninstall(&inventory, globals.json, printer)
execute_uninstall(&inventory, globals.json, printer).await
}
fn build_inventory(home_root: &Path, storage_dir: &Path) -> Inventory {
@ -226,7 +226,7 @@ struct UninstallResult {
binary_hint: Option<String>,
}
fn execute_uninstall(inventory: &Inventory, json: bool, printer: Printer) -> Result<()> {
async fn execute_uninstall(inventory: &Inventory, json: bool, printer: Printer) -> Result<()> {
let green = console::Style::new().green();
let dim = console::Style::new().dim();
let bold = console::Style::new().bold();
@ -242,7 +242,7 @@ fn execute_uninstall(inventory: &Inventory, json: bool, printer: Printer) -> Res
// Unit 3a: Server stop
if inventory.server_running {
server::stop::execute(&inventory.storage_dir, Duration::from_secs(5), printer);
server::stop::execute(&inventory.storage_dir, Duration::from_secs(5), printer).await;
result.server_stopped = true;
}

View file

@ -113,6 +113,7 @@ pub(crate) async fn connect_server_with_settings(
async fn connect_api_client_bundle(storage_dir: &Path) -> Result<ServerStoreClient> {
let config_path = user_config::active_settings_path(None);
let bind = start::ensure_server_running_for_storage(storage_dir, &config_path)
.await
.with_context(|| format!("Failed to start fabro server for {}", storage_dir.display()))?;
match bind {
Bind::Unix(path) => connect_unix_socket_api_client_bundle(&path).await,
@ -145,6 +146,7 @@ async fn connect_target_api_client_bundle(
&runtime.active_config_path,
&runtime.storage_dir,
)
.await
.with_context(|| format!("Failed to start fabro server for {}", path.display()))?;
connect_unix_socket_api_client_bundle(path).await
}

View file

@ -180,6 +180,10 @@ fn attach_replays_completed_detached_run() {
}
#[test]
#[expect(
clippy::disallowed_methods,
reason = "This sync integration test uses a dedicated stderr reader thread so the child process can stream output concurrently."
)]
fn attach_before_completion_streams_to_finished_state() {
let context = test_context!();
let gate = write_gated_workflow(&context.temp_dir.join("slow.fabro"), "slow", "Run slowly");
@ -284,6 +288,10 @@ fn attach_before_completion_streams_to_finished_state() {
}
#[test]
#[expect(
clippy::disallowed_methods,
reason = "This sync integration test polls logs for a human gate without creating a Tokio runtime."
)]
fn attach_json_errors_without_prompting_for_human_input() {
let context = test_context!();
let workflow = context.temp_dir.join("human-gate.fabro");

View file

@ -66,6 +66,10 @@ fn spawn_worker_process(
cmd.spawn().expect("worker should spawn")
}
#[expect(
clippy::disallowed_methods,
reason = "This sync integration helper polls child exit without requiring a Tokio runtime."
)]
fn wait_for_child_exit(child: &mut Child, timeout: Duration) -> ExitStatus {
let deadline = Instant::now() + timeout;
loop {

View file

@ -337,6 +337,10 @@ fn isolated_server_switches_context_to_separate_daemon() {
}
#[test]
#[expect(
clippy::disallowed_methods,
reason = "This sync integration test uses OS threads to exercise concurrent CLI auto-start behavior across separate processes."
)]
fn concurrent_autostart_converges_on_one_shared_daemon_and_cleans_up() {
fn run_ps_json(
home_dir: &std::path::Path,

View file

@ -219,6 +219,10 @@ fn fast_simple_workflow(context: &TestContext) -> PathBuf {
workflow
}
#[expect(
clippy::disallowed_methods,
reason = "This sync integration helper polls run artifacts after spawning a detached CLI process."
)]
pub(crate) fn setup_detached_dry_run(context: &TestContext) -> RunSetup {
let workflow = fixture("simple.fabro");
let run_id = unique_run_id();
@ -457,6 +461,10 @@ pub(crate) fn write_gated_workflow(path: &Path, name: &str, goal: &str) -> Workf
WorkflowGate { gate_path }
}
#[expect(
clippy::disallowed_methods,
reason = "This sync integration helper polls stored run status without requiring a Tokio runtime."
)]
pub(crate) fn wait_for_status(run_dir: &Path, expected: &[&str]) -> String {
let deadline = Instant::now() + COMMAND_TIMEOUT;
loop {
@ -518,6 +526,10 @@ pub(crate) fn git_filters(context: &TestContext) -> Vec<(String, String)> {
filters
}
#[expect(
clippy::disallowed_methods,
reason = "This sync integration helper polls for the run directory to appear without requiring a Tokio runtime."
)]
pub(crate) fn resolve_run(context: &TestContext, run_id: &str) -> RunSetup {
let deadline = Instant::now() + COMMAND_TIMEOUT;
loop {
@ -674,6 +686,10 @@ pub(crate) fn run_events(run_dir: &Path) -> Vec<EventEnvelope> {
crate::support::parse_event_envelopes(&response)
}
#[expect(
clippy::disallowed_methods,
reason = "This sync integration helper polls stored events without requiring a Tokio runtime."
)]
pub(crate) fn wait_for_event_names(run_dir: &Path, expected: &[&str]) {
let deadline = std::time::Instant::now() + COMMAND_TIMEOUT;

View file

@ -73,6 +73,10 @@ fn timeline_node_names(repo_dir: &Path, run_id: &str) -> Vec<String> {
.collect()
}
#[expect(
clippy::disallowed_methods,
reason = "This sync git integration helper polls until metadata commits become readable without requiring Tokio."
)]
fn build_timeline_when_ready(repo_dir: &Path, run_id: &str) -> RunTimeline {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
@ -91,6 +95,10 @@ fn build_timeline_when_ready(repo_dir: &Path, run_id: &str) -> RunTimeline {
}
}
#[expect(
clippy::disallowed_methods,
reason = "This sync git integration helper retries metadata branch deletion until libgit2 releases its lock."
)]
fn delete_metadata_branch_when_ready(repo_dir: &Path, run_id: &str) {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {

View file

@ -8094,6 +8094,10 @@ timeout = "30s"
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[expect(
clippy::disallowed_methods,
reason = "This test intentionally blocks inside a sync registry factory to simulate slow startup before cancellation."
)]
async fn cancel_before_run_transitions_to_running_returns_empty_attach_stream() {
let state = create_app_state_with_registry_factory(|interviewer| {
std::thread::sleep(std::time::Duration::from_millis(200));

View file

@ -148,6 +148,10 @@ mod tests {
}
#[test]
#[expect(
clippy::disallowed_methods,
reason = "This sync test needs a dedicated OS thread to let the blocking consumer loop flush on its own timer."
)]
fn flushes_on_time_threshold() {
let (tx, rx) = mpsc::channel();
let mid_flushes: Arc<Mutex<Vec<Vec<String>>>> = Arc::new(Mutex::new(Vec::new()));

View file

@ -71,6 +71,10 @@ pub fn init_server() {
init_inner(level, anonymous_id);
}
#[expect(
clippy::disallowed_methods,
reason = "Telemetry uses a long-lived dedicated OS thread for buffered blocking delivery and shutdown joins."
)]
fn init_inner(level: TelemetryLevel, anonymous_id: String) {
let ctx = context::build_context();
let (tx, rx) = mpsc::channel();

View file

@ -150,6 +150,10 @@ mod tests {
#[cfg(unix)]
#[test]
#[expect(
clippy::disallowed_methods,
reason = "This sync test polls for detached child output without creating a Tokio runtime."
)]
fn spawn_detached_unix_creates_marker_file() {
fn wait_for_file(path: &std::path::Path) {
for _ in 0..100 {

View file

@ -255,6 +255,10 @@ fn ensure_parent_dir(path: &Path) {
}
}
#[expect(
clippy::disallowed_methods,
reason = "This sync test helper polls filesystem and flock state without requiring a Tokio runtime."
)]
fn with_session_lock<T>(root: &Path, f: impl FnOnce() -> T) -> T {
let lock_path = session_lock_path(root);
// Retry create-dir + create-file as a unit: another process's
@ -573,6 +577,10 @@ fn server_running(server: &ServerPaths) -> bool {
server_record_pid(&server.storage_dir).is_some_and(fabro_proc::process_alive)
}
#[expect(
clippy::disallowed_methods,
reason = "This sync test helper polls a child server process without requiring a Tokio runtime."
)]
fn wait_for_server_running(server: &ServerPaths) {
let poll = std::time::Duration::from_millis(50);
let timeout = std::time::Duration::from_secs(5);
@ -629,6 +637,10 @@ fn ensure_server_running(fabro_bin: &Path, server: &ServerPaths, config_path: &P
wait_for_server_running(server);
}
#[expect(
clippy::disallowed_methods,
reason = "This sync test helper polls child shutdown during cleanup without requiring a Tokio runtime."
)]
fn stop_test_server(server: &ServerPaths) {
let record_path = server_record_path(&server.storage_dir);
let Some(pid) = server_record_pid(&server.storage_dir) else {

View file

@ -93,6 +93,10 @@ impl EngineServices {
/// Test-only default: empty registry, no hooks, local sandbox at cwd.
#[cfg(test)]
#[expect(
clippy::disallowed_methods,
reason = "This test helper must initialize a current-thread runtime safely from both sync tests and #[tokio::test]."
)]
pub fn test_default() -> Self {
let store = Arc::new(Database::new(
Arc::new(InMemory::new()),

View file

@ -45,6 +45,10 @@ fn test_run_id(label: &str) -> RunId {
RunId::from(Ulid(u128::from(hasher.finish())))
}
#[expect(
clippy::disallowed_methods,
reason = "This helper spins up a dedicated current-thread runtime when called from inside an existing Tokio runtime."
)]
fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error::Error>> {
let run_dir = run_dir.to_path_buf();
let uses_shared_store = run_dir

View file

@ -80,6 +80,10 @@ fn load_checkpoint(path: &Path) -> Result<Checkpoint, Box<dyn std::error::Error>
Ok(serde_json::from_str(&data)?)
}
#[expect(
clippy::disallowed_methods,
reason = "This helper spins up a dedicated current-thread runtime when called from inside an existing Tokio runtime."
)]
fn load_run_checkpoint(run_dir: &Path) -> Result<Checkpoint, Box<dyn std::error::Error>> {
let run_dir = run_dir.to_path_buf();
let uses_shared_store = run_dir