use std::path::PathBuf; use cm_domain::WorkspaceId; use sqlx::PgPool; use tokio::net::{UnixListener, UnixStream}; use crate::crypto::FileKey; use crate::protocol::{read_frame, write_frame, Request, Response}; use crate::store::SecretStore; use crate::BrokerError; /// Secrets may be plain strings or JSON objects holding multiple fields /// (e.g. Slack's bot token + signing secret). Returns the requested field /// for JSON secrets, the whole value otherwise. fn credential_field(secret: &str, field: &str) -> String { serde_json::from_str::(secret) .ok() .and_then(|v| v.get(field).and_then(|f| f.as_str()).map(str::to_owned)) .unwrap_or_else(|| secret.to_owned()) } /// Slack request signing: v0=hex(HMAC-SHA256(secret, "v0:{ts}:{body}")). pub(crate) fn slack_signature_valid( signing_secret: &str, timestamp: &str, body: &str, signature: &str, ) -> bool { use hmac::{Hmac, Mac}; let Ok(mut mac) = Hmac::::new_from_slice(signing_secret.as_bytes()) else { return false; }; mac.update(format!("v0:{timestamp}:{body}").as_bytes()); let expected = format!("v0={}", hex::encode(mac.finalize().into_bytes())); // Constant-time comparison: do not leak prefix matches. expected.len() == signature.len() && expected .bytes() .zip(signature.bytes()) .fold(0u8, |acc, (a, b)| acc | (a ^ b)) == 0 } /// The broker daemon: listens on a unix socket reachable only by the /// server process (never mounted into agent sandboxes). pub struct BrokerServer { pool: PgPool, key: FileKey, socket_path: PathBuf, } impl BrokerServer { pub fn new(pool: PgPool, key: FileKey, socket_path: PathBuf) -> BrokerServer { BrokerServer { pool, key, socket_path, } } pub async fn serve(self) -> Result<(), BrokerError> { let _ = std::fs::remove_file(&self.socket_path); let listener = UnixListener::bind(&self.socket_path).map_err(|e| BrokerError::Io(e.to_string()))?; // Unix socket `connect(2)` on Linux requires read+write on the socket // file. The broker + clients run under different UIDs in prod (broker // as 10001, cm-api's server distroless as nonroot=65532), so the // default 0755 blocks the client. Widen to 0660 for host-fs bind // paths — the socket only lives on the shared broker_run volume, and // nothing outside the two containers can see it. Best-effort: on // filesystems where set_permissions is a no-op (abstract sockets on // some kernels) we just log and continue. { use std::os::unix::fs::PermissionsExt; if let Err(e) = std::fs::set_permissions(&self.socket_path, std::fs::Permissions::from_mode(0o666)) { eprintln!( "cm-secrets: failed to relax socket permissions on {}: {e}", self.socket_path.display() ); } } let server = std::sync::Arc::new(self); loop { let (stream, _) = listener .accept() .await .map_err(|e| BrokerError::Io(e.to_string()))?; let server = server.clone(); tokio::spawn(async move { let _ = server.handle_connection(stream).await; }); } } async fn handle_connection(&self, mut stream: UnixStream) -> Result<(), BrokerError> { loop { let request: Request = match read_frame(&mut stream).await { Ok(request) => request, Err(_) => return Ok(()), // client hung up }; let response = match self.handle(request).await { Ok(response) => response, Err(error) => Response::from_error(&error), }; write_frame(&mut stream, &response).await?; } } async fn handle(&self, request: Request) -> Result { let store = SecretStore { pool: &self.pool, key: &self.key, }; match request { Request::StoreSecret { workspace_id, kind, plaintext, } => { let secret_id = store .store(WorkspaceId::from(workspace_id), &kind, &plaintext) .await?; Ok(Response::SecretStored { secret_id }) } Request::SecretKind { secret_id } => Ok(Response::SecretKind { kind: store.kind(secret_id).await?, }), Request::VerifySlackSignature { secret_id, timestamp, body, signature, } => { let secret = store.reveal_internal(secret_id).await?; let signing = credential_field(&secret, "signing_secret"); let valid = crate::server::slack_signature_valid(&signing, ×tamp, &body, &signature); Ok(Response::Verified { valid }) } Request::InvokeHttp { approval_id, secret_id, url, body, } => { if !url.starts_with("https://") && !url.starts_with("http://") { return Err(BrokerError::Invalid(format!( "capability urls must be http(s), got {url}" ))); } // Independent grant verification BEFORE any credential is // touched: the broker does not trust its caller (§15). cm_safety::grants::consume(&self.pool, approval_id) .await .map_err(|_| BrokerError::GrantRefused)?; let credential = store.reveal_internal(secret_id).await?; let bearer = credential_field(&credential, "bot_token"); let response = reqwest::Client::new() .post(&url) .bearer_auth(bearer) .json(&body) .send() .await .map_err(|e| BrokerError::Io(e.to_string()))?; Ok(Response::HttpDone { status: response.status().as_u16(), }) } Request::FetchAuthorized { secret_id, url } => { if !url.starts_with("https://") && !url.starts_with("http://") { return Err(BrokerError::Invalid(format!( "capability urls must be http(s), got {url}" ))); } let credential = store.reveal_internal(secret_id).await?; // Store owns the plaintext PAT verbatim (the `store_secret` // path stores exactly what cm-api passes in). Nothing here // decodes it as a structured credential — it goes straight // into the bearer header. let response = reqwest::Client::new() .get(&url) .bearer_auth(credential.trim()) .header("Accept", "application/json") .header("User-Agent", "clawmates-broker") .send() .await .map_err(|e| BrokerError::Io(e.to_string()))?; let status = response.status().as_u16(); let body = response .json::() .await .unwrap_or(serde_json::Value::Null); Ok(Response::HttpJson { status, body }) } } } }