use degenbot_cli_core::CancelHandle;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[must_use]
pub enum Action {
Cancel,
Abort,
}
pub const fn action(already_cancelled: bool) -> Action {
if already_cancelled {
Action::Abort
} else {
Action::Cancel
}
}
#[derive(Debug)]
#[must_use = "the guard must outlive the command run"]
pub struct Guard {
listener: Option<tokio::task::AbortHandle>,
}
impl Drop for Guard {
fn drop(&mut self) {
if let Some(listener) = self.listener.take() {
listener.abort();
}
}
}
#[cfg(unix)]
const fn census_entry() -> degenbot_core::worker_census::WorkerCensusEntry {
degenbot_core::worker_census::WorkerCensusEntry {
resource: "cli_sigint_listener",
kind: "std signal-listener thread (one current-thread runtime; SIGINT -> cooperative CancelHandle)",
count: 1,
thread_name: "degenbot-cli-sigint",
sizing: "exactly one per run (fixed; owned by the run guard)",
binding: "pinned",
}
}
#[cfg(unix)]
pub fn install(cancel: CancelHandle) -> Guard {
degenbot_core::worker_census::register(census_entry());
let Ok(runtime) = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
else {
degenbot_core::op_warn!(
domain = pump,
"could not build the SIGINT listener runtime; cooperative cancel is unavailable"
);
return Guard { listener: None };
};
let (tx, rx) = std::sync::mpsc::channel();
let spawned = std::thread::Builder::new()
.name("degenbot-cli-sigint".to_owned())
.spawn(move || {
let task = runtime.spawn(listen(cancel));
let _ = tx.send(task.abort_handle());
let _ = runtime.block_on(task);
});
if let Err(error) = spawned {
degenbot_core::op_warn!(
domain = pump,
error = %error,
"could not spawn the SIGINT listener thread; cooperative cancel is unavailable"
);
return Guard { listener: None };
}
let listener = rx
.recv()
.map_err(|_| {
degenbot_core::op_warn!(
domain = pump,
"SIGINT listener thread died before installing; cooperative cancel is unavailable"
);
})
.ok();
Guard { listener }
}
#[cfg(not(unix))]
pub fn install(_cancel: CancelHandle) -> Guard {
degenbot_core::op_warn!(
domain = pump,
"SIGINT handling is Unix-only; cooperative cancel is unavailable"
);
Guard { listener: None }
}
#[cfg(unix)]
async fn listen(cancel: CancelHandle) {
let Ok(mut interrupts) =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt())
else {
degenbot_core::op_warn!(
domain = pump,
"could not install the SIGINT handler; cooperative cancel is unavailable"
);
return;
};
loop {
if interrupts.recv().await.is_none() {
return;
}
match action(cancel.is_cancelled()) {
Action::Cancel => {
cancel.cancel();
degenbot_core::op_warn!(
domain = pump,
"interrupt received: finishing the in-flight chunk, then stopping (press Ctrl+C again to abort)"
);
}
Action::Abort => {
degenbot_core::op_warn!(domain = pump, "second interrupt received: aborting");
std::process::abort();
}
}
}
}
#[cfg(test)]
mod tests {
use super::{action, Action};
#[test]
fn first_interrupt_cancels_and_second_aborts() {
assert_eq!(action(false), Action::Cancel);
assert_eq!(action(true), Action::Abort);
}
}