From 65587b940b7490a1f6eb95866399754e7738316d Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Fri, 18 Sep 2026 19:44:01 +0000 Subject: [PATCH] =?UTF-8?q?=E7=A1=AE=E8=AE=A4=E5=AE=8C=E6=88=90=E4=BA=8B?= =?UTF-8?q?=E4=BB=B6=E5=90=8E=20ACK=20=E6=8C=81=E4=B9=85=E6=B6=88=E6=81=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/gitea_dynamic_runner/controller.py | 20 ++++++++---- tests/test_controller.py | 43 ++++++++++++++++++++++++-- 2 files changed, 55 insertions(+), 8 deletions(-) diff --git a/src/gitea_dynamic_runner/controller.py b/src/gitea_dynamic_runner/controller.py index cccc368..a867326 100644 --- a/src/gitea_dynamic_runner/controller.py +++ b/src/gitea_dynamic_runner/controller.py @@ -176,12 +176,20 @@ async def run_message(message: object, scheduler: OpenSandboxScheduler) -> None: runner_request = RunnerRequest.from_json(message.data) await scheduler.create(runner_request) task = scheduler.active[runner_request.job_id] - while not task.done(): - try: - await asyncio.wait_for(asyncio.shield(task), timeout=30) - except asyncio.TimeoutError: - await message.in_progress() - await task + try: + while not task.done(): + try: + await asyncio.wait_for(asyncio.shield(task), timeout=30) + except asyncio.TimeoutError: + await message.in_progress() + await task + except asyncio.CancelledError: + # A completed webhook cancels the lifecycle monitor after its + # Lifecycle DELETE succeeds. That is successful message handling. + # During controller shutdown the scheduler still owns the job, so + # preserve the unacked message for redelivery. + if runner_request.job_id in scheduler.active: + raise except ValueError as error: LOG.info("runner request already active: %s", error) await message.nak(delay=5) diff --git a/tests/test_controller.py b/tests/test_controller.py index 4daa2a4..163de7d 100644 --- a/tests/test_controller.py +++ b/tests/test_controller.py @@ -1,3 +1,4 @@ +import asyncio import hashlib import hmac import json @@ -31,8 +32,14 @@ def test_rejects_other_actions_labels_and_boolean_id(): completed = queued_job() completed["action"] = "completed" assert controller.accepts(completed) == (False, None) - assert controller.accepts(queued_job(labels=["self-hosted", "other"])) == (False, None) - assert controller.accepts(queued_job(labels=["self-hosted", "pod", "vm"])) == (False, None) + assert controller.accepts(queued_job(labels=["self-hosted", "other"])) == ( + False, + None, + ) + assert controller.accepts(queued_job(labels=["self-hosted", "pod", "vm"])) == ( + False, + None, + ) assert controller.accepts(queued_job(id=True)) == (False, None) @@ -95,3 +102,35 @@ def test_signature_accepts_gitea_and_prefixed_forms(tmp_path, monkeypatch): assert controller.valid_signature(body, digest) assert controller.valid_signature(body, f"sha256={digest}") assert not controller.valid_signature(body, "bad") + + +async def test_completed_lifecycle_cancellation_acks_message(): + request = RunnerRequest.from_webhook(queued_job()) + assert request is not None + + class Message: + data = request.to_json() + acked = False + + async def ack(self): + self.acked = True + + async def nak(self, **kwargs): + raise AssertionError(f"unexpected NAK: {kwargs}") + + async def in_progress(self): + pass + + class Scheduler: + active = {} + + async def create(self, runner_request): + async def completed(): + self.active.pop(runner_request.job_id, None) + raise asyncio.CancelledError + + self.active[runner_request.job_id] = asyncio.create_task(completed()) + + message = Message() + await controller.run_message(message, Scheduler()) + assert message.acked is True