Files
Agentswarm/orchestrator/task_queue.py
T
Songhaoz666andClaude Opus 4.8 d487923646 benchmark: 落地决策层(τ/η/P)、质量(Q_quality)、通信遥测;关闭 #10 #23
四块互相交织的 benchmark 覆盖增量,统一提交:

1) 通信遥测(#23):orchestrator 路由 peer 消息时按 correlation_id 计请求/应答到
   SwarmRun.collaboration(内部状态,不进 Manager 事件流);collector 算 s_communication。
   治理计数由 run.approvals 派生(合规/总数)→ s_governance。

2) Q_quality 掩码归一(v2.1 裁定):metrics.quality_score 改为对 present 输入加权归一,
   非编码任务自动忽略 TestPassRate,全缺 → NaN(不伪造)。

3) 质量插桩 / Group B:新增 Pod 内代码测试沙箱(orchestrator/sandbox.py,环境清洗 +
   超时强杀 + 资源限额 + 路径越界校验,门控 ENABLE_QUALITY_EVAL)与 held-out fixture
   (benchmark/fixtures/);run 完成时用留出测试评分得 TestPassRate → Q_quality →
   collector 合成 reward。安全边界见 docs/integration/security-boundary.md §8.1。

4) 决策引擎 / Group A(#10,Option A score-at-pull):新增 orchestrator/decision_engine.py
   —— 信息素 τ(Redis 持久、(role,agent) 键控、冷启动 0.5、ρ 蒸发、夹紧、学习常开)+
   η 启发式评分 + ε-greedy 概率采样;每次 dispatch 产一条 DecisionTrace →
   SwarmRun.decisions;collector 算 tau/eta/p_decision。概率选择门控 ENABLE_ACO_DISPATCH
   (默认关,CI 用 ACO_SEED 固定)。

覆盖:单次 run 真实可算字段由 4 提升至最多 10/15(新增 communication/reward/tau/eta/
p_decision,外加 governance 有条件)。

测试:新增 test-sandbox / test-quality / test-decision-engine;扩充 collector/metrics 用例;
CI 纳入全部 benchmark 套件 + flag-on 的 ACO e2e。本地 11 项 gate 全绿。

诚实边界(未越界声称):
- Group A 为单边匹配(Option B 待 Group C);概率派发优于贪心未证;默认关闭。
- reward 的 CodeReview/UserAcceptance 未采集(掩码忽略);P_risk 为审批派生低估。
- s_gain/s_swarm/g_e/g_e_cost/benchmark 仍 NaN —— 需基线(#21/#13),本 PR 不动验收。

影响范围:Swarm(orchestrator + benchmark + docs + CI)。不改 Manager↔Swarm 事件契约
(遥测均为运行时内部状态);不影响 Client/计费/密钥/发布链路。新增 ENABLE_QUALITY_EVAL /
ENABLE_ACO_DISPATCH 两个开关,默认关闭。

Closes #10
Closes #23

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 12:51:32 +08:00

542 lines
19 KiB
Python

"""Task queue management with dependency-aware failure recovery."""
import json
import time
import uuid
import logging
from typing import Dict, List, Optional
from enum import Enum
from pydantic import BaseModel, Field
from .redis_client import redis_client
from .agent_registry import agent_registry, AgentStatus
logger = logging.getLogger(__name__)
class TaskStatus(str, Enum):
"""Task status enumeration."""
PENDING = "pending"
ASSIGNED = "assigned"
IN_PROGRESS = "in_progress"
BLOCKED = "blocked"
COMPLETED = "completed"
FAILED = "failed"
CANCELLED = "cancelled"
class Task(BaseModel):
"""Task model."""
task_id: str
title: Optional[str] = None
description: str
status: TaskStatus
agent_role: str = "general"
required_capabilities: List[str] = Field(default_factory=list)
depends_on: List[str] = Field(default_factory=list)
parent_task_id: Optional[str] = None
root_task_id: Optional[str] = None
source: str = "manual"
assigned_agent_id: Optional[str] = None
created_at: float
started_at: Optional[float] = None
completed_at: Optional[float] = None
result: Optional[str] = None
blocked_reason: Optional[str] = None
child_task_ids: List[str] = Field(default_factory=list)
retry_count: int = 0
max_retries: int = 3
context: Dict = Field(default_factory=dict)
class TaskQueue:
"""Manages task assignment and reassignment with failure recovery."""
TASK_KEY_PREFIX = "task:"
PENDING_QUEUE_KEY = "queue:pending"
AGENT_TASK_KEY_PREFIX = "agent_task:"
def __init__(self):
pass
async def _save_task(self, task: Task):
"""Persist a task snapshot."""
key = f"{self.TASK_KEY_PREFIX}{task.task_id}"
await redis_client.set(key, task.model_dump_json())
async def create_task(
self,
description: str,
context: Optional[Dict] = None,
max_retries: int = 3,
task_id: Optional[str] = None,
title: Optional[str] = None,
agent_role: str = "general",
required_capabilities: Optional[List[str]] = None,
depends_on: Optional[List[str]] = None,
parent_task_id: Optional[str] = None,
root_task_id: Optional[str] = None,
source: str = "manual",
enqueue: bool = True,
) -> Task:
"""Create a new task and add to pending queue."""
task = Task(
task_id=task_id or str(uuid.uuid4()),
title=title,
description=description,
status=TaskStatus.PENDING,
agent_role=agent_role,
required_capabilities=required_capabilities or [],
depends_on=depends_on or [],
parent_task_id=parent_task_id,
root_task_id=root_task_id,
source=source,
created_at=time.time(),
max_retries=max_retries,
context=context or {},
)
await self._save_task(task)
if enqueue:
await redis_client.lpush(self.PENDING_QUEUE_KEY, task.task_id)
logger.info(f"Created task {task.task_id}: {description}")
return task
async def add_child_task(self, parent_task_id: str, child_task_id: str):
"""Register a child task on a parent task."""
parent = await self.get_task(parent_task_id)
if not parent:
return
if child_task_id not in parent.child_task_ids:
parent.child_task_ids.append(child_task_id)
await self._save_task(parent)
async def get_ready_pending_task(
self,
agent_capabilities: Optional[List[str]] = None,
) -> Optional[Task]:
"""Return and dequeue the next dispatchable task for an agent."""
pending_ids = await redis_client.lrange(self.PENDING_QUEUE_KEY, 0, -1)
capabilities = set(agent_capabilities or [])
for task_id in pending_ids:
task = await self.get_task(task_id)
if not task:
await self.remove_pending_task(task_id)
continue
if task.status != TaskStatus.PENDING:
await self.remove_pending_task(task_id)
continue
if not await self.is_task_ready(task):
continue
if not self.can_agent_run_task(task, capabilities):
continue
await self.remove_pending_task(task_id)
return task
return None
async def get_ready_pending_tasks(
self,
agent_capabilities: Optional[List[str]] = None,
) -> List[Task]:
"""Return ALL dispatchable tasks for an agent WITHOUT dequeuing any.
Candidate enumeration for the ACO decision engine (score-at-pull): the caller
scores/samples one and removes it via remove_pending_task. Dead/stale queue
entries are cleaned up the same way get_ready_pending_task does.
"""
pending_ids = await redis_client.lrange(self.PENDING_QUEUE_KEY, 0, -1)
capabilities = set(agent_capabilities or [])
candidates: List[Task] = []
for task_id in pending_ids:
task = await self.get_task(task_id)
if not task:
await self.remove_pending_task(task_id)
continue
if task.status != TaskStatus.PENDING:
await self.remove_pending_task(task_id)
continue
if not await self.is_task_ready(task):
continue
if not self.can_agent_run_task(task, capabilities):
continue
candidates.append(task)
return candidates
async def is_task_ready(self, task: Task) -> bool:
"""Return True when all dependencies are terminal and successful."""
if task.status != TaskStatus.PENDING:
return False
for dependency_id in task.depends_on:
dependency = await self.get_task(dependency_id)
if not dependency or dependency.status != TaskStatus.COMPLETED:
return False
return True
def can_agent_run_task(self, task: Task, capabilities: set[str]) -> bool:
"""Return whether an agent capability set satisfies task requirements."""
required = set(task.required_capabilities or [])
if not required:
return True
return required.issubset(capabilities)
async def assign_task(self, task_id: str, agent_id: str) -> bool:
"""Assign a task to an agent."""
# Get task
task = await self.get_task(task_id)
if not task:
logger.error(f"Task {task_id} not found")
return False
if task.status not in [TaskStatus.PENDING, TaskStatus.FAILED]:
logger.error(f"Task {task_id} cannot be assigned (status: {task.status})")
return False
if task.status == TaskStatus.PENDING and not await self.is_task_ready(task):
logger.error(f"Task {task_id} cannot be assigned before dependencies complete")
return False
# Verify agent is idle
agent = await agent_registry.get_agent(agent_id)
if not agent or agent.status != AgentStatus.IDLE:
logger.error(f"Agent {agent_id} is not available for task assignment")
return False
if not self.can_agent_run_task(task, set(agent.capabilities)):
logger.error(f"Agent {agent_id} does not satisfy task {task_id} capabilities")
return False
# Update task
task.status = TaskStatus.ASSIGNED
task.assigned_agent_id = agent_id
task.started_at = time.time()
task.blocked_reason = None
await self._save_task(task)
# Track agent's current task
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{agent_id}"
await redis_client.set(agent_task_key, task_id)
# Update agent status
await agent_registry.update_status(agent_id, AgentStatus.BUSY, task_id)
logger.info(f"Assigned task {task_id} to agent {agent_id}")
return True
async def start_task(self, task_id: str) -> bool:
"""Mark task as in progress."""
task = await self.get_task(task_id)
if not task:
return False
task.status = TaskStatus.IN_PROGRESS
task.blocked_reason = None
await self._save_task(task)
logger.info(f"Task {task_id} started")
return True
async def block_task(
self,
task_id: str,
reason: str = "",
release_agent: bool = True,
) -> Optional[Task]:
"""Mark a task blocked while waiting for delegated child work."""
task = await self.get_task(task_id)
if not task:
return None
if task.status in {TaskStatus.COMPLETED, TaskStatus.FAILED, TaskStatus.CANCELLED}:
return task
previous_agent_id = task.assigned_agent_id
task.status = TaskStatus.BLOCKED
task.blocked_reason = reason or "Waiting for delegated handoff work"
task.assigned_agent_id = None if release_agent else task.assigned_agent_id
await self.remove_pending_task(task_id)
await self._save_task(task)
if release_agent and previous_agent_id:
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{previous_agent_id}"
await redis_client.delete(agent_task_key)
await agent_registry.update_status(previous_agent_id, AgentStatus.IDLE)
logger.info(f"Task {task_id} blocked: {task.blocked_reason}")
return task
async def complete_task(self, task_id: str, result: str = None) -> bool:
"""Mark task as completed."""
task = await self.get_task(task_id)
if not task:
return False
if task.status == TaskStatus.CANCELLED:
logger.warning(f"Ignoring completion for cancelled task {task_id}")
return False
task.status = TaskStatus.COMPLETED
task.completed_at = time.time()
task.blocked_reason = None
if result:
task.result = result
await self._save_task(task)
# Clear agent's current task
if task.assigned_agent_id:
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{task.assigned_agent_id}"
await redis_client.delete(agent_task_key)
# Update agent to idle
await agent_registry.update_status(task.assigned_agent_id, AgentStatus.IDLE)
logger.info(f"Task {task_id} completed")
return True
async def fail_task(self, task_id: str, reason: str = "") -> bool:
"""Mark task as failed and handle retry logic."""
task = await self.get_task(task_id)
if not task:
return False
if task.status == TaskStatus.CANCELLED:
logger.warning(f"Ignoring failure for cancelled task {task_id}: {reason}")
return False
task.retry_count += 1
# Clear agent's current task
if task.assigned_agent_id:
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{task.assigned_agent_id}"
await redis_client.delete(agent_task_key)
# Update agent to idle
await agent_registry.update_status(task.assigned_agent_id, AgentStatus.IDLE)
# Check if we should retry
if task.retry_count < task.max_retries:
task.status = TaskStatus.PENDING
task.assigned_agent_id = None
task.started_at = None
# Re-add to pending queue
await redis_client.lpush(self.PENDING_QUEUE_KEY, task_id)
logger.warning(
f"Task {task_id} failed (retry {task.retry_count}/{task.max_retries}): {reason}"
)
else:
task.status = TaskStatus.FAILED
task.completed_at = time.time()
logger.error(
f"Task {task_id} permanently failed after {task.retry_count} retries: {reason}"
)
task.blocked_reason = None if task.status == TaskStatus.PENDING else task.blocked_reason
await self._save_task(task)
return True
async def cancel_task(self, task_id: str, reason: str = "") -> bool:
"""Cancel a task without retrying it."""
task = await self.get_task(task_id)
if not task:
return False
if task.status in {TaskStatus.COMPLETED, TaskStatus.FAILED, TaskStatus.CANCELLED}:
logger.info(f"Task {task_id} already terminal ({task.status}); skip cancellation")
return False
if task.assigned_agent_id:
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{task.assigned_agent_id}"
await redis_client.delete(agent_task_key)
await self.remove_pending_task(task_id)
task.status = TaskStatus.CANCELLED
task.completed_at = time.time()
task.blocked_reason = None
if reason:
task.result = json.dumps({"cancelled": True, "reason": reason})
await self._save_task(task)
logger.info(f"Task {task_id} cancelled: {reason}")
return True
async def release_task(self, task_id: str, agent_id: Optional[str] = None) -> bool:
"""Return an assigned/in-progress task to the pending queue (e.g. agent rejected it).
Unlike fail_task this does not increment retry_count: a capacity rejection is not
a task failure, just a dispatch that needs to find a different agent.
"""
task = await self.get_task(task_id)
if not task:
return False
if task.status not in {TaskStatus.ASSIGNED, TaskStatus.IN_PROGRESS}:
return False
if agent_id and task.assigned_agent_id and task.assigned_agent_id != agent_id:
return False
previous_agent_id = task.assigned_agent_id
task.status = TaskStatus.PENDING
task.assigned_agent_id = None
task.started_at = None
task.blocked_reason = None
await self._save_task(task)
await self.requeue_task(task_id)
if previous_agent_id:
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{previous_agent_id}"
await redis_client.delete(agent_task_key)
await agent_registry.update_status(previous_agent_id, AgentStatus.IDLE)
logger.info(f"Released task {task_id} back to pending queue")
return True
async def reopen_task(self, task_id: str) -> bool:
"""Re-open a completed task for another round (review loop rejected its result).
Resets the task to PENDING and requeues it without touching retry_count (a review
rejection is a quality decision, not a failure). The review-cycle budget in the
orchestrator bounds how many times this can happen.
"""
task = await self.get_task(task_id)
if not task:
return False
if task.status not in {TaskStatus.COMPLETED, TaskStatus.FAILED}:
return False
if task.assigned_agent_id:
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{task.assigned_agent_id}"
await redis_client.delete(agent_task_key)
task.status = TaskStatus.PENDING
task.assigned_agent_id = None
task.started_at = None
task.completed_at = None
task.blocked_reason = None
await self._save_task(task)
await self.requeue_task(task_id)
logger.info(f"Re-opened task {task_id} for another review cycle")
return True
async def reassign_agent_tasks(self, failed_agent_id: str) -> List[str]:
"""Reassign all tasks from a failed agent back to pending queue."""
# Get agent's current task
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{failed_agent_id}"
task_id = await redis_client.get(agent_task_key)
reassigned_tasks = []
if task_id:
await self.fail_task(task_id, f"Agent {failed_agent_id} failed")
reassigned_tasks.append(task_id)
logger.info(
f"Reassigned {len(reassigned_tasks)} tasks from failed agent {failed_agent_id}"
)
return reassigned_tasks
async def recover_orphaned_tasks(
self,
active_agent_ids: set[str],
stale_after_seconds: int = 30,
) -> List[str]:
"""Recover assigned/in-progress tasks whose agents are no longer connected."""
current_time = time.time()
recoverable_statuses = {TaskStatus.ASSIGNED, TaskStatus.IN_PROGRESS}
recovered_tasks = []
for task in await self.get_all_tasks():
if task.status not in recoverable_statuses:
continue
if not task.assigned_agent_id:
continue
if task.assigned_agent_id in active_agent_ids:
continue
if task.started_at and current_time - task.started_at < stale_after_seconds:
continue
old_agent_id = task.assigned_agent_id
agent_task_key = f"{self.AGENT_TASK_KEY_PREFIX}{old_agent_id}"
await redis_client.delete(agent_task_key)
task.status = TaskStatus.PENDING
task.assigned_agent_id = None
task.started_at = None
if await self.is_task_ready(task):
await self._save_task(task)
await redis_client.lpush(self.PENDING_QUEUE_KEY, task.task_id)
else:
await self._save_task(task)
recovered_tasks.append(task.task_id)
logger.warning(
f"Recovered orphaned task {task.task_id} from inactive agent {old_agent_id}"
)
return recovered_tasks
async def get_task(self, task_id: str) -> Optional[Task]:
"""Get task by ID."""
key = f"{self.TASK_KEY_PREFIX}{task_id}"
data = await redis_client.get(key)
if not data:
return None
return Task.model_validate_json(data)
async def get_next_pending_task(self) -> Optional[Task]:
"""Get next pending task from queue."""
task_id = await redis_client.rpop(self.PENDING_QUEUE_KEY)
if not task_id:
return None
return await self.get_task(task_id)
async def requeue_task(self, task_id: str):
"""Put a task back on the pending queue."""
await redis_client.rpush(self.PENDING_QUEUE_KEY, task_id)
async def remove_pending_task(self, task_id: str):
"""Remove a task from the pending queue if present."""
await redis_client.lrem(self.PENDING_QUEUE_KEY, 0, task_id)
async def get_pending_count(self) -> int:
"""Get count of pending tasks."""
return await redis_client.llen(self.PENDING_QUEUE_KEY)
async def get_dependents(self, task_id: str) -> List[Task]:
"""Return tasks that directly depend on a task."""
return [
task for task in await self.get_all_tasks()
if task_id in task.depends_on
]
async def get_all_tasks(self, status: Optional[TaskStatus] = None) -> List[Task]:
"""Get all tasks, optionally filtered by status."""
pattern = f"{self.TASK_KEY_PREFIX}*"
keys = await redis_client.keys(pattern)
tasks = []
for key in keys:
data = await redis_client.get(key)
if data:
task = Task.model_validate_json(data)
if status is None or task.status == status:
tasks.append(task)
return tasks
# Global task queue instance
task_queue = TaskQueue()