//! Delivery half of the harness-local reminders store. Polls [`Reminders`] //! for due rows and pushes each as a //! [`hive_sh4re::inbox::DeliveredMessage`] down an mpsc channel the serve loop //! races against the broker long-poll — mirrors `todo_server`'s `Notify` //! wake, but a reminder fire carries real per-row data (message/id), so a //! bare `Notify` doesn't fit; the channel carries the finished message //! instead. //! //! File-path semantics (large-body auto-save, delivery-time persist) //! port the old `hive-c0re::workers::reminder_scheduler` / //! `socket_server::reminders` logic, simplified: this runs *inside* the //! agent's own container now, so the symlink-escape defense that guarded //! `hive-c0re` (running outside the container, writing into a path an //! agent claimed was its own) is moot — the harness IS the agent, //! already sandboxed by the container boundary. Still keeps the cheap //! belt-and-suspenders checks: `file_path` must resolve under //! [`crate::paths::state_dir`] with no `..`/absolute components in the //! relative tail. use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Duration; use tokio::sync::mpsc; use crate::reminders::Reminders; /// Per-tick cap on reminders delivered — mirrors the old c0re /// `REMINDER_BATCH_LIMIT`, same rationale (bound a deep backlog's /// per-tick cost). const REMINDER_BATCH_LIMIT: u64 = 100; /// Poll interval — matches the old c0re `POLL_INTERVAL`. const POLL_INTERVAL: Duration = Duration::from_secs(5); /// Same cap the broker used to enforce on `send`/`ask`/`remind` bodies /// (`hive-c0re`'s `agent_config::limits::MESSAGE_MAX_BYTES`, not /// reachable from here — hive-agent doesn't depend on hive-c0re). /// Duplicated rather than shared: this is the last remaining reminder /// caller of that constant once the c0re-side store is deleted (commit 6). const REMINDER_BODY_MAX_BYTES: usize = 4096; /// Maximum pending (undelivered) reminders this agent may hold at once. /// Exceeding this makes `store` return an error so the caller backs off /// instead of silently flooding the table. Override via /// `HIVE_REMIND_MAX_PENDING_PER_AGENT`; `0` disables the cap. const DEFAULT_REMIND_MAX_PENDING: u64 = 50; fn remind_max_pending() -> u64 { std::env::var("HIVE_REMIND_MAX_PENDING_PER_AGENT") .ok() .and_then(|s| s.trim().parse::().ok()) .unwrap_or(DEFAULT_REMIND_MAX_PENDING) } /// Run the delivery timer. `store: None` (reminders db failed to open at /// boot) parks forever instead of looping — keeps `tx` alive so the /// serve loop's receiver never observes a closed channel (which would /// otherwise resolve immediately every iteration and busy-loop the /// select), while cleanly disabling delivery. pub async fn run( store: Option>, tx: mpsc::UnboundedSender, ) { let Some(store) = store else { tracing::error!("reminders db unavailable — reminder delivery disabled"); std::future::pending::<()>().await; return; }; loop { tick(&store, &tx); tokio::time::sleep(POLL_INTERVAL).await; } } fn tick(store: &Reminders, tx: &mpsc::UnboundedSender) { let now = chrono::Utc::now().timestamp(); let due = match store.due(now, REMINDER_BATCH_LIMIT) { Ok(rows) => rows, Err(e) => { tracing::warn!(error = ?e, "failed to query due reminders"); return; } }; for r in due { let body = prepare_body(&r.message, r.file_path.as_deref()); let dm = hive_sh4re::inbox::DeliveredMessage { from: "reminder".into(), body, id: 0, redelivered: false, in_reply_to: None, }; if tx.send(dm).is_err() { // Receiver (serve loop) is gone — process is shutting down. // Leave the row pending; nothing else to do here. tracing::warn!(reminder_id = r.id, "reminder delivery channel closed"); return; } if let Err(e) = store.mark_delivered(r.id) { tracing::warn!(reminder_id = r.id, error = ?e, "failed to mark reminder delivered"); } } } /// Store a new reminder, applying the pending cap + large-body auto-save /// dance. Returns the new row id, or a caller-ready error string (used /// directly as `Response::Err.message` by the in-agent socket dispatch). pub fn store( store: &Reminders, message: &str, timing: &hive_sh4re::inbox::ReminderTiming, file_path: Option<&str>, ) -> Result { let max = remind_max_pending(); if max > 0 { let pending = store.count_pending().unwrap_or(0); if pending >= max { return Err(format!( "reminder rejected: already {pending} pending reminders (cap {max}). \ Cancel some via `cancel_loose_end` or wait for them to fire before \ scheduling more. Override the cap with `HIVE_REMIND_MAX_PENDING_PER_AGENT`." )); } } let due_at = resolve_due_at(timing).map_err(|e| format!("invalid reminder timing: {e:#}"))?; let (stored_message, stored_path) = prepare_remind_storage(message, file_path)?; store .store(&stored_message, stored_path.as_deref(), due_at) .map_err(|e| format!("failed to store reminder: {e:#}")) } /// Decide what to actually persist in the reminder row (see the old /// `prepare_remind_storage` this ports — same three outcomes: verbatim /// under the cap, auto-saved-to-a-generated-path over cap with no /// caller path, or auto-saved-to-the-caller's-path over cap). fn prepare_remind_storage( message: &str, file_path: Option<&str>, ) -> Result<(String, Option), String> { if message.len() <= REMINDER_BODY_MAX_BYTES { return Ok((message.to_owned(), file_path.map(str::to_owned))); } let req_path = match file_path { Some(p) => p.to_owned(), None => auto_reminder_path(), }; let path = resolve_state_path(&req_path) .map_err(|reason| format!("auto-save path `{req_path}` rejected: {reason}"))?; write_payload(&path, message).map_err(|reason| { format!("auto-save of large reminder body to `{req_path}` failed: {reason}") })?; let hint = format!( "[reminder body of {} bytes auto-saved to `{req_path}`; read with your filesystem tools]", message.len() ); Ok((hint, None)) } /// Generate a fresh auto-save path under this agent's own /// `state/reminders/` dir. fn auto_reminder_path() -> String { let ts_ns = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| d.as_nanos()); crate::paths::state_dir() .join("reminders") .join(format!("auto-{ts_ns}.md")) .to_string_lossy() .into_owned() } /// Build the delivered body for a due reminder: verbatim when /// `file_path` is unset, otherwise persist to that path and return a /// short pointer (falling back to inline delivery with a warning note /// on any rejection/write failure, so the reminder still fires). fn prepare_body(message: &str, file_path: Option<&str>) -> String { let Some(req_path) = file_path else { return message.to_owned(); }; let path = match resolve_state_path(req_path) { Ok(p) => p, Err(reason) => return inline_fallback(req_path, &format!("rejected: {reason}"), message), }; match write_payload(&path, message) { Ok(()) => format!( "reminder body persisted to `{req_path}` ({} bytes); read with your filesystem tools", message.len() ), Err(reason) => inline_fallback(req_path, &reason, message), } } fn inline_fallback(req_path: &str, reason: &str, message: &str) -> String { format!("[reminder file_path '{req_path}' {reason}; delivering body inline]\n\n{message}") } /// Validate `req_path` is absolute, lives under this agent's own /// `state_dir()`, and its relative tail has no traversal / absolute /// components. Returns the (already-real, no host/container /// translation needed in-container) path. fn resolve_state_path(req_path: &str) -> Result { let base = crate::paths::state_dir(); let path = Path::new(req_path); if !path.is_absolute() { return Err(format!("must be absolute (got `{req_path}`)")); } let Ok(rel) = path.strip_prefix(&base) else { return Err(format!( "must live under `{}` (got `{req_path}`)", base.display() )); }; if rel.as_os_str().is_empty() { return Err("file_path must include a filename, not just the state dir".to_owned()); } for comp in rel.components() { match comp { std::path::Component::Normal(_) => {} other => { return Err(format!( "path component `{other:?}` not allowed (no traversal / absolute / root)" )); } } } Ok(path.to_path_buf()) } /// Write `message` to `path`, creating parent dirs as needed. fn write_payload(path: &Path, message: &str) -> Result<(), String> { let Some(parent) = path.parent() else { return Err("internal: path has no parent".to_owned()); }; std::fs::create_dir_all(parent).map_err(|e| format!("parent dir create failed: {e}"))?; std::fs::write(path, message).map_err(|e| format!("write failed: {e}")) } /// Resolve the `due_at` unix timestamp for a `StoreReminder` request. fn resolve_due_at(timing: &hive_sh4re::inbox::ReminderTiming) -> anyhow::Result { use hive_sh4re::inbox::ReminderTiming; match timing { ReminderTiming::InSeconds { seconds } => { let now = std::time::SystemTime::now(); let future = now .checked_add(std::time::Duration::from_secs(*seconds)) .ok_or_else(|| { anyhow::anyhow!("InSeconds overflow: {seconds}s exceeds system time range") })?; let duration = future .duration_since(std::time::UNIX_EPOCH) .map_err(|e| anyhow::anyhow!("system time before UNIX_EPOCH: {e}"))?; i64::try_from(duration.as_secs()) .map_err(|e| anyhow::anyhow!("unix timestamp exceeds i64 range: {e}")) } ReminderTiming::At { unix_timestamp } => Ok(*unix_timestamp), } } #[cfg(test)] mod tests { use super::*; #[test] fn resolve_state_path_rejects_non_absolute() { assert!(resolve_state_path("relative.md").is_err()); } #[test] fn resolve_due_at_in_seconds_is_close_to_now_plus_n() { let due = resolve_due_at(&hive_sh4re::inbox::ReminderTiming::InSeconds { seconds: 60 }).unwrap(); let now = chrono::Utc::now().timestamp(); assert!((due - now - 60).abs() <= 2, "due={due} now={now}"); } #[test] fn resolve_due_at_at_passes_through() { let due = resolve_due_at(&hive_sh4re::inbox::ReminderTiming::At { unix_timestamp: 123_456, }) .unwrap(); assert_eq!(due, 123_456); } #[test] fn prepare_remind_storage_passthrough_under_cap() { let (msg, fp) = prepare_remind_storage("small body", None).unwrap(); assert_eq!(msg, "small body"); assert_eq!(fp, None); } #[test] fn prepare_body_passthrough_when_no_file_path() { assert_eq!(prepare_body("hello world", None), "hello world"); } #[test] fn prepare_body_falls_back_inline_on_bad_path() { let s = prepare_body("payload", Some("/etc/passwd")); assert!(s.starts_with("[reminder file_path '/etc/passwd' rejected:")); assert!(s.contains("payload")); } }