From 69cd85c164e3527859c0ce7081d5d32d03ea7ee9 Mon Sep 17 00:00:00 2001 From: Maja Massarini Date: Mon, 28 Sep 2026 14:48:36 +0200 Subject: [PATCH 1/2] Route error requeues through priority without maintainer noise --- openshift/Makefile | 7 ++-- openshift/scripts/requeue_error.py | 41 +++++++++++-------- .../scripts/tests/unit/test_requeue_error.py | 37 +++++++++++++++++ ymir/agents/reproducer_agent.py | 1 + ymir/agents/triage_agent.py | 1 + ymir/common/models.py | 5 +++ 6 files changed, 73 insertions(+), 19 deletions(-) create mode 100644 openshift/scripts/tests/unit/test_requeue_error.py diff --git a/openshift/Makefile b/openshift/Makefile index 27f7113a5..732330689 100644 --- a/openshift/Makefile +++ b/openshift/Makefile @@ -85,14 +85,15 @@ show-error-list: @oc exec deployment/valkey -- valkey-cli --no-raw LRANGE error_list 0 -1 | python3 scripts/error_list.py --file - # Requeue a specific error_list entry (by the [id=...] shown in show-error-list) -# back onto its original trigger queue, with attempts reset to 0. +# onto its priority trigger queue, with attempts reset to 0. By default it +# avoids user-triggered maintainer feedback; set USER_TRIGGERED=true to opt in. # Required: ERROR_ID requeue-error: @if [ -z "$(ERROR_ID)" ]; then \ - echo "Usage: make requeue-error ERROR_ID=42"; \ + echo "Usage: make requeue-error ERROR_ID=42 [USER_TRIGGERED=true]"; \ exit 1; \ fi - python3 scripts/requeue_error.py $(ERROR_ID) + python3 scripts/requeue_error.py $(ERROR_ID) $(if $(filter true TRUE 1 yes YES,$(USER_TRIGGERED)),--user-triggered) logs-triage: oc logs -f deployment/triage-agent diff --git a/openshift/scripts/requeue_error.py b/openshift/scripts/requeue_error.py index 754f29b55..963bab888 100644 --- a/openshift/scripts/requeue_error.py +++ b/openshift/scripts/requeue_error.py @@ -7,15 +7,12 @@ or entries recorded for a payload that never parsed into a `Task` in the first place, have no `queue`/`task` to requeue and are reported as non-requeueable. -The task's `attempts` counter is reset to 0 and `user_triggered` is forced to -True before requeuing, on the assumption that whoever runs this has already -fixed the underlying issue and wants a fresh retry budget. Forcing -`user_triggered` also matters functionally: triage and reproducer skip -processing outright when a terminal `ymir_*_errored`/similar label is still -on the issue and the task isn't user-triggered — exactly the state a -just-requeued task is in until it gets a chance to run and clear that label -itself. It also routes the task onto the priority (`_todo`) twin of its -queue, same as any other maintainer-triggered run. +The task's `attempts` counter is reset to 0 before requeuing. By default, +`user_triggered` is set to False to keep acknowledgement and result comments +for maintainers low. Requeued tasks still go to the priority (`_todo`) twin +of their queue. A separate `requeued_from_error_list` marker lets triage and +reproducer process them even while the issue has a terminal errored label. +Pass `--user-triggered` to opt into maintainer-triggered comments and labels. The entry is located and its replacement task computed here, client-side (a plain LRANGE + local JSON parsing) rather than inside Redis: scanning and @@ -34,8 +31,10 @@ size limits than a Redis value. Usage: - make requeue-error ERROR_ID=42 # from openshift/ (preferred) + make requeue-error ERROR_ID=42 # priority requeue, fewer comments + make requeue-error ERROR_ID=42 USER_TRIGGERED=true # priority/user-triggered requeue python3 scripts/requeue_error.py 42 # same, run directly + python3 scripts/requeue_error.py 42 --user-triggered # priority/user-triggered python3 scripts/requeue_error.py 42 --dry-run # show the plan, don't mutate anything """ @@ -168,6 +167,15 @@ def atomic_requeue( return int(out.strip()) +def prepare_requeue_task(task: dict, source_queue: str, user_triggered: bool) -> tuple[str, str]: + """Reset a task and send it to the priority queue, with optional maintainer feedback.""" + task["attempts"] = 0 + task["user_triggered"] = user_triggered + task["requeued_from_error_list"] = True + target_queue = source_queue if source_queue.endswith("_todo") else f"{source_queue}_todo" + return target_queue, json.dumps(task) + + def main() -> None: ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) ap.add_argument( @@ -180,6 +188,11 @@ def main() -> None: ap.add_argument( "--dry-run", action="store_true", help="Print the plan without pushing or removing anything" ) + ap.add_argument( + "--user-triggered", + action="store_true", + help="Set user_triggered=true for maintainer comments and labels (priority queue either way)", + ) args = ap.parse_args() blob = fetch_via_oc(args.queue, args.deployment) @@ -198,15 +211,11 @@ def main() -> None: ) old_attempts = task.get("attempts", 0) - task["attempts"] = 0 - task["user_triggered"] = True - if not target_queue.endswith("_todo"): - target_queue = f"{target_queue}_todo" - task_json = json.dumps(task) + target_queue, task_json = prepare_requeue_task(task, target_queue, args.user_triggered) print( f"Requeuing error_id={args.error_id} ({jira_issue}) onto '{target_queue}' " - f"(attempts {old_attempts} -> 0, user_triggered -> true)" + f"(attempts {old_attempts} -> 0, user_triggered -> {str(args.user_triggered).lower()})" ) if args.dry_run: print(f"[dry-run] Would atomically remove from {args.queue} and LPUSH {target_queue}: {task_json}") diff --git a/openshift/scripts/tests/unit/test_requeue_error.py b/openshift/scripts/tests/unit/test_requeue_error.py new file mode 100644 index 000000000..b6191133c --- /dev/null +++ b/openshift/scripts/tests/unit/test_requeue_error.py @@ -0,0 +1,37 @@ +import json + +import pytest + +from openshift.scripts.requeue_error import prepare_requeue_task +from ymir.common.models import Task + + +@pytest.mark.parametrize("source_queue", ["backport_queue_c9s", "backport_queue_c9s_todo"]) +def test_prepare_requeue_task_defaults_to_priority_queue_without_maintainer_feedback(source_queue) -> None: + task = {"metadata": {"issue": "RHEL-1"}, "attempts": 3, "user_triggered": True} + + queue, task_json = prepare_requeue_task(task, source_queue, user_triggered=False) + + assert queue == "backport_queue_c9s_todo" + assert json.loads(task_json) == { + "metadata": {"issue": "RHEL-1"}, + "attempts": 0, + "user_triggered": False, + "requeued_from_error_list": True, + } + restored = Task.model_validate_json(task_json) + assert restored.requeued_from_error_list is True + assert restored.user_triggered is False + + +def test_prepare_requeue_task_uses_priority_queue_when_requested() -> None: + task = {"attempts": 2, "user_triggered": False} + + queue, task_json = prepare_requeue_task(task, "backport_queue_c10s", user_triggered=True) + + assert queue == "backport_queue_c10s_todo" + assert json.loads(task_json) == { + "attempts": 0, + "user_triggered": True, + "requeued_from_error_list": True, + } diff --git a/ymir/agents/reproducer_agent.py b/ymir/agents/reproducer_agent.py index ecb5712b3..1fda0f2d2 100644 --- a/ymir/agents/reproducer_agent.py +++ b/ymir/agents/reproducer_agent.py @@ -1235,6 +1235,7 @@ async def process_task(payload): terminal_ymir_labels and JiraLabels.REPRODUCER_IN_PROGRESS.value not in current_labels and not user_triggered + and not task.requeued_from_error_list ): logger.info( f"Skipping duplicate reproducer for {input_data.jira_issue} — " diff --git a/ymir/agents/triage_agent.py b/ymir/agents/triage_agent.py index 223637a3d..f481dd184 100644 --- a/ymir/agents/triage_agent.py +++ b/ymir/agents/triage_agent.py @@ -1527,6 +1527,7 @@ async def _process_triage_locked(task, input): and JiraLabels.RETRY_NEEDED.value not in current_labels and JiraLabels.TRIAGE_IN_PROGRESS.value not in current_labels and not user_triggered + and not task.requeued_from_error_list ): logger.info( f"Skipping duplicate triage for {input.issue} — " diff --git a/ymir/common/models.py b/ymir/common/models.py index 3bf76ab4b..b3bda388d 100644 --- a/ymir/common/models.py +++ b/ymir/common/models.py @@ -115,6 +115,11 @@ class Task(BaseModel): "failure labels that are otherwise suppressed. Error comments are " "posted regardless of this flag.", ) + requeued_from_error_list: bool = Field( + default=False, + description="True when an operator requeued this task from error_list, so an existing " + "error label does not cause triage or reproducer to skip it as a duplicate.", + ) def to_json(self) -> str: """Convert to JSON string for Redis queue storage.""" From 40d8e94e6afc100709203fbdd6bea080dbdc8815 Mon Sep 17 00:00:00 2001 From: Maja Massarini Date: Mon, 28 Sep 2026 14:48:58 +0200 Subject: [PATCH 2/2] Clear resolved workflow failures from error list --- ymir/agents/backport_agent.py | 8 ++ ymir/agents/rebase_agent.py | 27 ++++++ ymir/agents/rebuild_agent.py | 27 ++++++ ymir/agents/reproducer_agent.py | 27 ++++++ .../tests/unit/test_rebase_consolidation.py | 36 +++++++- .../tests/unit/test_rebuild_consolidation.py | 35 ++++++++ .../tests/unit/test_reproducer_agent.py | 87 ++++++++++++++++--- ymir/agents/triage_agent.py | 6 ++ ymir/common/error_list.py | 55 ++++++++++++ ymir/common/tests/unit/test_error_list.py | 77 ++++++++++++++++ 10 files changed, 371 insertions(+), 14 deletions(-) create mode 100644 ymir/common/error_list.py create mode 100644 ymir/common/tests/unit/test_error_list.py diff --git a/ymir/agents/backport_agent.py b/ymir/agents/backport_agent.py index 4ca9a1070..42c943e35 100644 --- a/ymir/agents/backport_agent.py +++ b/ymir/agents/backport_agent.py @@ -90,6 +90,7 @@ run_task_loop, ) from ymir.common.constants import JiraLabels, RedisQueues +from ymir.common.error_list import clear_resolved_errors from ymir.common.issue_lock import issue_lock from ymir.common.logging_setup import configure_logging, current_jira_issue, get_trajectory_writeable from ymir.common.mock_repos import get_mock_local_tool_env @@ -2201,6 +2202,13 @@ async def retry( state.backport_result.model_dump_json(), ) ) + await clear_resolved_errors( + redis, + backport_data.jira_issue, + backport_queue, + target_branch=dist_git_branch, + dry_run=dry_run, + ) else: logger.warning( f"Backport failed for {backport_data.jira_issue}: {state.backport_result.error}" diff --git a/ymir/agents/rebase_agent.py b/ymir/agents/rebase_agent.py index 77625ee0f..93587b5d1 100644 --- a/ymir/agents/rebase_agent.py +++ b/ymir/agents/rebase_agent.py @@ -56,6 +56,7 @@ ) from ymir.common.base_utils import fix_await, install_shutdown_handler, redis_client, run_task_loop from ymir.common.constants import JiraLabels, RedisQueues +from ymir.common.error_list import clear_resolved_errors from ymir.common.issue_lock import issue_lock from ymir.common.logging_setup import configure_logging, current_jira_issue, get_trajectory_writeable from ymir.common.mock_repos import get_mock_local_tool_env @@ -99,6 +100,25 @@ def _consolidated_issue_keys( return list(dict.fromkeys([primary_issue] + [item.issue_key for item in consolidated_issues])) +async def _clear_rebase_resolved_errors( + redis_conn, + rebase_data: RebaseData, + queue: str, + *, + target_branch: str, + dry_run: bool, +) -> None: + """Clear prior rebase errors for the primary issue and each resolved sibling.""" + for issue_key in _consolidated_issue_keys(rebase_data.jira_issue, rebase_data.consolidated_issues): + await clear_resolved_errors( + redis_conn, + issue_key, + queue, + target_branch=target_branch, + dry_run=dry_run, + ) + + def get_instructions() -> str: return render_template("rebase/instructions.j2") @@ -920,6 +940,13 @@ async def retry( state.rebase_result.model_dump_json(), ) ) + await _clear_rebase_resolved_errors( + redis, + rebase_data, + rebase_queue, + target_branch=dist_git_branch, + dry_run=dry_run, + ) else: logger.warning(f"Rebase failed for {rebase_data.jira_issue}: {state.rebase_result.error}") # Label all consolidated issues with failure status diff --git a/ymir/agents/rebuild_agent.py b/ymir/agents/rebuild_agent.py index 9b4332ef2..dcf8e7410 100644 --- a/ymir/agents/rebuild_agent.py +++ b/ymir/agents/rebuild_agent.py @@ -35,6 +35,7 @@ ) from ymir.common.base_utils import fix_await, install_shutdown_handler, redis_client, run_task_loop from ymir.common.constants import JiraLabels, RedisQueues +from ymir.common.error_list import clear_resolved_errors from ymir.common.issue_lock import issue_lock from ymir.common.logging_setup import configure_logging, current_jira_issue from ymir.common.mock_repos import get_mock_local_tool_env @@ -56,6 +57,25 @@ redis_logger = logging.getLogger("agent.redis") +async def _clear_rebuild_resolved_errors( + redis_conn, + rebuild_data: RebuildData, + queue: str, + *, + target_branch: str, + dry_run: bool, +) -> None: + """Clear prior rebuild errors for the primary issue and each resolved sibling.""" + for issue_key in dict.fromkeys(rebuild_data.all_jira_issues): + await clear_resolved_errors( + redis_conn, + issue_key, + queue, + target_branch=target_branch, + dry_run=dry_run, + ) + + async def main() -> None: init_sentry() @@ -669,6 +689,13 @@ async def retry( ).model_dump_json(), ) ) + await _clear_rebuild_resolved_errors( + redis, + rebuild_data, + rebuild_queue, + target_branch=dist_git_branch, + dry_run=dry_run, + ) else: logger.warning(f"Rebuild failed for {rebuild_data.jira_issue}: {state.rebuild_error}") for issue_key in dict.fromkeys(rebuild_data.all_jira_issues): diff --git a/ymir/agents/reproducer_agent.py b/ymir/agents/reproducer_agent.py index 1fda0f2d2..7856c4447 100644 --- a/ymir/agents/reproducer_agent.py +++ b/ymir/agents/reproducer_agent.py @@ -38,6 +38,7 @@ from ymir.common.base_utils import fix_await, redis_client, run_task_loop from ymir.common.constants import JiraLabels, RedisQueues from ymir.common.delayed_queue import promote_due_tasks, schedule_task +from ymir.common.error_list import clear_resolved_errors from ymir.common.logging_setup import configure_logging, current_jira_issue from ymir.common.mock_repos import get_mock_local_tool_env from ymir.common.models import ( @@ -265,6 +266,31 @@ def _should_finalize_jira(result: OutputSchema) -> bool: return not result.retryable_error and not result.lock_deferred +async def _clear_finalized_reproducer_errors( + redis_conn, input_data: InputSchema, result: OutputSchema, *, dry_run: bool +) -> int: + """Clear earlier failures only when the final Jira label resolves this work.""" + if not _should_finalize_jira(result): + return 0 + # An adapted test still needs its MR update. If that step failed, the + # original test's presence must not make this run count as resolved. + if result.adapted_existing and not result.success: + return 0 + if _determine_result_label(result) not in { + JiraLabels.REPRODUCER_CREATED, + JiraLabels.REPRODUCER_ALREADY_EXISTS, + JiraLabels.REPRODUCER_NOT_REPRODUCIBLE, + }: + return 0 + return await clear_resolved_errors( + redis_conn, + input_data.jira_issue, + RedisQueues.REPRODUCER_QUEUE.value, + target_branch=input_data.target_branch, + dry_run=dry_run, + ) + + def _needs_merge_request(result: OutputSchema) -> bool: """Whether orchestration should commit/push an MR for this result.""" if result.lock_deferred or result.retryable_error: @@ -1440,6 +1466,7 @@ async def retry( logger.info( f"Pushed {input_data.jira_issue} to {RedisQueues.COMPLETED_REPRODUCER_LIST.value}" ) + await _clear_finalized_reproducer_errors(redis, input_data, output, dry_run=dry_run) finally: try: await release_reproducer_lock( diff --git a/ymir/agents/tests/unit/test_rebase_consolidation.py b/ymir/agents/tests/unit/test_rebase_consolidation.py index 608fcc04e..783caea1d 100644 --- a/ymir/agents/tests/unit/test_rebase_consolidation.py +++ b/ymir/agents/tests/unit/test_rebase_consolidation.py @@ -1,6 +1,6 @@ import pytest -from ymir.agents.rebase_agent import _consolidated_issue_keys +from ymir.agents.rebase_agent import _clear_rebase_resolved_errors, _consolidated_issue_keys from ymir.agents.rebase_consolidation import ( add_jira_tickets_to_latest_changelog_entry, build_rebase_siblings_jql, @@ -9,7 +9,7 @@ has_new_latest_changelog_entry, uses_autochangelog, ) -from ymir.common.models import ConsolidatedIssue +from ymir.common.models import ConsolidatedIssue, RebaseData from ymir.common.utils import extract_text_from_adf @@ -151,6 +151,38 @@ def test_consolidated_issue_keys_deduplicates_siblings_and_primary(): ) == ["RHEL-100", "RHEL-200"] +@pytest.mark.asyncio +@pytest.mark.parametrize("dry_run", [False, True]) +async def test_successful_consolidated_rebase_clears_primary_and_distinct_siblings(monkeypatch, dry_run): + calls = [] + + async def mock_clear(redis_conn, issue, queue, *, target_branch, dry_run): + calls.append((redis_conn, issue, queue, target_branch, dry_run)) + + monkeypatch.setattr("ymir.agents.rebase_agent.clear_resolved_errors", mock_clear) + redis = object() + data = RebaseData( + package="expat", + version="2.7.0", + jira_issue="RHEL-100", + consolidated_issues=[ + ConsolidatedIssue(issue_key="RHEL-200"), + ConsolidatedIssue(issue_key="RHEL-100"), + ConsolidatedIssue(issue_key="RHEL-200"), + ConsolidatedIssue(issue_key="RHEL-300"), + ], + ) + + await _clear_rebase_resolved_errors( + redis, data, "rebase_queue_c10s", target_branch="rhel-10.3", dry_run=dry_run + ) + + assert calls == [ + (redis, issue, "rebase_queue_c10s", "rhel-10.3", dry_run) + for issue in ("RHEL-100", "RHEL-200", "RHEL-300") + ] + + def test_build_rebase_siblings_jql(): jql = build_rebase_siblings_jql("RHEL-100", "dotnet10.0", "rhel-10.2") assert 'component = "dotnet10.0"' in jql diff --git a/ymir/agents/tests/unit/test_rebuild_consolidation.py b/ymir/agents/tests/unit/test_rebuild_consolidation.py index 296440d02..394c174f8 100644 --- a/ymir/agents/tests/unit/test_rebuild_consolidation.py +++ b/ymir/agents/tests/unit/test_rebuild_consolidation.py @@ -1,4 +1,8 @@ +import pytest + +from ymir.agents.rebuild_agent import _clear_rebuild_resolved_errors from ymir.agents.rebuild_consolidation import build_rebuild_siblings_jql +from ymir.common.models import ConsolidatedIssue, RebuildData def test_build_rebuild_siblings_jql(): @@ -29,6 +33,37 @@ def test_build_rebuild_siblings_jql_no_filter_for_nonmodular(): assert "cf[10669]" not in jql +@pytest.mark.asyncio +@pytest.mark.parametrize("dry_run", [False, True]) +async def test_successful_consolidated_rebuild_clears_primary_and_distinct_siblings(monkeypatch, dry_run): + calls = [] + + async def mock_clear(redis_conn, issue, queue, *, target_branch, dry_run): + calls.append((redis_conn, issue, queue, target_branch, dry_run)) + + monkeypatch.setattr("ymir.agents.rebuild_agent.clear_resolved_errors", mock_clear) + redis = object() + data = RebuildData( + package="git-lfs", + jira_issue="RHEL-100", + consolidated_issues=[ + ConsolidatedIssue(issue_key="RHEL-200"), + ConsolidatedIssue(issue_key="RHEL-100"), + ConsolidatedIssue(issue_key="RHEL-200"), + ConsolidatedIssue(issue_key="RHEL-300"), + ], + ) + + await _clear_rebuild_resolved_errors( + redis, data, "rebuild_queue_c10s", target_branch="rhel-10.3", dry_run=dry_run + ) + + assert calls == [ + (redis, issue, "rebuild_queue_c10s", "rhel-10.3", dry_run) + for issue in ("RHEL-100", "RHEL-200", "RHEL-300") + ] + + def test_build_rebuild_siblings_jql_escapes_component_quotes(): jql = build_rebuild_siblings_jql("RHEL-100", 'comp"name', "rhel-9.8.z") assert r'component = "comp\"name"' in jql diff --git a/ymir/agents/tests/unit/test_reproducer_agent.py b/ymir/agents/tests/unit/test_reproducer_agent.py index a12daa498..a32b6e1e5 100644 --- a/ymir/agents/tests/unit/test_reproducer_agent.py +++ b/ymir/agents/tests/unit/test_reproducer_agent.py @@ -13,6 +13,7 @@ PreparedTestsClone, _bootstrap_tests_clone, _build_mr_title, + _clear_finalized_reproducer_errors, _cve_only_needles, _determine_comment_resolution, _determine_result_label, @@ -32,8 +33,15 @@ ) from ymir.agents.tasks import InvalidReproducerConfigError, fetch_reproducer_config from ymir.common.base_utils import check_subprocess -from ymir.common.constants import JiraLabels -from ymir.common.models import MergeRequestDetails, ReproducerInputSchema, ReproducerOutputSchema, Task +from ymir.common.constants import JiraLabels, RedisQueues +from ymir.common.models import ( + ErrorData, + ErrorListEntry, + MergeRequestDetails, + ReproducerInputSchema, + ReproducerOutputSchema, + Task, +) def _output(**overrides) -> ReproducerOutputSchema: @@ -106,6 +114,63 @@ def test_should_finalize_jira_false_for_retryable_error(): assert _should_finalize_jira(_output(success=True)) is True +@pytest.mark.asyncio +@pytest.mark.parametrize( + ("result_overrides", "dry_run", "expected_removed"), + [ + ({"success": False, "not_reproducible_reason": "race"}, False, 1), + ({"success": False, "test_already_exists": True}, False, 1), + pytest.param( + { + "success": False, + "test_already_exists": True, + "adapted_existing": True, + "summary": "MR creation failed", + }, + False, + 0, + id="adapted-existing-mr-update-failed", + ), + ({"success": True}, False, 1), + ({"success": False}, False, 0), + ({"success": False, "retryable_error": True, "not_reproducible_reason": "race"}, False, 0), + ({"success": False, "lock_deferred": True, "test_already_exists": True}, False, 0), + ({"success": False, "not_reproducible_reason": "race"}, True, 0), + ], +) +async def test_finalized_reproducer_clears_only_resolved_errors(result_overrides, dry_run, expected_removed): + input_data = ReproducerInputSchema(jira_issue="RHEL-12345", package="libfoo", target_branch="c10s") + old_error = ( + ErrorListEntry( + error_id=1, + queue=RedisQueues.REPRODUCER_QUEUE.value, + task=Task(metadata=input_data.model_dump()), + error=ErrorData(jira_issue=input_data.jira_issue, details="prior failure"), + ) + .model_dump_json() + .encode() + ) + + class ErrorListRedis: + def __init__(self): + self.entries = [old_error] + self.lrem_calls = 0 + + async def lrange(self, *_args): + return list(self.entries) + + async def lrem(self, _key, _count, raw): + self.lrem_calls += 1 + self.entries.remove(raw) + return 1 + + redis = ErrorListRedis() + await _clear_finalized_reproducer_errors(redis, input_data, _output(**result_overrides), dry_run=dry_run) + + assert len(redis.entries) == 1 - expected_removed + assert redis.lrem_calls == expected_removed + + def test_needs_merge_request(): assert _needs_merge_request(_output(success=True)) is True assert _needs_merge_request(_output(success=True, test_already_exists=True)) is False @@ -788,7 +853,11 @@ async def _mock_run_task_loop(_redis, _queues, process_fn, **_kw): async def _mock_redis_client(*_args, **_kwargs): redis_mock = flexmock() redis_mock.should_receive("lpush").replace_with(_async_noop) - redis_mock.should_receive("model_dump_json").replace_with(_async_noop) + + async def _empty_error_list(*_args): + return [] + + redis_mock.should_receive("lrange").replace_with(_empty_error_list) yield redis_mock flexmock(r_agent).should_receive("init_sentry") @@ -841,9 +910,7 @@ async def _mock_jira_metadata(*_args, **_kwargs): return ["ymir_reproducer_created"], "New" async def _mock_workflow(*_args, **_kwargs): - result = flexmock(success=True, retryable_error=False, lock_deferred=False, summary="ok") - result.should_receive("model_dump_json").and_return("{}") - return flexmock(result=result) + return flexmock(result=_output()) flexmock(agent_tasks).should_receive("get_jira_issue_metadata").replace_with(_mock_jira_metadata) flexmock(r_agent).should_receive("run_workflow").replace_with(_mock_workflow).once() @@ -861,9 +928,7 @@ async def _mock_jira_metadata(*_args, **_kwargs): return ["ymir_reproducer_created", "ymir_reproducer_in_progress"], "New" async def _mock_workflow(*_args, **_kwargs): - result = flexmock(success=True, retryable_error=False, lock_deferred=False, summary="ok") - result.should_receive("model_dump_json").and_return("{}") - return flexmock(result=result) + return flexmock(result=_output()) flexmock(agent_tasks).should_receive("get_jira_issue_metadata").replace_with(_mock_jira_metadata) flexmock(r_agent).should_receive("run_workflow").once().replace_with(_mock_workflow) @@ -880,9 +945,7 @@ async def _mock_jira_metadata(*_args, **_kwargs): return [], "New" async def _mock_workflow(*_args, **_kwargs): - result = flexmock(success=True, retryable_error=False, lock_deferred=False, summary="ok") - result.should_receive("model_dump_json").and_return("{}") - return flexmock(result=result) + return flexmock(result=_output()) flexmock(agent_tasks).should_receive("get_jira_issue_metadata").replace_with(_mock_jira_metadata) flexmock(r_agent).should_receive("run_workflow").once().replace_with(_mock_workflow) diff --git a/ymir/agents/triage_agent.py b/ymir/agents/triage_agent.py index f481dd184..6432f47fe 100644 --- a/ymir/agents/triage_agent.py +++ b/ymir/agents/triage_agent.py @@ -43,6 +43,7 @@ from ymir.common.base_utils import fix_await, install_shutdown_handler, redis_client, run_task_loop from ymir.common.config import load_rhel_config from ymir.common.constants import YMIR_COMMENT_MARKER, JiraLabels, RedisQueues +from ymir.common.error_list import clear_resolved_errors from ymir.common.issue_lock import issue_lock from ymir.common.logging_setup import configure_logging, current_jira_issue, get_trajectory_writeable from ymir.common.mock_repos import get_mock_local_tool_env @@ -1959,6 +1960,11 @@ async def retry( input.issue, ) + if output.resolution != Resolution.ERROR: + await clear_resolved_errors( + redis, input.issue, RedisQueues.TRIAGE_QUEUE.value, dry_run=dry_run + ) + shutdown_event = asyncio.Event() install_shutdown_handler(asyncio.get_running_loop(), shutdown_event) await run_task_loop( diff --git a/ymir/common/error_list.py b/ymir/common/error_list.py new file mode 100644 index 000000000..9743f83ad --- /dev/null +++ b/ymir/common/error_list.py @@ -0,0 +1,55 @@ +"""Remove resolved workflow failures from the active error list.""" + +import json +import logging + +from ymir.common.base_utils import fix_await +from ymir.common.constants import RedisQueues + +logger = logging.getLogger(__name__) + + +async def clear_resolved_errors( + redis_conn, + jira_issue: str, + queue: str, + *, + target_branch: str | None = None, + dry_run: bool = False, +) -> int: + """Clear earlier failures for this issue, workflow queue, and target branch. + + The priority and ordinary queues represent the same workflow. Match the + branch as well because one issue can have independent work on several RHEL + streams. Remove the original bytes with LREM so concurrent list changes are + preserved. Legacy entries without a queue or task cannot be matched safely. + A dry run does not read or modify the production error list. + """ + if dry_run: + return 0 + + resolved_queue = queue.removesuffix("_todo") + removed = 0 + try: + entries = await fix_await(redis_conn.lrange(RedisQueues.ERROR_LIST.value, 0, -1)) + for raw in entries: + try: + entry = json.loads(raw) + task = entry.get("task") + error = entry.get("error") + if not isinstance(task, dict) or not isinstance(error, dict): + continue + if entry.get("queue", "").removesuffix("_todo") != resolved_queue: + continue + if error.get("jira_issue") != jira_issue: + continue + if task.get("metadata", {}).get("target_branch") != target_branch: + continue + except (TypeError, ValueError, AttributeError): + continue + removed += await fix_await(redis_conn.lrem(RedisQueues.ERROR_LIST.value, 1, raw)) + except Exception: + # The workflow has already succeeded. A cleanup failure must not turn + # that success into a retry or create another error-list entry. + logger.exception("Failed to clear resolved errors for %s from %s", jira_issue, resolved_queue) + return removed diff --git a/ymir/common/tests/unit/test_error_list.py b/ymir/common/tests/unit/test_error_list.py new file mode 100644 index 000000000..63c2bf470 --- /dev/null +++ b/ymir/common/tests/unit/test_error_list.py @@ -0,0 +1,77 @@ +import pytest + +from ymir.common.constants import RedisQueues +from ymir.common.error_list import clear_resolved_errors +from ymir.common.models import ErrorData, ErrorListEntry, Task + + +class ListRedis: + def __init__(self, entries): + self.entries = entries + self.lrange_calls = 0 + self.lrem_calls = 0 + + async def lrange(self, key, start, stop): + self.lrange_calls += 1 + assert (key, start, stop) == (RedisQueues.ERROR_LIST.value, 0, -1) + return list(self.entries) + + async def lrem(self, key, count, value): + self.lrem_calls += 1 + assert (key, count) == (RedisQueues.ERROR_LIST.value, 1) + self.entries.remove(value) + return 1 + + +def entry(issue, queue, branch=None): + task = Task(metadata={"target_branch": branch}) + return ( + ErrorListEntry( + error_id=1, + queue=queue, + task=task, + error=ErrorData(jira_issue=issue, details="failed"), + ) + .model_dump_json() + .encode() + ) + + +@pytest.mark.asyncio +async def test_success_clears_prior_errors_for_same_issue_workflow_and_branch(): + matching = entry("RHEL-1", "backport_queue_c9s_todo", "rhel-9.9.0") + another_matching = entry("RHEL-1", "backport_queue_c9s", "rhel-9.9.0") + other_branch = entry("RHEL-1", "backport_queue_c9s", "rhel-9.8.0") + other_workflow = entry("RHEL-1", "rebase_queue_c9s", "rhel-9.9.0") + other_issue = entry("RHEL-2", "backport_queue_c9s", "rhel-9.9.0") + legacy = b'{"jira_issue":"RHEL-1","details":"old failure"}' + redis = ListRedis([matching, another_matching, other_branch, other_workflow, other_issue, legacy]) + + removed = await clear_resolved_errors(redis, "RHEL-1", "backport_queue_c9s", target_branch="rhel-9.9.0") + + assert removed == 2 + assert redis.entries == [other_branch, other_workflow, other_issue, legacy] + + +@pytest.mark.asyncio +async def test_successful_dry_run_keeps_matching_production_error_entry(): + matching = entry("RHEL-1", "backport_queue_c9s", "rhel-9.9.0") + redis = ListRedis([matching]) + + removed = await clear_resolved_errors( + redis, "RHEL-1", "backport_queue_c9s", target_branch="rhel-9.9.0", dry_run=True + ) + + assert removed == 0 + assert redis.entries == [matching] + assert redis.lrange_calls == 0 + assert redis.lrem_calls == 0 + + +@pytest.mark.asyncio +async def test_cleanup_failure_does_not_fail_completed_workflow(): + class BrokenRedis: + async def lrange(self, *_args): + raise ConnectionError("Redis unavailable") + + assert await clear_resolved_errors(BrokenRedis(), "RHEL-1", "triage_queue") == 0