use std::sync::Arc;
use std::sync::mpsc::Receiver;
use cargoless_proto::Diagnostic;
use crate::config::{FleetConfig, FleetConfigError};
pub mod discovery;
pub mod http;
pub mod inproc;
pub mod unix;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CrateVerdict {
pub name: String,
pub verdict: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WorktreeStatus {
pub worktree: String,
pub verdict: String,
pub crates: Vec<CrateVerdict>,
pub red_diagnostics: u32,
pub heartbeat_age_secs: u64,
pub published_at: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WorktreeSummary {
pub worktree: String,
pub verdict: String,
pub red_diagnostics: u32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TransitionEvent {
pub worktree: String,
pub verdict: String,
pub red_diagnostics: u32,
pub published_at: u64,
}
pub trait VerdictService: Send + Sync {
fn get_status(&self, worktree: &str) -> Option<WorktreeStatus>;
fn get_verdict(&self, worktree: &str) -> Option<String>;
fn get_diagnostics(&self, worktree: &str) -> Vec<Diagnostic>;
fn list_worktrees(&self) -> Vec<WorktreeSummary>;
fn subscribe(&self) -> Receiver<TransitionEvent>;
}
pub trait TransportClient {
fn get_status(&self, worktree: &str) -> Result<Option<WorktreeStatus>, TransportError>;
fn get_verdict(&self, worktree: &str) -> Result<Option<String>, TransportError>;
fn get_diagnostics(&self, worktree: &str) -> Result<Vec<Diagnostic>, TransportError>;
fn list_worktrees(&self) -> Result<Vec<WorktreeSummary>, TransportError>;
fn subscribe(&self) -> Result<Receiver<TransitionEvent>, TransportError>;
}
pub trait Authorizer: Send + Sync {
fn authorize(&self, token: Option<&str>) -> bool;
}
#[derive(Debug, Clone, Copy, Default)]
pub struct AllowAll;
impl Authorizer for AllowAll {
fn authorize(&self, _token: Option<&str>) -> bool {
true
}
}
pub struct BearerToken {
secret: Vec<u8>,
}
impl BearerToken {
pub fn new(secret: impl Into<String>) -> Self {
Self {
secret: secret.into().into_bytes(),
}
}
}
#[inline(never)]
fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
if a.len() != b.len() {
return false;
}
let mut diff: u8 = 0;
for (x, y) in a.iter().zip(b.iter()) {
diff |= x ^ y;
}
diff == 0
}
impl Authorizer for BearerToken {
fn authorize(&self, token: Option<&str>) -> bool {
if self.secret.iter().all(u8::is_ascii_whitespace) {
return false;
}
match token {
None => false,
Some(presented) => constant_time_eq(presented.as_bytes(), &self.secret),
}
}
}
pub fn authorizer_for(cfg: &FleetConfig) -> Result<Arc<dyn Authorizer>, FleetConfigError> {
cfg.security_check()?;
Ok(match cfg.effective_auth_token() {
Some(secret) => Arc::new(BearerToken::new(secret)),
None => Arc::new(AllowAll),
})
}
#[derive(Debug)]
pub enum TransportError {
Io(std::io::Error),
Protocol(String),
Unauthorized,
}
impl std::fmt::Display for TransportError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
TransportError::Io(e) => write!(f, "transport I/O: {e}"),
TransportError::Protocol(m) => write!(f, "transport protocol: {m}"),
TransportError::Unauthorized => write!(f, "transport: unauthorized"),
}
}
}
impl std::error::Error for TransportError {}
impl From<std::io::Error> for TransportError {
fn from(e: std::io::Error) -> Self {
TransportError::Io(e)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Request {
GetStatus(String),
GetVerdict(String),
GetDiagnostics(String),
ListWorktrees,
Subscribe,
}
impl Request {
pub fn from_json(text: &str) -> Option<Request> {
let v: serde_json::Value = serde_json::from_str(text).ok()?;
let op = v.get("op")?.as_str()?;
let wt = || {
v.get("worktree")
.and_then(serde_json::Value::as_str)
.unwrap_or("")
.to_string()
};
match op {
"get_status" => Some(Request::GetStatus(wt())),
"get_verdict" => Some(Request::GetVerdict(wt())),
"get_diagnostics" => Some(Request::GetDiagnostics(wt())),
"list_worktrees" => Some(Request::ListWorktrees),
"subscribe" => Some(Request::Subscribe),
_ => None,
}
}
pub fn to_json(&self) -> String {
let v = match self {
Request::GetStatus(w) => serde_json::json!({"op":"get_status","worktree":w}),
Request::GetVerdict(w) => serde_json::json!({"op":"get_verdict","worktree":w}),
Request::GetDiagnostics(w) => {
serde_json::json!({"op":"get_diagnostics","worktree":w})
}
Request::ListWorktrees => serde_json::json!({"op":"list_worktrees"}),
Request::Subscribe => serde_json::json!({"op":"subscribe"}),
};
v.to_string()
}
}
fn crate_verdicts_json(crates: &[CrateVerdict]) -> serde_json::Value {
serde_json::Value::Array(
crates
.iter()
.map(|c| serde_json::json!({"name": c.name, "verdict": c.verdict}))
.collect(),
)
}
fn crate_verdicts_from_json(v: Option<&serde_json::Value>) -> Vec<CrateVerdict> {
let Some(serde_json::Value::Array(items)) = v else {
return Vec::new();
};
items
.iter()
.filter_map(|c| {
Some(CrateVerdict {
name: c.get("name")?.as_str()?.to_string(),
verdict: c
.get("verdict")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown")
.to_string(),
})
})
.collect()
}
pub fn status_to_json(s: &WorktreeStatus) -> String {
serde_json::json!({
"worktree": s.worktree,
"verdict": s.verdict,
"crates": crate_verdicts_json(&s.crates),
"red_diagnostics": s.red_diagnostics,
"heartbeat_age_secs": s.heartbeat_age_secs,
"published_at": s.published_at,
})
.to_string()
}
pub fn status_from_json(text: &str) -> Option<WorktreeStatus> {
let v: serde_json::Value = serde_json::from_str(text).ok()?;
Some(WorktreeStatus {
worktree: v.get("worktree")?.as_str()?.to_string(),
verdict: v
.get("verdict")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown")
.to_string(),
crates: crate_verdicts_from_json(v.get("crates")),
red_diagnostics: v
.get("red_diagnostics")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0) as u32,
heartbeat_age_secs: v
.get("heartbeat_age_secs")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0),
published_at: v
.get("published_at")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0),
})
}
pub fn summaries_to_json(list: &[WorktreeSummary]) -> String {
serde_json::Value::Array(
list.iter()
.map(|s| {
serde_json::json!({
"worktree": s.worktree,
"verdict": s.verdict,
"red_diagnostics": s.red_diagnostics,
})
})
.collect(),
)
.to_string()
}
pub fn summaries_from_json(text: &str) -> Vec<WorktreeSummary> {
let Ok(serde_json::Value::Array(items)) = serde_json::from_str::<serde_json::Value>(text)
else {
return Vec::new();
};
items
.iter()
.filter_map(|s| {
Some(WorktreeSummary {
worktree: s.get("worktree")?.as_str()?.to_string(),
verdict: s
.get("verdict")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown")
.to_string(),
red_diagnostics: s
.get("red_diagnostics")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0) as u32,
})
})
.collect()
}
pub fn event_to_json(e: &TransitionEvent) -> String {
serde_json::json!({
"worktree": e.worktree,
"verdict": e.verdict,
"red_diagnostics": e.red_diagnostics,
"published_at": e.published_at,
})
.to_string()
}
pub fn event_from_json(text: &str) -> Option<TransitionEvent> {
let v: serde_json::Value = serde_json::from_str(text).ok()?;
Some(TransitionEvent {
worktree: v.get("worktree")?.as_str()?.to_string(),
verdict: v
.get("verdict")
.and_then(serde_json::Value::as_str)
.unwrap_or("unknown")
.to_string(),
red_diagnostics: v
.get("red_diagnostics")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0) as u32,
published_at: v
.get("published_at")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::FleetConfig;
#[test]
fn bearer_token_accepts_exact_denies_wrong_and_none() {
let a = BearerToken::new("s3cr3t-abc");
assert!(a.authorize(Some("s3cr3t-abc")), "exact match ⇒ allow");
assert!(!a.authorize(Some("s3cr3t-abd")), "1-byte-off ⇒ deny");
assert!(!a.authorize(Some("s3cr3t-ab")), "prefix (shorter) ⇒ deny");
assert!(
!a.authorize(Some("s3cr3t-abcd")),
"superstring (longer) ⇒ deny"
);
assert!(!a.authorize(Some("")), "empty presented ⇒ deny");
assert!(!a.authorize(None), "no credential ⇒ deny (→ adapter 401)");
}
#[test]
fn constant_time_eq_is_correct_total_and_length_safe() {
assert!(constant_time_eq(b"", b""));
assert!(constant_time_eq(b"abc", b"abc"));
assert!(!constant_time_eq(b"abc", b"abd")); assert!(!constant_time_eq(b"abc", b"Xbc")); assert!(!constant_time_eq(b"abc", b"ab")); assert!(!constant_time_eq(b"ab", b"abc"));
assert_eq!(
constant_time_eq(b"\x00xxxxxxxx", b"\xffxxxxxxxx"),
constant_time_eq(b"xxxxxxxx\x00", b"xxxxxxxx\xff"),
"mismatch position must not change the result path"
);
}
fn cfg_bind_token(bind: Option<&str>, token: Option<&str>) -> FleetConfig {
let mut c = FleetConfig::defaults();
c.bind = bind.map(|b| b.parse().expect("test bind addr"));
c.auth_token = token.map(str::to_string);
c
}
#[test]
fn authorizer_for_loopback_no_token_is_allowall_open_posture() {
let c = cfg_bind_token(Some("127.0.0.1:8080"), None);
let a = authorizer_for(&c).expect("loopback no-token must not error");
assert!(a.authorize(None), "AllowAll ⇒ no-token allowed on loopback");
assert!(a.authorize(Some("whatever")));
}
#[test]
fn authorizer_for_non_loopback_no_token_fails_closed() {
let c = cfg_bind_token(Some("0.0.0.0:8080"), None);
let r = authorizer_for(&c);
assert!(
matches!(r, Err(FleetConfigError::BadBind { .. })),
"non-loopback + no token MUST be a refused config error \
(Ok would mean a public socket got a silent AllowAll)"
);
}
#[test]
fn authorizer_for_token_present_is_bearer_even_on_loopback() {
let c = cfg_bind_token(Some("127.0.0.1:8080"), Some("tok-XYZ"));
let a = authorizer_for(&c).expect("token present ⇒ ok");
assert!(a.authorize(Some("tok-XYZ")), "correct token allowed");
assert!(!a.authorize(Some("tok-xyz")), "wrong token denied");
assert!(!a.authorize(None), "no token denied when policy is bearer");
}
#[test]
fn authorizer_for_non_loopback_with_token_is_bearer_enforced() {
let c = cfg_bind_token(Some("0.0.0.0:8080"), Some("net-secret"));
let a = authorizer_for(&c).expect("non-loopback + token ⇒ ok");
assert!(a.authorize(Some("net-secret")));
assert!(!a.authorize(Some("net-secre")));
assert!(
!a.authorize(None),
"public bind w/ bearer ⇒ no-token denied"
);
}
#[test]
fn authorizer_for_non_loopback_blank_token_fails_closed() {
for blank in ["", " ", "\t "] {
let c = cfg_bind_token(Some("0.0.0.0:8080"), Some(blank));
let r = authorizer_for(&c);
assert!(
matches!(r, Err(FleetConfigError::BadBind { .. })),
"non-loopback + blank {blank:?} MUST refuse (got Ok ⇒ \
unauthenticated public socket)"
);
}
}
#[test]
fn bearer_with_empty_or_blank_secret_authorizes_nothing() {
for blank in ["", " ", "\t"] {
let bt = BearerToken::new(blank);
assert!(!bt.authorize(None), "blank-secret bearer denies None");
assert!(
!bt.authorize(Some("")),
"blank-secret bearer denies empty presented"
);
assert!(
!bt.authorize(Some(blank)),
"blank-secret bearer denies the blank itself"
);
assert!(
!bt.authorize(Some("anything")),
"blank-secret bearer denies any token"
);
}
let c = cfg_bind_token(Some("127.0.0.1:8080"), Some(" "));
let a = authorizer_for(&c).expect("loopback blank ⇒ AllowAll, not Err");
assert!(a.authorize(None), "loopback no-effective-token ⇒ AllowAll");
}
#[test]
fn authorizer_for_no_bind_defaults_open_v0_compat() {
let c = FleetConfig::defaults();
let a = authorizer_for(&c).expect("no bind ⇒ no auth required");
assert!(a.authorize(None));
}
#[test]
fn request_roundtrips_and_rejects_unknown_op() {
for r in [
Request::GetStatus("w1".into()),
Request::GetVerdict("w2".into()),
Request::GetDiagnostics("w3".into()),
Request::ListWorktrees,
Request::Subscribe,
] {
assert_eq!(Request::from_json(&r.to_json()), Some(r.clone()), "{r:?}");
}
assert_eq!(Request::from_json(r#"{"op":"nope"}"#), None);
assert_eq!(Request::from_json("not json"), None);
assert_eq!(Request::from_json("{}"), None);
}
#[test]
fn status_roundtrips_including_empty_crates_honesty_case() {
let s = WorktreeStatus {
worktree: "tf-mv-flat".into(),
verdict: "red".into(),
crates: vec![],
red_diagnostics: 3,
heartbeat_age_secs: 2,
published_at: 1234567890,
};
assert_eq!(status_from_json(&status_to_json(&s)), Some(s));
let s2 = WorktreeStatus {
worktree: "tf-mv-check".into(),
verdict: "red".into(),
crates: vec![
CrateVerdict {
name: "isolation".into(),
verdict: "green".into(),
},
CrateVerdict {
name: "physics".into(),
verdict: "red".into(),
},
],
red_diagnostics: 1,
heartbeat_age_secs: 0,
published_at: 42,
};
assert_eq!(status_from_json(&status_to_json(&s2)), Some(s2));
}
#[test]
fn status_from_json_is_best_effort_never_panics() {
assert_eq!(status_from_json(""), None);
assert_eq!(status_from_json("garbage"), None);
assert_eq!(status_from_json("{}"), None); let s = status_from_json(r#"{"worktree":"w"}"#).unwrap();
assert_eq!(s.verdict, "unknown");
assert_eq!(s.red_diagnostics, 0);
assert!(s.crates.is_empty());
}
#[test]
fn summaries_roundtrip_and_tolerate_malformed_elements() {
let list = vec![
WorktreeSummary {
worktree: "a".into(),
verdict: "green".into(),
red_diagnostics: 0,
},
WorktreeSummary {
worktree: "b".into(),
verdict: "red".into(),
red_diagnostics: 2,
},
];
assert_eq!(summaries_from_json(&summaries_to_json(&list)), list);
assert_eq!(
summaries_from_json(
r#"[{"verdict":"green"},{"worktree":"ok","verdict":"red","red_diagnostics":1}]"#
),
vec![WorktreeSummary {
worktree: "ok".into(),
verdict: "red".into(),
red_diagnostics: 1
}]
);
}
#[test]
fn allow_all_authorizes_with_or_without_token() {
assert!(AllowAll.authorize(None));
assert!(AllowAll.authorize(Some("anything")));
}
}