//! Publishing this agent's icon into the swarm's `agent-icons` bucket. //! //! The bytes are the file the harness's own `GET /icon` serves //! ([`ICON_PATH`]), unmodified. The key is the agent name, so the swarm can //! show the icon while the agent is stopped or after it has moved hives. //! //! **Once per harness start.** The icon comes from the agent's NixOS config, //! and a config change is applied by stopping the container, switching it and //! starting it again, so a new icon always arrives with a new harness process. //! An agent with no icon deletes its key, so an icon removed from the config //! also disappears from the swarm. //! //! **Only under the agent's own credential.** That credential is granted //! exactly this agent's key. The hive's shared client is granted none of the //! bucket, so under it this publishes nothing. //! //! The write is a `JetStream` publish straight to the key's subject, the same //! message `kv::Store::put` sends, without opening the bucket first. That //! keeps the grant to the one subject. swarm-controller creates the bucket, //! so a write can arrive before it exists; that write and any other failed //! one are retried with backoff until one is acked. use std::path::Path; use std::time::Duration; use async_nats::jetstream::context::PublishErrorKind; use crate::swarm_queue::{Connection, Presented}; /// Where the agent's NixOS config puts the icon /// (`services.hyperhive.agent.icon`), and the file `GET /icon` serves. Absent /// when none is configured. const ICON_PATH: &str = "/etc/hyperhive/icon.svg"; /// First retry delay after a failed write, doubled per failure up to /// [`MAX_RETRY`]. const FIRST_RETRY: Duration = Duration::from_secs(5); /// Longest wait between retries. const MAX_RETRY: Duration = Duration::from_mins(5); /// The KV header and value that mark an entry deleted. `kv::Store::get` /// answers `None` for an entry carrying them. const KV_OPERATION: &str = "KV-Operation"; const KV_OPERATION_DELETE: &str = "DEL"; /// What to write under this agent's key. #[derive(Debug, PartialEq, Eq)] enum Icon { /// The SVG's bytes. Set(Vec), /// No icon is configured, so the key is deleted. Unset, } /// Read the icon at `path`. A missing file is [`Icon::Unset`]; any other read /// error is returned, and the key is left as it is. fn read_icon(path: &Path) -> std::io::Result { match std::fs::read(path) { Ok(bytes) => Ok(Icon::Set(bytes)), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Icon::Unset), Err(e) => Err(e), } } /// The subject this agent's icon is written to, or `None` under a credential /// that is not granted it. fn subject(presented: &Presented, agent: &str) -> Option { match presented { Presented::Agent => Some(swarm_queue_client::agent_icon::subject(agent)), Presented::Hive { .. } => None, } } /// Start the publish task, if this agent has a queue credential and a label. pub fn spawn() { if !crate::swarm_queue::configured() { return; } let agent = crate::identity::label(); if agent.is_empty() { tracing::warn!("this agent has no label; not publishing its icon upward"); return; } tokio::spawn(run(agent)); } async fn run(agent: String) { let Some(Connection { client, presented }) = crate::swarm_queue::client().await else { return; }; let Some(subject) = subject(&presented, &agent) else { tracing::info!( "connected with the hive's shared client, which may not write an agent's \ icon; not publishing it upward" ); return; }; let icon = match read_icon(Path::new(ICON_PATH)) { Ok(icon) => icon, Err(e) => { tracing::warn!(path = ICON_PATH, error = %e, "reading the agent icon failed; not publishing it"); return; } }; let js = async_nats::jetstream::new(client.clone()); let mut delay = FIRST_RETRY; loop { match write(&client, &js, &subject, &icon).await { Ok(()) => { tracing::info!( subject, set = matches!(icon, Icon::Set(_)), "published this agent's icon" ); return; } Err(Failure::TooLarge(e)) => { tracing::warn!(error = %e, "the agent icon is over the queue's payload limit; not publishing it"); return; } Err(Failure::Retry(e)) => { tracing::warn!(error = e, retry_in = ?delay, "publishing the agent icon failed"); } } tokio::time::sleep(delay).await; delay = (delay * 2).min(MAX_RETRY); } } /// Why a write did not land. enum Failure { /// The payload can never fit; retrying cannot help. TooLarge(async_nats::jetstream::context::PublishError), /// Anything else, rendered for the log. Retry(String), } /// One acked write of `icon` to `subject`. async fn write( client: &async_nats::Client, js: &async_nats::jetstream::Context, subject: &str, icon: &Icon, ) -> Result<(), Failure> { // An unconnected client does not fail a publish, it buffers it, and the // payload limit it checks against is the library default until the // server has announced its own. swarm_queue_client::ensure_connected(client) .map_err(|e| Failure::Retry(swarm_queue_client::chain(&e)))?; let sent = match icon { Icon::Set(bytes) => js.publish(subject.to_owned(), bytes.clone().into()).await, Icon::Unset => { let mut headers = async_nats::HeaderMap::new(); headers.insert(KV_OPERATION, KV_OPERATION_DELETE); js.publish_with_headers(subject.to_owned(), headers, Vec::new().into()) .await } }; let ack = sent.map_err(|e| match e.kind() { PublishErrorKind::MaxPayloadExceeded => Failure::TooLarge(e), _ => Failure::Retry(e.to_string()), })?; // `StreamNotFound` here is the bucket not existing yet. ack.await.map_err(|e| Failure::Retry(e.to_string()))?; Ok(()) } #[cfg(test)] mod tests { use super::{Icon, Presented, read_icon, subject}; #[test] fn under_its_own_credential_an_agent_writes_its_own_key() { assert_eq!( subject(&Presented::Agent, "atlas").as_deref(), Some("$KV.agent-icons.atlas") ); } /// The hive's shared client is granted no icon key, so a publish under it /// would be refused. #[test] fn under_the_hives_shared_client_nothing_is_written() { let hive = Presented::Hive { client_id: "hive-alpha-agent".to_owned(), }; assert_eq!(subject(&hive, "atlas"), None); } #[test] fn a_configured_icon_is_published_byte_for_byte() { let dir = tempfile::tempdir().expect("tempdir"); let path = dir.path().join("icon.svg"); std::fs::write(&path, b"\n").expect("write"); assert_eq!( read_icon(&path).expect("read"), Icon::Set(b"\n".to_vec()) ); } /// No file is the unconfigured state, which clears the key rather than /// leaving a removed icon in the swarm. #[test] fn no_icon_file_clears_the_key() { let dir = tempfile::tempdir().expect("tempdir"); assert_eq!( read_icon(&dir.path().join("icon.svg")).expect("read"), Icon::Unset ); } /// Any other read failure is an error, so a transient one does not delete /// an icon the agent still has. #[test] fn an_unreadable_icon_is_an_error_not_an_unset_icon() { let dir = tempfile::tempdir().expect("tempdir"); assert!(read_icon(dir.path()).is_err()); } }