Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
182 changes: 171 additions & 11 deletions .agents/scripts/pulse-dispatch-engine.sh
Original file line number Diff line number Diff line change
Expand Up @@ -1185,13 +1185,157 @@ _pulse_run_blocker_refresh_catchup() {
return 0
}

#######################################
# Return the post-dispatch housekeeping stage-progress state file path.
#
# Lives next to (not inside) the lock directory so a stale-lock reclaim, which
# removes the lock directory, keeps the interrupted run's progress (GH#34287).
#
# Returns 0 always; emits path on stdout.
#######################################
_pulse_post_dispatch_housekeeping_statefile() {
printf '%s\n' "${HOME}/.aidevops/logs/pulse-post-dispatch-housekeeping.state"
return 0
}

#######################################
# Return 0 when a housekeeping stage completed within the resume window.
#
# The state file holds `stage=epoch` lines written after each stage finishes
# and is removed by a clean, complete run. Fresh entries therefore only exist
# after an interrupted run (deploy reconciliation SIGTERM, crash, reboot).
# AIDEVOPS_PULSE_HOUSEKEEPING_RESUME_WINDOW_S=0 disables resume.
#
# Args:
# $1 - stage name
# $2 - state file path
# Returns 0 if recently completed, 1 otherwise; emits age seconds on stdout.
#######################################
_pulse_housekeeping_stage_recently_done() {
local stage_name="$1"
local state_file="$2"
local window="${AIDEVOPS_PULSE_HOUSEKEEPING_RESUME_WINDOW_S:-1800}"
[[ "$window" =~ ^[0-9]+$ ]] || window=1800
[[ "$window" -gt 0 && -f "$state_file" ]] || return 1

local line="" done_epoch=""
while IFS= read -r line || [[ -n "$line" ]]; do
if [[ "$line" == "${stage_name}="* ]]; then
done_epoch="${line#*=}"
fi
done <"$state_file"
[[ "$done_epoch" =~ ^[0-9]+$ ]] || return 1

local now=""
now=$(date +%s 2>/dev/null) || return 1
[[ "$now" =~ ^[0-9]+$ ]] || return 1
local age=$((now - done_epoch))
[[ "$age" -ge 0 && "$age" -lt "$window" ]] || return 1
printf '%s\n' "$age"
return 0
}

#######################################
# Record a finished housekeeping stage in the state file (stage=epoch).
#
# Args:
# $1 - stage name
# $2 - state file path
# Returns 0 always (progress persistence is best-effort).
#######################################
_pulse_housekeeping_mark_stage_done() {
local stage_name="$1"
local state_file="$2"
local now=""
now=$(date +%s 2>/dev/null) || return 0
[[ "$now" =~ ^[0-9]+$ ]] || return 0

local tmp_file="${state_file}.tmp.${BASHPID:-$$}"
local line=""
{
if [[ -f "$state_file" ]]; then
while IFS= read -r line || [[ -n "$line" ]]; do
[[ -n "$line" && "$line" != "${stage_name}="* ]] && printf '%s\n' "$line"
done <"$state_file"
fi
printf '%s=%s\n' "$stage_name" "$now"
} >"$tmp_file" 2>/dev/null || {
rm -f "$tmp_file" 2>/dev/null || true
return 0
}
mv -f "$tmp_file" "$state_file" 2>/dev/null || rm -f "$tmp_file" 2>/dev/null || true
return 0
}

#######################################
# Run one resumable housekeeping stage.
#
# Skips the stage when an interrupted earlier run already completed it within
# the resume window; otherwise records it as the current stage (for the
# signal trap), runs it, and persists its completion.
#
# Args:
# $1 - stage name (state key)
# $2 - state file path
# $@ - stage command and arguments
# Returns 0 always.
#######################################
_pulse_run_resumable_housekeeping_stage() {
local stage_name="$1"
local state_file="$2"
shift 2
local age=""
if age=$(_pulse_housekeeping_stage_recently_done "$stage_name" "$state_file"); then
echo "[pulse-wrapper] Async post-dispatch housekeeping: resume — skipping ${stage_name} (completed ${age}s ago by an interrupted run)" >>"$LOGFILE"
return 0
fi
_PULSE_HOUSEKEEPING_CURRENT_STAGE="$stage_name"
local stage_rc=0
"$@" || stage_rc=$?
# A stage child killed by a signal (e.g. the same deploy SIGTERM that is
# about to reach this subshell) did not finish: leave it for the resume.
if [[ "$stage_rc" -gt 128 ]]; then
echo "[pulse-wrapper] Async post-dispatch housekeeping: ${stage_name} ended by signal (rc=${stage_rc}) — not recorded as complete" >>"$LOGFILE"
else
_pulse_housekeeping_mark_stage_done "$stage_name" "$state_file"
fi
_PULSE_HOUSEKEEPING_CURRENT_STAGE="between_stages"
return 0
}

#######################################
# Log and exit when the async housekeeping subshell receives a signal.
#
# Setup reconciliation SIGTERMs every Pulse runtime process on deploy
# (intended: old-bundle code must be replaced). Logging the signal and the
# active stage makes those deaths attributable from pulse.log (GH#34287).
#
# Args:
# $1 - signal name (TERM or HUP)
# Exits 143 for TERM, 129 for HUP.
#######################################
_pulse_housekeeping_on_signal() {
local signal_name="$1"
trap - TERM HUP
echo "[pulse-wrapper] Async post-dispatch housekeeping: terminated by SIG${signal_name} during ${_PULSE_HOUSEKEEPING_CURRENT_STAGE:-unknown}" >>"$LOGFILE"
if [[ "$signal_name" == "HUP" ]]; then
exit 129
fi
exit 143
}

#######################################
# Run non-dispatch post-dispatch housekeeping stages.
#
# These stages are intentionally after early dispatch and do not protect the
# immediate worker claim/ledger safety boundary. They can therefore run under a
# separate lock while the main pulse proceeds to prefetch + the next refill.
#
# GH#34287: deploy reconciliation kills long runs, so per-stage completion is
# persisted and a following run resumes at the first incomplete stage. The
# cheap terminal reconciles (preflight_ownership_reconcile) run before the
# slow GH#33246 blocker-refresh catch-up so they are not starved behind it.
#
# Args:
# $1 - per-stage timeout seconds
# Returns 0 always.
Expand All @@ -1200,22 +1344,34 @@ _pulse_run_post_dispatch_housekeeping_stages() {
local stage_timeout="${1:-${PREFLIGHT_GROUP_TIMEOUT:-${PRE_RUN_STAGE_TIMEOUT:-600}}}"
[[ "$stage_timeout" =~ ^[0-9]+$ ]] || stage_timeout=600

local lockdir
local lockdir="" state_file=""
lockdir="$(_pulse_post_dispatch_housekeeping_lockdir)"
state_file="$(_pulse_post_dispatch_housekeeping_statefile)"
if ! _pulse_acquire_post_dispatch_housekeeping_lock "$lockdir"; then
return 0
fi

_PULSE_HOUSEKEEPING_CURRENT_STAGE="startup"
echo "[pulse-wrapper] Async post-dispatch housekeeping: started (timeout=${stage_timeout}s)" >>"$LOGFILE"
_pulse_run_optional_stage_with_timeout "coderabbit_review" "$stage_timeout" run_daily_codebase_review || true
_pulse_run_optional_stage_with_timeout "post_merge_scanner" "$stage_timeout" _run_post_merge_review_scanner || true
_pulse_run_optional_stage_with_timeout "pr_review_thread_response" "$stage_timeout" _run_pr_review_thread_response_scanner || true
_pulse_run_optional_stage_with_timeout "auto_decomposer_scanner" "$stage_timeout" _run_auto_decomposer_scanner || true
_pulse_run_optional_stage_with_timeout "dedup_cleanup" "$stage_timeout" run_simplification_dedup_cleanup || true
_pulse_run_optional_stage_with_timeout "fast_fail_prune_expired" "$stage_timeout" fast_fail_prune_expired || true
_pulse_run_blocker_refresh_catchup "$stage_timeout" || true
run_stage_with_timeout "preflight_ownership_reconcile" "$stage_timeout" \
_preflight_ownership_reconcile "$stage_timeout" || true
_pulse_run_resumable_housekeeping_stage "coderabbit_review" "$state_file" \
_pulse_run_optional_stage_with_timeout "coderabbit_review" "$stage_timeout" run_daily_codebase_review
_pulse_run_resumable_housekeeping_stage "post_merge_scanner" "$state_file" \
_pulse_run_optional_stage_with_timeout "post_merge_scanner" "$stage_timeout" _run_post_merge_review_scanner
_pulse_run_resumable_housekeeping_stage "pr_review_thread_response" "$state_file" \
_pulse_run_optional_stage_with_timeout "pr_review_thread_response" "$stage_timeout" _run_pr_review_thread_response_scanner
_pulse_run_resumable_housekeeping_stage "auto_decomposer_scanner" "$state_file" \
_pulse_run_optional_stage_with_timeout "auto_decomposer_scanner" "$stage_timeout" _run_auto_decomposer_scanner
_pulse_run_resumable_housekeeping_stage "dedup_cleanup" "$state_file" \
_pulse_run_optional_stage_with_timeout "dedup_cleanup" "$stage_timeout" run_simplification_dedup_cleanup
_pulse_run_resumable_housekeeping_stage "fast_fail_prune_expired" "$state_file" \
_pulse_run_optional_stage_with_timeout "fast_fail_prune_expired" "$stage_timeout" fast_fail_prune_expired
_pulse_run_resumable_housekeeping_stage "preflight_ownership_reconcile" "$state_file" \
run_stage_with_timeout "preflight_ownership_reconcile" "$stage_timeout" \
_preflight_ownership_reconcile "$stage_timeout"
_pulse_run_resumable_housekeeping_stage "blocker_refresh_catchup" "$state_file" \
_pulse_run_blocker_refresh_catchup "$stage_timeout"
_PULSE_HOUSEKEEPING_CURRENT_STAGE="finishing"
rm -f "$state_file" 2>/dev/null || true
echo "[pulse-wrapper] Async post-dispatch housekeeping: complete" >>"$LOGFILE"

_pulse_release_post_dispatch_housekeeping_lock "$lockdir"
Expand Down Expand Up @@ -1251,7 +1407,11 @@ _pulse_start_post_dispatch_housekeeping() {
set -m 2>/dev/null || true
(
set +m 2>/dev/null || true
trap - EXIT INT TERM
trap - EXIT INT
# GH#34287: attribute deploy-reconciliation kills in pulse.log.
_PULSE_HOUSEKEEPING_CURRENT_STAGE="launch"
trap '_pulse_housekeeping_on_signal TERM' TERM
trap '_pulse_housekeeping_on_signal HUP' HUP
AIDEVOPS_PULSE_STAGE_CYCLE_CLAMP=0
_pulse_run_post_dispatch_housekeeping_stages "$stage_timeout"
) >>"$LOGFILE" 2>&1 &
Expand Down
138 changes: 137 additions & 1 deletion .agents/scripts/tests/test-pulse-post-dispatch-housekeeping.sh
Original file line number Diff line number Diff line change
Expand Up @@ -273,7 +273,8 @@ test_housekeeping_blocker_refresh_catchup() {
touch -t 202601010000 "$DEP_GRAPH_CACHE_FILE"
_pulse_run_post_dispatch_housekeeping_stages 7
got=$(_catchup_stages)
if [[ "$got" != "dep_graph:1 blocked_refresh brief_hold_release ownership_reconcile " ]]; then
# GH#34287: the cheap terminal reconcile runs before the slow catch-up.
if [[ "$got" != "ownership_reconcile dep_graph:1 blocked_refresh brief_hold_release " ]]; then
failures=$((failures + 1))
failmsg="${failmsg} | stale cache order: ${got}"
fi
Expand Down Expand Up @@ -306,11 +307,146 @@ test_housekeeping_blocker_refresh_catchup() {
return 0
}

_dead_pid() {
sleep 0 &
local pid=$!
wait "$pid" 2>/dev/null || true
printf '%s\n' "$pid"
return 0
}

# GH#34287: a run that reclaims a killed run's lock resumes at the first
# incomplete stage; a clean complete run clears the progress state.
test_housekeeping_resumes_after_interrupted_run() {
setup_test_env
local lockdir state_file now failures=0 failmsg=""
lockdir="$(_pulse_post_dispatch_housekeeping_lockdir)"
state_file="$(_pulse_post_dispatch_housekeeping_statefile)"
now=$(date +%s)
mkdir -p "$lockdir"
_dead_pid >"${lockdir}/pid"
printf 'coderabbit_review=%s\npost_merge_scanner=%s\npr_review_thread_response=%s\n' \
"$((now - 60))" "$((now - 30))" "$((now - 4000))" >"$state_file"

_pulse_run_post_dispatch_housekeeping_stages 7

local skipped
for skipped in coderabbit post_merge; do
if grep -q "^${skipped}$" "$STAGE_LOG" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | fresh completed stage ${skipped} re-ran"
fi
done
local expected
for expected in pr_review_thread_response auto_decomposer dedup_cleanup fast_fail_prune ownership_reconcile; do
if ! grep -q "^${expected}$" "$STAGE_LOG" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | incomplete or stale stage ${expected} skipped"
fi
done
if ! grep -q 'reclaiming stale lock' "$LOGFILE" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | stale lock not reclaimed"
fi
if ! grep -q 'resume — skipping coderabbit_review' "$LOGFILE" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | resume skip log missing"
fi
if [[ -e "$state_file" ]]; then
failures=$((failures + 1))
failmsg="${failmsg} | complete run left progress state"
fi

# Resume disabled: every stage runs despite fresh state.
: >"$STAGE_LOG"
printf 'coderabbit_review=%s\n' "$now" >"$state_file"
AIDEVOPS_PULSE_HOUSEKEEPING_RESUME_WINDOW_S=0 _pulse_run_post_dispatch_housekeeping_stages 7
if ! grep -q '^coderabbit$' "$STAGE_LOG" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | disabled resume skipped coderabbit"
fi

if [[ "$failures" -eq 0 ]]; then
print_result "post-dispatch housekeeping resumes at first incomplete stage" 0
else
print_result "post-dispatch housekeeping resumes at first incomplete stage" 1 "$failmsg"
fi
teardown_test_env
return 0
}

_stage_rc_143() { return 143; }
_stage_rc_1() { return 1; }

# GH#34287: a stage killed by a signal is not recorded as complete.
test_housekeeping_signal_killed_stage_not_recorded() {
setup_test_env
local state_file failures=0 failmsg=""
state_file="$(_pulse_post_dispatch_housekeeping_statefile)"
_pulse_run_resumable_housekeeping_stage "killed_stage" "$state_file" _stage_rc_143
_pulse_run_resumable_housekeeping_stage "failed_stage" "$state_file" _stage_rc_1
if grep -q '^killed_stage=' "$state_file" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | signal-killed stage recorded"
fi
if ! grep -q '^failed_stage=[0-9][0-9]*$' "$state_file" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | finished (failed) stage not recorded"
fi
if [[ "$failures" -eq 0 ]]; then
print_result "post-dispatch housekeeping does not record signal-killed stages" 0
else
print_result "post-dispatch housekeeping does not record signal-killed stages" 1 "$failmsg"
fi
teardown_test_env
return 0
}

# GH#34287: SIGTERM (deploy reconciliation) is attributed in pulse.log.
test_async_housekeeping_logs_termination() {
setup_test_env
export AIDEVOPS_PULSE_ASYNC_POST_DISPATCH_HOUSEKEEPING=1
export TEST_STAGE_SLEEP_ONCE=2
local failures=0 failmsg="" child_pid="" attempts=0
_pulse_start_post_dispatch_housekeeping 7
child_pid=$(sed -n 's/.*Async post-dispatch housekeeping: launched pid=\([0-9][0-9]*\).*/\1/p' "$LOGFILE" 2>/dev/null | tail -n 1)
while [[ "$attempts" -lt 20 ]] && ! grep -q '^optional:coderabbit_review:' "$STAGE_LOG" 2>/dev/null; do
sleep 0.25
attempts=$((attempts + 1))
done
if [[ -n "$child_pid" ]]; then
kill -TERM "$child_pid" 2>/dev/null || true
fi
attempts=0
while [[ "$attempts" -lt 12 ]] && ! grep -q 'terminated by SIGTERM' "$LOGFILE" 2>/dev/null; do
sleep 1
attempts=$((attempts + 1))
done
if ! grep -q 'Async post-dispatch housekeeping: terminated by SIGTERM during coderabbit_review' "$LOGFILE" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | termination log missing: $(grep -o 'terminated by.*' "$LOGFILE" 2>/dev/null | tr '\n' ' ')"
fi
if grep -q 'Async post-dispatch housekeeping: complete' "$LOGFILE" 2>/dev/null; then
failures=$((failures + 1))
failmsg="${failmsg} | killed run reported complete"
fi
if [[ "$failures" -eq 0 ]]; then
print_result "post-dispatch housekeeping logs SIGTERM with active stage" 0
else
print_result "post-dispatch housekeeping logs SIGTERM with active stage" 1 "$failmsg"
fi
teardown_test_env
return 0
}

main() {
test_sync_housekeeping_runs_all_stages
test_async_housekeeping_returns_before_slow_stage
test_housekeeping_lock_skips_live_duplicate
test_housekeeping_blocker_refresh_catchup
test_housekeeping_resumes_after_interrupted_run
test_housekeeping_signal_killed_stage_not_recorded
test_async_housekeeping_logs_termination

printf '\n============================================\n'
printf 'Tests run: %d\n' "$TESTS_RUN"
Expand Down
Loading