//! Who is acting on this machine, which of them has been told its name, and //! the file that keeps both across a restart. //! //! What a context *is* comes from the harness and what its name is comes from //! the registrar. This decides which of them still needs one, and writes the //! answer down: a daemon that starts with the file present serves the same //! identity to the same context rather than minting it a second account. //! //! The file holds an account's own credential, so it is a `0600` file in the //! `0700` state directory, beside the host's key and under the lock //! [`crate::node`] holds there for the daemon's life. It is replaced whole or //! not at all, through a synced temporary and a rename, and a file that does //! not parse refuses the start rather than being written over. //! //! A binding does not last forever. Every context is dropped once it has //! been quiet for [`Store::lifetime`], which bounds both the map and the //! secrets in it. The account stays where it is at the server: a name is //! never returned to the pool, and nothing here writes to the ledger. use std::collections::{BTreeSet, HashMap}; use std::io; use std::path::{Path, PathBuf}; use serde::{Deserialize, Serialize}; use time::{Duration, OffsetDateTime}; use tracing::{info, warn}; use crate::protocol::{Observed, Report}; use crate::secret::Secret; /// The file, in the daemon's state directory, that bindings are kept in. const CONTEXTS_FILE: &str = "contexts.json"; /// How long a context lasts after its last report unless a deployment says /// otherwise, in days. See [`LIFETIME_DAYS`] for the range. pub const DEFAULT_LIFETIME_DAYS: u64 = 30; /// The day counts `DIDBOT_CONTEXT_TTL_DAYS` accepts. /// /// The floor is a day because a shorter one starts expiring sessions that /// are still running. The ceiling is a year because that is /// `didbot_pds::DEFAULT_AGENT_TOKEN_TTL`: a binding kept past the life of the /// token in it names an account this daemon can no longer act as. pub const LIFETIME_DAYS: std::ops::RangeInclusive = 1..=365; /// How far the clock may move before a report on its own is worth a write. /// /// Every report touches a context, and most reports change nothing else. This /// is the granularity the lifetime above is measured at on disk, not a /// setting: at any value far below that lifetime a restart reads back the /// same answer, and the cost of a smaller one is an fsync per hook. const TOUCH_GRANULARITY: Duration = Duration::hours(1); /// How often the map is swept for contexts past their lifetime. /// /// The granularity the lifetime is reclaimed in, for the same reason as /// `TOUCH_GRANULARITY` above: a sweep walks every context, and one per hour /// is nothing beside a lifetime measured in days. const SWEEP_INTERVAL: Duration = Duration::hours(1); /// A context, named the way the harness names it. /// /// A session with no subagent is a context in its own right, so the pair is /// the key rather than the subagent id alone. Two harnesses could hand out the /// same session id; that is a problem for the day a second one exists. #[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] pub struct Key { /// The harness's identifier for the session. pub session: String, /// Its identifier for the acting context, absent for the session itself. pub context: Option, } /// What to call an asker that did not name itself. /// /// One unnamed asker is the ordinary single-plugin case and behaves as it /// always did; two of them would share a slot, which is the old bug and is /// why a plugin should say who it is. const ANONYMOUS: &str = "-"; /// Which plugin a report came from. fn asker_of(report: &Report) -> String { report.asker.clone().unwrap_or_else(|| ANONYMOUS.to_owned()) } impl Key { fn of(report: &Report) -> Self { Self { session: report.session.clone(), context: report.context.clone(), } } } /// What is known about one context. #[derive(Debug, Clone)] pub struct Context { /// How the harness names it. pub key: Key, /// The harness's word for what kind of context this is. pub kind: Option, /// Its identity, once a registrar has handed one over. pub did: Option, /// The account's own credential, handed over with the identity and kept /// for as long as the context lasts. /// /// The daemon presents it to sign in on this context's behalf and to /// fetch the sign-in decisions the context is being asked about; /// `plan/cred-delivery.md` says the model holds nothing, and this is the /// component holding it instead. It cannot reach an /// [`crate::protocol::Answer`] -- see [`crate::secret`] for how that is /// enforced rather than intended. It does reach the file this module /// writes, which is what a restart reads it back out of. pub token: Option, /// Which askers have put the identity in front of this context. /// /// A set rather than a flag: a machine runs more than one plugin on the /// same events, and each has to be told once. A flag would answer the /// first and silently withhold from the second. pub told: BTreeSet, /// Whether the harness has said this one is finished. Kept rather than /// removed: a name is never returned to the pool, so neither is the row /// that records who held it, until the row's own lifetime runs out. pub done: bool, /// When the harness last said anything about it, which is what its /// lifetime is measured from. pub last_reported: OffsetDateTime, } impl Context { fn new(key: Key, kind: Option, now: OffsetDateTime) -> Self { Self { key, kind, did: None, token: None, told: BTreeSet::new(), done: false, last_reported: now, } } } /// What the daemon should do about a report, decided before anything is done. #[derive(Debug, Clone, PartialEq, Eq)] pub enum Next { /// Nothing. The context is known and has already been told who it is. Nothing, /// This context has no identity yet and is about to need one. Provision(Key), /// It has one and has not heard it. Tell(String), } /// Every context this daemon has seen, and the file it keeps them in. #[derive(Debug)] pub struct Store { contexts: HashMap, /// Where the bindings are written, or nothing for a store that is asked /// to remember only for as long as this process runs. file: Option, lifetime: Duration, /// What the file was written from, so a report that only moves the clock /// costs a write once per `TOUCH_GRANULARITY` rather than every time. written_at: OffsetDateTime, swept_at: OffsetDateTime, } impl Store { /// A store that keeps its contexts for this process only. /// /// For a caller with no state directory to write in. A daemon has one, /// and uses [`Store::open`]. #[must_use] pub fn in_memory(lifetime: Duration) -> Self { let now = OffsetDateTime::now_utc(); Self { contexts: HashMap::new(), file: None, lifetime, written_at: now, swept_at: now, } } /// The bindings in `dir`, and the store that will keep writing them there. /// /// No file is an empty store; the first provisioning creates it. A file /// that does not parse is refused with its path, rather than replaced: /// it holds the credentials for accounts that exist at the server, and /// writing a fresh one over it would abandon every one of them. pub fn open(dir: impl AsRef, lifetime: Duration) -> io::Result { let path = dir.as_ref().join(CONTEXTS_FILE); let mut store = Self::in_memory(lifetime); if let Some(text) = crate::shut::read(&path)? { let disk: OnDisk = serde_json::from_str(&text).map_err(|error| { io::Error::new( io::ErrorKind::InvalidData, format!("{}: {error}", path.display()), ) })?; for binding in disk.contexts { let context = Context::from(binding); store.contexts.insert(context.key.clone(), context); } info!( path = %path.display(), contexts = store.contexts.len(), "read the contexts an earlier run left" ); } store.file = Some(path); store.sweep(OffsetDateTime::now_utc()); Ok(store) } /// How long a context lasts after its last report. #[must_use] pub fn lifetime(&self) -> Duration { self.lifetime } /// Record a report and say what it calls for. /// /// A context is created by whichever report mentions it first. That is /// usually the one announcing it began, but a harness that does not /// announce every context — or announces one late — must not leave an /// acting context unnamed, so any report will do it. pub fn observe(&mut self, report: &Report, now: OffsetDateTime) -> Next { let key = Key::of(report); let mut changed = false; { let context = self.contexts.entry(key.clone()).or_insert_with(|| { changed = true; Context::new(key.clone(), report.kind.clone(), now) }); if context.kind.is_none() && report.kind.is_some() { context.kind.clone_from(&report.kind); changed = true; } context.last_reported = now; // Resting is not ending: a resumed context reports it and then // carries on under the same id. if report.observed == Observed::Ended { changed |= !context.done; context.done = true; } } if changed || now - self.written_at >= TOUCH_GRANULARITY { self.persist(now); } if !matches!(report.observed, Observed::Began | Observed::Acted) { return Next::Nothing; } let asker = asker_of(report); let context = &self.contexts[&key]; match &context.did { None => Next::Provision(key), Some(_) if context.told.contains(&asker) => Next::Nothing, Some(did) => Next::Tell(did.clone()), } } /// Make a binding for `key` if none stands, as a report would, so an /// account can be written into it. A session is created for its first /// subagent this way when the session has not reported yet. pub fn ensure(&mut self, key: &Key, now: OffsetDateTime) { if !self.contexts.contains_key(key) { self.contexts .insert(key.clone(), Context::new(key.clone(), None, now)); self.persist(now); } } /// Write down the identity a registrar handed over, and its credential. pub fn provisioned(&mut self, key: &Key, did: impl Into, token: Secret) { if let Some(context) = self.contexts.get_mut(key) { context.did = Some(did.into()); context.token = Some(token); self.persist(OffsetDateTime::now_utc()); } } /// The context holding one account. /// /// How the daemon gets from a decision -- which names an account and /// never a context -- to the credential it has to present to act on one. /// A DID is issued to exactly one context, so this is a lookup rather /// than a search for the best match. pub fn by_did(&self, did: &str) -> Option<&Context> { self.contexts .values() .find(|context| context.did.as_deref() == Some(did)) } /// Note that one asker has now been handed a context's identity, so that /// asker is answered with silence on its next tool call and any other is /// not. /// /// Called once the answer carrying the identity has left the daemon. Until /// then the asker has not seen it, and a mark made earlier is a mark that /// can be wrong. pub fn told(&mut self, key: &Key, asker: &str) { if let Some(context) = self.contexts.get_mut(key) { if context.told.insert(asker.to_owned()) { self.persist(OffsetDateTime::now_utc()); } } } /// [`Store::sweep`], but no more often than once an hour. /// /// What the report path calls. A sweep walks every context, and a machine /// under load reports many times a second. pub fn sweep_due(&mut self, now: OffsetDateTime) -> Vec { if now - self.swept_at < SWEEP_INTERVAL { return Vec::new(); } self.sweep(now) } /// Drop every context that has been quiet for longer than the lifetime, /// and say which were dropped. /// /// The accounts stay standing at the server. pub fn sweep(&mut self, now: OffsetDateTime) -> Vec { self.swept_at = now; let lifetime = self.lifetime; let mut gone = Vec::new(); self.contexts.retain(|key, context| { let keep = now - context.last_reported < lifetime; if !keep { gone.push(key.clone()); } keep }); if !gone.is_empty() { info!( contexts = gone.len(), days = lifetime.whole_days(), "dropped contexts that have been quiet longer than their lifetime" ); self.persist(now); } gone } /// What is known about one context. pub fn get(&self, key: &Key) -> Option<&Context> { self.contexts.get(key) } /// Every context, in a stable order. The invariant is that they can be /// listed. pub fn all(&self) -> Vec<&Context> { let mut all: Vec<_> = self.contexts.values().collect(); all.sort_by(|a, b| a.key.cmp(&b.key)); all } /// Replace the file with what is in memory now. /// /// A write that fails is said and not raised. The account it names exists /// at the server either way, and refusing the report that provisioned it /// would leave that account with nobody holding it at all; what a failed /// write costs is the restart, which re-provisions as it did before. fn persist(&mut self, now: OffsetDateTime) { let Some(path) = self.file.clone() else { return; }; self.written_at = now; let disk = OnDisk { version: DISK_VERSION, contexts: self.all().iter().map(|c| Binding::from(*c)).collect(), }; let mut text = match serde_json::to_string_pretty(&disk) { Ok(text) => text, Err(error) => { warn!(%error, "could not render the contexts"); return; } }; text.push('\n'); if let Err(error) = crate::shut::write(&path, text.as_bytes()) { warn!( path = %path.display(), %error, "could not write the contexts; a restart will provision again" ); } } } /// The version this crate writes into the file. /// /// Read back but not acted on: every field is optional or defaulted, so a /// file one version old still reads. It is here so that a change which is /// not readable has somewhere to say so. const DISK_VERSION: u32 = 1; /// The file, as written. #[derive(Debug, Serialize, Deserialize)] struct OnDisk { version: u32, contexts: Vec, } /// One context's binding, as written. /// /// A separate type from [`Context`] because of the token. [`Secret`] has no /// `Serialize` on purpose, so that no message can carry one by accident, and /// this is the one place in the crate that takes the value out of it and /// writes it somewhere. See [`crate::secret`]. #[derive(Debug, Serialize, Deserialize)] struct Binding { session: String, #[serde(default, skip_serializing_if = "Option::is_none")] context: Option, #[serde(default, skip_serializing_if = "Option::is_none")] kind: Option, #[serde(default, skip_serializing_if = "Option::is_none")] did: Option, /// The account's own credential, in the clear. The file's mode and the /// directory's are what guard it, as they guard the host's key beside it. #[serde(default, skip_serializing_if = "Option::is_none")] token: Option, #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] told: BTreeSet, #[serde(default, skip_serializing_if = "std::ops::Not::not")] done: bool, #[serde(with = "time::serde::rfc3339")] last_reported: OffsetDateTime, } impl From<&Context> for Binding { fn from(context: &Context) -> Self { Self { session: context.key.session.clone(), context: context.key.context.clone(), kind: context.kind.clone(), did: context.did.clone(), token: context.token.as_ref().map(|t| t.reveal().to_owned()), told: context.told.clone(), done: context.done, last_reported: context.last_reported, } } } impl From for Context { fn from(binding: Binding) -> Self { Self { key: Key { session: binding.session, context: binding.context, }, kind: binding.kind, did: binding.did, token: binding.token.map(Secret::new), told: binding.told, done: binding.done, last_reported: binding.last_reported, } } } #[cfg(test)] mod tests { use super::*; use crate::protocol::VERSION; use crate::scratch::Scratch; use std::os::unix::fs::PermissionsExt; /// Long enough that nothing expires unless a test moves the clock itself. fn lifetime() -> Duration { Duration::days(DEFAULT_LIFETIME_DAYS as i64) } fn store() -> Store { Store::in_memory(lifetime()) } fn report(observed: Observed, context: Option<&str>) -> Report { Report { version: VERSION, observed, session: "s1".into(), context: context.map(Into::into), kind: None, asker: None, call: None, seen_request_uris: Vec::new(), browser_sign_in: None, operator: None, } } fn now() -> OffsetDateTime { OffsetDateTime::now_utc() } #[test] fn a_context_is_told_its_name_once_however_many_calls_it_makes() { let mut store = store(); let key = Key { session: "s1".into(), context: Some("a1".into()), }; assert_eq!( store.observe(&report(Observed::Began, Some("a1")), now()), Next::Provision(key.clone()) ); store.provisioned(&key, "did:web:one.example", Secret::new("tok")); // Provisioning is not telling: until the adapter has actually put the // name in front of the context, every report still owes it one. assert_eq!( store.observe(&report(Observed::Acted, Some("a1")), now()), Next::Tell("did:web:one.example".into()) ); store.told(&key, ANONYMOUS); for _ in 0..3 { assert_eq!( store.observe(&report(Observed::Acted, Some("a1")), now()), Next::Nothing ); } } #[test] fn a_context_that_never_announced_itself_still_gets_a_name() { let mut store = store(); let acted = report(Observed::Acted, Some("quiet")); assert!(matches!(store.observe(&acted, now()), Next::Provision(_))); } #[test] fn a_session_and_its_subagent_are_different_contexts() { let mut store = store(); store.observe(&report(Observed::Began, None), now()); store.observe(&report(Observed::Began, Some("a1")), now()); assert_eq!(store.all().len(), 2); } #[test] fn a_context_keeps_the_credential_it_was_issued_with() { let mut store = store(); let key = Key { session: "s1".into(), context: Some("a1".into()), }; store.observe(&report(Observed::Began, Some("a1")), now()); store.provisioned(&key, "did:web:one.example", Secret::new("agent-token")); // Found by the account, which is all a sign-in decision names. let found = store.by_did("did:web:one.example").expect("the context"); assert_eq!(found.key, key); assert_eq!( found.token.as_ref().map(Secret::reveal), Some("agent-token") ); assert!(store.by_did("did:web:somebody.else").is_none()); // And not printed by the listing that exists to be printed. let printed = format!("{:?}", store.all()); assert!(!printed.contains("agent-token"), "{printed}"); } #[test] fn resting_leaves_a_context_able_to_act_again() { let mut store = store(); let key = Key { session: "s1".into(), context: None, }; store.observe(&report(Observed::Began, None), now()); store.provisioned(&key, "did:web:one.example", Secret::new("tok")); store.told(&key, ANONYMOUS); store.observe(&report(Observed::Rested, None), now()); assert!(!store.get(&key).unwrap().done); store.observe(&report(Observed::Ended, None), now()); assert!(store.get(&key).unwrap().done); // Ended, but still listable and still holding its name. assert_eq!(store.all().len(), 1); assert!(store.get(&key).unwrap().did.is_some()); } /// The restart the memory-only map could not survive: the same context /// reporting again is the one the earlier run named, not a new one to /// mint a second account for. #[test] fn a_restart_serves_the_identity_the_last_run_issued() { let scratch = Scratch::new("contexts-restart"); std::fs::create_dir_all(&scratch.0).unwrap(); let key = Key { session: "s1".into(), context: Some("a1".into()), }; { let mut store = Store::open(&scratch.0, lifetime()).unwrap(); store.observe(&report(Observed::Began, Some("a1")), now()); store.provisioned(&key, "did:web:one.example", Secret::new("agent-token")); store.told(&key, ANONYMOUS); } let mut store = Store::open(&scratch.0, lifetime()).unwrap(); assert_eq!( store.observe(&report(Observed::Acted, Some("a1")), now()), Next::Nothing, "already named, and already told" ); let held = store.by_did("did:web:one.example").expect("the context"); assert_eq!(held.token.as_ref().map(Secret::reveal), Some("agent-token")); // A different asker on the same context is still owed the name. let mut second = report(Observed::Acted, Some("a1")); second.asker = Some("another-plugin".into()); assert_eq!( store.observe(&second, now()), Next::Tell("did:web:one.example".into()) ); } #[test] fn the_file_holding_a_credential_is_shut() { let scratch = Scratch::new("contexts-mode"); std::fs::create_dir_all(&scratch.0).unwrap(); let mut store = Store::open(&scratch.0, lifetime()).unwrap(); store.observe(&report(Observed::Began, Some("a1")), now()); let path = scratch.0.join(CONTEXTS_FILE); let mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777; assert_eq!(mode, 0o600); assert!( !crate::shut::temporary_for(&path).exists(), "nothing half-written is left beside it" ); } /// The file is a credential for every account in it. A daemon that read /// half of one and carried on would abandon the accounts in the half it /// could not read, so it refuses and says which file. #[test] fn a_file_that_does_not_parse_refuses_the_start_by_name() { let scratch = Scratch::new("contexts-corrupt"); std::fs::create_dir_all(&scratch.0).unwrap(); let path = scratch.0.join(CONTEXTS_FILE); std::fs::write(&path, "{\"version\":1,\"contexts\":[{\"session\":}]}").unwrap(); let refused = Store::open(&scratch.0, lifetime()).unwrap_err(); assert_eq!(refused.kind(), io::ErrorKind::InvalidData); assert!( refused.to_string().contains(CONTEXTS_FILE), "{refused}: names the file" ); assert!(path.exists(), "and leaves it alone"); } #[test] fn a_context_past_its_lifetime_is_gone_and_its_secret_with_it() { let scratch = Scratch::new("contexts-expiry"); std::fs::create_dir_all(&scratch.0).unwrap(); let key = Key { session: "s1".into(), context: Some("a1".into()), }; let lifetime = Duration::days(2); let long_ago = now() - Duration::days(9); { let mut store = Store::open(&scratch.0, lifetime).unwrap(); store.observe(&report(Observed::Began, Some("a1")), long_ago); store.provisioned(&key, "did:web:one.example", Secret::new("agent-token")); } // The start sweeps, so the binding is gone before the daemon serves // anything, and gone from the file with it. let store = Store::open(&scratch.0, lifetime).unwrap(); assert!(store.get(&key).is_none()); assert!(store.by_did("did:web:one.example").is_none()); let text = std::fs::read_to_string(scratch.0.join(CONTEXTS_FILE)).unwrap(); assert!(!text.contains("agent-token"), "{text}"); assert!(Store::open(&scratch.0, lifetime).unwrap().all().is_empty()); } /// A daemon that stays up sweeps as it goes, rather than only at a start /// nothing makes it take. #[test] fn a_running_daemon_sweeps_without_being_restarted() { let mut store = Store::in_memory(Duration::days(2)); let key = Key { session: "s1".into(), context: Some("a1".into()), }; store.observe( &report(Observed::Began, Some("a1")), now() - Duration::days(9), ); store.provisioned(&key, "did:web:one.example", Secret::new("agent-token")); // Not on every report: the last sweep was a moment ago. assert!(store.sweep_due(now()).is_empty()); assert_eq!(store.sweep(now()), vec![key.clone()]); assert!(store.get(&key).is_none()); } /// A context that is still reporting is not swept, however old its first /// report was: the lifetime runs from the last one. #[test] fn a_context_still_being_reported_on_outlives_the_lifetime() { let mut store = Store::in_memory(Duration::days(2)); let long_ago = now() - Duration::days(9); store.observe(&report(Observed::Began, Some("a1")), long_ago); store.observe(&report(Observed::Acted, Some("a1")), now()); assert!(store.sweep(now()).is_empty()); assert_eq!(store.all().len(), 1); } }