use anyhow::Result;
use async_nats::jetstream::kv::Store;
use futures::StreamExt;
use kanade_shared::ExecResult;
use kanade_shared::default_paths;
use kanade_shared::kv::{BUCKET_SCRIPT_CURRENT, BUCKET_SCRIPT_STATUS, SCRIPT_STATUS_REVOKED};
use kanade_shared::wire::{
Command, EXIT_REJECTED_UNSIGNED, EXIT_SKIP_DEADLINE, EXIT_SKIP_REVOKED, EXIT_SKIP_STALENESS,
EXIT_SKIP_VERSION_PIN, signature_refusal_result_id,
};
use tracing::{debug, error, info, warn};
use uuid::Uuid;
use crate::admission_ledger::{Ledger, LedgerError, Ticket};
use crate::command_intake::{Delivery, Intake, NoAck, process_delivery};
use crate::outbox;
use crate::process::{ExecOutcome, apply_jitter, run_command_with_start_deadline};
use crate::script_cache::ScriptCache;
use crate::staleness::{StalenessDecision, Tracker, decide as staleness_decide};
pub(crate) struct Sink {
ticket: Option<Ticket>,
reported: std::sync::atomic::AtomicBool,
}
impl Sink {
fn new(ticket: Option<Ticket>) -> Self {
Self {
ticket,
reported: std::sync::atomic::AtomicBool::new(false),
}
}
fn is_ledgered(&self) -> bool {
self.ticket.is_some()
}
fn is_reported(&self) -> bool {
self.reported.load(std::sync::atomic::Ordering::SeqCst)
}
fn result_id(&self) -> String {
match &self.ticket {
Some(t) => t.result_id().to_string(),
None => Uuid::new_v4().to_string(),
}
}
async fn emit(&self, result: ExecResult, note: &'static str) {
let Some(ticket) = &self.ticket else {
enqueue_result_best_effort(result, note);
return;
};
let t = ticket.clone();
let request_id = result.request_id.clone();
self.reported
.store(true, std::sync::atomic::Ordering::SeqCst);
match tokio::task::spawn_blocking(move || t.finish(result)).await {
Ok(Ok(())) => debug!(request_id = %request_id, "{note}"),
Ok(Err(e)) => {
error!(request_id = %request_id, error = %e, "outcome could not be fully recorded")
}
Err(e) => error!(request_id = %request_id, error = %e, "outcome recording task failed"),
}
}
}
fn enqueue_result_best_effort(result: ExecResult, note: &'static str) {
drop(enqueue_result_best_effort_in(
default_paths::data_dir().join("outbox"),
result,
note,
));
}
pub(crate) fn enqueue_result_best_effort_in(
outbox_dir: std::path::PathBuf,
result: ExecResult,
note: &'static str,
) -> tokio::task::JoinHandle<()> {
tokio::task::spawn_blocking(move || {
match outbox::enqueue(&outbox_dir, &result) {
Ok(path) => debug!(
request_id = %result.request_id,
exit_code = result.exit_code,
outbox = %path.display(),
"{note}",
),
Err(e) => warn!(
request_id = %result.request_id,
error = %e,
"outbox enqueue failed (run still completed)",
),
}
})
}
#[allow(clippy::too_many_arguments)]
pub async fn command_loop(
client: async_nats::Client,
pc_id: String,
ledger: std::sync::Arc<Ledger>,
staleness: Tracker,
mut sub: async_nats::Subscriber,
script_cache: ScriptCache,
check_sink: crate::check_cache::CheckSink,
verifier: std::sync::Arc<crate::command_verify::Verifier>,
) {
let jetstream = async_nats::jetstream::new(client.clone());
let script_current = jetstream.get_key_value(BUCKET_SCRIPT_CURRENT).await.ok();
let script_status = jetstream.get_key_value(BUCKET_SCRIPT_STATUS).await.ok();
if script_current.is_none() {
warn!(
bucket = BUCKET_SCRIPT_CURRENT,
"KV bucket missing — version-pinning skipped (the backend creates it at startup; no backend may have started against this broker yet — see `kanade jetstream status`)"
);
}
if script_status.is_none() {
warn!(
bucket = BUCKET_SCRIPT_STATUS,
"KV bucket missing — revoke check skipped (the backend creates it at startup; no backend may have started against this broker yet — see `kanade jetstream status`)"
);
}
while let Some(msg) = sub.next().await {
let delivery = Delivery {
subject: msg.subject.as_str(),
payload: &msg.payload,
headers: crate::command_verify::headers_of(&msg),
addressed: true,
};
let admitted = match process_delivery(&delivery, &verifier, &ledger, &NoAck).await {
Intake::Launch(a) => a,
Intake::Settled(why) => {
debug!(?why, "command delivery settled without launching");
continue;
}
Intake::Held => continue,
};
let client = client.clone();
let pc_id = pc_id.clone();
let cur = script_current.clone();
let sta = script_status.clone();
let staleness = staleness.clone();
let script_cache = script_cache.clone();
let check_sink = check_sink.clone();
tokio::spawn(async move {
if let Err(e) = handle_command(
client,
pc_id,
admitted.cmd,
cur,
sta,
staleness,
script_cache,
check_sink,
CommandSource::Nats,
admitted.envelope_deadline,
Some(admitted.ticket),
)
.await
{
error!(error = %e, "command handler failed");
}
});
}
}
fn outcome_is_retryable(outcome: &ExecOutcome) -> bool {
match outcome {
ExecOutcome::Completed { exit_code, .. } => *exit_code != 0,
ExecOutcome::Timeout { .. } => true,
ExecOutcome::Killed { .. } => false,
}
}
fn retry_note(attempt: u32, exit_code: i32, killed: bool) -> Option<String> {
(attempt > 0).then(|| {
let plural = if attempt == 1 { "retry" } else { "retries" };
if killed {
format!("stopped by remote kill after {attempt} {plural} (#418 on_failure.retry)")
} else if exit_code == 0 {
format!("succeeded after {attempt} {plural} (#418 on_failure.retry)")
} else {
format!("failed after {attempt} {plural} exhausted (#418 on_failure.retry)")
}
})
}
async fn wait_or_killed(
kill: &crate::kill::KillSwitch,
exec_id: Option<&str>,
backoff: std::time::Duration,
) -> bool {
if exec_id.is_none() {
tokio::time::sleep(backoff).await;
return false;
}
tokio::select! {
_ = tokio::time::sleep(backoff) => false,
_ = kill.killed() => true,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CommandOutcome {
Ran { exit_code: i32 },
Skipped,
}
impl CommandOutcome {
pub fn is_success(self) -> bool {
matches!(self, CommandOutcome::Ran { exit_code: 0 })
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CommandSource {
Nats,
LocalScheduler,
}
#[allow(clippy::too_many_arguments)]
async fn command_is_gated(
client: &async_nats::Client,
pc_id: &str,
cmd: &Command,
script_current: Option<&Store>,
script_status: Option<&Store>,
staleness: &Tracker,
source: CommandSource,
envelope_deadline: Option<chrono::DateTime<chrono::Utc>>,
sink: &Sink,
) -> Result<bool> {
match staleness_decide(&cmd.staleness, staleness.staleness(client)) {
StalenessDecision::Proceed => {}
StalenessDecision::Skip { observed, allowed } => {
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
observed_s = observed.as_secs(),
allowed_s = allowed.as_secs(),
"skip: staleness policy (mode=strict) exceeded — broker view too old",
);
publish_staleness_skipped(pc_id, cmd, observed, allowed, sink).await?;
return Ok(true);
}
}
if source == CommandSource::Nats
&& let Some(cur) = script_current
&& let Ok(Some(entry)) = cur.get(&cmd.id).await
{
let expected = String::from_utf8_lossy(&entry).to_string();
if version_pin_rejects(source, Some(&expected), &cmd.version) {
warn!(
cmd_id = %cmd.id,
expected = %expected,
got = %cmd.version,
request_id = %cmd.request_id,
"skip stale command (version mismatch)",
);
publish_version_mismatch_skipped(pc_id, cmd, &expected, sink).await?;
return Ok(true);
}
}
if let Some(sta) = script_status
&& let Ok(Some(entry)) = sta.get(&cmd.id).await
&& String::from_utf8_lossy(&entry) == SCRIPT_STATUS_REVOKED
{
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
"skip revoked command",
);
publish_revoked_skipped(pc_id, cmd, sink).await?;
return Ok(true);
}
#[cfg(not(any(target_os = "windows", target_os = "macos")))]
if !matches!(cmd.run_as, kanade_shared::wire::RunAs::System) {
let now = chrono::Utc::now();
let stderr = format!(
"skipped: run_as {:?} is not supported on Linux agents",
cmd.run_as
);
warn!(cmd_id = %cmd.id, run_as = ?cmd.run_as, "skip: run_as unsupported on this OS");
sink.emit(
skip_result(
pc_id,
cmd,
kanade_shared::wire::EXIT_SKIP_UNSUPPORTED,
stderr,
now,
),
"unsupported run_as skip result enqueued to outbox",
)
.await;
return Ok(true);
}
let now = chrono::Utc::now();
if let Some(deadline) = start_deadline_blocking(cmd.deadline_at, envelope_deadline, now) {
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
%deadline,
%now,
"skip: starting deadline expired",
);
publish_skipped(client, pc_id, cmd, deadline, now, sink).await?;
return Ok(true);
}
Ok(false)
}
#[allow(clippy::too_many_arguments)]
pub async fn handle_command(
client: async_nats::Client,
pc_id: String,
cmd: Command,
script_current: Option<Store>,
script_status: Option<Store>,
staleness: Tracker,
script_cache: ScriptCache,
check_sink: crate::check_cache::CheckSink,
source: CommandSource,
envelope_deadline: Option<chrono::DateTime<chrono::Utc>>,
ticket: Option<Ticket>,
) -> Result<CommandOutcome> {
let sink = Sink::new(ticket);
let for_report = sink.is_ledgered().then(|| (pc_id.clone(), cmd.clone()));
let outcome = handle_command_inner(
client,
pc_id,
cmd,
script_current,
script_status,
staleness,
script_cache,
check_sink,
source,
envelope_deadline,
&sink,
)
.await;
if let (Err(e), Some((pc_id, cmd))) = (&outcome, for_report)
&& !sink.is_reported()
&& !e.is::<LedgerError>()
{
sink.emit(
failure_result(
&pc_id,
&cmd,
&format!("command failed before it produced an outcome: {e:#}"),
),
"failure outcome recorded",
)
.await;
}
outcome
}
fn failure_result(pc_id: &str, cmd: &Command, message: &str) -> ExecResult {
ExecResult {
skipped: Some(false),
..skip_result(pc_id, cmd, -1, message.to_string(), chrono::Utc::now())
}
}
#[allow(clippy::too_many_arguments)]
async fn handle_command_inner(
client: async_nats::Client,
pc_id: String,
mut cmd: Command,
script_current: Option<Store>,
script_status: Option<Store>,
staleness: Tracker,
script_cache: ScriptCache,
check_sink: crate::check_cache::CheckSink,
source: CommandSource,
envelope_deadline: Option<chrono::DateTime<chrono::Utc>>,
sink: &Sink,
) -> Result<CommandOutcome> {
if command_is_gated(
&client,
&pc_id,
&cmd,
script_current.as_ref(),
script_status.as_ref(),
&staleness,
source,
envelope_deadline,
sink,
)
.await?
{
return Ok(CommandOutcome::Skipped);
}
let kill = crate::kill::KillSwitch::arm(Some(&client), cmd.exec_id.as_deref()).await;
tokio::select! {
_ = apply_jitter(&cmd) => {}
_ = kill.killed() => {
sink.emit(
admission_cancelled_result(
&pc_id,
&cmd,
ExecOutcome::Killed {
stdout: String::new(),
stderr: "killed during start jitter".into(),
},
),
"local admission cancellation enqueued",
)
.await;
return Ok(CommandOutcome::Skipped);
}
}
let _local_slot = match crate::concurrency::admit(&kill, &cmd).await {
Ok(permit) => permit,
Err(outcome) => {
sink.emit(
admission_cancelled_result(&pc_id, &cmd, outcome),
"local admission cancellation enqueued",
)
.await;
return Ok(CommandOutcome::Skipped);
}
};
if command_is_gated(
&client,
&pc_id,
&cmd,
script_current.as_ref(),
script_status.as_ref(),
&staleness,
source,
envelope_deadline,
sink,
)
.await?
{
return Ok(CommandOutcome::Skipped);
}
if cmd.script.is_empty()
&& let Some(key) = cmd.script_object.as_deref()
{
let sha = cmd.script_object_sha256.as_deref().ok_or_else(|| {
anyhow::anyhow!(
"Command {request_id} has script_object={key} but no script_object_sha256 \
— wire builder bug",
request_id = cmd.request_id,
)
})?;
match script_cache.resolve(key, sha).await {
Ok(body) => {
debug!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
%key,
sha256 = %sha,
size = body.len(),
"script_object resolved",
);
cmd.script = body;
}
Err(e) => {
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
%key,
sha256 = %sha,
error = %e,
"script_object resolve failed — aborting run",
);
return Err(e);
}
}
}
if envelope_deadline.is_some() {
let now = chrono::Utc::now();
if let Some(deadline) = start_deadline_blocking(None, envelope_deadline, now) {
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
%deadline,
%now,
"skip: envelope start deadline expired before launch",
);
publish_skipped(&client, &pc_id, &cmd, deadline, now, sink).await?;
return Ok(CommandOutcome::Skipped);
}
}
if let Some(ticket) = &sink.ticket {
ticket.mark_launching().map_err(anyhow::Error::from)?;
}
info!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
version = %cmd.version,
exec_id = ?cmd.exec_id,
"executing command",
);
let started_at = chrono::Utc::now();
let result_id = sink.result_id();
let live_handle = crate::live_tail::register(&result_id);
if let Some(exec_id) = cmd.exec_id.as_deref() {
let event = kanade_shared::wire::EventStarted {
result_id: result_id.clone(),
request_id: cmd.request_id.clone(),
exec_id: exec_id.to_string(),
pc_id: pc_id.clone(),
started_at,
manifest_id: cmd.id.clone(),
version: cmd.version.clone(),
};
let events_outbox_dir = default_paths::data_dir().join("events-outbox");
match crate::events_outbox::enqueue(&events_outbox_dir, &event) {
Ok(p) => debug!(
result_id = %result_id,
events_outbox = %p.display(),
"started event enqueued (drain task delivers via JetStream)",
),
Err(e) => warn!(
error = %e,
result_id = %result_id,
"events_outbox enqueue failed; in-flight view will not show this row until ExecResult lands",
),
}
}
let max_retries = cmd.retry.map(|r| r.max).unwrap_or(0);
let backoff = cmd
.retry
.map(|r| std::time::Duration::from_secs(r.backoff_secs));
let mut attempt: u32 = 0;
let mut retry_not_started: Option<String> = None;
let mut prior: Option<ExecOutcome> = None;
let outcome = loop {
let outcome = match run_command_with_start_deadline(
&kill,
&cmd,
Some(live_handle.tail()),
envelope_deadline,
)
.await
{
Ok(o) => o,
Err(e) if e.is::<crate::process::StartDeadlineExpired>() => {
let deadline = envelope_deadline.unwrap_or_else(chrono::Utc::now);
match prior.take() {
Some(previous) => {
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
attempt,
%deadline,
"envelope start deadline expired at launch — retry not started",
);
attempt -= 1;
retry_not_started = Some(format!(
"retry not started: envelope start deadline {deadline} passed"
));
break previous;
}
None => {
let now = chrono::Utc::now();
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
%deadline,
"skip: envelope start deadline expired at launch",
);
publish_skipped(&client, &pc_id, &cmd, deadline, now, sink).await?;
return Ok(CommandOutcome::Skipped);
}
}
}
Err(e) => return Err(e),
};
if !outcome_is_retryable(&outcome) || attempt >= max_retries {
break outcome;
}
attempt += 1;
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
attempt,
max_retries,
backoff_secs = backoff.map(|b| b.as_secs()).unwrap_or(0),
"fire failed; retrying after backoff (#418 on_failure.retry)",
);
if let Some(b) = backoff
&& wait_or_killed(&kill, cmd.exec_id.as_deref(), b).await
{
info!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
attempt,
"remote kill during retry backoff — aborting retries (#418 on_failure.retry)",
);
break ExecOutcome::Killed {
stdout: String::new(),
stderr: String::new(),
};
}
if let Some(deadline) = start_deadline_blocking(None, envelope_deadline, chrono::Utc::now())
{
warn!(
cmd_id = %cmd.id,
request_id = %cmd.request_id,
attempt,
%deadline,
"envelope start deadline expired — retry not started",
);
attempt -= 1;
retry_not_started = Some(format!(
"retry not started: envelope start deadline {deadline} passed"
));
break outcome;
}
prior = Some(outcome);
};
let finished_at = chrono::Utc::now();
let final_killed = matches!(outcome, ExecOutcome::Killed { .. });
drop(live_handle);
let (exit_code, stdout, stderr, status_note) = match outcome {
ExecOutcome::Completed {
exit_code,
stdout,
stderr,
} => (exit_code, stdout, stderr, None),
ExecOutcome::Killed { stdout, stderr } => {
let eid = cmd.exec_id.as_deref().unwrap_or("?");
(
-1,
stdout,
stderr,
Some(format!("killed by remote signal (kill.{eid})")),
)
}
ExecOutcome::Timeout { stdout, stderr } => (
-1,
stdout,
stderr,
Some(format!("timeout after {}s", cmd.timeout_secs)),
),
};
let stderr = [
status_note,
retry_note(attempt, exit_code, final_killed),
retry_not_started,
]
.into_iter()
.flatten()
.fold(stderr, |acc, note| {
if acc.is_empty() {
note
} else {
format!("{acc}\n{note}")
}
});
if let Some(check_hint) = &cmd.check {
let check = if exit_code == 0 {
crate::check_cache::build_check(check_hint, &stdout)
} else {
crate::check_cache::build_check_failed(check_hint, exit_code, &stderr)
};
check_sink.record(check);
}
let bundles = if exit_code == 0 && cmd.collect.is_some() {
let js = async_nats::jetstream::new(client.clone());
crate::collect::maybe_collect(&js, &client, &cmd, &pc_id, &result_id, &stdout, finished_at)
.await
} else {
Vec::new()
};
let finalize_json = cmd
.collect
.as_ref()
.map(|_| crate::finalize::collect_result_json(&bundles));
let collect_object = bundles.first().map(|b| b.key.clone());
let stdout = if exit_code == 0
&& matches!(
cmd.emit.as_ref().map(|e| e.kind),
Some(kanade_shared::manifest::EmitKind::Events),
) {
forward_obs_events(stdout, pc_id.clone()).await;
String::new()
} else {
stdout
};
let result = ExecResult {
result_id: result_id.clone(),
request_id: cmd.request_id.clone(),
exec_id: cmd.exec_id.clone(),
parent_result_id: None,
pc_id: pc_id.clone(),
exit_code,
skipped: Some(false),
stdout,
stderr,
started_at,
finished_at,
stdout_object: None,
stderr_object: None,
manifest_id: Some(cmd.id.clone()),
collect_object,
};
sink.emit(
result,
"result enqueued to outbox (drain task delivers via JetStream)",
)
.await;
if exit_code == 0 {
crate::local_scheduler::record_job_success(&cmd.id, finished_at).await;
}
if exit_code == 0
&& let Some(fin) = cmd.finalize.as_ref()
&& !fin.on_each_bundle
{
crate::finalize::run_finalize(
&client,
&cmd,
fin,
&pc_id,
&result_id,
None,
finalize_json.as_deref(),
)
.await;
}
let _ = client;
Ok(CommandOutcome::Ran { exit_code })
}
async fn forward_obs_events(stdout: String, pc_id: String) {
use kanade_shared::wire::ObsEvent;
let obs_outbox_dir = default_paths::data_dir().join("obs-outbox");
if let Err(e) = crate::obs_outbox::ensure_outbox_dir(&obs_outbox_dir) {
warn!(error = %e, "obs: ensure_outbox_dir failed; aborting forward");
return;
}
let pc_id_log = pc_id.clone();
let (ok, bad) = tokio::task::spawn_blocking(move || {
let mut ok = 0usize;
let mut bad = 0usize;
for (i, raw) in stdout.lines().enumerate() {
let trimmed = raw.trim();
if trimmed.is_empty() {
continue;
}
let mut event: ObsEvent = match serde_json::from_str(trimmed) {
Ok(e) => e,
Err(e) => {
warn!(
line_no = i + 1,
error = %e,
"obs: stdout line is not a valid ObsEvent JSON; skipping",
);
bad += 1;
continue;
}
};
event.pc_id = pc_id.clone();
if let Err(e) = crate::obs_outbox::enqueue(&obs_outbox_dir, &event) {
warn!(
line_no = i + 1,
error = %e,
"obs: enqueue to outbox failed; line dropped",
);
bad += 1;
} else {
ok += 1;
}
}
(ok, bad)
})
.await
.unwrap_or_else(|e| {
warn!(error = %e, "obs: forwarder task panicked / cancelled");
(0, 0)
});
debug!(ok, bad, pc_id = %pc_id_log, "obs: forwarded NDJSON stdout to obs-outbox");
}
fn start_deadline_blocking(
command: Option<chrono::DateTime<chrono::Utc>>,
envelope: Option<chrono::DateTime<chrono::Utc>>,
now: chrono::DateTime<chrono::Utc>,
) -> Option<chrono::DateTime<chrono::Utc>> {
let effective = match (command, envelope) {
(Some(a), Some(b)) => Some(a.min(b)),
(a, b) => a.or(b),
};
effective.filter(|d| should_skip_for_deadline(*d, now))
}
fn should_skip_for_deadline(
deadline: chrono::DateTime<chrono::Utc>,
now: chrono::DateTime<chrono::Utc>,
) -> bool {
now > deadline
}
fn version_pin_rejects(source: CommandSource, pinned: Option<&str>, cmd_version: &str) -> bool {
source == CommandSource::Nats && matches!(pinned, Some(p) if p != cmd_version)
}
async fn publish_staleness_skipped(
pc_id: &str,
cmd: &Command,
observed: std::time::Duration,
allowed: std::time::Duration,
sink: &Sink,
) -> Result<()> {
let now = chrono::Utc::now();
let stderr = format!(
"skipped: staleness policy (mode=strict) exceeded — agent has been disconnected for {}, max allowed {}",
humantime::format_duration(observed),
humantime::format_duration(allowed),
);
sink.emit(
skip_result(pc_id, cmd, EXIT_SKIP_STALENESS, stderr, now),
"staleness-skip result enqueued to outbox",
)
.await;
Ok(())
}
fn skip_result(
pc_id: &str,
cmd: &Command,
exit_code: i32,
stderr: String,
now: chrono::DateTime<chrono::Utc>,
) -> ExecResult {
ExecResult {
result_id: Uuid::new_v4().to_string(),
request_id: cmd.request_id.clone(),
exec_id: cmd.exec_id.clone(),
parent_result_id: None,
pc_id: pc_id.to_string(),
exit_code,
skipped: Some(true),
stdout: String::new(),
stderr,
started_at: now,
finished_at: now,
stdout_object: None,
stderr_object: None,
manifest_id: Some(cmd.id.clone()),
collect_object: None,
}
}
fn admission_cancelled_result(pc_id: &str, cmd: &Command, outcome: ExecOutcome) -> ExecResult {
let now = chrono::Utc::now();
match outcome {
ExecOutcome::Completed {
exit_code, stderr, ..
} => skip_result(pc_id, cmd, exit_code, stderr, now),
ExecOutcome::Killed { stderr, .. } | ExecOutcome::Timeout { stderr, .. } => ExecResult {
skipped: Some(false),
..skip_result(pc_id, cmd, -1, stderr, now)
},
}
}
pub(crate) fn publish_signature_refused(
outbox_dir: std::path::PathBuf,
pc_id: &str,
cmd: &Command,
reason: &str,
) -> tokio::task::JoinHandle<()> {
enqueue_result_best_effort_in(
outbox_dir,
signature_refusal_result(pc_id, cmd, reason, chrono::Utc::now()),
"signature-refusal result enqueued to outbox",
)
}
pub(crate) fn signature_refusal_result(
pc_id: &str,
cmd: &Command,
reason: &str,
now: chrono::DateTime<chrono::Utc>,
) -> ExecResult {
ExecResult {
result_id: signature_refusal_result_id(&cmd.request_id, pc_id),
skipped: Some(false),
..skip_result(
pc_id,
cmd,
EXIT_REJECTED_UNSIGNED,
format!("refused: {reason}"),
now,
)
}
}
async fn publish_skipped(
_client: &async_nats::Client,
pc_id: &str,
cmd: &Command,
deadline: chrono::DateTime<chrono::Utc>,
now: chrono::DateTime<chrono::Utc>,
sink: &Sink,
) -> Result<()> {
let lateness = now - deadline;
let stderr = format!(
"skipped: starting deadline expired {} ago (deadline {}, received {})",
humantime::format_duration(
lateness
.to_std()
.unwrap_or(std::time::Duration::from_secs(0))
),
deadline,
now,
);
sink.emit(
skip_result(pc_id, cmd, EXIT_SKIP_DEADLINE, stderr, now),
"synthetic skipped-result enqueued to outbox",
)
.await;
Ok(())
}
async fn publish_version_mismatch_skipped(
pc_id: &str,
cmd: &Command,
expected: &str,
sink: &Sink,
) -> Result<()> {
let now = chrono::Utc::now();
let stderr = format!(
"skipped: version-pin mismatch — script_current[{}] = {expected}, command brought {}",
cmd.id, cmd.version,
);
sink.emit(
skip_result(pc_id, cmd, EXIT_SKIP_VERSION_PIN, stderr, now),
"version-mismatch skip result enqueued to outbox",
)
.await;
Ok(())
}
async fn publish_revoked_skipped(pc_id: &str, cmd: &Command, sink: &Sink) -> Result<()> {
let now = chrono::Utc::now();
let stderr = format!(
"skipped: command was revoked (script_status[{}] = revoked)",
cmd.id,
);
sink.emit(
skip_result(pc_id, cmd, EXIT_SKIP_REVOKED, stderr, now),
"revoked skip result enqueued to outbox",
)
.await;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeZone;
fn at(secs: i64) -> chrono::DateTime<chrono::Utc> {
chrono::Utc
.timestamp_opt(1_700_000_000 + secs, 0)
.single()
.unwrap()
}
#[test]
fn a_repeated_refusal_lands_on_the_same_result_row() {
let a = signature_refusal_result_id("req-1", "PC1");
let b = signature_refusal_result_id("req-1", "PC1");
assert_eq!(a, b, "the projector collapses on result_id; it must repeat");
assert_ne!(a, signature_refusal_result_id("req-2", "PC1"));
assert_ne!(a, signature_refusal_result_id("req-1", "PC2"));
assert_eq!(
a,
Uuid::new_v5(&Uuid::NAMESPACE_OID, b"req-1|PC1|signature-refused").to_string()
);
assert_eq!(Uuid::parse_str(&a).unwrap().get_version_num(), 5);
}
#[test]
fn not_now_results_are_skipped_but_a_refusal_or_queued_kill_is_not() {
let cmd = Command {
id: "job".into(),
version: "1".into(),
request_id: "req-1".into(),
exec_id: Some("exec-1".into()),
shell: kanade_shared::wire::Shell::Powershell,
script: String::new(),
script_object: None,
script_object_sha256: None,
timeout_secs: 60,
bypass_local_limit: false,
jitter_secs: None,
run_as: Default::default(),
cwd: None,
deadline_at: None,
staleness: Default::default(),
emit: None,
check: None,
collect: None,
retry: None,
finalize: None,
};
for code in [
EXIT_SKIP_VERSION_PIN,
EXIT_SKIP_DEADLINE,
EXIT_SKIP_REVOKED,
EXIT_SKIP_STALENESS,
] {
let r = skip_result("PC1", &cmd, code, "why".into(), at(0));
assert_eq!(r.skipped, Some(true), "exit {code} is published as skipped");
assert_eq!(r.exit_code, code);
}
let refused = signature_refusal_result("PC1", &cmd, "unsigned", at(0));
assert_eq!(refused.skipped, Some(false));
assert_eq!(refused.exit_code, EXIT_REJECTED_UNSIGNED);
assert_eq!(refused.stderr, "refused: unsigned");
assert!(refused.is_signature_refusal());
let expired = admission_cancelled_result(
"PC1",
&cmd,
ExecOutcome::Completed {
exit_code: EXIT_SKIP_DEADLINE,
stdout: String::new(),
stderr: "deadline".into(),
},
);
assert_eq!(expired.skipped, Some(true));
assert_eq!(expired.exit_code, EXIT_SKIP_DEADLINE);
let killed = admission_cancelled_result(
"PC1",
&cmd,
ExecOutcome::Killed {
stdout: String::new(),
stderr: "killed".into(),
},
);
assert_eq!(killed.skipped, Some(false));
assert_eq!(killed.exit_code, -1);
}
#[test]
fn command_outcome_is_success_only_for_clean_run() {
assert!(CommandOutcome::Ran { exit_code: 0 }.is_success());
assert!(!CommandOutcome::Ran { exit_code: 1 }.is_success());
assert!(!CommandOutcome::Ran { exit_code: 124 }.is_success());
assert!(!CommandOutcome::Skipped.is_success());
}
#[test]
fn version_pin_rejects_stale_nats_command() {
assert!(version_pin_rejects(
CommandSource::Nats,
Some("0.2.0"),
"0.2.1"
));
}
#[test]
fn version_pin_allows_matching_nats_command() {
assert!(!version_pin_rejects(
CommandSource::Nats,
Some("0.2.1"),
"0.2.1"
));
}
#[test]
fn version_pin_allows_nats_command_with_no_kv_entry() {
assert!(!version_pin_rejects(CommandSource::Nats, None, "0.2.1"));
}
#[test]
fn version_pin_never_rejects_local_scheduler_fire() {
assert!(!version_pin_rejects(
CommandSource::LocalScheduler,
Some("0.2.0"),
"0.2.1"
));
assert!(!version_pin_rejects(
CommandSource::LocalScheduler,
Some("0.2.1"),
"0.2.1"
));
assert!(!version_pin_rejects(
CommandSource::LocalScheduler,
None,
"0.2.1"
));
}
#[test]
fn the_earlier_of_the_two_start_deadlines_blocks() {
assert_eq!(start_deadline_blocking(None, None, at(1_000)), None);
assert_eq!(start_deadline_blocking(Some(at(100)), None, at(99)), None);
assert_eq!(
start_deadline_blocking(Some(at(100)), None, at(101)),
Some(at(100))
);
assert_eq!(
start_deadline_blocking(Some(at(500)), Some(at(100)), at(101)),
Some(at(100))
);
assert_eq!(
start_deadline_blocking(Some(at(100)), Some(at(500)), at(101)),
Some(at(100))
);
}
#[test]
fn an_envelope_that_expires_during_the_jitter_wait_is_caught_by_the_second_check() {
let deadline = Some(at(100));
assert_eq!(start_deadline_blocking(None, deadline, at(50)), None);
assert_eq!(
start_deadline_blocking(None, deadline, at(110)),
Some(at(100))
);
assert_eq!(start_deadline_blocking(None, deadline, at(100)), None);
}
#[test]
fn now_strictly_before_deadline_runs() {
assert!(!should_skip_for_deadline(at(100), at(99)));
}
#[test]
fn now_one_second_before_deadline_runs() {
assert!(!should_skip_for_deadline(at(100), at(99)));
}
#[test]
fn now_exactly_at_deadline_still_runs() {
assert!(!should_skip_for_deadline(at(100), at(100)));
}
#[test]
fn now_one_second_past_deadline_skips() {
assert!(should_skip_for_deadline(at(100), at(101)));
}
#[test]
fn now_long_past_deadline_skips() {
assert!(should_skip_for_deadline(at(100), at(86400)));
}
fn completed(exit_code: i32) -> ExecOutcome {
ExecOutcome::Completed {
exit_code,
stdout: String::new(),
stderr: String::new(),
}
}
#[test]
fn clean_exit_is_not_retryable() {
assert!(!outcome_is_retryable(&completed(0)));
}
#[test]
fn nonzero_exit_is_retryable() {
assert!(outcome_is_retryable(&completed(1)));
assert!(outcome_is_retryable(&completed(-1)));
}
#[test]
fn timeout_is_retryable() {
assert!(outcome_is_retryable(&ExecOutcome::Timeout {
stdout: String::new(),
stderr: String::new(),
}));
}
#[tokio::test]
async fn local_kill_interrupts_retry_backoff_without_a_broker() {
let kill = crate::kill::KillSwitch::arm(None, Some("cmd-backoff-kill")).await;
let waiter = tokio::spawn(async move {
wait_or_killed(
&kill,
Some("cmd-backoff-kill"),
std::time::Duration::from_secs(300),
)
.await
});
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
assert!(crate::kill::trigger_local("cmd-backoff-kill"));
let killed = tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
.await
.expect("backoff interrupted")
.unwrap();
assert!(killed);
}
#[tokio::test]
async fn backoff_without_a_kill_runs_its_full_length() {
let kill = crate::kill::KillSwitch::inert();
assert!(!wait_or_killed(&kill, None, std::time::Duration::from_millis(10)).await);
}
#[tokio::test]
#[ignore = "requires a live nats-server"]
async fn remote_kill_interrupts_retry_backoff() {
let client = crate::kill::broker_test::connect().await;
let kill = crate::kill::KillSwitch::arm(Some(&client), Some("cmd-backoff-remote")).await;
let waiter = tokio::spawn(async move {
wait_or_killed(
&kill,
Some("cmd-backoff-remote"),
std::time::Duration::from_secs(300),
)
.await
});
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
crate::kill::broker_test::publish_kill(&client, "cmd-backoff-remote").await;
let killed = tokio::time::timeout(std::time::Duration::from_secs(2), waiter)
.await
.expect("backoff interrupted")
.unwrap();
assert!(killed);
}
#[test]
fn remote_kill_is_never_retried() {
assert!(!outcome_is_retryable(&ExecOutcome::Killed {
stdout: String::new(),
stderr: String::new(),
}));
}
#[test]
fn no_retry_emits_no_note() {
assert_eq!(retry_note(0, 0, false), None);
assert_eq!(retry_note(0, 1, false), None);
}
#[test]
fn retry_note_reports_eventual_success() {
let note = retry_note(2, 0, false).expect("a retry happened");
assert!(note.contains("succeeded after 2 retries"), "got: {note}");
}
#[test]
fn retry_note_reports_exhaustion() {
let note = retry_note(3, 1, false).expect("a retry happened");
assert!(
note.contains("failed after 3 retries exhausted"),
"got: {note}"
);
}
#[test]
fn retry_note_killed_is_not_exhausted() {
let note = retry_note(2, -1, true).expect("a retry happened");
assert!(
note.contains("stopped by remote kill after 2 retries"),
"got: {note}"
);
assert!(!note.contains("exhausted"), "got: {note}");
}
#[test]
fn retry_note_singular_for_one_retry() {
let note = retry_note(1, 0, false).expect("a retry happened");
assert!(note.contains("after 1 retry"), "got: {note}");
assert!(!note.contains("retries"), "got: {note}");
}
}