Repository navigation
feat(llc): close out heartbeat runs whose agent stopped reporting (#16817) #16821
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Closed
Closed
Changes from all commits
Commits
Show all changes
11 commits
Select commit
Hold shift + click to select a range
8d2972a
feat(llc): close out heartbeat runs whose agent stopped reporting (#1…
mrveiss be4bc5e
fix(llc): declare stalled_run_sweep as an eager task module in the sp…
mrveiss a563c4d
Merge branch 'main' into issue-16817-stalled-run-sweep
mrveiss a01653a
fix(llc): read the stall timeout via env_int so a malformed value can…
mrveiss 0f96137
Merge branch 'main' into issue-16817-stalled-run-sweep
mrveiss 3b54a0f
Merge branch 'main' into issue-16817-stalled-run-sweep
mrveiss 376c6a7
Merge branch 'main' into issue-16817-stalled-run-sweep
mrveiss 9de1c7a
Merge branch 'main' into issue-16817-stalled-run-sweep
mrveiss b17d7a8
Merge branch 'main' into issue-16817-stalled-run-sweep
mrveiss 3dd4bbb
Merge branch 'main' into issue-16817-stalled-run-sweep
mrveiss 78d0d36
Merge branch 'main' into issue-16817-stalled-run-sweep
github-actions[bot] File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Some comments aren't visible on the classic Files Changed page.
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,118 @@ | ||
| # Copyright 2025-2026 mrveiss | ||
| # SPDX-License-Identifier: Apache-2.0 | ||
| # AutoBot - AI-Powered Automation Platform | ||
| # Author: mrveiss | ||
| """Celery-beat sweep: close out heartbeat runs whose agent stopped reporting (#16817). | ||
|
|
||
| ``llc_heartbeat_runs`` has always recorded that a run started. Nothing read it back | ||
| to decide a run had *stopped*. A run whose agent died mid-work therefore stayed | ||
| ``running`` for ever, and the only thing that ever noticed was a person spotting an | ||
| inconsistency — which is how two abandoned runs were found on 2026-09-16, one of them | ||
| holding three already-merged worktrees that no live session could release. | ||
|
|
||
| The rule is deliberately narrow: a run that has been non-terminal for longer than | ||
| ``LLC_RUN_STALL_TIMEOUT_SECONDS`` is marked ``TIMEOUT`` with an error that says the | ||
| sweep decided it, not the adapter. ``TIMEOUT`` is the existing status for "ran out of | ||
| time" and reusing it keeps the state machine as it is; the distinction that matters — | ||
| *we lost contact* versus *the adapter reported a timeout* — lives in ``error``, which | ||
| is what the reader needs to tell them apart. | ||
|
|
||
| What this sweep deliberately does NOT do is release what the run held. Claims, | ||
| assignments and workspace leases are released in #16818, which models the lease this | ||
| sweep will then have something to release. Marking a run stalled and leaving its | ||
| holdings is an improvement over never noticing, and it is not the finished job: a | ||
| status nobody acts on is close to what we have today. | ||
| """ | ||
|
|
||
| import asyncio | ||
| import logging | ||
| from datetime import datetime, timedelta, timezone | ||
|
|
||
| from celery import shared_task | ||
| from sqlalchemy import select | ||
|
|
||
| from autobot_shared.env_utils import env_int | ||
| from llc.models.enums import LLCRunStatus | ||
| from llc.models.heartbeat_run import LLCHeartbeatRun | ||
| from user_management.database import get_async_session_factory | ||
| from utils.celery_reliability import ( | ||
| CELERY_MAX_RETRIES, | ||
| CELERY_RETRY_BACKOFF_MAX, | ||
| CELERY_TRANSIENT_ERRORS, | ||
| DeadLetterTask, | ||
| ) | ||
|
|
||
| logger = logging.getLogger(__name__) | ||
|
|
||
| # Env-var-backed, never a literal at the call site. Six hours is longer than any | ||
| # adapter run observed to date and short enough that an abandoned run is noticed | ||
| # the same working day. Lower it and long legitimate runs get killed; raise it and | ||
| # the failure it exists to catch stays invisible for longer. | ||
| STALL_TIMEOUT_SECONDS = env_int("LLC_RUN_STALL_TIMEOUT_SECONDS", 6 * 60 * 60) | ||
|
|
||
| #: Statuses a run can sit in while still believed to be alive. | ||
| NON_TERMINAL_STATUSES = (LLCRunStatus.QUEUED.value, LLCRunStatus.RUNNING.value) | ||
|
|
||
| #: Written to ``error`` so a swept run is never mistaken for an adapter-reported | ||
| #: timeout. The text is asserted by the tests — it is the only thing that tells a | ||
| #: reader which of the two happened. | ||
| STALL_ERROR = "run stalled: no completion reported within {seconds}s; closed out by the stalled-run sweep (#16817)" | ||
|
|
||
|
|
||
| @shared_task( | ||
| name="llc.scheduler.stalled_run_sweep.run_stalled_run_sweep", | ||
| bind=True, | ||
| base=DeadLetterTask, | ||
| autoretry_for=CELERY_TRANSIENT_ERRORS, | ||
| retry_backoff=True, | ||
| retry_jitter=True, | ||
| retry_backoff_max=CELERY_RETRY_BACKOFF_MAX, | ||
| max_retries=CELERY_MAX_RETRIES, | ||
| ) | ||
| def run_stalled_run_sweep(self: object) -> dict: # type: ignore[type-arg] | ||
| """Sync Celery entry point — closes out runs that stopped reporting.""" | ||
| try: | ||
| loop = asyncio.get_event_loop() | ||
| except RuntimeError: | ||
| loop = asyncio.new_event_loop() | ||
| asyncio.set_event_loop(loop) | ||
| stalled = loop.run_until_complete(_async_sweep()) | ||
| return {"stalled": stalled} | ||
|
|
||
|
|
||
| def _cutoff(now: datetime | None = None) -> datetime: | ||
| return (now or datetime.now(timezone.utc)) - timedelta(seconds=STALL_TIMEOUT_SECONDS) | ||
|
|
||
|
|
||
| async def _async_sweep() -> int: | ||
| """Select non-terminal runs older than the cutoff and close them out.""" | ||
| factory = get_async_session_factory() | ||
| cutoff = _cutoff() | ||
| stalled = 0 | ||
| async with factory() as session: | ||
| # started_at is NULL for a run that never got picked up; fall back to | ||
| # created_at so a run that was queued and abandoned is swept too. A bare | ||
| # ``started_at <= cutoff`` would exclude those rows through SQL | ||
| # three-valued logic and strand them exactly as the disposal sweep found. | ||
| result = await session.execute( | ||
| select(LLCHeartbeatRun).where( | ||
| LLCHeartbeatRun.status.in_(NON_TERMINAL_STATUSES), | ||
| LLCHeartbeatRun.finished_at.is_(None), | ||
| ) | ||
| ) | ||
| for run in result.scalars().all(): | ||
| age_anchor = run.started_at or run.created_at | ||
| if age_anchor is None or age_anchor > cutoff: | ||
| continue | ||
| run.status = LLCRunStatus.TIMEOUT.value | ||
| run.finished_at = datetime.now(timezone.utc) | ||
| run.error = STALL_ERROR.format(seconds=STALL_TIMEOUT_SECONDS) | ||
| stalled += 1 | ||
| await session.commit() | ||
| # Logged unconditionally: a sweep that found nothing and a sweep that did not | ||
| # run must not look the same in the logs. | ||
| logger.info("Stalled-run sweep closed out %d run(s) older than %ds", stalled, STALL_TIMEOUT_SECONDS) | ||
| return stalled | ||
|
|
||
|
|
||
| __all__ = ["run_stalled_run_sweep", "STALL_TIMEOUT_SECONDS", "STALL_ERROR", "NON_TERMINAL_STATUSES"] |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,69 @@ | ||
| # Copyright 2025-2026 mrveiss | ||
| # SPDX-License-Identifier: Apache-2.0 | ||
| # AutoBot - AI-Powered Automation Platform | ||
| # Author: mrveiss | ||
| """The stalled-run sweep must sweep, and must not sweep what is still alive (#16817). | ||
|
|
||
| The defect this closes was not a wrong result — it was no result at all: nothing read | ||
| ``llc_heartbeat_runs`` back to decide a run had stopped. So the tests that matter are | ||
| the ones that fail if the sweep ever stops selecting: a sweep that quietly matches | ||
| nothing is indistinguishable from the state before it existed. | ||
| """ | ||
|
|
||
| import importlib | ||
| from datetime import datetime, timedelta, timezone | ||
|
|
||
| import pytest | ||
|
|
||
|
|
||
| @pytest.fixture | ||
| def sweep(monkeypatch): | ||
| """Import with a short timeout so the tests do not depend on the six-hour default.""" | ||
| monkeypatch.setenv("LLC_RUN_STALL_TIMEOUT_SECONDS", "60") | ||
| module = importlib.import_module("llc.scheduler.stalled_run_sweep") | ||
| return importlib.reload(module) | ||
|
|
||
|
|
||
| def test_the_timeout_is_env_var_backed_not_a_literal(sweep): | ||
| assert sweep.STALL_TIMEOUT_SECONDS == 60 | ||
|
|
||
|
|
||
| def test_the_default_is_used_when_the_env_var_is_absent(monkeypatch): | ||
| monkeypatch.delenv("LLC_RUN_STALL_TIMEOUT_SECONDS", raising=False) | ||
| module = importlib.reload(importlib.import_module("llc.scheduler.stalled_run_sweep")) | ||
| assert module.STALL_TIMEOUT_SECONDS == 6 * 60 * 60 | ||
|
|
||
|
|
||
| def test_a_run_older_than_the_cutoff_is_past_it(sweep): | ||
| old = datetime.now(timezone.utc) - timedelta(seconds=120) | ||
| assert old <= sweep._cutoff() | ||
|
|
||
|
|
||
| def test_a_fresh_run_is_not_past_the_cutoff(sweep): | ||
| fresh = datetime.now(timezone.utc) - timedelta(seconds=5) | ||
| assert fresh > sweep._cutoff() | ||
|
|
||
|
|
||
| def test_only_non_terminal_statuses_are_candidates(sweep): | ||
| """A completed or failed run is finished; sweeping it would rewrite history.""" | ||
| assert set(sweep.NON_TERMINAL_STATUSES) == {"queued", "running"} | ||
| for terminal in ("completed", "failed", "interrupted", "timeout"): | ||
| assert terminal not in sweep.NON_TERMINAL_STATUSES | ||
|
|
||
|
|
||
| def test_the_error_says_the_sweep_decided_it(sweep): | ||
| """A swept run must be tellable from an adapter-reported timeout. | ||
|
|
||
| Both end as status ``timeout``. The status alone cannot answer "did the adapter | ||
| time out, or did we lose contact", and that is the question an operator asks | ||
| first. The error text is the only place that distinguishes them. | ||
| """ | ||
| message = sweep.STALL_ERROR.format(seconds=60) | ||
| assert "sweep" in message | ||
| assert "60" in message | ||
| assert "16817" in message | ||
|
|
||
|
|
||
| def test_the_sweep_is_registered_as_a_named_celery_task(sweep): | ||
| """An unregistered task is a sweep that never runs — the state before this.""" | ||
| assert sweep.run_stalled_run_sweep.name == "llc.scheduler.stalled_run_sweep.run_stalled_run_sweep" |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,67 @@ | ||
| # Copyright 2025-2026 mrveiss | ||
| # SPDX-License-Identifier: Apache-2.0 | ||
| # AutoBot - AI-Powered Automation Platform | ||
| # Author: mrveiss | ||
| """Every @shared_task under llc/scheduler must actually register (#16817 review). | ||
|
|
||
| `celery_beat_registration_test.py` checks that every *scheduled* task resolves to a | ||
| registered one. An unscheduled task is invisible to it -- and so is a task that is | ||
| neither scheduled nor imported, which is a decorator that runs at no point in the | ||
| process's life. | ||
|
|
||
| That is how #16821 shipped its first head: `run_stalled_run_sweep` carried | ||
| `@shared_task`, was absent from `llc/scheduler/__init__.py`'s eager-import block and | ||
| from `beat_schedule`, and therefore never registered. A sweep written to detect work | ||
| that silently never runs, which silently never ran. Celery reports nothing, because | ||
| from its side the task does not exist. | ||
|
|
||
| This guard reads the source rather than the registry: every module under | ||
| `llc/scheduler/` that defines a `@shared_task` must be imported by the package | ||
| `__init__`, because `autodiscover_tasks(related_name=None)` imports only that file. | ||
| """ | ||
|
|
||
| import re | ||
|
|
||
| from repo_tests._paths import repo_root | ||
|
|
||
| SCHEDULER = repo_root() / "autobot-backend" / "llc" / "scheduler" | ||
| _SHARED_TASK = re.compile(r"^@shared_task", re.M) | ||
|
|
||
|
|
||
| def _modules_defining_a_shared_task() -> set[str]: | ||
| out = set() | ||
| for path in SCHEDULER.glob("*.py"): | ||
| if path.name.startswith("__") or path.name.endswith("_test.py"): | ||
| continue | ||
| if _SHARED_TASK.search(path.read_text(encoding="utf-8")): | ||
| out.add(path.stem) | ||
| return out | ||
|
|
||
|
|
||
| def _modules_eagerly_imported() -> set[str]: | ||
| init = (SCHEDULER / "__init__.py").read_text(encoding="utf-8") | ||
| return set(re.findall(r"^from \.(\w+) import", init, re.M)) | ||
|
|
||
|
|
||
| def test_every_shared_task_module_is_eagerly_imported(): | ||
| defining = _modules_defining_a_shared_task() | ||
| assert defining, "no @shared_task found under llc/scheduler — this guard has gone blind" | ||
| missing = sorted(defining - _modules_eagerly_imported()) | ||
| assert not missing, ( | ||
| f"these modules define a @shared_task but are not imported by " | ||
| f"llc/scheduler/__init__.py, so their decorators never run and Celery never " | ||
| f"registers the task: {missing}. autodiscover_tasks(related_name=None) imports " | ||
| f"only the package __init__ (GH#12318)." | ||
| ) | ||
|
|
||
|
|
||
| def test_the_guard_can_see_a_module_it_would_fail_on(): | ||
| """Negative control: the detector must actually detect. | ||
|
|
||
| A guard whose finder silently matches nothing passes for ever and proves nothing. | ||
| This asserts the two halves disagree when they should -- that a module defining a | ||
| task but absent from the import set is reported, rather than swallowed. | ||
| """ | ||
| defining = {"a_task_module", "already_imported"} | ||
| imported = {"already_imported"} | ||
| assert sorted(defining - imported) == ["a_task_module"] | ||
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win
Exercise the detector with source fixtures.
This negative control bypasses
_modules_defining_a_shared_task()and_modules_eagerly_imported(). A broken glob, regex, or source parser can therefore pass this test.Create temporary scheduler fixtures and monkeypatch
SCHEDULER. Add one module that defines@shared_taskwithout an eager import, and one with the matching import.As per path instructions: “Every detector needs a contrast pair: a fixture that SHOULD trip it and one that should not.”
🤖 Prompt for AI Agents
Source: Path instructions