Files
Agentswarm/orchestrator/task_queue.py
T
gongzhiyongandClaude Opus 4.8 b394ce99d5 fix(swarm): reap orphan agents on terminal run + retry backoff + surface task errors
孤儿 agent 事故修复(2026-06-15 swarm-69e470561bd7-seed 跨 workspace 执行)。
不改去中心化认领逻辑(swarm_dispatch 不动),仅治理 agent pod 生命周期与可观测性:

- B1 终态回收: refresh_swarm_run_status 终态后 best-effort stop_launched 回收 pod+secret
  (此前仅 Manager stop 才回收,自然完成/失败的 run agent 残留 → 孤儿留在共享池抢别的 swarm 任务)
- B2 防驱逐: agent pod 加 karpenter.sh/do-not-disrupt(AGENT_POD_ALLOW_DISRUPTION=1 可关)
- stop_launched: backend=kubernetes 时按标签删,不再被内存集合 _k8s_swarms 门控(跨重启可靠)
- C 重试退避: Task.next_retry_at + fail_task 指数退避 5→30→180s(cap 300, TASK_RETRY_BACKOFF_*),
  is_task_ready 门控;TASK_MAX_RETRIES 可配
- D 错误上浮: task_executor 失败时聚合 subtask error 到顶层 error(根治通用 "Task failed"),
  _execute_subtask 打印 LLM 响应片段
- 配额硬上限 16: _clamp_user_cap

契约测试全过: runtime-contract / merge-smoke / workflow-e2e / contract-freeze /
max-agents-per-user / security-boundary。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-16 04:33:01 +08:00

575 lines
21 KiB
Python

"""Task queue management with dependency-aware failure recovery."""
import json
import os
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__)
def _retry_backoff_seconds(retry_count: int) -> float:
"""Delay before a failed task becomes dispatchable again (exponential, capped).
Defaults: 5s → 30s → 180s … capped at 300s. Spreads retries over minutes instead of burning
the whole budget in seconds, so a swarm's own agents (which may be cold-starting / waiting on
node scale-up) have time to register before the seed exhausts its retries (incident 2026-06-15).
Tunable via TASK_RETRY_BACKOFF_{BASE,FACTOR,CAP}.
"""
try:
base = float(os.getenv("TASK_RETRY_BACKOFF_BASE", "5") or 5)
factor = float(os.getenv("TASK_RETRY_BACKOFF_FACTOR", "6") or 6)
cap = float(os.getenv("TASK_RETRY_BACKOFF_CAP", "300") or 300)
except ValueError:
base, factor, cap = 5.0, 6.0, 300.0
exponent = max(0, retry_count - 1)
return min(cap, base * (factor ** exponent))
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
next_retry_at: Optional[float] = None
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
# Respect retry backoff: a task re-queued after a failure is not dispatchable until its
# next_retry_at has passed (see fail_task / _retry_backoff_seconds).
if task.next_retry_at and time.time() < task.next_retry_at:
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
# Backoff: gate re-dispatch until next_retry_at (is_task_ready enforces it) so retries
# spread over minutes rather than all firing within seconds.
delay = _retry_backoff_seconds(task.retry_count)
task.next_retry_at = time.time() + delay
# 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} "
f"(next retry in {delay:.0f}s)"
)
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
task.next_retry_at = None # a capacity release is not a failure — no backoff
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
task.next_retry_at = None # a review reopen is not a failure — no backoff
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()