diff --git a/backend/app/api/admin.py b/backend/app/api/admin.py index d0a3601..a7e23ff 100644 --- a/backend/app/api/admin.py +++ b/backend/app/api/admin.py @@ -107,6 +107,22 @@ async def agent_available( return {"ok": True} +@router.delete("/agents/{agent_id}", summary="删除 Agent(无心跳的僵尸 Agent)") +async def agent_delete( + agent_id: str, + svc: AgentService = Depends(get_agent_service), + _admin: None = Depends(require_admin_dep), +): + """删除 Agent 池中的 Agent。 + + 用于清理无心跳的僵尸 Agent。若 Agent 进程仍存活,删除后其心跳/注册会 + 触发重新注册,自动重新连接回池中。 + """ + if not await svc.unregister(agent_id): + raise HTTPException(status_code=404, detail="agent not found") + return {"ok": True} + + @router.post("/agents/{agent_id}/priority", response_model=AgentInfo, summary="设置 Agent 优先级") async def agent_priority( agent_id: str, diff --git a/backend/app/repository/agent_repo.py b/backend/app/repository/agent_repo.py index a74376b..da5b4d9 100644 --- a/backend/app/repository/agent_repo.py +++ b/backend/app/repository/agent_repo.py @@ -23,12 +23,21 @@ class AgentRepo: # ---------- 写入 ---------- async def upsert(self, agent: AgentInfo, *, heartbeat: bool = False) -> bool: - """写入 Agent 信息,返回是否新建。非心跳时为全量注册。""" + """写入 Agent 信息,返回是否新建。非心跳时为全量注册。 + + 非心跳(全量注册)时会重建标签索引(删旧增新),避免 Agent 换标签后 + 索引陈旧导致 by_tags 匹配不到。 + """ key = self._info_key(agent.agent_id) existed = await self.redis.exists(key) if not existed: await self.redis.sadd(C.AGENT_ALL_SET, agent.agent_id) - # 建立标签索引 + if not heartbeat: + # 全量注册:先清理旧标签索引,再建立当前标签索引 + old = await self.get(agent.agent_id) + if old and old.agent_tags: + for tag in old.agent_tags: + await self.redis.srem(self._tag_key(tag), agent.agent_id) for tag in agent.agent_tags: await self.redis.sadd(self._tag_key(tag), agent.agent_id) mapping: dict[str, Any] = { @@ -93,12 +102,16 @@ class AgentRepo: await self.redis.hset(key, "status", AgentStatus.READY.value) return True - async def beat(self, agent_id: str, current_load: int) -> bool: - """更新心跳时间与负载,按状态机规则流转状态。 + async def beat(self, agent_id: str, current_load: int | None = None) -> bool: + """更新心跳时间,按状态机规则流转状态。 - 不可用(unavailable):保持不可用,不因心跳复活。 - 离线(offline,心跳失联):恢复 ready。 - ready/processing/stopping:保持原状态。 + + 注意:负载(current_load)一律由网关调度器通过 adjust_load 记账, + 心跳不再覆写,避免与调度器计数互相覆盖导致负载错乱。 + current_load 参数仅作兼容保留,不再写入。 """ key = self._info_key(agent_id) if not await self.redis.exists(key): @@ -117,7 +130,6 @@ class AgentRepo: key, mapping={ "last_heartbeat": str(now), - "current_load": str(current_load), "status": status.value, }, ) diff --git a/backend/app/services/relay.py b/backend/app/services/relay.py index c052f5b..a1259b3 100644 --- a/backend/app/services/relay.py +++ b/backend/app/services/relay.py @@ -22,19 +22,19 @@ class RelayService: self.task_repo = task_repo or TaskRepo(redis) self.agent_repo = AgentRepo(redis) - async def dispatch_command(self, agent_id: str, request_id: str, payload: dict) -> bool: + async def dispatch_command(self, agent_id: str, request_id: str, payload: dict) -> tuple[bool, str]: """正向:向 Agent 真实 HTTP 推送任务指令(POST {agent.endpoint}/tasks/{request_id})。 - 成功返回 True;Agent 不存在 / 端点缺失 / 网络失败 / 非 202 均返回 False, - 由调度器回退任务状态并释放负载。 + 返回 (ok, reason):成功时 (True, "");Agent 不存在 / 端点缺失 / 网络失败 / + 非 202 均返回 (False, 原因描述),由调度器回退任务状态并释放负载。 """ agent = await self.agent_repo.get(agent_id) if not agent: logger.warning("relay dispatch failed: agent not found agent=%s request=%s", agent_id, request_id) - return False + return False, "Agent 不存在或已注销" if not agent.endpoint: logger.warning("relay dispatch failed: endpoint empty agent=%s request=%s", agent_id, request_id) - return False + return False, "Agent 未配置 endpoint" task = await self.task_repo.get(request_id) url = f"{agent.endpoint.rstrip('/')}/tasks/{request_id}" body = { @@ -52,15 +52,15 @@ class RelayService: except httpx.HTTPError as e: logger.error("relay dispatch network error agent=%s request=%s err=%s", agent_id, request_id, e) await self._log(agent_id, request_id, "dispatch", f"push failed: {e}") - return False + return False, f"推送失败: {e}" if resp.status_code != 202: logger.warning("relay dispatch http %s agent=%s request=%s body=%s", resp.status_code, agent_id, request_id, resp.text[:200]) await self._log(agent_id, request_id, "dispatch", f"push failed http {resp.status_code}") - return False + return False, f"推送失败: Agent 返回 HTTP {resp.status_code}(期望 202)" await self._log(agent_id, request_id, "dispatch", f"command pushed to {agent.endpoint}") logger.info("relay dispatch ok agent=%s request=%s url=%s", agent_id, request_id, url) - return True + return True, "" _TERMINAL = {TaskStatus.SUCCESS, TaskStatus.FAILED} diff --git a/backend/app/services/scheduler.py b/backend/app/services/scheduler.py index 9f71bb5..4fff234 100644 --- a/backend/app/services/scheduler.py +++ b/backend/app/services/scheduler.py @@ -37,13 +37,23 @@ class Scheduler: return candidates[0] async def dispatch(self, request_id: str) -> bool: - """将单个 pending 任务下发给匹配 Agent(绑定 → 推送 payload → 失败回退)。""" + """将单个 pending 任务下发给匹配 Agent(绑定 → 推送 payload → 失败回退)。 + + 成功 / 无可用 Agent / 推送失败时,都会在任务上记录 error_info(成功时清空), + 便于前端与日志定位“为什么没分配”。 + """ task = await self.task_repo.get(request_id) if not task or task.status != TaskStatus.PENDING: return False agent = await self.pick_candidate(task.task_tags) if not agent: - return False # 暂无可用 Agent,保持 pending,等待下次调度 + # 暂无可用 Agent(标签不匹配 / 离线 / 满载),保持 pending,等待下次调度 + await self.task_repo.update( + request_id, + error_info="无可用 Agent:标签不匹配或全部离线/满载", + ) + logger.info("task no candidate request_id=%s tags=%s", request_id, task.task_tags) + return False # 绑定 Agent,标记 running await self.task_repo.update( request_id, @@ -52,20 +62,33 @@ class Scheduler: progress=0, ) await self.task_repo.mark_running(request_id) + # 记录调度前状态,供推送失败回退时恢复 + prev_status = agent.status # 标记 Agent 为处理中 await self.agent_repo.update_status(agent.agent_id, AgentStatus.PROCESSING) # 更新 Agent 负载 await self.agent_repo.adjust_load(agent.agent_id, 1) # 真实推送 payload 到 Agent 端点 - pushed = await self.relay.dispatch_command(agent.agent_id, request_id, task.payload or {}) + pushed, reason = await self.relay.dispatch_command(agent.agent_id, request_id, task.payload or {}) if not pushed: - # 推送失败:回退 pending 并恢复 agent 就绪、释放负载,等待下次调度(避免重复扣负载) - await self.task_repo.update(request_id, status=TaskStatus.PENDING, agent_id="", progress=0) + # 推送失败:回退 pending 并释放负载,按原状态恢复 Agent,等待下次调度 + await self.task_repo.update( + request_id, + status=TaskStatus.PENDING, + agent_id="", + progress=0, + error_info=f"推送失败: {reason}", + ) await self.task_repo.mark_pending(request_id) await self.agent_repo.adjust_load(agent.agent_id, -1) - await self.agent_repo.update_status(agent.agent_id, AgentStatus.READY) - logger.warning("task dispatch push failed, reverted request_id=%s agent=%s", request_id, agent.agent_id) + # 若该 Agent 本为 PROCESSING(还有并发容量),保持 PROCESSING;否则置回 READY + restore = prev_status if prev_status == AgentStatus.PROCESSING else AgentStatus.READY + await self.agent_repo.update_status(agent.agent_id, restore) + logger.warning("task dispatch push failed, reverted request_id=%s agent=%s reason=%s", + request_id, agent.agent_id, reason) return False + # 推送成功:清空 error_info + await self.task_repo.update(request_id, error_info="") logger.info("task dispatched request_id=%s -> agent=%s", request_id, agent.agent_id) return True diff --git a/frontend/src/api/index.ts b/frontend/src/api/index.ts index d06d622..42228fd 100644 --- a/frontend/src/api/index.ts +++ b/frontend/src/api/index.ts @@ -92,6 +92,7 @@ export const api = { listAgents: () => request('/agents'), unavailableAgent: (id: string) => request<{ ok: boolean }>(`/agents/${id}/unavailable`, { method: 'POST' }), availableAgent: (id: string) => request<{ ok: boolean }>(`/agents/${id}/available`, { method: 'POST' }), + deleteAgent: (id: string) => request<{ ok: boolean }>(`/agents/${id}`, { method: 'DELETE' }), setAgentPriority: (id: string, priority: number) => request(`/agents/${id}/priority`, { method: 'POST', diff --git a/frontend/src/views/AgentView.vue b/frontend/src/views/AgentView.vue index eebd93f..75f6ef5 100644 --- a/frontend/src/views/AgentView.vue +++ b/frontend/src/views/AgentView.vue @@ -1,5 +1,6 @@