diff --git a/lib/hive/config.rb b/lib/hive/config.rb index c686876bc..e4236e2d6 100644 --- a/lib/hive/config.rb +++ b/lib/hive/config.rb @@ -324,6 +324,11 @@ module Hive # long codex-backed scan never consumes a task-dispatch slot. "max_concurrent_patrol_scans" => 1, "transient_retry_backoff_sec" => 60, + # Probe-gated auto-retry of the two v1 recoverable terminal-error + # markers (`implementer_failed` Codex-401 + `claude_launch_failed`). + # `enabled` is the global kill-switch; per-reason tuning is deferred + # (the retry limit / backoff are hardcoded in StaleAgentHealer). + "auto_retry" => { "enabled" => true }, "shutdown_grace_sec" => 600, # R-02: per-child wall-clock timeout for daemon-spawned `hive` # children. `child_timeout_sec` is the default cap (0 disables — diff --git a/lib/hive/daemon/dispatcher.rb b/lib/hive/daemon/dispatcher.rb index 1f272c697..dde38583e 100644 --- a/lib/hive/daemon/dispatcher.rb +++ b/lib/hive/daemon/dispatcher.rb @@ -12,6 +12,8 @@ require "hive/daemon/concurrency_controller" require "hive/daemon/child_supervisor" require "hive/daemon/status_consumer" require "hive/daemon/stale_agent_healer" +require "hive/daemon/health_probe" +require "hive/daemon/health_signals" require "hive/daemon/display_name_backfiller" require "hive/daemon/task_id_backfiller" require "hive/daemon/dispatch_request_queue" @@ -100,11 +102,17 @@ module Hive "agent_marker_grace_sec", Hive::TaskAction::DEFAULT_AGENT_MARKER_GRACE_SEC ) - @stale_agent_healer = StaleAgentHealer.new( - controller: @controller, - logger: @logger, - grace_sec: agent_marker_grace_sec - ) + @health_probe = Hive::Daemon::HealthProbe.new(config: @config) + # Per-tick probe memo shared with the healer; cleared each tick so + # parked markers never run the same expensive probe twice per poll + # (mirrors @enabled_cache). + @signal_cache = Hive::Daemon::HealthSignalCache.new + # Baseline SHA-256 of the file that defines SCHEMA_VERSIONS. Captured + # before the healer is built so the probe-gated auto-retry can verify + # the loaded code matches the on-disk lib/hive.rb (the same drift + # signal `run_forever` uses to re-exec a stale daemon). + @code_fingerprint = compute_code_fingerprint + @stale_agent_healer = build_stale_agent_healer(agent_marker_grace_sec) # Additive self-heal for tasks whose one-shot name generation at # `hive new` never landed (agent/codex outage). Re-spawns # `hive generate-name ` on later ticks; never touches @@ -124,16 +132,15 @@ module Hive @shutdown = false @reload = false @reexec_requested = false - # Baseline SHA-256 of the file that defines SCHEMA_VERSIONS. The - # daemon is a long-running process whose in-memory constants - # freeze at load time, while shelled-out `hive` subprocesses load - # fresh code on every invocation. After a `git pull` or gem - # upgrade that bumps a schema, the in-process consumer rejects - # every envelope (e.g. 8946 `got 2, want 1` events were logged - # over ~3 days between 2026-05-15 PR #78 and the next restart). - # Capturing the source digest here lets `run_forever` detect the - # drift and re-exec instead of hard-failing forever. - @code_fingerprint = compute_code_fingerprint + # `@code_fingerprint` (captured above, before the healer build) is the + # loaded-code baseline. The daemon is a long-running process whose + # in-memory constants freeze at load time, while shelled-out `hive` + # subprocesses load fresh code on every invocation; after a `git pull` + # or gem upgrade that bumps a schema, the in-process consumer rejects + # every envelope (e.g. 8946 `got 2, want 1` events were logged over + # ~3 days between 2026-05-15 PR #78 and the next restart). + # `run_forever` compares this baseline against a fresh digest and + # re-execs instead of hard-failing forever. @last_reexec_at = nil @started_at = nil @last_tick_at = nil @@ -189,6 +196,9 @@ module Hive # cache populated on first sight stuck for the daemon's # lifetime and the only way to honour a disable was SIGHUP. @enabled_cache.clear + # Clear the per-tick health-probe memo so probe results recompute on + # the next tick (mirrors the enable cache's one-poll-interval TTL). + @signal_cache.clear @logger.event(:tick_begin, now: now.utc.iso8601) reset_active_agent_snapshot @@ -1818,15 +1828,14 @@ module Hive "child_kill_grace_sec", ChildSupervisor::DEFAULT_KILL_GRACE_SEC ) ) - # Rebuild the healer so an operator tuning - # daemon.agent_marker_grace_sec via SIGHUP takes effect within - # one tick. Without this rebuild the healer keeps the grace it - # captured at boot and only a full daemon restart applies new - # values. - @stale_agent_healer = StaleAgentHealer.new( - controller: @controller, - logger: @logger, - grace_sec: @daemon_cfg.fetch( + # Rebuild the healer (and its probe/kill-switch collaborators) so an + # operator tuning daemon.agent_marker_grace_sec or the auto_retry + # kill-switch via SIGHUP takes effect within one tick. Without this + # rebuild the healer keeps the values it captured at boot. + @health_probe = Hive::Daemon::HealthProbe.new(config: @config) + @signal_cache = Hive::Daemon::HealthSignalCache.new + @stale_agent_healer = build_stale_agent_healer( + @daemon_cfg.fetch( "agent_marker_grace_sec", Hive::TaskAction::DEFAULT_AGENT_MARKER_GRACE_SEC ) @@ -1851,6 +1860,29 @@ module Hive keeping_previous: true) end + # Build the StaleAgentHealer with the probe-gated auto-retry seams + # wired: the health probe, the per-tick signal cache, the kill-switch + # flag, and the loaded-code fingerprint. Per-project config is NOT + # injected here — the healer resolves it through its own + # `Hive::Config.find_project` + `Config.load` fallback (the injectable + # `project_config_resolver` seam is exercised only by tests), which is + # functionally equivalent to the dispatcher's own `project_enabled?` + # resolution. + def build_stale_agent_healer(grace_sec) + auto_retry_enabled = @daemon_cfg.dig("auto_retry", "enabled") != false + StaleAgentHealer.new( + controller: @controller, + logger: @logger, + grace_sec: grace_sec, + probe: @health_probe, + signal_cache: @signal_cache, + auto_retry_enabled: auto_retry_enabled, + state_home: Hive::Paths.state_home, + env: ENV.to_h, + code_fingerprint: @code_fingerprint + ) + end + def install_signal_handlers! Signal.trap("TERM") { @shutdown = true } Signal.trap("INT") { @shutdown = true } diff --git a/lib/hive/daemon/health_probe.rb b/lib/hive/daemon/health_probe.rb new file mode 100644 index 000000000..16fa1df14 --- /dev/null +++ b/lib/hive/daemon/health_probe.rb @@ -0,0 +1,377 @@ +require "digest" +require "open3" +require "stringio" +require "timeout" + +require "hive/agent_profiles" +require "hive/claude_launcher" +require "hive/commands/doctor" + +module Hive + module Daemon + # Bounded, auditable answer to "is this dependency healthy right now?" + # for the StaleAgentHealer's probe-gated auto-retry of the two v1 + # recoverable terminal-error markers: + # + # - `implementer_failed` (Codex auth failure only) → `probe_codex_auth` + # - `claude_launch_failed` → `probe_claude_launcher` + # + # plus a universal `hive doctor` / agent-health precondition that must + # pass before EITHER reason is cleared. + # + # Every probe returns a `Result` — it never raises. A shell-out timeout, + # `Errno::ENOENT`, a non-zero exit (where applicable), or any unexpected + # exception degrades to `healthy: false` with a structured `detail`, so a + # wedged probe is always fail-closed ("not healthy, do not retry"). + # + # Shelling out is isolated behind an injectable `runner:` seam so unit + # tests never fork a real CLI. Production uses `ShellRunner`, which + # wraps `Open3.capture3` in a per-command `Timeout`. + class HealthProbe + # Normalized result of a single probe. `healthy` is the single boolean + # gate; `status` is "healthy"/"unhealthy"; `detail` is the human- and + # audit-log-readable reason; `stdout`/`stderr` carry captured command + # output (empty for in-process probes); `elapsed_sec` bounds the audit + # trail. + Result = Struct.new( + :name, :healthy, :status, :detail, :stdout, :stderr, :elapsed_sec, + keyword_init: true + ) do + def initialize(name:, healthy:, status: nil, detail: nil, stdout: "", stderr: "", elapsed_sec: 0.0) + super + self.healthy = !!healthy + self.status = status.nil? ? (healthy ? "healthy" : "unhealthy") : status.to_s + self.detail = detail.to_s + self.stdout = stdout.to_s + self.stderr = stderr.to_s + self.elapsed_sec = elapsed_sec.to_f + end + + def healthy? + healthy + end + end + + # Normalized single-command result returned by a runner. + CommandResult = Struct.new(:stdout, :stderr, :exitstatus, keyword_init: true) do + def success? + exitstatus == 0 + end + + def output + [ stdout, stderr ].map { |part| part.to_s.strip }.reject(&:empty?).join("\n") + end + end + + # Production runner: Open3.capture3 under a per-command Timeout. A + # missing binary becomes a 127 result; a timeout becomes a 124 result — + # both fail-closed via `success? == false` and never raise. + class ShellRunner + def capture3(argv, env:, timeout_sec:) + out, err, status = Timeout.timeout(timeout_sec) do + Open3.capture3(env.to_h, *argv) + end + CommandResult.new(stdout: out.to_s, stderr: err.to_s, exitstatus: status.exitstatus) + rescue Errno::ENOENT => e + CommandResult.new(stdout: "", stderr: e.message, exitstatus: 127) + rescue Timeout::Error + CommandResult.new(stdout: "", stderr: "timed out after #{timeout_sec}s", exitstatus: 124) + end + end + + CODEX_LOGIN_TIMEOUT_SEC = 15 + CODEX_SMOKE_TIMEOUT_SEC = 20 + # The in-process `hive doctor` precondition can shell out (e.g. tmux + # `-V` via `ClaudeLauncher.tmux_status`, qmd `--version`, skill + # verifications) on unbounded `Open3.capture3` calls, so it gets its + # own wall-clock bound in addition to the reason probes' per-command + # runner timeouts (R5). + UNIVERSAL_PROBE_TIMEOUT_SEC = 30 + CODEX_SMOKE_PROMPT = "reply ok".freeze + + # Doctor rows that mean "this dependency is unusable", shared with the + # universal probe's failing-row check and the stage-skill check. + FAILING_STATUSES = %w[missing version_too_old].freeze + + attr_reader :config, :project_root + + # Canonical on-disk path to the Claude interactive wrapper script. A + # class method so both the probe (default) and the healer (fingerprint + # input) read one source of truth. + def self.wrapper_path + File.expand_path("../scripts/interactive_claude_wrapper.sh", __dir__) + end + + # `doctor_factory:` is an injectable seam for the in-process `hive + # doctor` precondition (defaults to `Hive::Commands::Doctor`). Tests + # inject a fake whose `.new(...)` returns an object responding to + # `#call` and `#rows`, so the universal probe never touches a real + # skill filesystem. `runner:` isolates shell-outs (see ShellRunner). + def initialize(config:, project_root: nil, runner: nil, doctor_factory: nil) + @config = config + @project_root = project_root + @runner = runner || ShellRunner.new + @doctor_factory = doctor_factory || Hive::Commands::Doctor + end + + # Universal precondition: `hive doctor` green plus the retried stage's + # configured skill resolving to a present install. Healthy iff no + # doctor row has status `missing`/`version_too_old` AND the stage's + # profile skill (if any) verifies present. + def probe_universal(config: nil, project_root: nil, stage: nil) + build("universal") do + Timeout.timeout(UNIVERSAL_PROBE_TIMEOUT_SEC) do + cfg = config || @config + root = project_root || @project_root + doctor = @doctor_factory.new( + config: cfg, project_root: root, json: true, output: StringIO.new + ) + doctor.call + if doctor.rows.nil? + next [ false, "doctor config error (no rows produced)" ] + end + failing = doctor.rows.select { |row| FAILING_STATUSES.include?(row[:status].to_s) } + stage_detail = stage_skill_failure(stage, cfg: cfg, project_root: root) + healthy = failing.empty? && stage_detail.nil? + detail = if healthy + "doctor green (no missing/version_too_old rows)" + else + [ failing_desc(failing), stage_detail ].compact.join("; ") + end + [ healthy, detail ] + end + end + end + + # Codex auth probe: the in-process credential check AND `codex login + # status` (falling back to `codex auth status` on unknown-command + # shapes) AND a tiny `codex exec` smoke test must all pass. + def probe_codex_auth(env: nil) + build("codex_auth") do + env = (env || ENV.to_h).to_h + home = codex_home(env) + unless Hive::AgentProfiles.logged_in?(:codex, home: home) + next [ false, "no codex credential at ~/.codex/auth.json" ] + end + + login = run_codex_login_status(env) + next [ false, login.detail ] unless login.ok + + smoke = run_codex_smoke_test(env) + next [ false, smoke.detail ] unless smoke.ok + + [ true, "codex login status + smoke test passed" ] + end + end + + # Claude launcher probe: the wrapper file must exist, the tmux ready + # detector (tmux mode) must report `:present`, the claude binary must + # meet its minimum version, the loaded code must match the on-disk + # `lib/hive.rb`, and `hive doctor` must be green. All must pass. + def probe_claude_launcher(config: nil, project_root: nil, env: nil, wrapper_path: nil, code_fingerprint: nil, stage: nil, universal_result: nil) + build("claude_launcher") do + cfg = config || @config + env = (env || ENV.to_h).to_h + wrapper_path ||= default_wrapper_path + unless File.file?(wrapper_path) + next [ false, "launcher wrapper missing: #{wrapper_path}" ] + end + + if Hive::Config.claude_mode(cfg) == :tmux + status, message = Hive::ClaudeLauncher.tmux_status + next [ false, message ] unless status == :present + end + + version = claude_version_check(env, cfg: cfg) + next [ false, version.detail ] unless version.ok + + fingerprint_detail = code_fingerprint_match(code_fingerprint) + next [ false, fingerprint_detail ] unless fingerprint_detail.nil? + + # Reuse a precomputed universal result (the healer already ran it + # once this cycle) instead of building a second `hive doctor`. + universal = universal_result || probe_universal(config: cfg, project_root: project_root, stage: stage) + next [ false, "doctor not green: #{universal.detail}" ] unless universal.healthy? + + [ true, "wrapper + ready detector + binary version + doctor all green" ] + end + end + + private + + # Shared result builder: measures wall-clock elapsed time, coerces the + # block's `[healthy, detail]` into a `Result`, and — crucially — + # rescues any unexpected exception into `healthy: false` so a probe bug + # can never crash a daemon tick. + def build(name) + start = monotonic_now + healthy, detail = yield + Result.new( + name: name, + healthy: !!healthy, + detail: detail.to_s, + elapsed_sec: monotonic_now - start + ) + rescue StandardError => e + Result.new( + name: name, + healthy: false, + status: "unhealthy", + detail: "#{e.class}: #{e.message}", + elapsed_sec: monotonic_now - start + ) + end + + def monotonic_now + Process.clock_gettime(Process::CLOCK_MONOTONIC) + end + + def codex_home(env) + home = env["CODEX_HOME"].to_s + home.empty? ? (env["HOME"] || Dir.home) : home + end + + def codex_bin(env) + bin = env["HIVE_CODEX_BIN"].to_s + bin.empty? ? "codex" : bin + end + + # In-process shell of the codex login status check. Returns a small + # `LoginStatus`-ish struct with `ok`/`detail` so the public probe reads + # declaratively. The unknown-command fallback is applied only when the + # primary subcommand's stderr looks like an unknown-command rejection — + # a genuine auth failure (exit non-zero for a real reason) must NOT be + # silently retried through `auth status`. + def run_codex_login_status(env) + primary = @runner.capture3( + [ codex_bin(env), "login", "status" ], env: env, timeout_sec: CODEX_LOGIN_TIMEOUT_SEC + ) + if login_status_ok?(primary) + return ProbeStep.new(ok: true, detail: "codex login status reported logged-in") + end + return ProbeStep.new(ok: false, detail: login_failure_detail(primary)) unless unknown_command?(primary) + + fallback = @runner.capture3( + [ codex_bin(env), "auth", "status" ], env: env, timeout_sec: CODEX_LOGIN_TIMEOUT_SEC + ) + if login_status_ok?(fallback) + ProbeStep.new(ok: true, detail: "codex auth status reported logged-in") + else + ProbeStep.new(ok: false, detail: login_failure_detail(fallback, fallback: true)) + end + end + + def run_codex_smoke_test(env) + result = @runner.capture3( + [ codex_bin(env), "exec", "--dangerously-bypass-approvals-and-sandbox", CODEX_SMOKE_PROMPT ], + env: env, timeout_sec: CODEX_SMOKE_TIMEOUT_SEC + ) + if result.success? + ProbeStep.new(ok: true, detail: "codex exec smoke test passed") + else + ProbeStep.new(ok: false, detail: "codex exec smoke test failed (exit #{result.exitstatus})") + end + end + + # A shell-out step in a probe; `ok` is the pass/fail bit, `detail` the + # audit-log-able reason. + ProbeStep = Struct.new(:ok, :detail, keyword_init: true) + + def login_status_ok?(result) + return false unless result.success? + + text = result.output.downcase + text.include?("logged in") && !text.include?("not logged in") + end + + def unknown_command?(result) + return false if result.success? + + result.stderr.to_s.downcase.match?(/unknown|unrecognized|no such|not a recognized/i) + end + + def login_failure_detail(result, fallback: false) + subcommand = fallback ? "codex auth status" : "codex login status" + "not logged in per #{subcommand} (exit #{result.exitstatus}): #{first_line(result.output)}" + end + + def first_line(text) + text.to_s.lines.map(&:strip).find { |line| !line.empty? }.to_s + end + + # `claude binary --version` meets the profile minimum. Delegates to the + # profile's own `check_version!` (which shells out and raises + # `Hive::AgentError` on missing/too-old/parse failure). + def claude_version_check(env, cfg: @config) + profile = Hive::AgentProfiles.lookup(:claude, cfg: cfg) + profile.check_version! + ProbeStep.new(ok: true, detail: "claude binary meets minimum version") + rescue Hive::AgentError => e + ProbeStep.new(ok: false, detail: e.message) + rescue StandardError => e + ProbeStep.new(ok: false, detail: "#{e.class}: #{e.message}") + end + + # Compare the daemon's loaded code fingerprint against a fresh SHA-256 + # of `lib/hive.rb`. Returns nil when they match (or when no expected + # fingerprint is supplied — the fresh digest merely must be computable); + # returns a failure detail otherwise. + def code_fingerprint_match(code_fingerprint) + fresh = fresh_code_digest + return "could not compute fresh lib/hive.rb digest" if fresh.nil? + return nil if code_fingerprint.nil? || code_fingerprint == fresh + + "loaded code fingerprint does not match on-disk lib/hive.rb" + end + + def fresh_code_digest + path = Hive::Schemas.method(:schema_path).source_location.first + ::Digest::SHA256.file(path).hexdigest + rescue StandardError + nil + end + + def default_wrapper_path + self.class.wrapper_path + end + + # Verify the stage's configured skill resolves for its agent profile. + # Returns nil when healthy; a human-readable failure string otherwise. + # `stage` may be a full stage dir or a bare stage name (e.g. `execute`); + # a nil/unknown stage contributes no extra check (doctor's + # brainstorm/plan/reviewer checks already run regardless). + def stage_skill_failure(stage, cfg: @config, project_root: @project_root) + key = stage_config_key(stage) + return nil if key.nil? + + agent_name = (cfg.dig(key, "agent") || "claude").to_s + skill = Hive::Config.stage_skill(cfg, key) + return nil if skill.nil? || skill.to_s.empty? + + profile = Hive::AgentProfiles.lookup(agent_name.to_sym) + invocation = profile.format_skill_invocation(skill) + status, message = profile.verify_skill(invocation, project_root: project_root) + return nil unless FAILING_STATUSES.include?(status.to_s) + + "#{key} skill #{invocation} #{status}: #{message}" + rescue StandardError => e + "#{key} skill check failed: #{e.class}: #{e.message}" + end + + def stage_config_key(stage) + return nil if stage.nil? || stage.to_s.empty? + + # Stage dirs use hyphens (`5-open-pr`) while the config keys use + # underscores (`open_pr`); normalize so the skill check reads the + # right key instead of silently no-op'ing on `cfg.fetch("open-pr")`. + stage.to_s.split("-", 2).last.tr("-", "_") + end + + def failing_desc(failing) + return nil if failing.empty? + + "doctor failing rows: #{failing.map { |row| "#{row[:label]}=#{row[:status]}" }.join(', ')}" + end + end + end +end diff --git a/lib/hive/daemon/health_signals.rb b/lib/hive/daemon/health_signals.rb new file mode 100644 index 000000000..df3d6d7f7 --- /dev/null +++ b/lib/hive/daemon/health_signals.rb @@ -0,0 +1,173 @@ +require "digest" +require "hive/agent_profiles" + +module Hive + module Daemon + # Cheap, stable "health signal" fingerprint for the probe-gated auto-retry + # reasons, plus a per-tick memo cache so N parked rows sharing a + # reason/project run at most one probe per daemon tick. + # + # The fingerprint is deliberately TIME-INDEPENDENT: it changes only when + # an input that could plausibly fix the failure changes (binary path / + # code digest, config slice, relevant env, Codex login state, wrapper + # mtime+size, plugin/skill inventory, agent profile). A time-varying + # digest would defeat the changed-signal throttle that keeps the daemon + # from re-probing an expensive CLI on every 30s tick. + module HealthSignals + SKILL_INVENTORY_CAP = 100 + ENV_KEYS = %w[HIVE_CODEX_BIN HIVE_CLAUDE_BIN HIVE_TMUX_BIN CODEX_HOME HIVE_QMD_BIN].freeze + STAGE_KEYS = %w[brainstorm plan execute open_pr artifacts finalize].freeze + + module_function + + # Canonical SHA-256 over the health-relevant inputs. Any component that + # fails to read degrades to a constant token (`nil` is dropped), so a + # transient read error cannot make the fingerprint flap. + def fingerprint(cfg:, project_root:, profile:, wrapper_path:, env:) + parts = [ + code_digest, + config_slice(cfg), + env_slice(env), + codex_login_state(env), + wrapper_state(wrapper_path), + skill_inventory(project_root, env: env), + profile_state(profile) + ] + ::Digest::SHA256.hexdigest(parts.compact.join("\u0000")) + end + + # SHA-256 of lib/hive.rb (the file holding SCHEMA_VERSIONS) — the same + # source the dispatcher uses for its re-exec drift signal. nil on any + # read failure (a nil component is dropped, not hashed). + def code_digest + path = Hive::Schemas.method(:schema_path).source_location.first + ::Digest::SHA256.file(path).hexdigest + rescue StandardError + nil + end + + # The merged `daemon` block plus the stage agent/skill keys the retried + # stage (or its siblings) actually consume. + def config_slice(cfg) + daemon = cfg["daemon"] || {} + stages = STAGE_KEYS.flat_map do |stage| + slice = cfg[stage] + slice.nil? ? [] : [ stage, slice["agent"], slice["skill"], slice["skill_by_agent"] ] + end + canonical([ daemon, stages ]) + end + + def env_slice(env) + canonical(ENV_KEYS.to_h { |key| [ key, env[key] ] }) + end + + def codex_login_state(env) + home = env["CODEX_HOME"].to_s + home = env["HOME"] || Dir.home if home.empty? + Hive::AgentProfiles.logged_in?(:codex, home: home) ? "logged-in" : "logged-out" + end + + # mtime + size token for the launcher wrapper; "missing" when absent. + def wrapper_state(wrapper_path) + stat = File.stat(wrapper_path) + "#{stat.mtime.to_i}:#{stat.size}" + rescue SystemCallError + "missing" + end + + # Bounded `Dir[]` listing of the installed skill markers under the + # configured skill roots (project-local first, then user-global). + def skill_inventory(project_root, env:) + home = env["HOME"].to_s + home = Dir.home if home.empty? + roots = [] + unless project_root.nil? || project_root.to_s.empty? + roots << File.join(project_root.to_s, ".claude", "skills") + roots << File.join(project_root.to_s, ".codex", "skills") + end + roots << File.join(home, ".claude", "skills") + roots << File.join(home, ".codex", "skills") + roots.flat_map { |root| skill_entries(root) }.join(",") + end + + def profile_state(profile) + return nil if profile.nil? + + "#{profile.name}:#{profile.bin}" + end + + # Recursive, key-sorted, deterministic serialization so equal inputs + # (independent of hash insertion order) produce equal tokens. + def canonical(value) + case value + when Hash + value.keys.sort.map { |key| "#{key}=#{canonical(value[key])}" }.join(",") + when Array + value.map { |entry| canonical(entry) }.join("|") + when nil + "" + else + value.to_s + end + end + + def skill_entries(root) + Dir[File.join(root, "**", "SKILL.md")] + .sort + .first(SKILL_INVENTORY_CAP) + .map { |path| "#{File.basename(File.dirname(path))}=#{stat_token(path)}" } + end + + def stat_token(path) + stat = File.stat(path) + "#{stat.mtime.to_i}:#{stat.size}" + rescue SystemCallError + "?" + end + end + + # Per-tick probe-result memo. Keyed by `(reason, project)`; cleared at + # the start of each daemon tick (mirroring Dispatcher's `@enabled_cache`) + # so two parked rows sharing a reason/project within one tick share one + # probe result, while the next tick recomputes. + class HealthSignalCache + def initialize + @entries = {} + @fingerprints = {} + end + + def clear + @entries.clear + @fingerprints.clear + end + + def fetch(key) + @entries[key] + end + + def store(key, fingerprint:, results:) + @entries[key] = { fingerprint: fingerprint, results: results } + end + + # Memoize an expensive fingerprint computation per (reason, project) + # per tick. The block runs at most once per key per tick; the value is + # dropped on `clear` so the next tick recomputes (and can detect a + # changed signal). This is what makes the plan's "cheap, cached" + # fingerprint intent actually cheap: N parked rows sharing a + # reason/project no longer repeat the SHA-256 + skill-inventory globs. + def fingerprint(key) + return @fingerprints[key] if @fingerprints.key?(key) + + @fingerprints[key] = yield + end + + # True when there is no recorded entry for `key` OR the recorded + # fingerprint differs from the freshly computed one. This is the + # healer's "should I re-probe?" gate. + def changed_since_last?(key, fingerprint) + entry = @entries[key] + entry.nil? || entry[:fingerprint] != fingerprint + end + end + end +end diff --git a/lib/hive/daemon/logger.rb b/lib/hive/daemon/logger.rb index 6af59dc93..a7c17bfca 100644 --- a/lib/hive/daemon/logger.rb +++ b/lib/hive/daemon/logger.rb @@ -50,6 +50,7 @@ module Hive marker_heal_failed marker_heal_exhausted marker_heal_observer_missing + auto_retry_blocked display_name_backfill update_available update_check_no_result diff --git a/lib/hive/daemon/recoverable_signatures.rb b/lib/hive/daemon/recoverable_signatures.rb new file mode 100644 index 000000000..5db15c790 --- /dev/null +++ b/lib/hive/daemon/recoverable_signatures.rb @@ -0,0 +1,110 @@ +require "hive/diagnostic_helpers" +require "hive/paths" + +module Hive + module Daemon + # Classifies the two v1 allowlisted recoverable terminal-error signatures + # from on-disk evidence (marker attrs + task/global log tails). Everything + # else — business-logic failures, test failures, review findings, merge + # conflicts, dirty-worktree failures, and unknown `exit_code=1` — returns + # false so the healer never auto-clears it. + # + # The Codex 401 shape is the ONLY admitted `implementer_failed` variant: + # the marker `message` or the captured implementer output must carry `401` + # together with `bearer`/`basic auth` (case-insensitive), or the literal + # Codex `Missing bearer/basic auth` string. `claude_launch_failed` needs no + # extra text classification — the reason itself is the signature; the + # probe (not text) is its gate. + module RecoverableSignatures + LOG_GLOB_CAP = 20 + + module_function + + # Three-state classification for the healer: `:codex_auth` when the + # implementer_failed marker carries the Codex-401 diagnostic signature, + # `:claude_launcher` for the claude_launch_failed reason, else nil + # (never auto-clear). + def classify(row, state_home: nil) + reason = marker_reason(row) + return :codex_auth if reason == "implementer_failed" && codex_auth_text?(evidence_text(row, state_home: state_home)) + + return :claude_launcher if reason == "claude_launch_failed" + + nil + end + + # The only admitted `implementer_failed` shape: the marker's message + # AND the captured implementer output (task logs + global per-slug log) + # carry the Codex `Missing bearer/basic auth` / `401 + auth` signature. + def codex_auth_failure?(row, state_home: nil) + marker_reason(row) == "implementer_failed" && + codex_auth_text?(evidence_text(row, state_home: state_home)) + end + + def claude_launch_signature?(row) + marker_reason(row) == "claude_launch_failed" + end + + # Text predicate shared by the classifier and its tests. Accepts the + # literal Codex `Missing bearer/basic auth` string, or a `401` token + # appearing together with a `bearer`/`basic auth` token. + def codex_auth_text?(text) + text = text.to_s + return true if text.match?(/Missing bearer\/basic auth/i) + + text.match?(/401/) && text.match?(/bearer|basic auth/i) + end + + def evidence_text(row, state_home:) + [ + marker_message(row), + task_log_tails(row), + state_log_tails(row, state_home: state_home) + ].flatten.compact.join("\n") + end + + def marker_reason(row) + marker_attrs(row)["reason"].to_s + end + + def marker_message(row) + marker_attrs(row)["message"].to_s + end + + def marker_attrs(row) + return {} unless row.respond_to?(:marker_attrs) + + attrs = row.marker_attrs + attrs.is_a?(Hash) ? attrs : {} + end + + # Bounded tail scan of `/logs/*.log`. Non-regular files + # (FIFOs, sockets) are skipped so the scan can never block a tick. + def task_log_tails(row) + return [] unless row.respond_to?(:folder) && row.folder + + log_tails(File.join(row.folder.to_s, "logs")) + end + + # Bounded tail scan of the global `/logs//*.log` dir. + def state_log_tails(row, state_home:) + return [] unless row.respond_to?(:slug) && row.slug + + home = state_home || Hive::Paths.state_home + log_tails(File.join(home.to_s, "logs", row.slug.to_s)) + end + + def log_tails(dir) + Dir[File.join(dir, "*.log")].sort.first(LOG_GLOB_CAP).filter_map { |path| tail_file(path) } + end + + def tail_file(path) + return nil unless File.file?(path) + + Hive::DiagnosticHelpers.tail_file(path) + rescue SystemCallError, IOError + nil + end + end + end +end diff --git a/lib/hive/daemon/stale_agent_healer.rb b/lib/hive/daemon/stale_agent_healer.rb index 02042f789..c2ab8a2e3 100644 --- a/lib/hive/daemon/stale_agent_healer.rb +++ b/lib/hive/daemon/stale_agent_healer.rb @@ -2,10 +2,18 @@ require "digest" require "open3" require "time" require "yaml" +require "hive/agent_profiles" +require "hive/brainstorm_parser" +require "hive/config" +require "hive/events" require "hive/lock" require "hive/markers" require "hive/workflows" +require "hive/worktree" require "hive/daemon/dispatch_request_queue" +require "hive/daemon/health_probe" +require "hive/daemon/health_signals" +require "hive/daemon/recoverable_signatures" module Hive module Daemon @@ -97,6 +105,17 @@ module Hive # because a timeout means "I ran and didn't finish", not "I was interrupted". TIMEOUT_RECOVERY_LIMIT = 1 TIMEOUT_RECOVERABLE_STAGES = %w[5-open-pr 7-artifacts].freeze # coding-scoped: coding stages whose timeout re-entry is idempotent + # Probe-gated auto-retry (v1): `implementer_failed` (Codex 401) and + # `claude_launch_failed` (launcher healthy) recover at most twice per + # task/reason, backoff immediate then 30 min, and only ever re-probe on + # a changed health signal plus a low-frequency fallback. + RECOVERABLE_MARKER_RETRY_LIMIT = 2 + RETRY_BACKOFF_SEC = 1800 + PROBE_FALLBACK_SEC = 1800 + WORKTREE_OWNING_STAGES = %w[4-execute 6-review 8-finalize].freeze # coding-scoped: stages whose retry safety gate must verify a clean worktree + # coding-scoped: stages whose state file holds user-answerable content; + # a from-scratch rerun must verify no answered content would be overwritten. + ANSWER_CONTENT_STAGES = %w[2-brainstorm 3-plan].freeze # Review fix-phase auto-commit failures that a bounded rerun can clear: # the fix agent left residue the scope check rejected, or a transient @@ -115,7 +134,14 @@ module Hive def initialize(controller:, logger:, grace_sec: 300, review_error_auto_recovery_limit: REVIEW_ERROR_AUTO_RECOVERY_LIMIT, error_auto_recovery_limit: ERROR_AUTO_RECOVERY_LIMIT, - request_queue: Hive::Daemon::DispatchRequestQueue) + request_queue: Hive::Daemon::DispatchRequestQueue, + probe: nil, + signal_cache: nil, + auto_retry_enabled: true, + state_home: nil, + env: nil, + project_config_resolver: nil, + code_fingerprint: nil) @controller = controller @logger = logger @request_queue = request_queue @@ -126,6 +152,23 @@ module Hive @error_auto_recoveries = Hash.new(0) @review_error_recovery_exhausted = {} @error_recovery_exhausted = {} + @probe = probe + @signal_cache = signal_cache || Hive::Daemon::HealthSignalCache.new + @auto_retry_enabled = auto_retry_enabled + @state_home = state_home + @env = (env || ENV.to_h).to_h + @project_config_resolver = project_config_resolver + # Loaded-code baseline (the daemon's startup digest of lib/hive.rb), + # passed through to the launcher probe so it can verify the loaded + # code matches the on-disk lib/hive.rb before retrying. + @code_fingerprint = code_fingerprint + # Per-process retry budget/backoff/fingerprint for the probe-gated + # reasons (documented `budget_scope=per_process`, matching the other + # healer budgets). Keyed by [project, slug, stage, reason]. + @probe_retry_state = {} + # One-shot throttled-negative dedup (once per fingerprint change). + @probe_blocked_logged = {} + @probe_exhausted_logged = {} end # Walk the row set, heal stale agent_working markers in place. @@ -180,6 +223,19 @@ module Hive def heal_error_if_auto_recoverable(row, now:) return if row.live_task_lock == true + + # Probe-gated auto-retry for the two v1 recoverable markers. This is + # a NARROW extension: only `implementer_failed` with a recognized + # Codex-401 signature, or `claude_launch_failed` — both require a + # health probe to pass, a changed health-signal fingerprint, a clean + # work area, and a stricter retry budget (2) with backoff. Everything + # else falls through to the existing marker-shape heuristics below. + probe_reason = probe_gated_reason(row) + if probe_reason + heal_probe_gated_error(row, now: now, probe_reason: probe_reason) + return + end + return unless auto_recoverable_error?(row, now: now) # The clear makes a markerless terminal-error row take the # edit-resume path; without a pre-clear mtime to seed as the @@ -285,6 +341,343 @@ module Hive remediation: "hive plan #{row.slug} --project #{row.project} --from 3-plan") # coding-scoped: healer re-enters coding plan verb end + # --- probe-gated auto-retry (v1) ------------------------------------ + + # Three-state classification for the two probe-gated reasons. nil when + # auto-retry is disabled, the probe seam is unwired, or the row is not + # one of the two allowlisted signatures. + def probe_gated_reason(row) + return nil unless @probe && @auto_retry_enabled + + Hive::Daemon::RecoverableSignatures.classify(row, state_home: @state_home) + end + + # Probe-first, changed-signal, safety-gated, budgeted retry. The outer + # `heal` loop already skips live_task_lock / running_task? / + # half-migrated projects; live_task_lock is re-checked by the caller. + def heal_probe_gated_error(row, now:, probe_reason:) + return if row.state_file_mtime.nil? + + retry_key = probe_retry_key(row, probe_reason) + cache_key = probe_cache_key(row, probe_reason) + state = (@probe_retry_state[retry_key] ||= { + attempts: 0, last_fingerprint: nil, last_attempt_at: nil, last_probe_at: nil + }) + context = probe_context_for(row, probe_reason) + # The fingerprint hashes lib/hive.rb plus recursive skill-inventory + # globs, so N parked rows sharing a (reason, project) must not repeat + # it N times per tick. Memoize it per tick alongside the probe + # results (R6: cheap, cached fingerprint). + fingerprint = @signal_cache.fingerprint(cache_key) do + Hive::Daemon::HealthSignals.fingerprint( + cfg: context[:cfg], project_root: context[:project_root], + profile: context[:profile], wrapper_path: context[:wrapper_path], + env: context[:env] + ) + end + + # Capture the changed-signal flag BEFORE `probe_results_for` runs: + # that method overwrites `state[:last_fingerprint]` with the current + # fingerprint (on both the probe and memo-hit paths), so recomputing + # it afterwards would always see "unchanged" and block a changed-signal + # retry as "backoff" — exactly the opposite of the plan's escape. + signal_changed = state[:last_fingerprint].nil? || state[:last_fingerprint] != fingerprint + results = probe_results_for(cache_key, retry_key, state, context, row, + fingerprint, probe_reason, now: now) + return if results.nil? # throttled negative already logged + + unless results_healthy?(results) + log_probe_blocked_once(retry_key, fingerprint, row, probe_reason, + rationale: "probe not healthy", probe_detail: probe_detail(results)) + return + end + + unless safe_to_retry?(row) + log_probe_blocked_once(retry_key, fingerprint, row, probe_reason, rationale: "unsafe work area") + return + end + + if state[:attempts] >= RECOVERABLE_MARKER_RETRY_LIMIT + log_probe_exhausted_once(retry_key, row, probe_reason, state) + return + end + + if state[:attempts].positive? && state[:last_attempt_at] && + (now - state[:last_attempt_at]) < RETRY_BACKOFF_SEC && !signal_changed + log_probe_blocked_once(retry_key, fingerprint, row, probe_reason, rationale: "backoff") + return + end + + marker_reason = probe_marker_reason(probe_reason) + return unless Hive::Markers.clear_current( + row.state_file, + expected_name: :error, + match_attrs: auto_recoverable_error_match_attrs(row, reason: marker_reason) + ) + + observe_pre_clear_mtime(row) + state[:attempts] += 1 + state[:last_attempt_at] = now + state[:last_fingerprint] = fingerprint + heal_label = probe_heal_label(probe_reason) + @logger.event(:marker_healed, + project: row.project, + slug: row.slug, + stage: row.stage, + prior_marker: row.marker, + reason: heal_label, + marker_reason: marker_reason, + state_file: row.state_file, + fingerprint: fingerprint, + probes: probe_summaries(results), + attempts: state[:attempts], + max_attempts: RECOVERABLE_MARKER_RETRY_LIMIT) + emit_auto_retry_event(row, probe_reason, results, fingerprint, state) + # 3-plan needs the explicit requeue exactly like the other heals — a + # markerless empty plan.md re-classifies straight back to :error. + requeue_plan_rerun(row) if Hive::Workflows.coding_row?(row) && row.stage.to_s == "3-plan" # coding-scoped: coding plan pause needs bespoke rerun after marker clear + rescue StandardError => e + @logger.event(:marker_heal_failed, + project: row.project, + slug: row.slug, + stage: row.stage, + reason: probe_heal_label(probe_reason), + error: "#{e.class}: #{e.message}") + end + + def probe_retry_key(row, probe_reason) + [ row.project.to_s, row.slug.to_s, row.stage.to_s, probe_reason.to_s ] + end + + def probe_cache_key(row, probe_reason) + [ probe_reason.to_s, row.project.to_s ] + end + + def probe_context_for(row, probe_reason) + cfg, project_root = resolve_project_config(row.project) + { + reason: probe_reason, + cfg: cfg || {}, + project_root: project_root, + profile: probe_profile(probe_reason, cfg), + wrapper_path: Hive::Daemon::HealthProbe.wrapper_path, + env: @env + } + end + + def resolve_project_config(project_name) + return @project_config_resolver.call(project_name) if @project_config_resolver + + entry = Hive::Config.find_project(project_name) + return [ nil, nil ] if entry.nil? + + [ Hive::Config.load(entry["path"]), entry["path"] ] + rescue Hive::ConfigError, StandardError + [ nil, nil ] + end + + def probe_profile(probe_reason, cfg) + name = probe_reason == :codex_auth ? :codex : :claude + Hive::AgentProfiles.lookup(name, cfg: cfg) + end + + # Universal precondition + reason-specific probe, both returned as + # `HealthProbe::Result` objects (never raised). The universal probe + # runs exactly once per probe cycle: `probe_claude_launcher` also runs + # it internally, so it receives the precomputed result instead of + # building a second `hive doctor` (R6: one doctor per cycle). + def run_probes(context, row) + universal = @probe.probe_universal( + config: context[:cfg], project_root: context[:project_root], stage: row.stage + ) + { universal: universal, reason: reason_probe_result(context, row, universal) } + end + + def reason_probe_result(context, row, universal) + if context[:reason] == :codex_auth + @probe.probe_codex_auth(env: context[:env]) + else + @probe.probe_claude_launcher( + config: context[:cfg], project_root: context[:project_root], + env: context[:env], wrapper_path: context[:wrapper_path], + code_fingerprint: @code_fingerprint, stage: row.stage, + universal_result: universal + ) + end + end + + # Per-tick memo + changed-signal/fallback throttle. Returns the probe + # results hash, or nil when this tick is throttled (negative logged). + def probe_results_for(cache_key, retry_key, state, context, row, fingerprint, probe_reason, now:) + memo = @signal_cache.fetch(cache_key) + if memo && memo[:fingerprint] == fingerprint + # A memo hit means another row sharing this (reason, project) already + # ran the probes this tick. Record this row's retry-state bookkeeping + # too, otherwise its `last_probe_at`/`last_fingerprint` stay nil and + # the next tick's fallback check re-runs the expensive probes. + state[:last_fingerprint] = fingerprint + state[:last_probe_at] = now + return memo[:results] + end + + signal_changed = state[:last_fingerprint].nil? || state[:last_fingerprint] != fingerprint + fallback_due = state[:last_probe_at].nil? || (now - state[:last_probe_at]) >= PROBE_FALLBACK_SEC + unless signal_changed || fallback_due + log_probe_blocked_once(retry_key, fingerprint, row, probe_reason, rationale: "signal unchanged") + return nil + end + + results = run_probes(context, row) + @signal_cache.store(cache_key, fingerprint: fingerprint, results: results) + state[:last_fingerprint] = fingerprint + state[:last_probe_at] = now + results + end + + def results_healthy?(results) + results.values.all? { |result| result.respond_to?(:healthy?) && result.healthy? } + end + + def probe_detail(results) + results.map { |name, result| "#{name}=#{result.detail}" }.join("; ") + end + + def probe_summaries(results) + results.map { |name, result| "#{name}:#{result.status}" } + end + + def probe_heal_label(probe_reason) + probe_reason == :codex_auth ? "codex_auth_recovered" : "claude_launcher_recovered" + end + + def probe_marker_reason(probe_reason) + probe_reason == :codex_auth ? "implementer_failed" : "claude_launch_failed" + end + + # Safety gate: never retry when it could clobber user work. For + # worktree-owning stages the worktree must be clean (porcelain empty); + # brainstorm/plan stages hold user-answerable content in their state + # file and must be verified answer-free before a from-scratch rerun + # (fail closed: uncertain ⇒ skip); other stages are safe. + def safe_to_retry?(row) + return worktree_clean?(row) if WORKTREE_OWNING_STAGES.include?(row.stage.to_s) + return no_answered_content?(row) if ANSWER_CONTENT_STAGES.include?(row.stage.to_s) + + true + end + + # Fail-closed "no answered content to overwrite" guard for the + # brainstorm/plan stages. A from-scratch rerun of those stages rewrites + # their state file, so we only retry when we can PROVE it holds no + # user-provided answers. A missing state file is the one certainly-safe + # case (nothing to overwrite); an unreadable file, a leftover + # WAITING/COMPLETE marker, or (for brainstorm) any answered `### A{n}.` + # section all fail closed. + def no_answered_content?(row) + path = row.state_file + return true if path.nil? || !File.exist?(path) + + content = File.read(path, encoding: "UTF-8").scrub + return false if content.match?(/)/) + return false if row.stage.to_s == "2-brainstorm" && # coding-scoped: coding brainstorm answer format + Hive::BrainstormParser.parse_text(content).any?(&:answered?) + + true + rescue StandardError + false + end + + def worktree_clean?(row) + path = worktree_path_for(row) + return false if path.nil? || !File.directory?(path) + + out, _err, status = Open3.capture3("git", "-C", path, "status", "--porcelain") + return false unless status.success? + + out.strip.empty? + rescue StandardError + false + end + + def worktree_path_for(row) + return nil unless row.respond_to?(:folder) && row.folder + + pointer = Hive::Worktree.read_pointer(row.folder) + pointer && pointer["path"] + rescue StandardError + nil + end + + # Throttled negative: logged once per distinct [project, slug, stage, + # reason, fingerprint] so a persistently-unhealthy marker does not spam + # every tick. + def log_probe_blocked_once(retry_key, fingerprint, row, probe_reason, rationale:, probe_detail: nil) + dedup_key = retry_key + [ fingerprint ] + return if @probe_blocked_logged[dedup_key] + + @probe_blocked_logged[dedup_key] = true + @logger.event(:auto_retry_blocked, + project: row.project, + slug: row.slug, + stage: row.stage, + reason: probe_heal_label(probe_reason), + rationale: rationale, + fingerprint: fingerprint, + probe_detail: probe_detail) + emit_auto_retry_blocked_event(row, probe_reason, fingerprint, rationale, probe_detail) + end + + def log_probe_exhausted_once(retry_key, row, probe_reason, state) + return if @probe_exhausted_logged[retry_key] + + @probe_exhausted_logged[retry_key] = true + @logger.event(:marker_heal_exhausted, + project: row.project, + slug: row.slug, + stage: row.stage, + prior_marker: row.marker, + state_file: row.state_file, + reason: probe_heal_label(probe_reason), + marker_reason: probe_marker_reason(probe_reason), + attempts: state[:attempts], + max_attempts: RECOVERABLE_MARKER_RETRY_LIMIT, + budget_scope: "per_process", + suggested_next_action: "manual_fix", + remediation: probe_recovery_remediation(row, probe_reason)) + end + + def probe_recovery_remediation(row, probe_reason) + command = "hive run #{row.slug} --project #{row.project} --stage #{row.stage}" + dependency = probe_reason == :codex_auth ? "Codex auth" : "the Claude launcher" + "#{dependency} did not recover after #{RECOVERABLE_MARKER_RETRY_LIMIT} auto-retries — " \ + "verify it manually, then rerun #{row.stage} (`#{command}`) or run `hive markers clear`" + end + + def emit_auto_retry_event(row, probe_reason, results, fingerprint, state) + Hive::Events.emit( + task_folder: row.folder.to_s, + slug: row.slug.to_s, + stage: row.stage.to_s, + event_type: :auto_retry, + agent: "daemon", + message: "auto-retried #{probe_marker_reason(probe_reason)} " \ + "(attempt #{state[:attempts]}/#{RECOVERABLE_MARKER_RETRY_LIMIT}); " \ + "probes=#{probe_summaries(results).join(' ')}; fingerprint=#{fingerprint}" + ) + end + + def emit_auto_retry_blocked_event(row, probe_reason, fingerprint, rationale, probe_detail) + Hive::Events.emit( + task_folder: row.folder.to_s, + slug: row.slug.to_s, + stage: row.stage.to_s, + event_type: :auto_retry_blocked, + agent: "daemon", + message: "auto-retry blocked (#{rationale}) for #{probe_marker_reason(probe_reason)}; " \ + "fingerprint=#{fingerprint}#{probe_detail ? "; #{probe_detail}" : ''}" + ) + end + def auto_recoverable_error?(row, now:) reason = marker_reason(row) return true if Hive::Workflows.coding_row?(row) && row.stage.to_s == "8-finalize" && reason == "unpushed_commits" # coding-scoped: unpushed commits are finalize/PR recovery diff --git a/lib/hive/events.rb b/lib/hive/events.rb index f8bca8160..1104c4981 100644 --- a/lib/hive/events.rb +++ b/lib/hive/events.rb @@ -15,6 +15,8 @@ module Hive round_complete clean_exit_auto_committed claude_completion_fallback + auto_retry + auto_retry_blocked ].freeze STATUS_TAIL_LINES = 20 diff --git a/test/integration/daemon_auto_retry_test.rb b/test/integration/daemon_auto_retry_test.rb new file mode 100644 index 000000000..2b00d4205 --- /dev/null +++ b/test/integration/daemon_auto_retry_test.rb @@ -0,0 +1,227 @@ +require "test_helper" +require "hive/commands/init" +require "hive/commands/status" +require "hive/daemon/health_probe" +require "hive/daemon/stale_agent_healer" +require "hive/daemon/status_consumer" +require "hive/daemon/logger" +require "hive/daemon/concurrency_controller" +require "hive/daemon/policy" +require "hive/markers" + +# End-to-end acceptance for the probe-gated auto-retry: drive the real +# `hive status --json` pipeline (so the marker `message`/`reason` attrs are +# proven to reach the classifier), then clear through the real Markers + +# Policy re-dispatch chain with an injected fake probe. +class DaemonAutoRetryTest < Minitest::Test + include HiveTestHelper + + Result = Hive::Daemon::HealthProbe::Result + HEADLESS_CFG = { "claude" => { "mode" => "headless" } }.freeze + + class FakeLogger + attr_reader :events + + def initialize + @events = [] + end + + def event(name, **attrs) + unless Hive::Daemon::Logger::EVENTS.include?(name) + raise ArgumentError, "FakeLogger rejected event #{name.inspect}" + end + @events << [ name, attrs ] + end + end + + class FakeProbe + attr_accessor :universal, :codex, :claude + + def initialize + @universal = Result.new(name: "universal", healthy: true, status: "healthy", detail: "ok") + @codex = Result.new(name: "codex_auth", healthy: true, status: "healthy", detail: "ok") + @claude = Result.new(name: "claude_launcher", healthy: true, status: "healthy", detail: "ok") + end + + def probe_universal(config: nil, project_root: nil, stage: nil) = @universal + def probe_codex_auth(env: nil) = @codex + def probe_claude_launcher(config: nil, project_root: nil, env: nil, wrapper_path: nil, code_fingerprint: nil, stage: nil, universal_result: nil) = @claude + end + + def setup + @logger = FakeLogger.new + @controller = Hive::Daemon::ConcurrencyController.new( + max_concurrent_runs: 4, max_concurrent_per_project: 2, max_runs_per_day_per_project: 50 + ) + end + + def build_healer(probe:, auto_retry_enabled: true, worktree_clean: true) + healer = Hive::Daemon::StaleAgentHealer.new( + controller: @controller, + logger: @logger, + grace_sec: 300, + probe: probe, + auto_retry_enabled: auto_retry_enabled, + project_config_resolver: ->(_name) { [ HEADLESS_CFG, nil ] } + ) + healer.define_singleton_method(:worktree_clean?, ->(_row) { worktree_clean }) unless worktree_clean == :real + healer + end + + def status_rows_via_consumer(dir) + out, _err = capture_io do + Hive::Commands::Status.new(json: true).call + rescue Hive::Error + # status exits non-zero when no tasks; that's fine here + end + doc = JSON.parse(out) + rows = [] + Array(doc["projects"]).each do |project| + Array(project["tasks"]).each do |task| + rows << Hive::Daemon::StatusConsumer::Row.new( + project: project["name"], slug: task["slug"], stage: task["stage"], + marker: task["marker"], marker_attrs: task["attrs"], folder: task["folder"], + state_file: task["state_file"], state_file_mtime: task["mtime"] ? Time.parse(task["mtime"]) : nil, + action: task["action"], suggested_command: task["suggested_command"], + claude_pid_alive: task["claude_pid_alive"], live_task_lock: task["live_task_lock"], + diagnostic: task["diagnostic"] + ) + end + end + rows + end + + def seed_execute_error(dir, slug, marker_line) + folder = File.join(dir, ".hive-state", "stages", "4-execute", slug) + FileUtils.mkdir_p(folder) + state_file = File.join(folder, "task.md") + File.write(state_file, "---\nslug: #{slug}\n---\n\n# #{slug}\n\n#{marker_line}\n") + pre_clear_mtime = Time.now - 1000 + File.utime(pre_clear_mtime, pre_clear_mtime, state_file) + [ folder, state_file ] + end + + def seed_plan_error(dir, slug, marker_line) + folder = File.join(dir, ".hive-state", "stages", "3-plan", slug) + FileUtils.mkdir_p(folder) + state_file = File.join(folder, "plan.md") + File.write(state_file, "# plan\n\n#{marker_line}\n") + pre_clear_mtime = Time.now - 1000 + File.utime(pre_clear_mtime, pre_clear_mtime, state_file) + [ folder, state_file ] + end + + def test_codex_401_auto_retry_clears_redispatches_and_emits_event + with_tmp_global_config do + with_tmp_git_repo do |dir| + capture_io { Hive::Commands::Init.new(dir).call } + slug = "codex-401-260701-aaaa" + folder, state_file = seed_execute_error( + dir, slug, + "" + ) + + row = status_rows_via_consumer(dir).find { |candidate| candidate.slug == slug } + assert row, "status pipeline must include the seeded task" + assert_equal "error", row.marker + assert_equal "implementer_failed", row.marker_attrs["reason"] + assert_match(/401 Missing bearer\/basic auth/, row.marker_attrs["message"].to_s) + + build_healer(probe: FakeProbe.new).heal([ row ], now: Time.now) + + assert Hive::Markers.current(state_file).none?, "the Codex-401 marker must clear" + assert(@logger.events.any? { |name, attrs| name == :marker_healed && attrs[:reason] == "codex_auth_recovered" }, + "expected codex_auth_recovered marker_healed, got: #{@logger.events.inspect}") + + # The task events.jsonl must carry the auto_retry audit line. + assert_includes File.read(File.join(folder, "events.jsonl")), "\"event_type\":\"auto_retry\"" + + # The clear + pre-clear mtime baseline must make Policy re-dispatch. + baseline = @controller.last_dispatched_state_file_mtime_for(project: File.basename(dir), slug: slug) + assert baseline, "heal must seed a dispatch baseline" + post_clear_row = status_rows_via_consumer(dir).find { |candidate| candidate.slug == slug } + decision = Hive::Daemon::Policy.decide( + action: post_clear_row.action, + stage: post_clear_row.stage, + command: post_clear_row.suggested_command, + state_file_mtime: post_clear_row.state_file_mtime, + last_dispatched_state_file_mtime: baseline, + now: post_clear_row.state_file_mtime + 3600, + edit_debounce_sec: 30 + ) + assert_equal :dispatch, decision, "a cleared markerless execute row must re-dispatch, not strand" + end + end + end + + def test_unknown_implementer_failed_stays_parked_through_status_pipeline + with_tmp_global_config do + with_tmp_git_repo do |dir| + capture_io { Hive::Commands::Init.new(dir).call } + slug = "impl-fail-260701-bbbb" + _folder, state_file = seed_execute_error( + dir, slug, + "" + ) + + row = status_rows_via_consumer(dir).find { |candidate| candidate.slug == slug } + assert row + probe = FakeProbe.new + build_healer(probe: probe).heal([ row ], now: Time.now) + + assert_match(/ERROR reason=implementer_failed/, File.read(state_file), + "an implementer_failed without a Codex-401 signature must stay parked") + assert_empty @logger.events.select { |name, _| name == :marker_healed } + end + end + end + + def test_kill_switch_disables_auto_retry + with_tmp_global_config do + with_tmp_git_repo do |dir| + capture_io { Hive::Commands::Init.new(dir).call } + slug = "codex-401-off-260701-cccc" + _folder, state_file = seed_execute_error( + dir, slug, + "" + ) + + row = status_rows_via_consumer(dir).find { |candidate| candidate.slug == slug } + assert row + build_healer(probe: FakeProbe.new, auto_retry_enabled: false).heal([ row ], now: Time.now) + + assert_match(/ERROR reason=implementer_failed/, File.read(state_file), + "kill-switch off must leave the marker parked") + assert_empty @logger.events.select { |name, _| name == :marker_healed } + end + end + end + + def test_3_plan_probe_gated_recovery_requeues_through_real_queue + with_tmp_global_config do + with_tmp_git_repo do |dir| + capture_io { Hive::Commands::Init.new(dir).call } + slug = "claude-plan-260701-dddd" + _folder, state_file = seed_plan_error( + dir, slug, + "" + ) + + row = status_rows_via_consumer(dir).find { |candidate| candidate.slug == slug } + assert row, "status pipeline must include the seeded 3-plan task" + assert_equal "3-plan", row.stage + assert_equal "claude_launch_failed", row.marker_attrs["reason"] + + build_healer(probe: FakeProbe.new).heal([ row ], now: Time.now) + + assert Hive::Markers.current(state_file).none?, "the probe-gated 3-plan marker must clear" + pending = Hive::Daemon::DispatchRequestQueue.pending + request = pending.find { |r| r.slug == slug } + refute_nil request, "3-plan recovery must land a rerun in the REAL queue" + assert_equal [ "hive", "plan", slug, "--project", File.basename(dir), "--from", "3-plan" ], + request.argv + assert_equal "healer", request.requestor + end + end + end +end diff --git a/test/unit/config_test.rb b/test/unit/config_test.rb index f9b9a92d6..c74a50665 100644 --- a/test/unit/config_test.rb +++ b/test/unit/config_test.rb @@ -2795,6 +2795,14 @@ class ConfigTest < Minitest::Test end end + def test_load_returns_daemon_auto_retry_enabled_by_default + with_tmp_dir do |dir| + cfg = Hive::Config.load(dir) + assert_equal true, cfg.dig("daemon", "auto_retry", "enabled"), + "probe-gated auto-retry must be enabled by default (kill-switch is opt-out)" + end + end + def test_load_rejects_negative_daemon_child_timeout_sec with_tmp_dir do |dir| FileUtils.mkdir_p(File.join(dir, ".hive-state")) diff --git a/test/unit/daemon/dispatcher_test.rb b/test/unit/daemon/dispatcher_test.rb index d5de27822..14a16016e 100644 --- a/test/unit/daemon/dispatcher_test.rb +++ b/test/unit/daemon/dispatcher_test.rb @@ -1888,6 +1888,27 @@ class HiveDaemonDispatcherTest < Minitest::Test "rebuilt healer must carry the reloaded grace value" end + def test_build_stale_agent_healer_gates_auto_retry_on_kill_switch + dispatcher, _sup, _ctrl, _logger, _mw = make_dispatcher + + # Key absent ⇒ enabled (the `!= false` gate defaults to on). + dispatcher.instance_variable_set(:@daemon_cfg, {}) + assert_equal true, + dispatcher.send(:build_stale_agent_healer, 300).instance_variable_get(:@auto_retry_enabled), + "an absent daemon.auto_retry.enabled key must default the kill-switch to enabled" + + # Explicit true ⇒ enabled. + dispatcher.instance_variable_set(:@daemon_cfg, { "auto_retry" => { "enabled" => true } }) + assert_equal true, + dispatcher.send(:build_stale_agent_healer, 300).instance_variable_get(:@auto_retry_enabled) + + # Explicit false ⇒ disabled. + dispatcher.instance_variable_set(:@daemon_cfg, { "auto_retry" => { "enabled" => false } }) + assert_equal false, + dispatcher.send(:build_stale_agent_healer, 300).instance_variable_get(:@auto_retry_enabled), + "daemon.auto_retry.enabled: false must disable probe-gated auto-retry" + end + def test_run_forever_reloads_ticks_and_shuts_down_cleanly dispatcher, supervisor, _ctrl, logger, _mw = make_dispatcher ticks = 0 diff --git a/test/unit/daemon/health_probe_test.rb b/test/unit/daemon/health_probe_test.rb new file mode 100644 index 000000000..8d452c278 --- /dev/null +++ b/test/unit/daemon/health_probe_test.rb @@ -0,0 +1,420 @@ +require "test_helper" +require "tmpdir" +require "timeout" +require "hive/daemon/health_probe" +require "hive/agent_profiles" + +# The health probe engine must answer "is this dependency healthy?" without +# ever raising, and must be fully driveable with injected fakes (no real +# codex/claude/tmux binary). These tests pin every branch: healthy and +# unhealthy paths for each probe, the shell-runner's fail-closed edges, and +# the exception-to-unhealthy degradation that protects a daemon tick. +class HiveDaemonHealthProbeTest < Minitest::Test + include HiveTestHelper + + HEADLESS_CFG = { "claude" => { "mode" => "headless" } }.freeze + TMUX_CFG = { "claude" => { "mode" => "tmux" } }.freeze + + CommandResult = Hive::Daemon::HealthProbe::CommandResult + + # A runner fake keyed by command prefix. `capture3` returns the first + # registered key that prefixes the argv, or a default failure. + class FakeRunner + attr_reader :calls + + def initialize + @calls = [] + @responses = {} + end + + def on(prefix, result) + @responses[prefix] = result + end + + def capture3(argv, env:, timeout_sec:) + @calls << { argv: argv, env: env, timeout_sec: timeout_sec } + key = @responses.keys.find { |prefix| argv[0, prefix.length] == prefix } + key ? @responses[key] : CommandResult.new(stdout: "", stderr: "", exitstatus: 1) + end + end + + def ok_result(stdout: "logged in", stderr: "", exitstatus: 0) + CommandResult.new(stdout: stdout, stderr: stderr, exitstatus: exitstatus) + end + + def fake_doctor(rows) + Class.new do + define_method(:initialize) { |**| } + define_method(:call) { } + define_method(:rows) { rows } + end + end + + def fake_claude_profile(check_raises: nil) + profile = Object.new + profile.define_singleton_method(:check_version!) do + raise check_raises if check_raises + + "2.1.118" + end + profile.define_singleton_method(:format_skill_invocation) { |skill| skill } + profile.define_singleton_method(:verify_skill) { |_invocation, project_root: nil| [ :present, "ok" ] } + profile + end + + def with_wrapper_file + Dir.mktmpdir do |dir| + path = File.join(dir, "wrapper.sh") + File.write(path, "#!/bin/sh\n") + yield path + end + end + + def make_probe(config: HEADLESS_CFG, runner: nil, doctor_factory: nil, project_root: nil) + Hive::Daemon::HealthProbe.new( + config: config, project_root: project_root, + runner: runner, doctor_factory: doctor_factory + ) + end + + # ── Result / CommandResult ───────────────────────────────────────────── + + def test_result_coerces_and_defaults_status + result = Hive::Daemon::HealthProbe::Result.new(name: "x", healthy: true) + assert result.healthy? + assert_equal "healthy", result.status + assert_equal "", result.detail + assert_equal 0.0, result.elapsed_sec + + falsey = Hive::Daemon::HealthProbe::Result.new(name: "y", healthy: false, status: "unhealthy", detail: "d") + refute falsey.healthy? + assert_equal "unhealthy", falsey.status + assert_equal "d", falsey.detail + end + + def test_command_result_success_and_output + ok = CommandResult.new(stdout: "a\n", stderr: "b", exitstatus: 0) + assert ok.success? + assert_equal "a\nb", ok.output + + fail = CommandResult.new(stdout: "", stderr: "", exitstatus: 1) + refute fail.success? + assert_equal "", fail.output + end + + # ── ShellRunner ──────────────────────────────────────────────────────── + + def test_shell_runner_captures_success + status = Struct.new(:exitstatus).new(0) + with_replaced_singleton_method(Open3, :capture3, ->(_env, *argv) { [ "out", "err", status ] }) do + result = Hive::Daemon::HealthProbe::ShellRunner.new.capture3([ "codex" ], env: {}, timeout_sec: 1) + assert result.success? + assert_equal "out", result.stdout + assert_equal "err", result.stderr + end + end + + def test_shell_runner_maps_missing_binary_to_127 + with_replaced_singleton_method(Open3, :capture3, ->(*_args) { raise Errno::ENOENT, "codex" }) do + result = Hive::Daemon::HealthProbe::ShellRunner.new.capture3([ "codex" ], env: {}, timeout_sec: 1) + assert_equal 127, result.exitstatus + refute result.success? + end + end + + def test_shell_runner_maps_timeout_to_124 + with_replaced_singleton_method(Timeout, :timeout, ->(_sec, &_block) { raise Timeout::Error }) do + result = Hive::Daemon::HealthProbe::ShellRunner.new.capture3([ "codex" ], env: {}, timeout_sec: 1) + assert_equal 124, result.exitstatus + refute result.success? + end + end + + # ── universal ────────────────────────────────────────────────────────── + + def test_universal_healthy_when_no_failing_doctor_rows + probe = make_probe(doctor_factory: fake_doctor([ { label: "a", status: "present" } ])) + result = probe.probe_universal + assert result.healthy? + assert_match(/doctor green/, result.detail) + end + + def test_universal_unhealthy_on_missing_row + probe = make_probe(doctor_factory: fake_doctor([ + { label: "brainstorm", status: "missing" }, + { label: "tmux", status: "version_too_old" } + ])) + result = probe.probe_universal + refute result.healthy? + assert_match(/brainstorm=missing/, result.detail) + assert_match(/tmux=version_too_old/, result.detail) + end + + def test_universal_unhealthy_when_stage_skill_missing + profile = Object.new + profile.define_singleton_method(:format_skill_invocation) { |skill| skill } + profile.define_singleton_method(:verify_skill) { |_invocation, project_root: nil| [ :missing, "nope" ] } + with_replaced_singleton_method(Hive::AgentProfiles, :lookup, ->(_name, cfg: nil) { profile }) do + probe = make_probe(doctor_factory: fake_doctor([])) + result = probe.probe_universal(stage: "3-plan") + refute result.healthy? + assert_match(/plan skill/, result.detail) + end + end + + def test_universal_rescues_doctor_exception_to_unhealthy + boom = Class.new do + define_method(:initialize) { |**| } + define_method(:call) { raise ArgumentError, "boom" } + define_method(:rows) { [] } + end + probe = make_probe(doctor_factory: boom) + result = probe.probe_universal + refute result.healthy? + assert_match(/ArgumentError: boom/, result.detail) + end + + def test_universal_unhealthy_when_doctor_returns_no_rows + # A doctor that hits its config-error rescue returns EXIT_CONFIG_ERROR + # without assigning @rows. The probe must treat nil rows as "probe error, + # not healthy" — never as "no failing rows ⇒ green". + nil_rows = Class.new do + define_method(:initialize) { |**| } + define_method(:call) { 78 } # Hive::Commands::Doctor::EXIT_CONFIG_ERROR + define_method(:rows) { nil } + end + probe = make_probe(doctor_factory: nil_rows) + result = probe.probe_universal + refute result.healthy? + assert_match(/config error/, result.detail) + end + + def test_universal_returns_unhealthy_on_timeout + # A hung in-process doctor (e.g. an unbounded tmux `-V` shell-out) must + # be bounded by UNIVERSAL_PROBE_TIMEOUT_SEC and degrade to unhealthy, + # never block the daemon tick (R5). + with_replaced_singleton_method(Timeout, :timeout, ->(_sec, &_block) { raise Timeout::Error }) do + probe = make_probe(doctor_factory: fake_doctor([])) + result = probe.probe_universal + refute result.healthy? + assert_match(/Timeout::Error/, result.detail) + end + end + + # ── codex auth ───────────────────────────────────────────────────────── + + def test_codex_auth_healthy_when_login_and_smoke_pass + runner = FakeRunner.new + runner.on([ "codex", "login", "status" ], ok_result) + runner.on([ "codex", "exec" ], ok_result(stdout: "ok")) + with_replaced_singleton_method(Hive::AgentProfiles, :logged_in?, ->(_name, home: Dir.home) { true }) do + result = make_probe(runner: runner).probe_codex_auth(env: {}) + assert result.healthy?, result.detail + assert_match(/smoke test passed/, result.detail) + end + end + + def test_codex_auth_unhealthy_without_credential + with_replaced_singleton_method(Hive::AgentProfiles, :logged_in?, ->(_name, home: Dir.home) { false }) do + result = make_probe.probe_codex_auth(env: {}) + refute result.healthy? + assert_match(/auth.json/, result.detail) + end + end + + def test_codex_auth_unhealthy_when_login_reports_not_logged_in + runner = FakeRunner.new + runner.on([ "codex", "login", "status" ], ok_result(stdout: "not logged in")) + with_replaced_singleton_method(Hive::AgentProfiles, :logged_in?, ->(_name, home: Dir.home) { true }) do + result = make_probe(runner: runner).probe_codex_auth(env: {}) + refute result.healthy? + assert_match(/not logged in/, result.detail) + end + end + + def test_codex_auth_falls_back_to_auth_status_on_unknown_command + runner = FakeRunner.new + runner.on([ "codex", "login", "status" ], CommandResult.new(stdout: "", stderr: "unknown command: login", exitstatus: 2)) + runner.on([ "codex", "auth", "status" ], ok_result) + runner.on([ "codex", "exec" ], ok_result(stdout: "ok")) + with_replaced_singleton_method(Hive::AgentProfiles, :logged_in?, ->(_name, home: Dir.home) { true }) do + result = make_probe(runner: runner).probe_codex_auth(env: {}) + assert result.healthy?, result.detail + assert_includes runner.calls.map { |c| c[:argv] }, [ "codex", "auth", "status" ], + "the unknown-command shape must trigger the auth status fallback" + end + end + + def test_codex_auth_unhealthy_when_fallback_also_fails + runner = FakeRunner.new + runner.on([ "codex", "login", "status" ], CommandResult.new(stdout: "", stderr: "unknown command", exitstatus: 2)) + runner.on([ "codex", "auth", "status" ], CommandResult.new(stdout: "", stderr: "no", exitstatus: 1)) + with_replaced_singleton_method(Hive::AgentProfiles, :logged_in?, ->(_name, home: Dir.home) { true }) do + result = make_probe(runner: runner).probe_codex_auth(env: {}) + refute result.healthy? + assert_match(/auth status/, result.detail) + end + end + + def test_codex_auth_unhealthy_when_smoke_fails + runner = FakeRunner.new + runner.on([ "codex", "login", "status" ], ok_result) + runner.on([ "codex", "exec" ], CommandResult.new(stdout: "", stderr: "boom", exitstatus: 1)) + with_replaced_singleton_method(Hive::AgentProfiles, :logged_in?, ->(_name, home: Dir.home) { true }) do + result = make_probe(runner: runner).probe_codex_auth(env: {}) + refute result.healthy? + assert_match(/smoke test failed/, result.detail) + end + end + + # ── claude launcher ──────────────────────────────────────────────────── + + def test_claude_launcher_unhealthy_when_wrapper_missing + result = make_probe.probe_claude_launcher(wrapper_path: "/nonexistent/wrapper.sh", code_fingerprint: nil) + refute result.healthy? + assert_match(/wrapper missing/, result.detail) + end + + def test_claude_launcher_unhealthy_when_tmux_not_present + with_wrapper_file do |wrapper| + with_replaced_singleton_method(Hive::ClaudeLauncher, :tmux_status, -> { [ :missing, "tmux not runnable" ] }) do + result = make_probe(config: TMUX_CFG).probe_claude_launcher(wrapper_path: wrapper, code_fingerprint: nil) + refute result.healthy? + assert_match(/tmux not runnable/, result.detail) + end + end + end + + def test_claude_launcher_unhealthy_when_version_check_raises + with_wrapper_file do |wrapper| + profile = fake_claude_profile(check_raises: Hive::AgentError.new("claude binary not runnable: claude")) + with_replaced_singleton_method(Hive::AgentProfiles, :lookup, ->(_name, cfg: nil) { profile }) do + result = make_probe.probe_claude_launcher(wrapper_path: wrapper, code_fingerprint: nil) + refute result.healthy? + assert_match(/not runnable/, result.detail) + end + end + end + + def test_claude_launcher_unhealthy_when_code_fingerprint_mismatches + with_wrapper_file do |wrapper| + profile = fake_claude_profile + with_replaced_singleton_method(Hive::AgentProfiles, :lookup, ->(_name, cfg: nil) { profile }) do + result = make_probe.probe_claude_launcher(wrapper_path: wrapper, code_fingerprint: "deadbeef") + refute result.healthy? + assert_match(/does not match/, result.detail) + end + end + end + + def test_claude_launcher_healthy_when_all_checks_pass + with_wrapper_file do |wrapper| + profile = fake_claude_profile + with_replaced_singleton_method(Hive::AgentProfiles, :lookup, ->(_name, cfg: nil) { profile }) do + probe = make_probe(doctor_factory: fake_doctor([])) + result = probe.probe_claude_launcher(wrapper_path: wrapper, code_fingerprint: nil) + assert result.healthy?, result.detail + assert_match(/all green/, result.detail) + end + end + end + + def test_claude_launcher_reuses_precomputed_universal_result + # When the healer already ran the universal probe this cycle, the claude + # probe must reuse it instead of building a second `hive doctor` (R6). + doctor_builds = [] + doctor_class = Class.new do + define_method(:initialize) { |**| } + define_method(:call) { } + define_method(:rows) { [] } + end + factory = Class.new(doctor_class) + factory.define_singleton_method(:new) do |*args, **kwargs| + doctor_builds << :doctor + super(*args, **kwargs) + end + + probe = make_probe(doctor_factory: factory) + universal = Hive::Daemon::HealthProbe::Result.new(name: "universal", healthy: true, status: "healthy", detail: "ok") + with_wrapper_file do |wrapper| + profile = fake_claude_profile + with_replaced_singleton_method(Hive::AgentProfiles, :lookup, ->(_name, cfg: nil) { profile }) do + result = probe.probe_claude_launcher(wrapper_path: wrapper, code_fingerprint: nil, universal_result: universal) + assert result.healthy?, result.detail + assert_empty doctor_builds, "a precomputed universal result must skip the doctor build" + end + end + end + + # ── helper edge coverage ─────────────────────────────────────────────── + + def test_codex_home_and_bin_env_resolution + probe = make_probe + assert_equal "/opt/codex", probe.send(:codex_home, { "CODEX_HOME" => "/opt/codex" }) + assert_equal "/home/u", probe.send(:codex_home, { "HOME" => "/home/u" }) + assert_equal Dir.home, probe.send(:codex_home, {}) + assert_equal "/bin/codex", probe.send(:codex_bin, { "HIVE_CODEX_BIN" => "/bin/codex" }) + assert_equal "codex", probe.send(:codex_bin, {}) + end + + def test_login_status_ok_and_unknown_command_predicates + probe = make_probe + assert probe.send(:login_status_ok?, ok_result(stdout: "Logged in as user@example.com")) + refute probe.send(:login_status_ok?, ok_result(stdout: "not logged in")) + refute probe.send(:login_status_ok?, CommandResult.new(stdout: "", stderr: "", exitstatus: 1)) + + assert probe.send(:unknown_command?, CommandResult.new(stdout: "", stderr: "unknown command", exitstatus: 2)) + refute probe.send(:unknown_command?, ok_result) + end + + def test_claude_version_check_returns_unhealthy_on_generic_error + profile = fake_claude_profile(check_raises: RuntimeError.new("kaboom")) + with_replaced_singleton_method(Hive::AgentProfiles, :lookup, ->(_name, cfg: nil) { profile }) do + step = make_probe.send(:claude_version_check, {}) + refute step.ok + assert_match(/RuntimeError: kaboom/, step.detail) + end + end + + def test_code_fingerprint_match_returns_nil_when_fresh_digest_unavailable + with_replaced_singleton_method(Digest::SHA256, :file, ->(_path) { raise Errno::ENOENT }) do + detail = make_probe.send(:code_fingerprint_match, "anything") + assert_match(/could not compute/, detail) + end + end + + def test_stage_skill_failure_returns_nil_for_skillless_stage + probe = make_probe + assert_nil probe.send(:stage_skill_failure, "4-execute") + assert_nil probe.send(:stage_skill_failure, nil) + end + + def test_stage_skill_failure_rescues_lookup_error + with_replaced_singleton_method(Hive::AgentProfiles, :lookup, ->(*_args, **_kwargs) { raise Hive::ConfigError, "nope" }) do + failure = make_probe.send(:stage_skill_failure, "3-plan") + assert_match(/plan skill check failed/, failure) + end + end + + def test_stage_config_key_strips_numeric_prefix + probe = make_probe + assert_equal "execute", probe.send(:stage_config_key, "4-execute") + assert_nil probe.send(:stage_config_key, nil) + end + + def test_stage_config_key_normalizes_hyphenated_stage_to_underscore_key + probe = make_probe + assert_equal "open_pr", probe.send(:stage_config_key, "5-open-pr") + assert_equal "artifacts", probe.send(:stage_config_key, "7-artifacts") + assert_equal "finalize", probe.send(:stage_config_key, "8-finalize") + end + + def test_failing_desc_empty_for_no_failing_rows + assert_nil make_probe.send(:failing_desc, []) + end + + def test_wrapper_path_class_and_instance_resolve_consistently + assert_equal Hive::Daemon::HealthProbe.wrapper_path, make_probe.send(:default_wrapper_path) + assert File.file?(Hive::Daemon::HealthProbe.wrapper_path) + end +end diff --git a/test/unit/daemon/health_signals_test.rb b/test/unit/daemon/health_signals_test.rb new file mode 100644 index 000000000..3a47f1206 --- /dev/null +++ b/test/unit/daemon/health_signals_test.rb @@ -0,0 +1,187 @@ +require "test_helper" +require "tmpdir" +require "hive/daemon/health_signals" + +# The health-signal fingerprint must be deterministic (stable for unchanged +# inputs, flips only when a health-relevant input changes) and the per-tick +# cache must memoize probe results within a tick but recompute after clear. +class HiveDaemonHealthSignalsTest < Minitest::Test + include HiveTestHelper + + def fake_profile(name: :claude, bin: "claude") + profile = Object.new + profile.define_singleton_method(:name) { name } + profile.define_singleton_method(:bin) { bin } + profile + end + + def cfg + { + "daemon" => { "enabled" => true, "poll_interval_sec" => 30 }, + "execute" => { "agent" => "codex" }, + "plan" => { "agent" => "claude", "skill_by_agent" => { "claude" => "/plan" } } + } + end + + def env(home: "/home/u") + { "HOME" => home, "HIVE_CODEX_BIN" => "/bin/codex", "HIVE_CLAUDE_BIN" => "/bin/claude" } + end + + def with_wrapper + Dir.mktmpdir do |dir| + path = File.join(dir, "wrapper.sh") + File.write(path, "#!/bin/sh\n") + yield path + end + end + + def test_fingerprint_is_deterministic_for_unchanged_inputs + with_wrapper do |wrapper| + first = Hive::Daemon::HealthSignals.fingerprint( + cfg: cfg, project_root: nil, profile: fake_profile, wrapper_path: wrapper, env: env + ) + second = Hive::Daemon::HealthSignals.fingerprint( + cfg: cfg, project_root: nil, profile: fake_profile, wrapper_path: wrapper, env: env + ) + assert_equal first, second + end + end + + def test_fingerprint_changes_when_wrapper_mtime_changes + with_wrapper do |wrapper| + first = Hive::Daemon::HealthSignals.fingerprint( + cfg: cfg, project_root: nil, profile: fake_profile, wrapper_path: wrapper, env: env + ) + File.utime(Time.now + 10, Time.now + 10, wrapper) + second = Hive::Daemon::HealthSignals.fingerprint( + cfg: cfg, project_root: nil, profile: fake_profile, wrapper_path: wrapper, env: env + ) + refute_equal first, second + end + end + + def test_fingerprint_changes_when_config_slice_changes + with_wrapper do |wrapper| + first = Hive::Daemon::HealthSignals.fingerprint( + cfg: cfg, project_root: nil, profile: fake_profile, wrapper_path: wrapper, env: env + ) + changed = cfg.merge("execute" => { "agent" => "claude" }) + second = Hive::Daemon::HealthSignals.fingerprint( + cfg: changed, project_root: nil, profile: fake_profile, wrapper_path: wrapper, env: env + ) + refute_equal first, second + end + end + + def test_fingerprint_drops_nil_components + # A nil profile and a nil code digest (read failure) must not crash — + # they are dropped from the digest rather than hashed. + with_replaced_singleton_method(Digest::SHA256, :file, ->(_path) { raise Errno::ENOENT }) do + fingerprint = Hive::Daemon::HealthSignals.fingerprint( + cfg: cfg, project_root: nil, profile: nil, wrapper_path: "/nope", env: env + ) + refute_nil fingerprint + assert_equal 64, fingerprint.length + end + end + + def test_code_digest_returns_nil_on_read_failure + with_replaced_singleton_method(Digest::SHA256, :file, ->(_path) { raise Errno::ENOENT }) do + assert_nil Hive::Daemon::HealthSignals.code_digest + end + refute_nil Hive::Daemon::HealthSignals.code_digest + end + + def test_config_slice_handles_missing_daemon_and_stages + refute_nil Hive::Daemon::HealthSignals.config_slice({}) + slice = Hive::Daemon::HealthSignals.config_slice(cfg) + assert_includes slice, "/plan" + assert_includes slice, "codex" + end + + def test_env_slice_includes_configured_keys + slice = Hive::Daemon::HealthSignals.env_slice(env) + assert_includes slice, "/bin/codex" + assert_includes slice, "/bin/claude" + end + + def test_codex_login_state_resolution + with_replaced_singleton_method(Hive::AgentProfiles, :logged_in?, ->(_name, home: Dir.home) { home == "/home/u" }) do + assert_equal "logged-out", Hive::Daemon::HealthSignals.codex_login_state({ "CODEX_HOME" => "/x" }) + assert_equal "logged-in", Hive::Daemon::HealthSignals.codex_login_state({ "HOME" => "/home/u" }) + assert_equal "logged-out", Hive::Daemon::HealthSignals.codex_login_state({}) + end + end + + def test_wrapper_state_missing_for_absent_file + assert_equal "missing", Hive::Daemon::HealthSignals.wrapper_state("/nonexistent/wrapper.sh") + with_wrapper do |wrapper| + assert_match(/\d+:\d+/, Hive::Daemon::HealthSignals.wrapper_state(wrapper)) + end + end + + def test_skill_inventory_lists_installed_skill_markers + Dir.mktmpdir do |dir| + skill_dir = File.join(dir, ".claude", "skills", "ce-work") + FileUtils.mkdir_p(skill_dir) + File.write(File.join(skill_dir, "SKILL.md"), "# skill\n") + inventory = Hive::Daemon::HealthSignals.skill_inventory(dir, env: env(home: dir)) + assert_includes inventory, "ce-work=" + end + end + + def test_skill_inventory_empty_for_nil_project_root + assert_equal "", Hive::Daemon::HealthSignals.skill_inventory(nil, env: env(home: "/nonexistent-home")) + end + + def test_profile_state_nil_and_present + assert_nil Hive::Daemon::HealthSignals.profile_state(nil) + assert_equal "claude:claude", Hive::Daemon::HealthSignals.profile_state(fake_profile) + end + + def test_canonical_serialization_is_order_independent + assert_equal Hive::Daemon::HealthSignals.canonical({ "b" => 1, "a" => 2 }), + Hive::Daemon::HealthSignals.canonical({ "a" => 2, "b" => 1 }) + assert_equal "1|2", Hive::Daemon::HealthSignals.canonical([ 1, 2 ]) + assert_equal "", Hive::Daemon::HealthSignals.canonical(nil) + assert_equal "x", Hive::Daemon::HealthSignals.canonical("x") + end + + def test_stat_token_rescues_missing_file + assert_equal "?", Hive::Daemon::HealthSignals.stat_token("/nonexistent/skill/SKILL.md") + end + + # ── cache lifecycle ──────────────────────────────────────────────────── + + def test_cache_memoizes_and_recomputes_after_clear + cache = Hive::Daemon::HealthSignalCache.new + key = [ "implementer_failed", "p" ] + assert cache.changed_since_last?(key, "fp-1") + + cache.store(key, fingerprint: "fp-1", results: [ :r1 ]) + refute cache.changed_since_last?(key, "fp-1") + assert_equal({ fingerprint: "fp-1", results: [ :r1 ] }, cache.fetch(key)) + + assert cache.changed_since_last?(key, "fp-2"), "a changed fingerprint must re-probe" + + cache.clear + assert cache.changed_since_last?(key, "fp-2"), "after clear the next tick must re-probe" + assert_nil cache.fetch(key) + end + + def test_cache_fingerprint_memoizes_per_key_and_recomputes_after_clear + cache = Hive::Daemon::HealthSignalCache.new + key = [ "claude_launch_failed", "p" ] + calls = 0 + + first = cache.fingerprint(key) { calls += 1; "fp-1" } + second = cache.fingerprint(key) { calls += 1; "fp-2" } + assert_equal "fp-1", first + assert_equal "fp-1", second + assert_equal 1, calls, "the fingerprint block must run once per key per tick" + + cache.clear + assert_equal "fp-3", cache.fingerprint(key) { "fp-3" }, + "after clear the next tick recomputes the fingerprint" + end +end diff --git a/test/unit/daemon/recoverable_signatures_test.rb b/test/unit/daemon/recoverable_signatures_test.rb new file mode 100644 index 000000000..25de13337 --- /dev/null +++ b/test/unit/daemon/recoverable_signatures_test.rb @@ -0,0 +1,140 @@ +require "test_helper" +require "tmpdir" +require "hive/daemon/recoverable_signatures" +require "hive/daemon/status_consumer" + +# The diagnostic-signature classifier must admit ONLY the two allowlisted +# recoverable shapes (Codex-401 implementer_failed + claude_launch_failed) +# and reject everything else, reading evidence from marker attrs and log +# tails without ever blocking on non-regular files. +class HiveDaemonRecoverableSignaturesTest < Minitest::Test + include HiveTestHelper + + Signatures = Hive::Daemon::RecoverableSignatures + Row = Hive::Daemon::StatusConsumer::Row + + def make_row(marker_attrs:, folder: nil, slug: "s", marker: "error") + Row.new( + project: "p", slug: slug, stage: "4-execute", workflow: nil, + marker: marker, marker_attrs: marker_attrs, folder: folder, + state_file: nil, state_file_mtime: nil, action: "error", + suggested_command: nil, claude_pid_alive: nil, live_task_lock: false, + diagnostic: nil + ) + end + + # ── text predicate ───────────────────────────────────────────────────── + + def test_codex_auth_text_matches_literal_and_401_auth_shapes + assert Signatures.codex_auth_text?("codex exec returned 401 Missing bearer/basic auth") + assert Signatures.codex_auth_text?("401 Unauthorized: missing bearer token") + assert Signatures.codex_auth_text?("error 401 with basic auth required") + assert Signatures.codex_auth_text?("MISSING BEARER/BASIC AUTH") # case-insensitive literal + assert Signatures.codex_auth_text?("HTTP 401 Bearer token expired") + refute Signatures.codex_auth_text?("some 401 error without auth words") + refute Signatures.codex_auth_text?("missing bearer token") # no 401 + refute Signatures.codex_auth_text?("") + end + + # ── classification ───────────────────────────────────────────────────── + + def test_classify_codex_auth_from_marker_message + row = make_row(marker_attrs: { "reason" => "implementer_failed", "message" => "codex exec returned 401 Missing bearer/basic auth" }) + assert_equal :codex_auth, Signatures.classify(row) + end + + def test_classify_codex_auth_from_log_tail + Dir.mktmpdir do |dir| + log_dir = File.join(dir, "logs") + FileUtils.mkdir_p(log_dir) + File.write(File.join(log_dir, "implementer.log"), "codex exec returned 401 missing bearer token\n") + row = make_row(marker_attrs: { "reason" => "implementer_failed", "message" => "exit_code=1" }, folder: dir) + assert_equal :codex_auth, Signatures.classify(row) + end + end + + def test_classify_claude_launch_failed + row = make_row(marker_attrs: { "reason" => "claude_launch_failed", "message" => "claude interactive prompt did not become ready" }) + assert_equal :claude_launcher, Signatures.classify(row) + end + + def test_classify_rejects_unknown_implementer_failed + row = make_row(marker_attrs: { "reason" => "implementer_failed", "message" => "exit_code=1" }) + assert_nil Signatures.classify(row) + end + + def test_classify_rejects_other_reasons + row = make_row(marker_attrs: { "reason" => "git_status_failed", "message" => "boom" }) + assert_nil Signatures.classify(row) + end + + def test_codex_auth_failure_and_claude_launch_signature_predicates + auth = make_row(marker_attrs: { "reason" => "implementer_failed", "message" => "401 Missing bearer/basic auth" }) + assert Signatures.codex_auth_failure?(auth) + refute Signatures.claude_launch_signature?(auth) + + launch = make_row(marker_attrs: { "reason" => "claude_launch_failed" }) + assert Signatures.claude_launch_signature?(launch) + refute Signatures.codex_auth_failure?(launch) + end + + # ── evidence / log tails ─────────────────────────────────────────────── + + def test_state_log_tails_reads_global_per_slug_logs + Dir.mktmpdir do |dir| + state_home = File.join(dir, "state") + log_dir = File.join(state_home, "logs", "s") + FileUtils.mkdir_p(log_dir) + File.write(File.join(log_dir, "run.log"), "401 Missing bearer/basic auth\n") + row = make_row(marker_attrs: { "reason" => "implementer_failed" }, slug: "s") + tails = Signatures.state_log_tails(row, state_home: state_home) + assert_includes tails.join, "401 Missing bearer/basic auth" + end + end + + def test_state_log_tails_returns_empty_without_slug + row = Row.new( + project: "p", slug: nil, stage: "4-execute", workflow: nil, marker: "error", + marker_attrs: {}, folder: nil, state_file: nil, state_file_mtime: nil, + action: "error", suggested_command: nil, claude_pid_alive: nil, + live_task_lock: false, diagnostic: nil + ) + assert_equal [], Signatures.state_log_tails(row, state_home: "/tmp/whatever") + end + + def test_task_log_tails_returns_empty_without_folder + assert_equal [], Signatures.task_log_tails(make_row(marker_attrs: {})) + end + + def test_tail_file_skips_non_regular_files + Dir.mktmpdir do |dir| + regular = File.join(dir, "regular.log") + File.write(regular, "hello") + assert_includes Signatures.tail_file(regular), "hello" + + # A directory named like a log is not a regular file and must be skipped. + dir_as_log = File.join(dir, "dir.log") + FileUtils.mkdir_p(dir_as_log) + assert_nil Signatures.tail_file(dir_as_log) + end + end + + def test_marker_attrs_handles_non_hash_and_missing + row = make_row(marker_attrs: nil) + assert_equal "", Signatures.marker_reason(row) + assert_equal "", Signatures.marker_message(row) + + plain = Object.new + assert_equal({}, Signatures.marker_attrs(plain)) + end + + def test_tail_file_rescues_read_failure + Dir.mktmpdir do |dir| + path = File.join(dir, "regular.log") + File.write(path, "hello") + with_replaced_singleton_method(Hive::DiagnosticHelpers, :tail_file, ->(_path) { raise Errno::EIO, "boom" }) do + assert_nil Signatures.tail_file(path) + end + end + end +end diff --git a/test/unit/daemon/stale_agent_healer_probe_retry_test.rb b/test/unit/daemon/stale_agent_healer_probe_retry_test.rb new file mode 100644 index 000000000..ae4043e49 --- /dev/null +++ b/test/unit/daemon/stale_agent_healer_probe_retry_test.rb @@ -0,0 +1,570 @@ +require "test_helper" +require "tmpdir" +require "time" +require "hive/markers" +require "hive/daemon/health_probe" +require "hive/daemon/health_signals" +require "hive/daemon/stale_agent_healer" +require "hive/daemon/status_consumer" + +# Probe-gated auto-retry of the two v1 recoverable terminal-error markers. +# The healer must clear `implementer_failed` (Codex 401) and +# `claude_launch_failed` ONLY when the universal + reason probes pass, the +# health fingerprint changed (or a fallback elapsed), the work area is safe, +# and the stricter retry budget (2) with backoff is respected. +class HiveDaemonStaleAgentHealerProbeRetryTest < Minitest::Test + include HiveTestHelper + + Row = Hive::Daemon::StatusConsumer::Row + Result = Hive::Daemon::HealthProbe::Result + + HEADLESS_CFG = { "claude" => { "mode" => "headless" } }.freeze + + class FakeController + attr_reader :observed_mtimes + + def initialize(running_pairs: []) + @running = running_pairs + @observed_mtimes = [] + end + + def running_task?(project:, slug:) + @running.include?([ project, slug ]) + end + + def observe_state_file_mtime(project:, slug:, mtime:) + @observed_mtimes << { project: project, slug: slug, mtime: mtime } + end + end + + class FakeLogger + attr_reader :events + + def initialize + @events = [] + end + + def event(name, **attrs) + @events << [ name, attrs ] + end + end + + class FakeRequestQueue + attr_reader :requests + + def initialize + @requests = [] + end + + def write_request!(**kwargs) + @requests << kwargs + "fake-req-#{@requests.size}" + end + end + + class FakeProbe + attr_accessor :universal, :codex, :claude + attr_reader :calls, :last_claude_universal + + def initialize + @calls = [] + @universal = Result.new(name: "universal", healthy: true, status: "healthy", detail: "ok") + @codex = Result.new(name: "codex_auth", healthy: true, status: "healthy", detail: "ok") + @claude = Result.new(name: "claude_launcher", healthy: true, status: "healthy", detail: "ok") + end + + def probe_universal(config: nil, project_root: nil, stage: nil) + @calls << :universal + @universal + end + + def probe_codex_auth(env: nil) + @calls << :codex + @codex + end + + def probe_claude_launcher(config: nil, project_root: nil, env: nil, wrapper_path: nil, code_fingerprint: nil, stage: nil, universal_result: nil) + @calls << :claude + @last_claude_universal = universal_result + @claude + end + end + + NOW = Time.utc(2026, 7, 1, 12, 0, 0) + + def setup + @logger = FakeLogger.new + @controller = FakeController.new + @request_queue = FakeRequestQueue.new + @probe = FakeProbe.new + end + + # `worktree_clean:` stubs the real worktree predicate so most tests can + # focus on the probe/budget/safety wiring; pass `:real` to exercise the + # actual git-porcelain predicate. + def build_healer(auto_retry_enabled: true, probe: @probe, signal_cache: nil, env: {}, worktree_clean: true, project_config_resolver: ->(_name) { [ HEADLESS_CFG, nil ] }) + healer = Hive::Daemon::StaleAgentHealer.new( + controller: @controller, + logger: @logger, + grace_sec: 300, + request_queue: @request_queue, + probe: probe, + signal_cache: signal_cache || Hive::Daemon::HealthSignalCache.new, + auto_retry_enabled: auto_retry_enabled, + project_config_resolver: project_config_resolver, + env: env + ) + healer.define_singleton_method(:worktree_clean?, ->(_row) { worktree_clean }) unless worktree_clean == :real + healer + end + + def with_marker_file(reason:, message:, stage: "4-execute", marker_id: "err-a") + Dir.mktmpdir do |dir| + state_file = File.join(dir, "task.md") + attrs = "reason=#{reason} marker_id=#{marker_id}" + attrs += " message=\"#{message}\"" if message + File.write(state_file, "# task\n\n\n") + yield state_file, dir + end + end + + def make_row(state_file, reason:, message:, stage: "4-execute", marker_id: "err-a", folder: nil, mtime: NOW - 1000, live_task_lock: false) + attrs = { "reason" => reason, "marker_id" => marker_id } + attrs["message"] = message if message + Row.new( + project: "p", slug: "s", stage: stage, workflow: nil, + marker: "error", marker_attrs: attrs, + folder: folder || File.dirname(state_file), state_file: state_file, + state_file_mtime: mtime, action: "error", suggested_command: nil, + claude_pid_alive: nil, live_task_lock: live_task_lock, diagnostic: nil + ) + end + + def heal(rows, now: NOW) + build_healer.heal(rows, now: now) + end + + # ── happy paths ──────────────────────────────────────────────────────── + + def test_codex_auth_failure_clears_and_emits_auto_retry + with_marker_file(reason: "implementer_failed", message: "401 Missing bearer/basic auth") do |state_file, dir| + row = make_row(state_file, reason: "implementer_failed", message: "401 Missing bearer/basic auth") + heal([ row ]) + + assert Hive::Markers.current(state_file).none?, "the Codex-401 marker must clear" + healed = @logger.events.find { |name, _| name == :marker_healed } + assert healed, "expected marker_healed, got: #{@logger.events.inspect}" + assert_equal "codex_auth_recovered", healed[1][:reason] + assert_equal "implementer_failed", healed[1][:marker_reason] + assert_equal 1, healed[1][:attempts] + + events_line = File.read(File.join(dir, "events.jsonl")) + assert_includes events_line, "\"event_type\":\"auto_retry\"" + assert_includes events_line, "implementer_failed" + end + end + + def test_claude_launch_failed_clears_when_probes_green + with_marker_file(reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") do |state_file, _dir| + row = make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") + heal([ row ]) + + assert Hive::Markers.current(state_file).none? + healed = @logger.events.find { |name, _| name == :marker_healed } + assert healed + assert_equal "claude_launcher_recovered", healed[1][:reason] + end + end + + def test_claude_probe_cycle_runs_universal_probe_once_and_reuses_result + with_marker_file(reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") do |state_file, _dir| + row = make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") + heal([ row ]) + + assert_equal 1, @probe.calls.count(:universal), + "the universal doctor probe must run once per claude probe cycle" + assert_equal 1, @probe.calls.count(:claude) + assert @probe.last_claude_universal, "claude probe must receive the precomputed universal result" + assert @probe.last_claude_universal.healthy? + end + end + + def test_plan_probe_gated_recovery_requeues_plan_rerun + Dir.mktmpdir do |dir| + state_file = File.join(dir, "plan.md") + File.write(state_file, <<~MD) + # Plan + + ## Overview + something + + + MD + row = make_row(state_file, reason: "claude_launch_failed", message: "x", stage: "3-plan", folder: dir) + build_healer.heal([ row ], now: NOW) + + assert Hive::Markers.current(state_file).none?, "a probe-gated 3-plan error must clear" + request = @request_queue.requests.find do |req| + req[:argv] == [ "hive", "plan", "s", "--project", "p", "--from", "3-plan" ] + end + assert request, + "3-plan recovery must enqueue an explicit plan rerun, got: #{@request_queue.requests.inspect}" + assert_equal "healer", request[:requestor] + end + end + + # ── classifier gating ────────────────────────────────────────────────── + + def test_unknown_implementer_failed_stays_parked + with_marker_file(reason: "implementer_failed", message: "exit_code=1") do |state_file, _dir| + row = make_row(state_file, reason: "implementer_failed", message: "exit_code=1") + heal([ row ]) + + refute @logger.events.any? { |name, _| name == :marker_healed }, + "an implementer_failed without a Codex-401 signature must stay parked" + assert_match(/ERROR/, File.read(state_file)) + end + end + + def test_other_reasons_never_enter_probe_path + with_marker_file(reason: "git_status_failed", message: "boom") do |state_file, _dir| + row = make_row(state_file, reason: "git_status_failed", message: "boom") + heal([ row ]) + + assert_empty @probe.calls, "non-recoverable reasons must not trigger probes" + assert_match(/ERROR reason=git_status_failed/, File.read(state_file)) + end + end + + # ── probe health gating ──────────────────────────────────────────────── + + def test_probe_not_healthy_blocks_retry_with_negative_event + @probe.codex = Result.new(name: "codex_auth", healthy: false, status: "unhealthy", detail: "smoke test failed") + with_marker_file(reason: "implementer_failed", message: "401 Missing bearer/basic auth") do |state_file, _dir| + row = make_row(state_file, reason: "implementer_failed", message: "401 Missing bearer/basic auth") + heal([ row ]) + + refute @logger.events.any? { |name, _| name == :marker_healed } + blocked = @logger.events.find { |name, _| name == :auto_retry_blocked } + assert blocked, "expected auto_retry_blocked, got: #{@logger.events.inspect}" + assert_equal "probe not healthy", blocked[1][:rationale] + assert_match(/ERROR reason=implementer_failed/, File.read(state_file)) + end + end + + # ── work area safety ─────────────────────────────────────────────────── + + def test_dirty_worktree_blocks_retry + healer = build_healer(worktree_clean: false) + with_marker_file(reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") do |state_file, _dir| + row = make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") + healer.heal([ row ], now: NOW) + + refute @logger.events.any? { |name, _| name == :marker_healed } + blocked = @logger.events.find { |name, _| name == :auto_retry_blocked } + assert blocked + assert_equal "unsafe work area", blocked[1][:rationale] + assert_match(/ERROR reason=claude_launch_failed/, File.read(state_file)) + end + end + + def test_safe_to_retry_true_for_non_worktree_stages + healer = build_healer(worktree_clean: :real) + row = make_row("/tmp/missing.md", reason: "claude_launch_failed", message: "x", stage: "3-plan") + assert healer.send(:safe_to_retry?, row) + end + + # ── answered-content safety (brainstorm/plan) ──────────────────────── + + def test_brainstorm_with_answered_content_blocks_retry + healer = build_healer(worktree_clean: :real) + Dir.mktmpdir do |dir| + state_file = File.join(dir, "brainstorm.md") + File.write(state_file, <<~MD) + ## Round 1 + ### Q1. What is the scope? + ### A1. + Build a CLI. + + + MD + row = make_row(state_file, reason: "claude_launch_failed", message: "x", stage: "2-brainstorm", folder: dir) + healer.heal([ row ], now: NOW) + + refute @logger.events.any? { |name, _| name == :marker_healed }, + "a brainstorm with answered content must never be auto-cleared" + blocked = @logger.events.find { |name, _| name == :auto_retry_blocked } + assert blocked, "expected auto_retry_blocked, got: #{@logger.events.inspect}" + assert_equal "unsafe work area", blocked[1][:rationale] + assert_match(/ERROR reason=claude_launch_failed/, File.read(state_file)) + end + end + + def test_brainstorm_without_answers_is_safe_to_retry + healer = build_healer(worktree_clean: :real) + Dir.mktmpdir do |dir| + state_file = File.join(dir, "brainstorm.md") + File.write(state_file, <<~MD) + ## Round 1 + ### Q1. What is the scope? + ### A1. + + + MD + row = make_row(state_file, reason: "claude_launch_failed", message: "x", stage: "2-brainstorm", folder: dir) + healer.heal([ row ], now: NOW) + + assert Hive::Markers.current(state_file).none?, "an answer-free brainstorm must clear" + end + end + + def test_plan_with_waiting_marker_blocks_retry + healer = build_healer(worktree_clean: :real) + Dir.mktmpdir do |dir| + state_file = File.join(dir, "plan.md") + File.write(state_file, <<~MD) + # Plan + + ## Overview + something + + + + MD + row = make_row(state_file, reason: "claude_launch_failed", message: "x", stage: "3-plan", folder: dir) + healer.heal([ row ], now: NOW) + + refute @logger.events.any? { |name, _| name == :marker_healed }, + "a plan with an in-flight WAITING answer must never be auto-cleared" + blocked = @logger.events.find { |name, _| name == :auto_retry_blocked } + assert_equal "unsafe work area", blocked[1][:rationale] + end + end + + # ── real worktree predicate ─────────────────────────────────────────── + + def test_worktree_clean_uses_git_porcelain + with_tmp_git_repo do |repo| + Dir.mktmpdir do |folder_dir| + File.write(File.join(folder_dir, "worktree.yml"), { "path" => repo }.to_yaml) + row = make_row(File.join(folder_dir, "task.md"), reason: "claude_launch_failed", message: "x", folder: folder_dir) + healer = build_healer(worktree_clean: :real) + assert healer.send(:worktree_clean?, row), "a clean git worktree must pass the safety gate" + assert healer.send(:safe_to_retry?, row) + + File.write(File.join(repo, "dirty.txt"), "residue") + refute healer.send(:worktree_clean?, row), "a dirty worktree must fail the safety gate" + end + end + end + + def test_worktree_clean_false_without_pointer + healer = build_healer(worktree_clean: :real) + row = make_row("/tmp/missing.md", reason: "claude_launch_failed", message: "x") + refute healer.send(:worktree_clean?, row) + refute healer.send(:safe_to_retry?, row) + end + + # ── retry budget / backoff ───────────────────────────────────────────── + + def test_third_failure_exhausts_budget + healer = build_healer + with_marker_file(reason: "implementer_failed", message: "401 Missing bearer/basic auth") do |state_file, _dir| + # Attempt 1 (immediate) + healer.heal([ make_row(state_file, reason: "implementer_failed", message: "401 Missing bearer/basic auth") ], now: NOW) + assert Hive::Markers.current(state_file).none? + File.write(state_file, "# task\n\n\n") + + # Attempt 2 (after backoff) + healer.heal([ make_row(state_file, reason: "implementer_failed", message: "401 Missing bearer/basic auth") ], now: NOW + 1800) + assert Hive::Markers.current(state_file).none? + File.write(state_file, "# task\n\n\n") + + # Attempt 3 → exhausted + healer.heal([ make_row(state_file, reason: "implementer_failed", message: "401 Missing bearer/basic auth") ], now: NOW + 3600) + + assert_match(/ERROR reason=implementer_failed/, File.read(state_file), + "the 3rd failure must park the marker for manual recovery") + assert_equal 2, @logger.events.count { |name, _| name == :marker_healed } + exhausted = @logger.events.find { |name, _| name == :marker_heal_exhausted } + assert exhausted + assert_equal "codex_auth_recovered", exhausted[1][:reason] + assert_equal 2, exhausted[1][:attempts] + assert_equal 2, exhausted[1][:max_attempts] + assert_equal "per_process", exhausted[1][:budget_scope] + assert_match(/hive markers clear/, exhausted[1][:remediation]) + end + end + + def test_second_retry_within_backoff_is_throttled_as_unchanged_signal + cache = Hive::Daemon::HealthSignalCache.new + healer = build_healer(signal_cache: cache) + with_marker_file(reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") do |state_file, _dir| + healer.heal([ make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") ], now: NOW) + assert Hive::Markers.current(state_file).none? + File.write(state_file, "# task\n\n\n") + + # Mirror the dispatcher's per-tick memo clear so this exercises the + # production path, not the in-process memo short-circuit. + cache.clear + # Same fingerprint, 10s later — within the 30 min backoff AND the + # fallback window, so the production path throttles as "signal + # unchanged" (the "backoff" rationale is unreachable in the + # unchanged-signal case). + healer.heal([ make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") ], now: NOW + 10) + + assert_match(/ERROR reason=claude_launch_failed/, File.read(state_file), + "a second retry within backoff must stay parked") + assert_equal 1, @logger.events.count { |name, _| name == :marker_healed } + blocked = @logger.events.select { |name, _| name == :auto_retry_blocked } + assert blocked.any? { |_, attrs| attrs[:rationale] == "signal unchanged" }, + "the unchanged-signal case must throttle as 'signal unchanged', got: #{@logger.events.inspect}" + refute blocked.any? { |_, attrs| attrs[:rationale] == "backoff" }, + "the 'backoff' rationale must not fire in the unchanged-signal case" + end + end + + def test_changed_signal_within_backoff_allows_immediate_retry + # A changed health signal (the operator fixed the dependency) must + # bypass the 30-min backoff and allow the second retry immediately. + cfg = [ { "claude" => { "mode" => "headless" } } ] + resolver = ->(_name) { [ cfg[0], nil ] } + cache = Hive::Daemon::HealthSignalCache.new + healer = build_healer(project_config_resolver: resolver, signal_cache: cache) + with_marker_file(reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") do |state_file, _dir| + healer.heal([ make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") ], now: NOW) + assert Hive::Markers.current(state_file).none? + File.write(state_file, "# task\n\n\n") + + # The operator fixes the dependency: the config slice changes, so the + # health fingerprint changes. 10s later is still within the 30-min + # backoff, but a changed signal must be retried immediately. Clear the + # per-tick memo first to mirror the dispatcher's `@signal_cache.clear` + # so the fingerprint is recomputed (and the change detected). + cfg[0] = { "claude" => { "mode" => "headless" }, "plan" => { "agent" => "claude" } } + cache.clear + healer.heal([ make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") ], now: NOW + 10) + + assert Hive::Markers.current(state_file).none?, + "a changed health signal within backoff must allow the second retry" + assert_equal 2, @logger.events.count { |name, _| name == :marker_healed } + end + end + + # ── unchanged-signal throttle / rescues ────────────────────────────── + + def test_unchanged_signal_within_fallback_is_throttled + cache = Hive::Daemon::HealthSignalCache.new + healer = build_healer(signal_cache: cache) + with_marker_file(reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") do |state_file, _dir| + healer.heal([ make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") ], now: NOW) + assert Hive::Markers.current(state_file).none? + File.write(state_file, "# task\n\n\n") + + # Simulate the next daemon tick: clear the per-tick memo, but keep the + # retry state (same fingerprint, no fallback elapsed). + cache.clear + healer.heal([ make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") ], now: NOW) + + assert_match(/ERROR reason=claude_launch_failed/, File.read(state_file)) + blocked = @logger.events.find { |name, attrs| name == :auto_retry_blocked && attrs[:rationale] == "signal unchanged" } + assert blocked, "an unchanged signal within the fallback window must throttle, got: #{@logger.events.inspect}" + end + end + + def test_fingerprint_computed_once_per_tick_for_shared_reason_project + fingerprint_calls = 0 + original = Hive::Daemon::HealthSignals.method(:fingerprint) + with_replaced_singleton_method(Hive::Daemon::HealthSignals, :fingerprint, + ->(**kwargs) { fingerprint_calls += 1; original.call(**kwargs) }) do + with_marker_file(reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") do |state_file, dir| + row1 = make_row(state_file, reason: "claude_launch_failed", message: "claude interactive prompt did not become ready") + row2 = make_row(File.join(dir, "other.md"), reason: "claude_launch_failed", message: "claude interactive prompt did not become ready", marker_id: "err-b") + build_healer.heal([ row1, row2 ], now: NOW) + + assert_equal 1, fingerprint_calls, + "the health fingerprint must be computed once per (reason, project) per tick" + end + end + end + + def test_probe_path_logs_marker_heal_failed_when_clear_raises + healer = build_healer + with_marker_file(reason: "implementer_failed", message: "401 Missing bearer/basic auth") do |state_file, _dir| + original = Hive::Markers.method(:clear_current) + Hive::Markers.define_singleton_method(:clear_current) do |path, *args| + raise Errno::ENOSPC, "no space left" if path == state_file + + original.call(path, *args) + end + + begin + healer.heal([ make_row(state_file, reason: "implementer_failed", message: "401 Missing bearer/basic auth") ], now: NOW) + ensure + Hive::Markers.define_singleton_method(:clear_current, &original) + end + + failure = @logger.events.find { |name, _| name == :marker_heal_failed } + assert failure, "expected marker_heal_failed, got: #{@logger.events.inspect}" + assert_equal "codex_auth_recovered", failure[1][:reason] + assert_match(/ENOSPC/, failure[1][:error]) + end + end + + def test_resolve_project_config_default_paths + healer = build_healer(project_config_resolver: nil) + with_replaced_singleton_method(Hive::Config, :find_project, ->(_name) { nil }) do + assert_equal [ nil, nil ], healer.send(:resolve_project_config, "unknown") + end + + entry = { "path" => "/tmp/proj" } + with_replaced_singleton_method(Hive::Config, :find_project, ->(_name) { entry }) do + with_replaced_singleton_method(Hive::Config, :load, ->(_path) { { "daemon" => {} } }) do + assert_equal [ { "daemon" => {} }, "/tmp/proj" ], healer.send(:resolve_project_config, "p") + end + end + + with_replaced_singleton_method(Hive::Config, :find_project, ->(_name) { raise Hive::ConfigError, "boom" }) do + assert_equal [ nil, nil ], healer.send(:resolve_project_config, "p") + end + end + + def test_worktree_clean_rescues_git_failure + with_tmp_git_repo do |repo| + Dir.mktmpdir do |folder_dir| + File.write(File.join(folder_dir, "worktree.yml"), { "path" => repo }.to_yaml) + row = make_row(File.join(folder_dir, "task.md"), reason: "claude_launch_failed", message: "x", folder: folder_dir) + healer = build_healer(worktree_clean: :real) + with_replaced_singleton_method(Open3, :capture3, ->(*_args) { raise Errno::EIO, "boom" }) do + refute healer.send(:worktree_clean?, row) + end + end + end + end + + def test_worktree_path_for_rescues_pointer_error + healer = build_healer(worktree_clean: :real) + row = make_row("/tmp/task.md", reason: "claude_launch_failed", message: "x") + with_replaced_singleton_method(Hive::Worktree, :read_pointer, ->(_folder) { raise Errno::EACCES, "nope" }) do + assert_nil healer.send(:worktree_path_for, row) + end + end + + # ── kill-switch ──────────────────────────────────────────────────────── + + def test_kill_switch_disables_probe_retry + @healer_off = Hive::Daemon::StaleAgentHealer.new( + controller: @controller, logger: @logger, grace_sec: 300, + request_queue: @request_queue, probe: @probe, auto_retry_enabled: false, + project_config_resolver: ->(_name) { [ HEADLESS_CFG, nil ] } + ) + with_marker_file(reason: "implementer_failed", message: "401 Missing bearer/basic auth") do |state_file, _dir| + row = make_row(state_file, reason: "implementer_failed", message: "401 Missing bearer/basic auth") + @healer_off.heal([ row ], now: NOW) + + refute @logger.events.any? { |name, _| name == :marker_healed } + assert_empty @probe.calls, "kill-switch off must skip the probe path entirely" + assert_match(/ERROR reason=implementer_failed/, File.read(state_file)) + end + end +end