diff --git a/rust/Cargo.lock b/rust/Cargo.lock index 710e30d2..f5093228 100644 --- a/rust/Cargo.lock +++ b/rust/Cargo.lock @@ -457,6 +457,7 @@ checksum = "3d52eff69cd5e647efe296129160853a42795992097e8af39800e1060caeea9b" name = "conversation-protocol" version = "0.1.0" dependencies = [ + "ed25519-dalek", "flate2", "git-locator", "serde", diff --git a/rust/crates/conversation-protocol/Cargo.toml b/rust/crates/conversation-protocol/Cargo.toml index c0ba3b2d..2678237c 100644 --- a/rust/crates/conversation-protocol/Cargo.toml +++ b/rust/crates/conversation-protocol/Cargo.toml @@ -9,6 +9,9 @@ serde = { version = "1", features = ["derive"] } serde_json = "1" sha1 = "0.10" sha2 = "0.10" +# Optional: only what signs or verifies a ref write needs it +# (design/ref-writers.md); a std tool that pushes with a run token does not. +ed25519-dalek = { version = "3", optional = true } [features] # On by default so the dependency artifact cargo builds for caos-cli (which asks @@ -17,6 +20,7 @@ sha2 = "0.10" default = ["git-cli"] git-cli = [] memory-store = [] +ed25519 = ["dep:ed25519-dalek"] [dependencies.git-locator] path = "../git-locator" diff --git a/rust/crates/conversation-protocol/src/v3/git_store.rs b/rust/crates/conversation-protocol/src/v3/git_store.rs index 001c36ca..fd323b26 100644 --- a/rust/crates/conversation-protocol/src/v3/git_store.rs +++ b/rust/crates/conversation-protocol/src/v3/git_store.rs @@ -36,10 +36,15 @@ pub struct GitStore { remote_tip: RefCell>, batch: RefCell>, batch_dirty: Cell, + push_auth: Option, #[cfg(test)] command_count: Cell, } +/// Proves a push to a governed ref (design/ref-writers.md): the `caos-auth` +/// push option for the commands being pushed. +pub type PushAuth = Box String>; + static NEXT_OBJECT_TEMP: AtomicU64 = AtomicU64::new(1); fn quote_config_value(value: &str) -> Result { @@ -144,7 +149,20 @@ impl GitStore { ); fs::write(dir.join("config"), config) .map_err(|error| format!("writing scratch config: {error}"))?; - GitStore::open(&dir, Some("origin")) + let mut store = GitStore::open(&dir, Some("origin"))?; + // A scratch store is a worker's; a job the server granted writes + // pushes with the token it injected, and any other job has none. + if let Some(token) = super::writers::injected_token() { + let option = super::writers::Auth::Run { token }.option(); + store.set_push_auth(Box::new(move |_: &[super::writers::Command]| { + option.clone() + })); + } + Ok(store) + } + + pub fn set_push_auth(&mut self, auth: PushAuth) { + self.push_auth = Some(auth); } pub fn open(dir: &Path, remote: Option<&str>) -> Result { @@ -157,6 +175,7 @@ impl GitStore { remote_tip: RefCell::new(None), batch: RefCell::new(None), batch_dirty: Cell::new(false), + push_auth: None, #[cfg(test)] command_count: Cell::new(0), }; @@ -304,6 +323,19 @@ impl GitStore { if updates.len() > 1 { arguments.push("--atomic".to_string()); } + if let Some(auth) = &self.push_auth { + let commands: Vec = updates + .iter() + .map(|u| { + super::writers::Command::new( + &u.refname, + u.expected.as_ref().map(Oid::as_str), + u.new.as_ref().map(Oid::as_str), + ) + }) + .collect(); + arguments.push(format!("--push-option={}", auth(&commands))); + } for update in updates { let expected = update.expected.as_ref().map(Oid::as_str).unwrap_or(""); arguments.push(format!("--force-with-lease={}:{expected}", update.refname)); diff --git a/rust/crates/conversation-protocol/src/v3/mod.rs b/rust/crates/conversation-protocol/src/v3/mod.rs index 0829e49b..e95f5cd6 100644 --- a/rust/crates/conversation-protocol/src/v3/mod.rs +++ b/rust/crates/conversation-protocol/src/v3/mod.rs @@ -19,9 +19,10 @@ pub use tasks::{TaskRecord, TaskStatus}; pub mod source_trees; pub mod validate; pub mod view; +pub mod writers; #[cfg(feature = "git-cli")] -pub use git_store::{GitStore, RefUpdate}; +pub use git_store::{GitStore, PushAuth, RefUpdate}; pub use kinds::Kind; pub use oid::Oid; pub use reconcile::{reconcile, CodeOps}; diff --git a/rust/crates/conversation-protocol/src/v3/writers.rs b/rust/crates/conversation-protocol/src/v3/writers.rs new file mode 100644 index 00000000..c4230aae --- /dev/null +++ b/rust/crates/conversation-protocol/src/v3/writers.rs @@ -0,0 +1,556 @@ +//! Ref writers (design/ref-writers.md): the formats a client, a worker and the +//! server's pre-receive hook must agree on. Signing and verifying live behind +//! the `ed25519` feature so std tools that only push with a run token do not +//! compile a signature crate. + +use super::oid::{object_id, ObjectKind, Oid}; +use super::tree::{ + encode_commit_bytes, encode_tree_bytes, CommitInfo, Mode, ObjectStore, Signature, TreeEntry, +}; + +pub const NAMESPACE_PREFIX: &str = "refs/caos/w/"; +pub const WRITERS_REF: &str = "writers"; +pub const WRITERS_PATH: &str = ".caos/writers"; +/// The push option carrying a write's proof: `caos-auth=sig:…` or `caos-auth=run:…`. +pub const PUSH_OPTION: &str = "caos-auth"; +/// Where the server injects a job's run token (`/secret/`). +pub const TOKEN_SECRET: &str = "caos-write"; +pub const TOKEN_PATH: &str = "/secret/caos-write"; +/// The request header that gives a top-level request its writer's authority. +pub const ADMIT_HEADER: &str = "X-Caos-Write"; +/// The ArgTree arg naming the namespaces a job asks to write. +pub const WRITES_ARG: &str = "writes"; +/// How long a client signature stays valid. +pub const SIGNATURE_TTL_SECS: u64 = 600; +pub const ZERO_OID: &str = "0000000000000000000000000000000000000000"; + +pub fn writers_ref(namespace: &str) -> String { + format!("{NAMESPACE_PREFIX}{namespace}/{WRITERS_REF}") +} + +/// `refs/caos/w//` → `(ns, rest)`, for a well-formed namespace id. +pub fn split_ref(refname: &str) -> Option<(&str, &str)> { + let (namespace, rest) = refname.strip_prefix(NAMESPACE_PREFIX)?.split_once('/')?; + (is_namespace(namespace) && !rest.is_empty()).then_some((namespace, rest)) +} + +pub fn is_namespace(value: &str) -> bool { + value.len() == 40 + && value + .bytes() + .all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f')) +} + +pub fn validate_namespace(value: &str) -> Result<(), String> { + if is_namespace(value) { + Ok(()) + } else { + Err(format!( + "{value:?} is not a namespace id (40 lowercase hex)" + )) + } +} + +pub fn is_key(value: &str) -> bool { + value.len() == 64 + && value + .bytes() + .all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f')) +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Writer { + pub key: String, + pub label: String, +} + +/// `.caos/writers`: one ` [label]` per line; `#` starts a comment. +pub fn parse_writers(text: &str) -> Result, String> { + let mut writers: Vec = Vec::new(); + for line in text.lines() { + let line = line.trim(); + if line.is_empty() || line.starts_with('#') { + continue; + } + let (key, label) = line.split_once(char::is_whitespace).unwrap_or((line, "")); + if !is_key(key) { + return Err(format!( + "{WRITERS_PATH}: {key:?} is not an ed25519 public key" + )); + } + if writers.iter().any(|w| w.key == key) { + return Err(format!("{WRITERS_PATH}: {key} is listed twice")); + } + writers.push(Writer { + key: key.to_string(), + label: label.trim().to_string(), + }); + } + Ok(writers) +} + +/// A list holding one writer, `key`. +pub fn sole_writer(key: &str) -> Vec { + vec![Writer { + key: key.to_string(), + label: String::new(), + }] +} + +pub fn format_writers(writers: &[Writer]) -> String { + let mut text = String::from("# \n"); + for writer in writers { + if writer.label.is_empty() { + text.push_str(&format!("{}\n", writer.key)); + } else { + text.push_str(&format!("{} {}\n", writer.key, writer.label)); + } + } + text +} + +/// One ref update as the hook sees it: absent is the zero oid. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Command { + pub old: String, + pub new: String, + pub refname: String, +} + +impl Command { + pub fn new(refname: &str, old: Option<&str>, new: Option<&str>) -> Command { + Command { + old: old.unwrap_or(ZERO_OID).to_string(), + new: new.unwrap_or(ZERO_OID).to_string(), + refname: refname.to_string(), + } + } +} + +/// What a client signs: the expiry and every command, sorted by ref. +pub fn update_message(expiry: u64, commands: &[Command]) -> Vec { + let mut sorted: Vec<&Command> = commands.iter().collect(); + sorted.sort_by(|a, b| a.refname.cmp(&b.refname)); + let mut text = format!("caos-ref-update v1\n{expiry}\n"); + for c in sorted { + text.push_str(&format!("{} {} {}\n", c.old, c.new, c.refname)); + } + text.into_bytes() +} + +/// What a client signs to give a top-level request its authority. +pub fn admit_message(expiry: u64, method: &str, target: &str) -> Vec { + format!("caos-write-admit v1\n{expiry}\n{method} {target}\n").into_bytes() +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Auth { + Signed { + key: String, + expiry: u64, + signature: String, + }, + Run { + token: String, + }, +} + +impl Auth { + /// The push option, `caos-auth=…`. + pub fn option(&self) -> String { + match self { + Auth::Signed { + key, + expiry, + signature, + } => format!("{PUSH_OPTION}=sig:{key}:{expiry}:{signature}"), + Auth::Run { token } => format!("{PUSH_OPTION}=run:{token}"), + } + } + + pub fn parse_option(option: &str) -> Option> { + let value = option.strip_prefix(PUSH_OPTION)?.strip_prefix('=')?; + Some(Auth::parse(value)) + } + + fn parse(value: &str) -> Result { + if let Some(token) = value.strip_prefix("run:") { + if token.is_empty() || !token.bytes().all(|b| b.is_ascii_hexdigit()) { + return Err("malformed run token".to_string()); + } + return Ok(Auth::Run { + token: token.to_string(), + }); + } + let signed = value + .strip_prefix("sig:") + .ok_or_else(|| format!("unknown {PUSH_OPTION} form"))?; + let mut parts = signed.splitn(3, ':'); + let (Some(key), Some(expiry), Some(signature)) = (parts.next(), parts.next(), parts.next()) + else { + return Err(format!("{PUSH_OPTION}=sig: needs key:expiry:signature")); + }; + if !is_key(key) { + return Err(format!("{key:?} is not an ed25519 public key")); + } + Ok(Auth::Signed { + key: key.to_string(), + expiry: expiry + .parse() + .map_err(|_| format!("bad signature expiry {expiry:?}"))?, + signature: signature.to_string(), + }) + } +} + +/// `X-Caos-Write: `. +pub fn admit_header(key: &str, expiry: u64, signature: &str) -> String { + format!("{key} {expiry} {signature}") +} + +pub fn parse_admit_header(value: &str) -> Result<(String, u64, String), String> { + let mut parts = value.split_whitespace(); + let (Some(key), Some(expiry), Some(signature), None) = + (parts.next(), parts.next(), parts.next(), parts.next()) + else { + return Err(format!("{ADMIT_HEADER} needs ")); + }; + if !is_key(key) { + return Err(format!( + "{ADMIT_HEADER}: {key:?} is not an ed25519 public key" + )); + } + let expiry = expiry + .parse() + .map_err(|_| format!("{ADMIT_HEADER}: bad expiry {expiry:?}"))?; + Ok((key.to_string(), expiry, signature.to_string())) +} + +/// What a job's `writes` arg asks for. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum Writes { + /// `*`: everything its creator was handed. For a job that only passes + /// writes on, such as a test harness above the step that pushes. + All, + Namespaces(Vec), +} + +/// The `writes` arg: namespace ids separated by whitespace, or `*`. +pub fn parse_writes(value: &str) -> Result { + if value.trim() == "*" { + return Ok(Writes::All); + } + let mut namespaces = Vec::new(); + for namespace in value.split_whitespace() { + validate_namespace(namespace)?; + if !namespaces.iter().any(|n| n == namespace) { + namespaces.push(namespace.to_string()); + } + } + Ok(Writes::Namespaces(namespaces)) +} + +pub fn now() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0) +} + +const GENESIS_SIGNATURE: &str = "caos"; + +/// A namespace's first `writers` commit. Its oid is the namespace id, so the +/// id is fixed by the initial list: nobody can claim a namespace before its +/// creator, or create one its creator is not in. The commit is deterministic +/// (`label` included), so the same writers and label always name the same +/// namespace. +pub fn genesis_commit( + store: &mut dyn ObjectStore, + writers: &[Writer], + label: &str, +) -> Result { + let (blob, info, id) = genesis_objects(writers, label); + write_writers_tree(store, &blob)?; + let oid = store + .write_commit(&info) + .map_err(|e| format!("writing the namespace's first commit: {e:?}"))?; + debug_assert_eq!(oid, id); + Ok(oid) +} + +fn write_writers_tree(store: &mut dyn ObjectStore, blob: &[u8]) -> Result { + let blob = store + .write_blob(blob) + .map_err(|e| format!("writing {WRITERS_PATH}: {e:?}"))?; + let caos = store + .write_tree(&[TreeEntry { + name: "writers".to_string(), + mode: Mode::Blob, + oid: blob, + }]) + .map_err(|e| format!("writing .caos: {e:?}"))?; + store + .write_tree(&[TreeEntry { + name: ".caos".to_string(), + mode: Mode::Tree, + oid: caos, + }]) + .map_err(|e| format!("writing the writers tree: {e:?}")) +} + +/// The namespace id `genesis_commit` would write, without writing anything. +pub fn genesis_id(writers: &[Writer], label: &str) -> Oid { + genesis_objects(writers, label).2 +} + +fn genesis_objects(writers: &[Writer], label: &str) -> (Vec, CommitInfo, Oid) { + let blob = format_writers(writers).into_bytes(); + let blob_oid = object_id(ObjectKind::Blob, &blob); + let caos = object_id( + ObjectKind::Tree, + &encode_tree_bytes(&[TreeEntry { + name: "writers".to_string(), + mode: Mode::Blob, + oid: blob_oid, + }]), + ); + let root = object_id( + ObjectKind::Tree, + &encode_tree_bytes(&[TreeEntry { + name: ".caos".to_string(), + mode: Mode::Tree, + oid: caos, + }]), + ); + let signature = Signature { + name: GENESIS_SIGNATURE.to_string(), + email: GENESIS_SIGNATURE.to_string(), + time: 0, + offset: "+0000".to_string(), + }; + let commit = CommitInfo { + tree: root, + parents: Vec::new(), + author: signature.clone(), + committer: signature, + extra_headers: Vec::new(), + message: format!("caos namespace\n\n{label}\n").into_bytes(), + }; + let oid = object_id(ObjectKind::Commit, &encode_commit_bytes(&commit)); + (blob, commit, oid) +} + +/// A later `writers` commit: the new list on top of `parent`. +pub fn writers_commit( + store: &mut dyn ObjectStore, + parent: &Oid, + writers: &[Writer], + author: &Signature, + message: &str, +) -> Result { + let root = write_writers_tree(store, format_writers(writers).as_bytes())?; + store + .write_commit(&CommitInfo { + tree: root, + parents: vec![parent.clone()], + author: author.clone(), + committer: author.clone(), + extra_headers: Vec::new(), + message: format!("{}\n", message.trim_end()).into_bytes(), + }) + .map_err(|e| format!("writing the writers commit: {e:?}")) +} + +/// The writers listed at `.caos/writers` in a `writers` commit. +pub fn read_writers(store: &dyn ObjectStore, commit: &Oid) -> Result, String> { + let info = store + .read_commit(commit) + .map_err(|e| format!("reading writers commit {commit}: {e:?}"))?; + let caos = store + .read_tree(&info.tree) + .map_err(|e| format!("reading writers tree: {e:?}"))? + .into_iter() + .find(|e| e.name == ".caos" && e.mode == Mode::Tree) + .ok_or_else(|| format!("writers commit {commit} has no .caos/"))?; + let blob = store + .read_tree(&caos.oid) + .map_err(|e| format!("reading .caos: {e:?}"))? + .into_iter() + .find(|e| e.name == "writers" && e.mode == Mode::Blob) + .ok_or_else(|| format!("writers commit {commit} has no {WRITERS_PATH}"))?; + let bytes = store + .read_blob(&blob.oid) + .map_err(|e| format!("reading {WRITERS_PATH}: {e:?}"))?; + parse_writers(&String::from_utf8(bytes).map_err(|_| format!("{WRITERS_PATH} is not UTF-8"))?) +} + +/// A run token, if this process is a job the server granted one. +pub fn injected_token() -> Option { + std::fs::read_to_string(TOKEN_PATH) + .ok() + .map(|t| t.trim().to_string()) + .filter(|t| !t.is_empty()) +} + +#[cfg(feature = "ed25519")] +pub use keys::*; + +#[cfg(feature = "ed25519")] +mod keys { + use super::super::oid::hex_lower; + use super::*; + use ed25519_dalek::{Signer, SigningKey, VerifyingKey}; + + /// A writer key: the 32-byte seed, as 64 hex. + pub struct WriterKey(SigningKey); + + impl WriterKey { + pub fn parse(seed_hex: &str) -> Result { + let seed = unhex::<32>(seed_hex.trim()) + .ok_or_else(|| "a ref writer key is 64 hex characters".to_string())?; + Ok(WriterKey(SigningKey::from_bytes(&seed))) + } + + pub fn from_seed(seed: [u8; 32]) -> WriterKey { + WriterKey(SigningKey::from_bytes(&seed)) + } + + pub fn seed_hex(&self) -> String { + hex_lower(&self.0.to_bytes()) + } + + pub fn public(&self) -> String { + hex_lower(&self.0.verifying_key().to_bytes()) + } + + pub fn sign(&self, message: &[u8]) -> String { + hex_lower(&self.0.sign(message).to_bytes()) + } + + /// The push option proving `commands` were made by this key. + pub fn sign_update(&self, commands: &[Command]) -> String { + let expiry = now() + SIGNATURE_TTL_SECS; + Auth::Signed { + key: self.public(), + expiry, + signature: self.sign(&update_message(expiry, commands)), + } + .option() + } + + /// The `X-Caos-Write` value for one request. + pub fn sign_admission(&self, method: &str, target: &str) -> String { + let expiry = now() + SIGNATURE_TTL_SECS; + admit_header( + &self.public(), + expiry, + &self.sign(&admit_message(expiry, method, target)), + ) + } + } + + pub fn verify(key: &str, message: &[u8], signature: &str) -> Result<(), String> { + let key = unhex::<32>(key).ok_or_else(|| "malformed public key".to_string())?; + let key = + VerifyingKey::from_bytes(&key).map_err(|_| "not an ed25519 public key".to_string())?; + let signature = unhex::<64>(signature).ok_or_else(|| "malformed signature".to_string())?; + key.verify_strict(message, &ed25519_dalek::Signature::from_bytes(&signature)) + .map_err(|_| "signature does not verify".to_string()) + } + + fn unhex(s: &str) -> Option<[u8; N]> { + if s.len() != N * 2 { + return None; + } + let mut out = [0u8; N]; + for (i, chunk) in s.as_bytes().chunks(2).enumerate() { + out[i] = u8::from_str_radix(std::str::from_utf8(chunk).ok()?, 16).ok()?; + } + Some(out) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const A: &str = "3b6a27bcceb6a42d62a3a8d02a6f0d73653215771de243a63ac048a18b59da29"; + + #[test] + fn writers_round_trip_and_reject_junk() { + let writers = vec![ + Writer { + key: A.into(), + label: "malcolm".into(), + }, + Writer { + key: "9".repeat(64), + label: String::new(), + }, + ]; + assert_eq!(parse_writers(&format_writers(&writers)).unwrap(), writers); + assert!(parse_writers("nothex\n").is_err()); + assert!(parse_writers(&format!("{A}\n{A} again\n")).is_err()); + } + + #[test] + fn refs_split_only_under_a_namespace() { + let ns = "a".repeat(40); + assert_eq!(split_ref(&writers_ref(&ns)), Some((ns.as_str(), "writers"))); + assert_eq!(split_ref("refs/caos/w/short/head"), None); + assert_eq!(split_ref(&format!("refs/caos/w/{ns}/")), None); + assert_eq!(split_ref("refs/heads/main"), None); + } + + #[test] + fn update_messages_ignore_command_order() { + let a = Command::new("refs/b", None, Some(&"1".repeat(40))); + let b = Command::new("refs/a", Some(&"2".repeat(40)), None); + assert_eq!( + update_message(5, &[a.clone(), b.clone()]), + update_message(5, &[b, a]) + ); + } + + #[test] + fn push_options_round_trip() { + let signed = Auth::Signed { + key: A.into(), + expiry: 9, + signature: "ab".into(), + }; + assert_eq!(Auth::parse_option(&signed.option()), Some(Ok(signed))); + let run = Auth::Run { + token: "00ff".into(), + }; + assert_eq!(Auth::parse_option(&run.option()), Some(Ok(run))); + assert_eq!(Auth::parse_option("other=1"), None); + assert!(matches!( + Auth::parse_option("caos-auth=run:x y"), + Some(Err(_)) + )); + } + + #[test] + fn genesis_is_deterministic_and_matches_what_it_writes() { + let writers = vec![Writer { + key: A.into(), + label: String::new(), + }]; + let mut store = crate::v3::MemoryStore::default(); + let written = genesis_commit(&mut store, &writers, "x").unwrap(); + assert_eq!(written, genesis_id(&writers, "x")); + assert_ne!(written, genesis_id(&writers, "y")); + assert_eq!(read_writers(&store, &written).unwrap(), writers); + } + + #[cfg(feature = "ed25519")] + #[test] + fn signatures_verify_only_their_message() { + let key = WriterKey::from_seed([7; 32]); + let message = update_message(1, &[Command::new("refs/x", None, Some(&"1".repeat(40)))]); + let signature = key.sign(&message); + assert!(verify(&key.public(), &message, &signature).is_ok()); + assert!(verify(&key.public(), b"other", &signature).is_err()); + } +} diff --git a/rust/crates/server/Cargo.toml b/rust/crates/server/Cargo.toml index 7722a034..249d8897 100644 --- a/rust/crates/server/Cargo.toml +++ b/rust/crates/server/Cargo.toml @@ -42,3 +42,4 @@ path = "../git-locator" [dependencies.conversation-protocol] path = "../conversation-protocol" +features = ["ed25519"] diff --git a/rust/crates/server/src/compute.rs b/rust/crates/server/src/compute.rs index 3c5db8fb..9ef286da 100644 --- a/rust/crates/server/src/compute.rs +++ b/rust/crates/server/src/compute.rs @@ -652,6 +652,9 @@ fn run_dispatch_inner( // Promise sub-runs see this computation as an ancestor. let mut child_stack: Vec = stack.to_vec(); child_stack.push(arg_tree.to_string()); + // What this job's continuations and children run with: its context, with + // only the writes it was granted (set at dispatch below). + let child_secrets; // Run the worker through the runner rendezvous: resolve the image to a // docker-pullable ref (always sent — a warm runner that pinned the image @@ -681,22 +684,37 @@ fn run_dispatch_inner( // out of band in the job payload — never in the ArgTree, so never in the // cache key — and the container runner drops them at `/secret/`. // Read from `arg_entries` before dispatch takes ownership of it. - let granted = crate::secrets::grant(config, secrets, &arg_entries); - crate::runner::dispatch( + let mut granted = crate::secrets::grant(config, secrets, &arg_entries); + // A job that asks to write (its `writes` arg) gets a run token for + // what it asked and was handed; what it creates is handed only that + // (design/ref-writers.md). + let (token, child_writes) = + crate::ref_writers::grant(config, secrets.writes(), &arg_entries).map_err(fail)?; + if let Some(token) = &token { + granted.push(( + conversation_protocol::v3::writers::TOKEN_SECRET.to_string(), + token.clone(), + )); + } + child_secrets = secrets.clone().with_writes(child_writes); + let result = crate::runner::dispatch( arg_tree, arg_entries, &image_ref, seeded, granted, - |sub_request| start_sub_run(config, sub_request, &child_stack, secrets), + |sub_request| start_sub_run(config, sub_request, &child_stack, &child_secrets), |note| match note { crate::runner::Note::Started => crate::status::started(config, arg_tree), crate::runner::Note::OutTrace(oid) => { crate::status::out_trace(config, arg_tree, &oid) } }, - ) - .map_err(fail)? + ); + if let Some(token) = &token { + crate::ref_writers::revoke(token); + } + result.map_err(fail)? }; if result_hash(&result).is_empty() { @@ -709,7 +727,8 @@ fn run_dispatch_inner( let (result, caught) = match result.split_once(' ') { Some((PROMISE_KIND, cont)) => { eprintln!("resolving promise: arg_tree={arg_tree} -> continuation {cont}"); - resolve_promise(config, arg_tree, cont, salt, &child_stack, secrets).map_err(fail)? + resolve_promise(config, arg_tree, cont, salt, &child_stack, &child_secrets) + .map_err(fail)? } _ => (result, false), }; diff --git a/rust/crates/server/src/main.rs b/rust/crates/server/src/main.rs index dc7ba60f..ed40b0fe 100644 --- a/rust/crates/server/src/main.rs +++ b/rust/crates/server/src/main.rs @@ -54,6 +54,7 @@ mod grant_history; mod import; mod locator; mod push; +mod ref_writers; mod remote_git; mod repair; mod runner; @@ -152,6 +153,10 @@ fn install_termination_handlers() { } fn main() { + // The repository's pre-receive hook is this binary (design/ref-writers.md). + if std::env::args().nth(1).as_deref() == Some(ref_writers::HOOK_ARG) { + std::process::exit(ref_writers::pre_receive()); + } install_termination_handlers(); let addr = std::env::var("SERVER_ADDR").unwrap_or_else(|_| DEFAULT_ADDR.to_string()); @@ -201,7 +206,24 @@ fn main() { eprintln!("fatal: {error}"); std::process::exit(1); }); - remove_managed_pre_receive_hook(&git_dir).unwrap_or_else(|error| { + // Who may write which ref is checked on every push, by a pre-receive hook + // that is this binary (design/ref-writers.md). Its proof travels as a push + // option, which receive-pack only accepts once advertised. + git(&[ + "-C", + &git_dir, + "config", + "receive.advertisePushOptions", + "true", + ]); + ensure_git_config_value( + &git_dir, + ref_writers::UNGUARDED_CONFIG, + ref_writers::DEFAULT_UNGUARDED, + ) + .and_then(|()| ref_writers::install_hook(&git_dir, &addr)) + .and_then(|()| ensure_hooks_path(&git_dir)) + .unwrap_or_else(|error| { eprintln!("fatal: {error}"); std::process::exit(1); }); @@ -574,49 +596,25 @@ fn remove_git_config_value(git_dir: &str, key: &str, value: &str) -> Result<(), } } -const MANAGED_HOOK_MARKERS: [&str; 2] = [ - "# managed by caos-server: append-only refs", - "# managed by caos-server: append-only conversation heads", -]; - -// TODO: Remove this one-time migration after every supported repository has -// been started by a server version that no longer installs the hook. -/// Remove the pre-receive hook installed by older CAOS servers. -/// -/// The old hook execs this binary with a validator mode that no longer exists, -/// so merely ceasing to install it would make every push to an upgraded -/// repository fail. Its marker is the ownership proof: an unmarked hook and its -/// configured path belong to the administrator and are left untouched. -fn remove_managed_pre_receive_hook(git_dir: &str) -> Result<(), String> { +/// The hook lives at `/hooks`, git's default. Older servers set +/// `core.hooksPath` to exactly that, which is equivalent; anything else is an +/// administrator's choice that would bypass the ref-writers hook. +fn ensure_hooks_path(git_dir: &str) -> Result<(), String> { let hooks = std::path::Path::new(git_dir).join("hooks"); - let hook = hooks.join("pre-receive"); - let hooks_value = hooks - .to_str() - .ok_or_else(|| format!("hooks path is not UTF-8: {}", hooks.display()))?; - let contents = match std::fs::read(&hook) { - Ok(contents) => contents, - // A previous cleanup may have removed the hook without clearing the - // absolute path our installer wrote. That exact value is redundant - // with Git's default `/hooks`, so clearing it cannot disable a - // surviving hook in this directory. - Err(error) if error.kind() == std::io::ErrorKind::NotFound => { - return remove_git_config_value(git_dir, "core.hooksPath", hooks_value); - } - Err(error) => return Err(format!("reading {}: {error}", hook.display())), - }; - let managed = MANAGED_HOOK_MARKERS.iter().any(|marker| { - contents - .windows(marker.len()) - .any(|window| window == marker.as_bytes()) - }); - if !managed { - return Ok(()); + let output = std::process::Command::new("git") + .args(["-C", git_dir, "config", "--get", "core.hooksPath"]) + .output() + .map_err(|e| format!("reading core.hooksPath: {e}"))?; + let configured = String::from_utf8_lossy(&output.stdout).trim().to_string(); + if configured.is_empty() || configured == "hooks" || std::path::Path::new(&configured) == hooks + { + Ok(()) + } else { + Err(format!( + "core.hooksPath is {configured:?}, so the ref-writers hook in {} would never run", + hooks.display() + )) } - - std::fs::remove_file(&hook).map_err(|error| format!("removing {}: {error}", hook.display()))?; - // The installer always wrote this absolute value. Remove only that exact - // value; preserve relative or alternate administrator-selected hook paths. - remove_git_config_value(git_dir, "core.hooksPath", hooks_value) } fn env_or(key: &str, default: &str) -> String { @@ -700,11 +698,17 @@ fn secret_context(config: &Config, request: &Request) -> Result, request: &mut Request) -> Result, HttpErr request.as_reader().read_to_string(&mut body)?; runner::sub_run(&body) } + Method::Post if path == ref_writers::TOKEN_PATH => { + let mut body = String::new(); + request.as_reader().read_to_string(&mut body)?; + ref_writers::token_endpoint(&body) + } Method::Post if path == "/trace/child" => { let mut body = String::new(); request.as_reader().read_to_string(&mut body)?; @@ -876,10 +885,7 @@ fn spawn_request_ref_pruner(git_dir: String) { #[cfg(test)] mod tests { - use super::{ - configure_ref_advertisements, remove_managed_pre_receive_hook, run_required_git, - MANAGED_HOOK_MARKERS, - }; + use super::{configure_ref_advertisements, ensure_hooks_path, run_required_git}; use std::process::Command; #[test] @@ -997,114 +1003,49 @@ mod tests { std::fs::remove_dir_all(dir).unwrap(); } - #[test] - fn managed_pre_receive_hooks_are_removed_on_upgrade() { - for (index, marker) in MANAGED_HOOK_MARKERS.iter().enumerate() { - let dir = std::env::temp_dir().join(format!( - "caos-managed-hook-test-{}-{index}", - std::process::id() - )); - std::fs::remove_dir_all(&dir).ok(); - run_required_git(&["init", "-q", "--bare", dir.to_str().unwrap()]).unwrap(); - let hooks = dir.join("hooks"); - let hook = hooks.join("pre-receive"); - std::fs::write(&hook, format!("#!/bin/sh\n{marker}\nexit 1\n")).unwrap(); - run_required_git(&[ - "-C", - dir.to_str().unwrap(), - "config", - "core.hooksPath", - hooks.to_str().unwrap(), - ]) - .unwrap(); - - remove_managed_pre_receive_hook(dir.to_str().unwrap()).unwrap(); - - assert!(!hook.exists()); - let configured = Command::new("git") - .args([ - "-C", - dir.to_str().unwrap(), - "config", - "--get", - "core.hooksPath", - ]) - .output() - .unwrap(); - assert_eq!(configured.status.code(), Some(1)); - std::fs::remove_dir_all(dir).unwrap(); - } + fn bare(name: &str) -> std::path::PathBuf { + let dir = std::env::temp_dir().join(format!("caos-{name}-{}", std::process::id())); + std::fs::remove_dir_all(&dir).ok(); + run_required_git(&["init", "-q", "--bare", dir.to_str().unwrap()]).unwrap(); + dir } #[test] - fn unmanaged_pre_receive_hook_and_path_are_preserved() { - let dir = - std::env::temp_dir().join(format!("caos-unmanaged-hook-test-{}", std::process::id())); - std::fs::remove_dir_all(&dir).ok(); - run_required_git(&["init", "-q", "--bare", dir.to_str().unwrap()]).unwrap(); - let hooks = dir.join("hooks"); - let hook = hooks.join("pre-receive"); - let contents = b"#!/bin/sh\necho administrator hook\n"; - std::fs::write(&hook, contents).unwrap(); - run_required_git(&[ - "-C", - dir.to_str().unwrap(), - "config", - "core.hooksPath", - hooks.to_str().unwrap(), - ]) + fn the_ref_writers_hook_replaces_an_older_managed_hook() { + let dir = bare("managed-hook"); + let hook = dir.join("hooks").join("pre-receive"); + std::fs::write( + &hook, + "#!/bin/sh\n# managed by caos-server: append-only refs\nexit 1\n", + ) .unwrap(); - - remove_managed_pre_receive_hook(dir.to_str().unwrap()).unwrap(); - - assert_eq!(std::fs::read(&hook).unwrap(), contents); - let configured = Command::new("git") - .args([ - "-C", - dir.to_str().unwrap(), - "config", - "--get", - "core.hooksPath", - ]) - .output() - .unwrap(); - assert!(configured.status.success()); - assert_eq!( - String::from_utf8(configured.stdout).unwrap().trim(), - hooks.to_str().unwrap() - ); + crate::ref_writers::install_hook(dir.to_str().unwrap(), "[::]:80").unwrap(); + let installed = std::fs::read_to_string(&hook).unwrap(); + assert!(installed.contains("--pre-receive"), "{installed}"); + assert!(installed.contains("http://127.0.0.1:80"), "{installed}"); + ensure_hooks_path(dir.to_str().unwrap()).unwrap(); std::fs::remove_dir_all(dir).unwrap(); } #[test] - fn stale_managed_hook_path_is_removed_when_the_hook_is_already_absent() { - let dir = - std::env::temp_dir().join(format!("caos-stale-hook-path-test-{}", std::process::id())); - std::fs::remove_dir_all(&dir).ok(); - run_required_git(&["init", "-q", "--bare", dir.to_str().unwrap()]).unwrap(); - let hooks = dir.join("hooks"); + fn an_administrators_hook_or_hook_path_is_refused_not_overwritten() { + let dir = bare("unmanaged-hook"); + let hook = dir.join("hooks").join("pre-receive"); + std::fs::write(&hook, "#!/bin/sh\necho administrator hook\n").unwrap(); + assert!(crate::ref_writers::install_hook(dir.to_str().unwrap(), "[::]:80").is_err()); + assert_eq!( + std::fs::read_to_string(&hook).unwrap(), + "#!/bin/sh\necho administrator hook\n" + ); run_required_git(&[ "-C", dir.to_str().unwrap(), "config", "core.hooksPath", - hooks.to_str().unwrap(), + "/elsewhere", ]) .unwrap(); - - remove_managed_pre_receive_hook(dir.to_str().unwrap()).unwrap(); - - let configured = Command::new("git") - .args([ - "-C", - dir.to_str().unwrap(), - "config", - "--get", - "core.hooksPath", - ]) - .output() - .unwrap(); - assert_eq!(configured.status.code(), Some(1)); + assert!(ensure_hooks_path(dir.to_str().unwrap()).is_err()); std::fs::remove_dir_all(dir).unwrap(); } } diff --git a/rust/crates/server/src/ref_writers.rs b/rust/crates/server/src/ref_writers.rs new file mode 100644 index 00000000..b36db59b --- /dev/null +++ b/rust/crates/server/src/ref_writers.rs @@ -0,0 +1,799 @@ +//! Who may write which ref (design/ref-writers.md). +//! +//! Three parts, one per moment a write is decided: +//! * admission: a top-level request's `X-Caos-Write` header names the writer +//! it acts for ([`admit`]); +//! * dispatch: a job asking for namespaces in its `writes` arg gets a run token +//! covering what it asked for and its creator held, and its children inherit +//! only that ([`grant`]); +//! * the push: the repository's pre-receive hook re-execs this binary +//! ([`pre_receive`]) and checks every command against the namespace's +//! `writers` list. + +use std::collections::HashMap; +use std::io::Read; +use std::process::Command; +use std::sync::{Arc, Mutex}; + +use conversation_protocol::v3::writers::{ + self, parse_writers, split_ref, update_message, Auth, Command as RefCommand, Writer, + WRITERS_PATH, WRITERS_REF, ZERO_OID, +}; + +use crate::{Config, HttpError}; + +/// The argv that makes this binary the pre-receive hook. +pub(crate) const HOOK_ARG: &str = "--pre-receive"; +/// Where the hook asks for a run token's grant. +pub(crate) const TOKEN_PATH: &str = "/ref-writers/token"; +const SERVER_ENV: &str = "CAOS_REF_WRITERS_SERVER"; +const HOOK_MARKER: &str = "# managed by caos-server: ref writers"; +/// Hooks older servers installed; ours replaces them. +const OLD_HOOK_MARKERS: [&str; 2] = [ + "# managed by caos-server: append-only refs", + "# managed by caos-server: append-only conversation heads", +]; +/// Refs named by their own content: each must point at its last component. +const CONTENT_NAMED: [&str; 2] = ["refs/caos/req/", "refs/heads/caos-test/"]; +/// `report` logs what enforcement would refuse and accepts it, for rollout. +const MODE_CONFIG: &str = "caos.refWriters"; +/// Refs any push may write. Each is a hole: `caosd up`'s dev publish is the +/// one the server keeps open by default. +pub(crate) const UNGUARDED_CONFIG: &str = "caos.unguardedRef"; +pub(crate) const DEFAULT_UNGUARDED: &str = "refs/caos/dev"; + +/// Which writer a run acts for, and the namespaces it may still hand on. +#[derive(Clone, Default)] +pub(crate) enum Writes { + #[default] + None, + As { + key: String, + /// `None`: whatever the writer may write (a top-level request). + scope: Option>>, + }, +} + +/// A top-level request's authority, from its `X-Caos-Write` header. +pub(crate) fn admit(header: Option<&str>, method: &str, target: &str) -> Result { + let Some(header) = header else { + return Ok(Writes::None); + }; + let refuse = |m: String| HttpError::new(403, format!("{}: {m}", writers::ADMIT_HEADER)); + let (key, expiry, signature) = writers::parse_admit_header(header).map_err(refuse)?; + if expiry < writers::now() { + return Err(refuse("expired".to_string())); + } + writers::verify( + &key, + &writers::admit_message(expiry, method, target), + &signature, + ) + .map_err(refuse)?; + Ok(Writes::As { key, scope: None }) +} + +// ---- run tokens --------------------------------------------------------------- + +#[derive(Clone)] +struct Grant { + key: String, + /// `None`: every namespace the writer may write. + namespaces: Option>, +} + +static TOKENS: Mutex>> = Mutex::new(None); + +/// A dispatched job's write grant: what it asked for in `writes`, within what +/// it was handed. Returns the token to inject (revoke it when the job's +/// container is done) and what the job's children may in turn be handed. +pub(crate) fn grant( + config: &Config, + writes: &Writes, + arg_entries: &std::collections::BTreeMap, +) -> Result<(Option, Writes), HttpError> { + let Writes::As { key, scope } = writes else { + return Ok((None, Writes::None)); + }; + let Some(oid) = arg_entries.get(writers::WRITES_ARG) else { + return Ok((None, Writes::None)); + }; + let asked = crate::storage::fetch_blob(config, oid) + .ok() + .and_then(|bytes| String::from_utf8(bytes).ok()) + .ok_or_else(|| HttpError::new(400, "the `writes` arg must be a text blob"))?; + let granted: Option> = + match writers::parse_writes(&asked).map_err(|e| HttpError::new(400, e))? { + writers::Writes::All => scope.as_ref().map(|scope| scope.to_vec()), + writers::Writes::Namespaces(asked) => Some( + asked + .into_iter() + .filter(|ns| scope.as_ref().is_none_or(|scope| scope.contains(ns))) + .collect(), + ), + }; + if granted.as_ref().is_some_and(Vec::is_empty) { + return Ok((None, Writes::None)); + } + let token = random_token()?; + TOKENS + .lock() + .unwrap_or_else(|e| e.into_inner()) + .get_or_insert_with(HashMap::new) + .insert( + token.clone(), + Grant { + key: key.clone(), + namespaces: granted.clone(), + }, + ); + Ok(( + Some(token), + Writes::As { + key: key.clone(), + scope: granted.map(Arc::new), + }, + )) +} + +pub(crate) fn revoke(token: &str) { + if let Some(tokens) = TOKENS.lock().unwrap_or_else(|e| e.into_inner()).as_mut() { + tokens.remove(token); + } +} + +/// `POST /ref-writers/token`: what a run token may write, for the hook. It +/// tells a caller nothing it could use: holding the token is already holding +/// the grant. +pub(crate) fn token_endpoint(body: &str) -> Result, HttpError> { + let grant = TOKENS + .lock() + .unwrap_or_else(|e| e.into_inner()) + .as_ref() + .and_then(|tokens| tokens.get(body.trim()).cloned()) + .ok_or_else(|| HttpError::new(404, "no such run token (its run has ended)"))?; + let all = grant.namespaces.is_none(); + Ok(serde_json::json!({ + "key": grant.key, + "namespaces": grant.namespaces.unwrap_or_default(), + "all": all, + }) + .to_string() + .into_bytes()) +} + +fn random_token() -> Result { + let mut bytes = [0u8; 32]; + std::fs::File::open("/dev/urandom") + .and_then(|mut f| f.read_exact(&mut bytes)) + .map_err(|e| HttpError::new(500, format!("reading /dev/urandom: {e}")))?; + Ok(bytes.iter().map(|b| format!("{b:02x}")).collect()) +} + +// ---- installing the hook -------------------------------------------------------- + +/// Point the repository's pre-receive hook at this binary. Rewritten on every +/// boot so it always names the binary that is running; an administrator's own +/// hook (no marker) is refused rather than overwritten. +pub(crate) fn install_hook(git_dir: &str, server_addr: &str) -> Result<(), String> { + let hooks = std::path::Path::new(git_dir).join("hooks"); + let hook = hooks.join("pre-receive"); + match std::fs::read_to_string(&hook) { + Ok(existing) + if !existing.contains(HOOK_MARKER) + && !OLD_HOOK_MARKERS.iter().any(|m| existing.contains(m)) => + { + return Err(format!( + "{} exists and is not caos's; ref writers cannot be enforced", + hook.display() + )); + } + _ => {} + } + let exe = std::env::current_exe().map_err(|e| format!("locating the server binary: {e}"))?; + let shell = find_on_path("sh") + .or_else(|| find_on_path("bash")) + .ok_or_else(|| "no sh or bash on PATH for the pre-receive hook".to_string())?; + let script = format!( + "#!{}\n{HOOK_MARKER}\n{SERVER_ENV}={} exec {} {HOOK_ARG}\n", + shell.display(), + shell_quote(&hook_server_url(server_addr)), + shell_quote(&exe.display().to_string()), + ); + std::fs::create_dir_all(&hooks).map_err(|e| format!("creating {}: {e}", hooks.display()))?; + let temp = hooks.join(format!("pre-receive.{}", std::process::id())); + std::fs::write(&temp, script).map_err(|e| format!("writing {}: {e}", temp.display()))?; + use std::os::unix::fs::PermissionsExt; + std::fs::set_permissions(&temp, std::fs::Permissions::from_mode(0o755)) + .map_err(|e| format!("chmod {}: {e}", temp.display()))?; + std::fs::rename(&temp, &hook).map_err(|e| format!("installing {}: {e}", hook.display()))?; + Ok(()) +} + +/// The address the hook reaches this server at: the listen address, with an +/// unspecified host made loopback (the hook runs beside the server). +fn hook_server_url(server_addr: &str) -> String { + let (host, port) = server_addr.rsplit_once(':').unwrap_or((server_addr, "80")); + let host = match host { + "" | "0.0.0.0" | "[::]" | "::" => "127.0.0.1", + host => host, + }; + format!("http://{host}:{port}") +} + +fn find_on_path(name: &str) -> Option { + std::env::var_os("PATH")? + .to_str()? + .split(':') + .map(|dir| std::path::Path::new(dir).join(name)) + .find(|path| path.is_file()) +} + +fn shell_quote(value: &str) -> String { + format!("'{}'", value.replace('\'', "'\\''")) +} + +// ---- the hook -------------------------------------------------------------------- + +/// Run as the repository's pre-receive hook: refuse the whole push unless every +/// command is allowed. Returns the process exit code. +pub(crate) fn pre_receive() -> i32 { + let mut input = String::new(); + if let Err(e) = std::io::stdin().read_to_string(&mut input) { + eprintln!("caos: reading pushed refs: {e}"); + return 1; + } + let commands: Vec = input + .lines() + .filter_map(|line| { + let mut fields = line.split_whitespace(); + Some(RefCommand { + old: fields.next()?.to_string(), + new: fields.next()?.to_string(), + refname: fields.next()?.to_string(), + }) + }) + .collect(); + let options: Vec = (0..std::env::var("GIT_PUSH_OPTION_COUNT") + .ok() + .and_then(|n| n.parse::().ok()) + .unwrap_or(0)) + .filter_map(|i| std::env::var(format!("GIT_PUSH_OPTION_{i}")).ok()) + .collect(); + let report = git(&["config", "--get", MODE_CONFIG]).is_ok_and(|mode| mode.trim() == "report"); + match check(&Repo, &commands, &options, &TokenServer) { + Ok(()) => 0, + Err(message) if report => { + eprintln!("caos: ref writers would refuse this push: {message} (report only)"); + 0 + } + Err(message) => { + eprintln!("caos: {message}"); + 1 + } + } +} + +/// What the check reads from the repository. A trait so it can be tested +/// without a git process. +trait Store { + fn ref_value(&self, refname: &str) -> Option; + /// `.caos/writers` of a commit, even one only in the push's quarantine. + fn writers_at(&self, commit: &str) -> Result, String>; + fn parents(&self, commit: &str) -> Result, String>; + fn is_ancestor(&self, ancestor: &str, descendant: &str) -> bool; + fn unguarded(&self) -> Vec; +} + +trait Tokens { + fn lookup(&self, token: &str) -> Result<(String, Scope), String>; +} + +/// Who signed or ran the push. +struct Pusher { + key: String, + /// A run token's grant; `None` for a writer's own signature. + token: Option, +} + +/// What a run token may write: every namespace its writer may, or these. +#[derive(Clone)] +enum Scope { + All, + Only(Vec), +} + +impl Scope { + fn covers(&self, namespace: &str) -> bool { + match self { + Scope::All => true, + Scope::Only(namespaces) => namespaces.iter().any(|n| n == namespace), + } + } +} + +fn check( + store: &dyn Store, + commands: &[RefCommand], + options: &[String], + tokens: &dyn Tokens, +) -> Result<(), String> { + // Identified once, and only if the push touches a governed ref. + let mut identified: Option> = None; + let mut pusher_of = || -> Result { + identified + .get_or_insert_with(|| identify(commands, options, tokens)) + .as_ref() + .map(|p| Pusher { + key: p.key.clone(), + token: p.token.clone(), + }) + .map_err(Clone::clone) + }; + let unguarded = store.unguarded(); + // Writers lists this push establishes, so a namespace created here governs + // the other refs created alongside it. + let mut created: HashMap> = HashMap::new(); + let mut ordered: Vec<&RefCommand> = commands.iter().collect(); + // `writers` refs first: the rest of the push is judged by them. + ordered.sort_by_key(|c| !split_ref(&c.refname).is_some_and(|(_, rest)| rest == WRITERS_REF)); + for command in ordered { + let refname = command.refname.as_str(); + let deny = |why: &str| Err(format!("{refname}: {why}")); + if unguarded.iter().any(|r| r == refname) { + continue; + } + if let Some(name) = CONTENT_NAMED.iter().find_map(|p| refname.strip_prefix(p)) { + // Deleting one only unpins it; the object stays. + if command.new == name || command.new == ZERO_OID { + continue; + } + return deny("a content-named ref must point at the object it names"); + } + // Conversations still live outside namespaces, so a ref outside one + // is let through until they move in (design/ref-writers.md). + let Some((namespace, rest)) = split_ref(refname) else { + continue; + }; + let pusher = pusher_of().map_err(|e| format!("{refname}: {e}"))?; + if rest == WRITERS_REF { + // A job may found a namespace for its writer (a test harness + // creating a conversation), but only a writer changes who writes. + let creating = command.old == ZERO_OID; + if let Some(scope) = &pusher.token { + if !creating { + return deny( + "only a writer key can change who writes; a job's run token cannot", + ); + } + if !scope.covers(namespace) { + return deny("this job's run token was not granted that namespace"); + } + } + if command.new == ZERO_OID { + return deny("a namespace's writers list cannot be deleted"); + } + let list = if command.old == ZERO_OID { + if command.new != namespace { + return deny("a namespace's id is the hash of its first writers commit"); + } + if !store.parents(&command.new)?.is_empty() { + return deny("a namespace's first writers commit has no parents"); + } + store.writers_at(&command.new)? + } else { + let current = store.writers_at(&command.old)?; + if !current.iter().any(|w| w.key == pusher.key) { + return deny(&format!("{} is not one of its writers", pusher.key)); + } + if !store.is_ancestor(&command.old, &command.new) { + return deny("the writers list only moves forward (fast-forward)"); + } + store.writers_at(&command.new)? + }; + if command.old == ZERO_OID && !list.iter().any(|w| w.key == pusher.key) { + return deny("the creator must be one of the first writers"); + } + created.insert(namespace.to_string(), list); + continue; + } + let list = match created.get(namespace) { + Some(list) => list.clone(), + None => match store.ref_value(&writers::writers_ref(namespace)) { + Some(tip) => store.writers_at(&tip)?, + None => return deny("its namespace has no writers ref"), + }, + }; + if !list.iter().any(|w| w.key == pusher.key) { + return deny(&format!( + "{} is not one of the namespace's writers", + pusher.key + )); + } + if let Some(scope) = &pusher.token { + if !scope.covers(namespace) { + return deny("this job's run token was not granted that namespace"); + } + } + } + Ok(()) +} + +fn identify( + commands: &[RefCommand], + options: &[String], + tokens: &dyn Tokens, +) -> Result { + let mut auths = options.iter().filter_map(|o| Auth::parse_option(o)); + let auth = match (auths.next(), auths.next()) { + (None, _) => { + return Err(format!( + "a governed ref needs a `{}` push option (a writer's signature or a run token)", + writers::PUSH_OPTION + )) + } + (Some(_), Some(_)) => { + return Err(format!( + "more than one `{}` push option", + writers::PUSH_OPTION + )) + } + (Some(auth), None) => auth?, + }; + match auth { + Auth::Signed { + key, + expiry, + signature, + } => { + if expiry < writers::now() { + return Err("the push's signature has expired".to_string()); + } + writers::verify(&key, &update_message(expiry, commands), &signature)?; + Ok(Pusher { key, token: None }) + } + Auth::Run { token } => { + let (key, scope) = tokens.lookup(&token)?; + Ok(Pusher { + key, + token: Some(scope), + }) + } + } +} + +struct Repo; + +fn git(args: &[&str]) -> Result { + let output = Command::new("git") + .args(args) + .output() + .map_err(|e| format!("running git: {e}"))?; + if output.status.success() { + String::from_utf8(output.stdout).map_err(|_| "git printed non-UTF-8".to_string()) + } else { + Err(format!( + "git {}: {}", + args.join(" "), + String::from_utf8_lossy(&output.stderr).trim() + )) + } +} + +impl Store for Repo { + fn ref_value(&self, refname: &str) -> Option { + git(&["rev-parse", "--verify", "--quiet", refname]) + .ok() + .map(|v| v.trim().to_string()) + } + + fn writers_at(&self, commit: &str) -> Result, String> { + let text = git(&["cat-file", "blob", &format!("{commit}:{WRITERS_PATH}")]) + .map_err(|_| format!("{commit} has no {WRITERS_PATH}"))?; + parse_writers(&text) + } + + fn parents(&self, commit: &str) -> Result, String> { + let text = git(&["cat-file", "commit", commit])?; + Ok(text + .lines() + .take_while(|l| !l.is_empty()) + .filter_map(|l| l.strip_prefix("parent ")) + .map(str::to_string) + .collect()) + } + + fn is_ancestor(&self, ancestor: &str, descendant: &str) -> bool { + Command::new("git") + .args(["merge-base", "--is-ancestor", ancestor, descendant]) + .status() + .is_ok_and(|s| s.success()) + } + + fn unguarded(&self) -> Vec { + git(&["config", "--get-all", UNGUARDED_CONFIG]) + .map(|v| v.lines().map(str::to_string).collect()) + .unwrap_or_default() + } +} + +struct TokenServer; + +impl Tokens for TokenServer { + fn lookup(&self, token: &str) -> Result<(String, Scope), String> { + let server = std::env::var(SERVER_ENV) + .map_err(|_| format!("{SERVER_ENV} is not set; cannot check a run token"))?; + let response = minreq::post(format!("{server}{TOKEN_PATH}")) + .with_body(token) + .with_timeout(10) + .send() + .map_err(|e| format!("asking the server about a run token: {e}"))?; + if response.status_code == 404 { + return Err("unknown run token (its run has ended?)".to_string()); + } + if response.status_code != 200 { + return Err(format!( + "the server answered {} about a run token", + response.status_code + )); + } + let body: serde_json::Value = serde_json::from_slice(response.as_bytes()) + .map_err(|e| format!("reading the run token's grant: {e}"))?; + let key = body["key"].as_str().unwrap_or_default().to_string(); + if body["all"].as_bool() == Some(true) { + return Ok((key, Scope::All)); + } + let namespaces = body["namespaces"] + .as_array() + .map(|a| { + a.iter() + .filter_map(|v| v.as_str().map(str::to_string)) + .collect() + }) + .unwrap_or_default(); + Ok((key, Scope::Only(namespaces))) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use conversation_protocol::v3::writers::{genesis_id, WriterKey}; + + #[derive(Default)] + struct Fake { + refs: HashMap, + writers: HashMap>, + parents: HashMap>, + ancestry: Vec<(String, String)>, + } + + impl Store for Fake { + fn ref_value(&self, refname: &str) -> Option { + self.refs.get(refname).cloned() + } + fn writers_at(&self, commit: &str) -> Result, String> { + self.writers + .get(commit) + .cloned() + .ok_or_else(|| format!("no {commit}")) + } + fn parents(&self, commit: &str) -> Result, String> { + Ok(self.parents.get(commit).cloned().unwrap_or_default()) + } + fn is_ancestor(&self, a: &str, d: &str) -> bool { + self.ancestry.iter().any(|(x, y)| x == a && y == d) + } + fn unguarded(&self) -> Vec { + vec![DEFAULT_UNGUARDED.to_string()] + } + } + + struct Grants(Vec<(String, String, Option>)>); + + impl Tokens for Grants { + fn lookup(&self, token: &str) -> Result<(String, Scope), String> { + self.0 + .iter() + .find(|(t, _, _)| t == token) + .map(|(_, k, n)| { + let scope = n.clone().map_or(Scope::All, Scope::Only); + (k.clone(), scope) + }) + .ok_or_else(|| "unknown".to_string()) + } + } + + fn writer(key: &WriterKey) -> Writer { + Writer { + key: key.public(), + label: String::new(), + } + } + + fn oid(c: char) -> String { + c.to_string().repeat(40) + } + + /// A store holding namespace `ns` whose writers are `keys`. + fn namespace(keys: &[&WriterKey]) -> (Fake, String) { + let list: Vec = keys.iter().map(|k| writer(k)).collect(); + let ns = genesis_id(&list, "test").to_string(); + let mut fake = Fake::default(); + fake.refs.insert(writers::writers_ref(&ns), ns.clone()); + fake.writers.insert(ns.clone(), list); + (fake, ns) + } + + fn signed(key: &WriterKey, commands: &[RefCommand]) -> Vec { + vec![key.sign_update(commands)] + } + + #[test] + fn content_named_and_unguarded_refs_need_no_proof() { + let none = Grants(vec![]); + let ok = RefCommand::new( + &format!("refs/caos/req/{}", oid('a')), + None, + Some(&oid('a')), + ); + assert!(check(&Fake::default(), &[ok], &[], &none).is_ok()); + let lie = RefCommand::new( + &format!("refs/caos/req/{}", oid('a')), + None, + Some(&oid('b')), + ); + assert!(check(&Fake::default(), &[lie], &[], &none).is_err()); + let dev = RefCommand::new("refs/caos/dev", None, Some(&oid('b'))); + assert!(check(&Fake::default(), &[dev], &[], &none).is_ok()); + let other = RefCommand::new("refs/heads/main", None, Some(&oid('b'))); + assert!(check(&Fake::default(), &[other], &[], &none).is_ok()); + } + + #[test] + fn a_writer_signature_writes_its_namespace_and_nobody_elses() { + let alice = WriterKey::from_seed([1; 32]); + let mallory = WriterKey::from_seed([2; 32]); + let (store, ns) = namespace(&[&alice]); + let head = [RefCommand::new( + &format!("refs/caos/w/{ns}/head"), + None, + Some(&oid('c')), + )]; + let none = Grants(vec![]); + assert!(check(&store, &head, &signed(&alice, &head), &none).is_ok()); + assert!(check(&store, &head, &signed(&mallory, &head), &none).is_err()); + assert!(check(&store, &head, &[], &none).is_err()); + // A signature covers exactly the commands it was made for. + let other = [RefCommand::new( + &format!("refs/caos/w/{ns}/head"), + None, + Some(&oid('d')), + )]; + assert!(check(&store, &other, &signed(&alice, &head), &none).is_err()); + } + + #[test] + fn a_namespace_is_created_only_at_its_own_hash_by_a_listed_writer() { + let alice = WriterKey::from_seed([1; 32]); + let bob = WriterKey::from_seed([3; 32]); + let list = vec![writer(&alice)]; + let ns = genesis_id(&list, "new").to_string(); + let mut store = Fake::default(); + store.writers.insert(ns.clone(), list); + let create = [ + RefCommand::new(&writers::writers_ref(&ns), None, Some(&ns)), + RefCommand::new(&format!("refs/caos/w/{ns}/head"), None, Some(&oid('c'))), + ]; + let none = Grants(vec![]); + assert!(check(&store, &create, &signed(&alice, &create), &none).is_ok()); + assert!(check(&store, &create, &signed(&bob, &create), &none).is_err()); + let squat = [RefCommand::new( + &writers::writers_ref(&oid('e')), + None, + Some(&ns), + )]; + assert!(check(&store, &squat, &signed(&alice, &squat), &none).is_err()); + } + + #[test] + fn writers_change_forward_only_and_removal_takes_effect() { + let alice = WriterKey::from_seed([1; 32]); + let bob = WriterKey::from_seed([3; 32]); + let (mut store, ns) = namespace(&[&alice, &bob]); + let next = oid('f'); + store.writers.insert(next.clone(), vec![writer(&alice)]); + store.ancestry.push((ns.clone(), next.clone())); + let remove_bob = [RefCommand::new( + &writers::writers_ref(&ns), + Some(&ns), + Some(&next), + )]; + let none = Grants(vec![]); + assert!(check(&store, &remove_bob, &signed(&bob, &remove_bob), &none).is_ok()); + store.refs.insert(writers::writers_ref(&ns), next.clone()); + let head = [RefCommand::new( + &format!("refs/caos/w/{ns}/head"), + None, + Some(&oid('c')), + )]; + assert!(check(&store, &head, &signed(&bob, &head), &none).is_err()); + assert!(check(&store, &head, &signed(&alice, &head), &none).is_ok()); + let rewind = [RefCommand::new( + &writers::writers_ref(&ns), + Some(&next), + Some(&ns), + )]; + assert!(check(&store, &rewind, &signed(&alice, &rewind), &none).is_err()); + } + + #[test] + fn a_run_token_writes_only_granted_namespaces_and_never_the_list() { + let alice = WriterKey::from_seed([1; 32]); + let (store, ns) = namespace(&[&alice]); + let grants = Grants(vec![ + ("aa".into(), alice.public(), Some(vec![ns.clone()])), + ("bb".into(), alice.public(), Some(vec![oid('9')])), + ("ff".into(), alice.public(), None), + ]); + let head = [RefCommand::new( + &format!("refs/caos/w/{ns}/head"), + None, + Some(&oid('c')), + )]; + let run = |t: &str| vec![Auth::Run { token: t.into() }.option()]; + assert!(check(&store, &head, &run("aa"), &grants).is_ok()); + assert!(check(&store, &head, &run("bb"), &grants).is_err()); + assert!(check(&store, &head, &run("cc"), &grants).is_err()); + assert!(check(&store, &head, &run("ff"), &grants).is_ok()); + let list = [RefCommand::new( + &writers::writers_ref(&ns), + Some(&ns), + Some(&oid('f')), + )]; + assert!(check(&store, &list, &run("aa"), &grants).is_err()); + assert!(check(&store, &list, &run("ff"), &grants).is_err()); + } + + #[test] + fn a_run_token_may_found_a_namespace_for_its_writer() { + let alice = WriterKey::from_seed([1; 32]); + let bob = WriterKey::from_seed([3; 32]); + let list = vec![writer(&alice)]; + let ns = genesis_id(&list, "job").to_string(); + let mut store = Fake::default(); + store.writers.insert(ns.clone(), list); + let create = [ + RefCommand::new(&writers::writers_ref(&ns), None, Some(&ns)), + RefCommand::new(&format!("refs/caos/w/{ns}/head"), None, Some(&oid('c'))), + ]; + let grants = Grants(vec![ + ("ff".into(), alice.public(), None), + ("b0b".into(), bob.public(), None), + ("eee".into(), alice.public(), Some(vec![oid('9')])), + ]); + let run = |t: &str| vec![Auth::Run { token: t.into() }.option()]; + assert!(check(&store, &create, &run("ff"), &grants).is_ok()); + assert!(check(&store, &create, &run("b0b"), &grants).is_err()); + assert!(check(&store, &create, &run("eee"), &grants).is_err()); + } + + #[test] + fn admission_takes_only_a_live_signature_over_this_request() { + let alice = WriterKey::from_seed([1; 32]); + let header = alice.sign_admission("GET", "/run?x"); + assert!(matches!( + admit(Some(&header), "GET", "/run?x"), + Ok(Writes::As { scope: None, .. }) + )); + assert!(admit(Some(&header), "GET", "/run?y").is_err()); + assert!(admit(Some(&header), "POST", "/run?x").is_err()); + assert!(matches!(admit(None, "GET", "/run?x"), Ok(Writes::None))); + } + + #[test] + fn hook_reaches_the_server_on_loopback() { + assert_eq!(hook_server_url("[::]:80"), "http://127.0.0.1:80"); + assert_eq!(hook_server_url("127.0.0.1:4567"), "http://127.0.0.1:4567"); + } +} diff --git a/rust/crates/server/src/repair.rs b/rust/crates/server/src/repair.rs index e928dbe7..8598dc09 100644 --- a/rust/crates/server/src/repair.rs +++ b/rust/crates/server/src/repair.rs @@ -155,11 +155,15 @@ enum IntegrityDepth { } const V3_PREFIX: &str = "refs/caos/v3/"; +/// Ref-writers namespaces (design/ref-writers.md): conversation heads and +/// `writers` lists, held to the same rules as v3 heads. +const NAMESPACE_PREFIX: &str = conversation_protocol::v3::writers::NAMESPACE_PREFIX; fn integrity_depth(refname: &str) -> IntegrityDepth { if refname.starts_with("refs/caos/req/") || refname.starts_with("refs/caos/res/") || refname.starts_with(V3_PREFIX) + || refname.starts_with(NAMESPACE_PREFIX) { IntegrityDepth::TargetOnly } else { @@ -170,9 +174,10 @@ fn integrity_depth(refname: &str) -> IntegrityDepth { // A v3 conversation head is never rewound: an older reflog value can name a // state whose write-ahead effects (a landed publication push, half of an atomic // spawn) already happened. The loose ref is dropped and the reflog kept for a -// human to restore from. +// human to restore from. Nor is a `writers` list: an older one can name a +// writer since removed. fn reflog_recoverable(refname: &str) -> bool { - !refname.starts_with(V3_PREFIX) + !refname.starts_with(V3_PREFIX) && !refname.starts_with(NAMESPACE_PREFIX) } #[derive(Default)] @@ -380,6 +385,18 @@ mod tests { static COUNTER: AtomicU32 = AtomicU32::new(0); + #[test] + fn a_namespaced_ref_is_checked_and_kept_like_a_v3_head() { + for name in [ + "refs/caos/v3/conversations/74616c6b/head", + "refs/caos/w/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa/writers", + ] { + assert!(integrity_depth(name) == IntegrityDepth::TargetOnly); + assert!(!reflog_recoverable(name)); + } + assert!(reflog_recoverable("refs/heads/main")); + } + fn temp_repo() -> (gix::ThreadSafeRepository, PathBuf) { let n = COUNTER.fetch_add(1, Ordering::Relaxed); let dir = std::env::temp_dir().join(format!("caos-repair-test-{}-{n}", std::process::id())); diff --git a/rust/crates/server/src/secrets.rs b/rust/crates/server/src/secrets.rs index bf034bf8..49ac998d 100644 --- a/rust/crates/server/src/secrets.rs +++ b/rust/crates/server/src/secrets.rs @@ -39,6 +39,10 @@ pub(crate) struct Context { conversation: Option, /// The conversation's head tree, when a `reader:@=` grant names it. head: Option, + /// Which writer this run acts for (design/ref-writers.md). Not a secret + /// and not in any key; it rides here because this is what reaches every + /// sub-run. + writes: crate::ref_writers::Writes, } impl Context { @@ -89,9 +93,19 @@ impl Context { pins: Arc::new(pins), conversation, head, + writes: crate::ref_writers::Writes::None, }) } + pub(crate) fn with_writes(mut self, writes: crate::ref_writers::Writes) -> Context { + self.writes = writes; + self + } + + pub(crate) fn writes(&self) -> &crate::ref_writers::Writes { + &self.writes + } + pub(crate) fn is_empty(&self) -> bool { self.stored.is_empty() } diff --git a/rust/crates/server/tests/pre_receive.rs b/rust/crates/server/tests/pre_receive.rs new file mode 100644 index 00000000..5787266a --- /dev/null +++ b/rust/crates/server/tests/pre_receive.rs @@ -0,0 +1,156 @@ +//! The ref-writers hook end to end: real `git push`es, with push options, into a +//! bare repository whose pre-receive hook is this crate's binary +//! (design/ref-writers.md). The unit tests in `ref_writers` cover the rules; +//! this covers reading a pushed `writers` list out of git's quarantine. + +use std::path::{Path, PathBuf}; +use std::process::Command; + +use conversation_protocol::v3::writers::{self, genesis_commit, Writer, WriterKey}; +use conversation_protocol::v3::{GitStore, ObjectStore, Oid, RefUpdate}; + +fn git(dir: &Path, args: &[&str]) { + let status = Command::new("git") + .arg("-C") + .arg(dir) + .args(args) + .status() + .unwrap(); + assert!(status.success(), "git {args:?}"); +} + +fn server_repo(name: &str) -> PathBuf { + let dir = std::env::temp_dir().join(format!("caos-pre-receive-{name}-{}", std::process::id())); + std::fs::remove_dir_all(&dir).ok(); + git( + Path::new("/"), + &["init", "-q", "--bare", dir.to_str().unwrap()], + ); + git(&dir, &["config", "receive.advertisePushOptions", "true"]); + let hook = dir.join("hooks").join("pre-receive"); + std::fs::write( + &hook, + format!( + "#!/bin/sh\nexec '{}' --pre-receive\n", + env!("CARGO_BIN_EXE_server") + ), + ) + .unwrap(); + use std::os::unix::fs::PermissionsExt; + std::fs::set_permissions(&hook, std::fs::Permissions::from_mode(0o755)).unwrap(); + dir +} + +fn client(name: &str, server: &Path, key: &'static WriterKey) -> (PathBuf, GitStore) { + let dir = std::env::temp_dir().join(format!( + "caos-pre-receive-client-{name}-{}", + std::process::id() + )); + std::fs::remove_dir_all(&dir).ok(); + git( + Path::new("/"), + &["init", "-q", "--bare", dir.to_str().unwrap()], + ); + git(&dir, &["remote", "add", "caos", server.to_str().unwrap()]); + let mut store = GitStore::open(&dir, Some("caos")).unwrap(); + store.set_push_auth(Box::new(move |commands: &[writers::Command]| { + key.sign_update(commands) + })); + (dir, store) +} + +fn update(refname: &str, expected: Option<&Oid>, new: &Oid) -> RefUpdate { + RefUpdate { + refname: refname.to_string(), + expected: expected.cloned(), + new: Some(new.clone()), + } +} + +#[test] +fn only_a_namespaces_writers_can_push_into_it() { + static ALICE: std::sync::LazyLock = + std::sync::LazyLock::new(|| WriterKey::from_seed([1; 32])); + static BOB: std::sync::LazyLock = + std::sync::LazyLock::new(|| WriterKey::from_seed([2; 32])); + let server = server_repo("writers"); + let (alice_dir, mut alice) = client("alice", &server, &ALICE); + let (bob_dir, bob) = client("bob", &server, &BOB); + + let list = vec![Writer { + key: ALICE.public(), + label: "alice".into(), + }]; + let ns = genesis_commit(&mut alice, &list, "test").unwrap(); + let tree = alice.write_tree(&[]).unwrap(); + let content = alice + .write_commit(&conversation_protocol::v3::CommitInfo { + tree, + parents: vec![], + author: sig(), + committer: sig(), + extra_headers: vec![], + message: b"content\n".to_vec(), + }) + .unwrap(); + let head = format!("refs/caos/w/{ns}/head"); + + // Creating the namespace and a ref in it, in one atomic push. + alice + .push(&[ + update(&writers::writers_ref(ns.as_str()), None, &ns), + update(&head, None, &ns), + ]) + .expect("alice creates her namespace"); + + // Bob is not a writer. + bob.fetch_object(&ns).unwrap(); + let refused = bob.push(&[update(&format!("refs/caos/w/{ns}/other"), None, &ns)]); + let refused = refused.expect_err("bob wrote alice's namespace"); + assert!( + refused.contains("not one of the namespace's writers"), + "{refused}" + ); + + // Alice adds Bob; now he can write. + let with_bob = writers::writers_commit( + &mut alice, + &ns, + &[ + list[0].clone(), + Writer { + key: BOB.public(), + label: "bob".into(), + }, + ], + &sig(), + "add bob", + ) + .unwrap(); + alice + .push(&[update( + &writers::writers_ref(ns.as_str()), + Some(&ns), + &with_bob, + )]) + .expect("alice adds bob"); + alice + .push(&[update(&head, Some(&ns), &content)]) + .expect("alice advances head"); + bob.fetch_object(&content).unwrap(); + bob.push(&[update(&format!("refs/caos/w/{ns}/other"), None, &content)]) + .expect("bob writes once added"); + + for dir in [server, alice_dir, bob_dir] { + std::fs::remove_dir_all(dir).ok(); + } +} + +fn sig() -> conversation_protocol::v3::Signature { + conversation_protocol::v3::Signature { + name: "t".into(), + email: "t@t".into(), + time: 1, + offset: "+0000".into(), + } +}