Skip to content

Recover tasks lost to shutdown and stuck in-flight labels - #675

Merged
opohorel merged 8 commits into
packit:mainfrom
opohorel:re-enqueue_tasks
Jul 14, 2026
Merged

opohorel merged 8 commits into
packit:mainfrom
opohorel:re-enqueue_tasks

Conversation

@opohorel

Copy link
Copy Markdown
Collaborator

Pods currently have no SIGTERM handler, so a redeployment silently drops whatever task was mid-processing. Separately, a hard crash (OOM, SIGKILL, node failure) leaves an issue's in-flight Jira label stuck forever, since the fetcher treats any Ymir label besides ymir_retry_needed as "already being handled." This PR adds two safety nets to cover both cases.

Cooperative SIGTERM shutdown (ymir/common/base_utils.py)

  • run_task_loop now accepts a shutdown_event. On shutdown it stops pulling new tasks, cancels whatever's still in flight, and RPUSHes the original (source_queue, payload) back to Redis so nothing is lost. No grace period — tasks run for minutes to hours, so waiting for them to finish naturally would rarely succeed and only delays recovery.
  • The semaphore acquire is raced against shutdown too, not just the poll — otherwise, once every max_concurrent slot is held by a long-running task, the loop would never notice shutdown fired.
  • A BRPOP/poll_fn call already in flight when shutdown fires is abandoned rather than cancelled (cancelling mid-read risks corrupting the connection for the next command drawn from the pool), but its eventual result is still awaited and re-pushed if it turns out to be a real task.
  • Wired into triage_agent, backport_agent, rebase_agent, rebuild_agent, mr_consolidation_agent. terminationGracePeriodSeconds bumped 30s → 45s on the affected deployments to give this room.

Stale in-flight-label safety net (ymir/jira_issue_fetcher/jira_issue_fetcher.py)

  • Added updated to the existing bulk Jira search fields (free — no extra API calls).
  • An in-flight label with no update in STALE_LABEL_THRESHOLD_HOURS (default 24h) and no other Ymir label alongside it is treated as abandoned, flipped to ymir_retry_needed, and re-triaged from scratch.
  • Guarded against re-enqueuing a duplicate when the "abandoned-looking" label's task is actually still sitting in a live Redis queue.

Known limitation

The 24h threshold can still false-positive for backport/rebase/rebuild tasks that are actively processing (not queued) for longer than that — e.g. COPR_BUILD_TIMEOUT (3h) × max_build_attempts (10 default) means a single legitimate task can run 30+ hours. Recommend raising the threshold (or adding a Redis heartbeat) as a follow-up before relying on this safety net for those stages.

@qodo-for-packit

Copy link
Copy Markdown

PR Summary by Qodo

Recover in-flight Redis tasks on SIGTERM; re-triage issues stuck with stale in-flight labels

🐞 Bug fix ✨ Enhancement 🧪 Tests ⚙️ Configuration changes 🕐 40+ Minutes

Grey Divider

AI Description

• Re-queue in-flight Redis tasks on SIGTERM to prevent redeploys dropping work.
• Detect and recover Jira issues stuck with stale in-flight Ymir labels.
• Add unit tests and OpenShift config updates for the new safety nets.
Diagram

graph TD
  OCP{{"OpenShift"}} --> SIG["SIGTERM/SIGINT"] --> SE["shutdown_event"] --> LOOP["run_task_loop"] --> R[("Redis")]
  AP["Agent pod"] --> LOOP
  LOOP -->|"RPUSH on shutdown"| R
  CFG["ConfigMap"] --> JF["JiraIssueFetcher"] --> JIRA{{"Jira API"}} --> JF
  JF -->|"LPUSH triage"| R
  subgraph Legend
    direction LR
    _db[("Database")] ~~~ _svc["Service/Worker"] ~~~ _ext{{"External/System"}}
  end
Loading
High-Level Assessment

The following are alternative approaches to this PR:

1. Redis Streams + consumer groups (pending entries)
  • ➕ Built-in claim/reclaim semantics for crashed consumers
  • ➕ Avoids relying on manual RPUSH bookkeeping
  • ➖ Migration complexity from existing BRPOP/list queues
  • ➖ Operational overhead (stream trimming, consumer group management)
2. Heartbeat/lease key per task (visibility timeout)
  • ➕ Prevents false positives for legitimately long-running work
  • ➕ Can be shared by both shutdown recovery and label-staleness recovery
  • ➖ Requires additional Redis writes during processing
  • ➖ Needs careful TTL and renewal behavior to avoid extending stuck leases forever
3. Use Jira label-add timestamp from changelog instead of fields.updated
  • ➕ More accurate staleness signal for the specific in-flight label
  • ➕ Less susceptible to unrelated comments resetting staleness
  • ➖ May require extra Jira API calls or larger payloads per issue
  • ➖ More complex and potentially slower for large sweeps

Recommendation: The PR’s incremental approach is pragmatic: cooperative SIGTERM handling fixes the common redeploy-loss case immediately, and the stale-label sweep provides a best-effort safety net for hard crashes without extra Jira calls. Given the stated limitation (tasks can legitimately run >24h), a follow-up heartbeat/lease mechanism (or a higher threshold per stage) is the best next step to reduce false-positive re-triage risk while keeping recovery guarantees strong.

Files changed (19) +925 / -53

Bug fix (7) +329 / -43
backport_agent.pyWire cooperative shutdown into backport agent loop +4/-1

Wire cooperative shutdown into backport agent loop

• Installs SIGTERM/SIGINT handler and passes a shutdown_event into run_task_loop so in-flight tasks are cancelled and requeued on termination.

ymir/agents/backport_agent.py

mr_consolidation_agent.pyWire cooperative shutdown into MR consolidation poll loop +13/-3

Wire cooperative shutdown into MR consolidation poll loop

• Updates custom poller to return (source_queue, payload) and passes shutdown_event into run_task_loop. Uses a sentinel queue name since jobs come from a Redis hash and are completed via complete_job rather than RPUSH requeue.

ymir/agents/mr_consolidation_agent.py

rebase_agent.pyWire cooperative shutdown into rebase agent loop +4/-1

Wire cooperative shutdown into rebase agent loop

• Installs SIGTERM/SIGINT handler and passes shutdown_event into run_task_loop so in-flight tasks are cancelled and requeued on termination.

ymir/agents/rebase_agent.py

rebuild_agent.pyWire cooperative shutdown into rebuild agent loop +4/-1

Wire cooperative shutdown into rebuild agent loop

• Installs SIGTERM/SIGINT handler and passes shutdown_event into run_task_loop so in-flight tasks are cancelled and requeued on termination.

ymir/agents/rebuild_agent.py

triage_agent.pyWire cooperative shutdown into triage agent loop +4/-1

Wire cooperative shutdown into triage agent loop

• Installs SIGTERM/SIGINT handler and passes shutdown_event into run_task_loop so in-flight tasks are cancelled and requeued on termination.

ymir/agents/triage_agent.py

base_utils.pyAdd cooperative shutdown + safe shutdown racing to run_task_loop +169/-36

Add cooperative shutdown + safe shutdown racing to run_task_loop

• Introduces install_shutdown_handler and a shutdown-aware run_task_loop that (1) stops pulling new work, (2) cancels active tasks, and (3) RPUSHes their original payloads back to the originating Redis queue. Adds careful shutdown racing to avoid cancelling in-flight Redis I/O while still accounting for late BRPOP results.

ymir/common/base_utils.py

jira_issue_fetcher.pyRe-triage issues with stale in-flight labels +131/-0

Re-triage issues with stale in-flight labels

• Adds a staleness-based safety net that treats abandoned in-flight labels as retry-needed and re-triages issues from scratch using the bulk-fetched fields.updated timestamp. Guards against duplicate enqueues by skipping the flip if the issue already has a live queued task in Redis.

ymir/jira_issue_fetcher/jira_issue_fetcher.py

Tests (2) +582 / -1
test_base_utils.pyAdd unit tests for shutdown-aware run_task_loop +275/-0

Add unit tests for shutdown-aware run_task_loop

• Adds a FakeRedis harness and comprehensive async tests covering shutdown before start, cancelling and repushing active tasks, orphaned BRPOP behavior, and interruptible idle sleep for custom pollers. Also tests install_shutdown_handler signal registration.

ymir/common/tests/unit/test_base_utils.py

test_jira_issue_fetcher.pyAdd tests for stale in-flight label detection and recovery +307/-1

Add tests for stale in-flight label detection and recovery

• Adds unit coverage for _is_label_stale and _find_stale_in_flight_label, plus push_issues_to_queue behavior for stale labels (reenqueue, dry-run behavior, queued-in-Redis dedup, and failure handling). Updates the mocked Jira search fields to include updated.

ymir/jira_issue_fetcher/tests/unit/test_jira_issue_fetcher.py

Other (10) +14 / -9
configmap-jira-issue-fetcher-env.ymlAdd env var for stale in-flight label threshold +5/-0

Add env var for stale in-flight label threshold

• Introduces STALE_LABEL_THRESHOLD_HOURS (default 24) with rationale comments. This config is shared across Jira issue fetcher cronjobs.

openshift/configmap-jira-issue-fetcher-env.yml

deployment-backport-agent-c10s.ymlIncrease backport agent termination grace period +1/-1

Increase backport agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to give the cooperative shutdown logic time to cancel and requeue.

openshift/deployment-backport-agent-c10s.yml

deployment-backport-agent-c9s.ymlIncrease backport agent termination grace period +1/-1

Increase backport agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to give the cooperative shutdown logic time to cancel and requeue.

openshift/deployment-backport-agent-c9s.yml

deployment-mr-consolidation-agent-c10s.ymlIncrease MR consolidation agent termination grace period +1/-1

Increase MR consolidation agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to support cooperative shutdown behavior.

openshift/deployment-mr-consolidation-agent-c10s.yml

deployment-mr-consolidation-agent-c9s.ymlIncrease MR consolidation agent termination grace period +1/-1

Increase MR consolidation agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to support cooperative shutdown behavior.

openshift/deployment-mr-consolidation-agent-c9s.yml

deployment-rebase-agent-c10s.ymlIncrease rebase agent termination grace period +1/-1

Increase rebase agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to support cooperative shutdown behavior.

openshift/deployment-rebase-agent-c10s.yml

deployment-rebase-agent-c9s.ymlIncrease rebase agent termination grace period +1/-1

Increase rebase agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to support cooperative shutdown behavior.

openshift/deployment-rebase-agent-c9s.yml

deployment-rebuild-agent-c10s.ymlIncrease rebuild agent termination grace period +1/-1

Increase rebuild agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to support cooperative shutdown behavior.

openshift/deployment-rebuild-agent-c10s.yml

deployment-rebuild-agent-c9s.ymlIncrease rebuild agent termination grace period +1/-1

Increase rebuild agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to support cooperative shutdown behavior.

openshift/deployment-rebuild-agent-c9s.yml

deployment-triage-agent.ymlIncrease triage agent termination grace period +1/-1

Increase triage agent termination grace period

• Bumps terminationGracePeriodSeconds from 30s to 45s to support cooperative shutdown behavior.

openshift/deployment-triage-agent.yml

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request implements a graceful shutdown mechanism for Ymir agents and a safety net for stale in-flight Jira labels. Specifically, it introduces signal handlers to catch SIGTERM/SIGINT, allowing the task loop to stop pulling new tasks, cancel active tasks, and safely re-push their payloads back to Redis. Additionally, the Jira issue fetcher is updated to detect and recover abandoned in-flight labels by flipping them to a retry-needed state. The review feedback is highly constructive, highlighting potential issues with task leaks during outer cancellation in _race_shutdown, platform compatibility on Windows for signal handlers, and the need for robust exception handling when re-pushing tasks to Redis during shutdown.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread ymir/common/base_utils.py
Comment thread ymir/common/base_utils.py
Comment thread ymir/common/base_utils.py Outdated
Comment thread ymir/common/base_utils.py
@qodo-for-packit

qodo-for-packit Bot commented Jul 10, 2026 •

Copy link
Copy Markdown

Code Review by Qodo

🐞 Bugs (0) 📘 Rule violations (0) 📜 Skill insights (0)

Context used
✅ Compliance rules (platform): 7 rules

Grey Divider


Action required

1. Consolidation job lost ✓ Resolved 🐞 Bug ≡ Correctness
Description
On SIGTERM, mr_consolidation_agent cancels the active job but still calls complete_job() (which
deletes the active entry from the Redis hash), while run_task_loop repushes the payload to a
sentinel Redis list key that no worker consumes, permanently dropping the consolidation job.
Code

ymir/agents/mr_consolidation_agent.py[R1348-1363]

+        async def poll_consolidation_queue() -> tuple[bytes, bytes] | None:
            job = await pick_next_job(redis_conn)
            if job is None:
                return None
            redis_logger.info("Picked job for %s/%s", job.package, job.target_branch)
-            return job.model_dump_json().encode()
+            # This queue is a Redis Hash (pick_next_job/complete_job), not a
+            # list run_task_loop can RPUSH back into on shutdown — re-push
+            # doesn't apply here (see process_task's finally: complete_job
+            # already drops cancelled-mid-flight jobs the same way it drops
+            # any other failure). The sentinel is only so a cancelled job
+            # still satisfies run_task_loop's generic (source_queue, payload)
+            # bookkeeping; nothing ever reads from it.
+            return b"mr_consolidation", job.model_dump_json().encode()

        async def process_task(payload: bytes) -> None:
            job = MergeConsolidationJob.model_validate_json(payload)
Relevance

⭐⭐⭐ High

Team actively hardens Redis queue correctness; run_task_loop recovery is key goal (PR630);
consolidation queue introduced recently (PR633).

PR-#630
PR-#633

ⓘ Recommendations generated based on similar findings in past PRs

Evidence
mr_consolidation_agent returns a fake source queue for run_task_loop bookkeeping, but run_task_loop
repushes to that queue name as if it were a real Redis list. Meanwhile, the underlying consolidation
queue is a Redis hash where pick_next_job moves a pending job into an active field and complete_job
deletes that active field, so cancellation results in permanent loss rather than recovery.

ymir/agents/mr_consolidation_agent.py[1348-1406]
ymir/common/base_utils.py[240-277]
ymir/common/merge_queue.py[74-113]
ymir/common/merge_queue.py[116-130]

Agent prompt
The issue below was found during a code review. Follow the provided context and guidance below and implement a solution

## Issue description
`mr_consolidation_agent` uses a custom poller backed by a Redis Hash (`merge_queue.py`). On shutdown, `run_task_loop` will repush active tasks to `source_queue`, but the consolidation poller returns a sentinel `b"mr_consolidation"` queue name (a Redis List key no one reads), while `process_task` always calls `complete_job(...)` in `finally`, deleting the active hash entry. This combination drops the job on routine redeployments.

## Issue Context
- `pick_next_job()` promotes `:pending` to `:active` by deleting the pending field and setting active.
- `complete_job()` deletes the `:active` field.
- Therefore, cancellation without a hash-level “demote active back to pending” loses the job.

## Fix Focus Areas
- ymir/agents/mr_consolidation_agent.py[1348-1416]
- ymir/common/merge_queue.py[74-130]
- ymir/common/base_utils.py[129-277]

## Suggested fix approach
1. Add a hash-aware abandonment path for consolidation jobs (e.g., `abandon_job(...)`) that demotes the active job back to pending (or otherwise preserves it) on cancellation/shutdown.
2. In `mr_consolidation_agent.process_task`, detect cancellation (`except asyncio.CancelledError`) and call the abandonment path instead of `complete_job`.
3. Prevent `run_task_loop` from blindly RPUSHing these hash-backed jobs (e.g., add an optional `repush_fn` callback / `enable_repush` flag, or allow `source_queue=None` and skip repush).

ⓘ Copy this prompt and use it to remediate the issue with your preferred AI generation tools



Remediation recommended

2. Repush may duplicate tasks ✓ Resolved 🐞 Bug ☼ Reliability
Description
run_task_loop snapshots active and repushes every entry on shutdown without filtering out tasks
that are already done but haven’t yet executed their done-callback, which can re-enqueue
already-completed work; additionally, there’s no shutdown check between receiving a poll result and
creating the processing task, so a task can begin briefly after shutdown is requested.
Code

ymir/common/base_utils.py[R233-256]

+        task_loop_logger.info("Received task from queue.")
+
+        source_queue, payload = result
+        t = asyncio.create_task(_run(payload))
+        active[t] = (source_queue, payload)
+        t.add_done_callback(lambda _t: active.pop(_t, None))
+
+    if active:
+        # Snapshot BEFORE cancelling: the done-callback above pops entries
+        # out of `active` as each task transitions to done, which happens
+        # *during* the gather below. Iterating `active` after the gather
+        # would see an already-emptied dict and silently re-push nothing.
+        to_repush = list(active.items())
+        task_loop_logger.info(
+            "Shutting down: cancelling %d active task(s) and re-pushing to Redis",
+            len(to_repush),
+        )
+        for t, _ in to_repush:
+            t.cancel()
+        await asyncio.gather(*(t for t, _ in to_repush), return_exceptions=True)
+
+        for _t, (source_queue, payload) in to_repush:
+            await fix_await(redis_conn.rpush(source_queue, payload))
+            task_loop_logger.info("Re-pushed task to %s on shutdown", source_queue)
Relevance

⭐⭐ Medium

Duplication fixes sometimes accepted (TOCTOU dedup PR459), but similar re-queue-dup concerns also
rejected (PR527).

PR-#459
PR-#527

ⓘ Recommendations generated based on similar findings in past PRs

Evidence
The shutdown path snapshots active.items() and repushes all of them, but entries are only removed
via a done-callback, so a finished task can still be present at snapshot time. Also, the worker task
is created immediately after polling without re-checking shutdown state.

ymir/common/base_utils.py[195-256]

Agent prompt
The issue below was found during a code review. Follow the provided context and guidance below and implement a solution

## Issue description
`run_task_loop` relies on a task done-callback to remove entries from `active`, but on shutdown it snapshots `active.items()` and repushes them all. A task can be `done()` but still present in `active` until its callback runs, so shutdown can repush work that already completed successfully (duplicate processing). Separately, shutdown can be set after a poll result is returned but before the worker task is created, allowing a brief start of processing after shutdown.

## Issue Context
- `active` is mutated asynchronously by a done-callback.
- Shutdown repush currently repushes the snapshot unconditionally.

## Fix Focus Areas
- ymir/common/base_utils.py[195-277]

## Suggested fix approach
1. Before building `to_repush`, filter by task state (e.g., only include entries where `not t.done()`).
2. After `result` is received but before `create_task(_run(payload))`, add a shutdown check; if set, immediately repush the polled `(source_queue, payload)`, release the semaphore slot, and exit the loop without starting the task.

ⓘ Copy this prompt and use it to remediate the issue with your preferred AI generation tools


Grey Divider

Qodo Logo

Comment thread ymir/agents/mr_consolidation_agent.py
Comment thread ymir/common/base_utils.py Outdated
@TomasTomecek
TomasTomecek self-requested a review July 13, 2026 14:30

@TomasTomecek TomasTomecek left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I just reviewed it with Cursor (Composer) and Sonnet 5 and both agree it's very well done and most edge cases are sorted out. One unsolved that stood out to me is:

mr_consolidation has no recovery path for its own stuck state (moderate, undocumented)
When a consolidation task is cancelled, process_task deliberately leaves the :active hash field in place rather than calling complete_job() — correctly matching pre-PR crash behavior. But pick_next_job's Lua
script refuses to promote a new pending job for that package/branch while :active is set (merge_queue.py:85), and nothing in this PR (or elsewhere) ever clears a stale :active entry — unlike
triage/backport/rebase/rebuild, which get the new Jira-label staleness sweep. So any redeploy that catches an in-flight consolidation job permanently blocks future consolidation for that package/branch until
someone manually HDELs the field. The PR's "Known limitation" section calls out the 24h Jira threshold false-positive risk but doesn't mention this — and this one is arguably worse (silent, permanent, no
self-healing), not just a threshold tuning problem. Worth a callout/follow-up ticket at minimum.

Hard for me to say if Sonnet is correct here.

opohorel added 7 commits July 13, 2026 18:11
Pods currently have no SIGTERM handler, so a redeployment silently
drops whatever task was popped off a Redis queue but not yet finished
processing. Track each in-flight task's (source_queue, payload); on
shutdown, stop pulling new work, cancel what's still running, and
RPUSH the original payloads back so nothing is lost. There's no grace
period — task processing runs for minutes to hours (e.g. build
polling), so waiting for it to finish naturally would rarely succeed
and only delays recovery.

Also races the semaphore acquire (not just the poll) against shutdown:
without it, once every max_concurrent slot is held by a long-running
task, sem.acquire() for the next iteration blocks forever waiting for a
slot that never frees up naturally, and the loop never notices
shutdown at all — exactly the scenario this exists to handle. Abandons
(rather than cancels) an in-flight BRPOP/poll_fn call on shutdown,
since redis-py never resets a connection on CancelledError and
cancelling mid-read risks a stale response corrupting the next command
drawn from the same pool (e.g. our own re-push RPUSH calls).

Assisted-by: Claude (Cursor)
Each agent's queue-mode main() now creates a shutdown_event, installs
the signal handler, and passes it through to run_task_loop so
redeployments stop pulling new work and re-push whatever was in flight
back to Redis instead of dropping it.

mr_consolidation_agent's job queue is a Redis Hash (pick_next_job /
complete_job), not a list run_task_loop can RPUSH back into, so re-push
doesn't apply there — its poll_fn now returns a sentinel source_queue
just to satisfy run_task_loop's generic bookkeeping. A cancelled
mid-flight job is dropped the same way process_task's own finally
already drops any other failure via complete_job; this is a pre-existing
gap, not a regression, and is out of scope here.

Assisted-by: Claude (Cursor)
SIGTERM handling (previous commits) recovers planned redeployments
instantly, but a hard crash (SIGKILL, OOM, node failure) still leaves
an issue's in-flight label (ymir_triage_in_progress /
ymir_triaged_backport / _rebase / _rebuild) stuck forever, since the
fetcher currently treats any Ymir label besides ymir_retry_needed as
"already being handled" and skips the issue on every sweep.

Add "updated" to the existing bulk Jira search fields (free — same
paginated response already being fetched, no extra API calls) and use
it to detect an in-flight label that's had no update in
STALE_LABEL_THRESHOLD_HOURS (default 24h, conservative pending real
task duration data from Phoenix traces) with no other Ymir label
alongside it — the same "no terminal outcome" signal triage_agent's
own dedup check already uses. Recovery flips the stuck label to
ymir_retry_needed and reuses the existing retry-needed re-triage path,
since the original Redis task payload is unrecoverable by this point
and re-triage is the only generically correct fallback.

Assisted-by: Claude (Cursor)
_race_shutdown abandons an in-flight BRPOP/poll_fn call on shutdown
rather than cancelling it (redis-py never resets a connection on
CancelledError, so cancelling mid-read risks corrupting the next
command drawn from the pool). But it never accounted for that
abandoned call's eventual result: if it resolved with a real task
after shutdown had already started winding down, that task was popped
from Redis and then silently dropped - never processed, never
re-pushed.

Have _race_shutdown hand the abandoned task back to the caller instead
of firing-and-forgetting it, and resolve it in run_task_loop's
shutdown tail before returning: await it (bounded by poll_timeout, the
same cap BRPOP itself already uses, so this isn't a grace period for
task processing) and RPUSH its result back if it turns out to be a
real task rather than a natural timeout.

Bump terminationGracePeriodSeconds from 30s to 45s for the affected
agent deployments to give this bounded wait room within the shutdown
window, alongside the (near-instant) active-task cancel-and-repush.

Assisted-by: Claude (Cursor)
The stale-label recovery path flipped an abandoned-looking in-flight
label straight to ymir_retry_needed and added the issue to
remove_issues_for_retry, which bypasses the existing_keys dedup check
in the push loop below - the same check every other path relies on to
avoid double-queuing an issue. If the label only looked abandoned
because its task is still legitimately sitting in a live Redis queue
(the SIGTERM handler already re-pushed it, or a downstream queue is
simply backed up - ymir_triaged_backport et al. persist for as long
as the task waits its turn, not just during a crash window), this
started a second, concurrent agent run on top of the one already
queued.

Check existing_keys before treating the label as recoverable: if the
issue's payload is still found in a live queue, leave the label alone
and let it drain naturally instead of flipping it.

Assisted-by: Claude (Cursor)
- Wrap _race_shutdown's asyncio.wait in try/except CancelledError so
  both internal tasks are cleaned up if the caller is itself cancelled
- Wrap each individual RPUSH in try/except during shutdown so a
  transient Redis failure for one task doesn't abort the rest
- Move orphan-poll RPUSH into the try body (was in else, which isn't
  covered by the preceding except)
- Filter already-done tasks out of the repush snapshot to avoid
  duplicating work that already ran

Assisted-by: Claude (Cursor)
Before this PR, SIGTERM killed the process outright so the finally
block never ran.  Now that the shutdown handler delivers a real
CancelledError, the finally: complete_job() call would silently
HDEL the :active hash entry, losing the job.  Catch CancelledError
and re-raise without cleanup so the entry stays in :active (matching
pre-PR behavior).  A proper requeue/sweep is follow-up work.

Assisted-by: Claude (Cursor)
@opohorel

Copy link
Copy Markdown
Collaborator Author

I just reviewed it with Cursor (Composer) and Sonnet 5 and both agree it's very well done and most edge cases are sorted out. One unsolved that stood out to me is:

mr_consolidation has no recovery path for its own stuck state (moderate, undocumented)
When a consolidation task is cancelled, process_task deliberately leaves the :active hash field in place rather than calling complete_job() — correctly matching pre-PR crash behavior. But pick_next_job's Lua
script refuses to promote a new pending job for that package/branch while :active is set (merge_queue.py:85), and nothing in this PR (or elsewhere) ever clears a stale :active entry — unlike
triage/backport/rebase/rebuild, which get the new Jira-label staleness sweep. So any redeploy that catches an in-flight consolidation job permanently blocks future consolidation for that package/branch until
someone manually HDELs the field. The PR's "Known limitation" section calls out the 24h Jira threshold false-positive risk but doesn't mention this — and this one is arguably worse (silent, permanent, no
self-healing), not just a threshold tuning problem. Worth a callout/follow-up ticket at minimum.

Hard for me to say if Sonnet is correct here.

I've added a commit which fixes the issue pointed out. but there is still one issue and that is the mr_consolidation tasks are not requed, they are just thrown away. the architecture of mr_consolidation differs to the other agents. I'm not sure if I want to expand the scope of this PR, so maybe I'll create a follow-up ticket for it? wdyt?

When a consolidation task is cancelled by shutdown, the :active hash
field is left in place to avoid silent data loss.  But pick_next_job
refuses to promote :pending while :active exists for the same
package/branch, so without cleanup a redeploy permanently blocks
future consolidation for that pair.

Add sweep_stale_active_jobs() to merge_queue.py: on every poll cycle
it scans :active entries and removes any whose activated_at is older
than a configurable threshold (default 6h, env STALE_ACTIVE_THRESHOLD_HOURS).

Staleness is measured from activated_at (set by pick_next_job at
promotion time), not submitted_at (set at initial queueing), to avoid
falsely sweeping jobs that waited in :pending behind a backlog.
Deletion uses an atomic compare-and-delete Lua script so a concurrent
complete_job + pick_next_job can't cause the sweep to remove a fresh
entry.

Assisted-by: Claude (Cursor)

@TomasTomecek TomasTomecek left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM, thanks!

agreed with opening a followup

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants