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
12 changes: 12 additions & 0 deletions src/agent_capacity/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -1444,6 +1444,10 @@ def cleanup_stale() -> dict[str, Any]:
continue
worker_pid = int(job.get("worker_pid", 0))
process_pid = int(job.get("process_pid", 0))
# Older/imported waiting records may not carry process evidence.
# Absence of evidence is not evidence that a managed worker died.
if worker_pid <= 0 and process_pid <= 0:
continue
if pid_alive(worker_pid) or pid_alive(process_pid):
continue
job["state"] = "failed"
Expand Down Expand Up @@ -1864,6 +1868,7 @@ def public_local_job(job: dict[str, Any]) -> dict[str, Any]:

def run_local_job_action(job: dict[str, Any], action: str) -> int:
if action == "status":
cleanup_stale()
current = find_job(job["id"]) or job
print_json({"job": public_local_job(current)})
return 0
Expand Down Expand Up @@ -2178,6 +2183,11 @@ def job_summary(job: dict[str, Any]) -> dict[str, Any]:


def queue_snapshot(history_limit: int = 20) -> dict[str, Any]:
# Queue is an operational truth surface, not a passive dump of the ledger.
# A managed local worker can disappear between polls (host shutdown, disk
# exhaustion, SIGKILL). Reconcile that state before reporting capacity so a
# dead job cannot continue to look active or hold admission hostage.
cleanup_stale()
with locked_jobs() as (job_data, _):
jobs = [dict(job) for job in job_data.get("jobs", [])]
with locked_state() as (lease_data, _):
Expand Down Expand Up @@ -3249,6 +3259,7 @@ def main() -> int:

if args.command == "acquire":
validate_count_and_ttl(args.count, args.ttl)
cleanup_stale()
code, value = acquire(args.workload, args.count, args.owner, args.ttl)
print_json(value)
return code
Expand All @@ -3273,6 +3284,7 @@ def main() -> int:
command = command[1:]
if not command:
raise SystemExit("run requires a command after --")
cleanup_stale()
code, value = acquire(args.workload, args.count, args.owner, args.ttl)
if code:
if args.queue:
Expand Down
30 changes: 30 additions & 0 deletions tests/test_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -386,6 +386,36 @@ def main() -> None:
assert stale_cleanup["jobs"][0]["memory_mb"] == 2500
assert stale_cleanup["released_reservations"] == ["stale-job-token"]

auto_directory = Path(directory) / "auto-stale"
auto_directory.mkdir()
auto_state = auto_directory / "leases.json"
auto_jobs = auto_directory / "jobs.json"
auto_state.write_text(json.dumps({
"version": 1, "leases": [{
"token": "auto-stale-token", "owner": "test:auto-stale", "owner_pid": 99999999,
"workload": "build", "count": 1, "reserved_mb": 4200,
"created_at": int(time.time()) - 120, "expires_at": int(time.time()) + 600,
}],
}))
auto_jobs.write_text(json.dumps({
"version": 1, "jobs": [{
"id": "auto-stale", "provider": "local", "state": "running",
"owner": "test:auto-stale", "workload": "build", "count": 1,
"worker_pid": 99999998, "process_pid": 99999999,
"lease_token": "auto-stale-token", "created_at": int(time.time()) - 120,
}],
}))
_, reconciled_queue = call(auto_state, "queue", "--json", level=80, host_metrics=host_metrics)
assert reconciled_queue["counts"]["local_running"] == 0
assert reconciled_queue["counts"]["reservations"] == 0
assert reconciled_queue["history"][0]["id"] == "auto-stale"
assert reconciled_queue["history"][0]["state"] == "failed"
_, auto_readmitted = call(
auto_state, "acquire", "--workload", "build", "--owner", "test:auto-readmitted",
level=80, host_metrics=host_metrics,
)
assert auto_readmitted["allowed"] is True

parsed_swap = parse_swap_usage("vm.swapusage: total = 8.00G used = 7.27G free = 747.12M")
assert parsed_swap["swap_known"] is True
assert parsed_swap["swap_total_mb"] == 8192
Expand Down