Task 007 add agent job admin runner follow-up
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import time
|
||||
from pathlib import Path
|
||||
@@ -90,6 +91,34 @@ def run_fake_runner_once(
|
||||
)
|
||||
|
||||
|
||||
def run_fake_runner_loop(
|
||||
*,
|
||||
backend_url: str,
|
||||
workspace_root: Path,
|
||||
poll_interval_seconds: float,
|
||||
) -> None:
|
||||
backend_client = HttpBackendJobClient(backend_url)
|
||||
while True:
|
||||
try:
|
||||
run_fake_runner_once(backend_client, workspace_root=workspace_root)
|
||||
except Exception:
|
||||
pass
|
||||
time.sleep(poll_interval_seconds)
|
||||
|
||||
|
||||
def get_backend_url() -> str:
|
||||
backend_url = os.environ.get("BACKEND_URL")
|
||||
if not backend_url:
|
||||
raise RuntimeError("BACKEND_URL is required for fake runner")
|
||||
return backend_url
|
||||
|
||||
|
||||
def get_workspace_root() -> Path:
|
||||
return Path(
|
||||
os.environ.get("RUNNER_WORKSPACE_ROOT", "/tmp/pipeline-runner-workspaces")
|
||||
)
|
||||
|
||||
|
||||
def materialize_fake_workspace(
|
||||
job: dict[str, Any],
|
||||
*,
|
||||
|
||||
@@ -6,7 +6,13 @@ import threading
|
||||
from http import HTTPStatus
|
||||
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||
|
||||
from src.application.fake_runner import get_backend_url, get_workspace_root, run_fake_runner_loop
|
||||
from src.application.fake_runner import (
|
||||
HttpBackendJobClient,
|
||||
get_backend_url,
|
||||
get_workspace_root,
|
||||
run_fake_runner_loop,
|
||||
run_fake_runner_once,
|
||||
)
|
||||
|
||||
|
||||
class HealthHandler(BaseHTTPRequestHandler):
|
||||
@@ -22,6 +28,32 @@ class HealthHandler(BaseHTTPRequestHandler):
|
||||
self.end_headers()
|
||||
self.wfile.write(body)
|
||||
|
||||
def do_POST(self) -> None:
|
||||
if self.path != "/run-once":
|
||||
self.send_error(HTTPStatus.NOT_FOUND)
|
||||
return
|
||||
|
||||
try:
|
||||
result = run_fake_runner_once(
|
||||
HttpBackendJobClient(get_backend_url()),
|
||||
workspace_root=get_workspace_root(),
|
||||
)
|
||||
except Exception as error:
|
||||
body = json.dumps({"status": "error", "detail": str(error)}).encode("utf-8")
|
||||
self.send_response(HTTPStatus.INTERNAL_SERVER_ERROR)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(body)))
|
||||
self.end_headers()
|
||||
self.wfile.write(body)
|
||||
return
|
||||
|
||||
body = json.dumps({"status": "ok", "job": result}, default=str).encode("utf-8")
|
||||
self.send_response(HTTPStatus.OK)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.send_header("Content-Length", str(len(body)))
|
||||
self.end_headers()
|
||||
self.wfile.write(body)
|
||||
|
||||
def log_message(self, format: str, *args: object) -> None:
|
||||
return
|
||||
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import tempfile
|
||||
import sys
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
|
||||
REPO_ROOT = Path(__file__).resolve().parents[3]
|
||||
RUNNER_ROOT = REPO_ROOT / "apps" / "runner"
|
||||
sys.path.insert(0, str(RUNNER_ROOT))
|
||||
|
||||
from src.application.fake_runner import run_fake_runner_once
|
||||
|
||||
|
||||
class FakeBackendJobClient:
|
||||
def __init__(self, job: dict[str, Any] | None) -> None:
|
||||
self.job = job
|
||||
self.completed: dict[str, Any] | None = None
|
||||
|
||||
def claim_job(self) -> dict[str, Any] | None:
|
||||
job = self.job
|
||||
self.job = None
|
||||
return job
|
||||
|
||||
def complete_job(self, job_id: str, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
self.completed = {"id": job_id, **payload}
|
||||
return self.completed
|
||||
|
||||
|
||||
class FakeRunnerTest(unittest.TestCase):
|
||||
def test_fake_runner_claims_job_writes_workspace_and_completes(self) -> None:
|
||||
backend = FakeBackendJobClient(
|
||||
{
|
||||
"id": "11111111-1111-1111-1111-111111111111",
|
||||
"agent_profile": "fake-codex",
|
||||
}
|
||||
)
|
||||
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
result = run_fake_runner_once(
|
||||
backend,
|
||||
workspace_root=Path(temp_dir),
|
||||
)
|
||||
self.assertIsNotNone(result)
|
||||
assert result is not None
|
||||
workspace_path = Path(result["workspace_path"])
|
||||
self.assertTrue((workspace_path / "inputs" / "job.json").is_file())
|
||||
self.assertTrue((workspace_path / "logs" / "stdout.log").is_file())
|
||||
self.assertTrue((workspace_path / "logs" / "stderr.log").is_file())
|
||||
self.assertTrue((workspace_path / "outputs" / "result.json").is_file())
|
||||
self.assertEqual(0, result["exit_code"])
|
||||
self.assertEqual("SUCCEEDED", result["output"]["status"])
|
||||
self.assertEqual([{"path": "outputs/result.json"}], result["output"]["output_files"])
|
||||
|
||||
def test_fake_runner_can_emit_invalid_schema_fixture(self) -> None:
|
||||
backend = FakeBackendJobClient(
|
||||
{
|
||||
"id": "22222222-2222-2222-2222-222222222222",
|
||||
"agent_profile": "fake-invalid-schema",
|
||||
}
|
||||
)
|
||||
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
result = run_fake_runner_once(
|
||||
backend,
|
||||
workspace_root=Path(temp_dir),
|
||||
)
|
||||
|
||||
self.assertIsNotNone(result)
|
||||
assert result is not None
|
||||
self.assertEqual({"status": "NOT_A_STATUS"}, result["output"])
|
||||
|
||||
def test_fake_runner_noops_without_queued_job(self) -> None:
|
||||
backend = FakeBackendJobClient(None)
|
||||
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
result = run_fake_runner_once(
|
||||
backend,
|
||||
workspace_root=Path(temp_dir),
|
||||
)
|
||||
|
||||
self.assertIsNone(result)
|
||||
self.assertIsNone(backend.completed)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user