From 8ebb5ad623d8558e529edd31d22c68d0e7df1627 Mon Sep 17 00:00:00 2001 From: "E.Gavrilov" Date: Thu, 21 May 2026 23:10:34 +0300 Subject: [PATCH] feat(task-012): orchestrate parallel section production jobs --- apps/backend/src/application/agent_jobs.py | 60 +++- apps/backend/src/application/articles.py | 2 + .../src/application/boundary_questions.py | 1 + apps/backend/src/application/evidence.py | 115 ++++++- apps/backend/src/application/research.py | 1 + apps/backend/src/domain/contracts/models.py | 1 + .../src/infrastructure/repositories.py | 24 ++ apps/backend/src/infrastructure/schema.py | 3 + .../test_evidence_matrix_public_api.py | 2 +- .../test_parallel_production_public_api.py | 321 ++++++++++++++++++ .../src/features/article-detail/model.ts | 74 +++- .../src/features/article-detail/ui.tsx | 41 +++ .../tests/article_detail.model.test.mjs | 104 ++++++ packages/shared/src/api-types.ts | 1 + tasks/012-parallel-production-jobs.md | 45 ++- 15 files changed, 756 insertions(+), 39 deletions(-) create mode 100644 apps/backend/tests/integration/test_parallel_production_public_api.py create mode 100644 apps/frontend/tests/article_detail.model.test.mjs diff --git a/apps/backend/src/application/agent_jobs.py b/apps/backend/src/application/agent_jobs.py index a0f6cb7..581c74f 100644 --- a/apps/backend/src/application/agent_jobs.py +++ b/apps/backend/src/application/agent_jobs.py @@ -61,6 +61,7 @@ def retry_agent_job(repository: object, job_id: UUID) -> AgentJobResponse: agent_profile=parent.agent_profile, status=AgentJobStatus.QUEUED, input_files=[file_ref.model_dump(mode="json") for file_ref in parent.input_files], + payload=parent.payload, queued_at=_now(), ) return AgentJobResponse(job=retry) @@ -88,6 +89,7 @@ def complete_agent_job( duration_ms: int | None, output: dict[str, Any], ) -> AgentJobResponse: + existing_job = repository.agent_jobs.get(job_id) try: validated_output = AgentJobOutput.model_validate(output) except ValidationError as error: @@ -96,6 +98,7 @@ def complete_agent_job( status=AgentJobStatus.FAILED, workspace_path=workspace_path, output_files=[], + payload=existing_job.payload, error_category=AgentJobErrorCategory.FAILED_SCHEMA_VALIDATION, error_message=str(error), stdout=stdout, @@ -106,23 +109,40 @@ def complete_agent_job( ) return AgentJobResponse(job=job) + merged_payload = _merge_payload( + base_payload=existing_job.payload, + new_payload=validated_output.payload, + ) + status = _normalize_completion_status( + status=validated_output.status, + exit_code=exit_code, + error_category=validated_output.error_category, + ) + error_category = _normalize_error_category( + initial=validated_output.error_category, + status=validated_output.status, + exit_code=exit_code, + ) + error_message = validated_output.error_message + unsupported_claims = _extract_unsupported_claims(merged_payload) + if existing_job.job_type == AgentJobType.SECTION_SCAFFOLD and unsupported_claims: + status = AgentJobStatus.FAILED + error_category = AgentJobErrorCategory.UNSUPPORTED_CLAIMS_FOUND + if not error_message: + error_message = ( + f"Unsupported claims introduced during scaffolding: {len(unsupported_claims)}" + ) + 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, - ), + status=status, 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, + payload=merged_payload, + error_category=error_category, + error_message=error_message, stdout=stdout, stderr=stderr, exit_code=exit_code, @@ -163,3 +183,21 @@ def _normalize_error_category( if status == AgentJobStatus.SUCCEEDED and exit_code not in (0, None): return AgentJobErrorCategory.CLI_EXIT_CODE_FAILURE return None + + +def _merge_payload( + *, + base_payload: dict[str, Any], + new_payload: dict[str, Any], +) -> dict[str, Any]: + merged = dict(base_payload) + for key, value in new_payload.items(): + merged[key] = value + return merged + + +def _extract_unsupported_claims(payload: dict[str, Any]) -> list[Any]: + claims = payload.get("unsupported_claims") + if isinstance(claims, list): + return claims + return [] diff --git a/apps/backend/src/application/articles.py b/apps/backend/src/application/articles.py index 72c7d99..b99f3ea 100644 --- a/apps/backend/src/application/articles.py +++ b/apps/backend/src/application/articles.py @@ -63,6 +63,7 @@ def get_article_detail( research_manifests = repository.research_manifests.list_for_article(article_id) evidence = repository.evidence_items.list_for_article(article_id) claims = repository.claims.list_for_article(article_id) + agent_jobs = repository.agent_jobs.list_for_article(article_id) return ArticleDetailResponse( article=article, target_site=target_site, @@ -71,6 +72,7 @@ def get_article_detail( plan=plans[-1] if plans else None, evidence=evidence, claims=claims, + agent_jobs=agent_jobs, research_manifests=research_manifests, ) diff --git a/apps/backend/src/application/boundary_questions.py b/apps/backend/src/application/boundary_questions.py index 79a4e00..dbea73a 100644 --- a/apps/backend/src/application/boundary_questions.py +++ b/apps/backend/src/application/boundary_questions.py @@ -106,6 +106,7 @@ def generate_boundary_questions( status=AgentJobStatus.SUCCEEDED, workspace_path=None, output_files=[{"path": "outputs/boundary-questions.json"}], + payload={}, error_category=None, error_message=None, stdout="fake boundary question fixture generated\n", diff --git a/apps/backend/src/application/evidence.py b/apps/backend/src/application/evidence.py index 703a729..e5b575e 100644 --- a/apps/backend/src/application/evidence.py +++ b/apps/backend/src/application/evidence.py @@ -15,6 +15,7 @@ from src.domain.contracts import ( EvidenceResponse, EvidenceCreateRequest, EvidenceUpdateRequest, + PlanReviewStatus, ) @@ -32,19 +33,40 @@ def get_evidence_matrix(repository: object, *, article_id: UUID) -> EvidenceMatr claims = repository.claims.list_for_article(article_id) reasons = _insufficient_reasons(claims) if reasons and article.status != ArticleWorkflowStatus.PLAN_REVISION_REQUIRED: + from_status = article.status + now = _now() article = repository.articles.update_status( article_id=article_id, status=ArticleWorkflowStatus.PLAN_REVISION_REQUIRED, - updated_at=_now(), + updated_at=now, ) repository.articles.create_workflow_event( article_id=article_id, event_type="INSUFFICIENT_EVIDENCE_FOUND", - from_status=ArticleWorkflowStatus.RESEARCH_RUNNING, + from_status=from_status, to_status=ArticleWorkflowStatus.PLAN_REVISION_REQUIRED, actor_user_id=None, payload={"reasons": reasons}, - created_at=_now(), + created_at=now, + ) + if not reasons and article.status == ArticleWorkflowStatus.RESEARCH_RUNNING: + now = _now() + article = repository.articles.update_status( + article_id=article_id, + status=ArticleWorkflowStatus.EVIDENCE_MATRIX_READY, + updated_at=now, + ) + repository.articles.create_workflow_event( + article_id=article_id, + event_type="EVIDENCE_MATRIX_READY", + from_status=ArticleWorkflowStatus.RESEARCH_RUNNING, + to_status=ArticleWorkflowStatus.EVIDENCE_MATRIX_READY, + actor_user_id=None, + payload={ + "evidence_count": len(evidence), + "claim_count": len(claims), + }, + created_at=now, ) return EvidenceMatrixResponse( @@ -133,21 +155,69 @@ def remove_evidence( def start_draft(repository: object, *, article_id: UUID) -> AgentJobListResponse: + article = repository.articles.get(article_id) + if article.status != ArticleWorkflowStatus.EVIDENCE_MATRIX_READY: + raise PermissionError("Evidence matrix must be ready before draft production") + matrix = get_evidence_matrix(repository, article_id=article_id) if matrix.insufficient_evidence_reasons: raise PermissionError("Acceptable evidence is required before draft production") + if matrix.article.status != ArticleWorkflowStatus.EVIDENCE_MATRIX_READY: + raise PermissionError("Evidence matrix must be ready before draft production") - job = repository.agent_jobs.create( + approved_plan = _approved_plan(repository, article_id=article_id) + if approved_plan is None or not approved_plan.sections: + raise PermissionError("Approved plan with sections is required before draft production") + + claim_evidence_by_section = _claim_evidence_ids_by_section(matrix.claims) + now = _now() + repository.articles.update_status( article_id=article_id, - parent_job_id=None, - attempt=1, - job_type=AgentJobType.DRAFT_ASSEMBLY, - agent_profile="fake-draft-assembly", - status=AgentJobStatus.QUEUED, - input_files=[{"path": "inputs/evidence-matrix.json", "content_hash": None}], - queued_at=_now(), + status=ArticleWorkflowStatus.PARALLEL_PRODUCTION_RUNNING, + updated_at=now, ) - return AgentJobListResponse(jobs=[job]) + repository.articles.create_workflow_event( + article_id=article_id, + event_type="PARALLEL_PRODUCTION_STARTED", + from_status=ArticleWorkflowStatus.EVIDENCE_MATRIX_READY, + to_status=ArticleWorkflowStatus.PARALLEL_PRODUCTION_RUNNING, + actor_user_id=None, + payload={"section_count": len(approved_plan.sections)}, + created_at=now, + ) + + jobs = [] + for section in approved_plan.sections: + used_evidence_ids = claim_evidence_by_section.get(str(section.id), []) + jobs.append( + repository.agent_jobs.create( + article_id=article_id, + parent_job_id=None, + attempt=1, + job_type=AgentJobType.SECTION_SCAFFOLD, + agent_profile="fake-section-scaffold", + status=AgentJobStatus.QUEUED, + input_files=[ + {"path": "inputs/evidence-matrix.json", "content_hash": None}, + {"path": "inputs/approved-plan.json", "content_hash": None}, + {"path": f"inputs/sections/{section.id}.json", "content_hash": None}, + ], + payload={ + "artifact_key": f"section:{section.id}", + "artifact_label": section.heading, + "artifact_type": "section_scaffold", + "section_id": str(section.id), + "heading": section.heading, + "used_evidence_ids": used_evidence_ids, + "unsupported_claims": [], + "suggested_visuals": [], + "draft_markdown": "", + }, + queued_at=now, + ) + ) + + return AgentJobListResponse(jobs=jobs) def approve_final(repository: object, *, article_id: UUID) -> ArticleCreateResponse: @@ -212,6 +282,27 @@ def _resolve_manifest_id(repository: object, *, article_id: UUID) -> UUID: raise ValueError("Research manifest is required before adding manual evidence") +def _approved_plan(repository: object, article_id: UUID) -> object | None: + for plan in reversed(repository.article_plans.list_for_article(article_id)): + if plan.status == PlanReviewStatus.APPROVED: + return plan + return None + + +def _claim_evidence_ids_by_section(claims: list[object]) -> dict[str, list[str]]: + mapping: dict[str, list[str]] = {} + for claim in claims: + if claim.section_id is None: + continue + section_id = str(claim.section_id) + bucket = mapping.setdefault(section_id, []) + for evidence_id in claim.evidence_item_ids: + encoded = str(evidence_id) + if encoded not in bucket: + bucket.append(encoded) + return mapping + + def _insufficient_reasons(claims: list[object]) -> list[str]: reasons: list[str] = [] for claim in claims: diff --git a/apps/backend/src/application/research.py b/apps/backend/src/application/research.py index 693daaf..a0aac1d 100644 --- a/apps/backend/src/application/research.py +++ b/apps/backend/src/application/research.py @@ -67,6 +67,7 @@ def start_research_run( status=AgentJobStatus.SUCCEEDED, workspace_path=None, output_files=[{"path": artifact["object_key"]} for artifact in artifacts], + payload={}, error_category=None, error_message=None, stdout="fake research fetch completed\n", diff --git a/apps/backend/src/domain/contracts/models.py b/apps/backend/src/domain/contracts/models.py index 28ce153..c793e8e 100644 --- a/apps/backend/src/domain/contracts/models.py +++ b/apps/backend/src/domain/contracts/models.py @@ -429,6 +429,7 @@ class AgentJobSummary(ContractModel): workspace_path: str | None = None input_files: list[RunnerFileRef] = Field(default_factory=list) output_files: list[RunnerFileRef] = Field(default_factory=list) + payload: JsonObject = Field(default_factory=dict) error_category: AgentJobErrorCategory | None = None error_message: str | None = None stdout: str = "" diff --git a/apps/backend/src/infrastructure/repositories.py b/apps/backend/src/infrastructure/repositories.py index 190c3a0..6450350 100644 --- a/apps/backend/src/infrastructure/repositories.py +++ b/apps/backend/src/infrastructure/repositories.py @@ -1413,6 +1413,7 @@ class AgentJobsRepository: agent_profile: str, status: AgentJobStatus, input_files: list[JsonObject], + payload: JsonObject | None = None, queued_at: datetime, ) -> AgentJobSummary: job_id = uuid4() @@ -1428,6 +1429,7 @@ class AgentJobsRepository: agent_profile, status, input_files, + payload, queued_at ) VALUES ( @@ -1439,6 +1441,7 @@ class AgentJobsRepository: {placeholder}, {placeholder}, {placeholder}{json_cast}, + {placeholder}{json_cast}, {placeholder} ) """ @@ -1454,6 +1457,7 @@ class AgentJobsRepository: agent_profile, status.value, _json_value(input_files), + _json_value(payload or {}), _datetime_value(queued_at), ), ) @@ -1472,6 +1476,21 @@ class AgentJobsRepository: return [_agent_job_summary_from_row(row) for row in rows] + def list_for_article(self, article_id: UUID) -> list[AgentJobSummary]: + placeholder = self._repository.placeholder() + with self._repository.connection() as connection: + rows = connection.execute( + f""" + SELECT {self._select_columns()} + FROM agent_jobs + WHERE article_id = {placeholder} + ORDER BY queued_at DESC + """, + (str(article_id),), + ).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: @@ -1528,6 +1547,7 @@ class AgentJobsRepository: status: AgentJobStatus, workspace_path: str | None, output_files: list[JsonObject], + payload: JsonObject | None, error_category: AgentJobErrorCategory | None, error_message: str | None, stdout: str, @@ -1550,6 +1570,7 @@ class AgentJobsRepository: status = {placeholder}, workspace_path = {placeholder}, output_files = {placeholder}{json_cast}, + payload = {placeholder}{json_cast}, error_category = {placeholder}, error_message = {placeholder}, stdout = {placeholder}, @@ -1563,6 +1584,7 @@ class AgentJobsRepository: status.value, workspace_path, _json_value(output_files), + _json_value(payload or {}), _agent_error_category_value(error_category), error_message, stdout, @@ -1614,6 +1636,7 @@ class AgentJobsRepository: workspace_path, input_files, output_files, + payload, error_category, error_message, stdout, @@ -2165,6 +2188,7 @@ def _agent_job_summary_from_row(row: Any) -> AgentJobSummary: workspace_path=_row_value(row, "workspace_path"), input_files=_json_from_row(row, "input_files"), output_files=_json_from_row(row, "output_files"), + payload=_json_from_row(row, "payload"), error_category=_row_value(row, "error_category"), error_message=_row_value(row, "error_message"), stdout=_row_value(row, "stdout"), diff --git a/apps/backend/src/infrastructure/schema.py b/apps/backend/src/infrastructure/schema.py index 3d12c2e..ac97cf8 100644 --- a/apps/backend/src/infrastructure/schema.py +++ b/apps/backend/src/infrastructure/schema.py @@ -260,6 +260,7 @@ POSTGRES_SCHEMA_STATEMENTS: tuple[str, ...] = ( workspace_path TEXT, input_files JSONB NOT NULL DEFAULT '[]'::jsonb, output_files JSONB NOT NULL DEFAULT '[]'::jsonb, + payload JSONB NOT NULL DEFAULT '{{}}'::jsonb, error_category TEXT CHECK ( error_category IS NULL OR error_category IN ({AGENT_JOB_ERROR_CATEGORY_VALUES}) @@ -280,6 +281,7 @@ POSTGRES_SCHEMA_STATEMENTS: tuple[str, ...] = ( "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", + "ALTER TABLE agent_jobs ADD COLUMN IF NOT EXISTS payload JSONB NOT NULL DEFAULT '{}'::jsonb", """ CREATE TABLE IF NOT EXISTS research_run_manifests ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), @@ -568,6 +570,7 @@ SQLITE_SCHEMA_STATEMENTS: tuple[str, ...] = ( workspace_path TEXT, input_files TEXT NOT NULL DEFAULT '[]', output_files TEXT NOT NULL DEFAULT '[]', + payload TEXT NOT NULL DEFAULT '{}', error_category TEXT, error_message TEXT, stdout TEXT NOT NULL DEFAULT '', diff --git a/apps/backend/tests/integration/test_evidence_matrix_public_api.py b/apps/backend/tests/integration/test_evidence_matrix_public_api.py index 1415e4f..ef808cd 100644 --- a/apps/backend/tests/integration/test_evidence_matrix_public_api.py +++ b/apps/backend/tests/integration/test_evidence_matrix_public_api.py @@ -90,7 +90,7 @@ class EvidenceMatrixPublicApiTest(unittest.TestCase): ) self.assertEqual(200, evidence_response.status_code, evidence_response.text) body = evidence_response.json() - self.assertEqual("RESEARCH_RUNNING", body["article"]["status"]) + self.assertEqual("EVIDENCE_MATRIX_READY", body["article"]["status"]) self.assertFalse(body["insufficient_evidence_reasons"]) self.assertTrue(body["evidence"]) self.assertTrue(body["claims"]) diff --git a/apps/backend/tests/integration/test_parallel_production_public_api.py b/apps/backend/tests/integration/test_parallel_production_public_api.py new file mode 100644 index 0000000..5878bbe --- /dev/null +++ b/apps/backend/tests/integration/test_parallel_production_public_api.py @@ -0,0 +1,321 @@ +from __future__ import annotations + +import os +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_EDITOR_EMAIL = "editor@example.com" +DEMO_ADMIN_EMAIL = "admin@example.com" +DEMO_USER_EMAIL_HEADER = "X-Demo-User-Email" + + +class ParallelProductionPublicApiTest(unittest.TestCase): + def setUp(self) -> None: + self.tmp_dir = tempfile.TemporaryDirectory() + os.environ["OBJECT_STORAGE_LOCAL_ROOT"] = str(Path(self.tmp_dir.name) / "objects") + dsn = f"sqlite:///{Path(self.tmp_dir.name) / 'parallel-production.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() + os.environ.pop("OBJECT_STORAGE_LOCAL_ROOT", None) + self.tmp_dir.cleanup() + + def test_start_production_requires_evidence_matrix_ready(self) -> None: + article_id, _ = self._create_article_with_approved_plan() + research_response = self.client.post( + f"/api/articles/{article_id}/research/start", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(201, research_response.status_code, research_response.text) + + start_response = self.client.post( + f"/api/articles/{article_id}/draft/start", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(409, start_response.status_code, start_response.text) + self.assertIn("Evidence matrix must be ready", start_response.text) + + def test_start_production_creates_one_section_job_per_approved_plan_section(self) -> None: + article_id, approved_section_count = self._prepare_article_for_parallel_production() + + start_production_response = self.client.post( + f"/api/articles/{article_id}/draft/start", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(202, start_production_response.status_code, start_production_response.text) + + jobs = start_production_response.json()["jobs"] + section_jobs = [job for job in jobs if job["job_type"] == "SECTION_SCAFFOLD"] + self.assertEqual( + approved_section_count, + len(section_jobs), + "Expected one SECTION_SCAFFOLD job per approved plan section", + ) + + def test_section_jobs_run_independently_and_can_fail_independently(self) -> None: + article_id, _ = self._prepare_article_for_parallel_production() + created_jobs = self._start_parallel_production(article_id) + section_jobs = [job for job in created_jobs if job["job_type"] == "SECTION_SCAFFOLD"] + self.assertGreaterEqual(len(section_jobs), 2) + + success_job = section_jobs[0] + failed_job = section_jobs[1] + complete_success = self._complete_job( + success_job["id"], + output={ + "status": "SUCCEEDED", + "output_files": [{"path": "outputs/section-success.md"}], + "payload": { + "used_evidence_ids": success_job["payload"]["used_evidence_ids"], + "unsupported_claims": [], + "draft_markdown": "Supported section draft", + }, + }, + ) + self.assertEqual("SUCCEEDED", complete_success["status"]) + self.assertIsNone(complete_success["error_category"]) + + unsupported_claims = [ + { + "claim_text": "Unverified benchmark introduced during scaffolding.", + "risk_level": "high", + } + ] + complete_failure = self._complete_job( + failed_job["id"], + output={ + "status": "SUCCEEDED", + "output_files": [{"path": "outputs/section-failed.md"}], + "payload": { + "used_evidence_ids": failed_job["payload"]["used_evidence_ids"], + "unsupported_claims": unsupported_claims, + "draft_markdown": "Draft includes unsupported claim", + }, + }, + ) + self.assertEqual("FAILED", complete_failure["status"]) + self.assertEqual("UNSUPPORTED_CLAIMS_FOUND", complete_failure["error_category"]) + self.assertEqual(unsupported_claims, complete_failure["payload"]["unsupported_claims"]) + + def test_retry_failed_section_does_not_rerun_successful_sections(self) -> None: + article_id, _ = self._prepare_article_for_parallel_production() + created_jobs = self._start_parallel_production(article_id) + section_jobs = [job for job in created_jobs if job["job_type"] == "SECTION_SCAFFOLD"] + self.assertGreaterEqual(len(section_jobs), 2) + + success_job = self._complete_job( + section_jobs[0]["id"], + output={ + "status": "SUCCEEDED", + "output_files": [{"path": "outputs/section-success.md"}], + "payload": {"unsupported_claims": []}, + }, + ) + failed_job = self._complete_job( + section_jobs[1]["id"], + output={ + "status": "FAILED", + "error_category": "CLI_EXIT_CODE_FAILURE", + "error_message": "Runner failed for one section", + "payload": {"unsupported_claims": []}, + }, + ) + + retry_response = self.client.post( + f"/api/agent-jobs/{failed_job['id']}/retry", + headers={DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}, + ) + self.assertEqual(201, retry_response.status_code, retry_response.text) + retry_job = retry_response.json()["job"] + self.assertEqual("QUEUED", retry_job["status"]) + self.assertEqual(failed_job["id"], retry_job["parent_job_id"]) + self.assertEqual(2, retry_job["attempt"]) + + detail = self.client.get( + f"/api/articles/{article_id}", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ).json() + successful_descendants = [ + job + for job in detail["agent_jobs"] + if job.get("parent_job_id") == success_job["id"] + ] + failed_descendants = [ + job + for job in detail["agent_jobs"] + if job.get("parent_job_id") == failed_job["id"] + ] + self.assertEqual([], successful_descendants) + self.assertEqual(1, len(failed_descendants)) + self.assertEqual("QUEUED", failed_descendants[0]["status"]) + + def test_each_scaffold_payload_lists_used_evidence_ids(self) -> None: + article_id, _ = self._prepare_article_for_parallel_production() + evidence = self.client.get( + f"/api/articles/{article_id}/evidence", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ).json()["evidence"] + evidence_ids = {item["id"] for item in evidence} + + jobs = self._start_parallel_production(article_id) + section_jobs = [job for job in jobs if job["job_type"] == "SECTION_SCAFFOLD"] + self.assertTrue(section_jobs) + self.assertTrue(any(job["payload"]["used_evidence_ids"] for job in section_jobs)) + for job in section_jobs: + self.assertIn("used_evidence_ids", job["payload"]) + self.assertIsInstance(job["payload"]["used_evidence_ids"], list) + for evidence_id in job["payload"]["used_evidence_ids"]: + self.assertIn(evidence_id, evidence_ids) + + def test_unsupported_claims_from_scaffold_are_captured_in_article_detail(self) -> None: + article_id, _ = self._prepare_article_for_parallel_production() + jobs = self._start_parallel_production(article_id) + section_job = [job for job in jobs if job["job_type"] == "SECTION_SCAFFOLD"][0] + unsupported_claims = [ + { + "claim_text": "Unverified migration timeline claim", + "risk_level": "medium", + } + ] + + completed_job = self._complete_job( + section_job["id"], + output={ + "status": "SUCCEEDED", + "output_files": [{"path": "outputs/section-with-unsupported.md"}], + "payload": { + "used_evidence_ids": section_job["payload"]["used_evidence_ids"], + "unsupported_claims": unsupported_claims, + "draft_markdown": "Draft with unsupported claim", + }, + }, + ) + self.assertEqual("FAILED", completed_job["status"]) + self.assertEqual("UNSUPPORTED_CLAIMS_FOUND", completed_job["error_category"]) + + detail_response = self.client.get( + f"/api/articles/{article_id}", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(200, detail_response.status_code, detail_response.text) + article_jobs = detail_response.json()["agent_jobs"] + refreshed_job = [job for job in article_jobs if job["id"] == section_job["id"]][0] + self.assertEqual(unsupported_claims, refreshed_job["payload"]["unsupported_claims"]) + self.assertIn("Unsupported claims introduced during scaffolding", refreshed_job["error_message"]) + + def _start_parallel_production(self, article_id: str) -> list[dict[str, Any]]: + response = self.client.post( + f"/api/articles/{article_id}/draft/start", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(202, response.status_code, response.text) + return response.json()["jobs"] + + def _complete_job(self, job_id: str, *, output: dict[str, Any]) -> dict[str, Any]: + response = self.client.post( + f"/internal/agent-jobs/{job_id}/complete", + json={ + "workspace_path": f"/tmp/{job_id}", + "stdout": "fake section scaffolding runner\n", + "stderr": "", + "exit_code": 0, + "duration_ms": 2, + "output": output, + }, + ) + self.assertEqual(200, response.status_code, response.text) + return response.json()["job"] + + def _prepare_article_for_parallel_production(self) -> tuple[str, int]: + article_id, approved_section_count = self._create_article_with_approved_plan() + + research_response = self.client.post( + f"/api/articles/{article_id}/research/start", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(201, research_response.status_code, research_response.text) + + evidence_response = self.client.get( + f"/api/articles/{article_id}/evidence", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(200, evidence_response.status_code, evidence_response.text) + self.assertFalse(evidence_response.json()["insufficient_evidence_reasons"]) + self.assertEqual("EVIDENCE_MATRIX_READY", evidence_response.json()["article"]["status"]) + return article_id, approved_section_count + + def _create_article_with_approved_plan(self) -> tuple[str, int]: + site = self.client.get( + "/api/sites", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ).json()[0]["site"] + article = self.client.post( + "/api/articles", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + json={ + "target_site_id": site["id"], + "brief_description": "Run parallel production jobs from approved sections.", + "working_title": "Parallel Production Jobs", + "content_type": "longform_guide", + "primary_keyword": "parallel production jobs", + }, + ).json()["article"] + article_id = article["id"] + + questions = self.client.post( + f"/api/articles/{article_id}/boundary-questions/generate", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ).json()["questions"] + for question in questions: + if question["is_required"]: + patch_response = self.client.patch( + f"/api/articles/{article_id}/boundary-questions/{question['id']}", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + json={"answer": f"Answer for {question['category']}"}, + ) + self.assertEqual(200, patch_response.status_code, patch_response.text) + + submit_response = self.client.post( + f"/api/articles/{article_id}/boundary-questions/submit", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(200, submit_response.status_code, submit_response.text) + + plan_response = self.client.post( + f"/api/articles/{article_id}/plan/generate", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(201, plan_response.status_code, plan_response.text) + plan = plan_response.json()["plan"] + self.assertGreaterEqual(len(plan["sections"]), 1) + + approve_response = self.client.post( + f"/api/articles/{article_id}/plans/{plan['id']}/approve", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(200, approve_response.status_code, approve_response.text) + return article_id, len(plan["sections"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/apps/frontend/src/features/article-detail/model.ts b/apps/frontend/src/features/article-detail/model.ts index b37c91f..67bb477 100644 --- a/apps/frontend/src/features/article-detail/model.ts +++ b/apps/frontend/src/features/article-detail/model.ts @@ -1,4 +1,4 @@ -import type { ArticleDetailResponse, WorkflowEventSummary } from "@pipeline/shared"; +import type { AgentJobSummary, ArticleDetailResponse, WorkflowEventSummary } from "@pipeline/shared"; export type ArticleTimelineItem = Pick< WorkflowEventSummary, @@ -15,6 +15,19 @@ export type DetailSummary = { targetSite: string; updatedAt: string; timeline: ArticleTimelineItem[]; + productionArtifacts: ProductionArtifactRow[]; +}; + +export type ProductionArtifactRow = { + jobId: string; + artifactKey: string; + artifactLabel: string; + status: string; + attempt: number; + jobType: string; + usedEvidenceIds: string[]; + unsupportedClaims: string[]; + errorMessage: string | null; }; export function buildArticleTimeline(detail: ArticleDetailResponse): ArticleTimelineItem[] { @@ -38,5 +51,64 @@ export function buildDetailSummary(detail: ArticleDetailResponse): DetailSummary targetSite: detail.target_site?.name ?? "Unknown", updatedAt: detail.article.updated_at, timeline: buildArticleTimeline(detail), + productionArtifacts: buildProductionArtifactRows(detail.agent_jobs ?? []), }; } + +export function buildProductionArtifactRows( + jobs: readonly AgentJobSummary[], +): ProductionArtifactRow[] { + return jobs + .filter((job) => job.job_type === "SECTION_SCAFFOLD") + .map((job) => { + const payload = (job.payload ?? {}) as Record; + const artifactKey = stringValue(payload.artifact_key) || `job:${job.id}`; + const artifactLabel = + stringValue(payload.artifact_label) + || stringValue(payload.heading) + || artifactKey; + return { + jobId: job.id, + artifactKey, + artifactLabel, + status: job.status, + attempt: job.attempt ?? 1, + jobType: job.job_type, + usedEvidenceIds: stringArray(payload.used_evidence_ids), + unsupportedClaims: unsupportedClaims(payload.unsupported_claims), + errorMessage: job.error_message ?? null, + }; + }) + .sort((left, right) => left.artifactLabel.localeCompare(right.artifactLabel)); +} + +function stringValue(value: unknown): string { + return typeof value === "string" ? value : ""; +} + +function stringArray(value: unknown): string[] { + if (!Array.isArray(value)) { + return []; + } + return value.filter((item): item is string => typeof item === "string"); +} + +function unsupportedClaims(value: unknown): string[] { + if (!Array.isArray(value)) { + return []; + } + const claims: string[] = []; + for (const item of value) { + if (typeof item === "string") { + claims.push(item); + continue; + } + if (typeof item === "object" && item !== null && "claim_text" in item) { + const claimText = item.claim_text; + if (typeof claimText === "string") { + claims.push(claimText); + } + } + } + return claims; +} diff --git a/apps/frontend/src/features/article-detail/ui.tsx b/apps/frontend/src/features/article-detail/ui.tsx index 436913d..19828e0 100644 --- a/apps/frontend/src/features/article-detail/ui.tsx +++ b/apps/frontend/src/features/article-detail/ui.tsx @@ -8,6 +8,7 @@ type DetailShellProps = { export function ArticleDetailShell({ summary }: DetailShellProps) { const timeline = summary.timeline; + const productionArtifacts = summary.productionArtifacts; return (
@@ -66,6 +67,46 @@ export function ArticleDetailShell({ summary }: DetailShellProps) { ))} {timeline.length === 0 ?
  • No workflow events yet.
  • : null} + +

    Production artifacts

    + + + + + + + + + + + + {productionArtifacts.map((artifact) => ( + + + + + + + + ))} + {productionArtifacts.length === 0 ? ( + + + + ) : null} + +
    ArtifactStatusAttemptUsed evidence IDsUnsupported claims
    +
    {artifact.artifactLabel}
    +
    {artifact.artifactKey}
    +
    {artifact.status}{artifact.attempt} + {artifact.usedEvidenceIds.length + ? artifact.usedEvidenceIds.join(", ") + : "—"} + + {artifact.unsupportedClaims.length + ? artifact.unsupportedClaims.join("; ") + : artifact.errorMessage ?? "—"} +
    No production artifacts yet.
    ); } diff --git a/apps/frontend/tests/article_detail.model.test.mjs b/apps/frontend/tests/article_detail.model.test.mjs new file mode 100644 index 0000000..62746fb --- /dev/null +++ b/apps/frontend/tests/article_detail.model.test.mjs @@ -0,0 +1,104 @@ +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; + +import ts from "typescript"; + +const testDir = path.dirname(fileURLToPath(import.meta.url)); +const sourcePath = path.resolve(testDir, "../src/features/article-detail/model.ts"); +const source = readFileSync(sourcePath, "utf8"); +const compiled = ts.transpileModule(source, { + compilerOptions: { module: ts.ModuleKind.CommonJS, target: ts.ScriptTarget.ES2022 }, +}); + +const moduleExports = {}; +new Function("exports", compiled.outputText)(moduleExports); + +const { buildProductionArtifactRows } = moduleExports; + +const rows = buildProductionArtifactRows([ + { + id: "job-section-1", + article_id: "a1", + parent_job_id: null, + attempt: 1, + job_type: "SECTION_SCAFFOLD", + agent_profile: "fake-section-scaffold", + status: "SUCCEEDED", + workspace_path: null, + input_files: [], + output_files: [], + payload: { + artifact_key: "section:1", + artifact_label: "Context and scope", + used_evidence_ids: ["e1", "e2"], + unsupported_claims: [], + }, + error_category: null, + error_message: null, + stdout: "", + stderr: "", + exit_code: 0, + duration_ms: 3, + queued_at: "2026-05-21T00:00:00Z", + started_at: "2026-05-21T00:00:01Z", + finished_at: "2026-05-21T00:00:02Z", + }, + { + id: "job-section-2", + article_id: "a1", + parent_job_id: null, + attempt: 1, + job_type: "SECTION_SCAFFOLD", + agent_profile: "fake-section-scaffold", + status: "FAILED", + workspace_path: null, + input_files: [], + output_files: [], + payload: { + artifact_key: "section:2", + artifact_label: "Implementation workflow", + used_evidence_ids: ["e2"], + unsupported_claims: [{ claim_text: "Unverified claim" }], + }, + error_category: "UNSUPPORTED_CLAIMS_FOUND", + error_message: "Unsupported claims introduced during scaffolding: 1", + stdout: "", + stderr: "", + exit_code: 0, + duration_ms: 3, + queued_at: "2026-05-21T00:00:00Z", + started_at: "2026-05-21T00:00:01Z", + finished_at: "2026-05-21T00:00:02Z", + }, + { + id: "job-non-artifact", + article_id: "a1", + parent_job_id: null, + attempt: 1, + job_type: "TEST_CODEX", + agent_profile: "fake-codex", + status: "SUCCEEDED", + workspace_path: null, + input_files: [], + output_files: [], + payload: {}, + error_category: null, + error_message: null, + stdout: "", + stderr: "", + exit_code: 0, + duration_ms: 0, + queued_at: "2026-05-21T00:00:00Z", + started_at: "2026-05-21T00:00:01Z", + finished_at: "2026-05-21T00:00:02Z", + }, +]); + +assert.equal(rows.length, 2); +assert.equal(rows[0].artifactLabel, "Context and scope"); +assert.equal(rows[0].status, "SUCCEEDED"); +assert.deepEqual(rows[0].usedEvidenceIds, ["e1", "e2"]); +assert.equal(rows[1].status, "FAILED"); +assert.deepEqual(rows[1].unsupportedClaims, ["Unverified claim"]); diff --git a/packages/shared/src/api-types.ts b/packages/shared/src/api-types.ts index de3fbd0..fec8c3b 100644 --- a/packages/shared/src/api-types.ts +++ b/packages/shared/src/api-types.ts @@ -35,6 +35,7 @@ export type AgentJobSummary = { input_files?: RunnerFileRef[]; job_type: AgentJobType; output_files?: RunnerFileRef[]; + payload?: Record; parent_job_id?: string | null; queued_at: string; started_at?: string | null; diff --git a/tasks/012-parallel-production-jobs.md b/tasks/012-parallel-production-jobs.md index 6601c28..236cd4f 100644 --- a/tasks/012-parallel-production-jobs.md +++ b/tasks/012-parallel-production-jobs.md @@ -32,14 +32,14 @@ Development description: Implement parallel section scaffolding, FAQ, SEO brief, ## Acceptance Criteria -- [ ] TDD pre-requirement: before implementation, write one failing behavior test that starts production and creates one job per plan section; proceed one artifact behavior at a time and record evidence in `Result`. -- [ ] Production cannot start before evidence matrix is ready. -- [ ] One section scaffold job is created per approved plan section. -- [ ] Jobs run independently and can fail independently. -- [ ] Retrying one failed section does not rerun successful sections. -- [ ] Each scaffold lists used evidence IDs. -- [ ] Unsupported claims introduced during scaffolding are captured and shown. -- [ ] UI displays per-artifact job state. +- [x] TDD pre-requirement: before implementation, write one failing behavior test that starts production and creates one job per plan section; proceed one artifact behavior at a time and record evidence in `Result`. +- [x] Production cannot start before evidence matrix is ready. +- [x] One section scaffold job is created per approved plan section. +- [x] Jobs run independently and can fail independently. +- [x] Retrying one failed section does not rerun successful sections. +- [x] Each scaffold lists used evidence IDs. +- [x] Unsupported claims introduced during scaffolding are captured and shown. +- [x] UI displays per-artifact job state. ## Verification @@ -49,9 +49,26 @@ Development description: Implement parallel section scaffolding, FAQ, SEO brief, ## 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: Completed (TDD RED -> GREEN). +- TDD plan: + 1. Keep the existing RED integration behavior for `/api/articles/{article_id}/draft/start` (one `SECTION_SCAFFOLD` job per approved section). + 2. Implement readiness gating and parallel section job orchestration in `start_draft`. + 3. Add integration tests for independent failures and retry lineage. + 4. Add frontend model/UI projection for per-artifact job state. +- Red evidence: + - Existing failing test: `test_start_production_creates_one_section_job_per_approved_plan_section`. + - Failure before implementation: `AssertionError: 4 != 0 : Expected one SECTION_SCAFFOLD job per approved plan section`. +- Green evidence: + - `start_draft` now requires `EVIDENCE_MATRIX_READY`, transitions article to `PARALLEL_PRODUCTION_RUNNING`, and creates one `SECTION_SCAFFOLD` job per approved plan section. + - Section jobs persist scaffold payload (`artifact_key`, `artifact_label`, `used_evidence_ids`, `unsupported_claims`, `draft_markdown`). + - Completing a section job with non-empty `unsupported_claims` forces job status to `FAILED` with `UNSUPPORTED_CLAIMS_FOUND`. + - Retry creates a child attempt only for the failed section job (successful siblings are untouched). + - Article detail API now includes article-specific `agent_jobs` with payload for UI rendering. +- Refactor notes: + - Added `payload` JSON field to `agent_jobs` contract + storage for deterministic scaffold metadata projection. + - Added repository helper `list_for_article` for article-level job timeline queries. +- Verification output: + - `PYTHONPATH=/private/tmp/pupline-backend-deps python3 -m unittest apps/backend/tests/integration/test_parallel_production_public_api.py` -> `Ran 6 tests ... OK` + - `PYTHONPATH=/private/tmp/pupline-backend-deps python3 -m unittest apps/backend/tests/integration/test_evidence_matrix_public_api.py` -> `Ran 2 tests ... OK` + - `node apps/frontend/tests/article_detail.model.test.mjs` -> `pass` + - `node apps/frontend/tests/agent_jobs.model.test.mjs` -> `pass`