Compare commits

...
Author SHA1 Message Date
panxiao81 65587b940b 确认完成事件后 ACK 持久消息
test / python (pull_request) Successful in 17s
test / shell (pull_request) Successful in 22s
2026-09-18 19:44:01 +00:00
2 changed files with 55 additions and 8 deletions
+8
View File
@@ -176,12 +176,20 @@ async def run_message(message: object, scheduler: OpenSandboxScheduler) -> None:
runner_request = RunnerRequest.from_json(message.data) runner_request = RunnerRequest.from_json(message.data)
await scheduler.create(runner_request) await scheduler.create(runner_request)
task = scheduler.active[runner_request.job_id] task = scheduler.active[runner_request.job_id]
try:
while not task.done(): while not task.done():
try: try:
await asyncio.wait_for(asyncio.shield(task), timeout=30) await asyncio.wait_for(asyncio.shield(task), timeout=30)
except asyncio.TimeoutError: except asyncio.TimeoutError:
await message.in_progress() await message.in_progress()
await task 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: except ValueError as error:
LOG.info("runner request already active: %s", error) LOG.info("runner request already active: %s", error)
await message.nak(delay=5) await message.nak(delay=5)
+41 -2
View File
@@ -1,3 +1,4 @@
import asyncio
import hashlib import hashlib
import hmac import hmac
import json import json
@@ -31,8 +32,14 @@ def test_rejects_other_actions_labels_and_boolean_id():
completed = queued_job() completed = queued_job()
completed["action"] = "completed" completed["action"] = "completed"
assert controller.accepts(completed) == (False, None) 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", "other"])) == (
assert controller.accepts(queued_job(labels=["self-hosted", "pod", "vm"])) == (False, None) False,
None,
)
assert controller.accepts(queued_job(labels=["self-hosted", "pod", "vm"])) == (
False,
None,
)
assert controller.accepts(queued_job(id=True)) == (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, digest)
assert controller.valid_signature(body, f"sha256={digest}") assert controller.valid_signature(body, f"sha256={digest}")
assert not controller.valid_signature(body, "bad") 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