feat(task-021): add admin workflow template api

This commit is contained in:
2026-05-22 03:45:02 +03:00
parent d1ae5c4437
commit cbd24d4a1c
12 changed files with 2131 additions and 15 deletions
@@ -0,0 +1,372 @@
from __future__ import annotations
from datetime import UTC, datetime
from uuid import UUID
from src.domain.contracts import (
CurrentUser,
Role,
WorkflowStageCreateRequest,
WorkflowStageResponse,
WorkflowStageUpdateRequest,
WorkflowTemplateAuditEventListResponse,
WorkflowTemplateCreateRequest,
WorkflowTemplateListResponse,
WorkflowTemplateResponse,
WorkflowTemplateStatus,
WorkflowTemplateSummary,
WorkflowTemplateUpdateRequest,
)
_TEMPLATE_CREATED_EVENT_TYPE = "WORKFLOW_TEMPLATE_CREATED"
_TEMPLATE_UPDATED_EVENT_TYPE = "WORKFLOW_TEMPLATE_UPDATED"
_STAGE_CREATED_EVENT_TYPE = "WORKFLOW_STAGE_CREATED"
_STAGE_UPDATED_EVENT_TYPE = "WORKFLOW_STAGE_UPDATED"
_STAGE_DELETED_EVENT_TYPE = "WORKFLOW_STAGE_DELETED"
_STAGES_REORDERED_EVENT_TYPE = "WORKFLOW_STAGES_REORDERED"
_TEMPLATE_ACTIVATED_EVENT_TYPE = "WORKFLOW_TEMPLATE_ACTIVATED"
_TEMPLATE_ARCHIVED_EVENT_TYPE = "WORKFLOW_TEMPLATE_ARCHIVED"
def create_workflow_template(
repository: object,
*,
request: WorkflowTemplateCreateRequest,
current_user: CurrentUser,
) -> WorkflowTemplateResponse:
now = _now()
workflow = repository.workflow_templates.create(
name=request.name,
slug=request.slug,
description=request.description,
status=WorkflowTemplateStatus.DRAFT,
version=1,
created_by=current_user.id,
updated_by=current_user.id,
created_at=now,
updated_at=now,
)
_create_event(
repository,
workflow_id=workflow.id,
event_type=_TEMPLATE_CREATED_EVENT_TYPE,
current_user=current_user,
payload={"slug": workflow.slug, "version": workflow.version},
created_at=now,
)
return WorkflowTemplateResponse(workflow=workflow)
def list_workflow_templates(
repository: object,
*,
current_user: CurrentUser,
) -> WorkflowTemplateListResponse:
if current_user.role == Role.ADMIN:
workflows = repository.workflow_templates.list()
else:
workflows = repository.workflow_templates.list_active()
return WorkflowTemplateListResponse(workflows=workflows)
def get_workflow_template(
repository: object,
*,
workflow_id: UUID,
current_user: CurrentUser,
) -> WorkflowTemplateResponse:
workflow = repository.workflow_templates.get(workflow_id)
if current_user.role != Role.ADMIN and workflow.status != WorkflowTemplateStatus.ACTIVE:
raise LookupError(f"Workflow template not found: {workflow_id}")
return WorkflowTemplateResponse(workflow=workflow)
def update_workflow_template(
repository: object,
*,
workflow_id: UUID,
request: WorkflowTemplateUpdateRequest,
current_user: CurrentUser,
) -> WorkflowTemplateResponse:
existing = repository.workflow_templates.get(workflow_id)
_require_draft(existing)
changes = request.model_dump(exclude_unset=True)
now = _now()
workflow = repository.workflow_templates.update(
workflow_id=workflow_id,
name=changes.get("name", existing.name),
slug=changes.get("slug", existing.slug),
description=changes.get("description", existing.description),
updated_by=current_user.id,
updated_at=now,
)
_create_event(
repository,
workflow_id=workflow.id,
event_type=_TEMPLATE_UPDATED_EVENT_TYPE,
current_user=current_user,
payload={"changed_fields": sorted(changes.keys())},
created_at=now,
)
return WorkflowTemplateResponse(workflow=workflow)
def add_workflow_stage(
repository: object,
*,
workflow_id: UUID,
request: WorkflowStageCreateRequest,
current_user: CurrentUser,
) -> WorkflowStageResponse:
workflow = repository.workflow_templates.get(workflow_id)
_require_draft(workflow)
now = _now()
stage = repository.workflow_template_stages.create(
workflow_id=workflow_id,
stable_key=request.stable_key,
display_name=request.display_name,
description=request.description,
position=request.position,
owner_role=request.owner_role,
runner_profile_key=request.runner_profile_key,
required_inputs=request.required_inputs,
expected_outputs=request.expected_outputs,
acceptance_criteria=request.acceptance_criteria,
requires_human_approval=request.requires_human_approval,
retry_policy=request.retry_policy,
parts=_parts_json(request.parts),
updated_by=current_user.id,
created_at=now,
updated_at=now,
)
_create_event(
repository,
workflow_id=workflow_id,
event_type=_STAGE_CREATED_EVENT_TYPE,
current_user=current_user,
payload={
"stage_id": str(stage.id),
"stable_key": stage.stable_key,
"position": stage.position,
},
created_at=now,
)
return WorkflowStageResponse(stage=stage)
def update_workflow_stage(
repository: object,
*,
workflow_id: UUID,
stage_id: UUID,
request: WorkflowStageUpdateRequest,
current_user: CurrentUser,
) -> WorkflowStageResponse:
workflow = repository.workflow_templates.get(workflow_id)
_require_draft(workflow)
existing = repository.workflow_template_stages.get(
workflow_id=workflow_id,
stage_id=stage_id,
)
changes = request.model_dump(exclude_unset=True)
now = _now()
stage = repository.workflow_template_stages.update(
workflow_id=workflow_id,
stage_id=stage_id,
stable_key=changes.get("stable_key", existing.stable_key),
display_name=changes.get("display_name", existing.display_name),
description=changes.get("description", existing.description),
position=changes.get("position", existing.position),
owner_role=request.owner_role or existing.owner_role,
runner_profile_key=changes.get(
"runner_profile_key",
existing.runner_profile_key,
),
required_inputs=changes.get("required_inputs", existing.required_inputs),
expected_outputs=changes.get("expected_outputs", existing.expected_outputs),
acceptance_criteria=changes.get(
"acceptance_criteria",
existing.acceptance_criteria,
),
requires_human_approval=changes.get(
"requires_human_approval",
existing.requires_human_approval,
),
retry_policy=changes.get("retry_policy", existing.retry_policy),
parts=(
_parts_json(request.parts)
if request.parts is not None
else _parts_json(existing.parts)
),
updated_by=current_user.id,
updated_at=now,
)
_create_event(
repository,
workflow_id=workflow_id,
event_type=_STAGE_UPDATED_EVENT_TYPE,
current_user=current_user,
payload={"stage_id": str(stage_id), "changed_fields": sorted(changes.keys())},
created_at=now,
)
return WorkflowStageResponse(stage=stage)
def delete_workflow_stage(
repository: object,
*,
workflow_id: UUID,
stage_id: UUID,
current_user: CurrentUser,
) -> WorkflowTemplateResponse:
workflow = repository.workflow_templates.get(workflow_id)
_require_draft(workflow)
existing = repository.workflow_template_stages.get(
workflow_id=workflow_id,
stage_id=stage_id,
)
now = _now()
repository.workflow_template_stages.delete(
workflow_id=workflow_id,
stage_id=stage_id,
updated_by=current_user.id,
updated_at=now,
)
_create_event(
repository,
workflow_id=workflow_id,
event_type=_STAGE_DELETED_EVENT_TYPE,
current_user=current_user,
payload={"stage_id": str(stage_id), "stable_key": existing.stable_key},
created_at=now,
)
return WorkflowTemplateResponse(workflow=repository.workflow_templates.get(workflow_id))
def reorder_workflow_stages(
repository: object,
*,
workflow_id: UUID,
stage_ids: list[UUID],
current_user: CurrentUser,
) -> WorkflowTemplateResponse:
workflow = repository.workflow_templates.get(workflow_id)
_require_draft(workflow)
now = _now()
repository.workflow_template_stages.reorder(
workflow_id=workflow_id,
stage_ids=stage_ids,
updated_by=current_user.id,
updated_at=now,
)
_create_event(
repository,
workflow_id=workflow_id,
event_type=_STAGES_REORDERED_EVENT_TYPE,
current_user=current_user,
payload={"stage_ids": [str(stage_id) for stage_id in stage_ids]},
created_at=now,
)
return WorkflowTemplateResponse(workflow=repository.workflow_templates.get(workflow_id))
def activate_workflow_template(
repository: object,
*,
workflow_id: UUID,
current_user: CurrentUser,
) -> WorkflowTemplateResponse:
workflow = repository.workflow_templates.get(workflow_id)
_require_draft(workflow)
now = _now()
activated = repository.workflow_templates.activate(
workflow_id=workflow_id,
updated_by=current_user.id,
activated_at=now,
)
_create_event(
repository,
workflow_id=workflow_id,
event_type=_TEMPLATE_ACTIVATED_EVENT_TYPE,
current_user=current_user,
payload={"version": activated.version},
created_at=now,
)
return WorkflowTemplateResponse(workflow=activated)
def archive_workflow_template(
repository: object,
*,
workflow_id: UUID,
current_user: CurrentUser,
) -> WorkflowTemplateResponse:
repository.workflow_templates.get(workflow_id)
now = _now()
archived = repository.workflow_templates.archive(
workflow_id=workflow_id,
updated_by=current_user.id,
archived_at=now,
)
_create_event(
repository,
workflow_id=workflow_id,
event_type=_TEMPLATE_ARCHIVED_EVENT_TYPE,
current_user=current_user,
payload={"version": archived.version},
created_at=now,
)
return WorkflowTemplateResponse(workflow=archived)
def list_workflow_template_audit_events(
repository: object,
*,
workflow_id: UUID,
) -> WorkflowTemplateAuditEventListResponse:
repository.workflow_templates.get(workflow_id)
return WorkflowTemplateAuditEventListResponse(
events=repository.workflow_template_events.list_for_workflow(workflow_id)
)
def _create_event(
repository: object,
*,
workflow_id: UUID,
event_type: str,
current_user: CurrentUser,
payload: dict[str, object],
created_at: datetime,
) -> None:
repository.workflow_template_events.create(
workflow_id=workflow_id,
event_type=event_type,
actor_user_id=current_user.id,
payload=payload,
created_at=created_at,
)
def _require_draft(workflow: WorkflowTemplateSummary) -> None:
if workflow.status != WorkflowTemplateStatus.DRAFT:
raise ValueError("Only draft workflow templates can be mutated.")
def _parts_json(parts: list[object]) -> list[dict[str, object]]:
return [
part.model_dump(mode="json") if hasattr(part, "model_dump") else dict(part)
for part in parts
]
def _now() -> datetime:
return datetime.now(UTC)
@@ -15,6 +15,7 @@ from .enums import (
ReviewType,
Role,
ScriptConfigVersionStatus,
WorkflowTemplateStatus,
)
from .models import (
AgentJobOutput,
@@ -33,6 +34,19 @@ from .models import (
BoundaryQuestionUpdateRequest,
WorkflowEventSummary,
ObservabilityTimelineEventSummary,
WorkflowStageCreateRequest,
WorkflowStagePart,
WorkflowStageReorderRequest,
WorkflowStageResponse,
WorkflowStageSummary,
WorkflowStageUpdateRequest,
WorkflowTemplateAuditEventListResponse,
WorkflowTemplateAuditEventSummary,
WorkflowTemplateCreateRequest,
WorkflowTemplateListResponse,
WorkflowTemplateResponse,
WorkflowTemplateSummary,
WorkflowTemplateUpdateRequest,
AssetGenerateSpecsResponse,
AssetListResponse,
AssetResponse,
@@ -133,6 +147,20 @@ __all__ = [
"BoundaryQuestionUpdateRequest",
"WorkflowEventSummary",
"ObservabilityTimelineEventSummary",
"WorkflowStageCreateRequest",
"WorkflowStagePart",
"WorkflowStageReorderRequest",
"WorkflowStageResponse",
"WorkflowStageSummary",
"WorkflowStageUpdateRequest",
"WorkflowTemplateAuditEventListResponse",
"WorkflowTemplateAuditEventSummary",
"WorkflowTemplateCreateRequest",
"WorkflowTemplateListResponse",
"WorkflowTemplateResponse",
"WorkflowTemplateStatus",
"WorkflowTemplateSummary",
"WorkflowTemplateUpdateRequest",
"AssetGenerateSpecsResponse",
"AssetListResponse",
"AssetResponse",
@@ -139,3 +139,9 @@ class ScriptConfigVersionStatus(str, Enum):
DRAFT = "DRAFT"
ACTIVE = "ACTIVE"
DEPRECATED = "DEPRECATED"
class WorkflowTemplateStatus(str, Enum):
DRAFT = "DRAFT"
ACTIVE = "ACTIVE"
ARCHIVED = "ARCHIVED"
+115
View File
@@ -23,6 +23,7 @@ from .enums import (
ReviewType,
Role,
ScriptConfigVersionStatus,
WorkflowTemplateStatus,
)
@@ -158,6 +159,120 @@ class ScriptConfigVersionAuditEventListResponse(ContractModel):
events: list[ScriptConfigVersionAuditEventSummary] = Field(default_factory=list)
class WorkflowStagePart(ContractModel):
key: str = Field(min_length=1)
type: str = Field(min_length=1)
title: str = Field(min_length=1)
payload: JsonObject = Field(default_factory=dict)
acceptance_criteria: list[str] = Field(default_factory=list)
class WorkflowStageSummary(ContractModel):
id: UUID
workflow_id: UUID
stable_key: str = Field(min_length=1)
display_name: str = Field(min_length=1)
description: str = Field(min_length=1)
position: int = Field(ge=1)
owner_role: Role
runner_profile_key: str = Field(min_length=1)
required_inputs: list[str] = Field(default_factory=list)
expected_outputs: list[str] = Field(default_factory=list)
acceptance_criteria: list[str] = Field(default_factory=list)
requires_human_approval: bool = False
retry_policy: JsonObject = Field(default_factory=dict)
parts: list[WorkflowStagePart] = Field(default_factory=list)
created_at: datetime
updated_at: datetime
class WorkflowTemplateSummary(ContractModel):
id: UUID
name: str = Field(min_length=1)
slug: str = Field(min_length=1)
description: str = Field(min_length=1)
status: WorkflowTemplateStatus
version: int = Field(ge=1)
created_by: UUID
updated_by: UUID | None = None
created_at: datetime
updated_at: datetime
activated_at: datetime | None = None
archived_at: datetime | None = None
stages: list[WorkflowStageSummary] = Field(default_factory=list)
class WorkflowTemplateCreateRequest(ContractModel):
name: str = Field(min_length=1)
slug: str = Field(min_length=1)
description: str = Field(min_length=1)
class WorkflowTemplateUpdateRequest(ContractModel):
name: str | None = Field(default=None, min_length=1)
slug: str | None = Field(default=None, min_length=1)
description: str | None = Field(default=None, min_length=1)
class WorkflowStageCreateRequest(ContractModel):
stable_key: str = Field(min_length=1)
display_name: str = Field(min_length=1)
description: str = Field(min_length=1)
position: int | None = Field(default=None, ge=1)
owner_role: Role
runner_profile_key: str = Field(min_length=1)
required_inputs: list[str] = Field(default_factory=list)
expected_outputs: list[str] = Field(default_factory=list)
acceptance_criteria: list[str] = Field(default_factory=list)
requires_human_approval: bool = False
retry_policy: JsonObject = Field(default_factory=dict)
parts: list[WorkflowStagePart] = Field(default_factory=list)
class WorkflowStageUpdateRequest(ContractModel):
stable_key: str | None = Field(default=None, min_length=1)
display_name: str | None = Field(default=None, min_length=1)
description: str | None = Field(default=None, min_length=1)
position: int | None = Field(default=None, ge=1)
owner_role: Role | None = None
runner_profile_key: str | None = Field(default=None, min_length=1)
required_inputs: list[str] | None = None
expected_outputs: list[str] | None = None
acceptance_criteria: list[str] | None = None
requires_human_approval: bool | None = None
retry_policy: JsonObject | None = None
parts: list[WorkflowStagePart] | None = None
class WorkflowStageReorderRequest(ContractModel):
stage_ids: list[UUID] = Field(min_length=1)
class WorkflowTemplateResponse(ContractModel):
workflow: WorkflowTemplateSummary
class WorkflowTemplateListResponse(ContractModel):
workflows: list[WorkflowTemplateSummary] = Field(default_factory=list)
class WorkflowStageResponse(ContractModel):
stage: WorkflowStageSummary
class WorkflowTemplateAuditEventSummary(ContractModel):
id: UUID
workflow_id: UUID
event_type: str = Field(min_length=1)
actor_user_id: UUID | None = None
payload: JsonObject = Field(default_factory=dict)
created_at: datetime
class WorkflowTemplateAuditEventListResponse(ContractModel):
events: list[WorkflowTemplateAuditEventSummary] = Field(default_factory=list)
class ArticleCreateRequest(ContractModel):
target_site_id: UUID
brief_description: str = Field(min_length=1)
@@ -22,6 +22,7 @@ from .enums import (
ReviewType,
Role,
ScriptConfigVersionStatus,
WorkflowTemplateStatus,
)
from .models import (
AgentJobOutput,
@@ -109,6 +110,19 @@ from .models import (
TargetSiteConfigCreateRequest,
TargetSiteConfigResponse,
WorkflowEventSummary,
WorkflowStageCreateRequest,
WorkflowStagePart,
WorkflowStageReorderRequest,
WorkflowStageResponse,
WorkflowStageSummary,
WorkflowStageUpdateRequest,
WorkflowTemplateAuditEventListResponse,
WorkflowTemplateAuditEventSummary,
WorkflowTemplateCreateRequest,
WorkflowTemplateListResponse,
WorkflowTemplateResponse,
WorkflowTemplateSummary,
WorkflowTemplateUpdateRequest,
TargetSiteConfigUpdateRequest,
UserSummary,
)
@@ -131,6 +145,7 @@ CONTRACT_ENUMS: tuple[type[Enum], ...] = (
ReviewType,
ContentReviewKind,
ScriptConfigVersionStatus,
WorkflowTemplateStatus,
)
CONTRACT_SCHEMA_MODELS: tuple[type[BaseModel], ...] = (
@@ -148,6 +163,19 @@ CONTRACT_SCHEMA_MODELS: tuple[type[BaseModel], ...] = (
ScriptConfigVersionListResponse,
ScriptConfigVersionAuditEventSummary,
ScriptConfigVersionAuditEventListResponse,
WorkflowStagePart,
WorkflowStageSummary,
WorkflowTemplateSummary,
WorkflowTemplateCreateRequest,
WorkflowTemplateUpdateRequest,
WorkflowStageCreateRequest,
WorkflowStageUpdateRequest,
WorkflowStageReorderRequest,
WorkflowTemplateResponse,
WorkflowTemplateListResponse,
WorkflowStageResponse,
WorkflowTemplateAuditEventSummary,
WorkflowTemplateAuditEventListResponse,
ArticleCreateRequest,
ArticleSummary,
ArticleCreateResponse,
+6
View File
@@ -4,6 +4,9 @@ from __future__ import annotations
USERS_TABLE = "users"
TARGET_SITES_TABLE = "target_sites"
SCRIPT_CONFIG_VERSIONS_TABLE = "script_config_versions"
WORKFLOW_TEMPLATES_TABLE = "workflow_templates"
WORKFLOW_TEMPLATE_STAGES_TABLE = "workflow_template_stages"
WORKFLOW_TEMPLATE_EVENTS_TABLE = "workflow_template_events"
ARTICLES_TABLE = "articles"
BOUNDARY_QUESTIONS_TABLE = "boundary_questions"
ARTICLE_PLANS_TABLE = "article_plans"
@@ -26,6 +29,9 @@ CORE_TABLES: tuple[str, ...] = (
USERS_TABLE,
TARGET_SITES_TABLE,
SCRIPT_CONFIG_VERSIONS_TABLE,
WORKFLOW_TEMPLATES_TABLE,
WORKFLOW_TEMPLATE_STAGES_TABLE,
WORKFLOW_TEMPLATE_EVENTS_TABLE,
ARTICLES_TABLE,
BOUNDARY_QUESTIONS_TABLE,
ARTICLE_PLANS_TABLE,
@@ -44,6 +44,11 @@ from src.domain.contracts import (
ScriptConfigVersionStatus,
TargetSiteConfig,
UserSummary,
WorkflowStagePart,
WorkflowStageSummary,
WorkflowTemplateAuditEventSummary,
WorkflowTemplateStatus,
WorkflowTemplateSummary,
)
from src.infrastructure.schema import setup_database
@@ -74,6 +79,9 @@ class BackendRepository:
self.agent_jobs = AgentJobsRepository(self)
self.script_config_versions = ScriptConfigVersionsRepository(self)
self.script_config_version_events = ScriptConfigVersionAuditEventsRepository(self)
self.workflow_templates = WorkflowTemplatesRepository(self)
self.workflow_template_stages = WorkflowTemplateStagesRepository(self)
self.workflow_template_events = WorkflowTemplateAuditEventsRepository(self)
self.schema = SchemaRepository(self)
def setup(self) -> None:
@@ -2735,6 +2743,722 @@ class AgentJobsRepository:
"""
class WorkflowTemplatesRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
name: str,
slug: str,
description: str,
status: WorkflowTemplateStatus,
version: int,
created_by: UUID,
updated_by: UUID,
created_at: datetime,
updated_at: datetime,
) -> WorkflowTemplateSummary:
workflow_id = uuid4()
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO workflow_templates (
id,
name,
slug,
description,
status,
version,
created_by,
updated_by,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}
)
""",
(
str(workflow_id),
name,
slug,
description,
status.value,
version,
str(created_by),
str(updated_by),
_datetime_value(created_at),
_datetime_value(updated_at),
),
)
return self.get(workflow_id)
def list(self) -> list[WorkflowTemplateSummary]:
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT {self._select_columns()}
FROM workflow_templates
ORDER BY updated_at DESC, slug
"""
).fetchall()
return [self._with_stages(row) for row in rows]
def list_active(self) -> list[WorkflowTemplateSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT {self._select_columns()}
FROM workflow_templates
WHERE status = {placeholder}
ORDER BY updated_at DESC, slug
""",
(WorkflowTemplateStatus.ACTIVE.value,),
).fetchall()
return [self._with_stages(row) for row in rows]
def get(self, workflow_id: UUID) -> WorkflowTemplateSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT {self._select_columns()}
FROM workflow_templates
WHERE id = {placeholder}
""",
(str(workflow_id),),
).fetchone()
if row is None:
raise LookupError(f"Workflow template not found: {workflow_id}")
return self._with_stages(row)
def update(
self,
*,
workflow_id: UUID,
name: str,
slug: str,
description: str,
updated_by: UUID,
updated_at: datetime,
) -> WorkflowTemplateSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE workflow_templates
SET
name = {placeholder},
slug = {placeholder},
description = {placeholder},
updated_by = {placeholder},
updated_at = {placeholder}
WHERE id = {placeholder}
""",
(
name,
slug,
description,
str(updated_by),
_datetime_value(updated_at),
str(workflow_id),
),
)
return self.get(workflow_id)
def activate(
self,
*,
workflow_id: UUID,
updated_by: UUID,
activated_at: datetime,
) -> WorkflowTemplateSummary:
existing = self.get(workflow_id)
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE workflow_templates
SET
status = {placeholder},
version = {placeholder},
updated_by = {placeholder},
updated_at = {placeholder},
activated_at = {placeholder}
WHERE id = {placeholder}
""",
(
WorkflowTemplateStatus.ACTIVE.value,
existing.version + 1,
str(updated_by),
_datetime_value(activated_at),
_datetime_value(activated_at),
str(workflow_id),
),
)
return self.get(workflow_id)
def archive(
self,
*,
workflow_id: UUID,
updated_by: UUID,
archived_at: datetime,
) -> WorkflowTemplateSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE workflow_templates
SET
status = {placeholder},
updated_by = {placeholder},
updated_at = {placeholder},
archived_at = {placeholder}
WHERE id = {placeholder}
""",
(
WorkflowTemplateStatus.ARCHIVED.value,
str(updated_by),
_datetime_value(archived_at),
_datetime_value(archived_at),
str(workflow_id),
),
)
return self.get(workflow_id)
def touch(
self,
*,
workflow_id: UUID,
updated_by: UUID,
updated_at: datetime,
) -> None:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE workflow_templates
SET updated_by = {placeholder}, updated_at = {placeholder}
WHERE id = {placeholder}
""",
(str(updated_by), _datetime_value(updated_at), str(workflow_id)),
)
def _with_stages(self, row: Any) -> WorkflowTemplateSummary:
workflow_id = _row_value(row, "id")
return _workflow_template_from_row(
row,
stages=self._stages_for_workflow(UUID(str(workflow_id))),
)
def _stages_for_workflow(self, workflow_id: UUID) -> list[WorkflowStageSummary]:
return self._repository.workflow_template_stages.list_for_workflow(workflow_id)
def _select_columns(self) -> str:
return """
id,
name,
slug,
description,
status,
version,
created_by,
updated_by,
created_at,
updated_at,
activated_at,
archived_at
"""
class WorkflowTemplateStagesRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
workflow_id: UUID,
stable_key: str,
display_name: str,
description: str,
position: int | None,
owner_role: Role,
runner_profile_key: str,
required_inputs: list[str],
expected_outputs: list[str],
acceptance_criteria: list[str],
requires_human_approval: bool,
retry_policy: JsonObject,
parts: list[dict[str, Any]],
updated_by: UUID,
created_at: datetime,
updated_at: datetime,
) -> WorkflowStageSummary:
stage_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
insert_position = self._clamped_insert_position(
workflow_id=workflow_id,
requested_position=position,
)
with self._repository.connection() as connection:
connection.execute(
f"""
UPDATE workflow_template_stages
SET position = position + 1
WHERE workflow_template_id = {placeholder}
AND position >= {placeholder}
""",
(str(workflow_id), insert_position),
)
connection.execute(
f"""
INSERT INTO workflow_template_stages (
id,
workflow_template_id,
stable_key,
display_name,
description,
position,
owner_role,
runner_profile_key,
required_inputs,
expected_outputs,
acceptance_criteria,
requires_human_approval,
retry_policy,
parts,
created_at,
updated_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder},
{placeholder}{json_cast},
{placeholder}{json_cast},
{placeholder},
{placeholder}
)
""",
(
str(stage_id),
str(workflow_id),
stable_key,
display_name,
description,
insert_position,
owner_role.value,
runner_profile_key,
_json_value(required_inputs),
_json_value(expected_outputs),
_json_value(acceptance_criteria),
_bool_value(requires_human_approval),
_json_value(retry_policy),
_json_value(parts),
_datetime_value(created_at),
_datetime_value(updated_at),
),
)
self._touch_parent(
connection,
workflow_id=workflow_id,
updated_by=updated_by,
updated_at=updated_at,
)
return self.get(workflow_id=workflow_id, stage_id=stage_id)
def list_for_workflow(self, workflow_id: UUID) -> list[WorkflowStageSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT {self._select_columns()}
FROM workflow_template_stages
WHERE workflow_template_id = {placeholder}
ORDER BY position, created_at, id
""",
(str(workflow_id),),
).fetchall()
return [_workflow_stage_from_row(row) for row in rows]
def get(self, *, workflow_id: UUID, stage_id: UUID) -> WorkflowStageSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT {self._select_columns()}
FROM workflow_template_stages
WHERE workflow_template_id = {placeholder}
AND id = {placeholder}
""",
(str(workflow_id), str(stage_id)),
).fetchone()
if row is None:
raise LookupError(f"Workflow stage not found: {stage_id}")
return _workflow_stage_from_row(row)
def update(
self,
*,
workflow_id: UUID,
stage_id: UUID,
stable_key: str,
display_name: str,
description: str,
position: int,
owner_role: Role,
runner_profile_key: str,
required_inputs: list[str],
expected_outputs: list[str],
acceptance_criteria: list[str],
requires_human_approval: bool,
retry_policy: JsonObject,
parts: list[dict[str, Any]],
updated_by: UUID,
updated_at: datetime,
) -> WorkflowStageSummary:
existing = self.get(workflow_id=workflow_id, stage_id=stage_id)
update_position = self._clamped_update_position(
workflow_id=workflow_id,
requested_position=position,
)
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
if update_position < existing.position:
connection.execute(
f"""
UPDATE workflow_template_stages
SET position = position + 1
WHERE workflow_template_id = {placeholder}
AND position >= {placeholder}
AND position < {placeholder}
""",
(str(workflow_id), update_position, existing.position),
)
elif update_position > existing.position:
connection.execute(
f"""
UPDATE workflow_template_stages
SET position = position - 1
WHERE workflow_template_id = {placeholder}
AND position <= {placeholder}
AND position > {placeholder}
""",
(str(workflow_id), update_position, existing.position),
)
connection.execute(
f"""
UPDATE workflow_template_stages
SET
stable_key = {placeholder},
display_name = {placeholder},
description = {placeholder},
position = {placeholder},
owner_role = {placeholder},
runner_profile_key = {placeholder},
required_inputs = {placeholder}{json_cast},
expected_outputs = {placeholder}{json_cast},
acceptance_criteria = {placeholder}{json_cast},
requires_human_approval = {placeholder},
retry_policy = {placeholder}{json_cast},
parts = {placeholder}{json_cast},
updated_at = {placeholder}
WHERE workflow_template_id = {placeholder}
AND id = {placeholder}
""",
(
stable_key,
display_name,
description,
update_position,
owner_role.value,
runner_profile_key,
_json_value(required_inputs),
_json_value(expected_outputs),
_json_value(acceptance_criteria),
_bool_value(requires_human_approval),
_json_value(retry_policy),
_json_value(parts),
_datetime_value(updated_at),
str(workflow_id),
str(stage_id),
),
)
self._touch_parent(
connection,
workflow_id=workflow_id,
updated_by=updated_by,
updated_at=updated_at,
)
return self.get(workflow_id=workflow_id, stage_id=stage_id)
def delete(
self,
*,
workflow_id: UUID,
stage_id: UUID,
updated_by: UUID,
updated_at: datetime,
) -> None:
existing = self.get(workflow_id=workflow_id, stage_id=stage_id)
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
connection.execute(
f"""
DELETE FROM workflow_template_stages
WHERE workflow_template_id = {placeholder}
AND id = {placeholder}
""",
(str(workflow_id), str(stage_id)),
)
connection.execute(
f"""
UPDATE workflow_template_stages
SET position = position - 1
WHERE workflow_template_id = {placeholder}
AND position > {placeholder}
""",
(str(workflow_id), existing.position),
)
self._touch_parent(
connection,
workflow_id=workflow_id,
updated_by=updated_by,
updated_at=updated_at,
)
def reorder(
self,
*,
workflow_id: UUID,
stage_ids: list[UUID],
updated_by: UUID,
updated_at: datetime,
) -> list[WorkflowStageSummary]:
existing = self.list_for_workflow(workflow_id)
existing_ids = {stage.id for stage in existing}
requested_ids = set(stage_ids)
if requested_ids != existing_ids or len(stage_ids) != len(existing):
raise ValueError("Stage reorder must include each workflow stage exactly once.")
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
for position, stage_id in enumerate(stage_ids, start=1):
connection.execute(
f"""
UPDATE workflow_template_stages
SET position = {placeholder}, updated_at = {placeholder}
WHERE workflow_template_id = {placeholder}
AND id = {placeholder}
""",
(
position,
_datetime_value(updated_at),
str(workflow_id),
str(stage_id),
),
)
self._touch_parent(
connection,
workflow_id=workflow_id,
updated_by=updated_by,
updated_at=updated_at,
)
return self.list_for_workflow(workflow_id)
def _clamped_insert_position(
self,
*,
workflow_id: UUID,
requested_position: int | None,
) -> int:
count = len(self.list_for_workflow(workflow_id))
if requested_position is None:
return count + 1
return max(1, min(requested_position, count + 1))
def _clamped_update_position(
self,
*,
workflow_id: UUID,
requested_position: int,
) -> int:
count = len(self.list_for_workflow(workflow_id))
return max(1, min(requested_position, count))
def _touch_parent(
self,
connection: Any,
*,
workflow_id: UUID,
updated_by: UUID,
updated_at: datetime,
) -> None:
placeholder = self._repository.placeholder()
connection.execute(
f"""
UPDATE workflow_templates
SET updated_by = {placeholder}, updated_at = {placeholder}
WHERE id = {placeholder}
""",
(str(updated_by), _datetime_value(updated_at), str(workflow_id)),
)
def _select_columns(self) -> str:
return """
id,
workflow_template_id,
stable_key,
display_name,
description,
position,
owner_role,
runner_profile_key,
required_inputs,
expected_outputs,
acceptance_criteria,
requires_human_approval,
retry_policy,
parts,
created_at,
updated_at
"""
class WorkflowTemplateAuditEventsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
def create(
self,
*,
workflow_id: UUID,
event_type: str,
actor_user_id: UUID | None,
payload: JsonObject,
created_at: datetime,
) -> WorkflowTemplateAuditEventSummary:
event_id = uuid4()
placeholder = self._repository.placeholder()
json_cast = self._repository.json_cast()
with self._repository.connection() as connection:
connection.execute(
f"""
INSERT INTO workflow_template_events (
id,
workflow_template_id,
event_type,
actor_user_id,
payload,
created_at
)
VALUES (
{placeholder},
{placeholder},
{placeholder},
{placeholder},
{placeholder}{json_cast},
{placeholder}
)
""",
(
str(event_id),
str(workflow_id),
event_type,
_uuid_value(actor_user_id),
_json_value(payload),
_datetime_value(created_at),
),
)
return self.get(event_id)
def get(self, event_id: UUID) -> WorkflowTemplateAuditEventSummary:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
row = connection.execute(
f"""
SELECT
id,
workflow_template_id,
event_type,
actor_user_id,
payload,
created_at
FROM workflow_template_events
WHERE id = {placeholder}
""",
(str(event_id),),
).fetchone()
if row is None:
raise LookupError(f"Workflow template event not found: {event_id}")
return _workflow_template_event_from_row(row)
def list_for_workflow(
self,
workflow_id: UUID,
) -> list[WorkflowTemplateAuditEventSummary]:
placeholder = self._repository.placeholder()
with self._repository.connection() as connection:
rows = connection.execute(
f"""
SELECT
id,
workflow_template_id,
event_type,
actor_user_id,
payload,
created_at
FROM workflow_template_events
WHERE workflow_template_id = {placeholder}
ORDER BY created_at, id
""",
(str(workflow_id),),
).fetchall()
return [_workflow_template_event_from_row(row) for row in rows]
class ScriptConfigVersionsRepository:
def __init__(self, repository: BackendRepository) -> None:
self._repository = repository
@@ -3407,6 +4131,63 @@ def _agent_job_summary_from_row(row: Any) -> AgentJobSummary:
)
def _workflow_template_from_row(
row: Any,
*,
stages: list[WorkflowStageSummary],
) -> WorkflowTemplateSummary:
return WorkflowTemplateSummary(
id=_row_value(row, "id"),
name=_row_value(row, "name"),
slug=_row_value(row, "slug"),
description=_row_value(row, "description"),
status=_row_value(row, "status"),
version=_row_value(row, "version"),
created_by=_row_value(row, "created_by"),
updated_by=_row_value(row, "updated_by"),
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
activated_at=_row_value(row, "activated_at"),
archived_at=_row_value(row, "archived_at"),
stages=stages,
)
def _workflow_stage_from_row(row: Any) -> WorkflowStageSummary:
return WorkflowStageSummary(
id=_row_value(row, "id"),
workflow_id=_row_value(row, "workflow_template_id"),
stable_key=_row_value(row, "stable_key"),
display_name=_row_value(row, "display_name"),
description=_row_value(row, "description"),
position=_row_value(row, "position"),
owner_role=_row_value(row, "owner_role"),
runner_profile_key=_row_value(row, "runner_profile_key"),
required_inputs=_json_from_row(row, "required_inputs"),
expected_outputs=_json_from_row(row, "expected_outputs"),
acceptance_criteria=_json_from_row(row, "acceptance_criteria"),
requires_human_approval=bool(_row_value(row, "requires_human_approval")),
retry_policy=_json_from_row(row, "retry_policy"),
parts=[
WorkflowStagePart.model_validate(item)
for item in _json_from_row(row, "parts")
],
created_at=_row_value(row, "created_at"),
updated_at=_row_value(row, "updated_at"),
)
def _workflow_template_event_from_row(row: Any) -> WorkflowTemplateAuditEventSummary:
return WorkflowTemplateAuditEventSummary(
id=_row_value(row, "id"),
workflow_id=_row_value(row, "workflow_template_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 _script_config_version_event_from_row(row: Any) -> dict[str, Any]:
return {
"id": _row_value(row, "id"),
@@ -3431,6 +4212,12 @@ def _plain_row(row: Any) -> dict[str, Any]:
"visual_rules",
"source_rules",
"publishing_rules",
"required_inputs",
"expected_outputs",
"acceptance_criteria",
"retry_policy",
"parts",
"payload",
):
if key in result:
result[key] = _json_decode(result[key])
+124
View File
@@ -17,6 +17,7 @@ from src.domain.contracts import (
PublishingStatus,
Role,
ScriptConfigVersionStatus,
WorkflowTemplateStatus,
)
@@ -30,6 +31,9 @@ PUBLISHING_STATUS_VALUES = _quoted_values(status.value for status in PublishingS
SCRIPT_CONFIG_STATUS_VALUES = _quoted_values(
status.value for status in ScriptConfigVersionStatus
)
WORKFLOW_TEMPLATE_STATUS_VALUES = _quoted_values(
status.value for status in WorkflowTemplateStatus
)
PLAN_REVIEW_STATUS_VALUES = _quoted_values(status.value for status in PlanReviewStatus)
CLAIM_SUPPORT_STATUS_VALUES = _quoted_values(
status.value for status in ClaimSupportStatus
@@ -94,6 +98,53 @@ POSTGRES_SCHEMA_STATEMENTS: tuple[str, ...] = (
)
""",
f"""
CREATE TABLE IF NOT EXISTS workflow_templates (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name TEXT NOT NULL,
slug TEXT NOT NULL UNIQUE,
description TEXT NOT NULL,
status TEXT NOT NULL CHECK (status IN ({WORKFLOW_TEMPLATE_STATUS_VALUES})),
version INTEGER NOT NULL,
created_by UUID NOT NULL REFERENCES users(id),
updated_by UUID REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
activated_at TIMESTAMPTZ,
archived_at TIMESTAMPTZ
)
""",
"""
CREATE TABLE IF NOT EXISTS workflow_template_stages (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
workflow_template_id UUID NOT NULL REFERENCES workflow_templates(id) ON DELETE CASCADE,
stable_key TEXT NOT NULL,
display_name TEXT NOT NULL,
description TEXT NOT NULL,
position INTEGER NOT NULL,
owner_role TEXT NOT NULL CHECK (owner_role IN ('ADMIN', 'EDITOR')),
runner_profile_key TEXT NOT NULL,
required_inputs JSONB NOT NULL DEFAULT '[]'::jsonb,
expected_outputs JSONB NOT NULL DEFAULT '[]'::jsonb,
acceptance_criteria JSONB NOT NULL DEFAULT '[]'::jsonb,
requires_human_approval BOOLEAN NOT NULL DEFAULT false,
retry_policy JSONB NOT NULL DEFAULT '{}'::jsonb,
parts JSONB NOT NULL DEFAULT '[]'::jsonb,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (workflow_template_id, stable_key)
)
""",
f"""
CREATE TABLE IF NOT EXISTS workflow_template_events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
workflow_template_id UUID NOT NULL REFERENCES workflow_templates(id) ON DELETE CASCADE,
event_type TEXT NOT NULL,
actor_user_id UUID REFERENCES users(id),
payload JSONB NOT NULL DEFAULT '{{}}'::jsonb,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
""",
f"""
CREATE TABLE IF NOT EXISTS articles (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
target_site_id UUID NOT NULL REFERENCES target_sites(id),
@@ -400,6 +451,19 @@ POSTGRES_SCHEMA_STATEMENTS: tuple[str, ...] = (
)
""",
"CREATE INDEX IF NOT EXISTS idx_target_sites_slug ON target_sites (slug)",
"CREATE INDEX IF NOT EXISTS idx_workflow_templates_slug ON workflow_templates (slug)",
"""
CREATE INDEX IF NOT EXISTS idx_workflow_templates_status_updated
ON workflow_templates (status, updated_at DESC)
""",
"""
CREATE INDEX IF NOT EXISTS idx_workflow_template_stages_order
ON workflow_template_stages (workflow_template_id, position)
""",
"""
CREATE INDEX IF NOT EXISTS idx_workflow_template_events_history
ON workflow_template_events (workflow_template_id, created_at)
""",
"""
CREATE INDEX IF NOT EXISTS idx_articles_dashboard
ON articles (target_site_id, status, publishing_status, updated_at DESC)
@@ -496,6 +560,53 @@ SQLITE_SCHEMA_STATEMENTS: tuple[str, ...] = (
)
""",
"""
CREATE TABLE IF NOT EXISTS workflow_templates (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
slug TEXT NOT NULL UNIQUE,
description TEXT NOT NULL,
status TEXT NOT NULL,
version INTEGER NOT NULL,
created_by TEXT NOT NULL,
updated_by TEXT,
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
activated_at TEXT,
archived_at TEXT
)
""",
"""
CREATE TABLE IF NOT EXISTS workflow_template_stages (
id TEXT PRIMARY KEY,
workflow_template_id TEXT NOT NULL,
stable_key TEXT NOT NULL,
display_name TEXT NOT NULL,
description TEXT NOT NULL,
position INTEGER NOT NULL,
owner_role TEXT NOT NULL,
runner_profile_key TEXT NOT NULL,
required_inputs TEXT NOT NULL DEFAULT '[]',
expected_outputs TEXT NOT NULL DEFAULT '[]',
acceptance_criteria TEXT NOT NULL DEFAULT '[]',
requires_human_approval INTEGER NOT NULL DEFAULT 0,
retry_policy TEXT NOT NULL DEFAULT '{}',
parts TEXT NOT NULL DEFAULT '[]',
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE (workflow_template_id, stable_key)
)
""",
"""
CREATE TABLE IF NOT EXISTS workflow_template_events (
id TEXT PRIMARY KEY,
workflow_template_id TEXT NOT NULL,
event_type TEXT NOT NULL,
actor_user_id TEXT,
payload TEXT NOT NULL DEFAULT '{}',
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
)
""",
"""
CREATE TABLE IF NOT EXISTS articles (
id TEXT PRIMARY KEY,
target_site_id TEXT NOT NULL,
@@ -767,6 +878,19 @@ SQLITE_SCHEMA_STATEMENTS: tuple[str, ...] = (
)
""",
"CREATE INDEX IF NOT EXISTS idx_target_sites_slug ON target_sites (slug)",
"CREATE INDEX IF NOT EXISTS idx_workflow_templates_slug ON workflow_templates (slug)",
"""
CREATE INDEX IF NOT EXISTS idx_workflow_templates_status_updated
ON workflow_templates (status, updated_at)
""",
"""
CREATE INDEX IF NOT EXISTS idx_workflow_template_stages_order
ON workflow_template_stages (workflow_template_id, position)
""",
"""
CREATE INDEX IF NOT EXISTS idx_workflow_template_events_history
ON workflow_template_events (workflow_template_id, created_at)
""",
"""
CREATE INDEX IF NOT EXISTS idx_articles_dashboard
ON articles (target_site_id, status, publishing_status, updated_at)
+2
View File
@@ -20,6 +20,7 @@ from src.presentation.routes.plans import router as plans_router
from src.presentation.routes.publishing import router as publishing_router
from src.presentation.routes.reviews import router as reviews_router
from src.presentation.routes.sites import router as sites_router
from src.presentation.routes.workflow_templates import router as workflow_templates_router
app = FastAPI(title="AI Content Pipeline Backend")
@@ -36,6 +37,7 @@ app.include_router(publishing_router)
app.include_router(agent_jobs_router)
app.include_router(internal_agent_jobs_router)
app.include_router(sites_router)
app.include_router(workflow_templates_router)
@app.get("/health")
@@ -0,0 +1,257 @@
from __future__ import annotations
from uuid import UUID
from fastapi import APIRouter, Depends, HTTPException, status
from src.application.workflow_templates import (
activate_workflow_template,
add_workflow_stage,
archive_workflow_template,
create_workflow_template,
delete_workflow_stage,
get_workflow_template,
list_workflow_template_audit_events,
list_workflow_templates,
reorder_workflow_stages,
update_workflow_stage,
update_workflow_template,
)
from src.domain.auth import ADMIN_ROLES, EDITOR_OR_ADMIN_ROLES
from src.domain.contracts import (
CurrentUser,
WorkflowStageCreateRequest,
WorkflowStageReorderRequest,
WorkflowStageResponse,
WorkflowStageUpdateRequest,
WorkflowTemplateAuditEventListResponse,
WorkflowTemplateCreateRequest,
WorkflowTemplateListResponse,
WorkflowTemplateResponse,
WorkflowTemplateUpdateRequest,
)
from src.infrastructure.repositories import BackendRepository
from src.presentation.dependencies import get_repository, require_roles
router = APIRouter(prefix="/api", tags=["workflow_templates"])
@router.post(
"/admin/workflows",
response_model=WorkflowTemplateResponse,
status_code=status.HTTP_201_CREATED,
)
def post_workflow_template(
request: WorkflowTemplateCreateRequest,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateResponse:
return create_workflow_template(
repository,
request=request,
current_user=current_user,
)
@router.get("/admin/workflows", response_model=WorkflowTemplateListResponse)
def get_workflow_templates(
current_user: CurrentUser = Depends(require_roles(EDITOR_OR_ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateListResponse:
return list_workflow_templates(repository, current_user=current_user)
@router.get(
"/admin/workflows/{workflow_id}",
response_model=WorkflowTemplateResponse,
)
def get_workflow_template_detail(
workflow_id: UUID,
current_user: CurrentUser = Depends(require_roles(EDITOR_OR_ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateResponse:
try:
return get_workflow_template(
repository,
workflow_id=workflow_id,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
@router.patch(
"/admin/workflows/{workflow_id}",
response_model=WorkflowTemplateResponse,
)
def patch_workflow_template(
workflow_id: UUID,
request: WorkflowTemplateUpdateRequest,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateResponse:
try:
return update_workflow_template(
repository,
workflow_id=workflow_id,
request=request,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
except ValueError as error:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(error)) from error
@router.post(
"/admin/workflows/{workflow_id}/stages",
response_model=WorkflowStageResponse,
status_code=status.HTTP_201_CREATED,
)
def post_workflow_stage(
workflow_id: UUID,
request: WorkflowStageCreateRequest,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowStageResponse:
try:
return add_workflow_stage(
repository,
workflow_id=workflow_id,
request=request,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
except ValueError as error:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(error)) from error
@router.patch(
"/admin/workflows/{workflow_id}/stages/{stage_id}",
response_model=WorkflowStageResponse,
)
def patch_workflow_stage(
workflow_id: UUID,
stage_id: UUID,
request: WorkflowStageUpdateRequest,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowStageResponse:
try:
return update_workflow_stage(
repository,
workflow_id=workflow_id,
stage_id=stage_id,
request=request,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
except ValueError as error:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(error)) from error
@router.delete(
"/admin/workflows/{workflow_id}/stages/{stage_id}",
response_model=WorkflowTemplateResponse,
)
def delete_workflow_template_stage(
workflow_id: UUID,
stage_id: UUID,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateResponse:
try:
return delete_workflow_stage(
repository,
workflow_id=workflow_id,
stage_id=stage_id,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
except ValueError as error:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(error)) from error
@router.post(
"/admin/workflows/{workflow_id}/stages/reorder",
response_model=WorkflowTemplateResponse,
)
def post_reorder_workflow_stages(
workflow_id: UUID,
request: WorkflowStageReorderRequest,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateResponse:
try:
return reorder_workflow_stages(
repository,
workflow_id=workflow_id,
stage_ids=request.stage_ids,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
except ValueError as error:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(error)) from error
@router.post(
"/admin/workflows/{workflow_id}/activate",
response_model=WorkflowTemplateResponse,
)
def post_activate_workflow_template(
workflow_id: UUID,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateResponse:
try:
return activate_workflow_template(
repository,
workflow_id=workflow_id,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
except ValueError as error:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(error)) from error
@router.post(
"/admin/workflows/{workflow_id}/archive",
response_model=WorkflowTemplateResponse,
)
def post_archive_workflow_template(
workflow_id: UUID,
current_user: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateResponse:
try:
return archive_workflow_template(
repository,
workflow_id=workflow_id,
current_user=current_user,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
@router.get(
"/admin/workflows/{workflow_id}/audit",
response_model=WorkflowTemplateAuditEventListResponse,
)
def get_workflow_template_audit_events(
workflow_id: UUID,
_: CurrentUser = Depends(require_roles(ADMIN_ROLES)),
repository: BackendRepository = Depends(get_repository),
) -> WorkflowTemplateAuditEventListResponse:
try:
return list_workflow_template_audit_events(
repository,
workflow_id=workflow_id,
)
except LookupError as error:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not Found") from error
@@ -0,0 +1,391 @@
from __future__ import annotations
import sys
import tempfile
import unittest
from pathlib import Path
from fastapi.testclient import TestClient
BACKEND_ROOT = Path(__file__).resolve().parents[2]
sys.path.insert(0, str(BACKEND_ROOT))
from src.application.seed_data import seed_reference_data # noqa: E402
from src.infrastructure.repositories import open_backend_repository # noqa: E402
from src.presentation.dependencies import get_repository # noqa: E402
from src.presentation.main import app # noqa: E402
DEMO_ADMIN_EMAIL = "admin@example.com"
DEMO_EDITOR_EMAIL = "editor@example.com"
DEMO_USER_EMAIL_HEADER = "X-Demo-User-Email"
class AdminWorkflowTemplatesPublicApiTest(unittest.TestCase):
def setUp(self) -> None:
self.tmp_dir = tempfile.TemporaryDirectory()
dsn = f"sqlite:///{Path(self.tmp_dir.name) / 'workflow-templates.db'}"
self.repository = open_backend_repository(dsn)
self.repository.setup()
seed_reference_data(self.repository)
app.dependency_overrides[get_repository] = lambda: self.repository
self.client = TestClient(app)
def tearDown(self) -> None:
app.dependency_overrides.clear()
self.tmp_dir.cleanup()
def test_admin_creates_workflow_template_with_two_editable_ordered_stages_and_structured_parts(
self,
) -> None:
admin_headers = {DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}
create_response = self.client.post(
"/api/admin/workflows",
headers=admin_headers,
json=self._workflow_payload("article-production-v1"),
)
self.assertEqual(201, create_response.status_code, create_response.text)
workflow = create_response.json()["workflow"]
self.assertEqual("article-production-v1", workflow["slug"])
self.assertEqual("DRAFT", workflow["status"])
self.assertEqual(1, workflow["version"])
self.assertEqual([], workflow["stages"])
self.assertTrue(workflow["created_at"])
self.assertTrue(workflow["updated_at"])
intake_response = self.client.post(
f"/api/admin/workflows/{workflow['id']}/stages",
headers=admin_headers,
json=self._intake_stage_payload(),
)
self.assertEqual(201, intake_response.status_code, intake_response.text)
intake_stage = intake_response.json()["stage"]
draft_response = self.client.post(
f"/api/admin/workflows/{workflow['id']}/stages",
headers=admin_headers,
json=self._draft_stage_payload(),
)
self.assertEqual(201, draft_response.status_code, draft_response.text)
draft_stage = draft_response.json()["stage"]
edited_draft_parts = [
{
"key": "draft-outline",
"type": "outline",
"title": "Edited draft outline",
"payload": {
"prompt": "Build a sourced outline before drafting.",
"config": {
"minimum_sections": 5,
"require_source_placeholders": True,
},
},
"acceptance_criteria": [
"Every section has a purpose.",
"Claims that need evidence are marked.",
],
},
{
"key": "draft-body",
"type": "generation",
"title": "Draft body",
"payload": {
"prompt": "Generate a complete longform draft.",
"config": {
"tone": "practical",
"include_evidence_markers": True,
},
},
"acceptance_criteria": [
"Draft includes all approved outline sections.",
"Unsupported claims remain marked for review.",
],
},
]
patch_response = self.client.patch(
f"/api/admin/workflows/{workflow['id']}/stages/{draft_stage['id']}",
headers=admin_headers,
json={
"display_name": "Draft and evidence assembly",
"description": "Create an evidence-aware article draft.",
"owner_role": "EDITOR",
"runner_profile_key": "draft-writer-v2",
"required_inputs": [
"approved_plan",
"evidence_matrix",
],
"expected_outputs": [
"article_draft",
"claim_evidence_map",
],
"acceptance_criteria": [
"Draft follows the approved plan.",
"Evidence markers are preserved.",
],
"requires_human_approval": True,
"retry_policy": {
"max_attempts": 2,
"backoff_seconds": 120,
},
"parts": edited_draft_parts,
},
)
self.assertEqual(200, patch_response.status_code, patch_response.text)
detail_response = self.client.get(
f"/api/admin/workflows/{workflow['id']}",
headers=admin_headers,
)
self.assertEqual(200, detail_response.status_code, detail_response.text)
detail = detail_response.json()["workflow"]
stages = detail["stages"]
self.assertEqual([1, 2], [stage["position"] for stage in stages])
self.assertEqual(
["intake-boundary-questions", "draft-assembly"],
[stage["stable_key"] for stage in stages],
)
self.assertEqual(intake_stage["id"], stages[0]["id"])
self.assertEqual(draft_stage["id"], stages[1]["id"])
self.assertEqual(self._intake_stage_payload()["parts"], stages[0]["parts"])
self.assertEqual("Draft and evidence assembly", stages[1]["display_name"])
self.assertEqual("draft-writer-v2", stages[1]["runner_profile_key"])
self.assertTrue(stages[1]["requires_human_approval"])
self.assertEqual(
{
"max_attempts": 2,
"backoff_seconds": 120,
},
stages[1]["retry_policy"],
)
self.assertEqual(edited_draft_parts, stages[1]["parts"])
def test_editor_mutation_is_denied_with_existing_authorization_error_shape(
self,
) -> None:
response = self.client.post(
"/api/admin/workflows",
headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL},
json=self._workflow_payload("editor-denied-workflow"),
)
self.assertEqual(403, response.status_code, response.text)
self.assertEqual({"detail": "Forbidden"}, response.json())
def test_admin_reorders_activates_archives_and_audits_editor_visible_workflow(
self,
) -> None:
admin_headers = {DEMO_USER_EMAIL_HEADER: DEMO_ADMIN_EMAIL}
editor_headers = {DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}
workflow = self.client.post(
"/api/admin/workflows",
headers=admin_headers,
json=self._workflow_payload("article-production-ops"),
).json()["workflow"]
workflow_id = workflow["id"]
intake_stage = self.client.post(
f"/api/admin/workflows/{workflow_id}/stages",
headers=admin_headers,
json=self._intake_stage_payload(),
).json()["stage"]
draft_stage = self.client.post(
f"/api/admin/workflows/{workflow_id}/stages",
headers=admin_headers,
json=self._draft_stage_payload(),
).json()["stage"]
reorder_response = self.client.post(
f"/api/admin/workflows/{workflow_id}/stages/reorder",
headers=admin_headers,
json={"stage_ids": [draft_stage["id"], intake_stage["id"]]},
)
self.assertEqual(200, reorder_response.status_code, reorder_response.text)
self.assertEqual(
["draft-assembly", "intake-boundary-questions"],
[
stage["stable_key"]
for stage in reorder_response.json()["workflow"]["stages"]
],
)
activate_response = self.client.post(
f"/api/admin/workflows/{workflow_id}/activate",
headers=admin_headers,
)
self.assertEqual(200, activate_response.status_code, activate_response.text)
activated = activate_response.json()["workflow"]
self.assertEqual("ACTIVE", activated["status"])
self.assertEqual(2, activated["version"])
editor_list_response = self.client.get(
"/api/admin/workflows",
headers=editor_headers,
)
self.assertEqual(200, editor_list_response.status_code, editor_list_response.text)
self.assertIn(
workflow_id,
[item["id"] for item in editor_list_response.json()["workflows"]],
)
editor_detail_response = self.client.get(
f"/api/admin/workflows/{workflow_id}",
headers=editor_headers,
)
self.assertEqual(200, editor_detail_response.status_code, editor_detail_response.text)
immutable_patch_response = self.client.patch(
f"/api/admin/workflows/{workflow_id}/stages/{draft_stage['id']}",
headers=admin_headers,
json={"display_name": "Should not mutate active template"},
)
self.assertEqual(409, immutable_patch_response.status_code, immutable_patch_response.text)
archive_response = self.client.post(
f"/api/admin/workflows/{workflow_id}/archive",
headers=admin_headers,
)
self.assertEqual(200, archive_response.status_code, archive_response.text)
self.assertEqual("ARCHIVED", archive_response.json()["workflow"]["status"])
archived_editor_list_response = self.client.get(
"/api/admin/workflows",
headers=editor_headers,
)
self.assertEqual(
200,
archived_editor_list_response.status_code,
archived_editor_list_response.text,
)
self.assertNotIn(
workflow_id,
[item["id"] for item in archived_editor_list_response.json()["workflows"]],
)
archived_editor_detail_response = self.client.get(
f"/api/admin/workflows/{workflow_id}",
headers=editor_headers,
)
self.assertEqual(
404,
archived_editor_detail_response.status_code,
archived_editor_detail_response.text,
)
audit_response = self.client.get(
f"/api/admin/workflows/{workflow_id}/audit",
headers=admin_headers,
)
self.assertEqual(200, audit_response.status_code, audit_response.text)
event_types = {
event["event_type"] for event in audit_response.json()["events"]
}
self.assertTrue(
{
"WORKFLOW_TEMPLATE_CREATED",
"WORKFLOW_STAGE_CREATED",
"WORKFLOW_STAGES_REORDERED",
"WORKFLOW_TEMPLATE_ACTIVATED",
"WORKFLOW_TEMPLATE_ARCHIVED",
}.issubset(event_types)
)
def _workflow_payload(self, slug: str) -> dict[str, str]:
return {
"name": "Article production workflow",
"slug": slug,
"description": "Reusable editorial workflow for longform article production.",
}
def _intake_stage_payload(self) -> dict[str, object]:
return {
"stable_key": "intake-boundary-questions",
"display_name": "Boundary question intake",
"description": "Collect required editorial context before planning.",
"position": 1,
"owner_role": "EDITOR",
"runner_profile_key": "boundary-question-agent-v1",
"required_inputs": [
"article_brief",
"target_site",
],
"expected_outputs": [
"answered_boundary_questions",
],
"acceptance_criteria": [
"All required boundary questions are answered.",
"Answers are tied to the source brief.",
],
"requires_human_approval": True,
"retry_policy": {
"max_attempts": 1,
"backoff_seconds": 0,
},
"parts": [
{
"key": "required-context-checklist",
"type": "checklist",
"title": "Required context checklist",
"payload": {
"prompt": "Identify missing audience, keyword, and source constraints.",
"config": {
"required_fields": [
"audience",
"primary_keyword",
"source_rules",
],
"block_on_missing": True,
},
},
"acceptance_criteria": [
"Missing required context is listed explicitly.",
"No free-text parsing is needed to read the checklist.",
],
}
],
}
def _draft_stage_payload(self) -> dict[str, object]:
return {
"stable_key": "draft-assembly",
"display_name": "Draft assembly",
"description": "Generate a first draft from the approved plan.",
"position": 2,
"owner_role": "EDITOR",
"runner_profile_key": "draft-writer-v1",
"required_inputs": [
"approved_plan",
],
"expected_outputs": [
"article_draft",
],
"acceptance_criteria": [
"Draft follows the approved plan.",
],
"requires_human_approval": False,
"retry_policy": {
"max_attempts": 2,
"backoff_seconds": 60,
},
"parts": [
{
"key": "draft-outline",
"type": "outline",
"title": "Draft outline",
"payload": {
"prompt": "Build a concise outline before drafting.",
"config": {
"minimum_sections": 4,
"require_source_placeholders": True,
},
},
"acceptance_criteria": [
"Every section has a purpose.",
],
}
],
}
if __name__ == "__main__":
unittest.main()