Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
65587b940b |
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user