diff --git a/hive-ag3nt/src/bin/hive-ag3nt.rs b/hive-ag3nt/src/bin/hive-ag3nt.rs index ee77d296..7c661a0f 100644 --- a/hive-ag3nt/src/bin/hive-ag3nt.rs +++ b/hive-ag3nt/src/bin/hive-ag3nt.rs @@ -56,16 +56,24 @@ async fn main() -> Result<()> { let initial = LoginState::from_dir(&claude_dir); tracing::info!(state = ?initial, claude_dir = %claude_dir.display(), "harness boot"); let login_state = Arc::new(Mutex::new(initial)); + let ui_state = login_state.clone(); let bus = Bus::new(); - let files = turn::TurnFiles::prepare(&cli.socket, &label, mcp::Flavor::Agent).await?; - tokio::spawn(web_ui::serve( - label, - port, - login_state.clone(), - bus.clone(), - cli.socket.clone(), - files.clone(), - )); + let ui_bus = bus.clone(); + let ui_socket = cli.socket.clone(); + tokio::spawn(async move { + if let Err(e) = web_ui::serve( + label, + port, + ui_state, + ui_bus, + ui_socket, + web_ui::Flavor::Agent, + ) + .await + { + tracing::error!(error = ?e, "web ui failed"); + } + }); match initial { LoginState::Online => { serve( @@ -73,7 +81,6 @@ async fn main() -> Result<()> { Duration::from_millis(poll_ms), login_state, bus, - &files, ) .await } @@ -87,7 +94,6 @@ async fn main() -> Result<()> { Duration::from_millis(poll_ms), login_state, bus, - &files, ) .await } @@ -102,10 +108,13 @@ async fn serve( interval: Duration, state: Arc>, bus: Bus, - files: &turn::TurnFiles, ) -> Result<()> { tracing::info!(socket = %socket.display(), "hive-ag3nt serve"); let _ = state; // reserved for future state transitions (turn-loop -> needs-login) + let mcp_config = turn::write_mcp_config(socket).await?; + let settings = turn::write_settings(socket).await?; + 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?; loop { let recv: Result = client::request(socket, &AgentRequest::Recv { wait_seconds: None }).await; @@ -120,7 +129,15 @@ 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 = turn::drive_turn( + &prompt, + &mcp_config, + &system_prompt, + &settings, + &bus, + mcp::Flavor::Agent, + ) + .await; turn::emit_turn_end(&bus, &outcome); bus.set_state(TurnState::Idle); } diff --git a/hive-ag3nt/src/bin/hive-m1nd.rs b/hive-ag3nt/src/bin/hive-m1nd.rs index 9e508062..bf997162 100644 --- a/hive-ag3nt/src/bin/hive-m1nd.rs +++ b/hive-ag3nt/src/bin/hive-m1nd.rs @@ -59,23 +59,29 @@ async fn main() -> Result<()> { let initial = LoginState::from_dir(&claude_dir); tracing::info!(state = ?initial, claude_dir = %claude_dir.display(), "hm1nd boot"); let login_state = Arc::new(Mutex::new(initial)); + let ui_state = login_state.clone(); let bus = Bus::new(); - let files = turn::TurnFiles::prepare(&cli.socket, &label, mcp::Flavor::Manager).await?; - tokio::spawn(web_ui::serve( - label, - port, - login_state.clone(), - bus.clone(), - cli.socket.clone(), - files.clone(), - )); - match initial { - LoginState::Online => { - serve(&cli.socket, Duration::from_millis(poll_ms), bus, &files).await + let ui_bus = bus.clone(); + let ui_socket = cli.socket.clone(); + tokio::spawn(async move { + if let Err(e) = web_ui::serve( + label, + port, + ui_state, + ui_bus, + ui_socket, + web_ui::Flavor::Manager, + ) + .await + { + tracing::error!(error = ?e, "web ui failed"); } + }); + match initial { + LoginState::Online => serve(&cli.socket, Duration::from_millis(poll_ms), bus).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).await } } } @@ -83,13 +89,12 @@ async fn main() -> Result<()> { } } -async fn serve( - socket: &Path, - interval: Duration, - bus: Bus, - files: &turn::TurnFiles, -) -> Result<()> { +async fn serve(socket: &Path, interval: Duration, bus: Bus) -> Result<()> { tracing::info!(socket = %socket.display(), "hive-m1nd serve"); + let mcp_config = turn::write_mcp_config(socket).await?; + let settings = turn::write_settings(socket).await?; + 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?; loop { let recv: Result = client::request(socket, &ManagerRequest::Recv { wait_seconds: None }).await; @@ -121,7 +126,15 @@ 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 = turn::drive_turn( + &prompt, + &mcp_config, + &system_prompt, + &settings, + &bus, + mcp::Flavor::Manager, + ) + .await; turn::emit_turn_end(&bus, &outcome); bus.set_state(TurnState::Idle); } diff --git a/hive-ag3nt/src/turn.rs b/hive-ag3nt/src/turn.rs index d85ef8cf..1c49981c 100644 --- a/hive-ag3nt/src/turn.rs +++ b/hive-ag3nt/src/turn.rs @@ -33,34 +33,6 @@ const CLAUDE_SETTINGS: &str = include_str!("../prompts/claude-settings.json"); /// claude exit with a useful error in the live view. const PROMPT_TOO_LONG_MARKER: &str = "Prompt is too long"; -/// The set of files claude reads on every invocation: the MCP server -/// config (`--mcp-config`), static settings (`--settings`), and the -/// pre-rendered role/tools system prompt (`--system-prompt-file`). -/// Materialised once at harness startup; shared between the turn loop -/// and the operator-driven `/compact` path so both invocations look -/// identical to claude (same MCP surface, same allowed tools, same -/// role prompt — only the stdin payload differs). -#[derive(Clone)] -pub struct TurnFiles { - pub mcp_config: PathBuf, - pub settings: PathBuf, - pub system_prompt: PathBuf, - pub flavor: mcp::Flavor, -} - -impl TurnFiles { - /// Write all three files into the per-agent runtime dir alongside - /// `socket`. Idempotent — overwrites whatever was there. - pub async fn prepare(socket: &Path, label: &str, flavor: mcp::Flavor) -> Result { - Ok(Self { - mcp_config: write_mcp_config(socket).await?, - settings: write_settings(socket).await?, - system_prompt: write_system_prompt(socket, label, flavor).await?, - flavor, - }) - } -} - /// Drop the MCP config blob claude reads from `--mcp-config `. /// `socket` is the hyperhive per-container socket (forwarded to the child /// as `--socket `); `binary_subcommand` is e.g. `"mcp"` for sub-agents @@ -127,14 +99,21 @@ pub enum TurnOutcome { /// Drive one turn end-to-end, transparently compacting + retrying once on /// `Prompt is too long`. Both the sub-agent and manager loops call this. -pub async fn drive_turn(prompt: &str, files: &TurnFiles, bus: &Bus) -> TurnOutcome { - match run_turn(prompt, files, bus).await { +pub async fn drive_turn( + prompt: &str, + mcp_config: &Path, + system_prompt: &Path, + settings: &Path, + bus: &Bus, + flavor: mcp::Flavor, +) -> TurnOutcome { + match run_turn(prompt, mcp_config, system_prompt, settings, bus, flavor).await { TurnOutcome::PromptTooLong => { - if let Err(e) = compact_session(files, bus).await { + if let Err(e) = compact_session(settings, bus).await { tracing::warn!(error = %format!("{e:#}"), "compact failed"); return TurnOutcome::Failed(e); } - run_turn(prompt, files, bus).await + run_turn(prompt, mcp_config, system_prompt, settings, bus, flavor).await } other => other, } @@ -187,8 +166,25 @@ pub async fn wait_for_login(claude_dir: &Path, state: Arc>, po /// prompt). The session is persistent across turns via `--continue` and /// claude's in-session auto-compact is disabled via `--settings` so it /// doesn't stall mid-turn — hyperhive owns compaction. -pub async fn run_turn(prompt: &str, files: &TurnFiles, bus: &Bus) -> TurnOutcome { - match run_claude(prompt, files, bus).await { +pub async fn run_turn( + prompt: &str, + mcp_config: &Path, + system_prompt: &Path, + settings: &Path, + bus: &Bus, + flavor: mcp::Flavor, +) -> TurnOutcome { + match run_claude( + prompt, + mcp_config, + Some(system_prompt), + settings, + bus, + flavor, + ClaudeMode::Turn, + ) + .await + { Ok(too_long) if too_long => TurnOutcome::PromptTooLong, Ok(_) => TurnOutcome::Ok, Err(e) => TurnOutcome::Failed(e), @@ -196,23 +192,49 @@ pub async fn run_turn(prompt: &str, files: &TurnFiles, bus: &Bus) -> TurnOutcome } /// Run claude's built-in `/compact` slash command on the persistent -/// session. Takes the *same* params as `run_turn` because compact -/// re-initialises claude with the full session shape — same MCP -/// surface, same system prompt, same allowed-tools — so the post- -/// compact state matches a normal turn's. Only the prompt over stdin -/// differs (`/compact` vs the wake-up payload). -pub async fn compact_session(files: &TurnFiles, bus: &Bus) -> Result<()> { +/// session so the next turn can fit. No MCP tools needed; we just feed +/// `/compact` over stdin and let claude rewrite its own history. +pub async fn compact_session(settings: &Path, bus: &Bus) -> Result<()> { bus.emit(LiveEvent::Note( "context overflow — running /compact on the persistent session".into(), )); - let _ = run_claude("/compact", files, bus).await?; + let _ = run_claude( + "/compact", + Path::new("/dev/null"), + None, + settings, + bus, + mcp::Flavor::Agent, // tool surface unused for /compact + ClaudeMode::Compact, + ) + .await?; bus.emit(LiveEvent::Note("/compact done".into())); Ok(()) } -async fn run_claude(prompt: &str, files: &TurnFiles, bus: &Bus) -> Result { +#[derive(Clone, Copy)] +enum ClaudeMode { + Turn, + Compact, +} + +async fn run_claude( + prompt: &str, + mcp_config: &Path, + system_prompt: Option<&Path>, + settings: &Path, + bus: &Bus, + flavor: mcp::Flavor, + mode: ClaudeMode, +) -> Result { let model = bus.model(); - let resume = !bus.take_skip_continue(); + // /compact must always run against the existing session — otherwise + // there's nothing to compact. Only normal turns honor the + // operator's "new session" one-shot flag. + let resume = match mode { + ClaudeMode::Turn => !bus.take_skip_continue(), + ClaudeMode::Compact => true, + }; if !resume { bus.emit(LiveEvent::Note( "fresh session (--continue suppressed for this turn)".into(), @@ -236,18 +258,22 @@ async fn run_claude(prompt: &str, files: &TurnFiles, bus: &Bus) -> Result .arg("--model") .arg(&model) .arg("--settings") - .arg(&files.settings); + .arg(settings); if resume { cmd.arg("--continue"); } - cmd.arg("--system-prompt-file").arg(&files.system_prompt); - cmd.arg("--mcp-config") - .arg(&files.mcp_config) - .arg("--strict-mcp-config") - .arg("--tools") - .arg(mcp::builtin_tools_arg()) - .arg("--allowedTools") - .arg(mcp::allowed_tools_arg(files.flavor)); + if let Some(p) = system_prompt { + cmd.arg("--system-prompt-file").arg(p); + } + if let ClaudeMode::Turn = mode { + cmd.arg("--mcp-config") + .arg(mcp_config) + .arg("--strict-mcp-config") + .arg("--tools") + .arg(mcp::builtin_tools_arg()) + .arg("--allowedTools") + .arg(mcp::allowed_tools_arg(flavor)); + } let mut child = cmd .stdin(Stdio::piped()) .stdout(Stdio::piped()) diff --git a/hive-ag3nt/src/web_ui.rs b/hive-ag3nt/src/web_ui.rs index 931ec2d6..c9d759e5 100644 --- a/hive-ag3nt/src/web_ui.rs +++ b/hive-ag3nt/src/web_ui.rs @@ -28,8 +28,6 @@ use crate::client; use crate::events::Bus; use crate::login::LoginState; use crate::login_session::{LoginSession, drop_if_finished}; -use crate::mcp; -use crate::turn::TurnFiles; /// Live login state for the web UI. The harness updates this in place as it /// transitions between `NeedsLogin` and `Online`; the UI reads on each @@ -43,25 +41,16 @@ struct AppState { session: Arc>>>, bus: Bus, socket: PathBuf, - /// Same `TurnFiles` the harness's turn loop uses. Shared so - /// `/api/compact` re-uses the exact MCP config / system prompt / - /// settings claude saw on the last regular turn — keeps the - /// session shape identical across compact + normal turns. - files: TurnFiles, -} - -impl AppState { - fn flavor(&self) -> Flavor { - self.files.flavor - } + flavor: Flavor, } /// Which wire protocol the per-agent UI's `/send` handler should speak. -/// Sub-agent → `AgentRequest::OperatorMsg`; manager → -/// `ManagerRequest::OperatorMsg`. Reuses the MCP-side enum so a -/// single value drives both the send protocol and (in -/// `post_compact`) the allowed-tools surface claude sees. -pub type Flavor = mcp::Flavor; +/// Sub-agent → `AgentRequest::OperatorMsg`; manager → `ManagerRequest::OperatorMsg`. +#[derive(Debug, Clone, Copy)] +pub enum Flavor { + Agent, + Manager, +} pub async fn serve( label: String, @@ -69,7 +58,7 @@ pub async fn serve( login: LoginStateCell, bus: Bus, socket: PathBuf, - files: TurnFiles, + flavor: Flavor, ) -> Result<()> { let state = AppState { label, @@ -77,7 +66,7 @@ pub async fn serve( session: Arc::new(Mutex::new(None)), bus, socket, - files, + flavor, }; let app = Router::new() .route("/", get(serve_index)) @@ -219,7 +208,7 @@ async fn api_state(State(state): State) -> axum::Json { .ok() .and_then(|s| s.parse::().ok()) .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 model = state.bus.model(); axum::Json(StateSnapshot { @@ -279,7 +268,7 @@ async fn post_send(State(state): State, Form(form): Form) -> if body.is_empty() { return error_response("send: `body` required"); } - let result = match state.flavor() { + let result = match state.flavor { Flavor::Agent => match client::request::<_, hive_sh4re::AgentResponse>( &state.socket, &hive_sh4re::AgentRequest::OperatorMsg { body }, @@ -407,13 +396,22 @@ async fn post_set_model(State(state): State, Form(form): Form) -> Response { let bus = state.bus.clone(); - let files = state.files.clone(); + let socket = state.socket.clone(); tokio::spawn(async move { bus.emit(crate::events::LiveEvent::Note( "operator: /compact — running on persistent session".into(), )); + let settings = match crate::turn::write_settings(&socket).await { + Ok(p) => p, + Err(e) => { + bus.emit(crate::events::LiveEvent::Note(format!( + "/compact failed: settings write — {e:#}" + ))); + return; + } + }; bus.set_state(crate::events::TurnState::Compacting); - let r = crate::turn::compact_session(&files, &bus).await; + let r = crate::turn::compact_session(&settings, &bus).await; bus.set_state(crate::events::TurnState::Idle); if let Err(e) = r { bus.emit(crate::events::LiveEvent::Note(format!( diff --git a/hive-c0re/src/meta.rs b/hive-c0re/src/meta.rs index afb30b9a..023ded47 100644 --- a/hive-c0re/src/meta.rs +++ b/hive-c0re/src/meta.rs @@ -73,18 +73,8 @@ pub async fn sync_agents( if initial { git(&dir, &["init", "--initial-branch=main"]).await?; } - // Stage flake.nix *before* running nix flake lock. When meta is - // a git repo, nix treats it as a `git+file://` self-reference; - // its dirty-tree fetcher includes index entries (tracked + - // staged) but skips untracked files, so without the stage step - // an untracked flake.nix surfaces as "source tree does not - // contain '/flake.nix'". Lock then commit once with both - // flake.nix and flake.lock — single commit per change. - git(&dir, &["add", "flake.nix"]).await?; nix(&dir, &["flake", "lock"]).await?; - if std::path::Path::new(&dir).join("flake.lock").exists() { - git(&dir, &["add", "flake.lock"]).await?; - } + git(&dir, &["add", "-A"]).await?; let msg = if initial { format!("seed meta from {} agent(s)", agents.len()) } else { @@ -96,60 +86,39 @@ pub async fn sync_agents( /// Phase 1 of an apply-commit deploy. Updates the locked rev of /// `agent-` to whatever `applied//main` currently points -/// at and **stages** the lock so `nixos-container update --flake -/// meta#` (which reads via `git+file://`) sees the new rev via -/// the index. Doesn't commit — `finalize_deploy` commits on build -/// success, `abort_deploy` drops the staged change on failure so -/// meta history only carries successful deploys. +/// at. **Doesn't commit** — caller must follow with +/// `finalize_deploy` on build success or `abort_deploy` on failure. #[allow(dead_code)] // wired up by actions::run_apply_commit in a later commit pub async fn prepare_deploy(name: &str) -> Result<()> { let dir = meta_dir(); let input = format!("agent-{name}"); - nix(&dir, &["flake", "update", &input]).await?; - // Stage the new lock — git+file://'s dirty-tree fetcher reads - // index entries, so the upcoming nixos-container update sees the - // bumped rev without a commit yet. - git(&dir, &["add", "flake.lock"]).await + nix(&dir, &["flake", "update", &input]).await } -/// Phase 2-success. Commit the staged lock with the deployed tag + -/// sha as the message. No-op when the rev was already at the right -/// place (nothing staged → nothing to commit). +/// Phase 2-success. Commits the staged `flake.lock` change with a +/// deploy-shaped message. No-op (clean working tree) is tolerated — +/// some lock-updates resolve to the same rev that's already locked. #[allow(dead_code)] pub async fn finalize_deploy(name: &str, sha: &str, tag: &str) -> Result<()> { let dir = meta_dir(); - if !has_staged_changes(&dir).await? { + if git_is_clean(&dir).await? { return Ok(()); } + git(&dir, &["add", "flake.lock"]).await?; let short = &sha[..sha.len().min(12)]; git_commit(&dir, &format!("deploy {name} {tag} {short}")).await } -/// Phase 2-failure. Unstage + restore the lock so meta returns to -/// the previously-committed shas. The failed proposal is still -/// captured in `applied/`'s annotated `failed/` tag. +/// Phase 2-failure. Drops the uncommitted `flake.lock` change so meta +/// stays pinned at the previously-deployed shas. The failed proposal +/// is still captured in `applied/`'s annotated `failed/` tag — +/// meta's history only carries successful deploys. #[allow(dead_code)] pub async fn abort_deploy() -> Result<()> { let dir = meta_dir(); - git(&dir, &["restore", "--staged", "flake.lock"]).await?; git(&dir, &["restore", "flake.lock"]).await } -async fn has_staged_changes(dir: &Path) -> Result { - let st = lifecycle::git_command() - .current_dir(dir) - .args(["diff", "--cached", "--quiet"]) - .status() - .await - .with_context(|| format!("git diff --cached in {}", dir.display()))?; - // exit 1 = differences present, 0 = no diff, other = error - match st.code() { - Some(0) => Ok(false), - Some(1) => Ok(true), - _ => bail!("git diff --cached exited unexpectedly"), - } -} - /// One-shot used by the manual-rebuild path: relock just one /// agent's input and commit the lock change if any. Single-phase /// (no separate finalize) because rebuild has no failure-revert @@ -159,11 +128,11 @@ pub async fn lock_update_for_rebuild(name: &str) -> Result<()> { let dir = meta_dir(); let input = format!("agent-{name}"); nix(&dir, &["flake", "update", &input]).await?; - if git_is_clean(&dir).await? { - return Ok(()); + if !git_is_clean(&dir).await? { + git(&dir, &["add", "flake.lock"]).await?; + git_commit(&dir, &format!("rebuild {name}: lock update")).await?; } - git(&dir, &["add", "flake.lock"]).await?; - git_commit(&dir, &format!("rebuild {name}: lock update")).await + Ok(()) } /// One-shot used by the auto-update path: pin the latest hyperhive @@ -173,11 +142,11 @@ pub async fn lock_update_for_rebuild(name: &str) -> Result<()> { pub async fn lock_update_hyperhive() -> Result<()> { let dir = meta_dir(); nix(&dir, &["flake", "update", "hyperhive"]).await?; - if git_is_clean(&dir).await? { - return Ok(()); + if !git_is_clean(&dir).await? { + git(&dir, &["add", "flake.lock"]).await?; + git_commit(&dir, "bump hyperhive").await?; } - git(&dir, &["add", "flake.lock"]).await?; - git_commit(&dir, "bump hyperhive").await + Ok(()) } fn render_flake(hyperhive_flake: &str, dashboard_port: u16, agents: &[AgentSpec]) -> String {