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
11 changes: 9 additions & 2 deletions WIRE.md
Original file line number Diff line number Diff line change
Expand Up @@ -105,8 +105,15 @@ single protocol.
ensemble on the `hum/hives/announce` topic.
- humd replies with `chi:"breath"` — a snapshot of any state relevant
to this nestler (today: `{}`; reserved for future state sync).
- A `protoVersion` mismatch is **a warning, not a hard error**. The
bumping rules:
- A hello with **no** `protoVersion` is **refused**: the bee is not
registered, not announced to the ensemble, and the connection is closed.
- A `protoVersion` whose **major** differs from humd's is **refused** the
same way.
- A `protoVersion` that differs only in minor or patch is admitted, with a
`thrum.hello.proto-drift` warning.
- An unparseable `protoVersion` is refused.

The bumping rules humd applies:
- **patch** — docstring tweaks, additive-optional fields
- **minor** — new chi value, new required field with compat path
- **major** — removed chi, renamed chi, semantics changed
Expand Down
52 changes: 51 additions & 1 deletion humd/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use parking_lot::RwLock;
use serde_json::Value;
use thrumd::{serve_with_hook as thrum_serve_with_hook, Thrum, Tone, ToneSink};
use thrum_core::{Chi, WaneTracker};
use tracing::{info, trace, warn};
use tracing::{error, info, trace, warn};

mod drone;
mod drift;
Expand Down Expand Up @@ -811,6 +811,56 @@ impl ToneSink for HumdSink {
match chi {
Some(Chi::Hello) => {
trace!(client_id, %chi_str, "thrum.recv.hello");

let rid = tone
.get("rid")
.and_then(Value::as_str)
.unwrap_or("hello-1")
.to_string();

match tone.get("protoVersion").and_then(Value::as_str) {
None => {
error!(
client_id,
"thrum.hello.proto-missing — hello declares no protoVersion. \
Not registered, not announced to the ensemble."
);
self.thrum.thrum_to(
client_id,
thrumd::echo_tone(&rid, false, Some("protoVersion is required")),
);
self.thrum.thrum_close(client_id);
return;
}
Some(declared) => {
match thrumd::proto_verdict(declared, thrum_core::THRUM_VERSION) {
thrumd::ProtoVerdict::Same => {}
thrumd::ProtoVerdict::SameMajor => {
warn!(
client_id,
declared,
current = thrum_core::THRUM_VERSION,
"thrum.hello.proto-drift — admitted within major"
);
}
thrumd::ProtoVerdict::Incompatible => {
error!(
client_id,
declared,
current = thrum_core::THRUM_VERSION,
"thrum.hello.proto-incompatible — rejected, not announced"
);
self.thrum.thrum_to(
client_id,
thrumd::echo_tone(&rid, false, Some("protoVersion incompatible")),
);
self.thrum.thrum_close(client_id);
return;
}
}
}
}

let breath = thrumd::breath_tone(serde_json::json!({}));
self.thrum.thrum_to(client_id, breath);

Expand Down
64 changes: 64 additions & 0 deletions humd/src/thrumd.rs
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,15 @@ impl Thrum {
self.inner.clients.write().remove(client_id);
}


/// Drop a client from the registry. Any tones already queued for it
/// still drain; once the queue empties the socket handler sees the
/// channel close and the connection ends. Used to refuse a bee that
/// failed admission instead of leaving it connected but unregistered.
pub fn thrum_close(&self, client_id: &str) {
self.inner.clients.write().remove(client_id);
}

/// Inject a tone as if it arrived from `client_id`. Bypasses the
/// NDJSON socket and goes straight to the installed sink. No
/// envelope validation, no auto-echo — sim drives the shape it
Expand Down Expand Up @@ -310,3 +319,58 @@ pub fn echo_tone(rid: &str, ok: bool, error: Option<&str>) -> Tone {
}
v
}


/// How a bee's declared `protoVersion` compares to the one humd speaks.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ProtoVerdict {
Same,
SameMajor,
Incompatible,
}

fn parse_proto(v: &str) -> Option<(u64, u64, u64)> {
let mut it = v.trim().split('.');
let major = it.next()?.parse().ok()?;
let minor = it.next().and_then(|s| s.parse().ok()).unwrap_or(0);
let patch = it.next().and_then(|s| s.parse().ok()).unwrap_or(0);
Some((major, minor, patch))
}

pub fn proto_verdict(declared: &str, current: &str) -> ProtoVerdict {
match (parse_proto(declared), parse_proto(current)) {
(Some(d), Some(c)) if d == c => ProtoVerdict::Same,
(Some(d), Some(c)) if d.0 == c.0 => ProtoVerdict::SameMajor,
_ => ProtoVerdict::Incompatible,
}
}

#[cfg(test)]
mod proto_tests {
use super::*;

#[test]
fn an_exact_match_is_same() {
assert_eq!(proto_verdict("0.7.0", "0.7.0"), ProtoVerdict::Same);
}

#[test]
fn drift_inside_a_major_is_admitted() {
assert_eq!(proto_verdict("0.6.2", "0.7.0"), ProtoVerdict::SameMajor);
assert_eq!(proto_verdict("0.9.9", "0.7.0"), ProtoVerdict::SameMajor);
assert_eq!(proto_verdict("0.7", "0.7.0"), ProtoVerdict::Same);
}

#[test]
fn a_different_major_is_incompatible() {
assert_eq!(proto_verdict("1.0.0", "0.7.0"), ProtoVerdict::Incompatible);
assert_eq!(proto_verdict("0.7.0", "1.0.0"), ProtoVerdict::Incompatible);
}

#[test]
fn an_unparseable_declaration_is_incompatible() {
assert_eq!(proto_verdict("banana", "0.7.0"), ProtoVerdict::Incompatible);
assert_eq!(proto_verdict("", "0.7.0"), ProtoVerdict::Incompatible);
assert_eq!(proto_verdict("x.y.z", "0.7.0"), ProtoVerdict::Incompatible);
}
}
95 changes: 95 additions & 0 deletions sim/tests/admission_gate.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
use std::time::Duration;

#[tokio::test(flavor = "multi_thread")]
async fn a_hello_with_no_proto_version_is_disconnected() {
let _ = tracing_subscriber::fmt::try_init();
let sim = sim::Sim::new();
let humd = sim.spawn_humd(ensemble::Hid::random_humd()).await;
sim.await_ready(humd.id).await.expect("humd ready");

let cid = hum_identity::HumId::mint().to_string();
let _rx = humd.thrum.register_synthetic(cid.clone());
assert!(humd.thrum.is_connected(&cid));

humd.thrum
.inject_tone(
&cid,
serde_json::json!({
"chi": "hello",
"bee": ["worker"],
"hive": "claude-cli",
"version": "0.1.0",
"chis": ["hello", "prompt"],
}),
)
.await;

tokio::time::sleep(Duration::from_millis(120)).await;
assert!(
!humd.thrum.is_connected(&cid),
"a bee declaring no protoVersion must not stay connected"
);
}

#[tokio::test(flavor = "multi_thread")]
async fn a_hello_from_another_major_is_disconnected() {
let _ = tracing_subscriber::fmt::try_init();
let sim = sim::Sim::new();
let humd = sim.spawn_humd(ensemble::Hid::random_humd()).await;
sim.await_ready(humd.id).await.expect("humd ready");

let cid = hum_identity::HumId::mint().to_string();
let _rx = humd.thrum.register_synthetic(cid.clone());
assert!(humd.thrum.is_connected(&cid));

humd.thrum
.inject_tone(
&cid,
serde_json::json!({
"chi": "hello",
"bee": ["worker"],
"hive": "claude-cli",
"version": "0.1.0",
"protoVersion": "9.0.0",
"chis": ["hello", "prompt"],
}),
)
.await;

tokio::time::sleep(Duration::from_millis(120)).await;
assert!(
!humd.thrum.is_connected(&cid),
"a bee speaking another major must not stay connected"
);
}

#[tokio::test(flavor = "multi_thread")]
async fn a_hello_within_the_same_major_stays_connected() {
let _ = tracing_subscriber::fmt::try_init();
let sim = sim::Sim::new();
let humd = sim.spawn_humd(ensemble::Hid::random_humd()).await;
sim.await_ready(humd.id).await.expect("humd ready");

let cid = hum_identity::HumId::mint().to_string();
let _rx = humd.thrum.register_synthetic(cid.clone());

humd.thrum
.inject_tone(
&cid,
serde_json::json!({
"chi": "hello",
"bee": ["worker"],
"hive": "claude-cli",
"version": "0.1.0",
"protoVersion": "0.1.0",
"chis": ["hello", "prompt"],
}),
)
.await;

tokio::time::sleep(Duration::from_millis(120)).await;
assert!(
humd.thrum.is_connected(&cid),
"drift inside the major must stay admitted"
);
}
Loading