//! Per-agent notification pollers — shared library half. //! //! Two binaries ship from this crate, one per notification host: //! //! - **`hive-forge-notify`** — the hive's internal Forgejo. Always //! deployed; the behaviour predates this split and is unchanged. //! - **`hive-github-notify`** — github.com, for agents that have a PAT. //! //! They are separate *binaries* rather than one process with two loops so //! the deployment can choose: an agent module installs the GitHub unit or //! it doesn't, and the decision lives in the module rather than in a cargo //! feature. A feature flag would unify across the workspace — enabling it //! for one consumer changes feature resolution for the whole graph and //! stops the two builds sharing any cached crate — which is a permanent //! cost for something a second binary expresses for free. //! //! Everything except the host-specific calls (list unread, mark read, //! resolve own login) is shared and lives here: classification, wake //! formatting, the tolerant parse, delivery dedupe, and the todo upsert. //! The host differences live behind [`source::Source`]. pub mod notify; pub mod source; /// Re-exported so a binary can spell it once at the crate root alongside /// the other things it needs to build a client; it is a property of the /// poller, not of the notification format. pub use notify::HTTP_TIMEOUT_SECS; /// Retry policy for the harness's in-agent socket. Deliberately fail-fast: /// both callers are inside the poll loop and both treat a failed request as /// "leave the thread unread and try again next tick", so the poll interval /// *is* the retry — a second, in-request backoff would only stack sleeps on /// top of it and delay the rest of the batch. That is the opposite /// trade-off from the serve loop's client, which rides out a hive-c0re /// restart because its callers have no natural retry of their own. pub const TODO_SOCKET_RETRY: hive_sock_client::Retry = hive_sock_client::Retry::None; /// Resolve the harness's in-agent todo socket from the environment, falling /// back to the well-known path. Shared by both binaries so they can't drift /// on where they deliver. #[must_use] pub fn agent_socket() -> std::path::PathBuf { std::env::var_os("HIVE_AGENT_SOCKET").map_or_else( || std::path::PathBuf::from(hive_agent_sock::DEFAULT_AGENT_SOCKET), std::path::PathBuf::from, ) } /// Install the tracing subscriber both binaries use: `RUST_LOG` when set, /// `info` otherwise. pub fn init_tracing() { tracing_subscriber::fmt() .with_env_filter( tracing_subscriber::EnvFilter::try_from_env("RUST_LOG") .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")), ) // This is a systemd-managed daemon — stdout always goes to journald, // never a human terminal, and journald doesn't strip ANSI escapes: // they land in victorialogs as raw byte-array spam otherwise. .with_ansi(false) .init(); } /// Read `HYPERHIVE_STATE_DIR`, the directory holding the agent's /// credentials (`forge-token`, `github-token`). Empty when unset, which /// makes the token paths relative and the read fail — the callers treat /// that as "not configured" and settle. #[must_use] pub fn state_dir() -> String { std::env::var("HYPERHIVE_STATE_DIR").unwrap_or_default() } /// Where this agent's forge token is read from, first match wins: /// `HIVE_FORGE_TOKEN_FILE` (the copy the agent fetched from the swarm secret /// store, set by `nix/agent-modules/forge-token.nix`), then /// `/forge-token` (the file the hive used to write, still the only /// copy on an agent without a store identity). #[must_use] pub fn forge_token_paths(state_dir: &str) -> Vec { let mut paths = Vec::with_capacity(2); if let Ok(fetched) = std::env::var("HIVE_FORGE_TOKEN_FILE") && !fetched.is_empty() { paths.push(std::path::PathBuf::from(fetched)); } paths.push(std::path::Path::new(state_dir).join("forge-token")); paths } /// The first non-empty token among `paths`, trimmed. `None` when none of /// them holds one. #[must_use] pub fn read_first_token(paths: &[std::path::PathBuf]) -> Option { paths.iter().find_map(|p| { let t = std::fs::read_to_string(p).ok()?; let t = t.trim(); (!t.is_empty()).then(|| t.to_owned()) }) } #[cfg(test)] mod token_tests { use super::read_first_token; fn scratch(tag: &str) -> std::path::PathBuf { let ts = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| d.as_nanos()); let dir = std::env::temp_dir().join(format!("hive-forge-notify-token-{tag}-{ts}")); std::fs::create_dir_all(&dir).expect("create scratch dir"); dir } #[test] fn the_fetched_token_wins_over_the_state_file() { let dir = scratch("wins"); let (fetched, state) = (dir.join("fetched"), dir.join("forge-token")); std::fs::write(&fetched, "new\n").expect("write"); std::fs::write(&state, "old\n").expect("write"); assert_eq!(read_first_token(&[fetched, state]).as_deref(), Some("new")); let _ = std::fs::remove_dir_all(&dir); } #[test] fn a_missing_or_empty_fetched_token_falls_back_to_the_state_file() { let dir = scratch("fallback"); let (fetched, state) = (dir.join("fetched"), dir.join("forge-token")); std::fs::write(&state, "old\n").expect("write"); assert_eq!( read_first_token(&[fetched.clone(), state.clone()]).as_deref(), Some("old") ); std::fs::write(&fetched, "\n").expect("write"); assert_eq!(read_first_token(&[fetched, state]).as_deref(), Some("old")); let _ = std::fs::remove_dir_all(&dir); } #[test] fn no_token_anywhere_is_none() { let dir = scratch("none"); assert_eq!(read_first_token(&[dir.join("a"), dir.join("b")]), None); let _ = std::fs::remove_dir_all(&dir); } }