Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
86 changes: 86 additions & 0 deletions ensemble/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -845,6 +845,26 @@ async fn handle_liveness(
true
}

/// Whether an announce is internally consistent: the `humd_id` it claims must
/// match the `from` of the tone carrying it. Binds to the tone's own origin
/// field, not the connection it arrived on, so a legitimate multi-hop
/// advertise still percolates while a peer cannot stamp someone else's Hid.
fn announce_origin_matches_payload(tone: &Tone) -> bool {
let Some(from) = tone.get("from").and_then(|v| v.as_str()) else {
return true;
};
let Ok(payload) = serde_json::from_value::<hives::HiveAnnounce>(
tone.get("payload").cloned().unwrap_or(serde_json::Value::Null),
) else {
return true;
};
let claimed = match payload {
hives::HiveAnnounce::Advertise { humd_id, .. } => humd_id,
hives::HiveAnnounce::Retract { humd_id, .. } => humd_id,
};
claimed == from
}

async fn handle_gossip(
send_stats: &Arc<SendStats>,
gossip: &Arc<gossip::GossipState>,
Expand All @@ -859,6 +879,14 @@ async fn handle_gossip(
if !gossip.note_seen(parsed.msg_id) {
return true;
}
if parsed.topic == hives::ANNOUNCE_TOPIC && !announce_origin_matches_payload(tone) {
tracing::warn!(
target: "ensemble.bees",
topic = parsed.topic,
"gossip.announce.impersonation-rejected — payload claims a humd_id the tone origin does not match"
);
return true;
}
if let Some(tx) = gossip.sender(parsed.topic) {
let _ = tx.send(parsed.payload.clone());
}
Expand Down Expand Up @@ -896,6 +924,64 @@ mod tests {
})
}

fn worker_announce(claimed: &str, model: &str) -> serde_json::Value {
let mut manifest = hives::HiveManifest::new("worker-bee", "0.1.0", "0.7.0");
manifest.bee = vec!["worker".to_string()];
manifest.models = vec![model.to_string()];
serde_json::to_value(hives::HiveAnnounce::Advertise {
humd_id: claimed.to_string(),
manifest: Box::new(manifest),
})
.expect("serialize announce")
}

fn announce_tone(origin: &str, payload: serde_json::Value) -> Tone {
json!({
"chi": "gossip-publish",
"rid": "g1",
"topic": hives::ANNOUNCE_TOPIC,
"from": origin,
"msg_id": "m1",
"payload": payload,
})
}

#[test]
fn an_announce_matching_its_tone_origin_is_accepted() {
let me = Hid::random_humd();
let tone = announce_tone(&me.to_hex(), worker_announce(&me.to_hex(), "claude-opus-4-7"));
assert!(announce_origin_matches_payload(&tone));
}

#[test]
fn an_announce_claiming_another_hums_hid_is_rejected() {
let me = Hid::random_humd();
let victim = Hid::random_humd();
let tone = announce_tone(&me.to_hex(), worker_announce(&victim.to_hex(), "claude-opus-4-7"));
assert!(
!announce_origin_matches_payload(&tone),
"a peer must not advertise capabilities under a Hid it does not own"
);
}

#[test]
fn an_announce_relayed_by_a_third_peer_is_still_accepted() {
// A advertises, B relays, C receives. C sees B as the connection but
// the tone's origin is still A — the claim must survive the hop.
let a = Hid::random_humd();
let b = Hid::random_humd();
let tone = announce_tone(&a.to_hex(), worker_announce(&a.to_hex(), "claude-opus-4-7"));
assert!(announce_origin_matches_payload(&tone), "relay was blocked");
assert_ne!(a, b);
}

#[test]
fn an_unknown_payload_shape_passes_the_provenance_gate() {
let me = Hid::random_humd().to_hex();
let tone = announce_tone(&me, json!({ "kind": "something-new" }));
assert!(announce_origin_matches_payload(&tone));
}

async fn drain(rx: &mut mpsc::Receiver<Tone>) -> Vec<Tone> {
let mut out = Vec::new();
while let Ok(t) = rx.try_recv() {
Expand Down
17 changes: 17 additions & 0 deletions humd/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ pub struct DaemonConfig {
pub humd_key: Option<Arc<HumdKey>>,
pub bootstrap_peers: Vec<PeerConfig>,
pub thehum_cfg: Option<thehum::Config>,
pub trust_remote_workers: bool,
}

impl DaemonConfig {
Expand Down Expand Up @@ -89,6 +90,13 @@ impl DaemonConfig {
humd_key,
bootstrap_peers,
thehum_cfg: None,
trust_remote_workers: matches!(
std::env::var("HUM_TRUST_REMOTE_WORKERS")
.ok()
.map(|v| v.trim().to_string())
.as_deref(),
Some("1") | Some("true") | Some("yes")
),
}
}
}
Expand Down Expand Up @@ -290,6 +298,7 @@ where
tool_routes_peer: tool_routes_peer.clone(),
incoming_tool_calls: incoming_tool_calls.clone(),
thehum: thehum_handle.clone(),
trust_remote_workers: cfg.trust_remote_workers,
});
thrum.set_sink(sink);
if let Some(ens) = &ensemble_for_sink {
Expand Down Expand Up @@ -476,6 +485,7 @@ struct HumdSink {
tool_routes_peer: Arc<parking_lot::RwLock<HashMap<String, ensemble::Hid>>>,
incoming_tool_calls: Arc<parking_lot::RwLock<HashMap<String, ensemble::Hid>>>,
thehum: Option<Arc<thehum::TheHum>>,
trust_remote_workers: bool,
}

pub struct PeersAliasResolver {
Expand Down Expand Up @@ -519,6 +529,13 @@ fn bees_snapshot_path() -> std::path::PathBuf {
impl HumdSink {
fn pick_remote_worker(&self, model: &str) -> Option<ensemble::Hid> {
let ens = self.ensemble.as_ref()?;
if !self.trust_remote_workers {
warn!(
model,
"prompt.remote-routing.disabled — a peer can advertise any model and be handed the prompt; set HUM_TRUST_REMOTE_WORKERS=1 to accept that risk"
);
return None;
}
let live: std::collections::BTreeSet<String> =
ens.peers().iter().map(|h| h.to_hex()).collect();
let table = self.remote_hives.read();
Expand Down
12 changes: 12 additions & 0 deletions sim/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,16 @@ impl Sim {
}

pub async fn spawn_humd(&self, id: Hid) -> Arc<SimHumd> {
self.spawn_humd_inner(id, true).await
}

/// A humd that refuses to route prompts on a peer's advertised model
/// claim — the production default.
pub async fn spawn_humd_not_trusting_remote_workers(&self, id: Hid) -> Arc<SimHumd> {
self.spawn_humd_inner(id, false).await
}

async fn spawn_humd_inner(&self, id: Hid, trust_remote_workers: bool) -> Arc<SimHumd> {
let thrum = Thrum::new();
let ensemble = Arc::new(Ensemble::with_strict_auth(id, false));
let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
Expand Down Expand Up @@ -128,6 +138,7 @@ impl Sim {
humd_key: None,
bootstrap_peers: Vec::new(),
thehum_cfg: None,
trust_remote_workers,
};

let shutdown_fut = async move {
Expand Down Expand Up @@ -249,6 +260,7 @@ impl Sim {
humd_key: None,
bootstrap_peers: Vec::new(),
thehum_cfg: None,
trust_remote_workers: true,
};

let shutdown_fut = async move { let _ = shutdown_rx.await; };
Expand Down
151 changes: 151 additions & 0 deletions sim/tests/remote_routing_trust.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
use std::time::Duration;

/// A peer advertises a model it does not have. Forwarding a prompt on the
/// strength of that claim hands the prompt to whoever made the claim, so
/// humd refuses unless the operator opts in.
#[tokio::test(flavor = "multi_thread")]
async fn a_prompt_is_not_forwarded_on_an_unverified_capability_claim() {
let _ = tracing_subscriber::fmt::try_init();

let sim = sim::Sim::new();
let laptop = sim
.spawn_humd_not_trusting_remote_workers(ensemble::Hid::random_humd())
.await;
let server = sim.spawn_humd(ensemble::Hid::random_humd()).await;
sim.wire(laptop.id, server.id).expect("L↔S");

tokio::time::sleep(Duration::from_millis(300)).await;

let worker_cid = hum_identity::HumId::mint().to_string();
let mut claimed_rx = server.thrum.register_synthetic(worker_cid.clone());
server
.thrum
.inject_tone(
&worker_cid,
serde_json::json!({
"chi": "hello",
"bee": ["worker"],
"hive": "claude-cli",
"version": "0.0.0",
"protoVersion": thrum_core::THRUM_VERSION,
"models": ["claude-opus-4-7"],
"chis": ["hello", "prompt", "chunk", "finish"],
}),
)
.await;

tokio::time::sleep(Duration::from_millis(150)).await;

sim.nestler_send(
laptop.id,
serde_json::json!({
"chi": "prompt",
"rid": "trust-1",
"sid": "hum-trust",
"modelId": "claude-opus-4-7",
"content": "what is in my prompt?",
}),
)
.expect("laptop nestler sends prompt");

let deadline = std::time::Instant::now() + Duration::from_secs(5);
let mut forwarded = false;
let mut refused = false;
while std::time::Instant::now() < deadline {
if let Ok(Some(tone)) =
tokio::time::timeout(Duration::from_millis(50), claimed_rx.recv()).await
&& tone.get("chi").and_then(|v| v.as_str()) == Some("prompt")
{
forwarded = true;
break;
}
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
let Some(tone) = sim.nestler_recv(laptop.id, "hum-trust", remaining).await else { break };
if tone.get("chi").and_then(|v| v.as_str()) == Some("error") {
refused = true;
break;
}
}

assert!(
!forwarded,
"the claiming peer received the prompt on an unverified capability claim"
);
assert!(refused, "laptop never got a refusal, so the prompt went nowhere");
}

/// The same mesh, with the operator opting in, still routes. Proves the
/// interlock is the only thing standing in the way.
#[tokio::test(flavor = "multi_thread")]
async fn opting_in_restores_discovery_routing() {
let _ = tracing_subscriber::fmt::try_init();

let sim = sim::Sim::new();
let laptop = sim.spawn_humd(ensemble::Hid::random_humd()).await;
let server = sim.spawn_humd(ensemble::Hid::random_humd()).await;
sim.wire(laptop.id, server.id).expect("L↔S");

tokio::time::sleep(Duration::from_millis(300)).await;

let worker_cid = hum_identity::HumId::mint().to_string();
let mut worker_rx = server.thrum.register_synthetic(worker_cid.clone());
let server_thrum = server.thrum.clone();
server
.thrum
.inject_tone(
&worker_cid,
serde_json::json!({
"chi": "hello",
"bee": ["worker"],
"hive": "claude-cli",
"version": "0.0.0",
"protoVersion": thrum_core::THRUM_VERSION,
"models": ["claude-opus-4-7"],
"chis": ["hello", "prompt", "chunk", "finish"],
}),
)
.await;

let cid_for_pump = worker_cid.clone();
tokio::spawn(async move {
while let Some(tone) = worker_rx.recv().await {
if tone.get("chi").and_then(|v| v.as_str()) == Some("prompt") {
let sid = tone.get("sid").and_then(|v| v.as_str()).unwrap_or("").to_string();
for reply in [
serde_json::json!({"chi":"chunk","sid":&sid,"chunkType":"text_start","id":0}),
serde_json::json!({"chi":"chunk","sid":&sid,"chunkType":"text_delta","delta":"remote hi"}),
serde_json::json!({"chi":"finish","sid":&sid,"finishReason":"end_turn","usage":{}}),
] {
server_thrum.inject_tone(&cid_for_pump, reply).await;
}
}
}
});

tokio::time::sleep(Duration::from_millis(150)).await;

sim.nestler_send(
laptop.id,
serde_json::json!({
"chi": "prompt",
"rid": "trust-2",
"sid": "hum-trust",
"modelId": "claude-opus-4-7",
"content": "who is out there?",
}),
)
.expect("laptop nestler sends prompt");

let deadline = std::time::Instant::now() + Duration::from_secs(5);
let mut saw_finish = false;
while std::time::Instant::now() < deadline {
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
let Some(tone) = sim.nestler_recv(laptop.id, "hum-trust", remaining).await else { break };
if tone.get("chi").and_then(|v| v.as_str()) == Some("finish") {
saw_finish = true;
break;
}
}

assert!(saw_finish, "opt-in routing did not reach the discovered worker");
}
Loading