studio / publisher /tasks /auto_publish.py
Ava2lon's picture
Upload 170 files
345855e verified
Raw
History Blame Contribute Delete
4.23 kB
import asyncio
from datetime import datetime
import uuid
# -------------------------------------------------
# CONTEXT NORMALIZER
# -------------------------------------------------
def normalize_context(context):
if isinstance(context, dict):
return {
"input_file": context.get("input_file"),
"url_input": context.get("url_input"),
"params": context.get("params", {}) or {},
"platforms": context.get("platforms", ["tiktok", "reels"])
}
return {
"input_file": getattr(context, "input_file", None),
"url_input": getattr(context, "url_input", None),
"params": getattr(context, "params", {}) or {},
"platforms": getattr(context, "platforms", ["tiktok", "reels"])
}
# -------------------------------------------------
# SAFE EXECUTOR WRAPPER (registry-first)
# -------------------------------------------------
async def run_task(task_name, payload):
"""
Uses V11 registry executor if available.
Falls back safely if not.
"""
try:
from core.execution.executor import execute_task as registry_execute
return await registry_execute(task_name, payload)
except Exception:
return {
"status": "failed",
"task": task_name,
"message": "registry executor unavailable"
}
# -------------------------------------------------
# MAIN AUTO-PUBLISH PIPELINE
# -------------------------------------------------
async def run(context):
ctx = normalize_context(context)
batch_id = str(uuid.uuid4())
started_at = datetime.utcnow().isoformat()
try:
# -------------------------------------------------
# STEP 1 — TRANSCRIBE
# -------------------------------------------------
transcript = await run_task("transcribe", ctx)
if transcript.get("status") != "success":
return {
"status": "error",
"stage": "transcribe",
"detail": transcript
}
# -------------------------------------------------
# STEP 2 — STRATEGY GENERATION
# -------------------------------------------------
strategy = await run_task("strategy", {
"text": transcript
})
# -------------------------------------------------
# STEP 3 — METADATA GENERATION
# -------------------------------------------------
metadata = await run_task("generate-metadata", {
"strategy": strategy
})
# -------------------------------------------------
# STEP 4 — VARIANT BUILDING
# -------------------------------------------------
platforms = ctx["platforms"]
variants = []
for platform in platforms:
variants.append({
"platform": platform,
"content": strategy,
"metadata": metadata
})
# -------------------------------------------------
# STEP 5 — PUBLISH
# -------------------------------------------------
publish_results = []
for v in variants:
result = await run_task("publish", {
"platform": v["platform"],
"content": v["content"],
"metadata": v["metadata"]
})
publish_results.append({
"platform": v["platform"],
"result": result
})
# -------------------------------------------------
# FINAL RESPONSE
# -------------------------------------------------
return {
"status": "success",
"pipeline": "auto_publish_v11",
"batch_id": batch_id,
"started_at": started_at,
"completed_at": datetime.utcnow().isoformat(),
"platforms": platforms,
"variants_created": len(variants),
"publish_results": publish_results
}
except Exception as e:
return {
"status": "error",
"pipeline": "auto_publish_v11",
"message": str(e),
"stage": "auto_publish_failed"
}