apps/whatsapp_control_api/app/services/task_registry.py
text
from __future__ import annotations
import json
from pathlib import Path
from typing import Any
from ..core.config import Settings
def compose_registry_task_id(raw_task_id: str) -> str:
return f"wa__{raw_task_id}"
def sync_store_task_to_registry(task: dict[str, Any] | None, logs: list[dict[str, Any]], settings: Settings) -> None:
if not task:
return
raw_task_id = str(task.get("id") or "").strip()
if not raw_task_id:
return
registry_id = compose_registry_task_id(raw_task_id)
source_name, agent_name = _classify_channel(str(task.get("channel") or "whatsapp"))
payload = {
"id": registry_id,
"raw_id": raw_task_id,
"channel": str(task.get("channel") or "whatsapp"),
"source_name": source_name,
"agent_name": agent_name,
"project_name": task.get("project_name") or None,
"status": str(task.get("status") or "queued"),
"current_step": str(task.get("current_step") or "-"),
"progress_percent": int(task.get("progress_percent") or 0),
"prompt": str(task.get("prompt") or ""),
"result_summary": str(task.get("result_summary") or ""),
"error_summary": str(task.get("error_summary") or ""),
"created_at": task.get("created_at"),
"updated_at": task.get("updated_at"),
"started_at": task.get("started_at"),
"finished_at": task.get("finished_at"),
"logs": [
{
"level": str(entry.get("level") or "INFO"),
"message": str(entry.get("message") or ""),
"created_at": entry.get("created_at"),
}
for entry in logs[-24:]
],
"task_log_path": str((settings.task_logs_dir / raw_task_id / "agent.log").resolve()),
}
registry_dir = settings.gateway_registry_tasks_dir / registry_id
registry_dir.mkdir(parents=True, exist_ok=True)
(registry_dir / "task.json").write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
def _classify_channel(channel: str) -> tuple[str, str]:
normalized = channel.strip().lower()
if normalized == "web":
return "Web Dashboard", "Hermes Gateway Dashboard"
return "WhatsApp", "AI Assistants Control"