Compare commits

..
6 changed files with 161 additions and 16 deletions

View file

@ -2,6 +2,8 @@ use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use hive_ag3nt::web_ui::TurnLock;
use anyhow::Result;
use clap::{Parser, Subcommand};
use hive_ag3nt::events::{Bus, LiveEvent, TurnState};
@ -71,14 +73,16 @@ 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?;
let turn_lock: TurnLock = Arc::new(tokio::sync::Mutex::new(()));
plugins::install_configured(&cli.socket, Some("manager")).await;
tokio::spawn(web_ui::serve(
label,
label.clone(),
port,
login_state.clone(),
bus.clone(),
cli.socket.clone(),
files.clone(),
turn_lock.clone(),
));
match initial {
LoginState::Online => {
@ -88,6 +92,8 @@ async fn main() -> Result<()> {
login_state,
bus,
&files,
turn_lock,
&label,
)
.await
}
@ -102,6 +108,8 @@ async fn main() -> Result<()> {
login_state,
bus,
&files,
turn_lock,
&label,
)
.await
}
@ -136,6 +144,8 @@ async fn serve(
state: Arc<Mutex<LoginState>>,
bus: Bus,
files: &turn::TurnFiles,
turn_lock: TurnLock,
label: &str,
) -> Result<()> {
tracing::info!(socket = %socket.display(), "hive-ag3nt serve");
let _ = state; // reserved for future state transitions (turn-loop -> needs-login)
@ -163,7 +173,10 @@ async fn serve(
});
bus.set_state(TurnState::Thinking);
let prompt = format_wake_prompt(&from, &body, unread);
let outcome = turn::drive_turn(&prompt, files, &bus).await;
let outcome = {
let _guard = turn_lock.lock().await;
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
@ -172,7 +185,7 @@ async fn serve(
// 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;
notify_manager_of_failure(socket, label, e).await;
}
// After turn completes, check if there are pending messages waiting.
@ -225,13 +238,15 @@ fn format_wake_prompt(from: &str, body: &str, unread: u64) -> String {
/// 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:#}");
/// as a system-style event; `label` is included explicitly in the
/// body so the manager can identify the failing agent without having
/// to look at the `from` field (which is broker-stamped and may
/// differ from what the operator sees in the dashboard). 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, label: &str, err: &anyhow::Error) {
let body = format!("[system] agent `{label}` claude turn failed:\n{err:#}");
let res = client::request::<_, AgentResponse>(
socket,
&AgentRequest::Send {

View file

@ -6,6 +6,8 @@ use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use hive_ag3nt::web_ui::TurnLock;
use anyhow::Result;
use clap::{Parser, Subcommand};
use hive_ag3nt::events::{Bus, LiveEvent, TurnState};
@ -61,6 +63,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?;
let turn_lock: TurnLock = Arc::new(tokio::sync::Mutex::new(()));
plugins::install_configured(&cli.socket, None).await;
tokio::spawn(web_ui::serve(
label,
@ -69,14 +72,15 @@ async fn main() -> Result<()> {
bus.clone(),
cli.socket.clone(),
files.clone(),
turn_lock.clone(),
));
match initial {
LoginState::Online => {
serve(&cli.socket, Duration::from_millis(poll_ms), bus, &files).await
serve(&cli.socket, Duration::from_millis(poll_ms), bus, &files, turn_lock).await
}
LoginState::NeedsLogin => {
turn::wait_for_login(&claude_dir, login_state, poll_ms).await;
serve(&cli.socket, Duration::from_millis(poll_ms), bus, &files).await
serve(&cli.socket, Duration::from_millis(poll_ms), bus, &files, turn_lock).await
}
}
}
@ -89,6 +93,7 @@ async fn serve(
interval: Duration,
bus: Bus,
files: &turn::TurnFiles,
turn_lock: TurnLock,
) -> Result<()> {
tracing::info!(socket = %socket.display(), "hive-m1nd serve");
loop {
@ -131,16 +136,30 @@ async fn serve(
});
let prompt = format_wake_prompt(&from, &body, unread);
bus.set_state(TurnState::Thinking);
let outcome = turn::drive_turn(&prompt, files, &bus).await;
let outcome = {
let _guard = turn_lock.lock().await;
turn::drive_turn(&prompt, files, &bus).await
};
turn::emit_turn_end(&bus, &outcome);
bus.set_state(TurnState::Idle);
// Check for messages that arrived during the turn and loop
// immediately if any are waiting — mirrors hive-ag3nt behaviour.
let pending = inbox_unread(socket).await;
if pending > 0 {
tracing::info!(%pending, "pending messages after turn; fetching next");
continue;
}
}
Ok(ManagerResponse::Empty) => {
// Idle: sleep briefly before next long-poll attempt.
tokio::time::sleep(interval).await;
}
Ok(ManagerResponse::Empty) => {}
Ok(
ManagerResponse::Ok
| ManagerResponse::Status { .. }
| ManagerResponse::QuestionQueued { .. }
| ManagerResponse::Recent { .. },
| ManagerResponse::Recent { .. }
| ManagerResponse::Logs { .. },
) => {
tracing::warn!("recv produced unexpected response kind");
}
@ -151,7 +170,6 @@ async fn serve(
tracing::warn!(error = ?e, "recv failed; retrying");
}
}
tokio::time::sleep(interval).await;
}
}

View file

@ -39,6 +39,7 @@ pub enum SocketReply {
Status(u64),
QuestionQueued(i64),
Recent(Vec<hive_sh4re::InboxRow>),
Logs(String),
}
impl From<hive_sh4re::AgentResponse> for SocketReply {
@ -65,6 +66,7 @@ impl From<hive_sh4re::ManagerResponse> for SocketReply {
hive_sh4re::ManagerResponse::Status { unread } => Self::Status(unread),
hive_sh4re::ManagerResponse::QuestionQueued { id } => Self::QuestionQueued(id),
hive_sh4re::ManagerResponse::Recent { rows } => Self::Recent(rows),
hive_sh4re::ManagerResponse::Logs { content } => Self::Logs(content),
}
}
}
@ -351,6 +353,15 @@ pub struct RequestApplyCommitArgs {
pub description: Option<String>,
}
#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
pub struct GetLogsArgs {
/// Logical name of the sub-agent container to fetch logs for.
pub agent: String,
/// How many journal lines to return (default: 50, max: 500).
#[serde(default)]
pub lines: Option<u32>,
}
#[derive(Debug, Clone)]
pub struct ManagerServer {
socket: PathBuf,
@ -580,6 +591,40 @@ impl ManagerServer {
})
.await
}
#[tool(
description = "Fetch recent journal log lines for a sub-agent container. Useful \
for diagnosing MCP server registration failures, startup crashes, plugin install \
errors, or any harness issue you can't see from inside the container. `lines` \
defaults to 50 (max capped at 500 on the host side)."
)]
async fn get_logs(&self, Parameters(args): Parameters<GetLogsArgs>) -> String {
let log = format!("{args:?}");
let agent = args.agent.clone();
run_tool_envelope("get_logs", log, async move {
let lines = args.lines.map(|n| n.min(500));
let (resp, retries) = self
.dispatch(hive_sh4re::ManagerRequest::GetLogs {
agent: agent.clone(),
lines,
})
.await;
let s = match resp {
Ok(SocketReply::Logs(content)) => {
if content.is_empty() {
format!("(no journal output for {agent})")
} else {
content
}
}
Ok(SocketReply::Err(m)) => format!("get_logs failed: {m}"),
Ok(other) => format!("get_logs unexpected response: {other:?}"),
Err(e) => format!("get_logs transport error: {e:#}"),
};
annotate_retries(s, retries)
})
.await
}
}
#[tool_handler(
@ -635,6 +680,7 @@ pub fn allowed_mcp_tools(flavor: Flavor) -> Vec<String> {
"update",
"request_apply_commit",
"ask_operator",
"get_logs",
],
};
let mut out: Vec<String> = names

View file

@ -36,6 +36,12 @@ use crate::turn::TurnFiles;
/// render.
pub type LoginStateCell = Arc<Mutex<LoginState>>;
/// Shared turn lock. The serve loop acquires this (as an async mutex) for the
/// duration of every `drive_turn` call. The `/api/compact` handler tries
/// `try_lock()` and rejects immediately if a turn is in flight, preventing
/// concurrent access to the claude session.
pub type TurnLock = Arc<tokio::sync::Mutex<()>>;
#[derive(Clone)]
struct AppState {
label: String,
@ -48,6 +54,8 @@ struct AppState {
/// settings claude saw on the last regular turn — keeps the
/// session shape identical across compact + normal turns.
files: TurnFiles,
/// Prevents `/api/compact` from racing with an in-flight normal turn.
turn_lock: TurnLock,
}
impl AppState {
@ -70,6 +78,7 @@ pub async fn serve(
bus: Bus,
socket: PathBuf,
files: TurnFiles,
turn_lock: TurnLock,
) -> Result<()> {
let state = AppState {
label,
@ -78,6 +87,7 @@ pub async fn serve(
bus,
socket,
files,
turn_lock,
};
let app = Router::new()
.route("/", get(serve_index))
@ -406,9 +416,21 @@ async fn post_set_model(State(state): State<AppState>, Form(form): Form<ModelFor
}
async fn post_compact(State(state): State<AppState>) -> Response {
// Clone the Arc before locking so the guard's lifetime is tied to the
// clone (which we can move into the spawn) rather than to `state`.
let lock = state.turn_lock.clone();
// Reject immediately if a normal turn is in flight — concurrent access
// to the claude session is unsafe and produces garbled output.
let guard = match lock.try_lock_owned() {
Ok(g) => g,
Err(_) => {
return error_response("turn in flight — wait for it to finish before compacting");
}
};
let bus = state.bus.clone();
let files = state.files.clone();
tokio::spawn(async move {
let _guard = guard; // keep lock alive for the duration of compaction
bus.emit(crate::events::LiveEvent::Note(
"operator: /compact — running on persistent session".into(),
));

View file

@ -273,6 +273,35 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc<Coordinator>) -> ManagerResp
},
}
}
ManagerRequest::GetLogs { agent, lines } => {
let n = lines.unwrap_or(50);
tracing::info!(%agent, %n, "manager: get_logs");
match tokio::process::Command::new("journalctl")
.args([
"-M",
agent,
"-n",
&n.to_string(),
"--no-pager",
"--output=short",
])
.output()
.await
{
Ok(out) => {
let content = if out.status.success() || !out.stdout.is_empty() {
String::from_utf8_lossy(&out.stdout).into_owned()
} else {
let stderr = String::from_utf8_lossy(&out.stderr);
format!("journalctl exited {}: {stderr}", out.status)
};
ManagerResponse::Logs { content }
}
Err(e) => ManagerResponse::Err {
message: format!("journalctl spawn failed: {e:#}"),
},
}
}
ManagerRequest::RequestApplyCommit {
agent,
commit_ref,

View file

@ -475,6 +475,17 @@ pub enum ManagerRequest {
#[serde(default)]
ttl_seconds: Option<u64>,
},
/// Fetch recent journal lines for a sub-agent container. hive-c0re
/// runs `journalctl -M <agent> -n <lines> --no-pager` and returns
/// the output as a string. Useful for diagnosing MCP registration
/// failures, startup crashes, and harness errors.
///
/// `lines` defaults to 50 when omitted.
GetLogs {
agent: String,
#[serde(default)]
lines: Option<u32>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@ -502,4 +513,8 @@ pub enum ManagerResponse {
Recent {
rows: Vec<InboxRow>,
},
/// `GetLogs` result: journal lines for the requested container.
Logs {
content: String,
},
}