from dataclasses import dataclass from datetime import timedelta from enum import StrEnum from temporalio import workflow from temporalio.common import RetryPolicy from temporalio.exceptions import ApplicationError from cartoonos.domain.content_attribution import ContentAttribution from cartoonos.domain.release_policy import PublicationIntent from cartoonos.ports.production_runtime import PublicationResult from cartoonos.ports.scene_to_screen import SceneToScreenPlanningResult from cartoonos.workflows.release_contract import FinalReleaseApproval, require_final_approval class ContentStage(StrEnum): IDEA = "IDEA" EVIDENCE_READY = "EVIDENCE_READY" SCRIPT_APPROVED = "SCRIPT_APPROVED" ANIMATIC_APPROVED = "ANIMATIC_APPROVED" CAST_READY = "CAST_READY" CHARACTER_PACKS_READY = "CHARACTER_PACKS_READY" SCENE_PLAN_APPROVED = "SCENE_PLAN_APPROVED" SHOTS_COMPILED = "SHOTS_COMPILED" GENERATED = "GENERATED" ROUGH_CUT_APPROVED = "ROUGH_CUT_APPROVED" PACKAGING_APPROVED = "PACKAGING_APPROVED" FINAL_QA = "FINAL_QA" PUBLISHED = "PUBLISHED" LEARNING_CAPTURED = "LEARNING_CAPTURED" @dataclass(frozen=True) class EpisodeWorkflowInput: """Generic project-scoped workflow input. Project-specific language, age, format and canon values intentionally have no global defaults. A project policy/frozen brief must resolve them before production can start. """ project_id: str episode_id: str content_brief_version: str series_id: str | None = None season_id: str | None = None target_language: str | None = None locale: str | None = None content_format: str | None = None target_age_min: int | None = None target_age_max: int | None = None sensitivity_tags: tuple[str, ...] = () fear_intensity: int = 0 imitation_risk: str | None = None publication_intent: PublicationIntent = PublicationIntent.INTERNAL_PREVIEW character_versions: tuple[str, ...] = () cast_mode: str = "" lead_character_id: str = "" # Character identity/profile versions and story roles are separate frozen inputs. cast_members: tuple[tuple[str, str], ...] = () cast_roles: tuple[tuple[str, str], ...] = () pairing_id: str | None = None character_bible_version: str | None = None relationship_bible_version: str | None = None memory_meaning_engine_version: str | None = None memory_meaning_block_ids: tuple[str, ...] = () content_pillar: str | None = None curiosity_hook_type: str | None = None viewer_participation_mechanic: str | None = None rewatch_mechanic: str | None = None next_curiosity_type: str | None = None def __post_init__(self) -> None: if not self.project_id.strip(): raise ValueError("project_id is required") if not self.episode_id.strip(): raise ValueError("episode_id is required") if self.target_age_min is not None and not 3 <= self.target_age_min <= 14: raise ValueError("target_age_min must be within 3..14") if self.target_age_max is not None and not 3 <= self.target_age_max <= 14: raise ValueError("target_age_max must be within 3..14") if ( self.target_age_min is not None and self.target_age_max is not None and self.target_age_min > self.target_age_max ): raise ValueError("target age band must be ordered") if self.fear_intensity not in {0, 1, 2, 3}: raise ValueError("fear_intensity must be 0..3") if self.imitation_risk is not None and self.imitation_risk not in {"low", "medium", "high"}: raise ValueError("imitation_risk must be low, medium or high") if self.memory_meaning_block_ids and not self.memory_meaning_engine_version: raise ValueError("memory_meaning_block_ids require memory_meaning_engine_version") def content_attribution(self) -> ContentAttribution: required = { "series_id": self.series_id, "season_id": self.season_id, "target_language": self.target_language, "locale": self.locale, "content_format": self.content_format, "target_age_min": self.target_age_min, "target_age_max": self.target_age_max, "imitation_risk": self.imitation_risk, "cast_mode": self.cast_mode, "lead_character_id": self.lead_character_id, "cast_members": self.cast_members, "character_bible_version": self.character_bible_version, "relationship_bible_version": self.relationship_bible_version, "content_pillar": self.content_pillar, "curiosity_hook_type": self.curiosity_hook_type, "viewer_participation_mechanic": self.viewer_participation_mechanic, "rewatch_mechanic": self.rewatch_mechanic, "next_curiosity_type": self.next_curiosity_type, } missing = sorted( key for key, value in required.items() if value is None or value == "" or value == () ) if missing: raise ValueError("production-locked attribution is missing: " + ", ".join(missing)) assert self.target_age_min is not None and self.target_age_max is not None return ContentAttribution( project_id=self.project_id, series_id=str(self.series_id), season_id=str(self.season_id), episode_id=self.episode_id, content_brief_version=self.content_brief_version, target_language=str(self.target_language), locale=str(self.locale), content_format=str(self.content_format), target_age_min=int(self.target_age_min), target_age_max=int(self.target_age_max), sensitivity_tags=self.sensitivity_tags, fear_intensity=self.fear_intensity, imitation_risk=str(self.imitation_risk), character_versions=self.character_versions, cast_mode=self.cast_mode, lead_character_id=self.lead_character_id, cast_members=self.cast_members, pairing_id=self.pairing_id, content_pillar=str(self.content_pillar), curiosity_hook_type=str(self.curiosity_hook_type), viewer_participation_mechanic=str(self.viewer_participation_mechanic), rewatch_mechanic=str(self.rewatch_mechanic), next_curiosity_type=str(self.next_curiosity_type), memory_meaning_engine_version=self.memory_meaning_engine_version, memory_meaning_block_ids=self.memory_meaning_block_ids, ) @dataclass(frozen=True) class EpisodeWorkflowResult: project_id: str episode_id: str final_stage: ContentStage SAFE_ACTIVITY_RETRY = RetryPolicy( initial_interval=timedelta(seconds=2), maximum_interval=timedelta(minutes=2), maximum_attempts=5, ) MUTATING_ACTIVITY_RETRY = RetryPolicy(maximum_attempts=1) METRIC_ACTIVITY_RETRY = RetryPolicy( initial_interval=timedelta(seconds=5), maximum_interval=timedelta(minutes=2), maximum_attempts=3, ) @workflow.defn(name="cartoonos.episode-production.v1") class EpisodeProductionWorkflow: """Durable project-scoped orchestration. External I/O belongs in activities.""" @workflow.run async def run(self, request: EpisodeWorkflowInput) -> EpisodeWorkflowResult: try: attribution = request.content_attribution() except ValueError as exc: raise ApplicationError( str(exc), type="InvalidProductionInput", non_retryable=True, ) from exc await workflow.execute_activity( "build_evidence_pack", request, start_to_close_timeout=timedelta(minutes=30), retry_policy=SAFE_ACTIVITY_RETRY, ) await workflow.execute_activity( "write_and_review_script", request, start_to_close_timeout=timedelta(hours=1), retry_policy=SAFE_ACTIVITY_RETRY, ) await workflow.execute_activity( "build_and_review_animatic", request, start_to_close_timeout=timedelta(hours=2), retry_policy=SAFE_ACTIVITY_RETRY, ) await workflow.execute_activity( "validate_cast_plan", request, start_to_close_timeout=timedelta(minutes=5), retry_policy=SAFE_ACTIVITY_RETRY, ) await workflow.execute_activity( "validate_character_production_readiness", request, start_to_close_timeout=timedelta(minutes=10), retry_policy=SAFE_ACTIVITY_RETRY, ) scene_plan = await workflow.execute_activity( "plan_scene_to_screen", request, result_type=SceneToScreenPlanningResult, start_to_close_timeout=timedelta(hours=1), retry_policy=SAFE_ACTIVITY_RETRY, ) compiled_requests = await workflow.execute_activity( "compile_shots_for_generation", (request, scene_plan), start_to_close_timeout=timedelta(minutes=30), retry_policy=SAFE_ACTIVITY_RETRY, ) generation_run_ids = await workflow.execute_activity( "generate_and_qa_media", (request, tuple(compiled_requests)), start_to_close_timeout=timedelta(hours=6), retry_policy=MUTATING_ACTIVITY_RETRY, ) await workflow.execute_activity( "build_and_review_rough_cut", (request, tuple(generation_run_ids)), start_to_close_timeout=timedelta(hours=2), retry_policy=SAFE_ACTIVITY_RETRY, ) await workflow.execute_activity( "build_and_review_platform_packaging", request, start_to_close_timeout=timedelta(hours=1), retry_policy=SAFE_ACTIVITY_RETRY, ) approval = await workflow.execute_activity( "review_final_release", request, result_type=FinalReleaseApproval, start_to_close_timeout=timedelta(minutes=10), retry_policy=SAFE_ACTIVITY_RETRY, ) publication_request = require_final_approval( approval, project_id=request.project_id, episode_id=request.episode_id, content_brief_version=request.content_brief_version, publication_intent=request.publication_intent, attribution=attribution, ) publication = await workflow.execute_activity( "publish_episode", publication_request, result_type=PublicationResult, start_to_close_timeout=timedelta(minutes=30), retry_policy=MUTATING_ACTIVITY_RETRY, ) await workflow.execute_activity( "schedule_metric_snapshots", publication, start_to_close_timeout=timedelta(minutes=10), retry_policy=METRIC_ACTIVITY_RETRY, ) return EpisodeWorkflowResult( project_id=request.project_id, episode_id=request.episode_id, final_stage=ContentStage.PUBLISHED, )