将 OpenSandbox 调度改为持久事件循环
This commit is contained in:
@@ -88,6 +88,7 @@ class OpenSandboxScheduler:
|
||||
self.tokens = tokens
|
||||
self.on_finished = on_finished
|
||||
self.active: dict[int, asyncio.Task[None]] = {}
|
||||
self.runners: dict[str, tuple[int, str, str]] = {}
|
||||
|
||||
async def create(self, request: RunnerRequest) -> str:
|
||||
if request.job_id in self.active:
|
||||
@@ -106,11 +107,12 @@ class OpenSandboxScheduler:
|
||||
raise
|
||||
|
||||
task = asyncio.create_task(
|
||||
self._monitor(request, sandbox_id, nonce),
|
||||
self._monitor(request, sandbox_id, nonce, runner_name),
|
||||
name=f"opensandbox-{sandbox_id}",
|
||||
)
|
||||
task.add_done_callback(self._report)
|
||||
self.active[request.job_id] = task
|
||||
self.runners[runner_name] = (request.job_id, sandbox_id, nonce)
|
||||
LOG.info(
|
||||
"OpenSandbox runner created sandbox=%s runner=%s job_id=%d "
|
||||
"repository=%s job_name=%r pool=ci-%s",
|
||||
@@ -124,7 +126,11 @@ class OpenSandboxScheduler:
|
||||
return sandbox_id
|
||||
|
||||
async def _monitor(
|
||||
self, request: RunnerRequest, sandbox_id: str, nonce: str
|
||||
self,
|
||||
request: RunnerRequest,
|
||||
sandbox_id: str,
|
||||
nonce: str,
|
||||
runner_name: str,
|
||||
) -> None:
|
||||
try:
|
||||
deadline = asyncio.get_running_loop().time() + SANDBOX_TIMEOUT
|
||||
@@ -139,13 +145,38 @@ class OpenSandboxScheduler:
|
||||
await asyncio.sleep(2)
|
||||
raise asyncio.TimeoutError(f"OpenSandbox {sandbox_id} timed out")
|
||||
finally:
|
||||
await self.tokens.revoke(nonce)
|
||||
try:
|
||||
await self.client.delete(sandbox_id)
|
||||
finally:
|
||||
self.active.pop(request.job_id, None)
|
||||
if self.on_finished is not None:
|
||||
self.on_finished(request.job_id)
|
||||
await self._cleanup(request.job_id, sandbox_id, nonce, runner_name)
|
||||
|
||||
async def complete(self, runner_name: str) -> bool:
|
||||
"""Delete the sandbox which actually ran a completed Gitea job."""
|
||||
state = self.runners.get(runner_name)
|
||||
if state is None:
|
||||
return False
|
||||
job_id, sandbox_id, nonce = state
|
||||
task = self.active.get(job_id)
|
||||
if task is not None:
|
||||
task.cancel()
|
||||
await asyncio.gather(task, return_exceptions=True)
|
||||
await self._cleanup(job_id, sandbox_id, nonce, runner_name)
|
||||
return True
|
||||
|
||||
async def _cleanup(
|
||||
self,
|
||||
job_id: int,
|
||||
sandbox_id: str,
|
||||
nonce: str,
|
||||
runner_name: str,
|
||||
) -> None:
|
||||
"""Revoke and delete once, including cancellation-before-start races."""
|
||||
if self.runners.pop(runner_name, None) is None:
|
||||
return
|
||||
await self.tokens.revoke(nonce)
|
||||
try:
|
||||
await self.client.delete(sandbox_id)
|
||||
finally:
|
||||
self.active.pop(job_id, None)
|
||||
if self.on_finished is not None:
|
||||
self.on_finished(job_id)
|
||||
|
||||
@staticmethod
|
||||
def _report(task: asyncio.Task[None]) -> None:
|
||||
|
||||
Reference in New Issue
Block a user