Código generado · orchestrator.py — LeadFlow Onboarding Swarm
# orchestrator.py — Enjambre jerárquico para onboarding LeadFlow
# Generado por skill: orquestacion-enjambre-agentes | session: lf-2026-0892
from dataclasses import dataclass, field
from enum import Enum
import asyncio, redis.asyncio as redis, json, logging
logger = logging.getLogger("leadflow.swarm")
class AgentRole(Enum):
SCRAPER = "web_scraper"
PROFILER = "icp_profiler"
DESIGNER = "campaign_designer"
REVIEWER = "quality_reviewer"
REPORTER = "report_generator"
@dataclass
class AgentTask:
id: str
role: AgentRole
input_data: dict
output_data: dict = field(default_factory=dict)
status: str = "pending"
retries: int = 0
confidence: float = 0.0
class SharedMemory:
def __init__(self, session_id: str):
self.r = redis.from_url("redis://localhost:6379")
self.key = f"swarm:{session_id}"
async def set(self, field: str, value: dict):
await self.r.hset(self.key, field, json.dumps(value))
async def get(self, field: str) -> dict:
raw = await self.r.hget(self.key, field)
return json.loads(raw) if raw else {}
class QualityGate:
def check(self, stage: str, output: dict) -> bool:
if stage == "icp":
return output.get("confidence", 0) >= 0.75 and "sector" in output
if stage == "campaigns":
return len(output.get("campaigns", [])) >= 3
return True
class LeadFlowOrchestrator:
"""Coordina el pipeline jerárquico de onboarding para nuevos clientes."""
def __init__(self, agents: dict, session_id: str):
self.agents = agents
self.memory = SharedMemory(session_id)
self.gate = QualityGate()
self.tasks: list[AgentTask] = []
async def run_onboarding(self, client_url: str) -> dict:
# Stage 1: Scraping
scrape = await self._run(AgentRole.SCRAPER, {"url": client_url})
await self.memory.set("scrape", scrape)
# Stage 2: ICP con quality gate y retry
icp = await self._run_with_gate(AgentRole.PROFILER, {"scrape": scrape}, "icp")
await self.memory.set("icp", icp)
# Stage 3: Campañas en paralelo (3 variantes)
ctx = await self.memory.get("icp")
camps = await self._run_with_gate(AgentRole.DESIGNER, {"icp": ctx}, "campaigns")
await self.memory.set("campaigns", camps)
# Stage 4: Revisión de calidad
review = await self._run(AgentRole.REVIEWER, {
"icp": icp, "campaigns": camps
})
# Stage 5: Informe final
report = await self._run(AgentRole.REPORTER, {
"scrape": scrape, "icp": icp,
"campaigns": camps, "review": review
})
return {"status": "ok", "report": report}
async def _run_with_gate(self, role, inp, stage, max_retries=3) -> dict:
for attempt in range(max_retries):
result = await self._run(role, inp)
if self.gate.check(stage, result):
return result
logger.warning(f"Gate failed [{stage}] attempt {attempt+1} conf={result.get('confidence')}")
inp["feedback"] = "Increase specificity and confidence score"
raise RuntimeError(f"Quality gate [{stage}] no superado tras {max_retries} intentos")
async def _run(self, role, inp) -> dict:
agent = self.agents[role]
task = AgentTask(id=f"{role.value}_{len(self.tasks)}", role=role, input_data=inp)
self.tasks.append(task)
task.status = "running"
logger.info(f"[HANDOFF] → {role.value} | task_id={task.id}", extra={"task": task.id})
try:
result = await agent.execute(inp)
task.output_data, task.status = result, "completed"
return result
except Exception as e:
task.status = "failed"
logger.error(f"Agent {role.value} failed: {e}")
raise