Skip to content
Draft
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
3 changes: 3 additions & 0 deletions crates/mesh-compute/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -57,3 +57,6 @@ pub mod publication;
pub mod catalog;
#[cfg(feature = "mesh")]
mod progress;

#[cfg(feature = "mesh")]
pub mod usage;
21 changes: 21 additions & 0 deletions crates/mesh-compute/src/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ pub enum Phase {
}

struct Slot {
generation: u64,
phase: Phase,
serving: bool,
joining: bool,
Expand All @@ -85,6 +86,7 @@ impl Default for Lifecycle {
Self {
progress: crate::progress::Progress::default(),
slot: Arc::new(Mutex::new(Slot {
generation: 0,
phase: Phase::Stopped,
serving: false,
joining: false,
Expand All @@ -97,6 +99,14 @@ impl Default for Lifecycle {
}
}
impl Lifecycle {
/// Atomically observe the worker epoch and phase for telemetry fencing.
pub fn activity_epoch(&self) -> anyhow::Result<(u64, Phase)> {
self.slot
.lock()
.map(|slot| (slot.generation, slot.phase.clone()))
.map_err(|_| anyhow::anyhow!("Mesh activity unavailable"))
}

pub fn phase(&self) -> Phase {
self.slot.lock().expect("mesh slot poisoned").phase.clone()
}
Expand Down Expand Up @@ -211,6 +221,10 @@ impl Lifecycle {
let (status, mut reads) =
mpsc::channel::<oneshot::Sender<anyhow::Result<mesh_llm_sdk::EmbeddedNodeStatus>>>(1);
let (stop, mut stopping) = watch::channel(false);
slot.generation = slot
.generation
.checked_add(1)
.ok_or_else(|| anyhow::anyhow!("Mesh worker epoch exhausted"))?;
slot.retryable = false;
slot.phase = Phase::Starting;
slot.serving = serving;
Expand Down Expand Up @@ -444,6 +458,7 @@ mod tests {
#[tokio::test]
async fn known_pre_node_failure_can_retry_but_unknown_startup_failure_cannot() {
let owner = Lifecycle::default();
assert_eq!(owner.activity_epoch().unwrap().0, 0);
owner
.launch(async {
Err::<FakeNode, _>(BeforeNode(anyhow::anyhow!("download failed")).into())
Expand All @@ -453,9 +468,15 @@ mod tests {
tokio::task::yield_now().await;
}
assert!(matches!(owner.phase(), Phase::Failed(_)));
assert_eq!(owner.activity_epoch().unwrap().0, 1);
owner.stop_and_wait().await.unwrap();
let (node, stopped, release, _) = fixture();
owner.launch(async { Ok(node) }).unwrap();
assert_eq!(owner.activity_epoch().unwrap().0, 2);
assert!(owner
.launch(async { Err::<FakeNode, _>(anyhow::anyhow!("not started")) })
.is_err());
assert_eq!(owner.activity_epoch().unwrap().0, 2);
owner.stop();
stopped.await.unwrap();
release.send(()).unwrap();
Expand Down
60 changes: 60 additions & 0 deletions crates/mesh-compute/src/usage.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
//! Adapted from Thomas Petersen's crates/community-compute/src/usage.rs at b4a910e7.
//! Allowlisted routing-observed session counters, not proof of local contribution.
use serde::Serialize;
use serde_json::Value;

/// Nullable, independently available counters from the app-owned SDK status.
#[derive(Clone, Debug, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct Usage {
pub tokens_served: Option<u64>,
pub inflight: Option<u64>,
pub tokens_per_second: Option<f64>,
pub peers: Option<usize>,
}
impl Usage {
/// Missing counters remain unknown; raw endpoints and credentials never leave the owner.
pub fn from_payload(payload: &Value) -> Self {
Self {
tokens_served: payload
.pointer("/routing_metrics/completion_tokens_observed")
.and_then(Value::as_u64),
inflight: payload
.pointer("/routing_metrics/local_node/current_inflight_requests")
.or_else(|| payload.get("inflight_requests"))
.and_then(Value::as_u64),
tokens_per_second: payload
.pointer("/routing_metrics/avg_tokens_per_second")
.and_then(Value::as_f64)
.filter(|v| v.is_finite() && *v >= 0.0),
peers: payload.get("peers").and_then(Value::as_array).map(Vec::len),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn projects_local_counters_without_inventing_missing_usage() {
let empty = serde_json::to_value(Usage::from_payload(&json!({}))).unwrap();
assert_eq!(
empty,
json!({"tokensServed":null,"inflight":null,"tokensPerSecond":null,"peers":null})
);
let payload = json!({"routing_metrics":{"completion_tokens_observed":123,"avg_tokens_per_second":-1,"local_node":{"current_inflight_requests":2}},"peers":[{},{}],"private":"not projected"});
let usage = Usage::from_payload(&payload);
assert_eq!(usage.tokens_served, Some(123));
assert_eq!(usage.inflight, Some(2));
assert_eq!(usage.peers, Some(2));
assert_eq!(usage.tokens_per_second, None);
assert!(serde_json::to_value(usage)
.unwrap()
.get("private")
.is_none());
let partial = Usage::from_payload(&json!({"peers":[],"inflight_requests":0}));
assert_eq!(partial.tokens_served, None);
assert_eq!(partial.inflight, Some(0));
assert_eq!(partial.peers, Some(0));
}
}
137 changes: 137 additions & 0 deletions dev/compute-widget.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
import { readFileSync } from "node:fs";
import vm from "node:vm";
import { JSDOM } from "jsdom";
import { expect, it } from "vitest";

it("keeps four designs tied to fresh native usage and resets totals on replacement", async () => {
const html = readFileSync(
new URL("../public/compute-widget.html", import.meta.url),
"utf8",
);
const script = readFileSync(
new URL("../public/compute-widget.js", import.meta.url),
"utf8",
);
const dom = new JSDOM(html, {
runScripts: "outside-only",
url: "http://localhost",
});
try {
let status = {
generation: 1,
state: "running",
usage: { tokensServed: 120, inflight: 0, tokensPerSecond: 5, peers: 2 },
};
const win = dom.window;
win.matchMedia = () => ({ matches: true });
win.HTMLCanvasElement.prototype.getContext = () =>
new Proxy(
{},
{
get: () => () => {},
},
);
win.requestAnimationFrame = () => 0;
win.setTimeout = () => 0;
win.clearTimeout = () => {};
const context = dom.getInternalVMContext();
const run = (code) => vm.runInContext(code, context);
run(script);
win.meshStatus = status;
run("acceptStatus(window.meshStatus)");
expect(run("model.total")).toBe(120);
expect(run("model.phase")).toBe("online");
// Exercise the real moving bee, not only the reduced-motion fallback.
const movingFrames = [0, 350, 700, 1100, 2200, 3800, 6500].map((elapsed) =>
run(
`draw("orbit", {phase:phases[3], index:3, elapsed:${elapsed}}, false)`,
),
);
for (const frame of movingFrames) {
expect(frame.sprite.length).toBeGreaterThan(0);
expect(frame.sprite.flat().every(Number.isFinite)).toBe(true);
expect(frame.streaks.flat().every(Number.isFinite)).toBe(true);
}
expect(movingFrames[3].sprite).not.toEqual(movingFrames[4].sprite);

for (const name of ["orbit", "signal", "original", "bee"]) {
win.dispatchEvent(
new win.KeyboardEvent("keydown", { key: "ArrowRight" }),
);
expect(run("variation")).toBe(name);
const frame = run(
"variation==='original'?drawOriginal({phase:phases[3],index:3,elapsed:1100},true):variation==='bee'?draw('orbit',{phase:phases[3],index:3,elapsed:1100},true):drawInstrument({phase:phases[3],index:3,elapsed:1100},true)",
);
expect(frame.white.every(Number.isFinite)).toBe(true);
run("render(performance.now())");
win.dispatchEvent(new win.KeyboardEvent("keydown", { key: "4" }));
expect(run("inspectedState")).toBeNull();
}
expect(win.document.querySelector("select")).toBeNull();
win.dispatchEvent(new win.KeyboardEvent("keydown", { key: "ArrowLeft" }));
expect(run("variation")).toBe("original");
status = { generation: 1, state: "running", usage: null };
win.meshStatus = status;
run("acceptStatus(window.meshStatus)");
expect(run("model.valid")).toBe(false);
expect(run("model.phase")).toBe("online");
status = { generation: 2, state: "starting", usage: null };
win.meshStatus = status;
run("acceptStatus(window.meshStatus)");
expect(run("model.total")).toBe(null);
expect(run("model.phase")).toBe("link");
status = { generation: 2, state: "failed", usage: null };
win.meshStatus = status;
run("acceptStatus(window.meshStatus)");
expect(run("model.phase")).toBe("error");
const timers = new Map();
win.setTimeout = (callback, delay) => {
timers.set(delay, callback);
return delay;
};
win.clearTimeout = (delay) => timers.delete(delay);
const update = (next, source = win.parent) =>
win.dispatchEvent(
new win.MessageEvent("message", {
source,
data: { type: "buzz-mesh-activity", status: next },
}),
);
update({
generation: 3,
state: "running",
usage: { tokensServed: 7, inflight: 1, peers: 2 },
});
expect(run("model.total")).toBe(7);
expect(run("model.phase")).toBe("active");
update(
{
generation: 3,
state: "running",
usage: { tokensServed: 999, inflight: 1 },
},
{},
);
expect(run("model.total")).toBe(7);
update({
generation: 4,
state: "running",
usage: { tokensServed: null, inflight: null, peers: 1 },
});
expect(run("model.total")).toBeNull();
expect(run("model.peers")).toBe(1);
expect(run("model.phase")).toBe("online");
expect(run("model.batch")).toBeNull();
timers.get(5000)();
expect(run("model.valid")).toBe(false);
expect(run("model.total")).toBeNull();
expect(run("model.phase")).toBe("offline");
update({ generation: 4, state: "running", usage: null });
expect(run("model.phase")).toBe("online");
expect(run("model.total")).toBeNull();
expect(script).not.toContain("community_compute_status");
expect(script).not.toContain("start_dragging");
} finally {
dom.window.close();
}
});
15 changes: 15 additions & 0 deletions docs/agent-control.md
Original file line number Diff line number Diff line change
Expand Up @@ -1085,3 +1085,18 @@ native artifact bytes. Stronger artifact policy (pinned hashes, bundling or
upstream signatures) is a separate follow-up, not part of this port. Optional failed
old-community retirement is best-effort; routing advertisements expire after
120 seconds. No live packaged acceptance is claimed by fixture tests.

### Shared-compute activity tile

Shared compute embeds Thomas Petersen's original four-design canvas from
`feat/community-compute-plugin` (`b4a910e7`). Its original history and notices are
preserved; see [the integration record](agent-tile-integration.md). Focus the tile
and use Left/Right to switch Bee, Orbit, Signal and LED. It is read-only: sharing
consent, model selection and native lifetime remain with the existing host.

The tile observes this app's worker epoch, lifecycle and allowlisted SDK counters.
Unknown counters stay unknown; stale samples clear. Token totals describe
routing-observed completions, not compute contributed by this machine, and peers
are known nodes rather than proven serving members. Request activity is not proof
of local model decoding. The tile does not access native IPC or register a separate
runtime/window. Numeric synthetic previews are disabled in the live tile.
Loading
Loading