确认完成事件后 ACK 持久消息
This commit is contained in:
@@ -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)
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user