Files
content-factory/apps/backend/src/infrastructure/repositories.py
T

3303 lines
109 KiB
Python

from __future__ import annotations
import json
import sqlite3
from collections.abc import Iterator
from contextlib import contextmanager
from datetime import datetime
from typing import Any
from uuid import UUID, uuid4
from src.domain.contracts import (
AgentJobErrorCategory,
AgentJobStatus,
AgentJobSummary,
AgentJobType,
ContentReviewKind,
ContentReviewIssueSummary,
ContentReviewReportSummary,
ContentReviewSuggestionSummary,
AssetRevisionSummary,
AssetStatus,
AssetSummary,
AssetType,
ArticleSummary,
BoundaryQuestionSummary,
ClaimRiskLevel,
ClaimSupportStatus,
ClaimSummary,
DraftFaqItem,
DraftSummary,
EvidenceSummary,
PlanReviewStatus,
PlanSectionSummary,
PlanSummary,
PublishingRules,
PublishingStatus,
ResearchArtifactManifestSummary,
ResearchArtifactSummary,
ArticleWorkflowStatus,
ReviewSuggestionStatus,
WorkflowEventSummary,
Role,
ScriptConfigVersionStatus,
TargetSiteConfig,
UserSummary,
)
from src.infrastructure.schema import setup_database
JsonObject = dict[str, Any]
def open_backend_repository(dsn: str) -> "BackendRepository":
return BackendRepository(dsn)
class BackendRepository:
def __init__(self, dsn: str) -> None:
self.dsn = dsn
self.dialect = "sqlite" if dsn.startswith("sqlite:///") else "postgres"
self.users = UsersRepository(self)
self.target_sites = TargetSitesRepository(self)
self.articles = ArticlesRepository(self)
self.boundary_questions = BoundaryQuestionsRepository(self)
self.article_plans = ArticlePlansRepository(self)
self.article_drafts = ArticleDraftsRepository(self)
self.content_reviews = ContentReviewsRepository(self)
self.assets = AssetsRepository(self)
self.research_manifests = ResearchManifestsRepository(self)
self.evidence_items = EvidenceItemsRepository(self)
self.claims = ClaimsRepository(self)
self.agent_jobs = AgentJobsRepository(self)
self.script_config_versions = ScriptConfigVersionsRepository(self)
self.script_config_version_events = ScriptConfigVersionAuditEventsRepository(self)
self.schema = SchemaRepository(self)
def setup(self) -> None:
setup_database(self.dsn)
@contextmanager
def connection(self) -> Iterator[Any]:
connection = self._connect()
try:
with connection:
yield connection
finally:
connection.close()
def _connect(self) -> Any:
if self.dialect == "sqlite":
sqlite_path = self.dsn.removeprefix("sqlite:///")
connection = sqlite3.connect(sqlite_path)
connection.row_factory = sqlite3.Row
return connection
import psycopg
from psycopg.rows import dict_row
return psycopg.connect(self.dsn, row_factory=dict_row)
def placeholder(self) -> str:
if self.dialect == "sqlite":
return "?"
return "%s"
def json_cast(self) -> str:
if self.dialect == "sqlite":
return ""
return "::jsonb"
class UsersRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def upsert(
self,
*,
user_id: UUID,
email: str,
display_name: str,
role: Role,
created_at: datetime,
updated_at: datetime,
) -> UserSummary:
placeholder = self._repository.placeholder()
sql = f"""
INSERT INTO users (
id, email, display_name, role, created_at, updated_at
)
VALUES (
{placeholder}, {placeholder}, {placeholder}, {placeholder},
{placeholder}, {placeholder}
)
ON CONFLICT (email) DO UPDATE SET
display_name = excluded.display_name,
role = excluded.role,
updated_at = excluded.updated_at
"""
params = (
str(user_id),
email,
display_name,
role.value,
_datetime_value(created_at),
_datetime_value(updated_at),
)
with self._repository.connection() as connection:
connection.execute(sql, params)
return self.get_by_email(email)
def get_by_email(self, email: str) -> UserSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT id, display_name, role
FROM users
WHERE email = {placeholder}
""",
(email,),
).fetchone()
if row is None:
raise LookupError(f"User not found: {email}")
return _user_from_row(row)
def list(self) -> list[UserSummary]:
with self._repository.connection() as connection:
rows = connection.execute(
"""
SELECT id, display_name, role
FROM users
ORDER BY email
"""
).fetchall()
return [_user_from_row(row) for row in rows]
def count_by_email(self, email: str) -> int:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"SELECT COUNT(*) AS count FROM users WHERE email = {placeholder}",
(email,),
).fetchone()
return int(_row_value(row, "count"))
class TargetSitesRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def upsert(
self,
*,
site_id: UUID,
name: str,
slug: str,
publishing_type: str,
default_language: str,
brand_voice: str,
audience: str,
seo_rules: JsonObject,
visual_rules: JsonObject,
source_rules: JsonObject,
publishing_rules: PublishingRules,
active_script_config_version_id: UUID | None,
created_at: datetime,
updated_at: datetime,
) -> TargetSiteConfig:
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
sql = f"""
INSERT INTO target_sites (
id,
name,
slug,
publishing_type,
default_language,
brand_voice,
audience,
seo_rules,
visual_rules,
source_rules,
publishing_rules,
active_script_config_version_id,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder}
)
ON CONFLICT (slug) DO UPDATE SET
name = excluded.name,
publishing_type = excluded.publishing_type,
default_language = excluded.default_language,
brand_voice = excluded.brand_voice,
audience = excluded.audience,
seo_rules = excluded.seo_rules,
visual_rules = excluded.visual_rules,
source_rules = excluded.source_rules,
publishing_rules = excluded.publishing_rules,
active_script_config_version_id = excluded.active_script_config_version_id,
updated_at = excluded.updated_at
"""
params = (
str(site_id),
name,
slug,
publishing_type,
default_language,
brand_voice,
audience,
_json_value(seo_rules),
_json_value(visual_rules),
_json_value(source_rules),
_json_value(publishing_rules.model_dump(mode="json")),
_uuid_value(active_script_config_version_id),
_datetime_value(created_at),
_datetime_value(updated_at),
)
with self._repository.connection() as connection:
connection.execute(sql, params)
return self.get_by_slug(slug)
def get_by_slug(self, slug: str) -> TargetSiteConfig:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
name,
slug,
publishing_type,
default_language,
brand_voice,
audience,
seo_rules,
visual_rules,
source_rules,
publishing_rules,
active_script_config_version_id,
created_at,
updated_at
FROM target_sites
WHERE slug = {placeholder}
""",
(slug,),
).fetchone()
if row is None:
raise LookupError(f"Target site not found: {slug}")
return _target_site_from_row(row)
def get_by_id(self, site_id: UUID) -> TargetSiteConfig:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
name,
slug,
publishing_type,
default_language,
brand_voice,
audience,
seo_rules,
visual_rules,
source_rules,
publishing_rules,
active_script_config_version_id,
created_at,
updated_at
FROM target_sites
WHERE id = {placeholder}
""",
(str(site_id),),
).fetchone()
if row is None:
raise LookupError(f"Target site not found: {site_id}")
return _target_site_from_row(row)
def update(
self,
*,
site_id: UUID,
name: str,
slug: str,
publishing_type: str,
default_language: str,
brand_voice: str,
audience: str,
seo_rules: JsonObject,
visual_rules: JsonObject,
source_rules: JsonObject,
publishing_rules: PublishingRules,
active_script_config_version_id: UUID | None,
updated_at: datetime,
) -> TargetSiteConfig:
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE target_sites
SET
name = {placeholder},
slug = {placeholder},
publishing_type = {placeholder},
default_language = {placeholder},
brand_voice = {placeholder},
audience = {placeholder},
seo_rules = {placeholder}{json_cast},
visual_rules = {placeholder}{json_cast},
source_rules = {placeholder}{json_cast},
publishing_rules = {placeholder}{json_cast},
active_script_config_version_id = {placeholder},
updated_at = {placeholder}
WHERE id = {placeholder}
""",
(
name,
slug,
publishing_type,
default_language,
brand_voice,
audience,
_json_value(seo_rules),
_json_value(visual_rules),
_json_value(source_rules),
_json_value(publishing_rules.model_dump(mode="json")),
_uuid_value(active_script_config_version_id),
_datetime_value(updated_at),
str(site_id),
),
)
return self.get_by_id(site_id)
def list(self) -> list[TargetSiteConfig]:
with self._repository.connection() as connection:
rows = connection.execute(
"""
SELECT
id,
name,
slug,
publishing_type,
default_language,
brand_voice,
audience,
seo_rules,
visual_rules,
source_rules,
publishing_rules,
active_script_config_version_id,
created_at,
updated_at
FROM target_sites
ORDER BY slug
"""
).fetchall()
return [_target_site_from_row(row) for row in rows]
class ArticlesRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
target_site_id: UUID,
status: ArticleWorkflowStatus,
publishing_status: PublishingStatus,
brief_description: str,
working_title: str | None,
language: str,
content_type: str,
primary_keyword: str | None,
assigned_editor_id: UUID | None,
created_at: datetime,
updated_at: datetime,
) -> ArticleSummary:
article_id = uuid4()
placeholder = self._repository.placeholder()
sql = f"""
INSERT INTO articles (
id,
target_site_id,
status,
publishing_status,
brief_description,
working_title,
language,
content_type,
primary_keyword,
assigned_editor_id,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
"""
with self._repository.connection() as connection:
connection.execute(
sql,
(
str(article_id),
str(target_site_id),
status.value,
publishing_status.value,
brief_description,
working_title,
language,
content_type,
primary_keyword,
_uuid_value(assigned_editor_id),
_datetime_value(created_at),
_datetime_value(updated_at),
),
)
return self.get(article_id)
def list(self) -> list[ArticleSummary]:
with self._repository.connection() as connection:
rows = connection.execute(
"""
SELECT
id,
target_site_id,
status,
publishing_status,
brief_description,
working_title,
language,
content_type,
primary_keyword,
assigned_editor_id,
created_at,
updated_at
FROM articles
ORDER BY updated_at DESC
"""
).fetchall()
return [_article_summary_from_row(row) for row in rows]
def get(self, article_id: UUID) -> ArticleSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
target_site_id,
status,
publishing_status,
brief_description,
working_title,
language,
content_type,
primary_keyword,
assigned_editor_id,
created_at,
updated_at
FROM articles
WHERE id = {placeholder}
""",
(str(article_id),),
).fetchone()
if row is None:
raise LookupError(f"Article not found: {article_id}")
return _article_summary_from_row(row)
def update_status(
self,
*,
article_id: UUID,
status: ArticleWorkflowStatus,
updated_at: datetime,
) -> ArticleSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE articles
SET status = {placeholder}, updated_at = {placeholder}
WHERE id = {placeholder}
""",
(status.value, _datetime_value(updated_at), str(article_id)),
)
return self.get(article_id)
def create_workflow_event(
self,
*,
article_id: UUID,
event_type: str,
from_status: ArticleWorkflowStatus | None,
to_status: ArticleWorkflowStatus | None,
actor_user_id: UUID | None,
payload: dict[str, Any],
created_at: datetime,
) -> None:
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
sql = f"""
INSERT INTO workflow_events (
id,
article_id,
event_type,
from_status,
to_status,
actor_user_id,
payload,
created_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}
)
"""
with self._repository.connection() as connection:
connection.execute(
sql,
(
str(uuid4()),
str(article_id),
event_type,
_article_status_value(from_status),
_article_status_value(to_status),
_uuid_value(actor_user_id),
_json_value(payload),
_datetime_value(created_at),
),
)
def list_workflow_events(self, article_id: UUID) -> list[WorkflowEventSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
event_type,
from_status,
to_status,
actor_user_id,
payload,
created_at
FROM workflow_events
WHERE article_id = {placeholder}
ORDER BY created_at
""",
(str(article_id),),
).fetchall()
return [_workflow_event_from_row(row) for row in rows]
class BoundaryQuestionsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def replace_for_article(
self,
*,
article_id: UUID,
questions: list[JsonObject],
created_at: datetime,
) -> list[BoundaryQuestionSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"DELETE FROM boundary_questions WHERE article_id = {placeholder}",
(str(article_id),),
)
for question in questions:
connection.execute(
f"""
INSERT INTO boundary_questions (
id,
article_id,
sort_order,
category,
question,
answer,
is_required,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
""",
(
str(uuid4()),
str(article_id),
question["sort_order"],
question["category"],
question["question"],
question.get("answer"),
_bool_value(question.get("is_required", True)),
_datetime_value(created_at),
_datetime_value(created_at),
),
)
return self.list_for_article(article_id)
def list_for_article(self, article_id: UUID) -> list[BoundaryQuestionSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
sort_order,
category,
question,
answer,
is_required,
created_at,
updated_at
FROM boundary_questions
WHERE article_id = {placeholder}
ORDER BY sort_order
""",
(str(article_id),),
).fetchall()
return [_boundary_question_from_row(row) for row in rows]
def get(self, *, article_id: UUID, question_id: UUID) -> BoundaryQuestionSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
article_id,
sort_order,
category,
question,
answer,
is_required,
created_at,
updated_at
FROM boundary_questions
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(str(article_id), str(question_id)),
).fetchone()
if row is None:
raise LookupError(f"Boundary question not found: {question_id}")
return _boundary_question_from_row(row)
def update_answer(
self,
*,
article_id: UUID,
question_id: UUID,
answer: str | None,
updated_at: datetime,
) -> BoundaryQuestionSummary:
self.get(article_id=article_id, question_id=question_id)
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE boundary_questions
SET answer = {placeholder}, updated_at = {placeholder}
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(
answer,
_datetime_value(updated_at),
str(article_id),
str(question_id),
),
)
return self.get(article_id=article_id, question_id=question_id)
class ArticlePlansRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create_version(
self,
*,
article_id: UUID,
version: int,
status: PlanReviewStatus,
title_options: list[str],
recommended_title: str | None,
reader_persona: str | None,
search_intent: str | None,
thesis: str | None,
claims_to_prove: list[str],
evidence_needs: list[str],
visual_needs: list[str],
seo_notes: list[str],
source_requirements: list[str],
excluded_sources: list[str],
tone: str | None,
audience: str | None,
risks: list[str],
sections: list[JsonObject],
created_at: datetime,
) -> PlanSummary:
plan_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO article_plans (
id,
article_id,
version,
status,
title_options,
recommended_title,
reader_persona,
search_intent,
thesis,
claims_to_prove,
evidence_needs,
visual_needs,
seo_notes,
source_requirements,
excluded_sources,
tone,
audience,
risks,
created_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}
)
""",
(
str(plan_id),
str(article_id),
version,
status.value,
_json_value(title_options),
recommended_title,
reader_persona,
search_intent,
thesis,
_json_value(claims_to_prove),
_json_value(evidence_needs),
_json_value(visual_needs),
_json_value(seo_notes),
_json_value(source_requirements),
_json_value(excluded_sources),
tone,
audience,
_json_value(risks),
_datetime_value(created_at),
),
)
for index, section in enumerate(sections, start=1):
connection.execute(
f"""
INSERT INTO plan_sections (
id,
article_plan_id,
sort_order,
heading,
purpose,
key_points,
evidence_needs,
claims_to_support,
target_word_count
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}
)
""",
(
str(uuid4()),
str(plan_id),
index,
section["heading"],
section.get("purpose"),
_json_value(section.get("key_points", [])),
_json_value(section.get("evidence_needs", [])),
_json_value(section.get("claims_to_support", [])),
section.get("target_word_count"),
),
)
return self.get(article_id=article_id, plan_id=plan_id)
def list_for_article(self, article_id: UUID) -> list[PlanSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT {self._select_columns()}
FROM article_plans
WHERE article_id = {placeholder}
ORDER BY version
""",
(str(article_id),),
).fetchall()
return [self._with_sections(row) for row in rows]
def get(self, *, article_id: UUID, plan_id: UUID) -> PlanSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT {self._select_columns()}
FROM article_plans
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(str(article_id), str(plan_id)),
).fetchone()
if row is None:
raise LookupError(f"Plan not found: {plan_id}")
return self._with_sections(row)
def latest_version(self, article_id: UUID) -> int:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT COALESCE(MAX(version), 0) AS version
FROM article_plans
WHERE article_id = {placeholder}
""",
(str(article_id),),
).fetchone()
return int(_row_value(row, "version"))
def update_status(
self,
*,
article_id: UUID,
plan_id: UUID,
status: PlanReviewStatus,
) -> PlanSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE article_plans
SET status = {placeholder}
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(status.value, str(article_id), str(plan_id)),
)
return self.get(article_id=article_id, plan_id=plan_id)
def _with_sections(self, row: Any) -> PlanSummary:
plan_id = _row_value(row, "id")
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
section_rows = connection.execute(
f"""
SELECT
id,
article_plan_id,
sort_order,
heading,
purpose,
key_points,
evidence_needs,
claims_to_support,
target_word_count
FROM plan_sections
WHERE article_plan_id = {placeholder}
ORDER BY sort_order
""",
(str(plan_id),),
).fetchall()
return _plan_summary_from_row(row, [_plan_section_from_row(item) for item in section_rows])
def _select_columns(self) -> str:
return """
id,
article_id,
version,
status,
title_options,
recommended_title,
reader_persona,
search_intent,
thesis,
claims_to_prove,
evidence_needs,
visual_needs,
seo_notes,
source_requirements,
excluded_sources,
tone,
audience,
risks,
created_at
"""
class ArticleDraftsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create_version(
self,
*,
article_id: UUID,
version: int,
title: str,
slug: str,
meta_title: str | None,
meta_description: str | None,
body_object_key: str | None,
body_markdown: str,
faq_items: list[JsonObject],
visual_placeholders: list[str],
evidence_references: list[UUID],
unsupported_claim_warnings: list[str],
based_on_draft_id: UUID | None,
status: ArticleWorkflowStatus,
created_at: datetime,
updated_at: datetime,
) -> DraftSummary:
draft_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO article_drafts (
id,
article_id,
version,
title,
slug,
meta_title,
meta_description,
body_object_key,
body_markdown,
faq_items,
visual_placeholders,
evidence_references,
unsupported_claim_warnings,
based_on_draft_id,
status,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
""",
(
str(draft_id),
str(article_id),
version,
title,
slug,
meta_title,
meta_description,
body_object_key,
body_markdown,
_json_value(faq_items),
_json_value(visual_placeholders),
_json_value([str(item_id) for item_id in evidence_references]),
_json_value(unsupported_claim_warnings),
_uuid_value(based_on_draft_id),
status.value,
_datetime_value(created_at),
_datetime_value(updated_at),
),
)
return self.get(article_id=article_id, draft_id=draft_id)
def list_for_article(self, article_id: UUID) -> list[DraftSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT {self._select_columns()}
FROM article_drafts
WHERE article_id = {placeholder}
ORDER BY version DESC
""",
(str(article_id),),
).fetchall()
return [_draft_summary_from_row(row) for row in rows]
def get(self, *, article_id: UUID, draft_id: UUID) -> DraftSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT {self._select_columns()}
FROM article_drafts
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(str(article_id), str(draft_id)),
).fetchone()
if row is None:
raise LookupError(f"Draft not found: {draft_id}")
return _draft_summary_from_row(row)
def latest_for_article(self, article_id: UUID) -> DraftSummary | None:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT {self._select_columns()}
FROM article_drafts
WHERE article_id = {placeholder}
ORDER BY version DESC
LIMIT 1
""",
(str(article_id),),
).fetchone()
if row is None:
return None
return _draft_summary_from_row(row)
def latest_version(self, article_id: UUID) -> int:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT COALESCE(MAX(version), 0) AS version
FROM article_drafts
WHERE article_id = {placeholder}
""",
(str(article_id),),
).fetchone()
return int(_row_value(row, "version"))
def _select_columns(self) -> str:
return """
id,
article_id,
version,
title,
slug,
meta_title,
meta_description,
body_object_key,
body_markdown,
faq_items,
visual_placeholders,
evidence_references,
unsupported_claim_warnings,
based_on_draft_id,
status,
created_at,
updated_at
"""
class ContentReviewsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create_report(
self,
*,
article_id: UUID,
review_kind: ContentReviewKind,
draft_id: UUID,
score: int,
recommended_slug: str,
recommended_title: str,
schema_json: JsonObject,
rules_snapshot: JsonObject,
created_at: datetime,
) -> ContentReviewReportSummary:
report_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO article_review_reports (
id,
article_id,
review_kind,
draft_id,
score,
issues,
recommended_slug,
recommended_title,
schema_json,
rules_snapshot,
created_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}
)
""",
(
str(report_id),
str(article_id),
review_kind.value,
str(draft_id),
score,
_json_value([]),
recommended_slug,
recommended_title,
_json_value(schema_json),
_json_value(rules_snapshot),
_datetime_value(created_at),
),
)
return self.get_report(article_id=article_id, report_id=report_id)
def list_reports_for_article(
self,
*,
article_id: UUID,
review_kind: ContentReviewKind,
) -> list[ContentReviewReportSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
review_kind,
draft_id,
score,
recommended_slug,
recommended_title,
schema_json,
rules_snapshot,
created_at
FROM article_review_reports
WHERE article_id = {placeholder}
AND review_kind = {placeholder}
ORDER BY created_at DESC
""",
(str(article_id), review_kind.value),
).fetchall()
return [self.get_report(article_id=article_id, report_id=_row_value(row, "id")) for row in rows]
def latest_report(
self,
*,
article_id: UUID,
review_kind: ContentReviewKind,
) -> ContentReviewReportSummary | None:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT id
FROM article_review_reports
WHERE article_id = {placeholder}
AND review_kind = {placeholder}
ORDER BY created_at DESC
LIMIT 1
""",
(str(article_id), review_kind.value),
).fetchone()
if row is None:
return None
return self.get_report(article_id=article_id, report_id=_row_value(row, "id"))
def get_report(self, *, article_id: UUID, report_id: UUID) -> ContentReviewReportSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
article_id,
review_kind,
draft_id,
score,
recommended_slug,
recommended_title,
schema_json,
rules_snapshot,
created_at
FROM article_review_reports
WHERE article_id = {placeholder}
AND id = {placeholder}
""",
(str(article_id), str(report_id)),
).fetchone()
if row is None:
raise LookupError(f"Content review report not found: {report_id}")
suggestions = self.list_suggestions_for_report(
article_id=article_id,
report_id=report_id,
)
issues = [
ContentReviewIssueSummary(
id=item.suggestion_key,
suggestion_id=item.id,
severity=item.severity,
location=item.location,
message=item.message,
suggested_fix=item.suggested_fix,
suggested_rewrite=item.suggested_rewrite,
status=item.status,
)
for item in suggestions
]
return _content_review_report_from_row(row, issues=issues)
def upsert_suggestion(
self,
*,
suggestion_id: UUID,
article_id: UUID,
review_kind: ContentReviewKind,
report_id: UUID,
suggestion_key: str,
severity: str,
location: str,
message: str,
suggested_fix: str | None,
suggested_rewrite: str | None,
patch: JsonObject,
created_at: datetime,
updated_at: datetime,
) -> ContentReviewSuggestionSummary:
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO article_review_suggestions (
id,
article_id,
review_kind,
report_id,
suggestion_key,
severity,
location,
message,
suggested_fix,
suggested_rewrite,
patch,
status,
applied_text,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
ON CONFLICT (article_id, review_kind, suggestion_key) DO UPDATE SET
report_id = excluded.report_id,
severity = excluded.severity,
location = excluded.location,
message = excluded.message,
suggested_fix = excluded.suggested_fix,
suggested_rewrite = excluded.suggested_rewrite,
patch = excluded.patch,
status = CASE
WHEN article_review_suggestions.status = 'PENDING' THEN excluded.status
ELSE article_review_suggestions.status
END,
applied_text = CASE
WHEN article_review_suggestions.status = 'PENDING' THEN excluded.applied_text
ELSE article_review_suggestions.applied_text
END,
updated_at = excluded.updated_at
""",
(
str(suggestion_id),
str(article_id),
review_kind.value,
str(report_id),
suggestion_key,
severity,
location,
message,
suggested_fix,
suggested_rewrite,
_json_value(patch),
ReviewSuggestionStatus.PENDING.value,
None,
_datetime_value(created_at),
_datetime_value(updated_at),
),
)
return self.get_suggestion(
article_id=article_id,
review_kind=review_kind,
suggestion_id=suggestion_id,
)
def get_suggestion(
self,
*,
article_id: UUID,
review_kind: ContentReviewKind,
suggestion_id: UUID,
) -> ContentReviewSuggestionSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
article_id,
review_kind,
report_id,
suggestion_key,
severity,
location,
message,
suggested_fix,
suggested_rewrite,
patch,
status,
applied_text,
created_at,
updated_at
FROM article_review_suggestions
WHERE article_id = {placeholder}
AND review_kind = {placeholder}
AND id = {placeholder}
""",
(
str(article_id),
review_kind.value,
str(suggestion_id),
),
).fetchone()
if row is None:
raise LookupError(f"Content review suggestion not found: {suggestion_id}")
return _content_review_suggestion_from_row(row)
def list_suggestions_for_report(
self,
*,
article_id: UUID,
report_id: UUID,
) -> list[ContentReviewSuggestionSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
review_kind,
report_id,
suggestion_key,
severity,
location,
message,
suggested_fix,
suggested_rewrite,
patch,
status,
applied_text,
created_at,
updated_at
FROM article_review_suggestions
WHERE article_id = {placeholder}
AND report_id = {placeholder}
ORDER BY suggestion_key
""",
(str(article_id), str(report_id)),
).fetchall()
return [_content_review_suggestion_from_row(row) for row in rows]
def list_unresolved_for_article(
self,
*,
article_id: UUID,
) -> list[ContentReviewSuggestionSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
review_kind,
report_id,
suggestion_key,
severity,
location,
message,
suggested_fix,
suggested_rewrite,
patch,
status,
applied_text,
created_at,
updated_at
FROM article_review_suggestions
WHERE article_id = {placeholder}
AND status = {placeholder}
ORDER BY review_kind, suggestion_key
""",
(str(article_id), ReviewSuggestionStatus.PENDING.value),
).fetchall()
return [_content_review_suggestion_from_row(row) for row in rows]
def update_suggestion(
self,
*,
article_id: UUID,
review_kind: ContentReviewKind,
suggestion_id: UUID,
status: ReviewSuggestionStatus,
applied_text: str | None,
updated_at: datetime,
) -> ContentReviewSuggestionSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
cursor = connection.execute(
f"""
UPDATE article_review_suggestions
SET
status = {placeholder},
applied_text = {placeholder},
updated_at = {placeholder}
WHERE article_id = {placeholder}
AND review_kind = {placeholder}
AND id = {placeholder}
""",
(
status.value,
applied_text,
_datetime_value(updated_at),
str(article_id),
review_kind.value,
str(suggestion_id),
),
)
if int(getattr(cursor, "rowcount", 0)) <= 0:
raise LookupError(f"Content review suggestion not found: {suggestion_id}")
return self.get_suggestion(
article_id=article_id,
review_kind=review_kind,
suggestion_id=suggestion_id,
)
class AssetsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
article_id: UUID,
section_id: UUID | None,
asset_type: AssetType,
title: str,
prompt: str | None,
object_key: str | None,
file_url: str | None,
alt_text: str | None,
caption: str | None,
status: AssetStatus,
created_at: datetime,
updated_at: datetime,
) -> AssetSummary:
asset_id = uuid4()
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO assets (
id,
article_id,
section_id,
asset_type,
title,
prompt,
object_key,
file_url,
alt_text,
caption,
status,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
""",
(
str(asset_id),
str(article_id),
_uuid_value(section_id),
asset_type.value,
title,
prompt,
object_key,
file_url,
alt_text,
caption,
status.value,
_datetime_value(created_at),
_datetime_value(updated_at),
),
)
return self.get(article_id=article_id, asset_id=asset_id)
def list_for_article(self, article_id: UUID) -> list[AssetSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
section_id,
asset_type,
title,
prompt,
object_key,
file_url,
alt_text,
caption,
status,
created_at,
updated_at
FROM assets
WHERE article_id = {placeholder}
ORDER BY created_at, id
""",
(str(article_id),),
).fetchall()
history_by_asset = self._history_by_asset(article_id)
assets: list[AssetSummary] = []
for row in rows:
asset = _asset_summary_from_row(
row,
history=history_by_asset.get(str(_row_value(row, "id")), []),
)
assets.append(asset)
return assets
def get(self, *, article_id: UUID, asset_id: UUID) -> AssetSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
article_id,
section_id,
asset_type,
title,
prompt,
object_key,
file_url,
alt_text,
caption,
status,
created_at,
updated_at
FROM assets
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(str(article_id), str(asset_id)),
).fetchone()
if row is None:
raise LookupError(f"Asset not found: {asset_id}")
history_by_asset = self._history_by_asset(article_id)
return _asset_summary_from_row(
row,
history=history_by_asset.get(str(asset_id), []),
)
def update(
self,
*,
article_id: UUID,
asset_id: UUID,
section_id: UUID | None,
title: str,
prompt: str | None,
object_key: str | None,
file_url: str | None,
alt_text: str | None,
caption: str | None,
status: AssetStatus,
updated_at: datetime,
) -> AssetSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE assets
SET
section_id = {placeholder},
title = {placeholder},
prompt = {placeholder},
object_key = {placeholder},
file_url = {placeholder},
alt_text = {placeholder},
caption = {placeholder},
status = {placeholder},
updated_at = {placeholder}
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(
_uuid_value(section_id),
title,
prompt,
object_key,
file_url,
alt_text,
caption,
status.value,
_datetime_value(updated_at),
str(article_id),
str(asset_id),
),
)
return self.get(article_id=article_id, asset_id=asset_id)
def next_revision_index(self, *, article_id: UUID, asset_id: UUID) -> int:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT COALESCE(MAX(revision_index), 0) + 1 AS next_revision
FROM asset_revisions
WHERE article_id = {placeholder} AND asset_id = {placeholder}
""",
(str(article_id), str(asset_id)),
).fetchone()
return int(_row_value(row, "next_revision"))
def create_revision(
self,
*,
article_id: UUID,
asset_id: UUID,
revision_index: int,
action: str,
actor_user_id: UUID | None,
payload: JsonObject,
created_at: datetime,
) -> AssetRevisionSummary:
revision_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO asset_revisions (
id,
article_id,
asset_id,
revision_index,
action,
actor_user_id,
payload,
created_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}
)
""",
(
str(revision_id),
str(article_id),
str(asset_id),
revision_index,
action,
_uuid_value(actor_user_id),
_json_value(payload),
_datetime_value(created_at),
),
)
return self.get_revision(
article_id=article_id,
asset_id=asset_id,
revision_id=revision_id,
)
def get_revision(
self,
*,
article_id: UUID,
asset_id: UUID,
revision_id: UUID,
) -> AssetRevisionSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
article_id,
asset_id,
revision_index,
action,
actor_user_id,
payload,
created_at
FROM asset_revisions
WHERE article_id = {placeholder}
AND asset_id = {placeholder}
AND id = {placeholder}
""",
(str(article_id), str(asset_id), str(revision_id)),
).fetchone()
if row is None:
raise LookupError(f"Asset revision not found: {revision_id}")
return _asset_revision_from_row(row)
def _history_by_asset(self, article_id: UUID) -> dict[str, list[AssetRevisionSummary]]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
asset_id,
revision_index,
action,
actor_user_id,
payload,
created_at
FROM asset_revisions
WHERE article_id = {placeholder}
ORDER BY revision_index, created_at, id
""",
(str(article_id),),
).fetchall()
history: dict[str, list[AssetRevisionSummary]] = {}
for row in rows:
asset_id = str(_row_value(row, "asset_id"))
history.setdefault(asset_id, []).append(_asset_revision_from_row(row))
return history
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 EvidenceItemsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
article_id: UUID,
source_title: str,
source_url: str,
source_type: str,
source_quality_score: float,
summary: str,
supports_claims: list[UUID],
artifact_manifest_id: UUID,
retrieved_at: datetime,
review_status: str = "PENDING",
) -> EvidenceSummary:
evidence_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO evidence_items (
id,
article_id,
source_title,
source_url,
source_type,
source_quality_score,
summary,
supports_claims,
artifact_manifest_id,
retrieved_at,
review_status
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder}
)
""",
(
str(evidence_id),
str(article_id),
source_title,
source_url,
source_type,
source_quality_score,
summary,
_json_value([str(claim_id) for claim_id in supports_claims]),
str(artifact_manifest_id),
_datetime_value(retrieved_at),
review_status,
),
)
return self.get(evidence_id)
def list_for_article(self, article_id: UUID) -> list[EvidenceSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT {self._select_columns()}
FROM evidence_items
WHERE article_id = {placeholder}
ORDER BY retrieved_at, id
""",
(str(article_id),),
).fetchall()
return [_evidence_from_row(row) for row in rows]
def get(self, evidence_id: UUID) -> EvidenceSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT {self._select_columns()}
FROM evidence_items
WHERE id = {placeholder}
""",
(str(evidence_id),),
).fetchone()
if row is None:
raise LookupError(f"Evidence not found: {evidence_id}")
return _evidence_from_row(row)
def update_review_status(
self,
*,
article_id: UUID,
evidence_id: UUID,
review_status: str,
) -> EvidenceSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE evidence_items
SET review_status = {placeholder}
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(review_status, str(article_id), str(evidence_id)),
)
return self.get(evidence_id)
def delete(
self,
*,
article_id: UUID,
evidence_id: UUID,
) -> None:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
DELETE FROM evidence_items
WHERE article_id = {placeholder} AND id = {placeholder}
""",
(str(article_id), str(evidence_id)),
)
def _select_columns(self) -> str:
return """
id,
article_id,
source_title,
source_url,
source_type,
source_quality_score,
summary,
supports_claims,
artifact_manifest_id,
retrieved_at,
review_status
"""
class ClaimsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
article_id: UUID,
section_id: UUID | None,
claim_text: str,
support_status: ClaimSupportStatus,
risk_level: ClaimRiskLevel,
evidence_item_ids: list[UUID],
created_at: datetime,
) -> ClaimSummary:
claim_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO claims (
id,
article_id,
section_id,
claim_text,
support_status,
risk_level,
evidence_item_ids,
created_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}
)
""",
(
str(claim_id),
str(article_id),
_uuid_value(section_id),
claim_text,
support_status.value,
risk_level.value,
_json_value([str(evidence_id) for evidence_id in evidence_item_ids]),
_datetime_value(created_at),
),
)
return self.get(claim_id)
def list_for_article(self, article_id: UUID) -> list[ClaimSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
article_id,
section_id,
claim_text,
support_status,
risk_level,
evidence_item_ids
FROM claims
WHERE article_id = {placeholder}
ORDER BY id
""",
(str(article_id),),
).fetchall()
return [_claim_from_row(row) for row in rows]
def get(self, claim_id: UUID) -> ClaimSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
article_id,
section_id,
claim_text,
support_status,
risk_level,
evidence_item_ids
FROM claims
WHERE id = {placeholder}
""",
(str(claim_id),),
).fetchone()
if row is None:
raise LookupError(f"Claim not found: {claim_id}")
return _claim_from_row(row)
def update_evidence_links(
self,
*,
claim_id: UUID,
evidence_item_ids: list[UUID],
support_status: ClaimSupportStatus,
) -> ClaimSummary:
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE claims
SET evidence_item_ids = {placeholder}{json_cast}, support_status = {placeholder}
WHERE id = {placeholder}
""",
(
_json_value([str(evidence_id) for evidence_id in evidence_item_ids]),
support_status.value,
str(claim_id),
),
)
return self.get(claim_id)
class AgentJobsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
article_id: UUID | None,
parent_job_id: UUID | None,
attempt: int,
job_type: AgentJobType,
agent_profile: str,
status: AgentJobStatus,
input_files: list[JsonObject],
payload: JsonObject | None = None,
queued_at: datetime,
) -> AgentJobSummary:
job_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
sql = f"""
INSERT INTO agent_jobs (
id,
article_id,
parent_job_id,
attempt,
job_type,
agent_profile,
status,
input_files,
payload,
queued_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}
)
"""
with self._repository.connection() as connection:
connection.execute(
sql,
(
str(job_id),
_uuid_value(article_id),
_uuid_value(parent_job_id),
attempt,
job_type.value,
agent_profile,
status.value,
_json_value(input_files),
_json_value(payload or {}),
_datetime_value(queued_at),
),
)
return self.get(job_id)
def list(self) -> list[AgentJobSummary]:
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT {self._select_columns()}
FROM agent_jobs
ORDER BY queued_at DESC
"""
).fetchall()
return [_agent_job_summary_from_row(row) for row in rows]
def 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:
row = connection.execute(
f"""
SELECT {self._select_columns()}
FROM agent_jobs
WHERE id = {placeholder}
""",
(str(job_id),),
).fetchone()
if row is None:
raise LookupError(f"Agent job not found: {job_id}")
return _agent_job_summary_from_row(row)
def claim_next_queued(self, *, started_at: datetime) -> AgentJobSummary | None:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT id
FROM agent_jobs
WHERE status = {placeholder}
ORDER BY queued_at
LIMIT 1
""",
(AgentJobStatus.QUEUED.value,),
).fetchone()
if row is None:
return None
job_id = _row_value(row, "id")
connection.execute(
f"""
UPDATE agent_jobs
SET status = {placeholder}, started_at = {placeholder}
WHERE id = {placeholder} AND status = {placeholder}
""",
(
AgentJobStatus.RUNNING.value,
_datetime_value(started_at),
str(job_id),
AgentJobStatus.QUEUED.value,
),
)
return self.get(UUID(str(job_id)))
def complete(
self,
*,
job_id: UUID,
status: AgentJobStatus,
workspace_path: str | None,
output_files: list[JsonObject],
payload: JsonObject | None,
error_category: AgentJobErrorCategory | None,
error_message: str | None,
stdout: str,
stderr: str,
exit_code: int | None,
duration_ms: int | None,
finished_at: datetime,
) -> AgentJobSummary:
existing = self.get(job_id)
if existing.status == AgentJobStatus.CANCELLED:
return existing
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE agent_jobs
SET
status = {placeholder},
workspace_path = {placeholder},
output_files = {placeholder}{json_cast},
payload = {placeholder}{json_cast},
error_category = {placeholder},
error_message = {placeholder},
stdout = {placeholder},
stderr = {placeholder},
exit_code = {placeholder},
duration_ms = {placeholder},
finished_at = {placeholder}
WHERE id = {placeholder}
""",
(
status.value,
workspace_path,
_json_value(output_files),
_json_value(payload or {}),
_agent_error_category_value(error_category),
error_message,
stdout,
stderr,
exit_code,
duration_ms,
_datetime_value(finished_at),
str(job_id),
),
)
return self.get(job_id)
def cancel(self, *, job_id: UUID, finished_at: datetime) -> AgentJobSummary:
existing = self.get(job_id)
if existing.status in {AgentJobStatus.SUCCEEDED, AgentJobStatus.FAILED}:
return existing
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE agent_jobs
SET
status = {placeholder},
error_message = {placeholder},
finished_at = {placeholder}
WHERE id = {placeholder}
""",
(
AgentJobStatus.CANCELLED.value,
"Cancelled",
_datetime_value(finished_at),
str(job_id),
),
)
return self.get(job_id)
def _select_columns(self) -> str:
return """
id,
article_id,
parent_job_id,
attempt,
job_type,
agent_profile,
status,
workspace_path,
input_files,
output_files,
payload,
error_category,
error_message,
stdout,
stderr,
exit_code,
duration_ms,
queued_at,
started_at,
finished_at
"""
class ScriptConfigVersionsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def upsert(
self,
*,
version_id: UUID,
target_site_id: UUID,
version: int,
status: ScriptConfigVersionStatus,
created_by: UUID,
created_at: datetime,
updated_at: datetime,
diff: JsonObject,
rollback_target_version_id: UUID | None,
activated_at: datetime | None,
publishing_yaml: str,
publishing_yaml_hash: str,
transform_script: str,
transform_script_hash: str,
) -> dict[str, Any]:
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
sql = f"""
INSERT INTO script_config_versions (
id,
target_site_id,
version,
status,
created_by,
created_at,
updated_at,
diff,
rollback_target_version_id,
activated_at,
publishing_yaml,
publishing_yaml_hash,
transform_script,
transform_script_hash
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
ON CONFLICT (target_site_id, version) DO UPDATE SET
status = excluded.status,
created_by = excluded.created_by,
updated_at = excluded.updated_at,
diff = excluded.diff,
rollback_target_version_id = excluded.rollback_target_version_id,
activated_at = excluded.activated_at,
publishing_yaml = excluded.publishing_yaml,
publishing_yaml_hash = excluded.publishing_yaml_hash,
transform_script = excluded.transform_script,
transform_script_hash = excluded.transform_script_hash
"""
params = (
str(version_id),
str(target_site_id),
version,
status.value,
str(created_by),
_datetime_value(created_at),
_datetime_value(updated_at),
_json_value(diff),
_uuid_value(rollback_target_version_id),
_datetime_value(activated_at),
publishing_yaml,
publishing_yaml_hash,
transform_script,
transform_script_hash,
)
with self._repository.connection() as connection:
connection.execute(sql, params)
return self.get_by_site_and_version(target_site_id=target_site_id, version=version)
def get_by_site_and_version(
self, *, target_site_id: UUID, version: int
) -> dict[str, Any]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT *
FROM script_config_versions
WHERE target_site_id = {placeholder} AND version = {placeholder}
""",
(str(target_site_id), version),
).fetchone()
if row is None:
raise LookupError(
f"Script config version not found: {target_site_id} v{version}"
)
return _plain_row(row)
def get_by_id(self, version_id: UUID) -> dict[str, Any]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT *
FROM script_config_versions
WHERE id = {placeholder}
""",
(str(version_id),),
).fetchone()
if row is None:
raise LookupError(f"Script config version not found: {version_id}")
return _plain_row(row)
def activate(
self,
*,
target_site_id: UUID,
version_id: UUID,
activated_at: datetime,
rollback_target_version_id: UUID | None = None,
) -> dict[str, Any]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
existing = connection.execute(
f"""
SELECT id
FROM script_config_versions
WHERE id = {placeholder} AND target_site_id = {placeholder}
""",
(str(version_id), str(target_site_id)),
).fetchone()
if existing is None:
raise LookupError(
f"Script config version not found: {target_site_id} {version_id}"
)
connection.execute(
f"""
UPDATE script_config_versions
SET
status = {placeholder},
updated_at = {placeholder}
WHERE target_site_id = {placeholder} AND id <> {placeholder}
""",
(
ScriptConfigVersionStatus.DEPRECATED.value,
_datetime_value(activated_at),
str(target_site_id),
str(version_id),
),
)
connection.execute(
f"""
UPDATE script_config_versions
SET
status = {placeholder},
activated_at = {placeholder},
updated_at = {placeholder},
rollback_target_version_id = COALESCE({placeholder}, rollback_target_version_id)
WHERE id = {placeholder}
""",
(
ScriptConfigVersionStatus.ACTIVE.value,
_datetime_value(activated_at),
_datetime_value(activated_at),
_uuid_value(rollback_target_version_id),
str(version_id),
),
)
connection.execute(
f"""
UPDATE target_sites
SET
active_script_config_version_id = {placeholder},
updated_at = {placeholder}
WHERE id = {placeholder}
""",
(str(version_id), _datetime_value(activated_at), str(target_site_id)),
)
return self.get_by_id(version_id)
def list_for_site(self, target_site_id: UUID) -> list[dict[str, Any]]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT *
FROM script_config_versions
WHERE target_site_id = {placeholder}
ORDER BY version
""",
(str(target_site_id),),
).fetchall()
return [_plain_row(row) for row in rows]
class ScriptConfigVersionAuditEventsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
target_site_id: UUID,
version_id: UUID,
event_type: str,
actor_user_id: UUID | None,
payload: JsonObject,
created_at: datetime,
) -> dict[str, Any]:
placeholder = self._repository.placeholder()
sql = f"""
INSERT INTO script_config_version_events (
id,
target_site_id,
version_id,
event_type,
actor_user_id,
payload,
created_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
"""
params = (
str(uuid4()),
str(target_site_id),
str(version_id),
event_type,
_uuid_value(actor_user_id),
_json_value(payload),
_datetime_value(created_at),
)
with self._repository.connection() as connection:
connection.execute(sql, params)
return self.get(event_type=event_type, target_site_id=target_site_id, version_id=version_id)
def list_for_site(self, target_site_id: UUID) -> list[dict[str, Any]]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT id, target_site_id, version_id, event_type, actor_user_id, payload, created_at
FROM script_config_version_events
WHERE target_site_id = {placeholder}
ORDER BY created_at
""",
(str(target_site_id),),
).fetchall()
return [_script_config_version_event_from_row(row) for row in rows]
def get(
self,
*,
target_site_id: UUID,
version_id: UUID,
event_type: str,
) -> dict[str, Any]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT id, target_site_id, version_id, event_type, actor_user_id, payload, created_at
FROM script_config_version_events
WHERE target_site_id = {placeholder}
AND version_id = {placeholder}
AND event_type = {placeholder}
""",
(str(target_site_id), str(version_id), event_type),
).fetchone()
if row is None:
raise LookupError(
f"Script config version event not found: {target_site_id} {version_id} {event_type}"
)
return _script_config_version_event_from_row(row)
class SchemaRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def list_tables(self) -> set[str]:
with self._repository.connection() as connection:
if self._repository.dialect == "sqlite":
rows = connection.execute(
"""
SELECT name
FROM sqlite_master
WHERE type = 'table' AND name NOT LIKE 'sqlite_%'
"""
).fetchall()
else:
rows = connection.execute(
"""
SELECT table_name AS name
FROM information_schema.tables
WHERE table_schema = 'public' AND table_type = 'BASE TABLE'
"""
).fetchall()
return {_row_value(row, "name") for row in rows}
def list_columns(self, table_name: str) -> set[str]:
with self._repository.connection() as connection:
if self._repository.dialect == "sqlite":
rows = connection.execute(f"PRAGMA table_info({table_name})").fetchall()
return {_row_value(row, "name") for row in rows}
rows = connection.execute(
"""
SELECT column_name
FROM information_schema.columns
WHERE table_schema = 'public' AND table_name = %s
""",
(table_name,),
).fetchall()
return {_row_value(row, "column_name") for row in rows}
def list_indexes(self, table_name: str) -> set[str]:
with self._repository.connection() as connection:
if self._repository.dialect == "sqlite":
rows = connection.execute(f"PRAGMA index_list({table_name})").fetchall()
return {_row_value(row, "name") for row in rows}
rows = connection.execute(
"""
SELECT indexname
FROM pg_indexes
WHERE schemaname = 'public' AND tablename = %s
""",
(table_name,),
).fetchall()
return {_row_value(row, "indexname") for row in rows}
def count_rows(self, table_name: str) -> int:
with self._repository.connection() as connection:
row = connection.execute(f"SELECT COUNT(*) AS count FROM {table_name}").fetchone()
return int(_row_value(row, "count"))
def _user_from_row(row: Any) -> UserSummary:
return UserSummary(
id=_row_value(row, "id"),
display_name=_row_value(row, "display_name"),
role=_row_value(row, "role"),
)
def _target_site_from_row(row: Any) -> TargetSiteConfig:
return TargetSiteConfig(
id=_row_value(row, "id"),
name=_row_value(row, "name"),
slug=_row_value(row, "slug"),
publishing_type=_row_value(row, "publishing_type"),
default_language=_row_value(row, "default_language"),
brand_voice=_row_value(row, "brand_voice"),
audience=_row_value(row, "audience"),
seo_rules=_json_from_row(row, "seo_rules"),
visual_rules=_json_from_row(row, "visual_rules"),
source_rules=_json_from_row(row, "source_rules"),
publishing_rules=_json_from_row(row, "publishing_rules"),
active_script_config_version_id=_row_value(
row, "active_script_config_version_id"
),
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
)
def _article_summary_from_row(row: Any) -> ArticleSummary:
return ArticleSummary(
id=_row_value(row, "id"),
target_site_id=_row_value(row, "target_site_id"),
status=_row_value(row, "status"),
publishing_status=_row_value(row, "publishing_status"),
brief_description=_row_value(row, "brief_description"),
working_title=_row_value(row, "working_title"),
language=_row_value(row, "language"),
content_type=_row_value(row, "content_type"),
primary_keyword=_row_value(row, "primary_keyword"),
assigned_editor_id=_row_value(row, "assigned_editor_id"),
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
)
def _workflow_event_from_row(row: Any) -> WorkflowEventSummary:
return WorkflowEventSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
event_type=_row_value(row, "event_type"),
from_status=_row_value(row, "from_status"),
to_status=_row_value(row, "to_status"),
actor_user_id=_row_value(row, "actor_user_id"),
payload=_json_from_row(row, "payload"),
created_at=_row_value(row, "created_at"),
)
def _boundary_question_from_row(row: Any) -> BoundaryQuestionSummary:
return BoundaryQuestionSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
sort_order=_row_value(row, "sort_order"),
category=_row_value(row, "category"),
question=_row_value(row, "question"),
answer=_row_value(row, "answer"),
is_required=bool(_row_value(row, "is_required")),
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
)
def _plan_section_from_row(row: Any) -> PlanSectionSummary:
return PlanSectionSummary(
id=_row_value(row, "id"),
article_plan_id=_row_value(row, "article_plan_id"),
sort_order=_row_value(row, "sort_order"),
heading=_row_value(row, "heading"),
purpose=_row_value(row, "purpose"),
key_points=_json_from_row(row, "key_points"),
evidence_needs=_json_from_row(row, "evidence_needs"),
claims_to_support=_json_from_row(row, "claims_to_support"),
target_word_count=_row_value(row, "target_word_count"),
)
def _plan_summary_from_row(
row: Any,
sections: list[PlanSectionSummary],
) -> PlanSummary:
return PlanSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
version=_row_value(row, "version"),
status=_row_value(row, "status"),
title_options=_json_from_row(row, "title_options"),
recommended_title=_row_value(row, "recommended_title"),
reader_persona=_row_value(row, "reader_persona"),
search_intent=_row_value(row, "search_intent"),
thesis=_row_value(row, "thesis"),
sections=sections,
claims_to_prove=_json_from_row(row, "claims_to_prove"),
evidence_needs=_json_from_row(row, "evidence_needs"),
visual_needs=_json_from_row(row, "visual_needs"),
seo_notes=_json_from_row(row, "seo_notes"),
source_requirements=_json_from_row(row, "source_requirements"),
excluded_sources=_json_from_row(row, "excluded_sources"),
tone=_row_value(row, "tone"),
audience=_row_value(row, "audience"),
risks=_json_from_row(row, "risks"),
created_at=_row_value(row, "created_at"),
)
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 _evidence_from_row(row: Any) -> EvidenceSummary:
return EvidenceSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
source_title=_row_value(row, "source_title"),
source_url=_row_value(row, "source_url"),
source_type=_row_value(row, "source_type"),
source_quality_score=float(_row_value(row, "source_quality_score")),
summary=_row_value(row, "summary"),
supports_claims=[UUID(value) for value in _json_from_row(row, "supports_claims")],
artifact_manifest_id=_row_value(row, "artifact_manifest_id"),
retrieved_at=_row_value(row, "retrieved_at"),
review_status=_row_value(row, "review_status"),
)
def _claim_from_row(row: Any) -> ClaimSummary:
return ClaimSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
section_id=_row_value(row, "section_id"),
claim_text=_row_value(row, "claim_text"),
support_status=_row_value(row, "support_status"),
risk_level=_row_value(row, "risk_level"),
evidence_item_ids=[
UUID(value) for value in _json_from_row(row, "evidence_item_ids")
],
)
def _draft_summary_from_row(row: Any) -> DraftSummary:
return DraftSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
version=_row_value(row, "version"),
title=_row_value(row, "title"),
slug=_row_value(row, "slug"),
meta_title=_row_value(row, "meta_title"),
meta_description=_row_value(row, "meta_description"),
body_object_key=_row_value(row, "body_object_key"),
body_markdown=_row_value(row, "body_markdown"),
faq_items=[
DraftFaqItem.model_validate(item)
for item in _json_from_row(row, "faq_items")
],
visual_placeholders=_json_from_row(row, "visual_placeholders"),
evidence_references=[
UUID(value) for value in _json_from_row(row, "evidence_references")
],
unsupported_claim_warnings=_json_from_row(
row, "unsupported_claim_warnings"
),
based_on_draft_id=_row_value(row, "based_on_draft_id"),
status=_row_value(row, "status"),
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
)
def _content_review_report_from_row(
row: Any,
*,
issues: list[ContentReviewIssueSummary],
) -> ContentReviewReportSummary:
return ContentReviewReportSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
review_kind=_row_value(row, "review_kind"),
draft_id=_row_value(row, "draft_id"),
score=_row_value(row, "score"),
issues=issues,
recommended_slug=_row_value(row, "recommended_slug"),
recommended_title=_row_value(row, "recommended_title"),
schema_json=_json_from_row(row, "schema_json"),
rules_snapshot=_json_from_row(row, "rules_snapshot"),
created_at=_row_value(row, "created_at"),
)
def _content_review_suggestion_from_row(row: Any) -> ContentReviewSuggestionSummary:
return ContentReviewSuggestionSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
review_kind=_row_value(row, "review_kind"),
report_id=_row_value(row, "report_id"),
suggestion_key=_row_value(row, "suggestion_key"),
severity=_row_value(row, "severity"),
location=_row_value(row, "location"),
message=_row_value(row, "message"),
suggested_fix=_row_value(row, "suggested_fix"),
suggested_rewrite=_row_value(row, "suggested_rewrite"),
patch=_json_from_row(row, "patch"),
status=_row_value(row, "status"),
applied_text=_row_value(row, "applied_text"),
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
)
def _asset_revision_from_row(row: Any) -> AssetRevisionSummary:
return AssetRevisionSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
asset_id=_row_value(row, "asset_id"),
revision_index=_row_value(row, "revision_index"),
action=_row_value(row, "action"),
actor_user_id=_row_value(row, "actor_user_id"),
payload=_json_from_row(row, "payload"),
created_at=_row_value(row, "created_at"),
)
def _asset_summary_from_row(
row: Any,
*,
history: list[AssetRevisionSummary] | None = None,
) -> AssetSummary:
return AssetSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
section_id=_row_value(row, "section_id"),
asset_type=_row_value(row, "asset_type"),
title=_row_value(row, "title"),
prompt=_row_value(row, "prompt"),
object_key=_row_value(row, "object_key"),
file_url=_row_value(row, "file_url"),
alt_text=_row_value(row, "alt_text"),
caption=_row_value(row, "caption"),
status=_row_value(row, "status"),
history=history or [],
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
)
def _agent_job_summary_from_row(row: Any) -> AgentJobSummary:
return AgentJobSummary(
id=_row_value(row, "id"),
article_id=_row_value(row, "article_id"),
parent_job_id=_row_value(row, "parent_job_id"),
attempt=_row_value(row, "attempt"),
job_type=_row_value(row, "job_type"),
agent_profile=_row_value(row, "agent_profile"),
status=_row_value(row, "status"),
workspace_path=_row_value(row, "workspace_path"),
input_files=_json_from_row(row, "input_files"),
output_files=_json_from_row(row, "output_files"),
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"),
stderr=_row_value(row, "stderr"),
exit_code=_row_value(row, "exit_code"),
duration_ms=_row_value(row, "duration_ms"),
queued_at=_row_value(row, "queued_at"),
started_at=_row_value(row, "started_at"),
finished_at=_row_value(row, "finished_at"),
)
def _script_config_version_event_from_row(row: Any) -> dict[str, Any]:
return {
"id": _row_value(row, "id"),
"target_site_id": _row_value(row, "target_site_id"),
"version_id": _row_value(row, "version_id"),
"event_type": _row_value(row, "event_type"),
"actor_user_id": _row_value(row, "actor_user_id"),
"payload": _json_from_row(row, "payload"),
"created_at": _row_value(row, "created_at"),
}
def _plain_row(row: Any) -> dict[str, Any]:
if isinstance(row, sqlite3.Row):
result = dict(row)
else:
result = dict(row)
for key in (
"diff",
"seo_rules",
"visual_rules",
"source_rules",
"publishing_rules",
):
if key in result:
result[key] = _json_decode(result[key])
return result
def _row_value(row: Any, key: str) -> Any:
return row[key]
def _json_from_row(row: Any, key: str) -> JsonObject:
return _json_decode(_row_value(row, key))
def _json_decode(value: Any) -> Any:
if isinstance(value, str):
return json.loads(value)
return value
def _json_value(value: Any) -> str:
return json.dumps(value, sort_keys=True, separators=(",", ":"))
def _uuid_value(value: UUID | None) -> str | None:
if value is None:
return None
return str(value)
def _datetime_value(value: datetime | None) -> str | None:
if value is None:
return None
return value.isoformat()
def _bool_value(value: bool) -> bool:
return bool(value)
def _article_status_value(value: ArticleWorkflowStatus | None) -> str | None:
if value is None:
return None
return value.value
def _agent_error_category_value(value: AgentJobErrorCategory | None) -> str | None:
if value is None:
return None
return value.value