use std::collections::HashSet;
use std::sync::LazyLock;
use std::time::Duration;
use anyhow::Context as _;
use anyhow::Result;
use chrono::{DateTime, Duration as ChronoDuration};
use crate::Role;
use crate::agent::message_router::{AgentJob, MessageKind};
use crate::db;
use crate::util::UnwrapPoison;
crate::define_store! {
pub(crate) static ALARMS: AlarmStore,
expect = "ALARMS not initialized — call init_all_stores() first",
}
crate::columns! {
ALARM_COLUMNS [ALARM] {
ID => "id",
SESSION_ID => "session_id",
USER_NAME => "user_name",
KIND => "kind",
TEXT => "text",
COMMAND => "command",
INTERVAL_SECONDS => "interval_seconds",
NEXT_FIRE_AT => "next_fire_at",
}
}
#[derive(Debug, Clone)]
pub(crate) struct Alarm {
pub id: String,
pub session_id: String,
pub user_name: String,
pub kind: String,
pub text: String,
pub command: Option<String>,
pub interval_seconds: Option<i64>,
pub next_fire_at: String,
}
fn alarm_from_row(row: &db::Row) -> Result<Alarm, ::turso::Error> {
Ok(Alarm {
id: row.get(COL_ALARM_ID)?,
session_id: row.get(COL_ALARM_SESSION_ID)?,
user_name: row.get(COL_ALARM_USER_NAME)?,
kind: row.get(COL_ALARM_KIND)?,
text: row.get(COL_ALARM_TEXT)?,
command: row.get(COL_ALARM_COMMAND)?,
interval_seconds: row.get(COL_ALARM_INTERVAL_SECONDS)?,
next_fire_at: row.get(COL_ALARM_NEXT_FIRE_AT)?,
})
}
const MAX_ACTIVE_ALARMS: i64 = 10;
const MAX_ALARM_COMMAND_CHARS: usize = 2000;
const MAX_PERIOD_SECS: i64 = i64::MAX / 1_000_000_000;
pub(crate) fn format_fire_time(timestamp: &str) -> Result<String> {
let dt = DateTime::parse_from_rfc3339(timestamp)?;
let local = dt
.with_timezone(&chrono::Local)
.format("%Y-%m-%d %H:%M:%S %Z");
let utc = dt.with_timezone(&chrono::Utc).format("%Y-%m-%d %H:%M:%S");
Ok(format!("{local} local time ({utc} UTC)"))
}
pub(crate) async fn add_alarm(
session_id: &str,
user_name: &str,
text: &str,
fire_at: Option<&str>,
interval_seconds: Option<u64>,
command: Option<&str>,
) -> Result<Alarm> {
if let Some(cmd) = command {
anyhow::ensure!(!cmd.trim().is_empty(), "Alarm command must not be empty");
anyhow::ensure!(
cmd.chars().count() <= MAX_ALARM_COMMAND_CHARS,
"Alarm command too long (maximum {MAX_ALARM_COMMAND_CHARS} characters)"
);
}
let normalized_fire_at = fire_at
.map(|f| {
db::parse_utc_timestamp(f)
.map(|dt| dt.to_rfc3339())
.with_context(|| format!("Invalid RFC3339/ISO-8601 fire time: {f}"))
})
.transpose()?;
let (kind, interval_secs, next_fire_at) = match (&normalized_fire_at, interval_seconds) {
(Some(fire), None) => {
anyhow::ensure!(
db::parse_utc_timestamp(fire)? > db::parse_utc_timestamp(&db::now())?,
"Cannot set an alarm for a time in the past"
);
("one-shot", None, fire.clone())
}
(None, Some(interval)) => {
anyhow::ensure!(interval >= 10, "Period must be at least 10 seconds");
let interval_secs = i64::try_from(interval)
.ok()
.filter(|v| *v <= MAX_PERIOD_SECS)
.with_context(|| "Period must be at most 292 years")?;
let start = db::parse_utc_timestamp(&db::now())?;
let next = start
.checked_add_signed(ChronoDuration::seconds(interval_secs))
.unwrap_or(DateTime::<chrono::Utc>::MAX_UTC);
("periodic", Some(interval_secs), next.to_rfc3339())
}
_ => anyhow::bail!("Exactly one of fire_at or interval_seconds must be provided"),
};
let rows = store()
.conn
.query(
"SELECT COUNT(*) FROM alarms WHERE session_id = ?1 AND status = 'active'",
db::params![session_id],
)
.await?;
let count: i64 = rows.first().map(|r| r.get(0)).transpose()?.unwrap_or(0);
anyhow::ensure!(
count < MAX_ACTIVE_ALARMS,
"Alarm limit reached (maximum 10 active alarms)"
);
let id = crate::generate_id();
let created_at = db::now();
store()
.conn
.execute(
"INSERT INTO alarms \
(id, session_id, user_name, kind, text, fire_at, interval_seconds, next_fire_at, status, created_at, command) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'active', ?9, ?10)",
db::params![
id.clone(),
session_id,
user_name,
kind,
text,
normalized_fire_at.clone(),
interval_secs,
next_fire_at.clone(),
created_at.clone(),
command,
],
)
.await?;
Ok(Alarm {
id,
session_id: session_id.to_string(),
user_name: user_name.to_string(),
kind: kind.to_string(),
text: text.to_string(),
command: command.map(str::to_string),
interval_seconds: interval_secs,
next_fire_at,
})
}
pub(crate) async fn list_alarms(session_id: &str) -> Result<Vec<Alarm>> {
let sql = format!(
"SELECT {ALARM_COLUMNS} FROM alarms \
WHERE session_id = ?1 AND status = 'active' \
ORDER BY next_fire_at ASC"
);
let rows = store().conn.query(&sql, db::params![session_id]).await?;
Ok(rows.iter().map(alarm_from_row).collect::<Result<_, _>>()?)
}
pub(crate) async fn list_user_alarms(user_name: &str) -> Result<Vec<Alarm>> {
let sql = format!(
"SELECT {ALARM_COLUMNS} FROM alarms \
WHERE user_name = ?1 AND status = 'active' \
ORDER BY next_fire_at ASC"
);
let rows = store().conn.query(&sql, db::params![user_name]).await?;
Ok(rows.iter().map(alarm_from_row).collect::<Result<_, _>>()?)
}
pub(crate) async fn remove_alarm(session_id: &str, id: &str) -> Result<Option<Alarm>> {
let sql = format!(
"SELECT {ALARM_COLUMNS} FROM alarms \
WHERE id = ?1 AND session_id = ?2 AND status = 'active'"
);
let rows = store()
.conn
.query(&sql, db::params![id, session_id])
.await?;
let Some(row) = rows.first() else {
return Ok(None);
};
let alarm = alarm_from_row(row)?;
store()
.conn
.execute(
"UPDATE alarms SET status = 'removed' WHERE id = ?1 AND session_id = ?2",
db::params![id, session_id],
)
.await?;
Ok(Some(alarm))
}
async fn delete_alarm_after_failure(alarm: &Alarm) -> Result<()> {
store()
.conn
.execute(
"UPDATE alarms SET status = 'removed' WHERE id = ?1 AND session_id = ?2",
db::params![alarm.id.as_str(), alarm.session_id.as_str()],
)
.await?;
Ok(())
}
async fn due_alarms(now: &str) -> Result<Vec<Alarm>> {
let sql = format!(
"SELECT {ALARM_COLUMNS} FROM alarms \
WHERE status = 'active' AND next_fire_at <= ?1 \
ORDER BY next_fire_at ASC"
);
let rows = store().conn.query(&sql, db::params![now]).await?;
Ok(rows.iter().map(alarm_from_row).collect::<Result<_, _>>()?)
}
async fn fire_alarm(alarm: &Alarm) -> Result<()> {
match &alarm.command {
Some(command) => fire_command_alarm(alarm, command).await,
None => fire_plain_alarm(alarm).await,
}
}
async fn fire_plain_alarm(alarm: &Alarm) -> Result<()> {
let now = db::now();
let content = crate::prompt::substitute(
&crate::prompt::load_prompt("alarm_notification.md"),
&[
("{{text}}", &alarm.text),
("{{fire_at}}", &alarm.next_fire_at),
],
);
deliver_alarm_notification(alarm, content).await?;
advance_alarm_state(alarm, &now).await?;
Ok(())
}
async fn deliver_alarm_notification(alarm: &Alarm, content: String) -> Result<()> {
let content = crate::util::scrub_credentials(&content);
let user = &alarm.user_name;
let workspace_name = crate::users::personal_workspace_name(user);
let agent_id = alarm.session_id.clone();
let mut job = AgentJob {
content,
workspace_name,
user_name: user.clone(),
channel: "gui".to_string(),
kind: MessageKind::UserMessage,
role: Role::Assistant,
reply_target: None,
pending_job_id: None,
};
let id = crate::generate_id();
let persisted = match crate::agent::message_router::persist_pending(&job, id.clone()).await {
Ok(()) => true,
Err(e) => {
tracing::warn!(alarm = %alarm.id, error = %e, "Failed to persist alarm delivery — routing best-effort");
false
}
};
if persisted {
job.pending_job_id = Some(id);
}
crate::agent::message_router::route(&agent_id, job);
Ok(())
}
async fn advance_alarm_state(alarm: &Alarm, now: &str) -> Result<()> {
if alarm.kind == "one-shot" {
store()
.conn
.execute(
"UPDATE alarms SET status = 'fired' WHERE id = ?1",
db::params![alarm.id.as_str()],
)
.await?;
} else {
let interval = alarm.interval_seconds.unwrap_or(0);
anyhow::ensure!(
interval > 0,
"Periodic alarm {} has no valid interval",
alarm.id
);
let next = next_periodic_fire(now, &alarm.next_fire_at, interval)?;
store()
.conn
.execute(
"UPDATE alarms SET next_fire_at = ?1 WHERE id = ?2",
db::params![next, alarm.id.as_str()],
)
.await?;
}
Ok(())
}
static COMMANDS_IN_FLIGHT: LazyLock<std::sync::Mutex<HashSet<String>>> =
LazyLock::new(|| std::sync::Mutex::new(HashSet::new()));
async fn fire_command_alarm(alarm: &Alarm, command: &str) -> Result<()> {
if !crate::users::is_admin(&alarm.user_name).await {
return fire_plain_alarm(alarm).await;
}
advance_alarm_state(alarm, &db::now()).await?;
if !claim_in_flight(&alarm.id) {
return Ok(());
}
tokio::spawn(run_alarm_command_task(alarm.clone(), command.to_string()));
Ok(())
}
fn claim_in_flight(alarm_id: &str) -> bool {
COMMANDS_IN_FLIGHT
.lock()
.unwrap_poison()
.insert(alarm_id.to_string())
}
struct InFlightGuard(String);
impl Drop for InFlightGuard {
fn drop(&mut self) {
COMMANDS_IN_FLIGHT.lock().unwrap_poison().remove(&self.0);
}
}
async fn run_alarm_command_task(alarm: Alarm, command: String) {
let _in_flight = InFlightGuard(alarm.id.clone());
let ws = crate::users::personal_workspace_struct(&alarm.user_name);
let outcome = crate::tools::shell::run_raw_command(&ws, &command).await;
let deletion = if outcome.success {
None
} else {
Some(delete_alarm_after_failure(&alarm).await)
};
match alarm_command_notification(&alarm, &command, &outcome, deletion.as_ref()) {
Some(content) => {
if let Err(e) = deliver_alarm_notification(&alarm, content).await {
tracing::warn!(alarm = %alarm.id, error = %e, "Failed to deliver alarm command notification");
}
}
None => {
tracing::info!(alarm = %alarm.id, "alarm command produced no output — staying silent");
}
}
}
fn alarm_command_notification(
alarm: &Alarm,
command: &str,
outcome: &crate::tools::shell::RawCommandOutcome,
deletion: Option<&Result<()>>,
) -> Option<String> {
if outcome.success && !outcome.has_output {
return None;
}
let output_display = if outcome.has_output {
crate::util::truncate_sandwich(
&outcome.output,
crate::util::TOOL_OUTPUT_BUDGET_BYTES,
"alarm command output",
)
} else {
"(no output)".to_string()
};
let status_line = if outcome.success {
format!("The command exited successfully ({}).", outcome.detail)
} else {
format!("The command FAILED ({}).", outcome.detail)
};
let deletion_line = match deletion {
None => String::new(),
Some(Ok(())) => crate::prompt::load_prompt("alarm_command_deletion.md"),
Some(Err(_)) => crate::prompt::load_prompt("alarm_command_deletion_failed.md"),
};
Some(crate::prompt::substitute(
&crate::prompt::load_prompt("alarm_command_notification.md"),
&[
("{{text}}", &alarm.text),
("{{fire_at}}", &alarm.next_fire_at),
("{{command}}", command),
("{{command_status}}", &status_line),
("{{command_deletion}}", &deletion_line),
("{{command_output}}", &output_display),
],
))
}
fn next_periodic_fire(now: &str, next_fire: &str, interval_secs: i64) -> Result<String> {
let now_dt = db::parse_utc_timestamp(now)?;
let next = db::parse_utc_timestamp(next_fire)?;
let elapsed = now_dt.signed_duration_since(next).num_seconds();
let to_skip = elapsed.div_euclid(interval_secs).saturating_add(1);
let add_secs = interval_secs.saturating_mul(to_skip).min(MAX_PERIOD_SECS);
let next = next
.checked_add_signed(ChronoDuration::seconds(add_secs))
.unwrap_or(DateTime::<chrono::Utc>::MAX_UTC);
Ok(next.to_rfc3339())
}
async fn run_alarm_sweep(batch_limit: usize) -> Result<()> {
let due = due_alarms(&db::now()).await?;
let mut fired = 0usize;
let mut failed = 0usize;
for alarm in due.into_iter().take(batch_limit) {
match fire_alarm(&alarm).await {
Ok(()) => fired += 1,
Err(e) => {
failed += 1;
tracing::warn!(alarm = %alarm.id, error = %e, "Failed to fire alarm");
}
}
}
if fired > 0 || failed > 0 {
tracing::info!(fired, failed, "alarm sweep complete");
}
Ok(())
}
pub async fn run_alarm_sweep_loop() {
loop {
if !crate::shutdown::sleep_or_shutdown_or_drain(Duration::from_secs(1)).await {
break;
}
if let Err(e) = run_alarm_sweep(50).await {
tracing::warn!(error = %e, "alarm sweep failed");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn format_fire_time_renders_local_and_utc() {
let out = format_fire_time("2026-08-28T07:30:00+00:00").unwrap();
assert!(
out.ends_with("local time (2026-08-28 07:30:00 UTC)"),
"got: {out}"
);
assert!(!out.contains("UTC UTC"), "double-UTC in output: {out}");
}
#[test]
fn format_fire_time_rejects_invalid() {
assert!(format_fire_time("not-a-time").is_err());
}
#[test]
fn next_periodic_fire_skips_missed_periods_in_one_step() {
let now = "2026-08-28T00:00:25Z";
let next_fire = "2026-08-28T00:00:00Z";
let out = next_periodic_fire(now, next_fire, 10).unwrap();
assert_eq!(out, "2026-08-28T00:00:30+00:00");
}
#[test]
fn next_periodic_fire_saturates_absurd_interval() {
let out =
next_periodic_fire("2026-08-28T00:00:00Z", "2026-08-28T00:00:00Z", i64::MAX).unwrap();
let parsed = db::parse_utc_timestamp(&out).unwrap();
assert!(
parsed > db::parse_utc_timestamp("2026-08-28T00:00:00Z").unwrap(),
"advance must move past now, got {out}"
);
}
#[tokio::test]
async fn add_alarm_rejects_past_one_shot() {
let err = add_alarm(
"session-a",
"alice",
"remind me",
Some("2020-01-01T00:00:00Z"),
None,
None,
)
.await
.unwrap_err();
assert!(err.to_string().contains("past"), "got: {err}");
}
#[tokio::test]
async fn add_alarm_rejects_short_periodic_interval() {
let err = add_alarm("session-a", "alice", "remind me", None, Some(5), None)
.await
.unwrap_err();
assert!(
err.to_string().contains("at least 10 seconds"),
"got: {err}"
);
}
#[tokio::test]
async fn add_alarm_rejects_absurd_periodic_interval() {
let err = add_alarm(
"session-a",
"alice",
"remind me",
None,
Some(u64::MAX),
None,
)
.await
.unwrap_err();
assert!(err.to_string().contains("at most 292 years"), "got: {err}");
}
#[tokio::test]
async fn add_alarm_requires_exactly_one_of_fire_at_or_interval() {
let err = add_alarm("session-a", "alice", "remind me", None, None, None)
.await
.unwrap_err();
assert!(err.to_string().contains("Exactly one"), "got: {err}");
let err = add_alarm(
"session-a",
"alice",
"remind me",
Some("2099-01-01T00:00:00Z"),
Some(60),
None,
)
.await
.unwrap_err();
assert!(err.to_string().contains("Exactly one"), "got: {err}");
}
#[tokio::test]
async fn add_alarm_enforces_active_cap() {
crate::util::test::init_test_stores().await;
let session = "cap-session";
for i in 0..10 {
let fire = format!("2099-01-01T00:00:{i:02}Z");
add_alarm(
session,
"alice",
&format!("reminder {i}"),
Some(&fire),
None,
None,
)
.await
.unwrap();
}
let err = add_alarm(
session,
"alice",
"eleventh",
Some("2099-01-01T00:01:00Z"),
None,
None,
)
.await
.unwrap_err();
assert!(err.to_string().contains("limit reached"), "got: {err}");
}
#[tokio::test]
async fn add_alarm_rejects_blank_command() {
let err = add_alarm(
"session-a",
"alice",
"remind me",
Some("2099-01-01T00:00:00Z"),
None,
Some(" "),
)
.await
.unwrap_err();
assert!(err.to_string().contains("must not be empty"), "got: {err}");
}
#[tokio::test]
async fn add_alarm_rejects_overlong_command() {
let long = "x".repeat(MAX_ALARM_COMMAND_CHARS + 1);
let err = add_alarm(
"session-a",
"alice",
"remind me",
Some("2099-01-01T00:00:00Z"),
None,
Some(long.as_str()),
)
.await
.unwrap_err();
assert!(err.to_string().contains("too long"), "got: {err}");
}
#[tokio::test]
async fn add_alarm_stores_command() {
crate::util::test::init_test_stores().await;
let session = "cmd-session";
let alarm = add_alarm(
session,
"alice",
"remind me",
Some("2099-01-01T00:00:00Z"),
None,
Some("echo hi"),
)
.await
.unwrap();
assert_eq!(alarm.command.as_deref(), Some("echo hi"));
let listed = list_alarms(session).await.unwrap();
assert_eq!(listed.len(), 1, "one command-armed alarm must be listed");
assert_eq!(listed[0].command.as_deref(), Some("echo hi"));
}
fn command_alarm() -> Alarm {
Alarm {
id: "alarm-cmd".to_string(),
session_id: "assistant:alice".to_string(),
user_name: "alice".to_string(),
kind: "periodic".to_string(),
text: "check the thing".to_string(),
interval_seconds: Some(60),
next_fire_at: "2026-09-04T12:00:00+00:00".to_string(),
command: Some("echo hi".to_string()),
}
}
fn raw_outcome(
success: bool,
has_output: bool,
output: &str,
) -> crate::tools::shell::RawCommandOutcome {
crate::tools::shell::RawCommandOutcome {
success,
detail: if success {
"exit status 0".to_string()
} else {
"exit status 2".to_string()
},
has_output,
output: output.to_string(),
}
}
#[test]
fn alarm_command_notification_stays_silent_on_clean_success() {
let alarm = command_alarm();
let outcome = raw_outcome(true, false, "");
assert!(alarm_command_notification(&alarm, "echo hi", &outcome, None).is_none());
}
#[test]
fn alarm_command_notification_wakes_on_output() {
let alarm = command_alarm();
let outcome = raw_outcome(true, true, "all good");
let content = alarm_command_notification(&alarm, "echo hi", &outcome, None).unwrap();
assert!(content.contains("<alarm-notification>"));
assert!(content.contains("check the thing"));
assert!(content.contains("`echo hi`"));
assert!(content.contains("exited successfully (exit status 0)"));
assert!(content.contains("all good"));
assert!(
!content.contains("DELETED"),
"success path must not report a deletion"
);
}
#[test]
fn alarm_command_notification_wakes_on_failure_even_when_empty() {
let alarm = command_alarm();
let outcome = raw_outcome(false, false, "");
let content =
alarm_command_notification(&alarm, "echo hi", &outcome, Some(&Ok(()))).unwrap();
assert!(content.contains("FAILED (exit status 2)"), "got: {content}");
assert!(content.contains("(no output)"), "got: {content}");
assert!(content.contains("DELETED"), "got: {content}");
}
#[test]
fn alarm_command_notification_reports_failed_deletion() {
let alarm = command_alarm();
let outcome = raw_outcome(false, false, "");
let content = alarm_command_notification(
&alarm,
"echo hi",
&outcome,
Some(&Err(anyhow::anyhow!("db down"))),
)
.unwrap();
assert!(content.contains("FAILED"), "got: {content}");
assert!(content.contains("delete"), "got: {content}");
assert!(
!content.contains("was DELETED because the command failed"),
"must not claim a successful deletion"
);
}
#[test]
fn in_flight_claim_guards_periodic_overlap() {
assert!(claim_in_flight("alarm-overlap"));
assert!(!claim_in_flight("alarm-overlap"));
drop(InFlightGuard("alarm-overlap".to_string()));
assert!(claim_in_flight("alarm-overlap"));
drop(InFlightGuard("alarm-overlap".to_string()));
}
#[tokio::test]
async fn fire_alarm_degrades_command_to_plain_for_non_admin_owner() {
crate::util::test::init_test_stores().await;
let _ = crate::agent::message_router::init_global();
let mut rx = crate::agent::message_router::register_agent("assistant:alarm-bob");
let mut alarm = add_alarm(
"assistant:alarm-bob",
"alarm-bob",
"check the deploy",
Some("2099-01-01T00:00:00Z"),
None,
Some("echo secret-run"),
)
.await
.unwrap();
store()
.conn
.execute(
"UPDATE alarms SET next_fire_at = '2020-01-01T00:00:00+00:00' WHERE id = ?1",
db::params![alarm.id.as_str()],
)
.await
.unwrap();
alarm.next_fire_at = "2020-01-01T00:00:00+00:00".to_string();
fire_alarm(&alarm).await.unwrap();
let job = rx.recv().await.expect("degraded plain reminder must route");
assert!(job.content.contains("<alarm-notification>"));
assert!(job.content.contains("check the deploy"));
assert!(!job.content.contains("Command:"), "got: {}", job.content);
assert!(!job.content.contains("secret-run"), "got: {}", job.content);
crate::agent::message_router::unregister_agent("assistant:alarm-bob");
let status: String = store()
.conn
.query(
"SELECT status FROM alarms WHERE id = ?1",
db::params![alarm.id.as_str()],
)
.await
.unwrap()
.first()
.map(|r| r.get(0))
.transpose()
.unwrap()
.expect("alarm row must exist");
assert_eq!(status, "fired", "one-shot must be terminalized");
}
#[tokio::test]
async fn fire_alarm_executes_command_and_delivers_output_for_admin_owner() {
crate::util::test::init_test_stores().await;
let _ = crate::agent::message_router::init_global();
let owner = "alarm-admin-it";
crate::users::USER_STORE
.get()
.expect("user store initialized")
.add_user(owner, Some("full"), crate::Role::Assistant)
.await
.unwrap();
let mut rx = crate::agent::message_router::register_agent("assistant:alarm-admin-it");
let alarm = Alarm {
id: "alarm-admin-run".to_string(),
session_id: "assistant:alarm-admin-it".to_string(),
user_name: owner.to_string(),
kind: "one-shot".to_string(),
text: "poll the thing".to_string(),
interval_seconds: None,
next_fire_at: "2026-09-04T12:00:00+00:00".to_string(),
command: Some("echo hello-alarm-run".to_string()),
};
fire_alarm(&alarm).await.unwrap();
let job = tokio::time::timeout(std::time::Duration::from_secs(30), rx.recv())
.await
.expect("command run must finish in time")
.expect("admin command notification must route");
assert!(job.content.contains("<alarm-notification>"));
assert!(job.content.contains("poll the thing"));
assert!(
job.content.contains("`echo hello-alarm-run`"),
"got: {}",
job.content
);
assert!(
job.content.contains("exited successfully"),
"got: {}",
job.content
);
assert!(
job.content.contains("hello-alarm-run"),
"got: {}",
job.content
);
crate::agent::message_router::unregister_agent("assistant:alarm-admin-it");
}
#[tokio::test]
async fn fire_alarm_deletes_alarm_and_notifies_on_command_failure() {
crate::util::test::init_test_stores().await;
let _ = crate::agent::message_router::init_global();
let owner = "alarm-admin-fail";
crate::users::USER_STORE
.get()
.expect("user store initialized")
.add_user(owner, Some("full"), crate::Role::Assistant)
.await
.unwrap();
let mut rx = crate::agent::message_router::register_agent("assistant:alarm-admin-fail");
store()
.conn
.execute(
"INSERT INTO alarms \
(id, session_id, user_name, kind, text, fire_at, interval_seconds, next_fire_at, status, created_at, command) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'active', ?9, ?10)",
db::params![
"alarm-fail-run",
"assistant:alarm-admin-fail",
owner,
"one-shot",
"poll the thing",
"2020-01-01T00:00:00+00:00",
None::<i64>,
"2020-01-01T00:00:00+00:00",
db::now(),
"exit 3",
],
)
.await
.unwrap();
let alarm = Alarm {
id: "alarm-fail-run".to_string(),
session_id: "assistant:alarm-admin-fail".to_string(),
user_name: owner.to_string(),
kind: "one-shot".to_string(),
text: "poll the thing".to_string(),
interval_seconds: None,
next_fire_at: "2020-01-01T00:00:00+00:00".to_string(),
command: Some("exit 3".to_string()),
};
fire_alarm(&alarm).await.unwrap();
let job = tokio::time::timeout(std::time::Duration::from_secs(30), rx.recv())
.await
.expect("command run must finish in time")
.expect("failed command notification must route");
assert!(job.content.contains("<alarm-notification>"));
assert!(job.content.contains("poll the thing"));
assert!(job.content.contains("FAILED"), "got: {}", job.content);
assert!(job.content.contains("DELETED"), "got: {}", job.content);
crate::agent::message_router::unregister_agent("assistant:alarm-admin-fail");
let status: String = store()
.conn
.query(
"SELECT status FROM alarms WHERE id = ?1",
db::params!["alarm-fail-run"],
)
.await
.unwrap()
.first()
.map(|r| r.get(0))
.transpose()
.unwrap()
.expect("alarm row must exist");
assert_eq!(status, "removed", "failed one-shot must be soft-deleted");
}
}