Task 010 implement research manifest flow
This commit is contained in:
@@ -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,
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
@@ -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"),
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user