remove request_next_turn: same-turn continuation is always worse than an external wake
This commit is contained in:
parent
b9aab7e923
commit
fffe0a2c29
8 changed files with 103 additions and 253 deletions
|
|
@ -322,7 +322,7 @@ binary flavor.
|
||||||
|---|---|
|
|---|---|
|
||||||
| `messaging` | `send`, `recv`, `ask`, `answer` |
|
| `messaging` | `send`, `recv`, `ask`, `answer` |
|
||||||
| `meta` | `get_agent_meta` (`set_status` is always-on, see below) |
|
| `meta` | `get_agent_meta` (`set_status` is always-on, see below) |
|
||||||
| `inbox` | `get_loose_ends`, `cancel_loose_end`, `remind`, `request_next_turn` |
|
| `inbox` | `get_loose_ends`, `cancel_loose_end`, `remind` |
|
||||||
| `execution` | vestigial — `mcp__bash__run` / `mcp__bash__status` are always available unconditionally via `extraMcpServers`; this group's entries expand to non-existent `mcp__hyperhive__run` / `mcp__hyperhive__status` and have no effect. See `docs/tools/bash.md`. |
|
| `execution` | vestigial — `mcp__bash__run` / `mcp__bash__status` are always available unconditionally via `extraMcpServers`; this group's entries expand to non-existent `mcp__hyperhive__run` / `mcp__hyperhive__status` and have no effect. See `docs/tools/bash.md`. |
|
||||||
| `lifecycle` | `kill`, `start`, `restart`, `update` *(privileged)* |
|
| `lifecycle` | `kill`, `start`, `restart`, `update` *(privileged)* |
|
||||||
| `approvals` | `request_init_config`, `request_update_meta_inputs` *(privileged)* |
|
| `approvals` | `request_init_config`, `request_update_meta_inputs` *(privileged)* |
|
||||||
|
|
|
||||||
|
|
@ -12,10 +12,8 @@ agents) runs:
|
||||||
loop does nothing but re-stat it every 5 s — no broker poll, no
|
loop does nothing but re-stat it every 5 s — no broker poll, no
|
||||||
claude process. Because step 1 is never reached, messages stay
|
claude process. Because step 1 is never reached, messages stay
|
||||||
queued and unacked, so a resume drains the backlog instead of
|
queued and unacked, so a resume drains the backlog instead of
|
||||||
losing it; reminders and todo wakes buffer in their channels. The
|
losing it; reminders and todo wakes buffer in their channels. Set
|
||||||
check runs before the self-continue slot is consumed, so a pending
|
it with `hivectl agent <name> pause` or the dashboard toggle; see
|
||||||
`request_next_turn` survives the pause. Set it with
|
|
||||||
`hivectl agent <name> pause` or the dashboard toggle; see
|
|
||||||
[persistence](persistence.md#-harnesspaused-per-agent).
|
[persistence](persistence.md#-harnesspaused-per-agent).
|
||||||
1. Long-poll `Recv` on its socket. The host-side broker
|
1. Long-poll `Recv` on its socket. The host-side broker
|
||||||
(`broker.rs::recv_blocking_batch`) returns immediately if there's
|
(`broker.rs::recv_blocking_batch`) returns immediately if there's
|
||||||
|
|
@ -141,21 +139,17 @@ only complete output silence for the window trips it. The harness sets the
|
||||||
window from `HIVE_TURN_IDLE_SECS` (`0` disables) and maps the driver's
|
window from `HIVE_TURN_IDLE_SECS` (`0` disables) and maps the driver's
|
||||||
`Error::IdleTimeout` onto `TurnError::ApiStall`.
|
`Error::IdleTimeout` onto `TurnError::ApiStall`.
|
||||||
|
|
||||||
After the outcome handler, the stats sink records a row and the
|
After the outcome handler, the stats sink records a row. `handle_turn`
|
||||||
`hyperhive-continue` sentinel (dropped by the `request_next_turn`
|
reports the result to `serve_loop` via `TurnControl { auth_failed }` —
|
||||||
MCP tool) is consumed if present. `handle_turn` reports the result
|
on auth failure the loop parks in `wait_for_login`; otherwise it loops
|
||||||
to `serve_loop` via `TurnControl { auth_failed, continue_requested,
|
straight back to the idle wait (step 1). There is no same-turn
|
||||||
pending }`. When a continue was requested, the turn did not
|
self-continue tool (removed — see forge #2777): every multi-step
|
||||||
auth-fail, and the inbox is empty (`pending == 0`), `serve_loop`
|
continuation rides an external wake instead — a new inbox message, a
|
||||||
drives the next turn in-process with a synthetic
|
`remind`, or an in-container todo wake (bash-task completion, forge
|
||||||
`{ from: "self", body: "continue" }` message (`synthetic_continue`)
|
notification, matrix activity). Ending the turn and letting one of
|
||||||
— it never goes through the broker, so the self-continue doesn't
|
those drive the next one is strictly better than parking in-process:
|
||||||
persist to sqlite or show up as a recv'able inbox message. If real
|
it checkpoints the session and observes wakes that only reach the
|
||||||
messages are already pending the continue is dropped: those messages
|
harness between turns.
|
||||||
drive the next turn(s) via `recv_next`, so an explicit self-wake
|
|
||||||
isn't needed (this is the `request_next_turn` contract — "no effect
|
|
||||||
if a new inbox message arrives before this turn ends"). The
|
|
||||||
`should_self_continue` predicate encodes exactly that decision.
|
|
||||||
|
|
||||||
## Sub-pages
|
## Sub-pages
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -80,7 +80,7 @@ shapes and routing logic in
|
||||||
|
|
||||||
**Inbox** (`inbox` group): `get_loose_ends(agent?)`,
|
**Inbox** (`inbox` group): `get_loose_ends(agent?)`,
|
||||||
`cancel_loose_end(kind, id)`, `remind(message, delay_seconds? |
|
`cancel_loose_end(kind, id)`, `remind(message, delay_seconds? |
|
||||||
at_unix_timestamp?)`, `request_next_turn()`.
|
at_unix_timestamp?)`.
|
||||||
|
|
||||||
- `get_loose_ends(agent?)` — list pending questions (asked/owed),
|
- `get_loose_ends(agent?)` — list pending questions (asked/owed),
|
||||||
scheduled reminders, and active local tasks published by external MCP
|
scheduled reminders, and active local tasks published by external MCP
|
||||||
|
|
@ -97,9 +97,15 @@ at_unix_timestamp?)`, `request_next_turn()`.
|
||||||
- `remind` — schedule a reminder in this agent's own inbox. Large
|
- `remind` — schedule a reminder in this agent's own inbox. Large
|
||||||
payloads spill to `/agents/<self>/state/reminders/`. Pending count
|
payloads spill to `/agents/<self>/state/reminders/`. Pending count
|
||||||
capped at 50 per agent (`HIVE_REMIND_MAX_PENDING_PER_AGENT`).
|
capped at 50 per agent (`HIVE_REMIND_MAX_PENDING_PER_AGENT`).
|
||||||
- `request_next_turn` — ask the harness to start another turn
|
|
||||||
immediately after this one ends, even if the inbox is empty.
|
There is no same-turn self-continue tool (`request_next_turn` was
|
||||||
Next turn fires with `from: "self"` and `body: "continue"`.
|
removed — unanimous consensus across every agent that used the harness
|
||||||
|
that ending the turn and letting an external wake drive the next one
|
||||||
|
is strictly better: it checkpoints the session and observes wakes that
|
||||||
|
only reach the harness between turns, forge #2777). Multi-step work
|
||||||
|
rides `remind` for a durable self-wake, or an in-container todo wake
|
||||||
|
(bash-task completion, forge notification, matrix activity) for
|
||||||
|
work already in flight.
|
||||||
|
|
||||||
**Meta** (`meta` group): `set_status(text)`, `get_agent_meta(name?)`.
|
**Meta** (`meta` group): `set_status(text)`, `get_agent_meta(name?)`.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -626,30 +626,6 @@ impl AgentServer {
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tool(
|
|
||||||
description = "Ask the harness to start another turn immediately after this one \
|
|
||||||
completes, even if the inbox is empty. Use this when you have ongoing work that \
|
|
||||||
spans multiple turns (long builds, multi-step tasks) and you want to continue \
|
|
||||||
without waiting for an external message. The next turn will start with \
|
|
||||||
`from: \"self\"` and `body: \"continue\"`. Has no effect if a new inbox message \
|
|
||||||
arrives before this turn ends — the harness already loops immediately on pending \
|
|
||||||
messages. No args."
|
|
||||||
)]
|
|
||||||
async fn request_next_turn(&self) -> String {
|
|
||||||
run_tool_envelope("request_next_turn", String::new(), async move {
|
|
||||||
let sentinel = crate::paths::state_dir().join("hyperhive-continue");
|
|
||||||
match std::fs::write(&sentinel, b"") {
|
|
||||||
Ok(()) => "ok — harness will start another turn immediately after this one",
|
|
||||||
Err(e) => {
|
|
||||||
tracing::warn!(error = %e, path = %sentinel.display(), "request_next_turn: write failed");
|
|
||||||
return format!("request_next_turn failed: {e}");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
.to_string()
|
|
||||||
})
|
|
||||||
.await
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tool(
|
#[tool(
|
||||||
description = "Compact the current session's context, mirroring the operator's \
|
description = "Compact the current session's context, mirroring the operator's \
|
||||||
dashboard `/compact` button. Gated: only honoured when this agent's last \
|
dashboard `/compact` button. Gated: only honoured when this agent's last \
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,7 @@ You are hyperhive agent `{label}` (qualified: `{qualified_label}`){hive_identity
|
||||||
|
|
||||||
Tools (hyperhive surface). Full signature + behavior for each comes from the tool's own MCP description (you already received it via the MCP tool schema) — this is just the map of what exists and which ones are gated, so you know where to look:
|
Tools (hyperhive surface). Full signature + behavior for each comes from the tool's own MCP description (you already received it via the MCP tool schema) — this is just the map of what exists and which ones are gated, so you know where to look:
|
||||||
|
|
||||||
- **Inbox / messaging** (always available): `mcp__hyperhive__recv`, `mcp__hyperhive__ack_until`, `mcp__hyperhive__send`, `mcp__hyperhive__ask`, `mcp__hyperhive__answer`, `mcp__hyperhive__get_loose_ends`, `mcp__hyperhive__cancel_loose_end`, `mcp__hyperhive__remind`, `mcp__hyperhive__set_status`, `mcp__hyperhive__get_agent_meta`, `mcp__hyperhive__request_next_turn`. Two habits worth internalizing beyond the tool descriptions themselves: prefer ending the turn over parking in `recv` when idle (only turn-boundaries observe in-container todo wakes — bash-task completions, matrix unread, forge activity — and ending the turn is also your checkpoint); and `ask`/`answer` are async — `ask` returns immediately with a question id, the reply lands later as a `question_answered` system event, never block a turn waiting on it inline.
|
- **Inbox / messaging** (always available): `mcp__hyperhive__recv`, `mcp__hyperhive__ack_until`, `mcp__hyperhive__send`, `mcp__hyperhive__ask`, `mcp__hyperhive__answer`, `mcp__hyperhive__get_loose_ends`, `mcp__hyperhive__cancel_loose_end`, `mcp__hyperhive__remind`, `mcp__hyperhive__set_status`, `mcp__hyperhive__get_agent_meta`. Two habits worth internalizing beyond the tool descriptions themselves: prefer ending the turn over parking in `recv` when idle (only turn-boundaries observe in-container todo wakes — bash-task completions, matrix unread, forge activity — and ending the turn is also your checkpoint); and `ask`/`answer` are async — `ask` returns immediately with a question id, the reply lands later as a `question_answered` system event, never block a turn waiting on it inline.
|
||||||
- **Extra MCP tools** (some agents only): `mcp__<server>__<tool>` — agent-specific (matrix client, scraper, db connector, etc.) declared in your `agent.nix` under `hyperhive.extraMcpServers`. First-class tools, already operator-approved at deploy time.
|
- **Extra MCP tools** (some agents only): `mcp__<server>__<tool>` — agent-specific (matrix client, scraper, db connector, etc.) declared in your `agent.nix` under `hyperhive.extraMcpServers`. First-class tools, already operator-approved at deploy time.
|
||||||
- **Lifecycle** (_requires `lifecycle` tool group_, direct children only, no approval needed): `restart`, `kill`, `start`, `update`, `list_containers`.
|
- **Lifecycle** (_requires `lifecycle` tool group_, direct children only, no approval needed): `restart`, `kill`, `start`, `update`, `list_containers`.
|
||||||
- **Approvals** (_requires `approvals` tool group_, queues an operator approval): `request_init_config`, `request_apply_commit`, `request_update_meta_inputs`.
|
- **Approvals** (_requires `approvals` tool group_, queues an operator approval): `request_init_config`, `request_apply_commit`, `request_update_meta_inputs`.
|
||||||
|
|
@ -42,6 +42,6 @@ Keep messages short — a few sentences each. For anything big (file listings, l
|
||||||
|
|
||||||
When your inbox has a message, handle it and stop. Don't narrate intent — act.
|
When your inbox has a message, handle it and stop. Don't narrate intent — act.
|
||||||
|
|
||||||
**Turns are your checkpoint.** The harness runs one claude turn per inbox message; when you stop, it acknowledges that message and your `--continue` session is saved to disk. Ending the turn is how you commit progress — both the session and the inbox acknowledgement. If the container restarts _while a turn is still running_, the message that drove it was never acknowledged, so it gets redelivered on the next boot, prefixed `[redelivered after harness restart — may already be handled]`. A long single turn that does step after step widens the window where a restart loses work and forces that redelivery, so prefer short turns: do a unit of work, write anything durable under `/agents/{label}/state/`, and end.
|
**Turns are your checkpoint.** The harness runs one claude turn per inbox message; when you stop, it acknowledges that message and your `--continue` session is saved to disk. Ending the turn is how you commit progress — both the session and the inbox acknowledgement. If the container restarts _while a turn is still running_, the message that drove it was never acknowledged, so it gets redelivered on the next boot, prefixed `[redelivered after harness restart — may already be handled]`. A long single turn that does step after step widens the window where a restart loses work and forces that redelivery, so prefer short turns: do a unit of work, write anything durable under `/agents/{label}/state/` **at every natural boundary, not just when a turn happens to end**, and stop.
|
||||||
|
|
||||||
**To keep working without waiting for a new message, call `request_next_turn()`** before you stop. The harness immediately starts a fresh turn with `from: "self"`, `body: "continue"` — the supported way to run multi-step work (long builds, sequential edits) as a series of checkpointed turns rather than one monolithic turn. Don't busy-wait inside a turn for a condition to resolve: end the turn and let the next wake drive the continuation — a `remind` you scheduled, an external event, a backgrounded bash task's completion, or `request_next_turn()`. (A long-poll `recv(wait_seconds: …)` blocks _within_ the current turn — it parks for new inbox messages but does not end the turn or checkpoint, so it isn't a substitute for ending the turn.)
|
For multi-step work (long builds, sequential edits) that spans more than one turn: end the turn and let an external wake drive the next one — a new inbox message, a `remind` you scheduled, or a backgrounded bash task's completion. Ending the turn is never a same-turn continuation: it's the checkpoint itself, and the only place an in-container todo wake (bash-task completion, matrix unread, forge activity) can reach you — a long-poll `recv(wait_seconds: …)` blocks _within_ the current turn instead and misses exactly those wakes, so it isn't a substitute for ending the turn. **Keep `set_status` current whenever the work changes** — a stale status is what actually loses context across a restart, not turn length; a status set at the start of a task and never touched again is a bug, not a shortcut.
|
||||||
|
|
|
||||||
|
|
@ -188,72 +188,19 @@ fn format_turn_failure(err: &anyhow::Error) -> String {
|
||||||
format!("[system] `{who}` claude turn failed:\n{err:#}")
|
format!("[system] `{who}` claude turn failed:\n{err:#}")
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Check for the `hyperhive-continue` sentinel under the state dir
|
/// What a finished turn tells the serve loop to do next.
|
||||||
/// (dropped by the `request_next_turn` MCP tool). Returns true and
|
|
||||||
/// consumes the file when present; false otherwise. Caller fires
|
|
||||||
/// the role-specific `Wake` request — the sentinel itself is wire-
|
|
||||||
/// agnostic so this helper lives outside both surfaces.
|
|
||||||
fn consume_continue_sentinel() -> bool {
|
|
||||||
let sentinel = crate::paths::state_dir().join("hyperhive-continue");
|
|
||||||
if !sentinel.exists() {
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
if let Err(e) = std::fs::remove_file(&sentinel) {
|
|
||||||
tracing::warn!(error = %e, "consume_continue_sentinel: remove sentinel failed");
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
true
|
|
||||||
}
|
|
||||||
|
|
||||||
/// What a finished turn tells the serve loop to do next. Replaces the
|
|
||||||
/// bare `auth_failed` bool so the loop can also act on a pending
|
|
||||||
/// `request_next_turn` without round-tripping a synthetic message
|
|
||||||
/// through the broker.
|
|
||||||
struct TurnControl {
|
struct TurnControl {
|
||||||
/// The turn ended in `AuthFailed` — caller parks on login.
|
/// The turn ended in `AuthFailed` — caller parks on login.
|
||||||
auth_failed: bool,
|
auth_failed: bool,
|
||||||
/// `request_next_turn` was called during the turn (the
|
|
||||||
/// `hyperhive-continue` sentinel was dropped + consumed).
|
|
||||||
continue_requested: bool,
|
|
||||||
/// Inbox unread count observed right after the turn. Used to
|
|
||||||
/// decide whether a self-continue is actually needed.
|
|
||||||
pending: u64,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Decide whether the serve loop should drive a self-continue turn
|
|
||||||
/// in-process. A continue is only "needed" when nothing else will
|
|
||||||
/// wake the agent: if real messages are already pending they drive
|
|
||||||
/// the next turn(s) and the continue is dropped (matches the
|
|
||||||
/// `request_next_turn` contract — "no effect if a new inbox message
|
|
||||||
/// arrives before this turn ends"). Auth-failed parks the loop on
|
|
||||||
/// login, so it suppresses the continue too.
|
|
||||||
fn should_self_continue(ctrl: &TurnControl) -> bool {
|
|
||||||
ctrl.continue_requested && !ctrl.auth_failed && ctrl.pending == 0
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Synthesize the `from: "self"` / `body: "continue"` message that a
|
|
||||||
/// `request_next_turn` self-continue drives. Built in-process rather
|
|
||||||
/// than fetched from the broker — it never touches the send/recv
|
|
||||||
/// path, so it doesn't persist to sqlite or pollute the inbox.
|
|
||||||
/// `id = 0` is a non-broker sentinel: the synthetic message
|
|
||||||
/// has no DB row, and `AckTurn` keys off the recipient's in-flight
|
|
||||||
/// list (which is empty here) rather than this id.
|
|
||||||
fn synthetic_continue() -> hive_sh4re::DeliveredMessage {
|
|
||||||
hive_sh4re::DeliveredMessage {
|
|
||||||
from: "self".into(),
|
|
||||||
body: "continue".into(),
|
|
||||||
id: 0,
|
|
||||||
redelivered: false,
|
|
||||||
in_reply_to: None,
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Synthesize the message that drives a turn when an in-container producer
|
/// Synthesize the message that drives a turn when an in-container producer
|
||||||
/// upserted a new/changed *todo* over the in-agent socket (loose-ends v2).
|
/// upserted a new/changed *todo* over the in-agent socket (loose-ends v2).
|
||||||
/// The harness owns the todo store locally and signals the serve loop
|
/// The harness owns the todo store locally and signals the serve loop
|
||||||
/// directly — so this wake never touches the broker (no long-poll, no
|
/// directly — so this wake never touches the broker (no long-poll, no
|
||||||
/// marker file). `id = 0` is the same non-broker sentinel as
|
/// marker file). `id = 0` is a non-broker sentinel: the synthetic message
|
||||||
/// [`synthetic_continue`].
|
/// has no DB row, and `AckTurn` keys off the recipient's in-flight list
|
||||||
|
/// (which is empty here) rather than this id.
|
||||||
fn synthetic_todo_message() -> hive_sh4re::DeliveredMessage {
|
fn synthetic_todo_message() -> hive_sh4re::DeliveredMessage {
|
||||||
hive_sh4re::DeliveredMessage {
|
hive_sh4re::DeliveredMessage {
|
||||||
from: "todo".into(),
|
from: "todo".into(),
|
||||||
|
|
@ -678,20 +625,12 @@ async fn serve_loop<S: Surface>(
|
||||||
// The durable claude session, built once and reused for every turn +
|
// The durable claude session, built once and reused for every turn +
|
||||||
// idle compaction below (it's effectively stateless).
|
// idle compaction below (it's effectively stateless).
|
||||||
let session = turn::make_session(&bus);
|
let session = turn::make_session(&bus);
|
||||||
// Set when a turn calls `request_next_turn` and no real work is
|
|
||||||
// pending — the next iteration drives this synthetic message
|
|
||||||
// in-process instead of long-polling the broker. Never
|
|
||||||
// persisted: it lives entirely in this loop's stack.
|
|
||||||
let mut self_continue: Option<hive_sh4re::DeliveredMessage> = None;
|
|
||||||
// Tracks the last observed pause state so the transitions get logged
|
// Tracks the last observed pause state so the transitions get logged
|
||||||
// once each instead of twelve lines a minute while parked.
|
// once each instead of twelve lines a minute while parked.
|
||||||
let mut was_paused = false;
|
let mut was_paused = false;
|
||||||
loop {
|
loop {
|
||||||
// Pause gate. While the marker is present this loop drives no
|
// Pause gate. While the marker is present this loop drives no
|
||||||
// turns at all — deliberately *before* the `self_continue.take()`
|
// turns at all.
|
||||||
// below, so a `request_next_turn` that raced the pause is still
|
|
||||||
// waiting when the agent resumes rather than being consumed by
|
|
||||||
// a turn that never runs.
|
|
||||||
//
|
//
|
||||||
// Nothing here touches the broker: not calling `S::recv_next` is
|
// Nothing here touches the broker: not calling `S::recv_next` is
|
||||||
// exactly the "messages queue unacked, resume drains the
|
// exactly the "messages queue unacked, resume drains the
|
||||||
|
|
@ -722,78 +661,75 @@ async fn serve_loop<S: Surface>(
|
||||||
});
|
});
|
||||||
was_paused = false;
|
was_paused = false;
|
||||||
}
|
}
|
||||||
let next = match self_continue.take() {
|
let next = match {
|
||||||
Some(msg) => msg,
|
// Idle wait: race the broker long-poll against a local
|
||||||
None => match {
|
// todo signal so an in-container producer's upsert drives a
|
||||||
// Idle wait: race the broker long-poll against a local
|
// turn without any broker round-trip. `biased` polls the
|
||||||
// todo signal so an in-container producer's upsert drives a
|
// broker recv first, so a genuinely-ready inbox message is
|
||||||
// turn without any broker round-trip. `biased` polls the
|
// never dropped in favour of the todo wake.
|
||||||
// broker recv first, so a genuinely-ready inbox message is
|
tokio::select! {
|
||||||
// never dropped in favour of the todo wake.
|
biased;
|
||||||
tokio::select! {
|
o = S::recv_next(socket) => o,
|
||||||
biased;
|
() = todo_wake.notified() => RecvOutcome::LocalTodo,
|
||||||
o = S::recv_next(socket) => o,
|
Some(dm) = reminder_rx.recv() => RecvOutcome::Message(dm),
|
||||||
() = todo_wake.notified() => RecvOutcome::LocalTodo,
|
}
|
||||||
Some(dm) = reminder_rx.recv() => RecvOutcome::Message(dm),
|
} {
|
||||||
}
|
RecvOutcome::Message(first) => first,
|
||||||
} {
|
RecvOutcome::LocalTodo => {
|
||||||
RecvOutcome::Message(first) => first,
|
// Gate on `has_any()` before spawning a turn: a burst of
|
||||||
RecvOutcome::LocalTodo => {
|
// same-turn upserts can arm a second `Notify` permit that
|
||||||
// Gate on `has_any()` before spawning a turn: a burst of
|
// outlives the turn which already drained its payload
|
||||||
// same-turn upserts can arm a second `Notify` permit that
|
// (the phantom-todo-wake issue) — `notify_one` doesn't
|
||||||
// outlives the turn which already drained its payload
|
// coalesce once the first permit's been consumed, so the surplus wake
|
||||||
// (the phantom-todo-wake issue) — `notify_one` doesn't
|
// fires the instant the loop is back here even though
|
||||||
// coalesce once the first permit's been consumed, so the surplus wake
|
// there's nothing left to show. Fail open (drive a turn
|
||||||
// fires the instant the loop is back here even though
|
// anyway) on a `has_any` error so a flaky sqlite read
|
||||||
// there's nothing left to show. Fail open (drive a turn
|
// never silently swallows a real wake.
|
||||||
// anyway) on a `has_any` error so a flaky sqlite read
|
let has_any = todos_store
|
||||||
// never silently swallows a real wake.
|
.as_ref()
|
||||||
let has_any = todos_store
|
.is_none_or(|store| store.has_any().unwrap_or(true));
|
||||||
.as_ref()
|
if !has_any {
|
||||||
.is_none_or(|store| store.has_any().unwrap_or(true));
|
tracing::debug!("todo wake fired against an empty store — stale, skipping");
|
||||||
if !has_any {
|
|
||||||
tracing::debug!("todo wake fired against an empty store — stale, skipping");
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
tracing::debug!("todo wake consumed, sending synthetic todo message");
|
|
||||||
synthetic_todo_message()
|
|
||||||
}
|
|
||||||
RecvOutcome::Empty => {
|
|
||||||
// Idle: no message this poll. Service a queued operator
|
|
||||||
// `/compact` here so it runs even when no turn is driving
|
|
||||||
// (the in-flight case is handled at the end of drive_turn).
|
|
||||||
let compacted = turn::run_pending_compact(files, &bus, &session).await;
|
|
||||||
if !compacted {
|
|
||||||
tokio::time::sleep(interval).await;
|
|
||||||
}
|
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
RecvOutcome::TransportError => {
|
tracing::debug!("todo wake consumed, sending synthetic todo message");
|
||||||
// `recv_next` already logged the detail; just retry.
|
synthetic_todo_message()
|
||||||
// No backoff: the long-poll wait is itself the throttle.
|
}
|
||||||
continue;
|
RecvOutcome::Empty => {
|
||||||
|
// Idle: no message this poll. Service a queued operator
|
||||||
|
// `/compact` here so it runs even when no turn is driving
|
||||||
|
// (the in-flight case is handled at the end of drive_turn).
|
||||||
|
let compacted = turn::run_pending_compact(files, &bus, &session).await;
|
||||||
|
if !compacted {
|
||||||
|
tokio::time::sleep(interval).await;
|
||||||
}
|
}
|
||||||
RecvOutcome::GracefulStop => {
|
continue;
|
||||||
// c0re fenced our inbox and wants a clean stop. Run one
|
}
|
||||||
// checkpoint turn so the agent flushes durable /state,
|
RecvOutcome::TransportError => {
|
||||||
// report completion, then exit the loop → the harness
|
// `recv_next` already logged the detail; just retry.
|
||||||
// process ends and the container can be stopped.
|
// No backoff: the long-poll wait is itself the throttle.
|
||||||
tracing::info!(
|
continue;
|
||||||
"graceful stop signalled — running stop-checkpoint turn, then exiting"
|
}
|
||||||
);
|
RecvOutcome::GracefulStop => {
|
||||||
let _ = handle_turn::<S>(
|
// c0re fenced our inbox and wants a clean stop. Run one
|
||||||
socket,
|
// checkpoint turn so the agent flushes durable /state,
|
||||||
&bus,
|
// report completion, then exit the loop → the harness
|
||||||
stats.as_ref(),
|
// process ends and the container can be stopped.
|
||||||
files,
|
tracing::info!(
|
||||||
&session,
|
"graceful stop signalled — running stop-checkpoint turn, then exiting"
|
||||||
graceful_stop_message(),
|
);
|
||||||
)
|
let _ = handle_turn::<S>(
|
||||||
.await;
|
socket,
|
||||||
S::graceful_stop_complete(socket).await;
|
&bus,
|
||||||
return Ok(());
|
stats.as_ref(),
|
||||||
}
|
files,
|
||||||
},
|
&session,
|
||||||
|
graceful_stop_message(),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
S::graceful_stop_complete(socket).await;
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
};
|
};
|
||||||
let ctrl = handle_turn::<S>(socket, &bus, stats.as_ref(), files, &session, next).await;
|
let ctrl = handle_turn::<S>(socket, &bus, stats.as_ref(), files, &session, next).await;
|
||||||
if ctrl.auth_failed {
|
if ctrl.auth_failed {
|
||||||
|
|
@ -805,19 +741,14 @@ async fn serve_loop<S: Surface>(
|
||||||
u64::try_from(interval.as_millis()).unwrap_or(2000),
|
u64::try_from(interval.as_millis()).unwrap_or(2000),
|
||||||
)
|
)
|
||||||
.await;
|
.await;
|
||||||
} else if should_self_continue(&ctrl) {
|
|
||||||
tracing::info!("request_next_turn: driving self-continue turn in-process");
|
|
||||||
self_continue = Some(synthetic_continue());
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Drive a single turn: emit boot-of-turn events, run claude, ack on
|
/// Drive a single turn: emit boot-of-turn events, run claude, ack on
|
||||||
/// success / requeue on rate-limit-or-401 / notify parent on failure,
|
/// success / requeue on rate-limit-or-401 / notify parent on failure,
|
||||||
/// record stats, then pick up the `request_next_turn` sentinel if it's
|
/// record stats. Returns a `TurnControl` carrying the auth-failed flag —
|
||||||
/// been dropped during the turn. Returns a `TurnControl` carrying the
|
/// the serve loop decides what to do next.
|
||||||
/// auth-failed flag, whether a self-continue was requested, and the
|
|
||||||
/// post-turn inbox count — the serve loop decides what to do next.
|
|
||||||
async fn handle_turn<S: Surface>(
|
async fn handle_turn<S: Surface>(
|
||||||
socket: &Path,
|
socket: &Path,
|
||||||
bus: &Bus,
|
bus: &Bus,
|
||||||
|
|
@ -933,54 +864,5 @@ async fn handle_turn<S: Surface>(
|
||||||
}
|
}
|
||||||
TurnControl {
|
TurnControl {
|
||||||
auth_failed: matches!(outcome, Err(turn::TurnError::AuthFailed)),
|
auth_failed: matches!(outcome, Err(turn::TurnError::AuthFailed)),
|
||||||
continue_requested: consume_continue_sentinel(),
|
|
||||||
pending,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod continue_tests {
|
|
||||||
use super::{TurnControl, should_self_continue, synthetic_continue};
|
|
||||||
|
|
||||||
fn ctrl(auth_failed: bool, continue_requested: bool, pending: u64) -> TurnControl {
|
|
||||||
TurnControl {
|
|
||||||
auth_failed,
|
|
||||||
continue_requested,
|
|
||||||
pending,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn self_continue_when_requested_and_inbox_empty() {
|
|
||||||
assert!(should_self_continue(&ctrl(false, true, 0)));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn no_self_continue_when_not_requested() {
|
|
||||||
assert!(!should_self_continue(&ctrl(false, false, 0)));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn no_self_continue_when_real_messages_pending() {
|
|
||||||
// A real message will drive the next turn via recv — the
|
|
||||||
// continue is superseded, not needed (request_next_turn contract).
|
|
||||||
assert!(!should_self_continue(&ctrl(false, true, 3)));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn no_self_continue_when_auth_failed() {
|
|
||||||
// Auth-failed parks the loop on login; a queued continue must
|
|
||||||
// not jump the gate.
|
|
||||||
assert!(!should_self_continue(&ctrl(true, true, 0)));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn synthetic_continue_shape() {
|
|
||||||
let m = synthetic_continue();
|
|
||||||
assert_eq!(m.from, "self");
|
|
||||||
assert_eq!(m.body, "continue");
|
|
||||||
assert_eq!(m.id, 0);
|
|
||||||
assert!(!m.redelivered);
|
|
||||||
assert!(m.in_reply_to.is_none());
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -362,7 +362,6 @@ fn tool_icon(name: &str) -> &'static str {
|
||||||
"mcp__hyperhive__cancel_loose_end" => "✂️",
|
"mcp__hyperhive__cancel_loose_end" => "✂️",
|
||||||
"mcp__hyperhive__ack_until" => "✅",
|
"mcp__hyperhive__ack_until" => "✅",
|
||||||
"mcp__hyperhive__get_agent_meta" => "ℹ️",
|
"mcp__hyperhive__get_agent_meta" => "ℹ️",
|
||||||
"mcp__hyperhive__request_next_turn" => "⏩",
|
|
||||||
"mcp__hyperhive__restart" => "↻",
|
"mcp__hyperhive__restart" => "↻",
|
||||||
"mcp__hyperhive__kill" => "⏹️",
|
"mcp__hyperhive__kill" => "⏹️",
|
||||||
"mcp__hyperhive__start" => "▶️",
|
"mcp__hyperhive__start" => "▶️",
|
||||||
|
|
|
||||||
|
|
@ -594,7 +594,7 @@ pub enum ToolGroup {
|
||||||
Messaging,
|
Messaging,
|
||||||
/// `get_agent_meta` (`set_status` is always-on — see `ALWAYS_ON_TOOLS`)
|
/// `get_agent_meta` (`set_status` is always-on — see `ALWAYS_ON_TOOLS`)
|
||||||
Meta,
|
Meta,
|
||||||
/// `get_loose_ends`, `cancel_loose_end`, `remind`, `request_next_turn`
|
/// `get_loose_ends`, `cancel_loose_end`, `remind`
|
||||||
Inbox,
|
Inbox,
|
||||||
/// `kill`, `start`, `restart`, `update` - *(privileged)*
|
/// `kill`, `start`, `restart`, `update` - *(privileged)*
|
||||||
Lifecycle,
|
Lifecycle,
|
||||||
|
|
@ -628,12 +628,7 @@ impl ToolGroup {
|
||||||
match self {
|
match self {
|
||||||
Self::Messaging => &["send", "recv", "ack_until", "ask", "answer"],
|
Self::Messaging => &["send", "recv", "ack_until", "ask", "answer"],
|
||||||
Self::Meta => &["get_agent_meta"],
|
Self::Meta => &["get_agent_meta"],
|
||||||
Self::Inbox => &[
|
Self::Inbox => &["get_loose_ends", "cancel_loose_end", "remind"],
|
||||||
"get_loose_ends",
|
|
||||||
"cancel_loose_end",
|
|
||||||
"remind",
|
|
||||||
"request_next_turn",
|
|
||||||
],
|
|
||||||
Self::Lifecycle => &["kill", "start", "restart", "update", "list_containers"],
|
Self::Lifecycle => &["kill", "start", "restart", "update", "list_containers"],
|
||||||
Self::Approvals => &["request_init_config", "request_update_meta_inputs"],
|
Self::Approvals => &["request_init_config", "request_update_meta_inputs"],
|
||||||
Self::Scheduling => &[
|
Self::Scheduling => &[
|
||||||
|
|
@ -735,9 +730,7 @@ impl ToolGroup {
|
||||||
Self::Meta => {
|
Self::Meta => {
|
||||||
"get_agent_meta — identity introspection (set_status is always available)"
|
"get_agent_meta — identity introspection (set_status is always available)"
|
||||||
}
|
}
|
||||||
Self::Inbox => {
|
Self::Inbox => "get_loose_ends, cancel_loose_end, remind — self-scheduling",
|
||||||
"get_loose_ends, cancel_loose_end, remind, request_next_turn — self-scheduling"
|
|
||||||
}
|
|
||||||
Self::Lifecycle => {
|
Self::Lifecycle => {
|
||||||
"kill, start, restart, update, list_containers — container lifecycle (privileged)"
|
"kill, start, restart, update, list_containers — container lifecycle (privileged)"
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue