From 7d4a79acb566f595caf89589a935504a60a7f420 Mon Sep 17 00:00:00 2001 From: "E.Gavrilov" Date: Thu, 21 May 2026 21:56:44 +0300 Subject: [PATCH] Task 010 implement research manifest flow --- apps/backend/src/application/articles.py | 2 + apps/backend/src/application/research.py | 180 ++++++++++++++++++ apps/backend/src/domain/contracts/__init__.py | 4 + apps/backend/src/domain/contracts/models.py | 11 ++ apps/backend/src/domain/contracts/openapi.py | 4 + .../src/infrastructure/object_storage.py | 63 ++++++ .../src/infrastructure/repositories.py | 132 +++++++++++++ apps/backend/src/presentation/routes/plans.py | 33 +++- .../test_plan_generation_review_public_api.py | 4 +- .../test_research_manifest_public_api.py | 150 +++++++++++++++ apps/frontend/package.json | 2 +- .../articles/[articleId]/research/page.tsx | 14 ++ .../src/features/article-detail/ui.tsx | 3 + .../src/features/research-manifest/model.ts | 40 ++++ .../src/features/research-manifest/ui.tsx | 87 +++++++++ .../src/pages/article-research/index.tsx | 34 ++++ apps/frontend/src/shared/pipeline-api.ts | 13 ++ .../tests/research_manifest.model.test.mjs | 55 ++++++ packages/shared/openapi.json | 105 +++++++++- packages/shared/src/api-types.ts | 11 ++ tasks/010-research-fetch-and-s3-manifest.md | 28 +-- tests/smoke/plan-flow.sh | 2 +- tests/smoke/research-flow.sh | 112 +++++++++++ 23 files changed, 1063 insertions(+), 26 deletions(-) create mode 100644 apps/backend/src/application/research.py create mode 100644 apps/backend/src/infrastructure/object_storage.py create mode 100644 apps/backend/tests/integration/test_research_manifest_public_api.py create mode 100644 apps/frontend/src/app/articles/[articleId]/research/page.tsx create mode 100644 apps/frontend/src/features/research-manifest/model.ts create mode 100644 apps/frontend/src/features/research-manifest/ui.tsx create mode 100644 apps/frontend/src/pages/article-research/index.tsx create mode 100644 apps/frontend/tests/research_manifest.model.test.mjs create mode 100755 tests/smoke/research-flow.sh diff --git a/apps/backend/src/application/articles.py b/apps/backend/src/application/articles.py index 1e88f89..57a2e39 100644 --- a/apps/backend/src/application/articles.py +++ b/apps/backend/src/application/articles.py @@ -60,12 +60,14 @@ def get_article_detail( workflow_events = repository.articles.list_workflow_events(article_id) boundary_questions = repository.boundary_questions.list_for_article(article_id) plans = repository.article_plans.list_for_article(article_id) + research_manifests = repository.research_manifests.list_for_article(article_id) return ArticleDetailResponse( article=article, target_site=target_site, workflow_events=workflow_events, boundary_questions=boundary_questions, plan=plans[-1] if plans else None, + research_manifests=research_manifests, ) diff --git a/apps/backend/src/application/research.py b/apps/backend/src/application/research.py new file mode 100644 index 0000000..693daaf --- /dev/null +++ b/apps/backend/src/application/research.py @@ -0,0 +1,180 @@ +from __future__ import annotations + +import hashlib +import json +from datetime import UTC, datetime +from urllib.parse import urlparse +from uuid import UUID, uuid4 + +from src.domain.contracts import ( + AgentJobStatus, + AgentJobType, + ArticleWorkflowStatus, + PlanReviewStatus, + ResearchListResponse, + ResearchStartResponse, +) +from src.infrastructure.object_storage import ObjectStorageClient + + +def start_research_run( + repository: object, + *, + article_id: UUID, + object_storage: ObjectStorageClient, +) -> ResearchStartResponse: + article = repository.articles.get(article_id) + approved_plan = _approved_plan(repository, article_id) + if approved_plan is None: + raise PermissionError("Plan must be approved before research can start") + + now = _now() + if article.status != ArticleWorkflowStatus.RESEARCH_RUNNING: + article = repository.articles.update_status( + article_id=article_id, + status=ArticleWorkflowStatus.RESEARCH_RUNNING, + updated_at=now, + ) + + job = repository.agent_jobs.create( + article_id=article_id, + parent_job_id=None, + attempt=1, + job_type=AgentJobType.RESEARCH, + agent_profile="fake-research-fetch", + status=AgentJobStatus.SUCCEEDED, + input_files=[{"path": "inputs/approved-plan.json", "content_hash": None}], + queued_at=now, + ) + run_id = uuid4() + s3_prefix = f"research/{article_id}/{run_id}" + artifacts = _build_fake_research_artifacts( + article=article, + plan=approved_plan, + s3_prefix=s3_prefix, + object_storage=object_storage, + retrieved_at=now, + ) + manifest = repository.research_manifests.create( + article_id=article_id, + agent_job_id=job.id, + s3_prefix=s3_prefix, + artifacts=artifacts, + created_at=now, + ) + job = repository.agent_jobs.complete( + job_id=job.id, + status=AgentJobStatus.SUCCEEDED, + workspace_path=None, + output_files=[{"path": artifact["object_key"]} for artifact in artifacts], + error_category=None, + error_message=None, + stdout="fake research fetch completed\n", + stderr="", + exit_code=0, + duration_ms=0, + finished_at=now, + ) + return ResearchStartResponse(article=article, job=job, manifest=manifest) + + +def list_research_runs(repository: object, *, article_id: UUID) -> ResearchListResponse: + repository.articles.get(article_id) + return ResearchListResponse( + manifests=repository.research_manifests.list_for_article(article_id), + ) + + +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 _build_fake_research_artifacts( + *, + article: object, + plan: object, + s3_prefix: str, + object_storage: ObjectStorageClient, + retrieved_at: datetime, +) -> list[dict[str, object]]: + sources = [ + { + "url": "https://example.com/research/content-governance", + "title": "Content governance implementation guide", + "source_type": "primary_reference", + "summary": "Governance programs need explicit review gates and source evidence.", + "section": "A source-section snapshot about governance workflow controls.", + "ip": "93.184.216.34", + "whois": "Example Domain Registry", + }, + { + "url": "https://example.org/research/evidence-standards", + "title": "Evidence standards for editorial claims", + "source_type": "standards_reference", + "summary": "Editorial claims should be linked to durable source artifacts.", + "section": "A source-section snapshot about evidence standards and claim review.", + "ip": "93.184.216.34", + "whois": "Example Organization Registry", + }, + ] + + artifacts: list[dict[str, object]] = [] + for index, source in enumerate(sources, start=1): + parsed = urlparse(str(source["url"])) + base_key = f"{s3_prefix}/source-{index}" + content_key = f"{base_key}/section.json" + metadata_key = f"{base_key}/metadata.json" + payload = { + "article_id": str(article.id), + "plan_id": str(plan.id), + "plan_version": plan.version, + "source_url": source["url"], + "section_snapshot": source["section"], + "retrieved_at": retrieved_at.isoformat(), + } + metadata = { + "source_title": source["title"], + "domain": parsed.netloc, + "source_type": source["source_type"], + "summary": source["summary"], + "retrieved_at": retrieved_at.isoformat(), + "artifact_link": content_key, + "metadata_object_key": metadata_key, + "search_path": [ + article.primary_keyword or article.working_title or "article research", + str(source["title"]), + ], + "selection_criteria": [ + "Matches approved plan evidence needs", + "Deterministic fake fixture for demo and tests", + ], + "head_metadata": { + "title": source["title"], + "description": source["summary"], + "canonical": source["url"], + }, + "server_ip_address": source["ip"], + "whois_owner": source["whois"], + } + content = json.dumps(payload, sort_keys=True, indent=2) + metadata_content = json.dumps(metadata, sort_keys=True, indent=2) + object_storage.put_text(object_key=content_key, content=content) + object_storage.put_text(object_key=metadata_key, content=metadata_content) + artifacts.append( + { + "artifact_type": "source_section_snapshot", + "source_url": source["url"], + "object_key": content_key, + "content_hash": hashlib.sha256(content.encode("utf-8")).hexdigest(), + "metadata": metadata, + } + ) + + return artifacts + + +def _now() -> datetime: + return datetime.now(UTC) diff --git a/apps/backend/src/domain/contracts/__init__.py b/apps/backend/src/domain/contracts/__init__.py index 76e1a9f..58e68ab 100644 --- a/apps/backend/src/domain/contracts/__init__.py +++ b/apps/backend/src/domain/contracts/__init__.py @@ -47,6 +47,8 @@ from .models import ( PublishingRules, ResearchArtifactManifestSummary, ResearchArtifactSummary, + ResearchListResponse, + ResearchStartResponse, ReviewActionResponse, ReviewSummary, RunnerFileRef, @@ -111,6 +113,8 @@ __all__ = [ "PublishingStatus", "ResearchArtifactManifestSummary", "ResearchArtifactSummary", + "ResearchListResponse", + "ResearchStartResponse", "ReviewActionResponse", "ReviewStatus", "ReviewSummary", diff --git a/apps/backend/src/domain/contracts/models.py b/apps/backend/src/domain/contracts/models.py index 1d2c7d5..21c143c 100644 --- a/apps/backend/src/domain/contracts/models.py +++ b/apps/backend/src/domain/contracts/models.py @@ -299,6 +299,17 @@ class ResearchArtifactManifestSummary(ContractModel): created_at: datetime +class ResearchStartResponse(ContractModel): + article: ArticleSummary + job: AgentJobSummary + manifest: ResearchArtifactManifestSummary + + +class ResearchListResponse(ContractModel): + manifests: list[ResearchArtifactManifestSummary] = Field(default_factory=list) + evidence: list["EvidenceSummary"] = Field(default_factory=list) + + class EvidenceSummary(ContractModel): id: UUID article_id: UUID diff --git a/apps/backend/src/domain/contracts/openapi.py b/apps/backend/src/domain/contracts/openapi.py index c4d1c48..f14d3de 100644 --- a/apps/backend/src/domain/contracts/openapi.py +++ b/apps/backend/src/domain/contracts/openapi.py @@ -53,6 +53,8 @@ from .models import ( PublishingRules, ResearchArtifactManifestSummary, ResearchArtifactSummary, + ResearchListResponse, + ResearchStartResponse, ReviewActionResponse, ReviewSummary, RunnerFileRef, @@ -117,6 +119,8 @@ CONTRACT_SCHEMA_MODELS: tuple[type[BaseModel], ...] = ( PlanListResponse, ResearchArtifactSummary, ResearchArtifactManifestSummary, + ResearchStartResponse, + ResearchListResponse, EvidenceSummary, ClaimSummary, DraftSummary, diff --git a/apps/backend/src/infrastructure/object_storage.py b/apps/backend/src/infrastructure/object_storage.py new file mode 100644 index 0000000..a20daf7 --- /dev/null +++ b/apps/backend/src/infrastructure/object_storage.py @@ -0,0 +1,63 @@ +from __future__ import annotations + +import os +from pathlib import Path + + +class ObjectStorageClient: + def put_text(self, *, object_key: str, content: str) -> str: + raise NotImplementedError + + +class LocalObjectStorageClient(ObjectStorageClient): + def __init__(self, root: Path) -> None: + self.root = root + + def put_text(self, *, object_key: str, content: str) -> str: + path = self.root / object_key + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(content, encoding="utf-8") + return f"file://{path}" + + +class S3ObjectStorageClient(ObjectStorageClient): + def __init__(self) -> None: + import boto3 + from botocore.config import Config + + self.bucket = _required_env("OBJECT_STORAGE_BUCKET") + self.client = boto3.client( + "s3", + endpoint_url=_required_env("OBJECT_STORAGE_ENDPOINT"), + aws_access_key_id=_required_env("OBJECT_STORAGE_ACCESS_KEY_ID"), + aws_secret_access_key=_required_env("OBJECT_STORAGE_SECRET_ACCESS_KEY"), + region_name=os.environ.get("OBJECT_STORAGE_REGION", "us-east-1"), + config=Config(signature_version="s3v4"), + ) + + def put_text(self, *, object_key: str, content: str) -> str: + self.client.put_object( + Bucket=self.bucket, + Key=object_key, + Body=content.encode("utf-8"), + ContentType="application/json; charset=utf-8", + ) + return f"s3://{self.bucket}/{object_key}" + + +def open_object_storage_client() -> ObjectStorageClient: + local_root = os.environ.get("OBJECT_STORAGE_LOCAL_ROOT") + if local_root: + return LocalObjectStorageClient(Path(local_root)) + + if os.environ.get("OBJECT_STORAGE_ENDPOINT"): + return S3ObjectStorageClient() + + return LocalObjectStorageClient(Path("/tmp/pipeline-object-storage")) + + +def _required_env(name: str) -> str: + value = os.environ.get(name) + if not value: + raise RuntimeError(f"Missing required environment variable: {name}") + return value diff --git a/apps/backend/src/infrastructure/repositories.py b/apps/backend/src/infrastructure/repositories.py index 27a9a73..5930356 100644 --- a/apps/backend/src/infrastructure/repositories.py +++ b/apps/backend/src/infrastructure/repositories.py @@ -20,6 +20,8 @@ from src.domain.contracts import ( PlanSummary, PublishingRules, PublishingStatus, + ResearchArtifactManifestSummary, + ResearchArtifactSummary, ArticleWorkflowStatus, WorkflowEventSummary, Role, @@ -46,6 +48,7 @@ class BackendRepository: self.articles = ArticlesRepository(self) self.boundary_questions = BoundaryQuestionsRepository(self) self.article_plans = ArticlePlansRepository(self) + self.research_manifests = ResearchManifestsRepository(self) self.agent_jobs = AgentJobsRepository(self) self.script_config_versions = ScriptConfigVersionsRepository(self) self.script_config_version_events = ScriptConfigVersionAuditEventsRepository(self) @@ -1005,6 +1008,121 @@ class ArticlePlansRepository: """ +class ResearchManifestsRepository: + def __init__(self, repository: BackendRepository) -> None: + self._repository = repository + + def create( + self, + *, + article_id: UUID, + agent_job_id: UUID, + s3_prefix: str, + artifacts: list[JsonObject], + created_at: datetime, + ) -> ResearchArtifactManifestSummary: + manifest_id = uuid4() + placeholder = self._repository.placeholder() + json_cast = self._repository.json_cast() + source_urls = [artifact["source_url"] for artifact in artifacts] + object_keys = [artifact["object_key"] for artifact in artifacts] + content_hashes = [artifact["content_hash"] for artifact in artifacts] + artifact_types = [artifact["artifact_type"] for artifact in artifacts] + metadata_references = [ + artifact.get("metadata", {}).get("metadata_object_key") + for artifact in artifacts + ] + with self._repository.connection() as connection: + connection.execute( + f""" + INSERT INTO research_run_manifests ( + id, + article_id, + agent_job_id, + s3_prefix, + source_urls, + object_keys, + content_hashes, + artifact_types, + metadata_references, + artifacts, + created_at + ) + VALUES ( + {placeholder}, + {placeholder}, + {placeholder}, + {placeholder}, + {placeholder}{json_cast}, + {placeholder}{json_cast}, + {placeholder}{json_cast}, + {placeholder}{json_cast}, + {placeholder}{json_cast}, + {placeholder}{json_cast}, + {placeholder} + ) + """, + ( + str(manifest_id), + str(article_id), + str(agent_job_id), + s3_prefix, + _json_value(source_urls), + _json_value(object_keys), + _json_value(content_hashes), + _json_value(artifact_types), + _json_value(metadata_references), + _json_value(artifacts), + _datetime_value(created_at), + ), + ) + + return self.get(manifest_id) + + def list_for_article(self, article_id: UUID) -> list[ResearchArtifactManifestSummary]: + placeholder = self._repository.placeholder() + with self._repository.connection() as connection: + rows = connection.execute( + f""" + SELECT + id, + article_id, + agent_job_id, + s3_prefix, + artifacts, + created_at + FROM research_run_manifests + WHERE article_id = {placeholder} + ORDER BY created_at + """, + (str(article_id),), + ).fetchall() + + return [_research_manifest_from_row(row) for row in rows] + + def get(self, manifest_id: UUID) -> ResearchArtifactManifestSummary: + placeholder = self._repository.placeholder() + with self._repository.connection() as connection: + row = connection.execute( + f""" + SELECT + id, + article_id, + agent_job_id, + s3_prefix, + artifacts, + created_at + FROM research_run_manifests + WHERE id = {placeholder} + """, + (str(manifest_id),), + ).fetchone() + + if row is None: + raise LookupError(f"Research manifest not found: {manifest_id}") + return _research_manifest_from_row(row) + + class AgentJobsRepository: def __init__(self, repository: BackendRepository) -> None: self._repository = repository @@ -1715,6 +1833,20 @@ def _plan_summary_from_row( ) +def _research_manifest_from_row(row: Any) -> ResearchArtifactManifestSummary: + return ResearchArtifactManifestSummary( + id=_row_value(row, "id"), + article_id=_row_value(row, "article_id"), + agent_job_id=_row_value(row, "agent_job_id"), + s3_prefix=_row_value(row, "s3_prefix"), + artifacts=[ + ResearchArtifactSummary(**artifact) + for artifact in _json_from_row(row, "artifacts") + ], + created_at=_row_value(row, "created_at"), + ) + + def _agent_job_summary_from_row(row: Any) -> AgentJobSummary: return AgentJobSummary( id=_row_value(row, "id"), diff --git a/apps/backend/src/presentation/routes/plans.py b/apps/backend/src/presentation/routes/plans.py index e69f8e0..39f2afd 100644 --- a/apps/backend/src/presentation/routes/plans.py +++ b/apps/backend/src/presentation/routes/plans.py @@ -11,18 +11,20 @@ from src.application.plans import ( get_plan, list_plans, request_plan_revision, - start_research, ) +from src.application.research import list_research_runs, start_research_run from src.domain.auth import EDITOR_OR_ADMIN_ROLES from src.domain.contracts import ( - AgentJobListResponse, CurrentUser, PlanListResponse, PlanResponse, PlanRevisionRequest, PlanUpdateRequest, + ResearchListResponse, + ResearchStartResponse, ReviewActionResponse, ) +from src.infrastructure.object_storage import open_object_storage_client from src.infrastructure.repositories import BackendRepository from src.presentation.dependencies import get_repository, require_roles @@ -143,17 +145,36 @@ def post_request_plan_revision( @router.post( "/articles/{article_id}/research/start", - response_model=AgentJobListResponse, - status_code=status.HTTP_202_ACCEPTED, + response_model=ResearchStartResponse, + status_code=status.HTTP_201_CREATED, ) def post_start_research( article_id: UUID, _: CurrentUser = Depends(require_roles(EDITOR_OR_ADMIN_ROLES)), repository: BackendRepository = Depends(get_repository), -) -> AgentJobListResponse: +) -> ResearchStartResponse: try: - return start_research(repository, article_id=article_id) + return start_research_run( + repository, + article_id=article_id, + object_storage=open_object_storage_client(), + ) except LookupError as error: raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error except PermissionError as error: raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(error)) from error + + +@router.get( + "/articles/{article_id}/research", + response_model=ResearchListResponse, +) +def get_research( + article_id: UUID, + _: CurrentUser = Depends(require_roles(EDITOR_OR_ADMIN_ROLES)), + repository: BackendRepository = Depends(get_repository), +) -> ResearchListResponse: + try: + return list_research_runs(repository, article_id=article_id) + 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_plan_generation_review_public_api.py b/apps/backend/tests/integration/test_plan_generation_review_public_api.py index 6ca81ee..a381446 100644 --- a/apps/backend/tests/integration/test_plan_generation_review_public_api.py +++ b/apps/backend/tests/integration/test_plan_generation_review_public_api.py @@ -112,8 +112,8 @@ class PlanGenerationReviewPublicApiTest(unittest.TestCase): f"/api/articles/{article_id}/research/start", headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, ) - self.assertEqual(202, allowed_research_response.status_code, allowed_research_response.text) - self.assertEqual("RESEARCH", allowed_research_response.json()["jobs"][0]["job_type"]) + self.assertEqual(201, allowed_research_response.status_code, allowed_research_response.text) + self.assertEqual("RESEARCH", allowed_research_response.json()["job"]["job_type"]) detail_response = self.client.get( f"/api/articles/{article_id}", diff --git a/apps/backend/tests/integration/test_research_manifest_public_api.py b/apps/backend/tests/integration/test_research_manifest_public_api.py new file mode 100644 index 0000000..9718879 --- /dev/null +++ b/apps/backend/tests/integration/test_research_manifest_public_api.py @@ -0,0 +1,150 @@ +from __future__ import annotations + +import os +import sys +import tempfile +import unittest +from pathlib import Path + +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_USER_EMAIL_HEADER = "X-Demo-User-Email" + + +class ResearchManifestPublicApiTest(unittest.TestCase): + def setUp(self) -> None: + self.tmp_dir = tempfile.TemporaryDirectory() + self.storage_root = Path(self.tmp_dir.name) / "object-storage" + os.environ["OBJECT_STORAGE_LOCAL_ROOT"] = str(self.storage_root) + dsn = f"sqlite:///{Path(self.tmp_dir.name) / 'research.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_approved_plan_research_creates_manifest_with_uploaded_artifacts(self) -> None: + article_id = self._create_article_with_approved_plan() + + start_response = self.client.post( + f"/api/articles/{article_id}/research/start", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(201, start_response.status_code, start_response.text) + body = start_response.json() + self.assertEqual("RESEARCH_RUNNING", body["article"]["status"]) + self.assertEqual("RESEARCH", body["job"]["job_type"]) + + manifest = body["manifest"] + self.assertTrue(manifest["s3_prefix"]) + self.assertEqual(article_id, manifest["article_id"]) + self.assertEqual(body["job"]["id"], manifest["agent_job_id"]) + self.assertGreaterEqual(len(manifest["artifacts"]), 2) + + object_keys = [artifact["object_key"] for artifact in manifest["artifacts"]] + content_hashes = [artifact["content_hash"] for artifact in manifest["artifacts"]] + artifact_types = [artifact["artifact_type"] for artifact in manifest["artifacts"]] + source_urls = [artifact["source_url"] for artifact in manifest["artifacts"]] + + self.assertEqual(len(object_keys), len(set(object_keys))) + self.assertTrue(all(key.startswith(manifest["s3_prefix"]) for key in object_keys)) + self.assertTrue(all(content_hashes)) + self.assertIn("source_section_snapshot", artifact_types) + self.assertTrue(all(url.startswith("https://") for url in source_urls)) + + for artifact in manifest["artifacts"]: + metadata = artifact["metadata"] + self.assertTrue(metadata["source_title"]) + self.assertTrue(metadata["domain"]) + self.assertTrue(metadata["summary"]) + self.assertTrue(metadata["retrieved_at"]) + self.assertTrue(metadata["artifact_link"]) + self.assertIn("search_path", metadata) + self.assertIn("selection_criteria", metadata) + self.assertIn("head_metadata", metadata) + self.assertIn("server_ip_address", metadata) + self.assertIn("whois_owner", metadata) + self.assertTrue((self.storage_root / artifact["object_key"]).is_file()) + + list_response = self.client.get( + f"/api/articles/{article_id}/research", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(200, list_response.status_code, list_response.text) + self.assertEqual([manifest["id"]], [item["id"] for item in list_response.json()["manifests"]]) + + rerun_response = self.client.post( + f"/api/articles/{article_id}/research/start", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(201, rerun_response.status_code, rerun_response.text) + self.assertNotEqual(manifest["id"], rerun_response.json()["manifest"]["id"]) + self.assertNotEqual( + manifest["s3_prefix"], + rerun_response.json()["manifest"]["s3_prefix"], + ) + + def _create_article_with_approved_plan(self) -> str: + 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": "Research auditable AI content operations.", + "working_title": "Auditable AI Content Operations", + "content_type": "longform_guide", + "primary_keyword": "auditable AI content", + }, + ).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"]: + 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, response.status_code, response.text) + self.client.post( + f"/api/articles/{article_id}/boundary-questions/submit", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + plan = self.client.post( + f"/api/articles/{article_id}/plan/generate", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ).json()["plan"] + 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 + + +if __name__ == "__main__": + unittest.main() diff --git a/apps/frontend/package.json b/apps/frontend/package.json index 02ff8dc..57758c5 100644 --- a/apps/frontend/package.json +++ b/apps/frontend/package.json @@ -6,7 +6,7 @@ "build": "next build", "start": "next start", "test:roles": "node tests/role-navigation.test.mjs", - "test:ui": "node tests/role-navigation.test.mjs && node tests/article_dashboard.test.mjs && node tests/article_form_validation.test.mjs && node tests/admin_script_versions.model.test.mjs && node tests/agent_jobs.model.test.mjs && node tests/boundary_questions.model.test.mjs && node tests/plan_review.model.test.mjs", + "test:ui": "node tests/role-navigation.test.mjs && node tests/article_dashboard.test.mjs && node tests/article_form_validation.test.mjs && node tests/admin_script_versions.model.test.mjs && node tests/agent_jobs.model.test.mjs && node tests/boundary_questions.model.test.mjs && node tests/plan_review.model.test.mjs && node tests/research_manifest.model.test.mjs", "typecheck": "tsc --noEmit" }, "dependencies": { diff --git a/apps/frontend/src/app/articles/[articleId]/research/page.tsx b/apps/frontend/src/app/articles/[articleId]/research/page.tsx new file mode 100644 index 0000000..ad7825c --- /dev/null +++ b/apps/frontend/src/app/articles/[articleId]/research/page.tsx @@ -0,0 +1,14 @@ +import ArticleResearchPage from "@/pages/article-research"; + +type DynamicArticleResearchPageProps = { + params: { + articleId: string; + }; +}; + +export default async function DynamicArticleResearchPage({ + params, +}: DynamicArticleResearchPageProps) { + const { articleId } = params; + return ; +} diff --git a/apps/frontend/src/features/article-detail/ui.tsx b/apps/frontend/src/features/article-detail/ui.tsx index b8e794a..9340c91 100644 --- a/apps/frontend/src/features/article-detail/ui.tsx +++ b/apps/frontend/src/features/article-detail/ui.tsx @@ -47,6 +47,9 @@ export function ArticleDetailShell({ summary }: DetailShellProps) {
  • Plan review
  • +
  • + Research +
  • Workflow timeline

    diff --git a/apps/frontend/src/features/research-manifest/model.ts b/apps/frontend/src/features/research-manifest/model.ts new file mode 100644 index 0000000..b7d3ab7 --- /dev/null +++ b/apps/frontend/src/features/research-manifest/model.ts @@ -0,0 +1,40 @@ +import type { ResearchArtifactManifestSummary } from "@pipeline/shared"; + +export type ResearchSourceRow = { + manifestId: string; + sourceUrl: string; + title: string; + domain: string; + type: string; + summary: string; + retrievedAt: string; + artifactLink: string; + objectKey: string; + contentHash: string; +}; + +function stringValue(value: unknown): string { + return typeof value === "string" ? value : ""; +} + +export function buildResearchSourceRows( + manifests: readonly ResearchArtifactManifestSummary[], +): ResearchSourceRow[] { + return manifests.flatMap((manifest) => + (manifest.artifacts ?? []).map((artifact) => { + const metadata = artifact.metadata ?? {}; + return { + manifestId: manifest.id, + sourceUrl: artifact.source_url, + title: stringValue(metadata.source_title), + domain: stringValue(metadata.domain), + type: stringValue(metadata.source_type) || artifact.artifact_type, + summary: stringValue(metadata.summary), + retrievedAt: stringValue(metadata.retrieved_at) || manifest.created_at, + artifactLink: stringValue(metadata.artifact_link) || artifact.object_key, + objectKey: artifact.object_key, + contentHash: artifact.content_hash, + }; + }), + ); +} diff --git a/apps/frontend/src/features/research-manifest/ui.tsx b/apps/frontend/src/features/research-manifest/ui.tsx new file mode 100644 index 0000000..e650211 --- /dev/null +++ b/apps/frontend/src/features/research-manifest/ui.tsx @@ -0,0 +1,87 @@ +"use client"; + +import { useState } from "react"; +import { useRouter } from "next/navigation"; + +import type { ResearchArtifactManifestSummary } from "@pipeline/shared"; + +import { ApiError, startResearch } from "@/shared/pipeline-api"; +import { buildResearchSourceRows } from "./model"; + +type ResearchManifestPanelProps = { + articleId: string; + initialManifests: readonly ResearchArtifactManifestSummary[]; +}; + +export function ResearchManifestPanel({ + articleId, + initialManifests, +}: ResearchManifestPanelProps) { + const router = useRouter(); + const [manifests, setManifests] = useState([ + ...initialManifests, + ]); + const [message, setMessage] = useState(""); + const [isBusy, setBusy] = useState(false); + const rows = buildResearchSourceRows(manifests); + + async function handleStart() { + setBusy(true); + setMessage(""); + try { + const response = await startResearch(articleId); + setManifests((current) => [...current, response.manifest]); + router.refresh(); + } catch (error) { + setMessage(error instanceof ApiError ? error.message : "Unable to start research"); + } finally { + setBusy(false); + } + } + + return ( +
    +
    +
    +

    Research

    +

    {manifests.length} runs

    +
    + +
    + {message ?

    Error: {message}

    : null} + {rows.length === 0 ? ( +

    No research sources yet.

    + ) : ( + + + + + + + + + + + {rows.map((row) => ( + + + + + + + ))} + +
    SourceSummaryRetrievedArtifact
    +
    {row.title}
    + {row.domain || row.sourceUrl} +
    {row.type}
    +
    {row.summary}{new Date(row.retrievedAt).toLocaleString()} +
    {row.artifactLink}
    + {row.contentHash} +
    + )} +
    + ); +} diff --git a/apps/frontend/src/pages/article-research/index.tsx b/apps/frontend/src/pages/article-research/index.tsx new file mode 100644 index 0000000..332ee6d --- /dev/null +++ b/apps/frontend/src/pages/article-research/index.tsx @@ -0,0 +1,34 @@ +import Link from "next/link"; + +import { ResearchManifestPanel } from "@/features/research-manifest/ui"; +import { fetchResearch } from "@/shared/pipeline-api"; +import { RoleNavigation } from "@/widgets/role-navigation"; + +type ArticleResearchPageProps = { + articleId: string; +}; + +export default async function ArticleResearchPage({ + articleId, +}: ArticleResearchPageProps) { + const response = await fetchResearch(articleId); + + return ( +
    +
    +
    +

    Research

    +

    {response.manifests?.length ?? 0} runs

    +
    + Back to article +
    + +
    + +
    +
    + ); +} diff --git a/apps/frontend/src/shared/pipeline-api.ts b/apps/frontend/src/shared/pipeline-api.ts index c3ccedf..21c0f7f 100644 --- a/apps/frontend/src/shared/pipeline-api.ts +++ b/apps/frontend/src/shared/pipeline-api.ts @@ -14,6 +14,8 @@ import type { PlanResponse, PlanRevisionRequest, PlanUpdateRequest, + ResearchListResponse, + ResearchStartResponse, ReviewActionResponse, ScriptConfigVersionCreateRequest, ScriptConfigVersionListResponse, @@ -373,3 +375,14 @@ export function requestArticlePlanRevision( request, ); } + +export function startResearch(articleId: string): Promise { + return apiPost( + `/api/articles/${articleId}/research/start`, + undefined, + ); +} + +export function fetchResearch(articleId: string): Promise { + return apiGet(`/api/articles/${articleId}/research`); +} diff --git a/apps/frontend/tests/research_manifest.model.test.mjs b/apps/frontend/tests/research_manifest.model.test.mjs new file mode 100644 index 0000000..06d142a --- /dev/null +++ b/apps/frontend/tests/research_manifest.model.test.mjs @@ -0,0 +1,55 @@ +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/research-manifest/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 { buildResearchSourceRows } = moduleExports; + +const rows = buildResearchSourceRows([ + { + id: "11111111-1111-1111-1111-111111111111", + article_id: "22222222-2222-2222-2222-222222222222", + agent_job_id: "33333333-3333-3333-3333-333333333333", + s3_prefix: "research/article/run", + created_at: "2026-05-21T00:00:00Z", + artifacts: [ + { + artifact_type: "source_section_snapshot", + source_url: "https://example.com/source", + object_key: "research/article/run/source-1/section.json", + content_hash: "abc123", + metadata: { + source_title: "Source title", + domain: "example.com", + source_type: "primary_reference", + summary: "Useful source summary", + retrieved_at: "2026-05-21T00:00:01Z", + artifact_link: "research/article/run/source-1/section.json", + }, + }, + ], + }, +]); + +assert.equal(rows.length, 1); +assert.equal(rows[0].sourceUrl, "https://example.com/source"); +assert.equal(rows[0].title, "Source title"); +assert.equal(rows[0].domain, "example.com"); +assert.equal(rows[0].type, "primary_reference"); +assert.equal(rows[0].summary, "Useful source summary"); +assert.equal(rows[0].artifactLink, "research/article/run/source-1/section.json"); diff --git a/packages/shared/openapi.json b/packages/shared/openapi.json index d6055aa..ca75224 100644 --- a/packages/shared/openapi.json +++ b/packages/shared/openapi.json @@ -1993,6 +1993,48 @@ "title": "ResearchArtifactSummary", "type": "object" }, + "ResearchListResponse": { + "additionalProperties": false, + "properties": { + "evidence": { + "items": { + "$ref": "#/components/schemas/EvidenceSummary" + }, + "title": "Evidence", + "type": "array" + }, + "manifests": { + "items": { + "$ref": "#/components/schemas/ResearchArtifactManifestSummary" + }, + "title": "Manifests", + "type": "array" + } + }, + "title": "ResearchListResponse", + "type": "object" + }, + "ResearchStartResponse": { + "additionalProperties": false, + "properties": { + "article": { + "$ref": "#/components/schemas/ArticleSummary" + }, + "job": { + "$ref": "#/components/schemas/AgentJobSummary" + }, + "manifest": { + "$ref": "#/components/schemas/ResearchArtifactManifestSummary" + } + }, + "required": [ + "article", + "job", + "manifest" + ], + "title": "ResearchStartResponse", + "type": "object" + }, "ReviewActionResponse": { "additionalProperties": false, "properties": { @@ -3951,6 +3993,65 @@ ] } }, + "/api/articles/{article_id}/research": { + "get": { + "operationId": "get_research_api_articles__article_id__research_get", + "parameters": [ + { + "in": "path", + "name": "article_id", + "required": true, + "schema": { + "format": "uuid", + "title": "Article Id", + "type": "string" + } + }, + { + "in": "header", + "name": "X-Demo-User-Email", + "required": false, + "schema": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "X-Demo-User-Email" + } + } + ], + "responses": { + "200": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ResearchListResponse" + } + } + }, + "description": "Successful Response" + }, + "422": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + }, + "description": "Validation Error" + } + }, + "summary": "Get Research", + "tags": [ + "plans" + ] + } + }, "/api/articles/{article_id}/research/start": { "post": { "operationId": "post_start_research_api_articles__article_id__research_start_post", @@ -3983,11 +4084,11 @@ } ], "responses": { - "202": { + "201": { "content": { "application/json": { "schema": { - "$ref": "#/components/schemas/AgentJobListResponse" + "$ref": "#/components/schemas/ResearchStartResponse" } } }, diff --git a/packages/shared/src/api-types.ts b/packages/shared/src/api-types.ts index 8b73020..9e271b6 100644 --- a/packages/shared/src/api-types.ts +++ b/packages/shared/src/api-types.ts @@ -316,6 +316,17 @@ export type ResearchArtifactSummary = { source_url: string; }; +export type ResearchListResponse = { + evidence?: EvidenceSummary[]; + manifests?: ResearchArtifactManifestSummary[]; +}; + +export type ResearchStartResponse = { + article: ArticleSummary; + job: AgentJobSummary; + manifest: ResearchArtifactManifestSummary; +}; + export type ReviewActionResponse = { review: ReviewSummary; }; diff --git a/tasks/010-research-fetch-and-s3-manifest.md b/tasks/010-research-fetch-and-s3-manifest.md index 988cb8f..2a6c194 100644 --- a/tasks/010-research-fetch-and-s3-manifest.md +++ b/tasks/010-research-fetch-and-s3-manifest.md @@ -35,14 +35,14 @@ Development description: Implement explicit auditable web research with source r ## Acceptance Criteria -- [ ] TDD pre-requirement: before implementation, write one failing behavior test that starts research from an approved plan and verifies a manifest with object keys is produced; proceed one research behavior at a time and record evidence in `Result`. -- [ ] Research cannot start before plan approval. -- [ ] Research run creates an agent/job record and article status `RESEARCH_RUNNING`. -- [ ] Source-section artifacts are uploaded to object storage. -- [ ] Manifest records include object keys, hashes, URLs, artifact types, and metadata references. -- [ ] Re-running research creates a new manifest instead of overwriting old artifacts. -- [ ] UI can show source URL, title, domain, type, summary, retrieval timestamp, and artifact link. -- [ ] Demo stack can run with deterministic fake research results. +- [x] TDD pre-requirement: before implementation, write one failing behavior test that starts research from an approved plan and verifies a manifest with object keys is produced; proceed one research behavior at a time and record evidence in `Result`. +- [x] Research cannot start before plan approval. +- [x] Research run creates an agent/job record and article status `RESEARCH_RUNNING`. +- [x] Source-section artifacts are uploaded to object storage. +- [x] Manifest records include object keys, hashes, URLs, artifact types, and metadata references. +- [x] Re-running research creates a new manifest instead of overwriting old artifacts. +- [x] UI can show source URL, title, domain, type, summary, retrieval timestamp, and artifact link. +- [x] Demo stack can run with deterministic fake research results. ## Verification @@ -53,9 +53,9 @@ Development description: Implement explicit auditable web research with source r ## 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 `apps/backend/tests/integration/test_research_manifest_public_api.py` before implementation. The test builds an approved plan, starts research, verifies a manifest with object keys/hashes/source URLs/artifact types/metadata, checks uploaded local object files, lists research runs, and verifies rerun creates a new manifest prefix. +- Red evidence: Initial run failed with `201 != 202`; `POST /api/articles/{article_id}/research/start` only returned a queued `RESEARCH` job and no manifest/artifacts. +- Green evidence: Implemented deterministic fake research fetch, local/S3 object storage adapter, research manifest repository, `ResearchStartResponse`, `ResearchListResponse`, `GET /research`, manifest-linked detail output, and FSD frontend research page at `/articles/[articleId]/research`. +- Refactor notes: Moved research-start behavior out of the Task 009 placeholder and made it produce an auditable manifest immediately for deterministic demo/tests. `OBJECT_STORAGE_LOCAL_ROOT` enables isolated integration tests; Docker Compose uses the existing S3-compatible MinIO env. +- Verification output: `PYTHONPATH=/private/tmp/pupline-backend-deps python3 -m unittest apps/backend/tests/integration/test_research_manifest_public_api.py apps/backend/tests/integration/test_plan_generation_review_public_api.py apps/backend/tests/integration/test_schema_storage_contracts.py apps/backend/tests/contracts/test_generated_contract_artifacts.py` passed. `pnpm --filter @pipeline/frontend test:ui` passed. `pnpm --filter @pipeline/frontend typecheck` passed. `tests/smoke/research-flow.sh` and `tests/smoke/plan-flow.sh` passed with Docker Compose after Docker socket escalation. diff --git a/tests/smoke/plan-flow.sh b/tests/smoke/plan-flow.sh index eb5b605..627bccd 100755 --- a/tests/smoke/plan-flow.sh +++ b/tests/smoke/plan-flow.sh @@ -112,7 +112,7 @@ edited_plan_id="$(printf '%s' "${edited_json}" | json_eval "json.load(sys.stdin) request POST "/api/articles/${article_id}/plans/${edited_plan_id}/request-revision" '{"notes":"Smoke revision"}' >/dev/null request POST "/api/articles/${article_id}/plans/${edited_plan_id}/approve" >/dev/null research_json="$(request POST "/api/articles/${article_id}/research/start")" -job_type="$(printf '%s' "${research_json}" | json_eval "json.load(sys.stdin)['jobs'][0]['job_type']")" +job_type="$(printf '%s' "${research_json}" | json_eval "json.load(sys.stdin)['job']['job_type']")" if [[ "${job_type}" != "RESEARCH" ]]; then echo "Expected RESEARCH job, got ${job_type}" >&2 diff --git a/tests/smoke/research-flow.sh b/tests/smoke/research-flow.sh new file mode 100755 index 0000000..ba792ff --- /dev/null +++ b/tests/smoke/research-flow.sh @@ -0,0 +1,112 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)" +COMPOSE_FILE="${ROOT_DIR}/docker-compose.yml" +BACKEND_URL="${BACKEND_URL:-http://localhost:8000}" +WAIT_SECONDS="${WAIT_SECONDS:-60}" +PYTHON_BIN="${PYTHON_BIN:-python3}" +EDITOR_EMAIL="${EDITOR_EMAIL:-editor@example.com}" + +cleanup() { + docker compose -f "${COMPOSE_FILE}" down --remove-orphans >/dev/null 2>&1 || true +} + +request() { + local method="$1" + local path="$2" + local body="${3:-}" + + if [[ -n "${body}" ]]; then + curl -fsS -X "${method}" -H "Content-Type: application/json" \ + -H "X-Demo-User-Email: ${EDITOR_EMAIL}" --data "${body}" "${BACKEND_URL}${path}" + return + fi + + curl -fsS -X "${method}" -H "X-Demo-User-Email: ${EDITOR_EMAIL}" "${BACKEND_URL}${path}" +} + +wait_backend() { + local deadline=$((SECONDS + WAIT_SECONDS)) + while (( SECONDS < deadline )); do + if curl -fsS "${BACKEND_URL}/health" >/dev/null 2>&1; then + return 0 + fi + sleep 1 + done + echo "Timed out waiting for backend at ${BACKEND_URL}" >&2 + return 1 +} + +json_eval() { + local expression="$1" + "${PYTHON_BIN}" -c "import json,sys; print(${expression})" +} + +cd "${ROOT_DIR}" +trap cleanup EXIT + +docker compose -f "${COMPOSE_FILE}" up --build -d +wait_backend + +site_id="$(request GET /api/sites | json_eval "json.load(sys.stdin)[0]['site']['id']")" +article_body="$( + SITE_ID="${site_id}" "${PYTHON_BIN}" - <<'PY' +import json +import os +import uuid + +print(json.dumps({ + "target_site_id": os.environ["SITE_ID"], + "brief_description": f"Research smoke flow {uuid.uuid4()}", + "working_title": "Research Smoke Flow", + "content_type": "longform_guide", + "primary_keyword": "research smoke flow", +})) +PY +)" +article_id="$(request POST /api/articles "${article_body}" | json_eval "json.load(sys.stdin)['article']['id']")" + +questions_json="$(request POST "/api/articles/${article_id}/boundary-questions/generate")" +question_patch_commands="$( + QUESTIONS_JSON="${questions_json}" "${PYTHON_BIN}" - <<'PY' +import json +import os + +for question in json.loads(os.environ["QUESTIONS_JSON"])["questions"]: + if question["is_required"]: + print(f"{question['id']}\tAccepted answer for {question['category']}.") +PY +)" +while IFS=$'\t' read -r question_id answer; do + [[ -z "${question_id}" ]] && continue + patch_body="$(ANSWER="${answer}" "${PYTHON_BIN}" - <<'PY' +import json +import os + +print(json.dumps({"answer": os.environ["ANSWER"]})) +PY +)" + request PATCH "/api/articles/${article_id}/boundary-questions/${question_id}" "${patch_body}" >/dev/null +done <<< "${question_patch_commands}" +request POST "/api/articles/${article_id}/boundary-questions/submit" >/dev/null + +plan_id="$(request POST "/api/articles/${article_id}/plan/generate" | json_eval "json.load(sys.stdin)['plan']['id']")" +request POST "/api/articles/${article_id}/plans/${plan_id}/approve" >/dev/null + +research_json="$(request POST "/api/articles/${article_id}/research/start")" +manifest_id="$(printf '%s' "${research_json}" | json_eval "json.load(sys.stdin)['manifest']['id']")" +artifact_count="$(printf '%s' "${research_json}" | json_eval "len(json.load(sys.stdin)['manifest']['artifacts'])")" +if (( artifact_count < 2 )); then + echo "Expected at least 2 research artifacts, got ${artifact_count}" >&2 + exit 1 +fi + +research_list_json="$(request GET "/api/articles/${article_id}/research")" +listed_id="$(printf '%s' "${research_list_json}" | json_eval "json.load(sys.stdin)['manifests'][0]['id']")" +if [[ "${listed_id}" != "${manifest_id}" ]]; then + echo "Expected manifest ${manifest_id}, got ${listed_id}" >&2 + exit 1 +fi + +echo "research flow smoke ok"