Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
76e4034e01 | ||
|
|
6e098fad29 | ||
|
|
0e2d26304e | ||
|
|
9585edef9b | ||
|
|
bcb3f580ff | ||
|
|
d890509be3 | ||
|
|
e772182724 |
8 changed files with 323 additions and 76 deletions
43
CLAUDE.md
43
CLAUDE.md
|
|
@ -45,9 +45,19 @@ hive-c0re/ host daemon + CLI (one binary, subcommand-dispatched)
|
||||||
src/crash_watch.rs poll every 10s; fire HelperEvent::ContainerCrash
|
src/crash_watch.rs poll every 10s; fire HelperEvent::ContainerCrash
|
||||||
when a previously-running container disappears
|
when a previously-running container disappears
|
||||||
without an operator-initiated transient
|
without an operator-initiated transient
|
||||||
|
src/container_view.rs ContainerView struct + build_all helper;
|
||||||
|
shared between dashboard.rs (cold-load via
|
||||||
|
/api/state) and coordinator.rs's
|
||||||
|
rescan_containers_and_emit
|
||||||
src/coordinator.rs shared state (broker/approvals/operator_questions/
|
src/coordinator.rs shared state (broker/approvals/operator_questions/
|
||||||
transient/sockets) + tombstone enumeration +
|
transient/sockets) + tombstone enumeration +
|
||||||
kick_agent + notify_agent (helper-event push)
|
kick_agent + notify_agent (helper-event push) +
|
||||||
|
last_containers cache + rescan_and_emit diff helper
|
||||||
|
src/open_threads.rs loose-ends aggregator (pending approvals +
|
||||||
|
unanswered questions) — for_agent (filtered) and
|
||||||
|
hive_wide (manager surface). Backs
|
||||||
|
AgentRequest::GetOpenThreads + ManagerRequest::
|
||||||
|
GetOpenThreads (the get_open_threads MCP tool).
|
||||||
src/actions.rs approve/deny/destroy (transient-aware)
|
src/actions.rs approve/deny/destroy (transient-aware)
|
||||||
src/auto_update.rs startup rebuild scan + ensure_manager +
|
src/auto_update.rs startup rebuild scan + ensure_manager +
|
||||||
meta::lock_update_hyperhive bump
|
meta::lock_update_hyperhive bump
|
||||||
|
|
@ -85,6 +95,9 @@ hive-ag3nt/ in-container harness crate; produces TWO binaries
|
||||||
src/client.rs generic JSON-line request/response over unix socket
|
src/client.rs generic JSON-line request/response over unix socket
|
||||||
src/web_ui.rs per-container axum HTTP page (incl /api/cancel,
|
src/web_ui.rs per-container axum HTTP page (incl /api/cancel,
|
||||||
/api/compact, /api/model, /events/history)
|
/api/compact, /api/model, /events/history)
|
||||||
|
src/turn_stats.rs per-turn analytics sink (one sqlite row per
|
||||||
|
turn at /state/hyperhive-turn-stats.sqlite);
|
||||||
|
schema + best-effort writer
|
||||||
src/events.rs LiveEvent + broadcast Bus + sqlite-backed history
|
src/events.rs LiveEvent + broadcast Bus + sqlite-backed history
|
||||||
(/state/hyperhive-events.sqlite) + TurnState +
|
(/state/hyperhive-events.sqlite) + TurnState +
|
||||||
model selection (persisted at /state/hyperhive-model)
|
model selection (persisted at /state/hyperhive-model)
|
||||||
|
|
@ -193,6 +206,34 @@ Prune freely.
|
||||||
domain tooling — the agent flake's `inputs` block pulls
|
domain tooling — the agent flake's `inputs` block pulls
|
||||||
the external flake, `agent.nix` references it via
|
the external flake, `agent.nix` references it via
|
||||||
`flakeInputs.<name>.packages.${pkgs.system}.default`.
|
`flakeInputs.<name>.packages.${pkgs.system}.default`.
|
||||||
|
- **Just landed:** per-turn analytics sink. New
|
||||||
|
`hive-ag3nt::turn_stats` writes one row per claude turn to
|
||||||
|
`/state/hyperhive-turn-stats.sqlite`: identity (model,
|
||||||
|
wake_from, result_kind), timing (started/ended_at,
|
||||||
|
duration_ms), cost (full token-usage breakdown), behaviour
|
||||||
|
(tool_call_count + per-tool JSON map), and post-turn snapshot
|
||||||
|
metrics (open_threads_count, open_reminders_count fetched via
|
||||||
|
the existing GetOpenThreads + new CountPendingReminders RPC).
|
||||||
|
Both ag3nt + m1nd bin loops capture, both Bus accumulates
|
||||||
|
tool_use blocks via observe_stream during the stdout pump.
|
||||||
|
Writes are best-effort. No host-side vacuum yet — TODO under
|
||||||
|
Telemetry; same shape as events_vacuum, target 90d retention.
|
||||||
|
- **Just landed:** agent web UI event-driven badges. New
|
||||||
|
`LiveEvent::StatusChanged / ModelChanged / TokenUsageChanged
|
||||||
|
/ TurnStateChanged` variants replace the per-agent page's
|
||||||
|
/api/state polling for the state row. Status/model/token/state
|
||||||
|
badges all update from SSE; /api/state only fetched on cold
|
||||||
|
load + during the login flow (session output isn't event-
|
||||||
|
shaped). Per-agent endpoints (`/api/cancel|compact|model|
|
||||||
|
new-session`, `/login/*`) all flip 303→200. New `alive-badge`
|
||||||
|
chip carries the harness reachability signal (replaces the
|
||||||
|
"● harness alive" paragraph); new `ctx-badge` mirrors Claude
|
||||||
|
Code's bottom-right "N tokens" indicator. Every chip carries
|
||||||
|
a `title=...` tooltip for hover detail.
|
||||||
|
- **Just landed:** events_vacuum simplified to age-only —
|
||||||
|
`KEEP_SECS = 7d`, no row cap. Chatty turn no longer evicts
|
||||||
|
a quiet day's history sooner than expected. Hourly sweep
|
||||||
|
unchanged.
|
||||||
- **Just landed:** Phase 6 container events. New
|
- **Just landed:** Phase 6 container events. New
|
||||||
`DashboardEvent::ContainerStateChanged { container }` +
|
`DashboardEvent::ContainerStateChanged { container }` +
|
||||||
`ContainerRemoved { name }` close the last refetch loop on the
|
`ContainerRemoved { name }` close the last refetch loop on the
|
||||||
|
|
|
||||||
1
TODO.md
1
TODO.md
|
|
@ -5,6 +5,7 @@
|
||||||
- Shared space for all agents to access documents/files without manager routing
|
- Shared space for all agents to access documents/files without manager routing
|
||||||
- Private git forge agents can push to and create new repos in
|
- Private git forge agents can push to and create new repos in
|
||||||
- Move bind mounts in agents to `/agents/<name>/state` so path for agent = path for manager
|
- Move bind mounts in agents to `/agents/<name>/state` so path for agent = path for manager
|
||||||
|
- **Split harness-internal state from agent-visible state**: the `/agents/<n>/state/` mount (host `/var/lib/hyperhive/agents/<n>/state/`) currently mixes the agent's durable notes with harness internals — `hyperhive-events.sqlite`, `hyperhive-turn-stats.sqlite`, `hyperhive-model`, future per-agent skill caches, etc. The agent can accidentally overwrite a harness file, the harness clutters what claude thinks is "my notes dir", and the host-side vacuum has to special-case filenames it owns. Move harness internals to a sibling dir, e.g. `/var/lib/hyperhive/agents/<n>/harness/`, bind-mounted RW into the container as `/agents/<n>/harness/` (same path inside + out, same convention as state). Container's `/agents/<n>/state/` becomes purely agent-owned. Touches: `paths.rs` (new `harness_dir()`), `events.rs`, `turn_stats.rs` (default paths flip), `events_vacuum.rs` (sweep root flips), `lifecycle.rs` (extra bind mount), and a migration that moves existing files on first boot under the new layout. Side benefit: makes the privsep TODO cheaper — the unprivileged web server only needs read access to `/agents/<n>/state/` (operator-meaningful files), not `/agents/<n>/harness/`. The legacy bare `/state` mount the manager still uses (`container_state_prefix("manager") == "/state/"`, manager bind in `lifecycle::set_nspawn_flags`) gets removed in the same pass — manager goes to `/agents/manager/state/` + `/agents/manager/harness/` like every other agent.
|
||||||
- **Broadcast messaging**: allow sending messages with recipient "*" to all agents; deliver with hint "this was a broadcast and may not need any action from you"
|
- **Broadcast messaging**: allow sending messages with recipient "*" to all agents; deliver with hint "this was a broadcast and may not need any action from you"
|
||||||
- **Multi-agent restart coordination**: when rebuilding all agents, manager should start first so it can coordinate post-restart confusion (notify agents, suppress unnecessary retries, etc)
|
- **Multi-agent restart coordination**: when rebuilding all agents, manager should start first so it can coordinate post-restart confusion (notify agents, suppress unnecessary retries, etc)
|
||||||
- **Shared docs/skills repo (RO)**: a single repo on the hive forge that every agent has read-only access to — common references, prompts, runbooks, "skills" the operator wants every agent to inherit without baking into the system prompt or `/shared`. Implementation likely: seed an `org-shared/docs` repo on first hive-forge boot, grant every per-agent user a read membership in the org. Agents `git clone` it (or use the API) to read; only the manager + operator can push.
|
- **Shared docs/skills repo (RO)**: a single repo on the hive forge that every agent has read-only access to — common references, prompts, runbooks, "skills" the operator wants every agent to inherit without baking into the system prompt or `/shared`. Implementation likely: seed an `org-shared/docs` repo on first hive-forge boot, grant every per-agent user a read membership in the org. Agents `git clone` it (or use the API) to read; only the manager + operator can push.
|
||||||
|
|
|
||||||
|
|
@ -41,17 +41,36 @@ One table:
|
||||||
harness emits during turn loop execution.
|
harness emits during turn loop execution.
|
||||||
|
|
||||||
The harness writes; the host vacuums. `hive-c0re::events_vacuum`
|
The harness writes; the host vacuums. `hive-c0re::events_vacuum`
|
||||||
runs hourly and sweeps every existing agent state dir, applying the
|
runs hourly and sweeps every existing agent state dir, deleting
|
||||||
same two-stage delete to each file: drop rows older than 7 days,
|
rows older than 7 days. Age-only — no row cap — so a chatty turn
|
||||||
then trim to the 2000 most-recent. Centralising retention on the
|
doesn't lose history sooner than a quiet one; disk pressure on a
|
||||||
host means a misbehaving harness can't disable its own vacuum and
|
sustained burst is the cheaper problem to have. Centralising
|
||||||
agents don't need any cleanup wiring of their own.
|
retention on the host means a misbehaving harness can't disable
|
||||||
|
its own vacuum and agents don't need any cleanup wiring of their
|
||||||
|
own.
|
||||||
|
|
||||||
Path overridable via `HYPERHIVE_EVENTS_DB` (for dev / no-`/state`
|
Path overridable via `HYPERHIVE_EVENTS_DB` (for dev / no-`/state`
|
||||||
setups). On open failure the `Bus` falls back to no-store mode
|
setups). On open failure the `Bus` falls back to no-store mode
|
||||||
rather than crashing the harness — events still broadcast over SSE,
|
rather than crashing the harness — events still broadcast over SSE,
|
||||||
just nothing persisted.
|
just nothing persisted.
|
||||||
|
|
||||||
|
### `/state/hyperhive-turn-stats.sqlite` (per agent)
|
||||||
|
|
||||||
|
Per-turn analytics sink. One row per claude turn captures
|
||||||
|
identity (`model`, `wake_from`, `result_kind`), timing
|
||||||
|
(`started_at`, `ended_at`, `duration_ms`), cost (input / output /
|
||||||
|
cache_read / cache_creation token counts), behaviour
|
||||||
|
(`tool_call_count` + `tool_call_breakdown_json`), and post-turn
|
||||||
|
snapshot metrics (`open_threads_count`,
|
||||||
|
`open_reminders_count` — fetched via the same socket the harness
|
||||||
|
already uses for `GetOpenThreads` + `CountPendingReminders`).
|
||||||
|
Bin-loop helpers `build_row` + `record` land each row at
|
||||||
|
`turn_end`; writes are best-effort, a sqlite hiccup logs + lets
|
||||||
|
the turn loop continue.
|
||||||
|
|
||||||
|
No host-side vacuum yet — tracked in `TODO.md` under Telemetry
|
||||||
|
(target retention ~90 days, age-only sweep like events_vacuum).
|
||||||
|
|
||||||
### `/state/hyperhive-model` (per agent)
|
### `/state/hyperhive-model` (per agent)
|
||||||
|
|
||||||
Single-line text file holding the claude model name currently
|
Single-line text file holding the claude model name currently
|
||||||
|
|
@ -68,8 +87,10 @@ Under `/var/lib/hyperhive/agents/<name>/`:
|
||||||
- `config/` — the proposed nix repo (manager-editable).
|
- `config/` — the proposed nix repo (manager-editable).
|
||||||
- `claude/` — claude OAuth credentials, bind-mounted RW to
|
- `claude/` — claude OAuth credentials, bind-mounted RW to
|
||||||
`/root/.claude` inside the container.
|
`/root/.claude` inside the container.
|
||||||
- `state/` — durable notes + the events.sqlite db, bind-mounted
|
- `state/` — durable notes, the events.sqlite db, and the
|
||||||
to `/state` inside the container.
|
turn-stats sqlite db. Bind-mounted to `/agents/<name>/state`
|
||||||
|
inside the container (the manager still uses the legacy
|
||||||
|
`/state` mount point — same host path either way).
|
||||||
|
|
||||||
Under `/var/lib/hyperhive/applied/<name>/` — the hive-c0re-only
|
Under `/var/lib/hyperhive/applied/<name>/` — the hive-c0re-only
|
||||||
applied repo. Tracks `flake.nix` (module-only boilerplate; never
|
applied repo. Tracks `flake.nix` (module-only boilerplate; never
|
||||||
|
|
|
||||||
130
docs/web-ui.md
130
docs/web-ui.md
|
|
@ -201,6 +201,29 @@ not ours.
|
||||||
a managed container.
|
a managed container.
|
||||||
- `GET /api/agent-config/{name}` — read-only view of the applied
|
- `GET /api/agent-config/{name}` — read-only view of the applied
|
||||||
`agent.nix`.
|
`agent.nix`.
|
||||||
|
- `GET /api/state-file?path=<host-or-container-path>` — bounded
|
||||||
|
text read of a file under the per-agent `state/` subtree or
|
||||||
|
the shared `/var/lib/hyperhive/shared/`. Accepts the
|
||||||
|
container-view forms (`/agents/<n>/state/...`, `/shared/...`)
|
||||||
|
and the host form. Canonicalises + verifies the path stays
|
||||||
|
inside the allow-list, refuses anything but a regular file,
|
||||||
|
refuses `/agents/<n>/claude` / `config` subtrees, truncates
|
||||||
|
bodies at 1 MiB. Click-time backing for the inline path-link
|
||||||
|
preview.
|
||||||
|
|
||||||
|
Detection of which tokens *are* path links is done
|
||||||
|
**server-side at broker-message ingest**, not client-side:
|
||||||
|
the broker forwarder calls `scan_validated_paths(body)` —
|
||||||
|
same allow-list helper the read endpoint uses — and attaches
|
||||||
|
the verified file tokens to the event as `file_refs: Vec<String>`.
|
||||||
|
The client trusts that list and linkifies only those tokens,
|
||||||
|
so directories, missing files, and forbidden subtrees never
|
||||||
|
become anchors. No probe endpoint, no client-side regex
|
||||||
|
heuristics. Historical messages get the same treatment on
|
||||||
|
`/dashboard/history` backfill.
|
||||||
|
- `GET /api/reminders` — list pending reminders for the
|
||||||
|
dashboard's queued-reminders panel.
|
||||||
|
- `POST /cancel-reminder/{id}` — hard-delete a pending reminder.
|
||||||
- `GET /dashboard/stream` — unified live event channel:
|
- `GET /dashboard/stream` — unified live event channel:
|
||||||
broker `sent` / `delivered`, plus the mutation events listed
|
broker `sent` / `delivered`, plus the mutation events listed
|
||||||
below. Each frame carries `seq`.
|
below. Each frame carries `seq`.
|
||||||
|
|
@ -223,21 +246,37 @@ payload):
|
||||||
queue + history mutations. Client mutates a derived store and
|
queue + history mutations. Client mutates a derived store and
|
||||||
re-renders only the approvals section.
|
re-renders only the approvals section.
|
||||||
- `question_added` (id, asker, question, options, multi,
|
- `question_added` (id, asker, question, options, multi,
|
||||||
asked_at, deadline_at) / `question_resolved` (id, answer,
|
asked_at, deadline_at, target) / `question_resolved` (id,
|
||||||
answerer, answered_at, cancelled) — operator-targeted
|
answer, answerer, answered_at, cancelled, target) — both
|
||||||
questions only (peer-to-peer questions never fire these). The
|
operator-targeted and peer (agent-to-agent) threads fire
|
||||||
ttl watchdog fires `question_resolved` with
|
these. The dashboard's questions pane surfaces both, with
|
||||||
`answerer = "ttl-watchdog"` on expiry.
|
filter chips (all / @operator / @peer / per-participant) and
|
||||||
|
an `0V3RR1D3` button on peer rows so the operator can
|
||||||
|
answer when an agent is stuck. The ttl watchdog fires
|
||||||
|
`question_resolved` with `answerer = "ttl-watchdog"` on
|
||||||
|
expiry.
|
||||||
- `transient_set` (name, transient_kind, since_unix) /
|
- `transient_set` (name, transient_kind, since_unix) /
|
||||||
`transient_cleared` (name) — lifecycle action spinners. The
|
`transient_cleared` (name) — lifecycle action spinners. The
|
||||||
client ticks the elapsed-seconds badge off `since_unix`
|
client ticks the elapsed-seconds badge off `since_unix`
|
||||||
client-side, no polling.
|
client-side, no polling.
|
||||||
|
- `container_state_changed` (container: ContainerView) /
|
||||||
|
`container_removed` (name) — per-row container mutations,
|
||||||
|
emitted by `Coordinator::rescan_containers_and_emit` from
|
||||||
|
every mutation site (`actions::approve` post-spawn,
|
||||||
|
`actions::destroy`, the lifecycle_action wrapper,
|
||||||
|
`auto_update::rebuild_agent`) and from the 10s
|
||||||
|
`crash_watch` poll. Client upserts/removes by name; the
|
||||||
|
pending overlay is read from `transientsState` since the
|
||||||
|
payload doesn't carry it.
|
||||||
|
|
||||||
`/api/state` still serves `approvals` / `approval_history` /
|
`/api/state` is **only fetched on cold-load and on the few
|
||||||
`questions` / `question_history` / `transients` for cold-start
|
forms that mutate non-event-derived state** (PURG3 +
|
||||||
on first page load and as a safety-net resync from the 5s poll;
|
meta-update, since tombstones + meta_inputs aren't event-
|
||||||
the client maintains the same arrays in derived stores and
|
shaped yet). Every other section — approvals, questions,
|
||||||
applies the events on top.
|
transients, containers, operator inbox, message flow —
|
||||||
|
derives from `/dashboard/stream` after the initial snapshot,
|
||||||
|
maintaining its own client-side store and applying events on
|
||||||
|
top. The 5s periodic poll is gone.
|
||||||
|
|
||||||
Generalised form helpers: `form[data-confirm="…"]` pops
|
Generalised form helpers: `form[data-confirm="…"]` pops
|
||||||
`confirm()` before submit; `form[data-prompt="…"]` pops
|
`confirm()` before submit; `form[data-prompt="…"]` pops
|
||||||
|
|
@ -250,16 +289,34 @@ Layout, top to bottom:
|
||||||
|
|
||||||
- Banner (gradient shimmer while state=thinking).
|
- Banner (gradient shimmer while state=thinking).
|
||||||
- Title with `↑ DASHB04RD` back-link (new tab) + `↻ R3BU1LD`.
|
- Title with `↑ DASHB04RD` back-link (new tab) + `↻ R3BU1LD`.
|
||||||
- Status section (online / needs login / login-in-progress).
|
- Status section: empty when online (alive-badge in the state
|
||||||
- **State row**: state badge + model chip + last-turn timing +
|
row carries the signal), populated with the login form /
|
||||||
cancel-turn button + new-session button.
|
OAuth URL when `status` is `needs_login_*`.
|
||||||
|
- **State row**: alive badge + state badge + model chip + ctx
|
||||||
|
badge + last-turn timing + cancel-turn button + new-session
|
||||||
|
button. Every chip carries a `title=...` tooltip with the
|
||||||
|
detailed breakdown.
|
||||||
|
- Alive badge: `● alive` (green) / `◌ needs login` (amber) /
|
||||||
|
`◌ logging in` / `○ offline` / `… connecting`. Driven by
|
||||||
|
`LiveEvent::StatusChanged`; replaces the old "harness alive
|
||||||
|
— turn loop running" paragraph so the state row carries
|
||||||
|
every reachability signal.
|
||||||
- State badge: `💤 idle` / `🧠 thinking` / `📦 compacting` /
|
- State badge: `💤 idle` / `🧠 thinking` / `📦 compacting` /
|
||||||
`○ offline` / `… booting`, with an age suffix (`12s`,
|
`○ offline` / `… booting`, with an age suffix (`12s`,
|
||||||
`2m 14s`). Driven from `/api/state.turn_state` +
|
`2m 14s`). Driven by `LiveEvent::TurnStateChanged`
|
||||||
`turn_state_since`; SSE turn_start/turn_end still flip it
|
(`{state, since_unix}`) — the bus emits on every
|
||||||
instantly between polls. Authoritative source is the
|
`Bus::set_state` so the badge updates without a /api/state
|
||||||
harness's `Bus::state_snapshot()`.
|
refetch. Cold-load via `/api/state.turn_state` +
|
||||||
- Model chip: `model · <name>` (e.g. `model · haiku`).
|
`turn_state_since`.
|
||||||
|
- Model chip: `model · <name>` (e.g. `model · haiku`). Driven
|
||||||
|
by `LiveEvent::ModelChanged`; emitted from `Bus::set_model`.
|
||||||
|
- Ctx badge: `ctx · 142k` — total prompt tokens in the
|
||||||
|
current context window (input + cache_read + cache_write),
|
||||||
|
mirroring claude code's bottom-right indicator. Hover for
|
||||||
|
the breakdown including output. Driven by
|
||||||
|
`LiveEvent::TokenUsageChanged`; emitted from
|
||||||
|
`Bus::record_usage` whenever the terminal `result` event
|
||||||
|
delivers a fresh usage block.
|
||||||
- Last-turn chip: `last turn 12.3s` appears after the first
|
- Last-turn chip: `last turn 12.3s` appears after the first
|
||||||
turn ends, computed from the state-since deltas.
|
turn ends, computed from the state-since deltas.
|
||||||
- `■ cancel turn` button: visible only while state=thinking,
|
- `■ cancel turn` button: visible only while state=thinking,
|
||||||
|
|
@ -269,6 +326,11 @@ Layout, top to bottom:
|
||||||
arm a one-shot Bus flag — the next turn drops
|
arm a one-shot Bus flag — the next turn drops
|
||||||
`--continue`, starting a fresh claude session. Subsequent
|
`--continue`, starting a fresh claude session. Subsequent
|
||||||
turns resume normal `--continue`.
|
turns resume normal `--continue`.
|
||||||
|
|
||||||
|
Polling: `/api/state` is fetched **once** on cold load, and
|
||||||
|
again while `status === 'needs_login_in_progress'` (login
|
||||||
|
session output isn't event-shaped yet). Every other badge
|
||||||
|
updates from SSE; no periodic refresh timer runs.
|
||||||
- Inbox `<details>` block (collapsed): `inbox · N` — last 30
|
- Inbox `<details>` block (collapsed): `inbox · N` — last 30
|
||||||
messages addressed to this agent, fetched via
|
messages addressed to this agent, fetched via
|
||||||
`AgentRequest::Recent { limit: 30 }`. (Separate from
|
`AgentRequest::Recent { limit: 30 }`. (Separate from
|
||||||
|
|
@ -345,14 +407,38 @@ Unknown `/foo` shows an error row instead of being silently sent.
|
||||||
|
|
||||||
### Per-agent endpoints
|
### Per-agent endpoints
|
||||||
|
|
||||||
|
All POSTs return 200 (no 303 redirects). The matching mutations
|
||||||
|
fire `LiveEvent` variants on the per-agent bus, so the client
|
||||||
|
doesn't refetch `/api/state` on submit — the SSE stream
|
||||||
|
delivers the new state faster anyway. Only the login flow still
|
||||||
|
polls (session output streams in updates that aren't event-
|
||||||
|
shaped).
|
||||||
|
|
||||||
- `POST /send` — operator-injected message into this agent's inbox.
|
- `POST /send` — operator-injected message into this agent's inbox.
|
||||||
- `POST /login/{start,code,cancel}` — claude OAuth login flow.
|
- `POST /login/{start,code,cancel}` — claude OAuth login flow.
|
||||||
- `POST /api/cancel` — SIGINT the in-flight claude turn.
|
Start/cancel emit `LiveEvent::StatusChanged` to flip the
|
||||||
|
badge to/from `needs_login_in_progress`.
|
||||||
|
- `POST /api/cancel` — SIGINT the in-flight claude turn. Emits a
|
||||||
|
`LiveEvent::Note`.
|
||||||
- `POST /api/compact` — run `/compact` on the persistent session
|
- `POST /api/compact` — run `/compact` on the persistent session
|
||||||
(same MCP config + system prompt + allowed tools as a normal
|
(same MCP config + system prompt + allowed tools as a normal
|
||||||
turn — only the stdin payload differs).
|
turn — only the stdin payload differs). Flips state to
|
||||||
|
`Compacting` via `Bus::set_state`, which emits
|
||||||
|
`TurnStateChanged`.
|
||||||
- `POST /api/model` (`model=<name>`) — switch the model for
|
- `POST /api/model` (`model=<name>`) — switch the model for
|
||||||
future turns.
|
future turns. `Bus::set_model` emits `ModelChanged`.
|
||||||
- `POST /api/new-session` — arm a one-shot for the next turn to
|
- `POST /api/new-session` — arm a one-shot for the next turn to
|
||||||
drop `--continue`.
|
drop `--continue`. Emits a `LiveEvent::Note`.
|
||||||
- `GET /events/history` — replay buffer for the terminal.
|
- `GET /events/history` — replay buffer for the terminal.
|
||||||
|
|
||||||
|
Bus events (new vocabulary on `/events/stream`):
|
||||||
|
|
||||||
|
- `status_changed { status }` — `online` /
|
||||||
|
`needs_login_idle` / `needs_login_in_progress`. Drives the
|
||||||
|
alive-badge.
|
||||||
|
- `model_changed { model }` — drives the model chip.
|
||||||
|
- `token_usage_changed { usage: TokenUsage }` — drives the
|
||||||
|
ctx-badge. Emitted from `Bus::record_usage` whenever the
|
||||||
|
stream-json `result` event delivers a fresh usage block.
|
||||||
|
- `turn_state_changed { state, since_unix }` — drives the
|
||||||
|
state badge (`idle`/`thinking`/`compacting`).
|
||||||
|
|
|
||||||
|
|
@ -48,7 +48,6 @@
|
||||||
// perspective (we'd need to know which agent the message is about
|
// perspective (we'd need to know which agent the message is about
|
||||||
// to translate it). Prefer `/agents/<name>/state/...` in agent
|
// to translate it). Prefer `/agents/<name>/state/...` in agent
|
||||||
// outputs and the link will resolve.
|
// outputs and the link will resolve.
|
||||||
const PATH_RE = /(\/var\/lib\/hyperhive\/agents\/[\w.-]+\/state\/[\w./-]+|\/var\/lib\/hyperhive\/shared\/[\w./-]+|\/agents\/[\w.-]+\/state\/[\w./-]+|\/shared\/[\w./-]+)/g;
|
|
||||||
async function fetchStateFile(path) {
|
async function fetchStateFile(path) {
|
||||||
const resp = await fetch('/api/state-file?path=' + encodeURIComponent(path));
|
const resp = await fetch('/api/state-file?path=' + encodeURIComponent(path));
|
||||||
const text = await resp.text();
|
const text = await resp.text();
|
||||||
|
|
@ -85,27 +84,54 @@
|
||||||
});
|
});
|
||||||
return { anchor, details };
|
return { anchor, details };
|
||||||
}
|
}
|
||||||
// Append `text` to `parent` as a mix of text nodes + path anchors.
|
// Append `text` to `parent` as a mix of text nodes + path
|
||||||
// Returns the array of generated `<details>` previews so the
|
// anchors. `refs` is the server-attached `file_refs` array
|
||||||
// caller can append them as block siblings under the row.
|
// (verified-file tokens that appear in `text`); each occurrence
|
||||||
function appendLinkified(parent, text) {
|
// of a ref in `text` becomes a clickable anchor + a sibling
|
||||||
|
// <details> preview that lazy-fetches from /api/state-file.
|
||||||
|
// Anything not in `refs` stays plain text. No client-side
|
||||||
|
// regex, no probe endpoint — the server saw the body first
|
||||||
|
// and made the call. When `refs` is empty/missing we just
|
||||||
|
// emit plain text.
|
||||||
|
function appendLinkified(parent, text, refs) {
|
||||||
const previews = [];
|
const previews = [];
|
||||||
if (text == null) return previews;
|
if (text == null) return previews;
|
||||||
const str = String(text);
|
const str = String(text);
|
||||||
let lastIdx = 0;
|
const tokens = (refs || []).slice();
|
||||||
PATH_RE.lastIndex = 0;
|
if (!tokens.length) {
|
||||||
let m;
|
if (str) parent.appendChild(document.createTextNode(str));
|
||||||
while ((m = PATH_RE.exec(str)) !== null) {
|
return previews;
|
||||||
if (m.index > lastIdx) {
|
}
|
||||||
parent.appendChild(document.createTextNode(str.slice(lastIdx, m.index)));
|
// Walk the string left-to-right, at each step looking for the
|
||||||
|
// next occurrence of any token. Longest-first tie-break so a
|
||||||
|
// ref like `/agents/foo/state/x.md` wins over a (hypothetical)
|
||||||
|
// shorter token that prefixes it. O(text * refs) worst case;
|
||||||
|
// refs is bounded server-side to whatever fits in a body, so
|
||||||
|
// this stays cheap.
|
||||||
|
tokens.sort((a, b) => b.length - a.length);
|
||||||
|
let i = 0;
|
||||||
|
while (i < str.length) {
|
||||||
|
let bestStart = -1;
|
||||||
|
let bestToken = null;
|
||||||
|
for (const t of tokens) {
|
||||||
|
const idx = str.indexOf(t, i);
|
||||||
|
if (idx === -1) continue;
|
||||||
|
if (bestStart === -1 || idx < bestStart || (idx === bestStart && t.length > bestToken.length)) {
|
||||||
|
bestStart = idx;
|
||||||
|
bestToken = t;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
const { anchor, details } = makePathPreview(m[0]);
|
if (bestStart === -1) {
|
||||||
|
parent.appendChild(document.createTextNode(str.slice(i)));
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
if (bestStart > i) {
|
||||||
|
parent.appendChild(document.createTextNode(str.slice(i, bestStart)));
|
||||||
|
}
|
||||||
|
const { anchor, details } = makePathPreview(bestToken);
|
||||||
parent.appendChild(anchor);
|
parent.appendChild(anchor);
|
||||||
previews.push(details);
|
previews.push(details);
|
||||||
lastIdx = m.index + m[0].length;
|
i = bestStart + bestToken.length;
|
||||||
}
|
|
||||||
if (lastIdx < str.length) {
|
|
||||||
parent.appendChild(document.createTextNode(str.slice(lastIdx)));
|
|
||||||
}
|
}
|
||||||
return previews;
|
return previews;
|
||||||
}
|
}
|
||||||
|
|
@ -917,7 +943,12 @@
|
||||||
const operatorInbox = [];
|
const operatorInbox = [];
|
||||||
function inboxAppendFromEvent(ev) {
|
function inboxAppendFromEvent(ev) {
|
||||||
if (ev.kind !== 'sent' || ev.to !== 'operator') return false;
|
if (ev.kind !== 'sent' || ev.to !== 'operator') return false;
|
||||||
operatorInbox.unshift({ from: ev.from, body: ev.body, at: ev.at });
|
operatorInbox.unshift({
|
||||||
|
from: ev.from,
|
||||||
|
body: ev.body,
|
||||||
|
at: ev.at,
|
||||||
|
file_refs: ev.file_refs || [],
|
||||||
|
});
|
||||||
if (operatorInbox.length > INBOX_LIMIT) operatorInbox.length = INBOX_LIMIT;
|
if (operatorInbox.length > INBOX_LIMIT) operatorInbox.length = INBOX_LIMIT;
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
@ -934,7 +965,7 @@
|
||||||
for (const m of operatorInbox) {
|
for (const m of operatorInbox) {
|
||||||
const li = el('li');
|
const li = el('li');
|
||||||
const body = el('span', { class: 'msg-body' });
|
const body = el('span', { class: 'msg-body' });
|
||||||
const previews = appendLinkified(body, m.body);
|
const previews = appendLinkified(body, m.body, m.file_refs);
|
||||||
li.append(
|
li.append(
|
||||||
el('span', { class: 'msg-ts' }, fmt(m.at)), ' ',
|
el('span', { class: 'msg-ts' }, fmt(m.at)), ' ',
|
||||||
el('span', { class: 'msg-from' }, m.from), ' ',
|
el('span', { class: 'msg-from' }, m.from), ' ',
|
||||||
|
|
@ -1460,7 +1491,7 @@
|
||||||
to.className = 'msg-to'; to.textContent = ev.to;
|
to.className = 'msg-to'; to.textContent = ev.to;
|
||||||
const body = document.createElement('span');
|
const body = document.createElement('span');
|
||||||
body.className = 'msg-body';
|
body.className = 'msg-body';
|
||||||
const previews = appendLinkified(body, ev.body);
|
const previews = appendLinkified(body, ev.body, ev.file_refs);
|
||||||
row.append(ts, ' ', arrow, ' ', from, ' ', sep, ' ', to, ' ', body);
|
row.append(ts, ' ', arrow, ' ', from, ' ', sep, ' ', to, ' ', body);
|
||||||
for (const d of previews) row.appendChild(d);
|
for (const d of previews) row.appendChild(d);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -634,21 +634,25 @@ async fn dashboard_history(State(state): State<AppState>) -> Response {
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.map(|m| match m {
|
.map(|m| match m {
|
||||||
crate::broker::MessageEvent::Sent { from, to, body, at } => {
|
crate::broker::MessageEvent::Sent { from, to, body, at } => {
|
||||||
|
let file_refs = scan_validated_paths(&body);
|
||||||
crate::dashboard_events::DashboardEvent::Sent {
|
crate::dashboard_events::DashboardEvent::Sent {
|
||||||
seq: 0,
|
seq: 0,
|
||||||
from,
|
from,
|
||||||
to,
|
to,
|
||||||
body,
|
body,
|
||||||
at,
|
at,
|
||||||
|
file_refs,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
crate::broker::MessageEvent::Delivered { from, to, body, at } => {
|
crate::broker::MessageEvent::Delivered { from, to, body, at } => {
|
||||||
|
let file_refs = scan_validated_paths(&body);
|
||||||
crate::dashboard_events::DashboardEvent::Delivered {
|
crate::dashboard_events::DashboardEvent::Delivered {
|
||||||
seq: 0,
|
seq: 0,
|
||||||
from,
|
from,
|
||||||
to,
|
to,
|
||||||
body,
|
body,
|
||||||
at,
|
at,
|
||||||
|
file_refs,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
@ -911,15 +915,17 @@ struct StateFileQuery {
|
||||||
/// traversal and symlink games can't escape the roots. Files larger
|
/// traversal and symlink games can't escape the roots. Files larger
|
||||||
/// than `MAX_BYTES` are truncated with a banner so a runaway log
|
/// than `MAX_BYTES` are truncated with a banner so a runaway log
|
||||||
/// can't OOM the browser.
|
/// can't OOM the browser.
|
||||||
async fn get_state_file(
|
/// Resolve a caller-supplied path string to a canonical host path
|
||||||
axum::extract::Query(q): axum::extract::Query<StateFileQuery>,
|
/// that has been verified against the allow-list. Returns `Err`
|
||||||
) -> Response {
|
/// with a human-readable reason for every failure mode (path
|
||||||
const MAX_BYTES: usize = 1 << 20; // 1 MiB
|
/// outside roots, canonicalize failure, escape via symlink,
|
||||||
|
/// per-agent subdir not `state`). Shared by `get_state_file` (read)
|
||||||
|
/// and `post_state_file_check` (existence probe) so both endpoints
|
||||||
|
/// apply identical security rules.
|
||||||
|
fn resolve_state_path(raw: &str) -> std::result::Result<std::path::PathBuf, String> {
|
||||||
const AGENTS_ROOT: &str = "/var/lib/hyperhive/agents";
|
const AGENTS_ROOT: &str = "/var/lib/hyperhive/agents";
|
||||||
const SHARED_ROOT: &str = "/var/lib/hyperhive/shared";
|
const SHARED_ROOT: &str = "/var/lib/hyperhive/shared";
|
||||||
let raw = q.path.trim();
|
let raw = raw.trim();
|
||||||
// Translate the container-view forms to host paths so the
|
|
||||||
// allow-list check has a single canonical shape to match.
|
|
||||||
let mapped: std::path::PathBuf = if let Some(rest) = raw.strip_prefix("/agents/") {
|
let mapped: std::path::PathBuf = if let Some(rest) = raw.strip_prefix("/agents/") {
|
||||||
std::path::PathBuf::from(format!("{AGENTS_ROOT}/{rest}"))
|
std::path::PathBuf::from(format!("{AGENTS_ROOT}/{rest}"))
|
||||||
} else if let Some(rest) = raw.strip_prefix("/shared/") {
|
} else if let Some(rest) = raw.strip_prefix("/shared/") {
|
||||||
|
|
@ -927,38 +933,84 @@ async fn get_state_file(
|
||||||
} else if raw.starts_with(AGENTS_ROOT) || raw.starts_with(SHARED_ROOT) {
|
} else if raw.starts_with(AGENTS_ROOT) || raw.starts_with(SHARED_ROOT) {
|
||||||
std::path::PathBuf::from(raw)
|
std::path::PathBuf::from(raw)
|
||||||
} else {
|
} else {
|
||||||
return error_response(&format!("state-file: path not in allow-list: {raw}"));
|
return Err(format!("path not in allow-list: {raw}"));
|
||||||
};
|
};
|
||||||
// Canonicalise so `..` / symlinks resolve before the prefix
|
let canonical = std::fs::canonicalize(&mapped).map_err(|e| format!("{}: {e}", mapped.display()))?;
|
||||||
// check. A failure here means the path doesn't exist on disk
|
if !(canonical.starts_with(AGENTS_ROOT) || canonical.starts_with(SHARED_ROOT)) {
|
||||||
// (or we can't reach it) — surface the underlying error.
|
return Err(format!(
|
||||||
let canonical = match std::fs::canonicalize(&mapped) {
|
"resolved path escapes allow-list: {}",
|
||||||
Ok(p) => p,
|
|
||||||
Err(e) => return error_response(&format!("state-file: {}: {e}", mapped.display())),
|
|
||||||
};
|
|
||||||
let allowed = canonical.starts_with(AGENTS_ROOT) || canonical.starts_with(SHARED_ROOT);
|
|
||||||
if !allowed {
|
|
||||||
return error_response(&format!(
|
|
||||||
"state-file: resolved path escapes allow-list: {}",
|
|
||||||
canonical.display()
|
canonical.display()
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
// For per-agent paths, also require the second-from-root
|
|
||||||
// component to be `state` (not `claude` or `config`). Claude
|
|
||||||
// creds shouldn't leak through this endpoint; config is the
|
|
||||||
// applied repo (already exposed via /api/agent-config). Reading
|
|
||||||
// `/var/lib/hyperhive/agents/<n>/state/...` is the intended use.
|
|
||||||
if let Ok(rel) = canonical.strip_prefix(AGENTS_ROOT) {
|
if let Ok(rel) = canonical.strip_prefix(AGENTS_ROOT) {
|
||||||
let mut components = rel.components();
|
let mut components = rel.components();
|
||||||
let _agent = components.next();
|
let _agent = components.next();
|
||||||
let dir = components.next().and_then(|c| c.as_os_str().to_str());
|
let dir = components.next().and_then(|c| c.as_os_str().to_str());
|
||||||
if dir != Some("state") {
|
if dir != Some("state") {
|
||||||
return error_response(&format!(
|
return Err(format!(
|
||||||
"state-file: only per-agent state/ is readable here ({} dir not allowed)",
|
"only per-agent state/ is readable here ({} dir not allowed)",
|
||||||
dir.unwrap_or("(root)")
|
dir.unwrap_or("(root)")
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
Ok(canonical)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Scan `body` for path-shaped tokens, validate each against the
|
||||||
|
/// allow-list, return the unique set of tokens that resolve to a
|
||||||
|
/// regular file. Called at broker-message ingest time so the
|
||||||
|
/// dashboard event already carries the verified set — no client-
|
||||||
|
/// side probe endpoint required, and historical messages get the
|
||||||
|
/// same treatment on `/dashboard/history` backfill.
|
||||||
|
///
|
||||||
|
/// Tokenisation: split on whitespace + a handful of trailing
|
||||||
|
/// punctuation chars (`,;:)]}`) that commonly follow paths in
|
||||||
|
/// natural-language text but aren't part of the path itself. Any
|
||||||
|
/// token starting with `/agents/`, `/shared/`, or
|
||||||
|
/// `/var/lib/hyperhive/{agents,shared}/` is a candidate. The
|
||||||
|
/// allow-list + is_file check happens via the same
|
||||||
|
/// `resolve_state_path` helper the read endpoint uses, so the
|
||||||
|
/// security rules can't drift.
|
||||||
|
pub(crate) fn scan_validated_paths(body: &str) -> Vec<String> {
|
||||||
|
const PREFIXES: [&str; 4] = [
|
||||||
|
"/agents/",
|
||||||
|
"/shared/",
|
||||||
|
"/var/lib/hyperhive/agents/",
|
||||||
|
"/var/lib/hyperhive/shared/",
|
||||||
|
];
|
||||||
|
let mut out = Vec::<String>::new();
|
||||||
|
for raw in body.split(|c: char| c.is_whitespace()) {
|
||||||
|
// Trim trailing natural-language punctuation that wouldn't
|
||||||
|
// be part of any real path. Inline rather than via a regex
|
||||||
|
// dep — the set is small and the call is hot.
|
||||||
|
let token = raw.trim_end_matches(|c: char| matches!(c, ',' | ';' | ':' | ')' | ']' | '}' | '.' | '\'' | '"'));
|
||||||
|
if token.is_empty() {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if !PREFIXES.iter().any(|p| token.starts_with(p)) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
// Cheap dedupe — typical message has 0-3 refs.
|
||||||
|
if out.iter().any(|s| s == token) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if let Ok(canonical) = resolve_state_path(token) {
|
||||||
|
if std::fs::metadata(&canonical).is_ok_and(|m| m.is_file()) {
|
||||||
|
out.push(token.to_owned());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
out
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn get_state_file(
|
||||||
|
axum::extract::Query(q): axum::extract::Query<StateFileQuery>,
|
||||||
|
) -> Response {
|
||||||
|
const MAX_BYTES: usize = 1 << 20; // 1 MiB
|
||||||
|
let canonical = match resolve_state_path(&q.path) {
|
||||||
|
Ok(p) => p,
|
||||||
|
Err(e) => return error_response(&format!("state-file: {e}")),
|
||||||
|
};
|
||||||
let meta = match std::fs::metadata(&canonical) {
|
let meta = match std::fs::metadata(&canonical) {
|
||||||
Ok(m) => m,
|
Ok(m) => m,
|
||||||
Err(e) => return error_response(&format!("state-file: stat {}: {e}", canonical.display())),
|
Err(e) => return error_response(&format!("state-file: stat {}: {e}", canonical.display())),
|
||||||
|
|
|
||||||
|
|
@ -31,20 +31,31 @@ use crate::container_view::ContainerView;
|
||||||
#[serde(rename_all = "snake_case", tag = "kind")]
|
#[serde(rename_all = "snake_case", tag = "kind")]
|
||||||
pub enum DashboardEvent {
|
pub enum DashboardEvent {
|
||||||
/// Broker `Sent` event mirrored onto the dashboard channel.
|
/// Broker `Sent` event mirrored onto the dashboard channel.
|
||||||
|
/// `file_refs` carries every path-shaped token in `body` that
|
||||||
|
/// hive-c0re verified is a regular file under the allow-listed
|
||||||
|
/// roots (per-agent `state/` + `shared/`). The forwarder
|
||||||
|
/// pre-validates so the dashboard doesn't need a probe
|
||||||
|
/// endpoint — the client renders anchors only for tokens that
|
||||||
|
/// appear in this list, everything else stays plain text.
|
||||||
Sent {
|
Sent {
|
||||||
seq: u64,
|
seq: u64,
|
||||||
from: String,
|
from: String,
|
||||||
to: String,
|
to: String,
|
||||||
body: String,
|
body: String,
|
||||||
at: i64,
|
at: i64,
|
||||||
|
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||||
|
file_refs: Vec<String>,
|
||||||
},
|
},
|
||||||
/// Broker `Delivered` event mirrored onto the dashboard channel.
|
/// Broker `Delivered` event mirrored onto the dashboard channel.
|
||||||
|
/// `file_refs` is the same shape as `Sent`.
|
||||||
Delivered {
|
Delivered {
|
||||||
seq: u64,
|
seq: u64,
|
||||||
from: String,
|
from: String,
|
||||||
to: String,
|
to: String,
|
||||||
body: String,
|
body: String,
|
||||||
at: i64,
|
at: i64,
|
||||||
|
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||||
|
file_refs: Vec<String>,
|
||||||
},
|
},
|
||||||
/// A new approval landed in the pending queue. Payload carries
|
/// A new approval landed in the pending queue. Payload carries
|
||||||
/// enough to render the dashboard row without a `/api/state`
|
/// enough to render the dashboard row without a `/api/state`
|
||||||
|
|
|
||||||
|
|
@ -226,21 +226,25 @@ fn spawn_broker_to_dashboard_forwarder(coord: Arc<Coordinator>) {
|
||||||
loop {
|
loop {
|
||||||
match rx.recv().await {
|
match rx.recv().await {
|
||||||
Ok(MessageEvent::Sent { from, to, body, at }) => {
|
Ok(MessageEvent::Sent { from, to, body, at }) => {
|
||||||
|
let file_refs = dashboard::scan_validated_paths(&body);
|
||||||
coord.emit_dashboard_event(DashboardEvent::Sent {
|
coord.emit_dashboard_event(DashboardEvent::Sent {
|
||||||
seq: coord.next_seq(),
|
seq: coord.next_seq(),
|
||||||
from,
|
from,
|
||||||
to,
|
to,
|
||||||
body,
|
body,
|
||||||
at,
|
at,
|
||||||
|
file_refs,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
Ok(MessageEvent::Delivered { from, to, body, at }) => {
|
Ok(MessageEvent::Delivered { from, to, body, at }) => {
|
||||||
|
let file_refs = dashboard::scan_validated_paths(&body);
|
||||||
coord.emit_dashboard_event(DashboardEvent::Delivered {
|
coord.emit_dashboard_event(DashboardEvent::Delivered {
|
||||||
seq: coord.next_seq(),
|
seq: coord.next_seq(),
|
||||||
from,
|
from,
|
||||||
to,
|
to,
|
||||||
body,
|
body,
|
||||||
at,
|
at,
|
||||||
|
file_refs,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
|
Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue