Ava2lon's picture
Upload 20 files
038574d verified
Raw
History Blame Contribute Delete
6.41 kB
from __future__ import annotations
import os
import threading
from datetime import datetime, timezone
from typing import Any
from ..config import settings
from .connectors import ConnectorError, publish_to_account, request_json, send_integration, sync_account_analytics
from .intelligence import score_analytics
from .models import AnalyticsEvent, AutomationJob, Integration, SocialAccount, Variant
from .operations import audit, notify
from .store import now_iso, store
def _save_job(data: dict[str, Any], state: str, result: dict[str, Any] | None = None, log: str | None = None) -> AutomationJob:
data["state"] = state
data["updated_at"] = now_iso()
if result is not None:
data["result"] = result
if log:
data.setdefault("logs", []).append(log)
job = AutomationJob.model_validate(data)
store.save("jobs", job)
return job
def _account(platform: str, workspace_id: str | None) -> SocialAccount:
accounts = [item for item in store.list("accounts", SocialAccount) if item.platform == platform and item.active]
if workspace_id:
accounts = [item for item in accounts if item.workspace_id == workspace_id]
if not accounts:
raise ConnectorError(f"No active {platform} account is connected")
return accounts[0]
def _execute_render(data: dict[str, Any]) -> dict[str, Any]:
payload = data["payload"]
base_url = os.getenv("MAESTER_INTERNAL_BASE_URL", "http://127.0.0.1:7860").rstrip("/")
result = request_json("POST", f"{base_url}/services/render/render", payload["render_payload"], token=None, allow_private=True)
variant = store.get("variants", payload["variant_id"], Variant)
metadata = {**variant.metadata, "render_job": result}
store.save("variants", variant.model_copy(update={"job_state": "rendering", "metadata": metadata, "updated_at": now_iso()}))
return result
def _execute_publish(data: dict[str, Any]) -> dict[str, Any]:
payload = data["payload"]
variant = store.get("variants", payload["variant_id"], Variant)
results = []
for platform in payload["platforms"]:
account = _account(platform, payload.get("workspace_id"))
results.append(publish_to_account(account, {**payload, "render": variant.metadata.get("render_job", {}), "variant": variant.model_dump()}))
published = all(item.get("status") == "published" for item in results)
state = "published" if published else variant.job_state
approval = "published" if published else variant.approval_state
store.save("variants", variant.model_copy(update={"job_state": state, "approval_state": approval, "updated_at": now_iso()}))
return {"platforms": results, "published": published}
def _execute_analytics(data: dict[str, Any]) -> dict[str, Any]:
payload = data["payload"]
account = store.get("accounts", payload["account_id"], SocialAccount)
incoming = sync_account_analytics(account)
saved = 0
for raw in incoming:
variant_id = raw.get("variant_id")
if not variant_id:
continue
try:
variant = store.get("variants", variant_id, Variant)
except KeyError:
continue
event = AnalyticsEvent.model_validate({**raw, "variant_id": variant_id, "platform": account.platform, "source": "connector"})
event_data = event.model_dump()
event_data.update({"id": f"analytics_{data['id']}_{saved}", "captured_at": event.captured_at or now_iso(), "score": score_analytics(event_data)})
store.save_data("analytics", event_data["id"], event_data)
store.save("variants", variant.model_copy(update={"analytics": {**variant.analytics, account.platform: event_data}, "updated_at": now_iso()}))
saved += 1
return {"account_id": account.id, "events_ingested": saved}
def _execute_notification(data: dict[str, Any]) -> dict[str, Any]:
payload = data["payload"]
integration = store.get("integrations", payload["integration_id"], Integration)
return send_integration(integration, payload.get("message", {}))
def execute_job(data: dict[str, Any]) -> AutomationJob:
try:
if data["kind"] == "render":
result = _execute_render(data)
return _save_job(data, "rendering", result, "Render submitted to the render engine.")
if data["kind"] == "publish":
result = _execute_publish(data)
state = "published" if result["published"] else "ready"
return _save_job(data, state, result, "Publishing connector dispatch completed.")
if data["kind"] == "analytics_sync":
return _save_job(data, "ready", _execute_analytics(data), "Analytics synchronization completed.")
if data["kind"] == "notification":
return _save_job(data, "ready", _execute_notification(data), "Notification delivered.")
raise RuntimeError(f"Unsupported job kind: {data['kind']}")
except Exception as exc:
attempts = int(data.get("attempts", 1))
retryable = attempts < int(data.get("max_attempts", settings.job_max_attempts)) and not isinstance(exc, ConnectorError)
state = "queued" if retryable else "failed"
job = _save_job(data, state, {"error": str(exc)}, f"Attempt {attempts} failed: {exc}")
if state == "failed":
notify("Automation job failed", f"{job.kind} job {job.id} failed: {exc}", metadata={"job_id": job.id})
audit("job.failed", "job", job.id, metadata={"error": str(exc)})
return job
def process_next(kinds: list[str] | None = None) -> AutomationJob | None:
data = store.claim_job(kinds)
return execute_job(data) if data else None
class JobWorker:
def __init__(self) -> None:
self._stop = threading.Event()
self._thread: threading.Thread | None = None
def start(self) -> None:
if self._thread and self._thread.is_alive():
return
self._stop.clear()
self._thread = threading.Thread(target=self._run, name="maester-automation-worker", daemon=True)
self._thread.start()
def stop(self) -> None:
self._stop.set()
if self._thread:
self._thread.join(timeout=5)
def _run(self) -> None:
while not self._stop.wait(2):
process_next()
@property
def running(self) -> bool:
return bool(self._thread and self._thread.is_alive())
worker = JobWorker()