137 lines
4.2 KiB
Python
137 lines
4.2 KiB
Python
import asyncio
|
|
import hashlib
|
|
import hmac
|
|
import json
|
|
|
|
import pytest
|
|
|
|
from gitea_dynamic_runner import controller
|
|
from gitea_dynamic_runner.models import IdentityBinding, RunnerRequest
|
|
|
|
|
|
def queued_job(**overrides):
|
|
job = {
|
|
"id": 47,
|
|
"run_id": 12,
|
|
"name": "publish-image",
|
|
"labels": ["self-hosted", "pod"],
|
|
}
|
|
job.update(overrides)
|
|
return {
|
|
"action": "queued",
|
|
"workflow_job": job,
|
|
"repository": {"full_name": "panxiao81/example"},
|
|
}
|
|
|
|
|
|
def test_accepts_matching_queued_job():
|
|
assert controller.accepts(queued_job()) == (True, "47")
|
|
|
|
|
|
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(id=True)) == (False, None)
|
|
|
|
|
|
def test_runner_request_contains_stable_identity_context():
|
|
request = RunnerRequest.from_webhook(queued_job())
|
|
assert request is not None
|
|
assert request.backend == "pod"
|
|
assert request.repository == "panxiao81/example"
|
|
assert request.job_name == "publish-image"
|
|
assert request.job_id == 47
|
|
assert json.loads(request.to_json()) == {
|
|
"backend": "pod",
|
|
"job_id": 47,
|
|
"job_name": "publish-image",
|
|
"labels": ["self-hosted", "pod"],
|
|
"repository": "panxiao81/example",
|
|
"run_id": 12,
|
|
"schema_version": 1,
|
|
}
|
|
|
|
|
|
def test_runner_request_requires_complete_identity_context():
|
|
assert RunnerRequest.from_webhook(queued_job(run_id=None)) is None
|
|
assert RunnerRequest.from_webhook(queued_job(name=" ")) is None
|
|
|
|
|
|
def test_runner_request_json_round_trip_and_validation():
|
|
request = RunnerRequest.from_webhook(queued_job())
|
|
assert request is not None
|
|
assert RunnerRequest.from_json(request.to_json()) == request
|
|
with pytest.raises(ValueError, match="unsupported"):
|
|
RunnerRequest.from_json(b'{"schema_version":2}')
|
|
|
|
|
|
def test_identity_binding_uses_actual_runner_assignment():
|
|
payload = queued_job(runner_name="gitea-pod-6c47d03d")
|
|
payload["action"] = "in_progress"
|
|
binding = IdentityBinding.from_webhook(payload)
|
|
assert binding is not None
|
|
assert binding.backend == "pod"
|
|
assert binding.runner_name == "gitea-pod-6c47d03d"
|
|
assert binding.repository == "panxiao81/example"
|
|
assert binding.job_name == "publish-image"
|
|
assert json.loads(binding.to_json())["job_id"] == 47
|
|
|
|
|
|
def test_identity_binding_rejects_runner_from_another_pool():
|
|
payload = queued_job(runner_name="gitea-vm-6c47d03d")
|
|
payload["action"] = "in_progress"
|
|
assert IdentityBinding.from_webhook(payload) is None
|
|
|
|
|
|
def test_signature_accepts_gitea_and_prefixed_forms(tmp_path, monkeypatch):
|
|
secret = b"test-secret"
|
|
body = b'{"action":"queued"}'
|
|
secret_file = tmp_path / "secret"
|
|
secret_file.write_bytes(secret)
|
|
monkeypatch.setattr(controller, "WEBHOOK_SECRET_FILE", secret_file)
|
|
digest = hmac.new(secret, body, hashlib.sha256).hexdigest()
|
|
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
|