feat(task-012): orchestrate parallel section production jobs

This commit is contained in:
2026-05-21 23:10:34 +03:00
parent 7e1b05067c
commit 8ebb5ad623
15 changed files with 756 additions and 39 deletions
+49 -11
View File
@@ -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 []
+2
View File
@@ -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,
)
@@ -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",
+103 -12
View File
@@ -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:
+1
View File
@@ -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",
@@ -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 = ""
@@ -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"),
@@ -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 '',