diff --git a/config.example.yml b/config.example.yml index 808c665a3..6f7b03af6 100644 --- a/config.example.yml +++ b/config.example.yml @@ -1,6 +1,26 @@ --- registered_projects: [] +# Daemon settings (ADR-024). Per-project `daemon.enabled` lives in each +# project's .hive-state/config.yml; the block below is the GLOBAL daemon +# config. See wiki/modules/daemon.md for the full wiring reference. +# +# daemon.auto_retry gates the daemon's health-gated auto-retry of parked +# terminal ERROR markers whose reason is on the fixed v1 allowlist: +# - implementer_failed where the failure was a Codex auth failure +# (401 / missing bearer / basic auth — or provider=codex on new markers) +# - claude_launch_failed (launcher/wrapper/readiness problems) +# A marker is only cleared when the health probes pass (codex login status +# + read-only exec smoke test; claude wrapper + ready-detector + --version) +# AND a safety guard proves no user work would be discarded. Max 2 +# auto-retries per task/reason (immediate first retry, 30-minute backoff +# before the second); exhaustion persists across restarts and parks the +# task until a manual `hive markers clear`. Set enabled: false to turn the +# whole behavior off (module inert, zero log events). +daemon: + auto_retry: + enabled: true + # Optional hosted screenshot links for 7-artifacts visual demos. # Run `hive connect screenote` to authorize uploads. HIVE_SCREENOTE_BASE_URL # can override the default service URL for staging/self-hosted deployments. diff --git a/lib/hive/config.rb b/lib/hive/config.rb index c686876bc..5b518dbb4 100644 --- a/lib/hive/config.rb +++ b/lib/hive/config.rb @@ -355,7 +355,20 @@ module Hive "child_kill_grace_sec" => 30, "child_verb_timeouts" => { "digest" => 3600, "answer-digest" => 3600 }, "log_max_bytes" => 10_485_760, - "log_max_files" => 5 + "log_max_files" => 5, + # Health-gated auto-retry of parked terminal ERROR markers whose + # reason is on the fixed v1 allowlist (implementer_failed with a + # Codex auth failure signature; claude_launch_failed). When the + # daemon's health probes (codex login+smoke, claude launcher checks) + # pass and a safety guard proves no user work would be discarded, + # the marker is cleared and the same stage re-run — mechanically + # identical to `hive markers clear` + rerun. Max 2 auto-retries per + # task/reason with a 30-minute second-attempt backoff; exhaustion is + # persisted and parks the task permanently until a manual clear. + # Set enabled: false to disable entirely (module inert, zero events). + "auto_retry" => { + "enabled" => true + } }, # Update flow (plan 2026-05-27-002). The daemon checks the latest # release on a throttled cadence and, on the install.sh channel, @@ -2225,6 +2238,8 @@ module Hive "(true / false); got #{autostart.inspect} (#{autostart.class})" end + validate_daemon_auto_retry!(daemon, source_path) + DAEMON_NUMERIC_BOUNDS.each do |key, min| value = daemon[key] next if value.nil? @@ -2311,6 +2326,28 @@ module Hive "#{describe_source(source_path)} and run `hive connect screenote`." end + # daemon.auto_retry is an optional block; today it carries exactly one + # knob, `enabled` (boolean kill switch for the health-gated auto-retry + # of parked recoverable ERROR markers). Nested-block shape keeps future + # per-reason tuning knobs addable without a schema break. + def validate_daemon_auto_retry!(daemon, source_path) + block = daemon["auto_retry"] + return if block.nil? + + unless block.is_a?(Hash) + raise ConfigError, + "daemon.auto_retry in #{describe_source(source_path)} must be a mapping " \ + "(e.g. { enabled: true }); got #{block.inspect} (#{block.class})" + end + + enabled = block["enabled"] + unless enabled.nil? || enabled == true || enabled == false + raise ConfigError, + "daemon.auto_retry.enabled in #{describe_source(source_path)} must be a boolean " \ + "(true / false); got #{enabled.inspect} (#{enabled.class})" + end + end + # R-02: `daemon.child_verb_timeouts` is an optional map of hive verb # (String) → timeout seconds (Integer >= 0). 0 disables the cap for # that verb. Reject anything else loudly so a typo'd YAML knob fails diff --git a/lib/hive/daemon/auto_retry_state.rb b/lib/hive/daemon/auto_retry_state.rb new file mode 100644 index 000000000..ae928c90e --- /dev/null +++ b/lib/hive/daemon/auto_retry_state.rb @@ -0,0 +1,251 @@ +require "json" +require "fileutils" +require "time" +require "digest" + +require "hive/paths" + +module Hive + module Daemon + # Persists the auto-retrier's budget/backoff/fingerprint ledger across + # daemon restarts, so "park permanently" after retry exhaustion means + # permanently — unlike the healer's in-process per-process budgets, + # which reset on every restart/SIGHUP (a deliberate upgrade enabled by + # the ledger; see [[modules/daemon]]). + # + # Persistence discipline mirrors DispatchBaselines exactly: JSON on + # disk with a `schema_version` envelope, atomic write (tempfile + + # fsync + rename), sibling `.lock` flock as defense in depth, and a + # FAIL-CLOSED load (corrupt / torn / newer-schema file degrades to an + # empty map with a warning event — never raises into boot). Worst case + # after a lost ledger is re-baselined budgets: at most 2 extra stage + # reruns per task/reason, the same accepted trade-off DispatchBaselines + # makes. + # + # Ledger keys are [project, slug, stage, reason_class] — the same keying + # rationale as the healer's error_auto_recovery_key: budget is keyed by + # REASON, not marker_id, so a fresh marker id cannot mint fresh budget. + class AutoRetryState + SCHEMA_VERSION = 1 + STALE_TMP_SEC = 60 + + # R5: max 2 auto-retries per task per reason. The FIRST healthy signal + # retries immediately; the SECOND requires SECOND_RETRY_BACKOFF_SEC to + # have elapsed since the first; a third failure parks permanently. + MAX_AUTO_RETRIES = 2 + SECOND_RETRY_BACKOFF_SEC = 1800 + + Entry = { + "attempt_count" => 0, + "last_fingerprint" => nil, + "last_cleared_at" => nil, + "exhausted" => false + }.freeze + + def self.default_path + File.join(Hive::Paths.state_home, "daemon_auto_retry_state.json") + end + + def self.key(project, slug, stage, reason_class) + [ project.to_s, slug.to_s, stage.to_s, reason_class.to_s ] + end + + def initialize(path: self.class.default_path, logger: nil) + @path = path ? File.expand_path(path) : nil + @logger = logger + @entries = {} + clean_orphaned_tmp_files! + end + + # Fail-closed load (DispatchBaselines discipline): missing / unreadable / + # non-JSON / wrong shape / unrecognized schema_version ⇒ empty map. + def load! + @loaded = true + @entries = {} + return @entries unless @path && File.exist?(@path) + + parsed = JSON.parse(File.read(@path)) + return warn_corrupt("root is not a Hash") unless parsed.is_a?(Hash) + + version = parsed["schema_version"] + return warn_corrupt("missing/invalid schema_version") unless version.is_a?(Integer) + + if version > SCHEMA_VERSION + # Newer hive wrote this: degrade to empty (re-baselined budgets) + # rather than crash boot. Same accepted trade-off as baselines. + return warn_corrupt("newer schema_version #{version} > supported #{SCHEMA_VERSION}") + end + + entries = parsed["entries"] + unless entries.is_a?(Array) + return warn_corrupt("entries is not an Array") + end + + entries.each do |entry| + next unless entry.is_a?(Hash) + + key = self.class.key( + entry["project"], entry["slug"], entry["stage"], entry["reason_class"] + ) + next if key.any?(&:empty?) + + @entries[key] = normalize_entry(entry) + end + @entries + rescue JSON::ParserError, TypeError, SystemCallError, IOError, ArgumentError => e + @logger&.event(:auto_retry_state_corrupt, + path: @path, error_class: e.class.name, message: e.message) + @entries = {} + @entries + end + + def entries + @entries + end + + # Decision helper encoding R5: + # - exhausted ⇒ false (permanent park until manual clear) + # - unchanged fingerprint since last attempt ⇒ false (no health + # signal changed — re-probing would just burn budget on the same + # environment) + # - first attempt (no prior fingerprint) ⇒ true (immediate retry on + # first healthy signal after failure) + # - second attempt ⇒ true only once SECOND_RETRY_BACKOFF_SEC has + # elapsed since last_cleared_at + def retry_allowed?(key:, current_fingerprint:, now: Time.now) + load_if_needed! + entry = @entries[key] + # No ledger history at all: this is the first healthy signal after + # the failure — retry immediately. + return true if entry.nil? || entry.empty? + return false if entry["exhausted"] + + last_fp = entry["last_fingerprint"].to_s + return true if last_fp.empty? + + return false if last_fp == current_fingerprint.to_s + + attempt_count = entry["attempt_count"].to_i + return true if attempt_count < 1 + + last_cleared_at = parse_time(entry["last_cleared_at"]) + return true if last_cleared_at.nil? + + (now - last_cleared_at) >= SECOND_RETRY_BACKOFF_SEC + end + + # Record one auto-retry for key. Increments attempt_count, stores the + # pre-clear health fingerprint and cleared-at timestamp, and flips + # `exhausted` at MAX_AUTO_RETRIES. Persists atomically. + def record_attempt!(key:, fingerprint:, now: Time.now) + load_if_needed! + entry = @entries[key] || fresh_entry + entry["attempt_count"] = entry["attempt_count"].to_i + 1 + entry["last_fingerprint"] = fingerprint.to_s + entry["last_cleared_at"] = now.utc.iso8601(6) + if entry["attempt_count"] >= MAX_AUTO_RETRIES + entry["exhausted"] = true + entry["exhausted_at"] = now.utc.iso8601(6) + end + @entries[key] = entry + persist! + entry + end + + def exhausted?(key) + load_if_needed! + entry = @entries[key] + !entry.nil? && entry["exhausted"] == true + end + + def attempt_count(key) + load_if_needed! + entry = @entries[key] + entry.nil? ? 0 : entry["attempt_count"].to_i + end + + private + + def load_if_needed! + load! unless @loaded + end + + def fresh_entry + Entry.dup + end + + def normalize_entry(entry) + { + "attempt_count" => entry["attempt_count"].to_i, + "last_fingerprint" => entry["last_fingerprint"].to_s, + "last_cleared_at" => entry["last_cleared_at"], + "exhausted" => entry["exhausted"] == true + } + end + + def persist! + return unless @path + + FileUtils.mkdir_p(File.dirname(@path)) + tmp_path = File.join( + File.dirname(@path), ".#{File.basename(@path)}.#{$$}.#{Thread.current.object_id}.tmp" + ) + File.open(tmp_path, "w") do |f| + f.write(JSON.generate(envelope)) + f.fsync + end + File.rename(tmp_path, @path) + fsync_dir(File.dirname(@path)) + ensure + FileUtils.rm_f(tmp_path) if defined?(tmp_path) && tmp_path && File.exist?(tmp_path) + end + + def envelope + entries_list = @entries.map do |key, entry| + project, slug, stage, reason_class = key + { + "project" => project, "slug" => slug, + "stage" => stage, "reason_class" => reason_class + }.merge(entry.slice("attempt_count", "last_fingerprint", + "last_cleared_at", "exhausted", "exhausted_at")) + end + { "schema_version" => SCHEMA_VERSION, "entries" => entries_list } + end + + def fsync_dir(dir) + Dir.open(dir) { |d| d.fsync } + rescue StandardError + nil + end + + def clean_orphaned_tmp_files! + return unless @path + + dir = File.dirname(@path) + return unless File.directory?(dir) + + Dir.glob(File.join(dir, ".#{File.basename(@path)}.*.tmp")).each do |orphan| + FileUtils.rm_f(orphan) if (Time.now - File.mtime(orphan)) > STALE_TMP_SEC + rescue SystemCallError + next + end + rescue StandardError + nil + end + + def warn_corrupt(reason) + @logger&.event(:auto_retry_state_corrupt, path: @path, reason: reason) + @entries = {} + @entries + end + + def parse_time(value) + return nil if value.nil? || value.to_s.empty? + + Time.parse(value.to_s) + rescue ArgumentError + nil + end + end + end +end diff --git a/lib/hive/daemon/dispatcher.rb b/lib/hive/daemon/dispatcher.rb index 1f272c697..3249e0ef3 100644 --- a/lib/hive/daemon/dispatcher.rb +++ b/lib/hive/daemon/dispatcher.rb @@ -12,6 +12,7 @@ require "hive/daemon/concurrency_controller" require "hive/daemon/child_supervisor" require "hive/daemon/status_consumer" require "hive/daemon/stale_agent_healer" +require "hive/daemon/recoverable_marker_retrier" require "hive/daemon/display_name_backfiller" require "hive/daemon/task_id_backfiller" require "hive/daemon/dispatch_request_queue" @@ -42,6 +43,16 @@ module Hive class Dispatcher attr_reader :controller, :supervisor, :logger + # Kill-switch read (R9/R10): `daemon.auto_retry.enabled`, defaulting to + # true via Config::DEFAULTS per the #255 convention. False ⇒ the + # retrier module is never constructed live and emits zero events. + def auto_retry_enabled + block = @daemon_cfg.fetch( + "auto_retry", Hive::Config::DEFAULTS.dig("daemon", "auto_retry") || {} + ) + block.is_a?(Hash) ? block.fetch("enabled", true) != false : true + end + # Stage dir whose `needs_input` rows carry a brainstorm Q&A file the # daemon gates auto-resume on (see `brainstorm_answers_pending?`). BRAINSTORM_STAGE_DIR = "2-brainstorm".freeze # coding-scoped: answer-pending daemon gate only parses coding brainstorm.md @@ -105,6 +116,18 @@ module Hive logger: @logger, grace_sec: agent_marker_grace_sec ) + # Health-gated auto-retry of parked terminal ERROR markers on the + # fixed v1 allowlist (codex-auth implementer failures, claude + # launcher failures). Global kill switch daemon.auto_retry.enabled; + # rebuilt alongside the healer on SIGHUP so a config flip takes + # effect within one tick. + @recoverable_marker_retrier = RecoverableMarkerRetrier.new( + controller: @controller, + logger: @logger, + enabled: auto_retry_enabled, + dispatch_request_state_home: dispatch_request_state_home, + dry_run: @dry_run + ) # 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 @@ -269,6 +292,23 @@ module Hive keeping_previous: true) end + # Health-gated auto-retry of parked recoverable ERROR markers, + # immediately AFTER the healer and BEFORE dispatch-request + # processing: a marker cleared + enqueued here is covered by the + # same-tick in-flight gate in process_dispatch_requests, so the row + # scan can't double-dispatch it. The retrier's allowlist is disjoint + # from every healer-handled reason (asserted by test), so the two + # passes never double-heal the same marker. + begin + @recoverable_marker_retrier.tick( + result.rows, now: now, legacy_layout_projects: @legacy_layout_projects + ) + rescue StandardError => e + @logger.event(:fatal, + message: "recoverable_marker_retrier raised: #{e.class}: #{e.message}", + keeping_previous: true) + end + # Self-heal tasks left showing their raw slug because name # generation never landed at `hive new`. Purely additive and # marker-free, so order relative to dispatch is irrelevant — but @@ -1831,6 +1871,17 @@ module Hive Hive::TaskAction::DEFAULT_AGENT_MARKER_GRACE_SEC ) ) + # Rebuild the retrier alongside the healer so a SIGHUP flip of + # daemon.auto_retry.enabled takes effect within one tick. The + # persisted AutoRetryState ledger is re-read from disk by the new + # instance, so exhaustion survives the rebuild. + @recoverable_marker_retrier = RecoverableMarkerRetrier.new( + controller: @controller, + logger: @logger, + enabled: auto_retry_enabled, + dispatch_request_state_home: @dispatch_request_state_home, + dry_run: @dry_run + ) # Rebuild alongside the healer on SIGHUP reload so a future # operator-tunable knob (e.g. max_per_tick) would take effect # within one tick; today it carries only the dry_run flag. diff --git a/lib/hive/daemon/health_probes.rb b/lib/hive/daemon/health_probes.rb new file mode 100644 index 000000000..04c7c6424 --- /dev/null +++ b/lib/hive/daemon/health_probes.rb @@ -0,0 +1,307 @@ +require "open3" +require "digest" + +require "hive/claude_launcher" +require "hive/agent_profiles" + +module Hive + module Daemon + # Timeout-bounded, memoized health probes gating every auto-retry of a + # parked recoverable marker (see RecoverableMarkerRetrier). + # + # Probe execution mix (R3): Hive-owned checks run in-process; external + # CLIs (`codex`) shell out with hard per-probe timeouts — the child is + # killed by process group on expiry so a hung CLI can never wedge a + # daemon tick. stdout/stderr are captured verbatim and carried in each + # result for audit records. Any timeout / spawn error / nonzero exit ⇒ + # `{healthy: false, detail: ...}` — probes NEVER raise into the tick. + # + # Results are memoized per probe name on this instance. The retrier + # constructs ONE fresh HealthProbes per full tick (that instance IS the + # opaque cache handle), so N parked markers evaluated in the same tick + # share a single probe set and no CLI call runs twice per tick. + class HealthProbes + CODEX_LOGIN_TIMEOUT_SEC = 15 + CODEX_SMOKE_TIMEOUT_SEC = 30 + CLAUDE_VERSION_TIMEOUT_SEC = 10 + + # Trivial, read-only single-token smoke prompt. Deliberately NOT + # something stateful — the probe must never mutate operator state. + CODEX_SMOKE_PROMPT = "Reply with the single word: ok".freeze + # `codex login status` prints this when a credential exists. + CODEX_LOGGED_IN_RE = /logged\s+in/i.freeze + + Result = Struct.new(:healthy, :probe, :detail, keyword_init: true) do + def healthy? + healthy == true + end + + def to_h + { healthy: healthy, probe: probe, detail: detail } + end + end + + # @param command_runner [Proc, nil] test seam: ->(argv, timeout_sec:) -> + # [stdout, stderr, status] where status is an Integer exit code, + # :timeout, or an Exception. Production uses #capture_with_timeout. + # @param wrapper_path [String, nil] override for the interactive + # claude wrapper location (tests). + # @param claude_bin [String] binary name/path probed for --version. + def initialize(command_runner: nil, wrapper_path: nil, claude_bin: "claude", + home: Dir.home) + @command_runner = command_runner + @wrapper_path = wrapper_path || begin + launcher_dir = File.dirname( + Hive::ClaudeLauncher.method(:wrapper_command).source_location.first + ) + File.expand_path("scripts/interactive_claude_wrapper.sh", launcher_dir) + end + @claude_bin = claude_bin + @home = home + @memo = {} + end + + # Lightweight doctor/agent-health precondition shared by both classes. + # In-process first: every required agent profile must be registered + # AND carry its login credential artifact. Shell-out to `hive doctor + # --json` is intentionally avoided here: the daemon already runs in a + # hive context, the credential check covers the auth dimension doctor + # cannot see, and keeping it in-process honors R3's cost budget. + def universal_healthy? + memoize(:universal) do + missing = %i[codex claude].filter_map do |name| + next name unless Hive::AgentProfiles.registered?(name) + next name unless Hive::AgentProfiles.logged_in?(name, home: @home) + + nil + end + if missing.empty? + Result.new(healthy: true, probe: :universal, detail: {}) + else + Result.new(healthy: false, probe: :universal, + detail: { reason: "agents_not_ready", agents: missing.map(&:to_s) }) + end + end + end + + # Codex auth health: `codex login status` reports logged in AND a tiny + # read-only `codex exec` smoke prompt exits 0. + def codex_auth_healthy? + memoize(:codex_auth) do + universal = universal_healthy? + return universal unless universal.healthy? + + out, err, status = run_command(%w[codex login status], CODEX_LOGIN_TIMEOUT_SEC) + unless status.is_a?(Integer) && status.zero? && out.to_s.match?(CODEX_LOGGED_IN_RE) + return Result.new( + healthy: false, probe: :codex_login_status, + detail: { exit_status: status_description(status), + stdout: clip(out), stderr: clip(err) } + ) + end + + out2, err2, status2 = run_command( + %w[codex exec] + [CODEX_SMOKE_PROMPT], CODEX_SMOKE_TIMEOUT_SEC + ) + if status2.is_a?(Integer) && status2.zero? + Result.new(healthy: true, probe: :codex_auth, + detail: { codex_exec_smoke: "ok", stdout: clip(out2) }) + else + Result.new( + healthy: false, probe: :codex_exec_smoke, + detail: { exit_status: status_description(status2), + stdout: clip(out2), stderr: clip(err2) } + ) + end + end + end + + # Claude launcher health: wrapper file exists at the active install + # path, the ready-detector classifies canonical ready-prompt fixtures, + # the active `claude` binary resolves and answers --version, plus the + # universal precondition. + def claude_launcher_healthy? + memoize(:claude_launcher) do + universal = universal_healthy? + return universal unless universal.healthy? + + return wrapper_missing_result unless File.file?(@wrapper_path) + + failure = ready_detector_failure + return failure if failure + + out, err, status = run_command([@claude_bin, "--version"], CLAUDE_VERSION_TIMEOUT_SEC) + if status.is_a?(Integer) && status.zero? + Result.new(healthy: true, probe: :claude_launcher, + detail: { version_output: clip(out, 200).strip }) + else + Result.new( + healthy: false, probe: :claude_version, + detail: { bin: @claude_bin, + exit_status: status_description(status), + stdout: clip(out), stderr: clip(err) } + ) + end + end + end + + # Health-signal inputs for AutoRetryState.fingerprint. Cheap stats + + # digests only — never reads credential CONTENTS (auth.json contributes + # mtime only; env vars contribute presence only). Sorted + SHA-256'd so + # any change to a contributing signal changes the fingerprint. + def fingerprint_inputs + { + wrapper_mtime: safe_stat(@wrapper_path)[:mtime].to_s, + wrapper_size: safe_stat(@wrapper_path)[:size].to_s, + codex_auth_mtime: safe_stat(File.join(@home, ".codex", "auth.json"))[:mtime].to_s, + claude_bin_path: which(@claude_bin).to_s, + claude_bin_mtime: safe_stat(which(@claude_bin))[:mtime].to_s, + openai_api_key_present: ENV.key?("OPENAI_API_KEY").to_s, + anthropic_api_key_present: ENV.key?("ANTHROPIC_API_KEY").to_s + } + end + + def self.fingerprint(inputs) + ::Digest::SHA256.hexdigest( + inputs.sort_by { |key, _| key.to_s } + .map { |key, value| "#{key}=#{value}" }.join("\n") + ) + end + + private + + def memoize(probe_name) + @memo.fetch(probe_name) { @memo[probe_name] = yield } + end + + def run_command(argv, timeout_sec) + if @command_runner + begin + @command_runner.call(argv, timeout_sec: timeout_sec) + rescue StandardError => e + ["", "", e] + end + else + self.class.capture_with_timeout(argv, timeout_sec: timeout_sec) + end + end + + # Spawn argv in its own process group; on expiry TERM then KILL the + # whole group so grandchildren die too. Returns [stdout, stderr, status] + # where status is Integer exit code, :timeout, or an Exception. + def self.capture_with_timeout(argv, timeout_sec:) + out = +"" + err = +"" + Open3.popen3(*argv, pgroup: true) do |stdin, stdout, stderr, waiter| + stdin.close rescue nil + readers = [ + Thread.new { stdout.each_line { |line| out << line } }, + Thread.new { stderr.each_line { |line| err << line } } + ] + if waiter.join(timeout_sec) + readers.each(&:join) + [out, err, waiter.value.exitstatus] + else + kill_group(waiter.pid) + readers.each(&:join) + [out, err, :timeout] + end + end + rescue StandardError => e + [out, err, e] + end + + def self.kill_group(pid) + Process.kill("TERM", -pid) + deadline = Time.now + 5 + sleep 0.05 while Time.now < deadline && group_alive?(pid) + Process.kill("KILL", -pid) if group_alive?(pid) + rescue StandardError + nil + end + + def self.group_alive?(pid) + Process.kill(0, -pid) + true + rescue Errno::ESRCH + false + rescue Errno::EPERM + true + rescue StandardError + false + end + + READY_FIXTURES = [ + ["ready_prompt_banner_and_caret", + "Claude Code v2.1.133\nTip: try refactor\n\n❯ Try \"refactor \""], + ["ready_prompt_with_stale_permission_scrollback", + "Claude Code v2.1.133\nDo you want to make this edit?\n❯ 1. Yes\n" \ + "Claude Code v2.1.133\n❯ Try \"refactor \""] + ].freeze + + NEGATIVE_FIXTURES = [ + ["trust_prompt_menu", + "Claude Code v2.1.133\nProceed with the action?\n❯ 1. Yes"], + ["stale_prompt_scrollback", + "Claude Code v2.1.133\n❯ Try \"refactor \"\n\nbackground indexing update"] + ].freeze + + # Reuse the production detector against canonical fixtures so a drift + # in the ready-prompt regex fails the launcher probe instead of + # silently launching into a dead pane. + def ready_detector_failure + READY_FIXTURES.each do |label, fixture| + unless Hive::ClaudeLauncher.claude_ready_prompt?(fixture) + return Result.new(healthy: false, probe: :claude_ready_detector, + detail: { fixture: label, expected: "ready" }) + end + end + NEGATIVE_FIXTURES.each do |label, fixture| + if Hive::ClaudeLauncher.claude_ready_prompt?(fixture) + return Result.new(healthy: false, probe: :claude_ready_detector, + detail: { fixture: label, expected: "not ready" }) + end + end + nil + end + + def wrapper_missing_result + Result.new(healthy: false, probe: :claude_wrapper_missing, + detail: { wrapper_path: @wrapper_path }) + end + + def safe_stat(path) + return {} unless path && File.exist?(path) + + { mtime: File.mtime(path).to_i, size: File.size(path) } + rescue SystemCallError + {} + end + + def which(bin) + return bin if bin.include?(File::SEPARATOR) + + ENV["PATH"].to_s.split(File::PATH_SEPARATOR).each do |dir| + candidate = File.join(dir, bin) + return candidate if File.file?(candidate) && File.executable?(candidate) + end + nil + rescue StandardError + nil + end + + def status_description(status) + case status + when Integer then status + when :timeout then "timeout" + when Exception then "#{status.class}: #{status.message}" + else status.inspect + end + end + + def clip(text, limit = 500) + text.to_s.strip.byteslice(0, limit).to_s + end + end + end +end diff --git a/lib/hive/daemon/logger.rb b/lib/hive/daemon/logger.rb index 6af59dc93..9eab8f219 100644 --- a/lib/hive/daemon/logger.rb +++ b/lib/hive/daemon/logger.rb @@ -50,6 +50,11 @@ module Hive marker_heal_failed marker_heal_exhausted marker_heal_observer_missing + auto_retry_cleared + auto_retry_blocked + auto_retry_exhausted + auto_retry_enqueue_failed + auto_retry_state_corrupt display_name_backfill update_available update_check_no_result diff --git a/lib/hive/daemon/recoverable_marker_retrier.rb b/lib/hive/daemon/recoverable_marker_retrier.rb new file mode 100644 index 000000000..3d4620c55 --- /dev/null +++ b/lib/hive/daemon/recoverable_marker_retrier.rb @@ -0,0 +1,300 @@ +require "hive/paths" +require "hive/markers" +require "hive/events" +require "hive/daemon/recovery_classifier" +require "hive/daemon/health_probes" +require "hive/daemon/auto_retry_state" +require "hive/daemon/retry_safety_guard" + +module Hive + module Daemon + # Health-gated auto-retry of parked terminal ERROR markers whose reason + # classifies onto the fixed v1 allowlist (RecoveryClassifier). On each + # full tick, parked markers are re-evaluated through a cheap→expensive + # gate chain; when every gate is green the marker is CLEARED (state- + # machine completion — never forward-advancement past human gates, per + # ADR-024) and the same stage is re-enqueued via the existing + # dispatch-request queue (`requestor=healer`), mechanically identical to + # a manual `hive markers clear` + stage rerun. No resume path exists: + # the rerun starts the stage from its beginning. + # + # Gate order (short-circuiting with recorded rationale): + # 1. Kill switch (`daemon.auto_retry.enabled`) — checked at + # construction; off ⇒ this module is never built / stays inert. + # 2. Healer-inherited guards: controller.running_task?, + # live_task_lock, legacy_stage_dirs. + # 3. Ledger exhaustion ⇒ silent skip (one-shot `auto_retry_exhausted` + # was emitted when the last attempt hit the limit). + # 4. Per-tick probe cache: a class already probed unhealthy this tick + # skips silently (throttling by construction). + # 5. Fingerprint compare: unchanged health signals ⇒ throttled + # negative event (at most once per (task, reason, fingerprint)). + # 6. Probes (universal, then class-specific) ⇒ unhealthy names the + # failing probe in a throttled negative event. + # 7. RetrySafetyGuard unsafe ⇒ negative event with its rationale. + # 8. All green ⇒ clear marker + enqueue rerun + record attempt + + # positive audit event on BOTH channels (daemon.log and the task's + # events.jsonl). + # + # Budgets (R5): max 2 auto-retries per (task, stage, reason); first + # healthy-signal change retries immediately, second requires a 30-min + # backoff; exhaustion persists across restarts via the AutoRetryState + # ledger — unlike the healer's in-process budgets, "park permanently" + # here really means permanently until `hive markers clear`. + class RecoverableMarkerRetrier + def initialize(controller:, logger:, request_queue: Hive::Daemon::DispatchRequestQueue, + enabled: true, state: nil, guard: nil, probes_factory: nil, + dispatch_request_state_home: nil, dry_run: false) + @controller = controller + @logger = logger + @request_queue = request_queue + @enabled = enabled + @state = state || Hive::Daemon::AutoRetryState.new(logger: logger) + @guard = guard || RetrySafetyGuard.new + # Production: a fresh HealthProbes per tick (the instance IS the + # per-tick memoization cache). Tests inject a factory returning + # stubbed probes. + @probes_factory = probes_factory || -> { HealthProbes.new(home: probes_home) } + @dispatch_request_state_home = dispatch_request_state_home + @dry_run = dry_run + # `[project, slug, stage, reason_class] => last-seen health + # fingerprint` — a throttled-negative event fires at most once per + # combination, and re-fires only when the fingerprint actually + # changes (a new daemon process starts empty, which is itself the + # low-frequency periodic fallback re-probe). + @negative_seen = {} + end + + def enabled? + @enabled == true + end + + # One full-tick evaluation sweep. A fresh HealthProbes instance IS the + # per-tick cache handle: N parked markers share one probe set. + def tick(rows, now: Time.now, legacy_layout_projects: {}) + return unless enabled? + + probes = @probes_factory.call + probed_unhealthy = {} + rows.each do |row| + evaluate_row(row, now: now, probes: probes, + probed_unhealthy: probed_unhealthy, + legacy_layout_projects: legacy_layout_projects) + end + rescue StandardError => e + # Never crash a tick — mirror the healer's defensive contract. + @logger.event(:fatal, + message: "recoverable_marker_retrier raised: #{e.class}: #{e.message}", + keeping_previous: true) + end + + private + + def evaluate_row(row, now:, probes:, probed_unhealthy:, legacy_layout_projects:) + return if legacy_layout_projects.key?(row.project) + return if @controller.running_task?(project: row.project, slug: row.slug) + return if row.live_task_lock == true + return unless row.marker.to_s.casecmp("error").zero? + + attrs = row.respond_to?(:marker_attrs) && row.marker_attrs.is_a?(Hash) ? row.marker_attrs : {} + reason_class = RecoveryClassifier.classify(row.marker, attrs) + return unless reason_class + + key = AutoRetryState.key(row.project, row.slug, row.stage, reason_class) + + # Exhaustion persists in the ledger; the exhausted event fired once + # at limit time, later ticks skip silently. + return if @state.exhausted?(key) + + # Per-tick probe cache: this class already came back unhealthy on + # this tick — skip silently rather than re-running CLI probes or + # spamming the log. + return if probed_unhealthy.key?(reason_class) + + fingerprint = HealthProbes.fingerprint(probes.fingerprint_inputs) + + unless @state.retry_allowed?(key: key, current_fingerprint: fingerprint, now: now) + throttled_negative(key, row, reason_class, fingerprint, + reason: "health_signal_unchanged_or_backoff") + return + end + + probe_result = run_probes(reason_class, probes) + unless probe_result.healthy? + probed_unhealthy[reason_class] = probe_result + throttled_negative(key, row, reason_class, fingerprint, + reason: "#{probe_result.probe} failed") + return + end + + safety = @guard.evaluate(stage: row.stage, project: row.project, + slug: row.slug, state_file: row.state_file) + unless safety.safe + throttled_negative(key, row, reason_class, fingerprint, + reason: "safety_guard: #{safety.rationale}") + return + end + + clear_marker_and_rerun(row, attrs, reason_class, key, fingerprint, probe_result, now) + rescue StandardError => e + @logger.event(:fatal, + message: "auto_retry row evaluation raised: #{e.class}: #{e.message}", + keeping_previous: true, + project: row.project, slug: row.slug, stage: row.stage) + end + + def run_probes(reason_class, probes) + universal = probes.universal_healthy? + return universal unless universal.healthy? + + case reason_class + when :codex_auth then probes.codex_auth_healthy? + when :claude_launcher then probes.claude_launcher_healthy? + else universal + end + end + + def clear_marker_and_rerun(row, attrs, reason_class, key, fingerprint, probe_result, now) + marker_reason = attrs["reason"].to_s + marker_id = attrs["marker_id"].to_s + match_attrs = { "reason" => marker_reason } + match_attrs["marker_id"] = marker_id.empty? ? nil : marker_id + + # Direct Markers rewrite under the markers lock — the healer's + # precedent. Match on observed marker_id where present so a stale + # status row can never clear a NEWER marker written between status + # snapshot and now. + return unless Hive::Markers.clear_current( + row.state_file, expected_name: :error, match_attrs: match_attrs + ) + + observe_pre_clear_mtime(row) + + entry = @state.record_attempt!(key: key, fingerprint: fingerprint, now: now) + attempts = entry["attempt_count"].to_i + exhausted_now = entry["exhausted"] == true + + enqueue_rerun(row, marker_reason: marker_reason, reason_class: reason_class, + marker_id: marker_id, fingerprint: fingerprint, + probe_result: probe_result, attempts: attempts, + exhausted_now: exhausted_now, now: now) + rescue StandardError => e + @logger.event(:fatal, + message: "auto_retry clear/rerun raised: #{e.class}: #{e.message}", + keeping_previous: true, + project: row.project, slug: row.slug, stage: row.stage) + end + + def enqueue_rerun(row, marker_reason:, reason_class:, marker_id:, fingerprint:, + probe_result:, attempts:, exhausted_now:, now:) + audit_attrs = { + project: row.project, slug: row.slug, stage: row.stage, + marker_id: marker_id, marker_reason: marker_reason, + reason_class: reason_class.to_s, + probes: probe_result.detail, + health_fingerprint: fingerprint, + attempt_count: attempts, + max_attempts: AutoRetryState::MAX_AUTO_RETRIES + } + + begin + request_id = @request_queue.write_request!( + project: row.project, slug: row.slug, + argv: request_argv(row), + requestor: "healer", + trigger: "auto_retry", + state_home: dispatch_request_state_home + ) + rescue StandardError => e + # Own rescue, NOT the caller's: the clear already SUCCEEDED, so a + # generic "retry next tick" would be a lie — the retrier can never + # re-match this row. Distinct event carries the manual re-entry + # command (web Retry button / bot Autofix also still work). + @logger.event(:auto_retry_enqueue_failed, + **audit_attrs, + error: "#{e.class}: #{e.message}", + remediation: manual_command(row)) + return + end + + @logger.event(:auto_retry_cleared, + **audit_attrs, + action: "cleared_and_enqueued", + request_id: request_id, + prior_marker: row.marker, + state_file: row.state_file) + Hive::Events.emit( + task_folder: row.folder, slug: row.slug, stage: row.stage, + event_type: :auto_retry, + agent: "daemon", + message: "auto-retry cleared ERROR reason=#{marker_reason} " \ + "(#{reason_class}); probes passed; stage re-queued " \ + "(attempt #{attempts}/#{AutoRetryState::MAX_AUTO_RETRIES})" + ) + + return unless exhausted_now + + @logger.event(:auto_retry_exhausted, + **audit_attrs, + budget_scope: "persistent", + suggested_next_action: "manual markers clear") + end + + # Same verbs the operator would type: `hive run ` for + # execute-class stages; the healer's explicit `--from 3-plan` re-entry + # where the empty-artifact plan stage needs the bespoke verb. + def request_argv(row) + if row.stage.to_s == "3-plan" # coding-scoped: healer re-enters coding plan verb + [ "hive", "plan", row.slug, "--project", row.project, "--from", "3-plan" ] # coding-scoped: healer re-enters coding plan verb + else + [ "hive", "run", row.slug, "--project", row.project ] + end + end + + def manual_command(row) + if row.stage.to_s == "3-plan" # coding-scoped: healer re-enters coding plan verb + "hive plan #{row.slug} --project #{row.project} --from 3-plan" + else + "hive run #{row.slug} --project #{row.project}" + end + end + + # Throttled negative decisions: at most one event per + # (task, reason, fingerprint); a changed fingerprint re-arms it. + def throttled_negative(key, row, reason_class, fingerprint, reason:) + seen = @negative_seen[key] + if seen == fingerprint + nil + else + @negative_seen[key] = fingerprint + @logger.event(:auto_retry_blocked, + project: row.project, slug: row.slug, stage: row.stage, + reason_class: reason_class.to_s, + health_fingerprint: fingerprint, + not_retried_because: reason) + end + end + + def observe_pre_clear_mtime(row) + return unless row.respond_to?(:state_file_mtime) + return unless @controller.respond_to?(:observe_state_file_mtime) + + @controller.observe_state_file_mtime( + project: row.project, slug: row.slug, mtime: row.state_file_mtime + ) + rescue StandardError + nil + end + + def dispatch_request_state_home + @dispatch_request_state_home || Hive::Paths.state_home + end + + def probes_home + ENV.fetch("HOME", Dir.home) + rescue StandardError + Dir.home + end + end + end +end diff --git a/lib/hive/daemon/recovery_classifier.rb b/lib/hive/daemon/recovery_classifier.rb new file mode 100644 index 000000000..4f828913e --- /dev/null +++ b/lib/hive/daemon/recovery_classifier.rb @@ -0,0 +1,93 @@ +module Hive + module Daemon + # Pure classifier for terminal ERROR markers: maps a marker to a v1 + # recoverable class (:codex_auth or :claude_launcher) or nil (manual). + # + # This is the fixed v1 allowlist behind the daemon's health-gated + # auto-retry (see [[modules/daemon]] and RecoverableMarkerRetrier). + # The allowlist is deliberately HARDCODED — no plugin surface, no + # config-driven reason list — so a new reason can never start + # self-clearing markers without a code change and a test asserting + # it stays disjoint from every StaleAgentHealer-owned reason. + # + # Classification contract: + # + # :codex_auth — `reason=implementer_failed` where the failure was + # a Codex auth problem. Two attribution paths: + # (a) NEW markers carry `provider=codex`, stamped by + # Stages::Execute#mark_implementer_failure at write + # time; (b) LEGACY markers (pre-stamping) must carry + # an explicit auth-failure signature in their + # `message` attr (401 / missing bearer / basic auth / + # unauthorized / MissingBearerToken). A legacy + # implementer_failed WITHOUT a recognized signature + # classifies nil — it stays manual, per R1. + # + # :claude_launcher — `reason=claude_launch_failed`. That reason is only + # ever written by the launcher rescue + # (Stages::Base#spawn_claude_with_tmux_marker!), so + # the reason itself is the diagnostic signature; no + # message matching is needed. + # + # Everything else — plain `implementer_failed` without a signature, + # business/test failures, review findings, merge conflicts, dirty-worktree + # reasons, exit_code=1 unknowns — classifies nil and is never touched by + # the retrier. + module RecoveryClassifier + # The complete v1 allowlist of recoverable classes. + REASON_CLASSES = %i[codex_auth claude_launcher].freeze + + MARKER_REASON_CLASSES = { + "implementer_failed" => :codex_auth, + "claude_launch_failed" => :claude_launcher + }.freeze + + # Codex auth-failure signatures observed in real ERROR messages + # (task-58-class failures). Matched case-insensitively against the + # marker's `message` attr for legacy markers that predate provider + # stamping. `MissingBearerToken`-family tokens are matched on the + # bare token so camelCase variants still hit. + CODEX_AUTH_SIGNATURE_RE = / + 401| + missing\ bearer| + basic\ auth| + unauthorized| + missingbearertoken + /ix.freeze + + module_function + + # @param marker_name [Symbol,String] the on-disk marker name + # (:error expected; anything else is never recoverable) + # @param attrs [Hash] parsed marker attributes (string or symbol keys) + # @return [Symbol, nil] :codex_auth, :claude_launcher, or nil + def classify(marker_name, attrs = {}) + return nil unless marker_name.to_s.casecmp("error").zero? + + attrs = (attrs || {}).to_h.transform_keys(&:to_s) + reason_class = MARKER_REASON_CLASSES[attrs["reason"].to_s] + return nil unless reason_class + + case reason_class + when :codex_auth then classify_codex_auth(attrs) + when :claude_launcher then :claude_launcher + end + end + + def classify_codex_auth(attrs) + # New markers carry provider=codex stamped at write time — trust it + # outright. Non-codex providers are excluded even if the message + # happens to contain a signature word. + return :codex_auth if attrs["provider"].to_s.casecmp("codex").zero? + return nil unless attrs["provider"].to_s.empty? + + # Legacy fallback: no provider stamp; require an explicit auth + # signature in the message. Markers without ANY recognized signature + # stay manual. + return :codex_auth if attrs["message"].to_s.match?(CODEX_AUTH_SIGNATURE_RE) + + nil + end + end + end +end diff --git a/lib/hive/daemon/retry_safety_guard.rb b/lib/hive/daemon/retry_safety_guard.rb new file mode 100644 index 000000000..a20fb21b1 --- /dev/null +++ b/lib/hive/daemon/retry_safety_guard.rb @@ -0,0 +1,187 @@ +require "open3" + +require "hive/brainstorm_parser" +require "hive/config" +require "hive/worktree" + +module Hive + module Daemon + # Fail-closed safety gate: refuses an auto-retry whenever replaying the + # stage could clobber authored (user) content. The retrier consults this + # AFTER health probes pass but BEFORE it clears the marker; any + # uncertainty resolves to UNSAFE — a missed auto-retry just stays parked + # red until a human clears, while a wrong auto-retry destroys work. + # + # Two stage families: + # + # Worktree-backed stages (4-execute and friends): locate the task + # worktree via the same canonical-root resolution `hive run` uses, + # then `git status --porcelain`. Clean ⇒ safe. Only files matching + # the conservative agent-residue allowlist ⇒ safe. Anything else, any + # git error, or an unresolvable worktree ⇒ unsafe. + # + # Brainstorm/plan stages (2-brainstorm / 3-plan): parse the state file + # with Hive::BrainstormParser; ANY answered `### A{n}` slot means a + # rerun would regenerate the file over the user's answers ⇒ unsafe. + # Zero questions / parse-empty ⇒ safe ("nothing authored exists" — + # mirrors the answers-pending fail-open semantics). + class RetrySafetyGuard + GIT_STATUS_TIMEOUT_SEC = 10 + + # Conservative residue allowlist for worktree-backed stages: files the + # agent pipeline itself generates and a rerun regenerates anyway. + # Deliberately tiny — anything not matching is unsafe. + RESIDUE_ALLOWLIST = [ + %r{\A\.hive(/|\z)}, # hive task internals (.hive/, .hive/…) + %r{\A\.hive-worktree-pointer\z} + ].freeze + + ANSWERED_STAGE_DIRS = %w[2-brainstorm 3-plan].freeze # coding-scoped: authored-content stages whose rerun would overwrite user answers + + # Coding stages whose rerun replays inside a git worktree. Explicit + # and conservative on purpose — an unrecognized stage fails closed to + # unsafe rather than guessing. + WORKTREE_BACKED_STAGES = %w[ + 4-execute 5-open-pr 6-review 7-artifacts 8-finalize + ].freeze + + Result = Struct.new(:safe, :rationale, keyword_init: true) + + def initialize(git_runner: nil, worktree_resolver: nil) + # Test seams: git_runner is ->(worktree_path) -> [stdout, stderr, + # status]; worktree_resolver is ->(project, slug) -> path-or-nil + # (production resolves via Hive::Config.find_project + + # Hive::Worktree.canonical_root). + @git_runner = git_runner + @worktree_resolver = worktree_resolver + end + + # @param stage [String] row.stage (stage dir name, e.g. "4-execute") # not-a-stage-ref: doc comment example + # @param project [String] project name (registry lookup for the root) + # @param slug [String] task slug + # @param state_file [String, nil] path to the stage's state file + # @return [Result] with safe: bool + rationale string + def evaluate(stage:, project:, slug:, state_file: nil) + stage = stage.to_s + if ANSWERED_STAGE_DIRS.include?(stage) + evaluate_answered_stage(stage, state_file) + elsif worktree_backed_stage?(stage) + evaluate_worktree_stage(project, slug) + else + Result.new(safe: false, rationale: "unknown_stage #{stage}") + end + rescue StandardError => e + Result.new(safe: false, rationale: "guard_error #{e.class}: #{e.message}") + end + + private + + def worktree_backed_stage?(stage) + # Only stages that actually own a git worktree are eligible for the + # porcelain check; anything unrecognized fails closed. Uses the + # runtime union of registered workflows' stage dirs (same source the + # dispatcher ranks against) intersected with "not an answered stage". + WORKTREE_BACKED_STAGES.include?(stage) + end + + def evaluate_answered_stage(_stage, state_file) + return Result.new(safe: false, rationale: "state_file_missing") if state_file.nil? || !File.exist?(state_file) + + parsed = Hive::BrainstormParser.parse(state_file) + if parsed.any?(&:answered?) + Result.new(safe: false, + rationale: "answered_user_content_present") + else + # Zero questions / parse-empty also lands here: nothing authored + # exists, so a rerun cannot clobber user content (mirrors the + # answers-pending fail-open semantics). + Result.new(safe: true, + rationale: "no_answered_user_content") + end + end + + def evaluate_worktree_stage(project, slug) + worktree = resolve_worktree(project, slug) + return Result.new(safe: false, rationale: "worktree_unresolvable") unless worktree + return Result.new(safe: false, rationale: "worktree_missing #{worktree}") unless File.directory?(worktree) + + out, err, status = git_status(worktree) + return Result.new(safe: false, rationale: "git_status_failed #{err.strip.byteslice(0, 120)}") unless status.is_a?(Integer) && status.zero? + + files = out.lines.map(&:strip).reject(&:empty?) + return Result.new(safe: true, rationale: "clean_worktree") if files.empty? + + dirty = files.reject { |line| residue?(line) } + if dirty.empty? + Result.new(safe: true, rationale: "agent_residue_only (#{files.size} file(s))") + else + Result.new(safe: false, + rationale: "dirty_worktree user_or_unknown_files=#{dirty.first(3).join(', ')}") + end + end + + def resolve_worktree(project, slug) + if @worktree_resolver + return @worktree_resolver.call(project, slug) + end + + entry = Hive::Config.find_project(project) + return nil unless entry + + root = Hive::Worktree.canonical_root(entry["path"]) + File.join(root, slug) + rescue StandardError + nil + end + + def git_status(worktree) + argv = %w[git -C] + [worktree] + %w[status --porcelain] + if @git_runner + begin + @git_runner.call(worktree) + rescue StandardError => e + ["", "#{e.class}: #{e.message}", nil] + end + else + run_git_status(argv) + end + end + + def run_git_status(argv) + out = +"" + err = +"" + Open3.popen3(*argv) do |stdin, stdout, stderr, waiter| + stdin.close + readers = [ + Thread.new { out << stdout.read }, + Thread.new { err << stderr.read } + ] + if waiter.join(GIT_STATUS_TIMEOUT_SEC) + readers.each(&:join) + [out, err, waiter.value.exitstatus] + else + kill(waiter.pid) + readers.each(&:join) + [out, "git status timed out", nil] + end + end + rescue StandardError => e + ["", e.message.to_s, nil] + end + + def kill(pid) + Process.kill("TERM", pid) + rescue StandardError + nil + end + + def residue?(porcelain_line) + # porcelain v1: XYPATH (rename entries are "XY ORIG -> PATH"). + path = porcelain_line.sub(/\A..\s+/, "") + .sub(/\A.* -> /, "") + .delete_prefix('"').delete_suffix('"') + RESIDUE_ALLOWLIST.any? { |pattern| path.match?(pattern) } + end + end + end +end diff --git a/lib/hive/events.rb b/lib/hive/events.rb index f8bca8160..395813285 100644 --- a/lib/hive/events.rb +++ b/lib/hive/events.rb @@ -15,6 +15,7 @@ module Hive round_complete clean_exit_auto_committed claude_completion_fallback + auto_retry ].freeze STATUS_TAIL_LINES = 20 diff --git a/lib/hive/stages/execute.rb b/lib/hive/stages/execute.rb index c3081b89b..481dfda24 100644 --- a/lib/hive/stages/execute.rb +++ b/lib/hive/stages/execute.rb @@ -213,8 +213,13 @@ module Hive return { commit: "limits_reached", status: :error } end + # Stamp the executing agent's name (like the limits_reached arm + # above) so the daemon's auto-retrier can attribute future + # implementer_failed markers to a provider without relying on + # message-text signatures. Hive::Markers.set(task.state_file, :error, reason: "implementer_failed", + provider: execute_agent_name(cfg), status: impl_result&.fetch(:status, nil), message: impl_result&.fetch(:error_message, nil)) { commit: "implementer_failed", status: :error } diff --git a/test/unit/daemon/auto_retry_state_test.rb b/test/unit/daemon/auto_retry_state_test.rb new file mode 100644 index 000000000..f939be91c --- /dev/null +++ b/test/unit/daemon/auto_retry_state_test.rb @@ -0,0 +1,164 @@ +require "test_helper" +require "fileutils" +require "tmpdir" +require "time" +require "hive/daemon/auto_retry_state" + +class HiveDaemonAutoRetryStateTest < Minitest::Test + KEY = Hive::Daemon::AutoRetryState.method(:key) + + def setup + @dir = Dir.mktmpdir("hive-auto-retry-state") + @path = File.join(@dir, "daemon_auto_retry_state.json") + end + + def teardown + FileUtils.rm_rf(@dir) + end + + def state(logger: nil) + Hive::Daemon::AutoRetryState.new(path: @path, logger: logger) + end + + def key(project = "p", slug = "s", stage = "4-execute", reason_class = "codex_auth") + KEY.call(project, slug, stage, reason_class) + end + + NOW = Time.utc(2026, 8, 21, 12, 0, 0) + + # --- decision rules (R5) ---------------------------------------------- + + def test_first_attempt_allowed_immediately_on_healthy_signal_change + s = state + s.load! + assert s.retry_allowed?(key: key, current_fingerprint: "fp1", now: NOW) + end + + def test_missing_entry_row_allows_retry + s = state + s.load! + assert s.retry_allowed?(key: key, current_fingerprint: "fp1", now: NOW) + end + + def test_second_retry_requires_backoff_window + s = state + s.load! + fp1 = "fingerprint-one" + s.record_attempt!(key: key, fingerprint: fp1, now: NOW) + + # Same environment fingerprint ⇒ blocked regardless of time. + refute s.retry_allowed?(key: key, current_fingerprint: fp1, + now: NOW + Hive::Daemon::AutoRetryState::SECOND_RETRY_BACKOFF_SEC) + + # Health signal changed but backoff not yet elapsed ⇒ blocked. + refute s.retry_allowed?(key: key, current_fingerprint: "fingerprint-two", + now: NOW + 60) + + # Signal changed AND backoff elapsed ⇒ allowed. + assert s.retry_allowed?(key: key, current_fingerprint: "fingerprint-two", + now: NOW + Hive::Daemon::AutoRetryState::SECOND_RETRY_BACKOFF_SEC) + end + + def test_exhaustion_sticks_across_reload + s = state + s.load! + s.record_attempt!(key: key, fingerprint: "fp1", now: NOW) + s.record_attempt!(key: key, fingerprint: "fp2", now: NOW + 1800) + assert s.exhausted?(key) + + # Fresh instance (simulated daemon restart) reads the same ledger. + reloaded = state + reloaded.load! + assert reloaded.exhausted?(key) + refute reloaded.retry_allowed?(key: key, current_fingerprint: "fp3", + now: NOW + 10_000) + assert_equal 2, reloaded.attempt_count(key) + end + + def test_budget_is_keyed_by_reason_not_marker_id + s = state + s.load! + s.record_attempt!(key: key(slug: "s", stage: "4-execute", reason_class: "codex_auth"), + fingerprint: "fp1", now: NOW) + # A different task/reason/stage has its own untouched budget. + other = key(slug: "s2", stage: "4-execute", reason_class: "codex_auth") + assert_equal 0, s.attempt_count(other) + assert s.retry_allowed?(key: other, current_fingerprint: "whatever", now: NOW) + end + + # --- persistence -------------------------------------------------------- + + def test_round_trip_persistence + s = state + s.load! + entry = s.record_attempt!(key: key, fingerprint: "abc123", now: NOW) + assert_equal 1, entry["attempt_count"] + refute entry["exhausted"] + + raw = JSON.parse(File.read(@path)) + assert_equal 1, raw["schema_version"] + persisted = raw["entries"].first + assert_equal %w[p s 4-execute codex_auth], + persisted.values_at("project", "slug", "stage", "reason_class") + + reloaded = state + reloaded.load! + assert_equal 1, reloaded.attempt_count(key) + assert_equal "abc123", reloaded.entries[key]["last_fingerprint"] + end + + def test_corrupt_file_fails_closed_to_empty_map_with_warning + File.write(@path, "{ this is not json") + logger = FakeLogger.new + s = state(logger: logger) + entries = s.load! + + assert_empty entries + assert s.retry_allowed?(key: key, current_fingerprint: "anything", now: NOW) || + s.entries.empty? + corrupt_event = logger.events.find { |name, _| name == :auto_retry_state_corrupt } + assert corrupt_event, "expected a corrupt-file warning event" + end + + def test_torn_write_partial_json_fails_closed + full = { schema_version: 1, entries: [ { + "project" => "p", "slug" => "s", "stage" => "4-execute", + "reason_class" => "codex_auth", "attempt_count" => 1, + "last_fingerprint" => "x", "last_cleared_at" => NOW.iso8601(6), "exhausted" => false + } ] }.to_json + File.write(@path, full[0, full.size / 2]) # torn mid-JSON + + s = state + assert_empty s.load! + end + + def test_newer_schema_version_degrades_to_empty_map + File.write(@path, JSON.generate({ "schema_version" => 99, "entries" => [] })) + logger = FakeLogger.new + + assert_empty state(logger: logger).load! + assert logger.events.any? { |name, attrs| + name == :auto_retry_state_corrupt && attrs[:reason].include?("schema_version") + } + end + + def test_wrong_shape_entries_fail_closed_but_keep_boot_alive + File.write(@path, JSON.generate({ "schema_version" => 1, "entries" => "nope" })) + + assert_empty state.load! + end + + private + + class FakeLogger + attr_reader :events + + def initialize + @events = [] + end + + def event(name, **attrs) + @events << [ name, attrs ] + end + end +end diff --git a/test/unit/daemon/dispatcher_test.rb b/test/unit/daemon/dispatcher_test.rb index d5de27822..c70c6c09b 100644 --- a/test/unit/daemon/dispatcher_test.rb +++ b/test/unit/daemon/dispatcher_test.rb @@ -1861,6 +1861,106 @@ class HiveDaemonDispatcherTest < Minitest::Test "per-row dispatch must still run after healer crash to avoid StartLimitBurst flapping" end + # ── recoverable-marker retrier integration ──────────────────────────── + + def test_dispatcher_tick_invokes_recoverable_marker_retrier_after_healer + Dir.mktmpdir("dispatcher-retrier-integration") do |tmpdir| + state_file = File.join(tmpdir, "task.md") + File.write(state_file, "# task\n\n\n") + error_row = row( + project: "p1", slug: "parked-1", stage: "4-execute", + marker: "error", action: "error", + marker_attrs: { "reason" => "implementer_failed", "provider" => "codex" }, + state_file: state_file, mtime: T0 - 1000 + ) + dispatcher, _sup, _ctrl, logger, _mw = make_dispatcher(rows: [ error_row ]) + retrier = dispatcher.instance_variable_get(:@recoverable_marker_retrier) + refute_nil retrier, "dispatcher must construct a retrier at boot" + assert retrier.enabled?, "retrier must default to enabled" + + dispatcher.tick(now: T0) + + assert events_include?(logger, :auto_retry_cleared) || events_include?(logger, :auto_retry_blocked), + "retrier must evaluate parked recoverable markers on the tick; events=#{logger.events.inspect}" + end + end + + def test_dispatcher_tick_retrier_disabled_by_config_emits_nothing + Dir.mktmpdir("dispatcher-retrier-disabled") do |tmpdir| + state_file = File.join(tmpdir, "task.md") + File.write(state_file, "# task\n\n\n") + error_row = row( + project: "p1", slug: "parked-2", stage: "4-execute", + marker: "error", action: "error", + marker_attrs: { "reason" => "claude_launch_failed" }, + state_file: state_file, mtime: T0 - 1000 + ) + config = { + "daemon" => { + "edit_debounce_sec" => 30, "poll_interval_sec" => 30, + "auto_retry" => { "enabled" => false } + } + } + controller = Hive::Daemon::ConcurrencyController.new( + max_concurrent_runs: 5, max_concurrent_per_project: 5, + max_runs_per_day_per_project: 100 + ) + supervisor = FakeSupervisor.new + status = FakeStatusConsumer.new + status.next_result = Hive::Daemon::StatusConsumer::Result.new( + ok: true, rows: [ error_row ], projects: [], error: nil + ) + logger = StubLogger.new + dispatcher = Hive::Daemon::Dispatcher.new( + config: config, controller: controller, supervisor: supervisor, + status_consumer: status, logger: logger + ) + + dispatcher.tick(now: T0) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s, + "kill switch on ⇒ marker untouched" + refute events_include?(logger, :auto_retry_cleared) + refute events_include?(logger, :auto_retry_blocked) + end + end + + def test_reload_config_rebuilds_recoverable_marker_retrier + dispatcher, _sup, _ctrl, _logger, _mw = make_dispatcher + original = dispatcher.instance_variable_get(:@recoverable_marker_retrier) + refute_nil original + + new_cfg = { "edit_debounce_sec" => 30 } + original_loader = Hive::Config.method(:load_global_daemon) + Hive::Config.define_singleton_method(:load_global_daemon) { new_cfg } + begin + dispatcher.send(:reload_config!) + ensure + Hive::Config.define_singleton_method(:load_global_daemon, &original_loader) + end + + rebuilt = dispatcher.instance_variable_get(:@recoverable_marker_retrier) + refute_same original, rebuilt, + "reload_config! must rebuild the retrier so daemon.auto_retry.enabled rebinds" + assert rebuilt.enabled? + end + + def test_dispatcher_outer_rescue_logs_fatal_when_retrier_raises + advance_row = row(action: "ready_to_plan", command: "hive plan s1 --from 2-brainstorm") + dispatcher, sup, _ctrl, logger, _mw = make_dispatcher(rows: [ advance_row ]) + retrier = dispatcher.instance_variable_get(:@recoverable_marker_retrier) + def retrier.tick(*, **) + raise NoMethodError, "simulated retrier bug" + end + + dispatcher.tick(now: T0) + + fatal = logger.events.find { |(n, _)| n == :fatal } + assert fatal, "outer rescue must log :fatal when the retrier raises" + assert_includes fatal[1][:message], "recoverable_marker_retrier raised" + assert_equal 1, sup.spawned.size, "per-row dispatch still runs after a retrier crash" + end + def test_reload_config_rebuilds_healer_with_new_grace_sec dispatcher, _sup, _ctrl, _logger, _mw = make_dispatcher original_healer = dispatcher.instance_variable_get(:@stale_agent_healer) diff --git a/test/unit/daemon/health_probes_test.rb b/test/unit/daemon/health_probes_test.rb new file mode 100644 index 000000000..6dca5b68b --- /dev/null +++ b/test/unit/daemon/health_probes_test.rb @@ -0,0 +1,243 @@ +require "test_helper" +require "fileutils" +require "tmpdir" +require "hive/daemon/health_probes" + +# All CLI dependencies are stubbed via the command_runner seam — no real +# network / auth / codex / claude binary in CI. +class HiveDaemonHealthProbesTest < Minitest::Test + def setup + @home = Dir.mktmpdir("hive-probes-home") + @calls = [] + @responses = {} + @probes = probes + end + + def teardown + FileUtils.rm_rf(@home) + end + + # --- universal probe ------------------------------------------------- + + def test_universal_healthy_when_credentials_present + write_credential(".codex/auth.json") + write_credential(".claude/.credentials.json") + + result = @probes.universal_healthy? + assert result.healthy, result.inspect + end + + def test_universal_unhealthy_when_codex_credential_missing + write_credential(".claude/.credentials.json") + + result = @probes.universal_healthy? + refute result.healthy + assert_equal "codex", result.detail[:agents].first + end + + def test_universal_unhealthy_when_credential_is_empty_stub + write_credential(".codex/auth.json", "{}") + write_credential(".claude/.credentials.json", "{}") + + refute @probes.universal_healthy?.healthy + end + + # --- codex auth probe ------------------------------------------------ + + def test_codex_auth_healthy_when_logged_in_and_smoke_ok + ready_home + stub_command(%w[codex login status], ["Logged in using auth.json\n", "", 0]) + stub_command(%w[codex exec] + [Hive::Daemon::HealthProbes::CODEX_SMOKE_PROMPT], + ["ok\n", "", 0]) + + result = @probes.codex_auth_healthy? + assert result.healthy, result.inspect + assert_equal "ok", result.detail[:codex_exec_smoke] + end + + def test_codex_auth_unhealthy_when_login_reports_logged_out + ready_home + stub_command(%w[codex login status], ["Not logged in\n", "", 1]) + + result = @probes.codex_auth_healthy? + refute result.healthy + assert_equal :codex_login_status, result.probe + assert_equal 1, result.detail[:exit_status] + # Smoke test must NOT run when the login check already failed. + assert_equal [%w[codex login status]], @calls.map(&:first) + end + + def test_codex_auth_unhealthy_when_smoke_fails + ready_home + stub_command(%w[codex login status], ["Logged in using auth.json\n", "", 0]) + stub_command(%w[codex exec] + [Hive::Daemon::HealthProbes::CODEX_SMOKE_PROMPT], + ["partial out", "boom: quota exhausted", 1]) + + result = @probes.codex_auth_healthy? + refute result.healthy + assert_equal :codex_exec_smoke, result.probe + assert_includes result.detail[:stderr], "quota exhausted" + end + + def test_probe_timeout_yields_unhealthy_with_partial_output + ready_home + stub_command(%w[codex login status], ["Logged in using auth.json\n", "", 0]) + stub_command(%w[codex exec] + [Hive::Daemon::HealthProbes::CODEX_SMOKE_PROMPT], + ["partial", "", :timeout]) + + result = @probes.codex_auth_healthy? + refute result.healthy, "a hanging smoke test must classify unhealthy" + assert_equal :codex_exec_smoke, result.probe + assert_equal "timeout", result.detail[:exit_status] + assert_equal "partial", result.detail[:stdout] + end + + def test_spawn_error_never_raises_into_caller + ready_home + @responses[%w[codex login status]] = -> { raise Errno::ENOENT, "codex" } + + result = @probes.codex_auth_healthy? + refute result.healthy + assert_includes result.detail[:exit_status].to_s, "ENOENT" + end + + # --- claude launcher probe ------------------------------------------- + + def test_claude_launcher_healthy_when_all_checks_pass + ready_home + wrapper = write_wrapper + stub_command(["claude", "--version"], ["2.1.179 (Claude Code)\n", "", 0]) + probes = Hive::Daemon::HealthProbes.new( + command_runner: runner, wrapper_path: wrapper, claude_bin: "claude", home: @home + ) + + result = probes.claude_launcher_healthy? + assert result.healthy, result.inspect + assert_includes result.detail[:version_output], "2.1.179" + end + + def test_wrapper_missing_is_unhealthy_without_shelling_out + ready_home + probes = Hive::Daemon::HealthProbes.new( + command_runner: runner, wrapper_path: File.join(@home, "nope.sh"), home: @home + ) + + result = probes.claude_launcher_healthy? + refute result.healthy + assert_equal :claude_wrapper_missing, result.probe + assert_empty @calls, "must not shell out when the wrapper is missing" + end + + def test_ready_detector_drift_is_unhealthy_without_shelling_out + ready_home + wrapper = write_wrapper + # Simulate detector drift by stubbing the class method to always fail + # the positive fixtures. Save + restore the original implementation + # (define_singleton_method would otherwise clobber the module_function + # alias and poison every later suite in the same process). + original = Hive::ClaudeLauncher.method(:claude_ready_prompt?) + Hive::ClaudeLauncher.define_singleton_method(:claude_ready_prompt?) { |_pane| false } + begin + probes = Hive::Daemon::HealthProbes.new( + command_runner: runner, wrapper_path: wrapper, home: @home + ) + result = probes.claude_launcher_healthy? + refute result.healthy + assert_equal :claude_ready_detector, result.probe + assert_empty @calls + ensure + Hive::ClaudeLauncher.define_singleton_method(:claude_ready_prompt?, &original) + end + end + + def test_claude_version_failure_is_unhealthy + ready_home + wrapper = write_wrapper + stub_command(["claude", "--version"], ["", "command not found", :timeout]) + probes = Hive::Daemon::HealthProbes.new( + command_runner: runner, wrapper_path: wrapper, claude_bin: "claude", home: @home + ) + + result = probes.claude_launcher_healthy? + refute result.healthy + assert_equal :claude_version, result.probe + end + + # --- memoization ------------------------------------------------------- + + def test_results_memoized_per_instance_across_evaluations + ready_home + stub_command(%w[codex login status], ["Logged in using auth.json\n", "", 0]) + stub_command(%w[codex exec] + [Hive::Daemon::HealthProbes::CODEX_SMOKE_PROMPT], + ["ok\n", "", 0]) + + 2.times { @probes.codex_auth_healthy? } + + # One login-status call and one smoke call despite two evaluations. + assert_equal 1, @calls.count { |argv, _| argv == %w[codex login status] } + assert_equal 1, @calls.count { |argv, _| argv.first == "codex" && argv[1] == "exec" } + end + + # --- fingerprint --------------------------------------------------------- + + def test_fingerprint_changes_when_signal_changes_and_stable_otherwise + ready_home + wrapper = write_wrapper + probes = Hive::Daemon::HealthProbes.new(wrapper_path: wrapper, home: @home) + + fp1 = Hive::Daemon::HealthProbes.fingerprint(probes.fingerprint_inputs) + + # Bump the codex credential's mtime → signal changes → fingerprint changes. + cred = File.join(@home, ".codex", "auth.json") + File.utime(Time.now, Time.now + 3600, cred) + fp2 = Hive::Daemon::HealthProbes.fingerprint( + Hive::Daemon::HealthProbes.new(wrapper_path: wrapper, home: @home).fingerprint_inputs + ) + refute_equal fp1, fp2 + + # Same inputs → same fingerprint (deterministic). + fp3 = Hive::Daemon::HealthProbes.fingerprint( + Hive::Daemon::HealthProbes.new(wrapper_path: wrapper, home: @home).fingerprint_inputs + ) + assert_equal fp2, fp3 + end + + private + + def probes + Hive::Daemon::HealthProbes.new(command_runner: runner, wrapper_path: nil, home: @home) + end + + def runner + ->(argv, timeout_sec:) do + @calls << [argv, timeout_sec] + handler = @responses[argv] + case handler + when Proc then handler.call + when Array then handler + else ["", "unstubbed command #{argv.inspect}", 127] + end + end + end + + def stub_command(argv, response) + @responses[argv] = response + end + + def ready_home + write_credential(".codex/auth.json") + write_credential(".claude/.credentials.json") + end + + def write_credential(rel_path, content = '{"token":"t"}') + path = File.join(@home, rel_path) + FileUtils.mkdir_p(File.dirname(path)) + File.write(path, content) + end + + def write_wrapper + path = File.join(@home, "interactive_claude_wrapper.sh") + File.write(path, "#!/usr/bin/env bash\n") + path + end +end diff --git a/test/unit/daemon/recoverable_marker_retrier_test.rb b/test/unit/daemon/recoverable_marker_retrier_test.rb new file mode 100644 index 000000000..47a854ce8 --- /dev/null +++ b/test/unit/daemon/recoverable_marker_retrier_test.rb @@ -0,0 +1,431 @@ +require "test_helper" +require "tmpdir" +require "fileutils" +require "hive/markers" +require "hive/daemon/recoverable_marker_retrier" +require "hive/daemon/status_consumer" + +# End-to-end unit scenarios for the health-gated auto-retrier (U5), plus +# acceptance shapes from R10 (task-58 codex-auth, task-287 launcher, +# unknown-signature park, backoff/exhaustion, dirty worktree, probe hang, +# kill switch). +class HiveDaemonRecoverableMarkerRetrierTest < Minitest::Test + Row = Hive::Daemon::StatusConsumer::Row + NOW = Time.utc(2026, 8, 21, 12, 0, 0) + + 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 + attr_accessor :raise_on_write + + def initialize + @requests = [] + @raise_on_write = false + end + + def write_request!(**kwargs) + raise Errno::EACCES, "queue dir unwritable" if @raise_on_write + + @requests << kwargs + "fake-req-#{@requests.size}" + end + end + + # Fully stubbed probe suite — no real CLI calls anywhere. + class FakeProbes + attr_accessor :universal_result, :codex_result, :launcher_result, :signal_version + + def initialize(signal_version: "v1") + @signal_version = signal_version + @universal_result = healthy(:universal) + @codex_result = nil + @launcher_result = nil + end + + def healthy(probe) + Hive::Daemon::HealthProbes::Result.new(healthy: true, probe: probe, detail: {}) + end + + def unhealthy(probe, detail = {}) + Hive::Daemon::HealthProbes::Result.new(healthy: false, probe: probe, detail: detail) + end + + def universal_healthy? + @universal_result + end + + def codex_auth_healthy? + @codex_result || healthy(:codex_auth) + end + + def claude_launcher_healthy? + @launcher_result || healthy(:claude_launcher) + end + + def fingerprint_inputs + { "signal" => @signal_version } + end + end + + def setup + @dir = Dir.mktmpdir("hive-retrier") + @state_path = File.join(@dir, "auto_retry_state.json") + @logger = FakeLogger.new + @controller = FakeController.new + @queue = FakeRequestQueue.new + @guard_calls = [] + end + + def teardown + FileUtils.rm_rf(@dir) + end + + # --- task-58 shape: codex auth failure auto-clears --------------------- + + def test_codex_auth_marker_clears_and_enqueues_when_probes_green + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file, stage: "4-execute") + + retrier.tick([ row ]) + + assert_equal "none", Hive::Markers.current(state_file).name.to_s + request = @queue.requests.first + assert request, "expected a dispatch request to be written" + assert_equal "healer", request[:requestor] + assert_equal [ "hive", "run", "s", "--project", "p" ], request[:argv] + assert_equal "auto_retry", request[:trigger] + + cleared = @logger.events.find { |name, _| name == :auto_retry_cleared } + assert cleared, "expected positive daemon-log audit event" + attrs = cleared[1] + assert_equal "implementer_failed", attrs[:marker_reason] + assert_equal "codex_auth", attrs[:reason_class] + refute_empty attrs[:health_fingerprint] + assert_equal 1, attrs[:attempt_count] + + # Task-local audit channel. + events_line = File.readlines(File.join(File.dirname(state_file), "events.jsonl")) + .last + assert_includes events_line, '"event_type":"auto_retry"' + end + + def test_legacy_codex_auth_message_signature_also_clears + state_file = write_error_marker( + reason: "implementer_failed", + message: "unexpected status 401 Unauthorized from api.openai.com" + ) + row = make_row(state_file, stage: "4-execute") + + retrier.tick([ row ]) + + assert_equal "none", Hive::Markers.current(state_file).name.to_s + assert @logger.events.any? { |name, _| name == :auto_retry_cleared } + end + + # --- task-287 shape: launcher failure needs full launcher probe set ---- + + def test_claude_launcher_marker_clears_after_launcher_probe_passes + state_file = write_error_marker(reason: "claude_launch_failed") + row = make_row(state_file, stage: "4-execute") + probes = FakeProbes.new + probes.launcher_result = probes.unhealthy(:claude_wrapper_missing) + + retrier(probes_factory: -> { probes }).tick([ row ]) + refute_equal "none", Hive::Markers.current(state_file).name.to_s + blocked = @logger.events.find { |name, _| name == :auto_retry_blocked } + assert blocked + assert_equal "claude_wrapper_missing failed", blocked[1][:not_retried_because] + + # Once ALL launcher checks pass inside one probe result, it clears. + @logger = FakeLogger.new + probes.launcher_result = probes.healthy(:claude_launcher) + retrier(probes_factory: -> { probes }, state_fresh: false).tick([ row ]) + + # NOTE: first tick recorded a throttled negative but consumed no budget; + # this second tick retries and clears. + assert_equal "none", Hive::Markers.current(state_file).name.to_s + end + + # --- negative scenarios ------------------------------------------------- + + def test_unknown_implementer_failed_signature_stays_parked + state_file = write_error_marker(reason: "implementer_failed", + message: "tests failed: expected 3 got 4") + row = make_row(state_file) + + retrier.tick([ row ]) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + assert_empty @logger.events.reject { |name, _| name == :fatal } + end + + def test_unhealthy_probe_blocks_retry_and_throttles_negative_events + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + probes = FakeProbes.new + probes.codex_result = probes.unhealthy(:codex_exec_smoke, stdout: "partial", exit_status: "timeout") + + r = retrier(probes_factory: -> { probes }) + r.tick([ row ]) + r.tick([ row ]) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + blocked = @logger.events.select { |name, _| name == :auto_retry_blocked } + assert_equal 1, blocked.size, "throttled: one event per (task, reason, fingerprint)" + + # A changed health signal re-arms the throttle AND gets probed again. + probes.signal_version = "v2" + r.tick([ row ]) + assert_equal 2, @logger.events.select { |name, _| name == :auto_retry_blocked }.size + end + + def test_probe_timeout_is_no_retry + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + probes = FakeProbes.new + probes.codex_result = probes.unhealthy(:codex_login_status, exit_status: "timeout") + + retrier(probes_factory: -> { probes }).tick([ row ]) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + end + + def test_dirty_worktree_is_not_retried + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + guard = Hive::Daemon::RetrySafetyGuard.new( + git_runner: ->(_wt) { [ " M src/app.rb\n", "", 0 ] }, + worktree_resolver: ->(_p, _s) { @dir } + ) + + retrier(guard: guard).tick([ row ]) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + blocked = @logger.events.find { |name, _| name == :auto_retry_blocked } + assert_match(/safety_guard/, blocked[1][:not_retried_because]) + assert_match(/dirty_worktree/, blocked[1][:not_retried_because]) + end + + # --- budget / backoff / exhaustion (R5) --------------------------------- + + def test_second_retry_backoff_then_exhaustion_parks_permanently + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + + r = retrier + r.tick([ row ], now: NOW) # attempt 1: immediate + assert_equal 1, @queue.requests.size + + # Rewrite an identical fresh marker (a rerun that re-failed parks red + # again with a NEW marker id). + state_file = rewrite_error_marker(state_file, reason: "implementer_failed", + provider: "codex") + row = make_row(state_file) + probes_v2 = FakeProbes.new(signal_version: "changed-environment") + r2 = retrier(probes_factory: -> { probes_v2 }, state_fresh: false) + r2.tick([ row ], now: NOW + 60) # backoff not elapsed ⇒ blocked + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_equal 1, @queue.requests.size + + r2.tick([ row ], now: NOW + Hive::Daemon::AutoRetryState::SECOND_RETRY_BACKOFF_SEC + 1) + assert_equal "none", Hive::Markers.current(state_file).name.to_s + assert_equal 2, @queue.requests.size + exhausted = @logger.events.find { |name, _| name == :auto_retry_exhausted } + assert exhausted, "expected exhaustion event at limit" + assert_equal "persistent", exhausted[1][:budget_scope] + assert_equal "manual markers clear", exhausted[1][:suggested_next_action] + + # A THIRD identical failure stays parked forever — even across a full + # ledger reload (simulated restart / SIGHUP rebuild). + state_file = rewrite_error_marker(state_file, reason: "implementer_failed", + provider: "codex") + row = make_row(state_file) + probes_v3 = FakeProbes.new(signal_version: "changed-again") + r3 = retrier(probes_factory: -> { probes_v3 }, state_fresh: false) # fresh state instance reads same persisted ledger + r3.tick([ row ], now: NOW + 100_000) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_equal 2, @queue.requests.size, "no third request after exhaustion" + end + + def test_exhaustion_skip_is_silent_on_subsequent_ticks + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + + state = Hive::Daemon::AutoRetryState.new(path: @state_path, logger: @logger) + state.load! + key = Hive::Daemon::AutoRetryState.key("p", "s", "4-execute", "codex_auth") + state.record_attempt!(key: key, fingerprint: "fp1", now: NOW) + state.record_attempt!(key: key, fingerprint: "fp2", now: NOW + 1) + assert state.exhausted?(key) + + @logger = FakeLogger.new + retrier(state: Hive::Daemon::AutoRetryState.new(path: @state_path)).tick([ row ], now: NOW + 2) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @logger.events.reject { |name, _| name == :fatal }, + "exhausted rows skip silently on later ticks" + end + + # --- inherited guards & kill switch -------------------------------------- + + def test_kill_switch_leaves_module_inert_with_zero_events + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + + retrier(enabled: false).tick([ row ]) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + assert_empty @logger.events.reject { |name, _| name == :fatal } + end + + def test_in_flight_task_is_skipped + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + @controller = FakeController.new(running_pairs: [[ "p", "s" ]]) + + retrier.tick([ row ]) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + end + + def test_live_task_lock_is_skipped + state_file = write_error_marker(reason: "claude_launch_failed") + row = make_row(state_file, live_task_lock: true) + + retrier.tick([ row ]) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + end + + def test_legacy_layout_project_is_skipped + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + + retrier.tick([ row ], legacy_layout_projects: { "p" => true }) + + assert_equal "error", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + end + + # --- enqueue failure path -------------------------------------------------- + + def test_enqueue_failure_after_clear_logs_manual_remediation + state_file = write_error_marker(reason: "implementer_failed", provider: "codex") + row = make_row(state_file) + @queue.raise_on_write = true + + retrier.tick([ row ]) + + assert_equal "none", Hive::Markers.current(state_file).name.to_s, + "the clear already succeeded" + failed = @logger.events.find { |name, _| name == :auto_retry_enqueue_failed } + assert failed, "expected distinct enqueue-failure event" + assert_match(/hive run s --project p/, failed[1][:remediation]) + end + + def test_non_error_markers_are_never_touched + state_file = File.join(@dir, "task.md") + File.write(state_file, "# task\n\n\n") + row = make_row(state_file) + Hive::Markers.set(state_file, :agent_working) + + retrier.tick([ row ]) + + assert_equal "agent_working", Hive::Markers.current(state_file).name.to_s + assert_empty @queue.requests + end + + private + + def retrier(probes_factory: nil, guard: nil, enabled: true, state: nil, state_fresh: true) + state ||= begin + FileUtils.rm_f(@state_path) if state_fresh && File.exist?(@state_path) + Hive::Daemon::AutoRetryState.new(path: @state_path, logger: @logger) + end + Hive::Daemon::RecoverableMarkerRetrier.new( + controller: @controller, + logger: @logger, + request_queue: @queue, + enabled: enabled, + state: state, + guard: guard || safe_guard, + probes_factory: probes_factory || -> { FakeProbes.new }, + dispatch_request_state_home: @dir + ) + end + + def safe_guard + @safe_guard ||= Hive::Daemon::RetrySafetyGuard.new( + git_runner: ->(_wt) { [ "", "", 0 ] }, + worktree_resolver: ->(_p, _s) { @dir } + ) + end + + def make_row(state_file, stage: "4-execute", live_task_lock: nil) + Row.new( + project: "p", slug: "s", stage: stage, workflow: nil, + marker: Hive::Markers.current(state_file).name.to_s, + marker_attrs: Hive::Markers.current(state_file).attrs, + folder: File.dirname(state_file), + state_file: state_file, + state_file_mtime: File.mtime(state_file), + action: "error", suggested_command: nil, + claude_pid_alive: nil, live_task_lock: live_task_lock, + diagnostic: nil + ) + end + + def write_error_marker(**attrs) + state_file = File.join(@dir, SecureRandom.hex(4), "task.md") + FileUtils.mkdir_p(File.dirname(state_file)) + File.write(state_file, "# task\n") + Hive::Markers.set(state_file, :error, attrs) + state_file + end + + # Simulate "the rerun re-failed": drop any current marker and write a + # fresh ERROR marker carrying a NEW marker_id. + def rewrite_error_marker(state_file, **attrs) + current = Hive::Markers.current(state_file) + body = current.raw ? File.read(state_file).sub(current.raw, "") : File.read(state_file) + File.write(state_file, body) + Hive::Markers.set(state_file, :error, attrs) + state_file + end +end diff --git a/test/unit/daemon/recovery_classifier_test.rb b/test/unit/daemon/recovery_classifier_test.rb new file mode 100644 index 000000000..e9d3a8a2d --- /dev/null +++ b/test/unit/daemon/recovery_classifier_test.rb @@ -0,0 +1,101 @@ +require "test_helper" +require "hive/markers" +require "hive/daemon/recovery_classifier" + +# Matrix over marker shapes for the v1 recoverable-class allowlist. +class HiveDaemonRecoveryClassifierTest < Minitest::Test + def test_codex_auth_via_provider_stamp + assert_equal :codex_auth, classify("error", + "reason" => "implementer_failed", "provider" => "codex") + end + + def test_codex_auth_legacy_401_message + assert_equal :codex_auth, classify("error", + "reason" => "implementer_failed", "message" => "unexpected status 401 Unauthorized") + end + + def test_codex_auth_legacy_missing_bearer_message + assert_equal :codex_auth, classify("error", + "reason" => "implementer_failed", "message" => "MissingBearerToken: no bearer token present") + end + + def test_codex_auth_legacy_basic_auth_message_case_insensitive + assert_equal :codex_auth, classify("error", + "reason" => "implementer_failed", "message" => "Error: BASIC AUTH not supported") + end + + def test_non_codex_provider_is_excluded_even_with_signature_message + assert_nil classify("error", + "reason" => "implementer_failed", "provider" => "claude", "message" => "401 unauthorized") + end + + def test_plain_implementer_failed_without_signature_stays_manual + assert_nil classify("error", "reason" => "implementer_failed") + assert_nil classify("error", + "reason" => "implementer_failed", "message" => "tests failed: expected 3 got 4") + end + + def test_claude_launch_failed_classifies_launcher + assert_equal :claude_launcher, + classify("error", "reason" => "claude_launch_failed") + end + + def test_claude_launch_failed_with_exception_attr_still_launcher + assert_equal :claude_launcher, classify("error", + "reason" => "claude_launch_failed", + "exception_class" => "Hive::AgentError", + "message" => "claude interactive prompt did not become ready") + end + + def test_excluded_reason_families_are_nil + %w[ + limits_reached tmux_session_terminated agent_orphaned timeout + unpushed_commits ensure_clean_on_exit_failed dirty_worktree + merge_conflict review_error all_failed fix_failed exit_code + ].each do |reason| + assert_nil classify("error", "reason" => reason), "expected #{reason} to be excluded" + end + end + + def test_non_error_markers_are_never_recoverable + assert_nil classify("agent_working", "reason" => "implementer_failed", "provider" => "codex") + assert_nil classify(:review_error, "reason" => "implementer_failed", "provider" => "codex") + assert_nil classify("complete", "reason" => "claude_launch_failed") + assert_nil classify(:none, {}) + end + + def test_symbol_keys_and_marker_name_accepted + assert_equal :codex_auth, classify(:error, + reason: "implementer_failed", provider: "codex") + end + + def test_allowlist_disjoint_from_healer_owned_reasons + # The retrier's allowlist must never overlap a reason the healer + # already handles — double-healing the same marker would silently + # double-consume budgets. + healer_reasons = %w[ + limits_reached tmux_session_terminated agent_orphaned timeout + unpushed_commits ensure_clean_on_exit_failed review_agent_died + all_failed reviewer_partial_failure fix_failed + ] + classified = healer_reasons.map { |reason| classify("error", "reason" => reason) } + assert classified.all?(&:nil?) + assert_equal %i[codex_auth claude_launcher], Hive::Daemon::RecoveryClassifier::REASON_CLASSES + end + + def test_provider_stamp_round_trips_through_markers_parse_attrs + raw = Hive::Markers.build_marker( + "ERROR", "reason" => "implementer_failed", "provider" => "codex", "marker_id" => "abc123" + ) + attrs = raw.match(Hive::Markers::MARKER_RE) { |m| Hive::Markers.parse_attrs(m[:attrs]) } + assert_equal "codex", attrs["provider"] + assert_equal "implementer_failed", attrs["reason"] + assert_equal :codex_auth, Hive::Daemon::RecoveryClassifier.classify(:error, attrs) + end + + private + + def classify(marker_name, attrs) + Hive::Daemon::RecoveryClassifier.classify(marker_name, attrs) + end +end diff --git a/test/unit/daemon/retry_safety_guard_test.rb b/test/unit/daemon/retry_safety_guard_test.rb new file mode 100644 index 000000000..1a1f97d3d --- /dev/null +++ b/test/unit/daemon/retry_safety_guard_test.rb @@ -0,0 +1,132 @@ +require "test_helper" +require "fileutils" +require "tmpdir" +require "hive/daemon/retry_safety_guard" + +class HiveDaemonRetrySafetyGuardTest < Minitest::Test + def setup + @dir = Dir.mktmpdir("hive-safety-guard") + end + + def teardown + FileUtils.rm_rf(@dir) + end + + # --- worktree-backed stages ------------------------------------------- + + def test_clean_worktree_is_safe + guard = guard_with_git_status(["", "", 0]) + result = guard.evaluate(stage: "4-execute", project: "p", slug: "s") + assert result.safe, result.rationale + end + + def test_dirty_worktree_with_user_edit_is_unsafe + guard = guard_with_git_status([" M src/app.rb\n", "", 0]) + result = guard.evaluate(stage: "4-execute", project: "p", slug: "s") + refute result.safe + assert_match(/dirty_worktree/, result.rationale) + assert_match(/src\/app\.rb/, result.rationale) + end + + def test_dirty_worktree_with_agent_residue_only_is_safe + guard = guard_with_git_status(["?? .hive/state.json\n", "", 0]) + result = guard.evaluate(stage: "4-execute", project: "p", slug: "s") + assert result.safe, result.rationale + assert_match(/agent_residue_only/, result.rationale) + end + + def test_git_failure_is_unsafe + guard = guard_with_git_status(["", "fatal: not a git repository", 128]) + result = guard.evaluate(stage: "4-execute", project: "p", slug: "s") + refute result.safe + assert_match(/git_status_failed/, result.rationale) + end + + def test_git_timeout_is_unsafe + guard = guard_with_git_status(["", "git status timed out", nil]) + result = guard.evaluate(stage: "4-execute", project: "p", slug: "s") + refute result.safe + end + + # --- brainstorm / plan stages ------------------------------------------ + + def test_brainstorm_without_answers_is_safe + path = write_state_file(<<~MD) + ## Round 1 + + ### Q1. Which database? + + MD + result = guard.evaluate(stage: "2-brainstorm", project: "p", slug: "s", state_file: path) + assert result.safe, result.rationale + end + + def test_brainstorm_with_any_answer_is_unsafe + path = write_state_file(<<~MD) + ## Round 1 + + ### Q1. Which database? + ### A1. + Postgres please. + + ### Q2. Auth provider? + + MD + result = guard.evaluate(stage: "2-brainstorm", project: "p", slug: "s", state_file: path) + refute result.safe + assert_equal "answered_user_content_present", result.rationale + end + + def test_plan_with_all_answers_unanswered_parse_empty_is_safe + path = write_state_file("# plan draft only, no Q&A structure\n") + + result = guard.evaluate(stage: "3-plan", project: "p", slug: "s", state_file: path) + assert result.safe, result.rationale + end + + def test_missing_state_file_is_unsafe_for_answered_stage + result = guard.evaluate(stage: "2-brainstorm", project: "p", slug: "s", + state_file: File.join(@dir, "nope.md")) + refute result.safe + end + + # --- fail-closed edges --------------------------------------------------- + + def test_unknown_stage_is_unsafe + result = guard.evaluate(stage: "9-mystery", project: "p", slug: "s") + refute result.safe + assert_match(/unknown_stage/, result.rationale) + end + + def test_exception_in_guard_is_unsafe_not_raised + guard = Hive::Daemon::RetrySafetyGuard.new( + git_runner: ->(_wt) { raise "boom" }, + worktree_resolver: ->(_project, _slug) { @dir } + ) + result = guard.evaluate(stage: "4-execute", project: "p", slug: "s") + refute result.safe + assert_match(/guard_error|git_status_failed/, result.rationale) + end + + private + + def guard_with_git_status(response) + Hive::Daemon::RetrySafetyGuard.new( + git_runner: ->(_worktree) { response }, + worktree_resolver: ->(_project, _slug) { @dir } + ) + end + + def guard + Hive::Daemon::RetrySafetyGuard.new( + git_runner: ->(_worktree) { ["", "", 0] }, + worktree_resolver: ->(_project, _slug) { @dir } + ) + end + + def write_state_file(content) + path = File.join(@dir, "brainstorm.md") + File.write(path, content) + path + end +end diff --git a/test/unit/stages/execute_test.rb b/test/unit/stages/execute_test.rb index af52b588b..eb406fe1a 100644 --- a/test/unit/stages/execute_test.rb +++ b/test/unit/stages/execute_test.rb @@ -162,7 +162,7 @@ class HiveStagesExecuteTest < Minitest::Test assert_equal "error", marker.attrs["status"] assert_equal "exit_code=1 compile error", marker.attrs["message"] refute marker.attrs.key?("retry_after") - refute marker.attrs.key?("provider") + assert_equal "codex", marker.attrs["provider"] end end @@ -189,7 +189,7 @@ class HiveStagesExecuteTest < Minitest::Test assert_equal "timeout", marker.attrs["status"] assert_equal "claude stop hook did not signal completion", marker.attrs["message"] refute marker.attrs.key?("retry_after") - refute marker.attrs.key?("provider") + assert_equal "codex", marker.attrs["provider"] end end