Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6db38cf70c | ||
|
|
7d93dd9db4 | ||
|
|
f65ee88269 |
17 changed files with 278 additions and 74 deletions
39
TODO.md
39
TODO.md
|
|
@ -16,21 +16,32 @@ Pick anything from here when relevant. Cross-cutting design notes live in
|
||||||
claude-code's `--allowedTools` extended grammar. Likely lives in
|
claude-code's `--allowedTools` extended grammar. Likely lives in
|
||||||
`agent.nix` so each agent can scope its own shell surface.
|
`agent.nix` so each agent can scope its own shell surface.
|
||||||
|
|
||||||
|
## Per-agent extension
|
||||||
|
|
||||||
|
- **Custom per-agent MCP tools.** Today every sub-agent gets the
|
||||||
|
same fixed MCP surface (`send`, `recv`). To move bitburner-agent
|
||||||
|
(and anything else with rich domain tooling) into hyperhive, an
|
||||||
|
agent needs a way to ship its own tools alongside hyperhive's.
|
||||||
|
Sketch: `agent.nix` declares a list of extra MCP servers
|
||||||
|
(command + args + env), each registered into the agent's
|
||||||
|
`--mcp-config` blob at flake-render time. The harness MCP server
|
||||||
|
remains the hyperhive surface; new servers slot in as additional
|
||||||
|
entries under `mcpServers.<name>` so claude sees them as
|
||||||
|
`mcp__<name>__<tool>`. Per-agent tool whitelist (`allowedTools`)
|
||||||
|
derived from the same config so the operator stays in control of
|
||||||
|
what's exposed.
|
||||||
|
|
||||||
## Per-agent settings
|
## Per-agent settings
|
||||||
|
|
||||||
- **Model override.** Hard-coded to `haiku` in the turn loop right now.
|
- **Model override persistence.** `/model <name>` already switches
|
||||||
Surface as a per-agent override: operator via dashboard, manager via
|
the model at runtime via `Bus::set_model`; the chip on the agent
|
||||||
`request_apply_commit` setting an attr on the agent's flake (most natural
|
page reflects the current value. Override is in-memory only and
|
||||||
place since the flake already carries per-agent env/identity). Pair with
|
resets on harness restart — by design for now, but consider
|
||||||
a **model status** indicator on the agent page (active / queued / last
|
optional persistence (`/state/model` file?) so an operator-set
|
||||||
switched) once the override is in place.
|
model survives a rebuild.
|
||||||
|
|
||||||
## UI / UX
|
## UI / UX
|
||||||
|
|
||||||
- **State badge: napping state.** Idle / thinking / compacting
|
|
||||||
already ship from server-side `TurnState`. Add `napping 😴`
|
|
||||||
once the `nap` tool exists — it just adds a new `TurnState`
|
|
||||||
variant the harness flips into for the duration of the nap.
|
|
||||||
- **Terminal: `/model` slash command.** Operator-typeable model
|
- **Terminal: `/model` slash command.** Operator-typeable model
|
||||||
override from the terminal. Depends on the model-override work
|
override from the terminal. Depends on the model-override work
|
||||||
above; once an override mechanism exists, wire a `/model <name>`
|
above; once an override mechanism exists, wire a `/model <name>`
|
||||||
|
|
@ -81,14 +92,6 @@ Pick anything from here when relevant. Cross-cutting design notes live in
|
||||||
|
|
||||||
## Loop substance
|
## Loop substance
|
||||||
|
|
||||||
- **`nap` tool.** Agent-side MCP tool `mcp__hyperhive__nap(seconds)` that
|
|
||||||
parks the turn loop for a short while before next-message processing.
|
|
||||||
Use cases: agent decides it has nothing useful to do, or wants to
|
|
||||||
throttle itself between rapid wake events. Implementation: harness
|
|
||||||
records a "wake-not-before" timestamp; `recv_blocking` skips the long
|
|
||||||
poll until that ts; the state badge reads `napping · MM:SS` during.
|
|
||||||
Operator can cancel via the same `/cancel` slash command or a
|
|
||||||
dashboard button.
|
|
||||||
- **Notes compaction.** `/state/` is bind-mounted persistently and agents
|
- **Notes compaction.** `/state/` is bind-mounted persistently and agents
|
||||||
are told (in the system prompt) to keep `/state/notes.md` for durable
|
are told (in the system prompt) to keep `/state/notes.md` for durable
|
||||||
knowledge — but we don't currently nudge them to compact when notes
|
knowledge — but we don't currently nudge them to compact when notes
|
||||||
|
|
|
||||||
|
|
@ -171,6 +171,15 @@ pre.diff {
|
||||||
font-size: 0.8em;
|
font-size: 0.8em;
|
||||||
letter-spacing: 0.05em;
|
letter-spacing: 0.05em;
|
||||||
}
|
}
|
||||||
|
.model-chip {
|
||||||
|
display: inline-block;
|
||||||
|
padding: 0.1em 0.6em;
|
||||||
|
border: 1px solid var(--purple-dim);
|
||||||
|
border-radius: 999px;
|
||||||
|
color: var(--cyan);
|
||||||
|
font-size: 0.78em;
|
||||||
|
letter-spacing: 0.04em;
|
||||||
|
}
|
||||||
.btn-dashlink {
|
.btn-dashlink {
|
||||||
color: var(--cyan);
|
color: var(--cyan);
|
||||||
border: 1px solid var(--cyan);
|
border: 1px solid var(--cyan);
|
||||||
|
|
|
||||||
|
|
@ -174,8 +174,31 @@
|
||||||
{ name: '/clear', desc: 'wipe the terminal panel (local-only)' },
|
{ name: '/clear', desc: 'wipe the terminal panel (local-only)' },
|
||||||
{ name: '/cancel', desc: 'SIGINT the in-flight claude turn' },
|
{ name: '/cancel', desc: 'SIGINT the in-flight claude turn' },
|
||||||
{ name: '/compact', desc: 'compact the persistent claude session' },
|
{ name: '/compact', desc: 'compact the persistent claude session' },
|
||||||
|
{ name: '/model', desc: '/model <name> — switch claude model for future turns' },
|
||||||
];
|
];
|
||||||
|
|
||||||
|
async function postModel(name) {
|
||||||
|
try {
|
||||||
|
const resp = await fetch('/api/model', {
|
||||||
|
method: 'POST',
|
||||||
|
headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
|
||||||
|
body: new URLSearchParams({ model: name }),
|
||||||
|
redirect: 'manual',
|
||||||
|
});
|
||||||
|
const ok = resp.ok || resp.type === 'opaqueredirect'
|
||||||
|
|| (resp.status >= 200 && resp.status < 400);
|
||||||
|
if (!ok && termAPI) {
|
||||||
|
const text = await resp.text().catch(() => '');
|
||||||
|
termAPI.row('turn-end-fail', '✗ /model failed: ' + resp.status
|
||||||
|
+ (text ? ' — ' + text : ''));
|
||||||
|
} else {
|
||||||
|
refreshState();
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
if (termAPI) termAPI.row('turn-end-fail', '✗ /model failed: ' + err);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async function postSimple(url, label) {
|
async function postSimple(url, label) {
|
||||||
try {
|
try {
|
||||||
const resp = await fetch(url, { method: 'POST', redirect: 'manual' });
|
const resp = await fetch(url, { method: 'POST', redirect: 'manual' });
|
||||||
|
|
@ -213,6 +236,16 @@
|
||||||
case '/compact':
|
case '/compact':
|
||||||
postCompact();
|
postCompact();
|
||||||
return true;
|
return true;
|
||||||
|
case '/model': {
|
||||||
|
const parts = trimmed.split(/\s+/);
|
||||||
|
if (parts.length < 2 || !parts[1]) {
|
||||||
|
termAPI.row('turn-end-fail',
|
||||||
|
'✗ /model needs a name (e.g. /model haiku, /model sonnet, /model opus)');
|
||||||
|
} else {
|
||||||
|
postModel(parts[1]);
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
}
|
||||||
default:
|
default:
|
||||||
termAPI.row('turn-end-fail', '✗ unknown slash command: ' + cmd + ' — try /help');
|
termAPI.row('turn-end-fail', '✗ unknown slash command: ' + cmd + ' — try /help');
|
||||||
return true;
|
return true;
|
||||||
|
|
@ -365,6 +398,13 @@
|
||||||
list.append(li);
|
list.append(li);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
function renderModelChip(model) {
|
||||||
|
const el_ = $('model-chip');
|
||||||
|
if (!el_) return;
|
||||||
|
if (!model) { el_.hidden = true; return; }
|
||||||
|
el_.hidden = false;
|
||||||
|
el_.textContent = 'model · ' + model;
|
||||||
|
}
|
||||||
function renderLastTurn(ms) {
|
function renderLastTurn(ms) {
|
||||||
const el_ = $('last-turn');
|
const el_ = $('last-turn');
|
||||||
if (!el_) return;
|
if (!el_) return;
|
||||||
|
|
@ -424,6 +464,7 @@
|
||||||
} else if (s.turn_state) {
|
} else if (s.turn_state) {
|
||||||
setStateAbs(s.turn_state, s.turn_state_since);
|
setStateAbs(s.turn_state, s.turn_state_since);
|
||||||
}
|
}
|
||||||
|
renderModelChip(s.model);
|
||||||
// Skip the re-render if nothing structurally changed. The most
|
// Skip the re-render if nothing structurally changed. The most
|
||||||
// common case is `online` polling itself — without this guard, the
|
// common case is `online` polling itself — without this guard, the
|
||||||
// operator's <input value> gets clobbered every cycle.
|
// operator's <input value> gets clobbered every cycle.
|
||||||
|
|
|
||||||
|
|
@ -15,6 +15,7 @@
|
||||||
|
|
||||||
<div id="state-row">
|
<div id="state-row">
|
||||||
<span id="state-badge" class="state-badge state-loading">… booting</span>
|
<span id="state-badge" class="state-badge state-loading">… booting</span>
|
||||||
|
<span id="model-chip" class="model-chip" hidden></span>
|
||||||
<span id="last-turn" class="last-turn" hidden></span>
|
<span id="last-turn" class="last-turn" hidden></span>
|
||||||
<button type="button" id="cancel-btn" class="btn-cancel-turn" hidden>■ cancel turn</button>
|
<button type="button" id="cancel-btn" class="btn-cancel-turn" hidden>■ cancel turn</button>
|
||||||
</div>
|
</div>
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,7 @@ You are hyperhive agent `{label}` in a multi-agent system.
|
||||||
|
|
||||||
Tools (hyperhive surface):
|
Tools (hyperhive surface):
|
||||||
|
|
||||||
- `mcp__hyperhive__recv()` — drain one more message from your inbox (returns `(empty)` if nothing pending).
|
- `mcp__hyperhive__recv(wait_seconds?)` — drain one more message from your inbox (returns `(empty)` if nothing pending after the wait). Without `wait_seconds` it long-polls 30s. To **wait** for work when you have nothing else useful to do this turn, call with a long wait (e.g. `wait_seconds: 180`, the max) — you'll be woken instantly when a message arrives, otherwise return after the timeout. That is strictly better than calling `recv` repeatedly with short waits: lower latency on new work, fewer turns, no busy-loop. Never use a fixed `sleep` shell command for the same purpose.
|
||||||
- `mcp__hyperhive__send(to, body)` — message a peer (by their name) or the operator (recipient `operator`, surfaces in the dashboard).
|
- `mcp__hyperhive__send(to, body)` — message a peer (by their name) or the operator (recipient `operator`, surfaces in the dashboard).
|
||||||
|
|
||||||
Need new packages, env vars, or other NixOS config for yourself? You can't edit your own config directly — message the manager (recipient `manager`) describing what you need + why. The manager evaluates the request (it doesn't rubber-stamp), edits `/agents/{label}/config/agent.nix` on your behalf, commits, and submits an approval that the operator can accept on the dashboard; on approve hive-c0re rebuilds your container with the new config.
|
Need new packages, env vars, or other NixOS config for yourself? You can't edit your own config directly — message the manager (recipient `manager`) describing what you need + why. The manager evaluates the request (it doesn't rubber-stamp), edits `/agents/{label}/config/agent.nix` on your behalf, commits, and submits an approval that the operator can accept on the dashboard; on approve hive-c0re rebuilds your container with the new config.
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,7 @@ You are the hyperhive manager `{label}` in a multi-agent system. You coordinate
|
||||||
|
|
||||||
Tools (hyperhive surface):
|
Tools (hyperhive surface):
|
||||||
|
|
||||||
- `mcp__hyperhive__recv()` — drain one more message from your inbox.
|
- `mcp__hyperhive__recv(wait_seconds?)` — drain one more message from your inbox. Without `wait_seconds` it long-polls 30s. To **wait** when you have nothing else to do, call with a long wait (e.g. `wait_seconds: 180`, the max) — you'll wake instantly on new work, otherwise return after the timeout. Use this instead of ending the turn or sleeping in a Bash command.
|
||||||
- `mcp__hyperhive__send(to, body)` — message an agent (by name), another peer, or the operator (`operator` surfaces in the dashboard).
|
- `mcp__hyperhive__send(to, body)` — message an agent (by name), another peer, or the operator (`operator` surfaces in the dashboard).
|
||||||
- `mcp__hyperhive__request_spawn(name)` — queue a brand-new sub-agent for operator approval (≤9 char name).
|
- `mcp__hyperhive__request_spawn(name)` — queue a brand-new sub-agent for operator approval (≤9 char name).
|
||||||
- `mcp__hyperhive__kill(name)` — graceful stop on a sub-agent. No approval required.
|
- `mcp__hyperhive__kill(name)` — graceful stop on a sub-agent. No approval required.
|
||||||
|
|
|
||||||
|
|
@ -116,7 +116,8 @@ async fn serve(
|
||||||
let label = std::env::var("HIVE_LABEL").unwrap_or_else(|_| "hive-ag3nt".into());
|
let label = std::env::var("HIVE_LABEL").unwrap_or_else(|_| "hive-ag3nt".into());
|
||||||
let system_prompt = turn::write_system_prompt(socket, &label, mcp::Flavor::Agent).await?;
|
let system_prompt = turn::write_system_prompt(socket, &label, mcp::Flavor::Agent).await?;
|
||||||
loop {
|
loop {
|
||||||
let recv: Result<AgentResponse> = client::request(socket, &AgentRequest::Recv).await;
|
let recv: Result<AgentResponse> =
|
||||||
|
client::request(socket, &AgentRequest::Recv { wait_seconds: None }).await;
|
||||||
match recv {
|
match recv {
|
||||||
Ok(AgentResponse::Message { from, body }) => {
|
Ok(AgentResponse::Message { from, body }) => {
|
||||||
tracing::info!(%from, %body, "inbox");
|
tracing::info!(%from, %body, "inbox");
|
||||||
|
|
|
||||||
|
|
@ -96,7 +96,8 @@ async fn serve(socket: &Path, interval: Duration, bus: Bus) -> Result<()> {
|
||||||
let label = std::env::var("HIVE_LABEL").unwrap_or_else(|_| "hm1nd".into());
|
let label = std::env::var("HIVE_LABEL").unwrap_or_else(|_| "hm1nd".into());
|
||||||
let system_prompt = turn::write_system_prompt(socket, &label, mcp::Flavor::Manager).await?;
|
let system_prompt = turn::write_system_prompt(socket, &label, mcp::Flavor::Manager).await?;
|
||||||
loop {
|
loop {
|
||||||
let recv: Result<ManagerResponse> = client::request(socket, &ManagerRequest::Recv).await;
|
let recv: Result<ManagerResponse> =
|
||||||
|
client::request(socket, &ManagerRequest::Recv { wait_seconds: None }).await;
|
||||||
match recv {
|
match recv {
|
||||||
Ok(ManagerResponse::Message { from, body }) => {
|
Ok(ManagerResponse::Message { from, body }) => {
|
||||||
if from == SYSTEM_SENDER {
|
if from == SYSTEM_SENDER {
|
||||||
|
|
|
||||||
|
|
@ -140,6 +140,13 @@ pub enum TurnState {
|
||||||
Compacting,
|
Compacting,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Default claude model when nothing's been set at runtime. The
|
||||||
|
/// operator can switch via `/model <name>` in the web terminal; the
|
||||||
|
/// chosen model lives in `Bus::model` for the rest of the harness
|
||||||
|
/// process's life (resets on restart, by design — operator overrides
|
||||||
|
/// shouldn't survive accidentally).
|
||||||
|
pub const DEFAULT_MODEL: &str = "haiku";
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
pub struct Bus {
|
pub struct Bus {
|
||||||
tx: Arc<broadcast::Sender<LiveEvent>>,
|
tx: Arc<broadcast::Sender<LiveEvent>>,
|
||||||
|
|
@ -149,6 +156,9 @@ pub struct Bus {
|
||||||
store: Option<Arc<EventStore>>,
|
store: Option<Arc<EventStore>>,
|
||||||
/// Current turn-loop state + since-when (unix seconds).
|
/// Current turn-loop state + since-when (unix seconds).
|
||||||
state: Arc<Mutex<(TurnState, i64)>>,
|
state: Arc<Mutex<(TurnState, i64)>>,
|
||||||
|
/// Model name passed to `claude --model`. Default `haiku`; the
|
||||||
|
/// operator can override at runtime via `POST /api/model`.
|
||||||
|
model: Arc<Mutex<String>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Bus {
|
impl Bus {
|
||||||
|
|
@ -171,9 +181,23 @@ impl Bus {
|
||||||
tx: Arc::new(tx),
|
tx: Arc::new(tx),
|
||||||
store,
|
store,
|
||||||
state: Arc::new(Mutex::new((TurnState::Idle, now_unix()))),
|
state: Arc::new(Mutex::new((TurnState::Idle, now_unix()))),
|
||||||
|
model: Arc::new(Mutex::new(DEFAULT_MODEL.to_owned())),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Currently-selected claude model name. Read on every turn so a
|
||||||
|
/// `/model <name>` flip takes effect on the next turn.
|
||||||
|
#[must_use]
|
||||||
|
pub fn model(&self) -> String {
|
||||||
|
self.model.lock().unwrap().clone()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Switch the model for future turns. The current turn (if any)
|
||||||
|
/// keeps the model it was already running.
|
||||||
|
pub fn set_model(&self, name: impl Into<String>) {
|
||||||
|
*self.model.lock().unwrap() = name.into();
|
||||||
|
}
|
||||||
|
|
||||||
/// Update the harness's authoritative turn-loop state. Records
|
/// Update the harness's authoritative turn-loop state. Records
|
||||||
/// the transition time so `state_snapshot` can return a since-age.
|
/// the transition time so `state_snapshot` can return a since-age.
|
||||||
pub fn set_state(&self, next: TurnState) {
|
pub fn set_state(&self, next: TurnState) {
|
||||||
|
|
|
||||||
|
|
@ -116,7 +116,14 @@ pub struct SendArgs {
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
|
#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
|
||||||
pub struct RecvArgs {}
|
pub struct RecvArgs {
|
||||||
|
/// How long to long-poll for a new message before returning the
|
||||||
|
/// empty marker. Capped at 60s server-side. Default (None) is
|
||||||
|
/// 30s. Useful when an agent wants to throttle wakes without
|
||||||
|
/// actually napping — pick a longer wait to coalesce bursts.
|
||||||
|
#[serde(default)]
|
||||||
|
pub wait_seconds: Option<u64>,
|
||||||
|
}
|
||||||
|
|
||||||
/// Per-agent tool surface. Holds the socket path so each tool call doesn't
|
/// Per-agent tool surface. Holds the socket path so each tool call doesn't
|
||||||
/// re-derive it; the socket itself is the per-container `/run/hive/mcp.sock`.
|
/// re-derive it; the socket itself is the per-container `/run/hive/mcp.sock`.
|
||||||
|
|
@ -158,13 +165,20 @@ impl AgentServer {
|
||||||
|
|
||||||
#[tool(
|
#[tool(
|
||||||
description = "Pop one message from this agent's inbox. Returns the sender and body, \
|
description = "Pop one message from this agent's inbox. Returns the sender and body, \
|
||||||
or an empty marker if nothing is waiting."
|
or an empty marker if nothing is waiting. Optional `wait_seconds` long-polls \
|
||||||
|
for that many seconds (capped at 180) before returning empty — default 30. \
|
||||||
|
Use a long wait_seconds (e.g. 120 or 180) when you have nothing else to do — \
|
||||||
|
it parks the turn until either a message arrives or the timeout fires, which \
|
||||||
|
is strictly better than a fixed sleep because incoming work wakes you instantly."
|
||||||
)]
|
)]
|
||||||
async fn recv(&self, Parameters(_args): Parameters<RecvArgs>) -> String {
|
async fn recv(&self, Parameters(args): Parameters<RecvArgs>) -> String {
|
||||||
run_tool_envelope("recv", String::new(), async move {
|
let log = format!("{args:?}");
|
||||||
|
run_tool_envelope("recv", log, async move {
|
||||||
let resp = client::request::<_, hive_sh4re::AgentResponse>(
|
let resp = client::request::<_, hive_sh4re::AgentResponse>(
|
||||||
&self.socket,
|
&self.socket,
|
||||||
&hive_sh4re::AgentRequest::Recv,
|
&hive_sh4re::AgentRequest::Recv {
|
||||||
|
wait_seconds: args.wait_seconds,
|
||||||
|
},
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map(SocketReply::from);
|
.map(SocketReply::from);
|
||||||
|
|
@ -303,11 +317,19 @@ impl ManagerServer {
|
||||||
|
|
||||||
#[tool(
|
#[tool(
|
||||||
description = "Pop one message from the manager inbox. Returns sender + body, or \
|
description = "Pop one message from the manager inbox. Returns sender + body, or \
|
||||||
empty."
|
empty. Optional `wait_seconds` long-polls (capped at 180, default 30) so the \
|
||||||
|
manager can sit on Recv when there's nothing to do without burning turns — \
|
||||||
|
prefer a long wait (120 or 180) over ending a turn early; you'll wake \
|
||||||
|
instantly when work arrives."
|
||||||
)]
|
)]
|
||||||
async fn recv(&self, Parameters(_args): Parameters<RecvArgs>) -> String {
|
async fn recv(&self, Parameters(args): Parameters<RecvArgs>) -> String {
|
||||||
run_tool_envelope("recv", String::new(), async move {
|
let log = format!("{args:?}");
|
||||||
let resp = self.dispatch(hive_sh4re::ManagerRequest::Recv).await;
|
run_tool_envelope("recv", log, async move {
|
||||||
|
let resp = self
|
||||||
|
.dispatch(hive_sh4re::ManagerRequest::Recv {
|
||||||
|
wait_seconds: args.wait_seconds,
|
||||||
|
})
|
||||||
|
.await;
|
||||||
format_recv(resp)
|
format_recv(resp)
|
||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
|
|
|
||||||
|
|
@ -227,13 +227,14 @@ async fn run_claude(
|
||||||
flavor: mcp::Flavor,
|
flavor: mcp::Flavor,
|
||||||
mode: ClaudeMode,
|
mode: ClaudeMode,
|
||||||
) -> Result<bool> {
|
) -> Result<bool> {
|
||||||
|
let model = bus.model();
|
||||||
let mut cmd = Command::new("claude");
|
let mut cmd = Command::new("claude");
|
||||||
cmd.arg("--print")
|
cmd.arg("--print")
|
||||||
.arg("--verbose")
|
.arg("--verbose")
|
||||||
.arg("--output-format")
|
.arg("--output-format")
|
||||||
.arg("stream-json")
|
.arg("stream-json")
|
||||||
.arg("--model")
|
.arg("--model")
|
||||||
.arg("haiku")
|
.arg(&model)
|
||||||
.arg("--continue")
|
.arg("--continue")
|
||||||
.arg("--settings")
|
.arg("--settings")
|
||||||
.arg(settings);
|
.arg(settings);
|
||||||
|
|
|
||||||
|
|
@ -81,6 +81,7 @@ pub async fn serve(
|
||||||
.route("/login/cancel", post(post_login_cancel))
|
.route("/login/cancel", post(post_login_cancel))
|
||||||
.route("/api/cancel", post(post_cancel_turn))
|
.route("/api/cancel", post(post_cancel_turn))
|
||||||
.route("/api/compact", post(post_compact))
|
.route("/api/compact", post(post_compact))
|
||||||
|
.route("/api/model", post(post_set_model))
|
||||||
.with_state(state);
|
.with_state(state);
|
||||||
let addr = SocketAddr::from(([0, 0, 0, 0], port));
|
let addr = SocketAddr::from(([0, 0, 0, 0], port));
|
||||||
let listener = bind_with_retry(addr, "web UI").await?;
|
let listener = bind_with_retry(addr, "web UI").await?;
|
||||||
|
|
@ -93,16 +94,17 @@ pub async fn serve(
|
||||||
// Static assets + state snapshot
|
// Static assets + state snapshot
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
/// Bind a TCP listener, retrying on `AddrInUse` for up to ~20s.
|
/// Bind a TCP listener with `SO_REUSEADDR` set, retrying on
|
||||||
/// nspawn restarts can race the previous harness's socket release;
|
/// `AddrInUse` for up to ~20s. nspawn restarts can race the previous
|
||||||
/// without retry the new harness fails to bind and systemd just
|
/// harness's socket release; `SO_REUSEADDR` lets us reclaim a port
|
||||||
/// keeps restarting it. `SO_REUSEADDR` would be the proper fix but
|
/// still in `TIME_WAIT` from a clean previous exit, and the retry
|
||||||
/// would require socket2; retry is good enough here.
|
/// covers the case where the previous process is genuinely still
|
||||||
|
/// alive (systemd restart-delay overlap).
|
||||||
async fn bind_with_retry(addr: SocketAddr, label: &str) -> Result<tokio::net::TcpListener> {
|
async fn bind_with_retry(addr: SocketAddr, label: &str) -> Result<tokio::net::TcpListener> {
|
||||||
let mut delay_ms = 250u64;
|
let mut delay_ms = 250u64;
|
||||||
let mut attempts = 0u32;
|
let mut attempts = 0u32;
|
||||||
loop {
|
loop {
|
||||||
match tokio::net::TcpListener::bind(addr).await {
|
match try_bind(addr) {
|
||||||
Ok(l) => return Ok(l),
|
Ok(l) => return Ok(l),
|
||||||
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse && attempts < 12 => {
|
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse && attempts < 12 => {
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
|
|
@ -120,6 +122,16 @@ async fn bind_with_retry(addr: SocketAddr, label: &str) -> Result<tokio::net::Tc
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn try_bind(addr: SocketAddr) -> std::io::Result<tokio::net::TcpListener> {
|
||||||
|
let sock = match addr {
|
||||||
|
SocketAddr::V4(_) => tokio::net::TcpSocket::new_v4()?,
|
||||||
|
SocketAddr::V6(_) => tokio::net::TcpSocket::new_v6()?,
|
||||||
|
};
|
||||||
|
sock.set_reuseaddr(true)?;
|
||||||
|
sock.bind(addr)?;
|
||||||
|
sock.listen(1024)
|
||||||
|
}
|
||||||
|
|
||||||
async fn serve_index() -> impl IntoResponse {
|
async fn serve_index() -> impl IntoResponse {
|
||||||
(
|
(
|
||||||
[("content-type", "text/html; charset=utf-8")],
|
[("content-type", "text/html; charset=utf-8")],
|
||||||
|
|
@ -158,6 +170,10 @@ struct StateSnapshot {
|
||||||
/// client-side off this rather than tracking it from SSE events.
|
/// client-side off this rather than tracking it from SSE events.
|
||||||
turn_state: crate::events::TurnState,
|
turn_state: crate::events::TurnState,
|
||||||
turn_state_since: i64,
|
turn_state_since: i64,
|
||||||
|
/// Currently-active claude model name. Reflected on the page so
|
||||||
|
/// the operator can see what they just switched to (and what's
|
||||||
|
/// in flight). Mutable at runtime via `POST /api/model`.
|
||||||
|
model: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Serialize)]
|
#[derive(Serialize)]
|
||||||
|
|
@ -193,6 +209,7 @@ async fn api_state(State(state): State<AppState>) -> axum::Json<StateSnapshot> {
|
||||||
.unwrap_or(7000);
|
.unwrap_or(7000);
|
||||||
let inbox = recent_inbox(&state.socket, state.flavor).await;
|
let inbox = recent_inbox(&state.socket, state.flavor).await;
|
||||||
let (turn_state, turn_state_since) = state.bus.state_snapshot();
|
let (turn_state, turn_state_since) = state.bus.state_snapshot();
|
||||||
|
let model = state.bus.model();
|
||||||
axum::Json(StateSnapshot {
|
axum::Json(StateSnapshot {
|
||||||
label: state.label.clone(),
|
label: state.label.clone(),
|
||||||
dashboard_port,
|
dashboard_port,
|
||||||
|
|
@ -201,6 +218,7 @@ async fn api_state(State(state): State<AppState>) -> axum::Json<StateSnapshot> {
|
||||||
inbox,
|
inbox,
|
||||||
turn_state,
|
turn_state,
|
||||||
turn_state_since,
|
turn_state_since,
|
||||||
|
model,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -351,6 +369,30 @@ async fn post_login_cancel(State(state): State<AppState>) -> Response {
|
||||||
/// the "/compact done" note) lands in the live event panel like any
|
/// the "/compact done" note) lands in the live event panel like any
|
||||||
/// other turn. If a regular turn is in flight, claude's own session
|
/// other turn. If a regular turn is in flight, claude's own session
|
||||||
/// lock will reject this one and we surface the error as a Note.
|
/// lock will reject this one and we surface the error as a Note.
|
||||||
|
#[derive(Deserialize)]
|
||||||
|
struct ModelForm {
|
||||||
|
model: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Switch the model for future turns. The current turn (if any)
|
||||||
|
/// keeps its model; `/model <name>` applies starting with the next
|
||||||
|
/// `recv` cycle. Empty / whitespace-only inputs are rejected. No
|
||||||
|
/// claude-side validation — we just hand the string through to
|
||||||
|
/// `claude --model <name>`; an unknown model surfaces as a turn
|
||||||
|
/// failure in the live panel and the operator can revert.
|
||||||
|
async fn post_set_model(State(state): State<AppState>, Form(form): Form<ModelForm>) -> Response {
|
||||||
|
let name = form.model.trim();
|
||||||
|
if name.is_empty() {
|
||||||
|
return error_response("model: name required");
|
||||||
|
}
|
||||||
|
state.bus.set_model(name);
|
||||||
|
state.bus.emit(crate::events::LiveEvent::Note(format!(
|
||||||
|
"operator: /model — claude model set to '{name}' for future turns"
|
||||||
|
)));
|
||||||
|
tracing::info!(%name, "operator set model");
|
||||||
|
Redirect::to("/").into_response()
|
||||||
|
}
|
||||||
|
|
||||||
async fn post_compact(State(state): State<AppState>) -> Response {
|
async fn post_compact(State(state): State<AppState>) -> Response {
|
||||||
let bus = state.bus.clone();
|
let bus = state.bus.clone();
|
||||||
let socket = state.socket.clone();
|
let socket = state.socket.clone();
|
||||||
|
|
|
||||||
|
|
@ -77,9 +77,21 @@ async fn serve(stream: UnixStream, agent: String, broker: Arc<Broker>) -> Result
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// How long the long-poll `Recv` holds a connection open waiting for new
|
/// Default and max long-poll window for `Recv`. Caller can request a
|
||||||
/// mail. Set well below typical TCP/proxy idle limits.
|
/// shorter (or longer up to `RECV_LONG_POLL_MAX`) wait via the
|
||||||
const RECV_LONG_POLL: std::time::Duration = std::time::Duration::from_secs(30);
|
/// `wait_seconds` field; values above the cap are clamped. 180s
|
||||||
|
/// max keeps us under typical TCP/proxy idle limits while letting
|
||||||
|
/// agents park their turn until a message lands instead of busy-
|
||||||
|
/// looping with short waits.
|
||||||
|
const RECV_LONG_POLL_DEFAULT: std::time::Duration = std::time::Duration::from_secs(30);
|
||||||
|
const RECV_LONG_POLL_MAX: std::time::Duration = std::time::Duration::from_secs(180);
|
||||||
|
|
||||||
|
fn recv_timeout(wait_seconds: Option<u64>) -> std::time::Duration {
|
||||||
|
match wait_seconds {
|
||||||
|
Some(s) => std::time::Duration::from_secs(s).min(RECV_LONG_POLL_MAX),
|
||||||
|
None => RECV_LONG_POLL_DEFAULT,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async fn dispatch(req: &AgentRequest, agent: &str, broker: &Broker) -> AgentResponse {
|
async fn dispatch(req: &AgentRequest, agent: &str, broker: &Broker) -> AgentResponse {
|
||||||
match req {
|
match req {
|
||||||
|
|
@ -95,7 +107,10 @@ async fn dispatch(req: &AgentRequest, agent: &str, broker: &Broker) -> AgentResp
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
AgentRequest::Recv => match broker.recv_blocking(agent, RECV_LONG_POLL).await {
|
AgentRequest::Recv { wait_seconds } => match broker
|
||||||
|
.recv_blocking(agent, recv_timeout(*wait_seconds))
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(Some(msg)) => AgentResponse::Message {
|
Ok(Some(msg)) => AgentResponse::Message {
|
||||||
from: msg.from,
|
from: msg.from,
|
||||||
body: msg.body,
|
body: msg.body,
|
||||||
|
|
|
||||||
|
|
@ -72,13 +72,13 @@ pub async fn serve(port: u16, coord: Arc<Coordinator>) -> Result<()> {
|
||||||
// `/messages/stream` for broker traffic.
|
// `/messages/stream` for broker traffic.
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
/// Retry-on-AddrInUse bind. Same shape as the per-agent variant —
|
/// `SO_REUSEADDR` bind with retry. Mirrors the per-agent variant —
|
||||||
/// hive-c0re restarts also race the previous process's socket release.
|
/// hive-c0re restarts also race the previous process's socket release.
|
||||||
async fn bind_with_retry(addr: SocketAddr) -> Result<tokio::net::TcpListener> {
|
async fn bind_with_retry(addr: SocketAddr) -> Result<tokio::net::TcpListener> {
|
||||||
let mut delay_ms = 250u64;
|
let mut delay_ms = 250u64;
|
||||||
let mut attempts = 0u32;
|
let mut attempts = 0u32;
|
||||||
loop {
|
loop {
|
||||||
match tokio::net::TcpListener::bind(addr).await {
|
match try_bind(addr) {
|
||||||
Ok(l) => return Ok(l),
|
Ok(l) => return Ok(l),
|
||||||
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse && attempts < 12 => {
|
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse && attempts < 12 => {
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
|
|
@ -96,6 +96,16 @@ async fn bind_with_retry(addr: SocketAddr) -> Result<tokio::net::TcpListener> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn try_bind(addr: SocketAddr) -> std::io::Result<tokio::net::TcpListener> {
|
||||||
|
let sock = match addr {
|
||||||
|
SocketAddr::V4(_) => tokio::net::TcpSocket::new_v4()?,
|
||||||
|
SocketAddr::V6(_) => tokio::net::TcpSocket::new_v6()?,
|
||||||
|
};
|
||||||
|
sock.set_reuseaddr(true)?;
|
||||||
|
sock.bind(addr)?;
|
||||||
|
sock.listen(1024)
|
||||||
|
}
|
||||||
|
|
||||||
async fn serve_index() -> impl IntoResponse {
|
async fn serve_index() -> impl IntoResponse {
|
||||||
Html(include_str!("../assets/index.html"))
|
Html(include_str!("../assets/index.html"))
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -46,13 +46,18 @@ const DEFAULT_MEMORY_MAX: &str = "2G";
|
||||||
const DEFAULT_CPU_QUOTA: &str = "50%";
|
const DEFAULT_CPU_QUOTA: &str = "50%";
|
||||||
|
|
||||||
/// Returns the per-agent web UI port. Manager is fixed at `MANAGER_PORT`.
|
/// Returns the per-agent web UI port. Manager is fixed at `MANAGER_PORT`.
|
||||||
/// For sub-agents the port is sticky once chosen: looked up from
|
/// For sub-agents the port is sticky once chosen:
|
||||||
/// `agent_state_root(name)/port` if present, otherwise derived from
|
///
|
||||||
/// the FNV-1a hash of the name and *probed forward* through the
|
/// - **Port file present** (`state_root/port`): use it. End of story.
|
||||||
/// allocated range to skip any port another sub-agent has already
|
/// - **Port file absent, applied flake present**: this is a legacy
|
||||||
/// claimed (birthday-paradox collisions are real even at 2–3
|
/// agent whose container is already bound to the bare
|
||||||
/// agents). The chosen port is written back so subsequent calls
|
/// `port_hash(name)`. Don't probe; just migrate by writing that
|
||||||
/// resolve to the same value without re-probing.
|
/// value to the port file. The container stays where it is and
|
||||||
|
/// subsequent renders agree with it.
|
||||||
|
/// - **Port file absent, no applied flake**: this is a fresh spawn.
|
||||||
|
/// Probe forward from `port_hash(name)` to skip any port another
|
||||||
|
/// sub-agent has already claimed (via port file or legacy hash).
|
||||||
|
/// Write the chosen port back.
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn agent_web_port(name: &str) -> u16 {
|
pub fn agent_web_port(name: &str) -> u16 {
|
||||||
if name == MANAGER_NAME {
|
if name == MANAGER_NAME {
|
||||||
|
|
@ -66,27 +71,36 @@ pub fn agent_web_port(name: &str) -> u16 {
|
||||||
{
|
{
|
||||||
return port;
|
return port;
|
||||||
}
|
}
|
||||||
let taken = scan_taken_ports(name);
|
let applied_exists = crate::coordinator::Coordinator::agent_applied_dir(name).exists();
|
||||||
let start = port_hash(name);
|
let chosen = if applied_exists {
|
||||||
let mut port = start;
|
// Legacy agent — container already running on the hashed
|
||||||
for _ in 0..WEB_PORT_RANGE {
|
// port. Don't move it; just persist the value so future
|
||||||
if !taken.contains(&port) {
|
// calls bypass this path.
|
||||||
break;
|
port_hash(name)
|
||||||
|
} else {
|
||||||
|
let taken = scan_taken_ports(name);
|
||||||
|
let start = port_hash(name);
|
||||||
|
let mut port = start;
|
||||||
|
for _ in 0..WEB_PORT_RANGE {
|
||||||
|
if !taken.contains(&port) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
port = next_port(port);
|
||||||
|
if port == start {
|
||||||
|
// Range fully exhausted (very unlikely — 900 slots) —
|
||||||
|
// give up and use the hashed value; collisions are
|
||||||
|
// surfaced as bind errors by the harness retry loop.
|
||||||
|
tracing::warn!(%name, "agent_web_port: range exhausted, returning hash");
|
||||||
|
break;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
port = next_port(port);
|
port
|
||||||
if port == start {
|
};
|
||||||
// Range fully exhausted (very unlikely — 900 slots) —
|
|
||||||
// give up and just use the hashed value; collisions are
|
|
||||||
// surfaced as bind errors by the harness retry loop.
|
|
||||||
tracing::warn!(%name, "agent_web_port: range exhausted, returning hash");
|
|
||||||
return start;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
let _ = std::fs::create_dir_all(&state_root);
|
let _ = std::fs::create_dir_all(&state_root);
|
||||||
if let Err(e) = std::fs::write(&port_file, format!("{port}\n")) {
|
if let Err(e) = std::fs::write(&port_file, format!("{chosen}\n")) {
|
||||||
tracing::warn!(error = ?e, file = %port_file.display(), "persisting agent port failed");
|
tracing::warn!(error = ?e, file = %port_file.display(), "persisting agent port failed");
|
||||||
}
|
}
|
||||||
port
|
chosen
|
||||||
}
|
}
|
||||||
|
|
||||||
fn port_hash(name: &str) -> u16 {
|
fn port_hash(name: &str) -> u16 {
|
||||||
|
|
|
||||||
|
|
@ -69,7 +69,17 @@ async fn serve(stream: UnixStream, coord: Arc<Coordinator>) -> Result<()> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
const MANAGER_RECV_LONG_POLL: std::time::Duration = std::time::Duration::from_secs(30);
|
/// Default and max long-poll window for manager `Recv`. Caller can
|
||||||
|
/// request a shorter or longer (up to MAX) wait via `wait_seconds`.
|
||||||
|
const MANAGER_RECV_LONG_POLL_DEFAULT: std::time::Duration = std::time::Duration::from_secs(30);
|
||||||
|
const MANAGER_RECV_LONG_POLL_MAX: std::time::Duration = std::time::Duration::from_secs(180);
|
||||||
|
|
||||||
|
fn manager_recv_timeout(wait_seconds: Option<u64>) -> std::time::Duration {
|
||||||
|
match wait_seconds {
|
||||||
|
Some(s) => std::time::Duration::from_secs(s).min(MANAGER_RECV_LONG_POLL_MAX),
|
||||||
|
None => MANAGER_RECV_LONG_POLL_DEFAULT,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[allow(clippy::too_many_lines)]
|
#[allow(clippy::too_many_lines)]
|
||||||
async fn dispatch(req: &ManagerRequest, coord: &Arc<Coordinator>) -> ManagerResponse {
|
async fn dispatch(req: &ManagerRequest, coord: &Arc<Coordinator>) -> ManagerResponse {
|
||||||
|
|
@ -106,9 +116,9 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc<Coordinator>) -> ManagerResp
|
||||||
message: format!("{e:#}"),
|
message: format!("{e:#}"),
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
ManagerRequest::Recv => match coord
|
ManagerRequest::Recv { wait_seconds } => match coord
|
||||||
.broker
|
.broker
|
||||||
.recv_blocking(MANAGER_AGENT, MANAGER_RECV_LONG_POLL)
|
.recv_blocking(MANAGER_AGENT, manager_recv_timeout(*wait_seconds))
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(Some(msg)) => ManagerResponse::Message {
|
Ok(Some(msg)) => ManagerResponse::Message {
|
||||||
|
|
|
||||||
|
|
@ -166,8 +166,13 @@ pub struct InboxRow {
|
||||||
pub enum AgentRequest {
|
pub enum AgentRequest {
|
||||||
/// Send a message to another agent.
|
/// Send a message to another agent.
|
||||||
Send { to: String, body: String },
|
Send { to: String, body: String },
|
||||||
/// Pop one pending message from this agent's inbox.
|
/// Pop one pending message from this agent's inbox. Long-polls
|
||||||
Recv,
|
/// up to `wait_seconds` (capped at 60s server-side, default 30s
|
||||||
|
/// when None) before returning `Empty`.
|
||||||
|
Recv {
|
||||||
|
#[serde(default)]
|
||||||
|
wait_seconds: Option<u64>,
|
||||||
|
},
|
||||||
/// Non-mutating: how many pending messages are addressed to me?
|
/// Non-mutating: how many pending messages are addressed to me?
|
||||||
/// Used by the harness to render a status line after each tool call.
|
/// Used by the harness to render a status line after each tool call.
|
||||||
Status,
|
Status,
|
||||||
|
|
@ -274,7 +279,12 @@ pub enum ManagerRequest {
|
||||||
to: String,
|
to: String,
|
||||||
body: String,
|
body: String,
|
||||||
},
|
},
|
||||||
Recv,
|
/// Same shape as `AgentRequest::Recv` — caller-tunable long-poll
|
||||||
|
/// duration, capped at 60s server-side, default 30s when None.
|
||||||
|
Recv {
|
||||||
|
#[serde(default)]
|
||||||
|
wait_seconds: Option<u64>,
|
||||||
|
},
|
||||||
/// Non-mutating: pending message count, used to render a status line
|
/// Non-mutating: pending message count, used to render a status line
|
||||||
/// after each MCP tool call (mirrors `AgentRequest::Status`).
|
/// after each MCP tool call (mirrors `AgentRequest::Status`).
|
||||||
Status,
|
Status,
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue