Something went wrong. Try again.
Identities for entities did.bot
agent llm did
Something went wrong. Try again.
27 kB · 727 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728//! 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<u64> = 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<String>,}
/// 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<String>, /// Its identity, once a registrar has handed one over. pub did: Option<String>, /// 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<Secret>, /// 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<String>, /// 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<String>, 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<Key, Context>, /// 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<PathBuf>, 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<Path>, lifetime: Duration) -> io::Result<Self> { 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<String>, 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<Key> { 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<Key> { 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<Binding>,}
/// 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<String>, #[serde(default, skip_serializing_if = "Option::is_none")] kind: Option<String>, #[serde(default, skip_serializing_if = "Option::is_none")] did: Option<String>, /// 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<String>, #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] told: BTreeSet<String>, #[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<Binding> 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); }}