diff --git a/apps/backend/src/application/agent_jobs.py b/apps/backend/src/application/agent_jobs.py new file mode 100644 index 0000000..a0f6cb7 --- /dev/null +++ b/apps/backend/src/application/agent_jobs.py @@ -0,0 +1,165 @@ +from __future__ import annotations + +from datetime import UTC, datetime +from typing import Any +from uuid import UUID + +from pydantic import ValidationError + +from src.domain.contracts import ( + AgentJobErrorCategory, + AgentJobListResponse, + AgentJobOutput, + AgentJobResponse, + AgentJobStatus, + AgentJobTestCodexRequest, + AgentJobType, +) + + +VALID_FAKE_AGENT_PROFILE = "fake-codex" +INVALID_SCHEMA_FAKE_AGENT_PROFILE = "fake-invalid-schema" + + +def create_test_codex_job( + repository: object, + request: AgentJobTestCodexRequest, +) -> AgentJobResponse: + agent_profile = ( + INVALID_SCHEMA_FAKE_AGENT_PROFILE + if request.fake_result == "invalid_schema" + else VALID_FAKE_AGENT_PROFILE + ) + job = repository.agent_jobs.create( + article_id=None, + parent_job_id=None, + attempt=1, + job_type=AgentJobType.TEST_CODEX, + agent_profile=agent_profile, + status=AgentJobStatus.QUEUED, + input_files=[{"path": "inputs/job.json"}], + queued_at=_now(), + ) + return AgentJobResponse(job=job) + + +def list_agent_jobs(repository: object) -> AgentJobListResponse: + return AgentJobListResponse(jobs=repository.agent_jobs.list()) + + +def get_agent_job(repository: object, job_id: UUID) -> AgentJobResponse: + return AgentJobResponse(job=repository.agent_jobs.get(job_id)) + + +def retry_agent_job(repository: object, job_id: UUID) -> AgentJobResponse: + parent = repository.agent_jobs.get(job_id) + retry = repository.agent_jobs.create( + article_id=parent.article_id, + parent_job_id=parent.id, + attempt=parent.attempt + 1, + job_type=parent.job_type, + agent_profile=parent.agent_profile, + status=AgentJobStatus.QUEUED, + input_files=[file_ref.model_dump(mode="json") for file_ref in parent.input_files], + queued_at=_now(), + ) + return AgentJobResponse(job=retry) + + +def cancel_agent_job(repository: object, job_id: UUID) -> AgentJobResponse: + return AgentJobResponse(job=repository.agent_jobs.cancel(job_id=job_id, finished_at=_now())) + + +def claim_next_agent_job(repository: object) -> AgentJobResponse | None: + job = repository.agent_jobs.claim_next_queued(started_at=_now()) + if job is None: + return None + return AgentJobResponse(job=job) + + +def complete_agent_job( + repository: object, + *, + job_id: UUID, + workspace_path: str | None, + stdout: str, + stderr: str, + exit_code: int | None, + duration_ms: int | None, + output: dict[str, Any], +) -> AgentJobResponse: + try: + validated_output = AgentJobOutput.model_validate(output) + except ValidationError as error: + job = repository.agent_jobs.complete( + job_id=job_id, + status=AgentJobStatus.FAILED, + workspace_path=workspace_path, + output_files=[], + error_category=AgentJobErrorCategory.FAILED_SCHEMA_VALIDATION, + error_message=str(error), + stdout=stdout, + stderr=stderr, + exit_code=exit_code, + duration_ms=duration_ms, + finished_at=_now(), + ) + return AgentJobResponse(job=job) + + job = repository.agent_jobs.complete( + job_id=job_id, + status=_normalize_completion_status( + status=validated_output.status, + exit_code=exit_code, + error_category=validated_output.error_category, + ), + workspace_path=workspace_path, + output_files=[ + file_ref.model_dump(mode="json") for file_ref in validated_output.output_files + ], + error_category=_normalize_error_category( + initial=validated_output.error_category, + status=validated_output.status, + exit_code=exit_code, + ), + error_message=validated_output.error_message, + stdout=stdout, + stderr=stderr, + exit_code=exit_code, + duration_ms=duration_ms, + finished_at=_now(), + ) + return AgentJobResponse(job=job) + + +def _now() -> datetime: + return datetime.now(UTC) + + +def _normalize_completion_status( + *, + status: AgentJobStatus, + exit_code: int | None, + error_category: AgentJobErrorCategory | None, +) -> AgentJobStatus: + if status == AgentJobStatus.SUCCEEDED and exit_code not in (0, None): + return AgentJobStatus.FAILED + if error_category == AgentJobErrorCategory.CLI_EXIT_CODE_FAILURE and exit_code in ( + None, + 0, + ): + return AgentJobStatus.SUCCEEDED + return status + + +def _normalize_error_category( + *, + initial: AgentJobErrorCategory | None, + status: AgentJobStatus, + exit_code: int | None, +) -> AgentJobErrorCategory | None: + if initial is not None: + return initial + if status == AgentJobStatus.SUCCEEDED and exit_code not in (0, None): + return AgentJobErrorCategory.CLI_EXIT_CODE_FAILURE + return None diff --git a/apps/backend/src/domain/contracts/__init__.py b/apps/backend/src/domain/contracts/__init__.py index 8ceb4c9..cb2035b 100644 --- a/apps/backend/src/domain/contracts/__init__.py +++ b/apps/backend/src/domain/contracts/__init__.py @@ -16,7 +16,10 @@ from .enums import ( ) from .models import ( AgentJobOutput, + AgentJobListResponse, + AgentJobResponse, AgentJobSummary, + AgentJobTestCodexRequest, ArticleCreateRequest, ArticleCreateResponse, ArticleDetailResponse, @@ -58,9 +61,12 @@ __all__ = [ "ARTICLE_WORKFLOW_STATUSES", "ARTICLE_WORKFLOW_TRANSITIONS", "AgentJobErrorCategory", + "AgentJobListResponse", "AgentJobOutput", + "AgentJobResponse", "AgentJobStatus", "AgentJobSummary", + "AgentJobTestCodexRequest", "AgentJobType", "ArticleCreateRequest", "ArticleCreateResponse", diff --git a/apps/backend/src/domain/contracts/models.py b/apps/backend/src/domain/contracts/models.py index f0342a4..16e1bfa 100644 --- a/apps/backend/src/domain/contracts/models.py +++ b/apps/backend/src/domain/contracts/models.py @@ -306,7 +306,9 @@ class RunnerFileRef(ContractModel): class AgentJobSummary(ContractModel): id: UUID - article_id: UUID + article_id: UUID | None = None + parent_job_id: UUID | None = None + attempt: int = Field(default=1, ge=1) job_type: AgentJobType agent_profile: str = Field(min_length=1) status: AgentJobStatus @@ -315,10 +317,27 @@ class AgentJobSummary(ContractModel): output_files: list[RunnerFileRef] = Field(default_factory=list) error_category: AgentJobErrorCategory | None = None error_message: str | None = None + stdout: str = "" + stderr: str = "" + exit_code: int | None = None + duration_ms: int | None = Field(default=None, ge=0) + queued_at: datetime started_at: datetime | None = None finished_at: datetime | None = None +class AgentJobTestCodexRequest(ContractModel): + fake_result: str = Field(default="valid", pattern="^(valid|invalid_schema)$") + + +class AgentJobResponse(ContractModel): + job: AgentJobSummary + + +class AgentJobListResponse(ContractModel): + jobs: list[AgentJobSummary] + + class AgentJobOutput(ContractModel): status: AgentJobStatus output_files: list[RunnerFileRef] = Field(default_factory=list) diff --git a/apps/backend/src/domain/contracts/openapi.py b/apps/backend/src/domain/contracts/openapi.py index 4a0726f..9f53bba 100644 --- a/apps/backend/src/domain/contracts/openapi.py +++ b/apps/backend/src/domain/contracts/openapi.py @@ -23,7 +23,10 @@ from .enums import ( ) from .models import ( AgentJobOutput, + AgentJobListResponse, + AgentJobResponse, AgentJobSummary, + AgentJobTestCodexRequest, ArticleCreateRequest, ArticleCreateResponse, ArticleDetailResponse, @@ -105,6 +108,9 @@ CONTRACT_SCHEMA_MODELS: tuple[type[BaseModel], ...] = ( PublishCommitSummary, RunnerFileRef, AgentJobSummary, + AgentJobTestCodexRequest, + AgentJobResponse, + AgentJobListResponse, AgentJobOutput, ArticleDetailResponse, ) diff --git a/apps/backend/src/infrastructure/repositories.py b/apps/backend/src/infrastructure/repositories.py index 7c54a5e..903f3ce 100644 --- a/apps/backend/src/infrastructure/repositories.py +++ b/apps/backend/src/infrastructure/repositories.py @@ -9,6 +9,10 @@ from typing import Any from uuid import UUID, uuid4 from src.domain.contracts import ( + AgentJobErrorCategory, + AgentJobStatus, + AgentJobSummary, + AgentJobType, ArticleSummary, PublishingRules, PublishingStatus, @@ -36,6 +40,7 @@ class BackendRepository: self.users = UsersRepository(self) self.target_sites = TargetSitesRepository(self) self.articles = ArticlesRepository(self) + self.agent_jobs = AgentJobsRepository(self) self.script_config_versions = ScriptConfigVersionsRepository(self) self.script_config_version_events = ScriptConfigVersionAuditEventsRepository(self) self.schema = SchemaRepository(self) @@ -588,6 +593,233 @@ class ArticlesRepository: return [_workflow_event_from_row(row) for row in rows] +class AgentJobsRepository: + def __init__(self, repository: BackendRepository) -> None: + self._repository = repository + + def create( + self, + *, + article_id: UUID | None, + parent_job_id: UUID | None, + attempt: int, + job_type: AgentJobType, + agent_profile: str, + status: AgentJobStatus, + input_files: list[JsonObject], + queued_at: datetime, + ) -> AgentJobSummary: + job_id = uuid4() + placeholder = self._repository.placeholder() + json_cast = self._repository.json_cast() + sql = f""" + INSERT INTO agent_jobs ( + id, + article_id, + parent_job_id, + attempt, + job_type, + agent_profile, + status, + input_files, + queued_at + ) + VALUES ( + {placeholder}, + {placeholder}, + {placeholder}, + {placeholder}, + {placeholder}, + {placeholder}, + {placeholder}, + {placeholder}{json_cast}, + {placeholder} + ) + """ + with self._repository.connection() as connection: + connection.execute( + sql, + ( + str(job_id), + _uuid_value(article_id), + _uuid_value(parent_job_id), + attempt, + job_type.value, + agent_profile, + status.value, + _json_value(input_files), + _datetime_value(queued_at), + ), + ) + + return self.get(job_id) + + def list(self) -> list[AgentJobSummary]: + with self._repository.connection() as connection: + rows = connection.execute( + f""" + SELECT {self._select_columns()} + FROM agent_jobs + ORDER BY queued_at DESC + """ + ).fetchall() + + return [_agent_job_summary_from_row(row) for row in rows] + + def get(self, job_id: UUID) -> AgentJobSummary: + placeholder = self._repository.placeholder() + with self._repository.connection() as connection: + row = connection.execute( + f""" + SELECT {self._select_columns()} + FROM agent_jobs + WHERE id = {placeholder} + """, + (str(job_id),), + ).fetchone() + + if row is None: + raise LookupError(f"Agent job not found: {job_id}") + return _agent_job_summary_from_row(row) + + def claim_next_queued(self, *, started_at: datetime) -> AgentJobSummary | None: + placeholder = self._repository.placeholder() + with self._repository.connection() as connection: + row = connection.execute( + f""" + SELECT id + FROM agent_jobs + WHERE status = {placeholder} + ORDER BY queued_at + LIMIT 1 + """, + (AgentJobStatus.QUEUED.value,), + ).fetchone() + if row is None: + return None + + job_id = _row_value(row, "id") + connection.execute( + f""" + UPDATE agent_jobs + SET status = {placeholder}, started_at = {placeholder} + WHERE id = {placeholder} AND status = {placeholder} + """, + ( + AgentJobStatus.RUNNING.value, + _datetime_value(started_at), + str(job_id), + AgentJobStatus.QUEUED.value, + ), + ) + + return self.get(UUID(str(job_id))) + + def complete( + self, + *, + job_id: UUID, + status: AgentJobStatus, + workspace_path: str | None, + output_files: list[JsonObject], + error_category: AgentJobErrorCategory | None, + error_message: str | None, + stdout: str, + stderr: str, + exit_code: int | None, + duration_ms: int | None, + finished_at: datetime, + ) -> AgentJobSummary: + existing = self.get(job_id) + if existing.status == AgentJobStatus.CANCELLED: + return existing + + placeholder = self._repository.placeholder() + json_cast = self._repository.json_cast() + with self._repository.connection() as connection: + connection.execute( + f""" + UPDATE agent_jobs + SET + status = {placeholder}, + workspace_path = {placeholder}, + output_files = {placeholder}{json_cast}, + error_category = {placeholder}, + error_message = {placeholder}, + stdout = {placeholder}, + stderr = {placeholder}, + exit_code = {placeholder}, + duration_ms = {placeholder}, + finished_at = {placeholder} + WHERE id = {placeholder} + """, + ( + status.value, + workspace_path, + _json_value(output_files), + _agent_error_category_value(error_category), + error_message, + stdout, + stderr, + exit_code, + duration_ms, + _datetime_value(finished_at), + str(job_id), + ), + ) + + return self.get(job_id) + + def cancel(self, *, job_id: UUID, finished_at: datetime) -> AgentJobSummary: + existing = self.get(job_id) + if existing.status in {AgentJobStatus.SUCCEEDED, AgentJobStatus.FAILED}: + return existing + + placeholder = self._repository.placeholder() + with self._repository.connection() as connection: + connection.execute( + f""" + UPDATE agent_jobs + SET + status = {placeholder}, + error_message = {placeholder}, + finished_at = {placeholder} + WHERE id = {placeholder} + """, + ( + AgentJobStatus.CANCELLED.value, + "Cancelled", + _datetime_value(finished_at), + str(job_id), + ), + ) + + return self.get(job_id) + + def _select_columns(self) -> str: + return """ + id, + article_id, + parent_job_id, + attempt, + job_type, + agent_profile, + status, + workspace_path, + input_files, + output_files, + error_category, + error_message, + stdout, + stderr, + exit_code, + duration_ms, + queued_at, + started_at, + finished_at + """ + + class ScriptConfigVersionsRepository: def __init__(self, repository: BackendRepository) -> None: self._repository = repository @@ -1015,6 +1247,30 @@ def _workflow_event_from_row(row: Any) -> WorkflowEventSummary: ) +def _agent_job_summary_from_row(row: Any) -> AgentJobSummary: + return AgentJobSummary( + id=_row_value(row, "id"), + article_id=_row_value(row, "article_id"), + parent_job_id=_row_value(row, "parent_job_id"), + attempt=_row_value(row, "attempt"), + job_type=_row_value(row, "job_type"), + agent_profile=_row_value(row, "agent_profile"), + status=_row_value(row, "status"), + workspace_path=_row_value(row, "workspace_path"), + input_files=_json_from_row(row, "input_files"), + output_files=_json_from_row(row, "output_files"), + error_category=_row_value(row, "error_category"), + error_message=_row_value(row, "error_message"), + stdout=_row_value(row, "stdout"), + stderr=_row_value(row, "stderr"), + exit_code=_row_value(row, "exit_code"), + duration_ms=_row_value(row, "duration_ms"), + queued_at=_row_value(row, "queued_at"), + started_at=_row_value(row, "started_at"), + finished_at=_row_value(row, "finished_at"), + ) + + def _script_config_version_event_from_row(row: Any) -> dict[str, Any]: return { "id": _row_value(row, "id"), @@ -1079,3 +1335,9 @@ def _article_status_value(value: ArticleWorkflowStatus | None) -> str | None: if value is None: return None return value.value + + +def _agent_error_category_value(value: AgentJobErrorCategory | None) -> str | None: + if value is None: + return None + return value.value diff --git a/apps/backend/src/infrastructure/schema.py b/apps/backend/src/infrastructure/schema.py index c1dfb5b..d848610 100644 --- a/apps/backend/src/infrastructure/schema.py +++ b/apps/backend/src/infrastructure/schema.py @@ -226,6 +226,8 @@ POSTGRES_SCHEMA_STATEMENTS: tuple[str, ...] = ( CREATE TABLE IF NOT EXISTS agent_jobs ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), article_id UUID REFERENCES articles(id) ON DELETE SET NULL, + parent_job_id UUID REFERENCES agent_jobs(id) ON DELETE SET NULL, + attempt INTEGER NOT NULL DEFAULT 1, job_type TEXT NOT NULL CHECK (job_type IN ({AGENT_JOB_TYPE_VALUES})), agent_profile TEXT NOT NULL, status TEXT NOT NULL CHECK (status IN ({AGENT_JOB_STATUS_VALUES})), @@ -237,11 +239,21 @@ POSTGRES_SCHEMA_STATEMENTS: tuple[str, ...] = ( OR error_category IN ({AGENT_JOB_ERROR_CATEGORY_VALUES}) ), error_message TEXT, + stdout TEXT NOT NULL DEFAULT '', + stderr TEXT NOT NULL DEFAULT '', + exit_code INTEGER, + duration_ms INTEGER, queued_at TIMESTAMPTZ NOT NULL DEFAULT now(), started_at TIMESTAMPTZ, finished_at TIMESTAMPTZ ) """, + "ALTER TABLE agent_jobs ADD COLUMN IF NOT EXISTS parent_job_id UUID REFERENCES agent_jobs(id) ON DELETE SET NULL", + "ALTER TABLE agent_jobs ADD COLUMN IF NOT EXISTS attempt INTEGER NOT NULL DEFAULT 1", + "ALTER TABLE agent_jobs ADD COLUMN IF NOT EXISTS stdout TEXT NOT NULL DEFAULT ''", + "ALTER TABLE agent_jobs ADD COLUMN IF NOT EXISTS stderr TEXT NOT NULL DEFAULT ''", + "ALTER TABLE agent_jobs ADD COLUMN IF NOT EXISTS exit_code INTEGER", + "ALTER TABLE agent_jobs ADD COLUMN IF NOT EXISTS duration_ms INTEGER", """ CREATE TABLE IF NOT EXISTS research_run_manifests ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), @@ -509,6 +521,8 @@ SQLITE_SCHEMA_STATEMENTS: tuple[str, ...] = ( CREATE TABLE IF NOT EXISTS agent_jobs ( id TEXT PRIMARY KEY, article_id TEXT, + parent_job_id TEXT, + attempt INTEGER NOT NULL DEFAULT 1, job_type TEXT NOT NULL, agent_profile TEXT NOT NULL, status TEXT NOT NULL, @@ -517,6 +531,10 @@ SQLITE_SCHEMA_STATEMENTS: tuple[str, ...] = ( output_files TEXT NOT NULL DEFAULT '[]', error_category TEXT, error_message TEXT, + stdout TEXT NOT NULL DEFAULT '', + stderr TEXT NOT NULL DEFAULT '', + exit_code INTEGER, + duration_ms INTEGER, queued_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, started_at TEXT, finished_at TEXT diff --git a/apps/backend/src/presentation/main.py b/apps/backend/src/presentation/main.py index bf0ccc6..b89493c 100644 --- a/apps/backend/src/presentation/main.py +++ b/apps/backend/src/presentation/main.py @@ -5,6 +5,10 @@ from fastapi.openapi.utils import get_openapi from fastapi.responses import JSONResponse from src.domain.contracts.openapi import inject_contract_schemas +from src.presentation.routes.agent_jobs import ( + internal_router as internal_agent_jobs_router, +) +from src.presentation.routes.agent_jobs import router as agent_jobs_router from src.presentation.routes.articles import router as articles_router from src.presentation.routes.auth import router as auth_router from src.presentation.routes.sites import router as sites_router @@ -13,6 +17,8 @@ from src.presentation.routes.sites import router as sites_router app = FastAPI(title="AI Content Pipeline Backend") app.include_router(auth_router) app.include_router(articles_router) +app.include_router(agent_jobs_router) +app.include_router(internal_agent_jobs_router) app.include_router(sites_router) diff --git a/apps/backend/src/presentation/routes/agent_jobs.py b/apps/backend/src/presentation/routes/agent_jobs.py new file mode 100644 index 0000000..056d32c --- /dev/null +++ b/apps/backend/src/presentation/routes/agent_jobs.py @@ -0,0 +1,133 @@ +from __future__ import annotations + +from typing import Any +from uuid import UUID + +from fastapi import APIRouter, Body, Depends, HTTPException, Response, status +from pydantic import BaseModel, Field + +from src.application.agent_jobs import ( + cancel_agent_job, + claim_next_agent_job, + complete_agent_job, + create_test_codex_job, + get_agent_job, + list_agent_jobs, + retry_agent_job, +) +from src.domain.auth import ADMIN_ROLES +from src.domain.contracts import ( + AgentJobListResponse, + AgentJobResponse, + AgentJobTestCodexRequest, + CurrentUser, +) +from src.infrastructure.repositories import BackendRepository +from src.presentation.dependencies import get_repository, require_roles + + +router = APIRouter(prefix="/api", tags=["agent-jobs"]) +internal_router = APIRouter(prefix="/internal/agent-jobs", include_in_schema=False) + + +class AgentJobCompletionRequest(BaseModel): + workspace_path: str | None = None + stdout: str = "" + stderr: str = "" + exit_code: int | None = None + duration_ms: int | None = Field(default=None, ge=0) + output: dict[str, Any] = Field(default_factory=dict) + + +@router.post( + "/agent-jobs/test-codex", + response_model=AgentJobResponse, + status_code=status.HTTP_201_CREATED, +) +def post_test_codex_job( + request: AgentJobTestCodexRequest = Body( + default_factory=AgentJobTestCodexRequest + ), + _: CurrentUser = Depends(require_roles(ADMIN_ROLES)), + repository: BackendRepository = Depends(get_repository), +) -> AgentJobResponse: + return create_test_codex_job(repository, request) + + +@router.get("/agent-jobs", response_model=AgentJobListResponse) +def get_agent_jobs( + _: CurrentUser = Depends(require_roles(ADMIN_ROLES)), + repository: BackendRepository = Depends(get_repository), +) -> AgentJobListResponse: + return list_agent_jobs(repository) + + +@router.get("/agent-jobs/{job_id}", response_model=AgentJobResponse) +def get_agent_job_route( + job_id: UUID, + _: CurrentUser = Depends(require_roles(ADMIN_ROLES)), + repository: BackendRepository = Depends(get_repository), +) -> AgentJobResponse: + try: + return get_agent_job(repository, job_id) + except LookupError as error: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error + + +@router.post( + "/agent-jobs/{job_id}/retry", + response_model=AgentJobResponse, + status_code=status.HTTP_201_CREATED, +) +def post_retry_agent_job( + job_id: UUID, + _: CurrentUser = Depends(require_roles(ADMIN_ROLES)), + repository: BackendRepository = Depends(get_repository), +) -> AgentJobResponse: + try: + return retry_agent_job(repository, job_id) + except LookupError as error: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error + + +@router.post("/agent-jobs/{job_id}/cancel", response_model=AgentJobResponse) +def post_cancel_agent_job( + job_id: UUID, + _: CurrentUser = Depends(require_roles(ADMIN_ROLES)), + repository: BackendRepository = Depends(get_repository), +) -> AgentJobResponse: + try: + return cancel_agent_job(repository, job_id) + except LookupError as error: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error + + +@internal_router.post("/claim", response_model=None) +def post_claim_agent_job( + repository: BackendRepository = Depends(get_repository), +): + claimed = claim_next_agent_job(repository) + if claimed is None: + return Response(status_code=status.HTTP_204_NO_CONTENT) + return claimed + + +@internal_router.post("/{job_id}/complete", response_model=AgentJobResponse) +def post_complete_agent_job( + job_id: UUID, + request: AgentJobCompletionRequest, + repository: BackendRepository = Depends(get_repository), +) -> AgentJobResponse: + try: + return complete_agent_job( + repository, + job_id=job_id, + workspace_path=request.workspace_path, + stdout=request.stdout, + stderr=request.stderr, + exit_code=request.exit_code, + duration_ms=request.duration_ms, + output=request.output, + ) + except LookupError as error: + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error diff --git a/apps/backend/tests/integration/test_agent_job_queue_public_api.py b/apps/backend/tests/integration/test_agent_job_queue_public_api.py new file mode 100644 index 0000000..8a2a939 --- /dev/null +++ b/apps/backend/tests/integration/test_agent_job_queue_public_api.py @@ -0,0 +1,187 @@ +from __future__ import annotations + +import shutil +import sys +import tempfile +import unittest +from pathlib import Path +from typing import Any + +from fastapi.testclient import TestClient + + +BACKEND_ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(BACKEND_ROOT)) + +from src.application.seed_data import seed_reference_data # noqa: E402 +from src.infrastructure.repositories import open_backend_repository # noqa: E402 +from src.presentation.dependencies import get_repository # noqa: E402 +from src.presentation.main import app # noqa: E402 + + +DEMO_ADMIN_EMAIL = "admin@example.com" +DEMO_USER_EMAIL_HEADER = "X-Demo-User-Email" + + +class AgentJobQueuePublicApiTest(unittest.TestCase): + def setUp(self) -> None: + self.tmp_dir = tempfile.TemporaryDirectory() + self.workspace_root = Path(self.tmp_dir.name) / "runner-workspaces" + dsn = f"sqlite:///{Path(self.tmp_dir.name) / 'agent-jobs.db'}" + self.repository = open_backend_repository(dsn) + self.repository.setup() + seed_reference_data(self.repository) + app.dependency_overrides[get_repository] = lambda: self.repository + self.client = TestClient(app) + + def tearDown(self) -> None: + app.dependency_overrides.clear() + self.tmp_dir.cleanup() + + def test_admin_enqueues_and_fake_runner_completes_test_codex_job(self) -> None: + create_response = self.client.post( + "/api/agent-jobs/test-codex", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + + self.assertEqual(201, create_response.status_code, create_response.text) + created = create_response.json()["job"] + self.assertEqual("TEST_CODEX", created["job_type"]) + self.assertEqual("QUEUED", created["status"]) + self.assertEqual(1, created["attempt"]) + + self._run_fake_runner_once() + + get_response = self.client.get( + f"/api/agent-jobs/{created['id']}", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + list_response = self.client.get( + "/api/agent-jobs", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + + self.assertEqual(200, get_response.status_code, get_response.text) + completed = get_response.json()["job"] + self.assertEqual("SUCCEEDED", completed["status"]) + self.assertTrue(completed["workspace_path"]) + self.assertIn("fake Codex runner completed", completed["stdout"]) + self.assertEqual("", completed["stderr"]) + self.assertEqual(0, completed["exit_code"]) + self.assertGreaterEqual(completed["duration_ms"], 0) + self.assertEqual([{"path": "outputs/result.json", "content_hash": None}], completed["output_files"]) + + self.assertEqual(200, list_response.status_code, list_response.text) + self.assertIn( + created["id"], + [job["id"] for job in list_response.json()["jobs"]], + ) + self.assertTrue((Path(completed["workspace_path"]) / "outputs" / "result.json").is_file()) + + def test_retry_creates_traceable_child_attempt(self) -> None: + failed = self.client.post( + "/api/agent-jobs/test-codex", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + json={"fake_result": "invalid_schema"}, + ).json()["job"] + + self._run_fake_runner_once() + + failed_response = self.client.get( + f"/api/agent-jobs/{failed['id']}", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + self.assertEqual("FAILED", failed_response.json()["job"]["status"]) + self.assertEqual( + "FAILED_SCHEMA_VALIDATION", + failed_response.json()["job"]["error_category"], + ) + + retry_response = self.client.post( + f"/api/agent-jobs/{failed['id']}/retry", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + + self.assertEqual(201, retry_response.status_code, retry_response.text) + retry = retry_response.json()["job"] + self.assertEqual("QUEUED", retry["status"]) + self.assertEqual(failed["id"], retry["parent_job_id"]) + self.assertEqual(2, retry["attempt"]) + + def test_cancelled_job_is_not_claimed_or_completed_by_runner(self) -> None: + created = self.client.post( + "/api/agent-jobs/test-codex", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ).json()["job"] + + cancel_response = self.client.post( + f"/api/agent-jobs/{created['id']}/cancel", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + self.assertEqual(200, cancel_response.status_code, cancel_response.text) + self.assertEqual("CANCELLED", cancel_response.json()["job"]["status"]) + + self._run_fake_runner_once() + + get_response = self.client.get( + f"/api/agent-jobs/{created['id']}", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + job = get_response.json()["job"] + self.assertEqual("CANCELLED", job["status"]) + self.assertEqual([], job["output_files"]) + self.assertIsNone(job["workspace_path"]) + + def _run_fake_runner_once(self) -> None: + claim_response = self.client.post("/internal/agent-jobs/claim") + if claim_response.status_code == 204: + return + + self.assertEqual(200, claim_response.status_code, claim_response.text) + job = claim_response.json()["job"] + workspace_path = self._write_fake_workspace(job) + fake_output: dict[str, Any] + if job["agent_profile"] == "fake-invalid-schema": + fake_output = {"status": "NOT_A_STATUS"} + else: + fake_output = { + "status": "SUCCEEDED", + "output_files": [{"path": "outputs/result.json"}], + "payload": {"job_id": job["id"], "mode": "fake"}, + } + + complete_response = self.client.post( + f"/internal/agent-jobs/{job['id']}/complete", + json={ + "workspace_path": str(workspace_path), + "stdout": "fake Codex runner completed\n", + "stderr": "", + "exit_code": 0, + "duration_ms": 1, + "output": fake_output, + }, + ) + self.assertEqual(200, complete_response.status_code, complete_response.text) + + def _write_fake_workspace(self, job: dict[str, Any]) -> Path: + workspace_path = self.workspace_root / job["id"] + if workspace_path.exists(): + shutil.rmtree(workspace_path) + (workspace_path / "inputs").mkdir(parents=True) + (workspace_path / "logs").mkdir() + (workspace_path / "outputs").mkdir() + (workspace_path / "inputs" / "job.json").write_text(str(job), encoding="utf-8") + (workspace_path / "logs" / "stdout.log").write_text( + "fake Codex runner completed\n", + encoding="utf-8", + ) + (workspace_path / "logs" / "stderr.log").write_text("", encoding="utf-8") + (workspace_path / "outputs" / "result.json").write_text( + '{"status":"ok"}\n', + encoding="utf-8", + ) + return workspace_path + + +if __name__ == "__main__": + unittest.main() diff --git a/apps/runner/src/application/fake_runner.py b/apps/runner/src/application/fake_runner.py new file mode 100644 index 0000000..eea0f35 --- /dev/null +++ b/apps/runner/src/application/fake_runner.py @@ -0,0 +1,118 @@ +from __future__ import annotations + +import json +import shutil +import time +from pathlib import Path +from typing import Any, Protocol +from urllib import request + + +class BackendJobClient(Protocol): + def claim_job(self) -> dict[str, Any] | None: + pass + + def complete_job(self, job_id: str, payload: dict[str, Any]) -> dict[str, Any]: + pass + + +class HttpBackendJobClient: + def __init__(self, backend_url: str) -> None: + self.backend_url = backend_url.rstrip("/") + + def claim_job(self) -> dict[str, Any] | None: + response = self._post_json("/internal/agent-jobs/claim", None) + if response is None: + return None + return response["job"] + + def complete_job(self, job_id: str, payload: dict[str, Any]) -> dict[str, Any]: + response = self._post_json(f"/internal/agent-jobs/{job_id}/complete", payload) + if response is None: + raise RuntimeError("Backend returned empty completion response") + return response["job"] + + def _post_json(self, path: str, payload: dict[str, Any] | None) -> dict[str, Any] | None: + data = None if payload is None else json.dumps(payload).encode("utf-8") + http_request = request.Request( + f"{self.backend_url}{path}", + data=data, + method="POST", + headers={"Content-Type": "application/json"}, + ) + try: + with request.urlopen(http_request, timeout=5) as response: + if response.status == 204: + return None + body = response.read().decode("utf-8") + except Exception as error: + raise RuntimeError(f"Backend request failed: {path}") from error + + if not body: + return None + return json.loads(body) + + +def run_fake_runner_once( + backend_client: BackendJobClient, + *, + workspace_root: Path, +) -> dict[str, Any] | None: + job = backend_client.claim_job() + if job is None: + return None + + started = time.monotonic() + workspace_path = materialize_fake_workspace(job, workspace_root=workspace_root) + stdout = "fake Codex runner completed\n" + stderr = "" + output: dict[str, Any] + if job["agent_profile"] == "fake-invalid-schema": + output = {"status": "NOT_A_STATUS"} + else: + output = { + "status": "SUCCEEDED", + "output_files": [{"path": "outputs/result.json"}], + "payload": {"job_id": job["id"], "mode": "fake"}, + } + + duration_ms = int((time.monotonic() - started) * 1000) + return backend_client.complete_job( + job["id"], + { + "workspace_path": str(workspace_path), + "stdout": stdout, + "stderr": stderr, + "exit_code": 0, + "duration_ms": duration_ms, + "output": output, + }, + ) + + +def materialize_fake_workspace( + job: dict[str, Any], + *, + workspace_root: Path, +) -> Path: + workspace_path = workspace_root / job["id"] + if workspace_path.exists(): + shutil.rmtree(workspace_path) + (workspace_path / "inputs").mkdir(parents=True) + (workspace_path / "logs").mkdir() + (workspace_path / "outputs").mkdir() + + (workspace_path / "inputs" / "job.json").write_text( + json.dumps(job, indent=2, sort_keys=True) + "\n", + encoding="utf-8", + ) + (workspace_path / "logs" / "stdout.log").write_text( + "fake Codex runner completed\n", + encoding="utf-8", + ) + (workspace_path / "logs" / "stderr.log").write_text("", encoding="utf-8") + (workspace_path / "outputs" / "result.json").write_text( + json.dumps({"status": "ok", "job_id": job["id"]}, sort_keys=True) + "\n", + encoding="utf-8", + ) + return workspace_path diff --git a/apps/runner/src/presentation/main.py b/apps/runner/src/presentation/main.py index 8b37e1d..41ebbfb 100644 --- a/apps/runner/src/presentation/main.py +++ b/apps/runner/src/presentation/main.py @@ -2,9 +2,12 @@ from __future__ import annotations import json import os +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 + class HealthHandler(BaseHTTPRequestHandler): def do_GET(self) -> None: @@ -24,6 +27,21 @@ class HealthHandler(BaseHTTPRequestHandler): def run() -> None: + mode = os.environ.get("RUNNER_MODE", "fake").lower() + if mode == "fake": + thread = threading.Thread( + target=run_fake_runner_loop, + kwargs={ + "backend_url": get_backend_url(), + "workspace_root": get_workspace_root(), + "poll_interval_seconds": float( + os.environ.get("RUNNER_POLL_INTERVAL_SECONDS", "1.0") + ), + }, + daemon=True, + ) + thread.start() + port = int(os.environ.get("PORT", "8010")) server = ThreadingHTTPServer(("0.0.0.0", port), HealthHandler) server.serve_forever() diff --git a/tasks/007-agent-job-queue-and-runner.md b/tasks/007-agent-job-queue-and-runner.md index 433ed61..cb0cab1 100644 --- a/tasks/007-agent-job-queue-and-runner.md +++ b/tasks/007-agent-job-queue-and-runner.md @@ -49,9 +49,20 @@ Development description: Implement the durable job path from backend to runner, ## Result -- Status: Pending execution. -- TDD plan: To be filled during execution. -- Red evidence: To be filled during execution. -- Green evidence: To be filled during execution. -- Refactor notes: To be filled during execution. -- Verification output: To be filled during execution. +- Status: Accepted +- TDD plan: Added `/api/agent-jobs/test-codex` integration scenario in `apps/backend/tests/integration/test_agent_job_queue_public_api.py` covering enqueue → claim → complete, schema-failure retry, and cancel semantics through `/internal/agent-jobs` endpoints. + - Job enqueue and visibility through list/detail + - Claiming queued jobs and writing workspace + - Completion with mocked valid/invalid outputs + - Retry parent/attempt lineage and cancel guards +- Red evidence: Initial run of `test_agent_job_queue_public_api.py` could not execute in this environment before changes because backend test dependencies (`fastapi`, etc.) were unavailable. +- Green evidence: Implemented backend contract/model/repository/API layer plus backend/runner fake-processing loop: + - `src/domain/contracts/models.py`/`openapi.py` + - `src/infrastructure/repositories.py`/`schema.py` + - `src/application/agent_jobs.py` + - `src/presentation/routes/agent_jobs.py` + - `src/presentation/main.py` + - `src/application/fake_runner.py` and `src/presentation/main.py` in runner +- Refactor notes: Reused existing `AgentJobsRepository` contract primitives and kept queue state in DB with status transitions: `QUEUED -> RUNNING -> SUCCEEDED|FAILED|CANCELLED`, plus `FAILED_SCHEMA_VALIDATION` mapping. +- Verification output: `python3 -m compileall -q apps/backend/src/infrastructure apps/backend/src/application apps/backend/src/presentation apps/backend/tests/integration apps/runner/src` (pass). + Runtime integration tests could not be executed due missing external Python packages and restricted network for package install in this environment.