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
40 changes: 26 additions & 14 deletions ensemble/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -478,27 +478,19 @@ impl Ensemble {
rx
}

pub fn hive_discover(
&self,
name: impl Into<String>,
) -> mpsc::Receiver<(Hid, hives::HiveManifest)> {
let needle = name.into();
pub fn hive_discover_all(&self) -> mpsc::Receiver<(Hid, hives::HiveManifest)> {
let mut raw = self.subscribe_topic(hives::ANNOUNCE_TOPIC);
let (tx, rx) = mpsc::channel(64);
tokio::spawn(async move {
loop {
match raw.recv().await {
Ok(v) => {
let parsed: Result<hives::HiveAnnounce, _> = serde_json::from_value(v);
if let Ok(hives::HiveAnnounce::Advertise { humd_id, manifest }) = parsed {
if manifest.name != needle {
continue;
}
if let Ok(id) = Hid::from_hex(&humd_id)
&& tx.send((id, *manifest)).await.is_err()
{
break;
}
if let Ok(hives::HiveAnnounce::Advertise { humd_id, manifest }) = parsed
&& let Ok(id) = Hid::from_hex(&humd_id)
&& tx.send((id, *manifest)).await.is_err()
{
break;
}
}
Err(broadcast::error::RecvError::Closed) => break,
Expand All @@ -509,6 +501,26 @@ impl Ensemble {
rx
}

pub fn hive_discover(
&self,
name: impl Into<String>,
) -> mpsc::Receiver<(Hid, hives::HiveManifest)> {
let needle = name.into();
let mut raw = self.hive_discover_all();
let (tx, rx) = mpsc::channel(64);
tokio::spawn(async move {
while let Some((id, manifest)) = raw.recv().await {
if manifest.name != needle {
continue;
}
if tx.send((id, manifest)).await.is_err() {
break;
}
}
});
rx
}

pub async fn kad_find(&self, target: Hid, timeout: Duration) -> Option<HumdAddr> {
if let Some(addr) = self.kad.get(&target) {
return Some(addr);
Expand Down
83 changes: 81 additions & 2 deletions humd/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,8 @@ where
let observers: Observers = Arc::new(RwLock::new(HashMap::new()));
let hive_tag = cfg.hum_cfg.nest.default.clone();
let manifests: Manifests = Arc::new(parking_lot::RwLock::new(HashMap::new()));
let remote_hives: RemoteHives =
Arc::new(parking_lot::RwLock::new(HashMap::new()));
if let Some(thehum) = thehum_handle.as_ref() {
let manifests_for_replay = manifests.clone();
if let Err(e) = thehum.replay(|event| {
Expand Down Expand Up @@ -279,6 +281,7 @@ where
capacity: cfg.capacity,
hive_tag: hive_tag.clone(),
manifests: manifests.clone(),
remote_hives: remote_hives.clone(),
bees_snapshot_path: bees_snapshot_path(),
sid_origins: sid_origins.clone(),
tool_routes,
Expand All @@ -289,6 +292,17 @@ where
thehum: thehum_handle.clone(),
});
thrum.set_sink(sink);
if let Some(ens) = &ensemble_for_sink {
let ens = ens.clone();
let remote = remote_hives.clone();
tokio::spawn(async move {
let mut seen = ens.hive_discover_all();
while let Some((humd_id, manifest)) = seen.recv().await {
let key = bee_key(&manifest);
remote.write().entry(humd_id).or_default().insert(key, manifest);
}
});
}
if bind_thrum {
let thrum = thrum.clone();
let path = cfg.thrum_path.clone();
Expand Down Expand Up @@ -453,6 +467,7 @@ struct HumdSink {
capacity: LocalCapacity,
hive_tag: String,
manifests: Manifests,
remote_hives: RemoteHives,
bees_snapshot_path: std::path::PathBuf,
sid_origins: Arc<parking_lot::RwLock<HashMap<String, ensemble::Hid>>>,
tool_routes: Arc<parking_lot::RwLock<HashMap<String, String>>>,
Expand Down Expand Up @@ -486,12 +501,41 @@ impl ensemble::AliasResolver for PeersAliasResolver {
}

type Manifests = Arc<parking_lot::RwLock<HashMap<String, ensemble::HiveManifest>>>;
type RemoteHives =
Arc<parking_lot::RwLock<HashMap<ensemble::Hid, HashMap<String, ensemble::HiveManifest>>>>;

fn bee_key(manifest: &ensemble::HiveManifest) -> String {
manifest
.hid
.map(|h| h.to_hex())
.or_else(|| manifest.nestler_id.clone())
.unwrap_or_else(|| manifest.name.clone())
}

fn bees_snapshot_path() -> std::path::PathBuf {
hum_paths::bees_snapshot()
}

impl HumdSink {
fn pick_remote_worker(&self, model: &str) -> Option<ensemble::Hid> {
let ens = self.ensemble.as_ref()?;
let live: std::collections::BTreeSet<String> =
ens.peers().iter().map(|h| h.to_hex()).collect();
let table = self.remote_hives.read();
let mut found: Vec<String> = table
.iter()
.filter(|(humd, bees)| {
live.contains(&humd.to_hex())
&& bees.values().any(|m| {
m.bee.iter().any(|b| b == "worker") && m.models.iter().any(|x| x == model)
})
})
.map(|(humd, _)| humd.to_hex())
.collect();
found.sort();
Hid::from_hex(found.first()?).ok()
}

fn snapshot_bees(&self) {
let json = {
let m = self.manifests.read();
Expand Down Expand Up @@ -1006,12 +1050,35 @@ impl ToneSink for HumdSink {
pick
};
let Some(worker_client) = worker_client else {
warn!(sid, model, "prompt.no-worker — no worker bee advertises this model");
if let Some(peer) = self.pick_remote_worker(&model)
&& let Some(ens) = self.ensemble.clone()
{
let mut forward = tone.clone();
if let Some(obj) = forward.as_object_mut() {
obj.insert("to".into(), Value::String(peer.to_hex()));
obj.insert("from".into(), Value::String(ens.me().to_hex()));
}
trace!(sid, model, peer = %peer.short(), "prompt.forward.remote");
match ens.route(forward).await {
Ok(()) => return,
Err(e) => warn!(sid, peer = %peer.short(), err = %e,
"prompt.forward.remote.failed"),
}
}
warn!(sid, model, "prompt.no-worker — no bee on this mesh advertises this model");
let err = serde_json::json!({
"chi": "error",
"sid": sid,
"message": format!("no worker bee advertises model '{}'", model),
});
if let (Some(origin), Some(ens)) = (origin, &self.ensemble) {
let mut out = err.clone();
if let Some(obj) = out.as_object_mut() {
obj.insert("to".into(), Value::String(origin.to_hex()));
obj.insert("from".into(), Value::String(ens.me().to_hex()));
}
let _ = ens.route(out).await;
}
self.thrum.thrum_broadcast(&sid, &self.hive_tag, err);
return;
};
Expand Down Expand Up @@ -1177,6 +1244,14 @@ impl ToneSink for HumdSink {
Some(Chi::PeerAdd) => {
let humd_id = tone.get("humd_id").and_then(Value::as_str).unwrap_or("");
trace!(client_id, humd_id, "ensemble.peer.add");
if let Some(ens) = self.ensemble.clone() {
let known = self.manifests.read().values().cloned().collect::<Vec<_>>();
tokio::spawn(async move {
for manifest in known {
ens.hive_advertise(manifest).await;
}
});
}
}
Some(Chi::PeerRemove) => {
let humd_id = tone.get("humd_id").and_then(Value::as_str).unwrap_or("");
Expand All @@ -1186,7 +1261,11 @@ impl ToneSink for HumdSink {
if bytes.len() == 32 {
let mut id = [0u8; 32];
id.copy_from_slice(&bytes);
ensemble.remove_peer(&ensemble::Hid::from(id));
let gone = ensemble::Hid::from(id);
ensemble.remove_peer(&gone);
if self.remote_hives.write().remove(&gone).is_some() {
trace!(peer = %gone.short(), "discovery.evicted");
}
}
}
}
Expand Down
9 changes: 8 additions & 1 deletion scenarios/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ the noise of test fixtures. The matching test under `sim/` asserts the
same story in Rust, against the in-memory ensemble transport
(`ensemble::InMemoryEndpoint`) wired by [`/root/hum/sim/`](../sim/).

Pairing is strict — one MD, one test:
Pairing is the rule, not yet the state of the directory:

| scenario | test |
|---|---|
Expand All @@ -29,6 +29,13 @@ Pairing is strict — one MD, one test:
| `overflow-inference.md` | `sim/tests/overflow_inference.rs` |
| `partition-and-heal.md` | `sim/tests/partition_and_heal.rs` |
| `eggs-on-the-hum.md` | `sim/tests/eggs_on_the_hum.rs` |
| `find-the-worker.md` | `sim/tests/remote_discovery.rs` |

Tests with no scenario prose yet: `delivery.rs`, `liveness.rs`,
`lossy_link.rs`, `mock_prompt.rs`, `overflow_inference.rs`,
`phone_laptop_roam.rs`, `smoke.rs`, `stalled_peer.rs`,
`tool_catalogue.rs`, `two_humds_ping_pong.rs`. `wifi-p2p-meetup.md`
has prose with no test.

Each MD covers five sections in the same order: **setup**, **happy
path**, **failure modes**, **success criteria**, **what this validates**.
Expand Down
95 changes: 95 additions & 0 deletions scenarios/find-the-worker.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
---
title: "find-the-worker"
description: "the prompt names a model; the mesh decides which humd runs it"
---

# find-the-worker

> _the prompt names a model; the mesh decides which humd runs it_

See `sim/tests/remote_discovery.rs` for the executable form.

## The setup

Trust tier **T3/T4** — federated or open mesh. Two humds, asymmetric
by kind rather than by capacity:

- **humd-L** — the laptop. Has a nestler attached and **no worker at
all**. It cannot answer a single model on its own.
- **humd-S** — the server. Hosts a `claude-cli` worker nest advertising
`models: ["claude-opus-4-7"]`. No nestler is attached; its role is
purely to run work humd-L asks for.

humd-S has already advertised its worker on `hum/hives/announce` when
that bee handshook with its own humd. humd-L subscribed at startup and
holds the manifest under humd-S's Hid.

## The happy path

1. A nestler on humd-L emits `chi:"prompt"` with a `modelId` and
**no `to`**. It is not naming a machine. It cannot; it has never
heard of humd-S.
2. humd-L finds no local worker advertising that model. It consults the
manifests it heard over gossip, keeps only humds still in
`ens.peers()`, and picks one advertising a worker for that model.
Emits a `prompt.forward.remote` trace naming the chosen peer.
3. humd-L sets `to: <humd-S Hid>` and `from: <humd-L Hid>` and routes
the tone. It does not run the session and does not claim the sigil.
4. humd-S sees the prompt with `client_id == "ensemble"`, reads
`from` as its `origin`, and runs its **own** local worker selection —
the same by-model lookup it would run for a local nestler.
5. The worker emits `chi:"chunk"` then `chi:"finish"`. humd-S routes
each one to `sid_origins[sid]` (humd-L) *and* broadcasts locally.
6. humd-L receives each tone with `client_id == "ensemble"`, matches on
`sid`, and broadcasts onto its nestler's stream. The nestler sees
`finish` and the turn closes.

## The failure modes

- **Nobody has the model.** No local bee, no advertised remote. The
nestler gets `chi:"error"` naming the model. If the caller was itself
remote, that error is routed back to its origin rather than dying in
the local session and leaving the origin to time out.
- **The manifest outlived its peer.** Gossip carries no timestamp and a
bee advertises only when it handshakes with its *own* humd, so a
manifest can outlive the peer that made it reachable. Selection is
gated on `ens.peers()`, and `PeerRemove` drops that humd's manifests,
so a dead humd is never chosen.
- **The peer dies between selection and route.** The forward fails and
**falls through** to the error reply. It must not return silently —
that would hang the caller until its own timeout.
- **The peer reconnects.** A reconnect does not re-trigger the bee's
handshake with its own humd, so nothing re-advertises it. `PeerAdd`
re-advertises what we already know, closing the window where a
reachable humd stays invisible.
- **Several humds could serve the model.** Selection is deterministic
(sorted by Hid) rather than capacity-aware. Capacity-aware selection
is `pick_overflow_peer`'s job; two mechanisms, deliberately not
unified here.

## The success criteria

- The nestler's tone carries no `to`, and humd-S still runs the turn.
- Exactly one `prompt.forward.remote` trace, naming the humd chosen.
- `chi:"finish"` reaches the originating nestler on humd-L.
- Every reply carries `to: <humd-L Hid>` — the return path is explicit,
not ambient broadcast.
- Without the discovery feed the same tone yields
`no worker bee advertises model 'claude-opus-4-7'` and no `finish`.

## What this scenario validates

- **Discovery feeding routing.** `hive_advertise` was already called on
every bee handshake and `hive_discover` existed with no production
caller. This is the seam joining them to `Ensemble::route`.
- **The remote table is humd-keyed, not bee-keyed.** A bee has no
ensemble presence of its own; only its humd is dialable, so the Hid
returned by discovery is the routing address and the bee is the value.
- **The return path was already load-bearing.** `from` parsing,
`sid_origins`, and reply routing all predate this. Only the first hop
was missing — which is why the fix is one branch.
- **Staleness discipline without a lease.** Manifests carry no clock, so
the peer set is the honest liveness signal and no TTL can be honest
here.
- **A failed forward still answers.** Degrade to a real error rather
than a hang.
17 changes: 15 additions & 2 deletions sim/tests/liveness.rs
Original file line number Diff line number Diff line change
Expand Up @@ -152,8 +152,21 @@ async fn eviction_only_touches_the_dead_peer() {
sim.probe(a, c).await.expect("probe");
sim.kill_link(b, a).expect("kill b");

await_dead(&sim, a, b).await;
assert_eq!(sim.evict_expired(a, TTL).expect("sweep"), vec![b]);
let deadline = std::time::Instant::now() + Duration::from_secs(5);
let reaped = loop {
sim.probe(a, c).await.expect("c stays live");
let reaped = sim.evict_expired(a, TTL).expect("sweep");
if reaped.contains(&b) {
break reaped;
}
assert!(
std::time::Instant::now() < deadline,
"a peer whose link was killed must be reaped"
);
tokio::time::sleep(Duration::from_millis(5)).await;
};

assert_eq!(reaped, vec![b]);
assert_eq!(liveness(&sim, a, c), Some(Liveness::Live), "c is untouched");
assert_eq!(sim.peer_count(a), 1, "only b was reaped");
}
Expand Down
Loading
Loading