Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
772fdd8320 | ||
|
|
3e040d5b16 |
4 changed files with 97 additions and 5 deletions
|
|
@ -71,7 +71,7 @@ async fn main() -> Result<()> {
|
|||
let login_state = Arc::new(Mutex::new(initial));
|
||||
let bus = Bus::new();
|
||||
let files = turn::TurnFiles::prepare(&cli.socket, &label, mcp::Flavor::Agent).await?;
|
||||
plugins::install_configured().await;
|
||||
plugins::install_configured(&cli.socket, Some("manager")).await;
|
||||
tokio::spawn(web_ui::serve(
|
||||
label,
|
||||
port,
|
||||
|
|
@ -166,6 +166,14 @@ async fn serve(
|
|||
let outcome = turn::drive_turn(&prompt, files, &bus).await;
|
||||
turn::emit_turn_end(&bus, &outcome);
|
||||
bus.set_state(TurnState::Idle);
|
||||
// Failures are unhandled by definition — PromptTooLong is
|
||||
// absorbed inside drive_turn via compaction, so anything
|
||||
// that reaches Failed here is a real crash. Notify the
|
||||
// manager so it can investigate / restart / page the
|
||||
// operator; best-effort, swallow the send error.
|
||||
if let turn::TurnOutcome::Failed(e) = &outcome {
|
||||
notify_manager_of_failure(socket, e).await;
|
||||
}
|
||||
|
||||
// After turn completes, check if there are pending messages waiting.
|
||||
// If so, immediately process them instead of blocking on recv().
|
||||
|
|
@ -214,6 +222,29 @@ fn format_wake_prompt(from: &str, body: &str, unread: u64) -> String {
|
|||
format!("Incoming message from `{from}`:\n---\n{body}\n---{pending}")
|
||||
}
|
||||
|
||||
/// Best-effort: tell the manager that this agent's last turn crashed
|
||||
/// (claude exited non-zero, compaction didn't help, etc.). Routed
|
||||
/// through the normal send path so the manager's inbox surfaces it
|
||||
/// like any other message; the agent's label is what the broker
|
||||
/// stamps as `from`, so the message body doesn't need to repeat it.
|
||||
/// Swallows transport errors — we just logged the failure, the worst
|
||||
/// case is the manager learns about the crash from the dashboard
|
||||
/// instead of inbox.
|
||||
async fn notify_manager_of_failure(socket: &Path, err: &anyhow::Error) {
|
||||
let body = format!("claude turn failed:\n{err:#}");
|
||||
let res = client::request::<_, AgentResponse>(
|
||||
socket,
|
||||
&AgentRequest::Send {
|
||||
to: "manager".into(),
|
||||
body,
|
||||
},
|
||||
)
|
||||
.await;
|
||||
if let Err(e) = res {
|
||||
tracing::warn!(error = ?e, "failed to notify manager of turn failure");
|
||||
}
|
||||
}
|
||||
|
||||
/// Best-effort: ask our own per-agent socket how many messages are still
|
||||
/// pending after the wake-up Recv. Returns 0 if anything goes wrong.
|
||||
async fn inbox_unread(socket: &Path) -> u64 {
|
||||
|
|
|
|||
|
|
@ -61,7 +61,7 @@ async fn main() -> Result<()> {
|
|||
let login_state = Arc::new(Mutex::new(initial));
|
||||
let bus = Bus::new();
|
||||
let files = turn::TurnFiles::prepare(&cli.socket, &label, mcp::Flavor::Manager).await?;
|
||||
plugins::install_configured().await;
|
||||
plugins::install_configured(&cli.socket, None).await;
|
||||
tokio::spawn(web_ui::serve(
|
||||
label,
|
||||
port,
|
||||
|
|
|
|||
|
|
@ -7,11 +7,21 @@
|
|||
//! recreate is fine. Failures log a warning but do not abort boot —
|
||||
//! we'd rather start without a plugin than refuse to serve.
|
||||
|
||||
use std::path::Path;
|
||||
|
||||
use tokio::process::Command;
|
||||
|
||||
use crate::client;
|
||||
|
||||
const PLUGINS_PATH: &str = "/etc/hyperhive/claude-plugins.json";
|
||||
|
||||
pub async fn install_configured() {
|
||||
/// Install every plugin in `/etc/hyperhive/claude-plugins.json`. When
|
||||
/// `notify_recipient` is `Some(name)`, install failures also get sent
|
||||
/// as a hyperhive message to that recipient (typically `"manager"` for
|
||||
/// sub-agents) so it surfaces in the inbox rather than being buried in
|
||||
/// journald. The manager itself passes `None` — there's nobody above
|
||||
/// it to notify.
|
||||
pub async fn install_configured(socket: &Path, notify_recipient: Option<&str>) {
|
||||
let raw = match tokio::fs::read_to_string(PLUGINS_PATH).await {
|
||||
Ok(s) => s,
|
||||
Err(_) => return,
|
||||
|
|
@ -33,16 +43,49 @@ pub async fn install_configured() {
|
|||
tracing::info!(spec = %spec, "claude plugin install ok");
|
||||
}
|
||||
Ok(out) => {
|
||||
let stderr = String::from_utf8_lossy(&out.stderr).into_owned();
|
||||
tracing::warn!(
|
||||
spec = %spec,
|
||||
status = ?out.status,
|
||||
stderr = %String::from_utf8_lossy(&out.stderr),
|
||||
stderr = %stderr,
|
||||
"claude plugin install failed",
|
||||
);
|
||||
if let Some(to) = notify_recipient {
|
||||
notify(
|
||||
socket,
|
||||
to,
|
||||
format!(
|
||||
"claude plugin install failed for `{spec}`:\n{}",
|
||||
stderr.trim()
|
||||
),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(spec = %spec, error = ?e, "claude plugin install spawn failed");
|
||||
if let Some(to) = notify_recipient {
|
||||
notify(
|
||||
socket,
|
||||
to,
|
||||
format!("claude plugin install spawn failed for `{spec}`: {e}"),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Best-effort hyperhive send. Swallows transport errors — the warn log
|
||||
/// is already in journald and the harness boot must not stall waiting
|
||||
/// for the broker to be reachable.
|
||||
async fn notify(socket: &Path, to: &str, body: String) {
|
||||
let req = hive_sh4re::AgentRequest::Send {
|
||||
to: to.to_owned(),
|
||||
body,
|
||||
};
|
||||
if let Err(e) = client::request::<_, hive_sh4re::AgentResponse>(socket, &req).await {
|
||||
tracing::warn!(error = ?e, "failed to notify {to} of plugin install failure");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@
|
|||
//! manager surface) and their wake-prompt wording; the spawn shape,
|
||||
//! arg-vector, stdin plumbing, and stream-json pumping are identical.
|
||||
|
||||
use std::collections::VecDeque;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::Stdio;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
|
|
@ -281,6 +282,13 @@ async fn run_claude(prompt: &str, files: &TurnFiles, bus: &Bus) -> Result<bool>
|
|||
}
|
||||
}
|
||||
});
|
||||
// Keep the last STDERR_TAIL_LINES of stderr so a non-zero exit can
|
||||
// include real context in the bail message (and downstream in the
|
||||
// failure notification to the manager) instead of just "exit 1".
|
||||
const STDERR_TAIL_LINES: usize = 20;
|
||||
let stderr_tail: Arc<Mutex<VecDeque<String>>> =
|
||||
Arc::new(Mutex::new(VecDeque::with_capacity(STDERR_TAIL_LINES)));
|
||||
let tail_clone = stderr_tail.clone();
|
||||
let pump_stderr = tokio::spawn(async move {
|
||||
let mut reader = BufReader::new(stderr).lines();
|
||||
while let Ok(Some(line)) = reader.next_line().await {
|
||||
|
|
@ -293,6 +301,11 @@ async fn run_claude(prompt: &str, files: &TurnFiles, bus: &Bus) -> Result<bool>
|
|||
// surfaces when claude exits non-zero.
|
||||
tracing::warn!(line = %line, "claude stderr");
|
||||
bus_err.emit(LiveEvent::Note(format!("stderr: {line}")));
|
||||
let mut t = tail_clone.lock().unwrap();
|
||||
if t.len() >= STDERR_TAIL_LINES {
|
||||
t.pop_front();
|
||||
}
|
||||
t.push_back(line);
|
||||
}
|
||||
});
|
||||
|
||||
|
|
@ -301,7 +314,12 @@ async fn run_claude(prompt: &str, files: &TurnFiles, bus: &Bus) -> Result<bool>
|
|||
let _ = pump_stderr.await;
|
||||
let too_long = prompt_too_long.load(Ordering::Relaxed);
|
||||
if !status.success() && !too_long {
|
||||
bail!("claude exited {status}");
|
||||
let tail = stderr_tail.lock().unwrap();
|
||||
if tail.is_empty() {
|
||||
bail!("claude exited {status} (no stderr)");
|
||||
}
|
||||
let tail_str = tail.iter().cloned().collect::<Vec<_>>().join("\n");
|
||||
bail!("claude exited {status}\nstderr tail:\n{tail_str}");
|
||||
}
|
||||
Ok(too_long)
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue