""" FastAPI Web服务 - AI Agent管理服务 支持平台Agent和自定义Agent两种类型 """ from fastapi import FastAPI, HTTPException, Depends from pydantic import BaseModel, Field from typing import Dict, List, Optional from sqlalchemy.orm import Session from datetime import datetime import logging from kubernetes.client.rest import ApiException from k8s_manager import K8sManager, sanitize_k8s_name from database import ( get_db, Template, Agent, Quota, AgentMetric, AgentType, AgentStatus, parse_resource_string, SessionLocal ) from template_manager import template_manager from tool_generator_api import router as tool_generator_router from external_tool_api import router as external_tool_router from tool_storage import tool_storage import os # 配置日志 logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) # 创建FastAPI应用 app = FastAPI( title="AI Agent Manager", description="Kubernetes AI Agent管理服务,支持动态工具生成和 CI/CD 自动构建", version="2.0.0" ) # 注册动态工具生成 Router app.include_router(tool_generator_router) # 注册外部工具 API Router(符合 MCP-Server 规范) app.include_router(external_tool_router) # 注册 Heicode sub-mode Runtime 主 API Router from api.agent.router import router as agent_router app.include_router(agent_router) # 注册 Heicode 兼容 API Router(旧 agnet 命名) from api.agnet.router import router as agnet_router app.include_router(agnet_router) # 注册 Heicode sub-mode Runtime 兼容 Router(旧 /api/swarms 入口) from api.swarm.router import swarms_router app.include_router(swarms_router) # 初始化K8s管理器 NAMESPACE = os.getenv("NAMESPACE", "ai-agents") KUBECONFIG_PATH = os.getenv("KUBECONFIG_PATH", None) # 可选:指定kubeconfig路径 k8s_manager = K8sManager(namespace=NAMESPACE, kubeconfig_path=KUBECONFIG_PATH) def _delete_stale_agent_record(db: Session, db_agent: Optional[Agent], reason: str) -> None: """删除数据库中的失效 Agent 记录。""" if not db_agent: return agent_name = db_agent.name try: db.delete(db_agent) db.commit() logger.info(f"🧹 已清理失效 Agent 记录: {agent_name}, reason={reason}") except Exception as e: db.rollback() logger.error(f"清理失效 Agent 记录失败: {agent_name}, error={e}") def _discover_agent_namespace(agent_name: str) -> Optional[str]: """尝试在 K8s 中发现 Agent 所在命名空间。""" try: namespaces = k8s_manager.v1.list_namespace( label_selector=f"agent-name={agent_name}" ) if namespaces.items: return namespaces.items[0].metadata.name except Exception as e: logger.warning(f"按标签查找命名空间失败: {agent_name}, error={e}") for ns_pattern in [f"agent-{agent_name}", f"agent-test-{agent_name}"]: try: k8s_manager.v1.read_namespace(name=ns_pattern) return ns_pattern except ApiException as e: if e.status != 404: logger.warning(f"检查命名空间失败: {ns_pattern}, error={e}") except Exception as e: logger.warning(f"检查命名空间异常: {ns_pattern}, error={e}") return None def _find_agent_pod( agent_name: str, db: Session, db_agent: Optional[Agent] = None, cleanup_stale: bool = False, ): """查找 Agent 对应的 Pod,必要时同步 namespace 或清理失效数据库记录。""" namespaces_to_try: List[str] = [] if db_agent and db_agent.namespace: namespaces_to_try.append(db_agent.namespace) discovered_namespace = _discover_agent_namespace(agent_name) if discovered_namespace and discovered_namespace not in namespaces_to_try: namespaces_to_try.append(discovered_namespace) for namespace in namespaces_to_try: try: temp_manager = K8sManager(namespace=namespace, kubeconfig_path=KUBECONFIG_PATH) pod = temp_manager.v1.read_namespaced_pod( name=agent_name, namespace=namespace ) if db_agent and db_agent.namespace != namespace: db_agent.namespace = namespace try: db.commit() except Exception as e: db.rollback() logger.warning(f"同步 Agent namespace 失败: {agent_name}, error={e}") return pod, namespace except ApiException as e: if e.status == 404: continue raise if cleanup_stale and db_agent: _delete_stale_agent_record(db, db_agent, "pod_or_namespace_not_found") return None, discovered_namespace def _build_agent_access_info(db_agent: Optional[Agent]) -> Dict: """Build the Manager-facing access info snapshot from the DB record.""" if not db_agent: return {} access_info = {} if db_agent.external_ip: access_info["external_ip"] = db_agent.external_ip if db_agent.ip_url: access_info["ip_url"] = db_agent.ip_url if db_agent.domain: access_info["domain"] = db_agent.domain if db_agent.domain_url: access_info["domain_url"] = db_agent.domain_url if db_agent.recommended_url: access_info["recommended_url"] = db_agent.recommended_url if db_agent.service_name: access_info["service_name"] = db_agent.service_name return access_info def _normalize_agent_runtime_status(raw_status: Optional[str], db_status: Optional[AgentStatus] = None) -> str: """Project k8s/runtime states onto the HM-compatible lifecycle vocabulary.""" status = (raw_status or "").strip().lower() status_map = { "pending": "pending", "accepted": "pending", "initializing": "pending", "containercreating": "pending", "running": "running", "waiting": "pending", "unknown": "pending", "succeeded": "stopped", "terminated": "stopped", "stopped": "stopped", "failed": "failed", "error": "failed", "crashloopbackoff": "failed", } if status in status_map: return status_map[status] if db_status: return db_status.value.lower() return "pending" def _build_agent_lifecycle_response( agent_name: str, db_agent: Optional[Agent], pod_status: Optional[Dict] = None, namespace: Optional[str] = None, ) -> Dict: """Build a flat lifecycle payload that HM can consume directly.""" access_info = _build_agent_access_info(db_agent) runtime_status = _normalize_agent_runtime_status( (pod_status or {}).get("status"), db_agent.status if db_agent else None, ) resolved_name = db_agent.name if db_agent else (pod_status or {}).get("name") or agent_name resolved_namespace = namespace or (pod_status or {}).get("namespace") or (db_agent.namespace if db_agent else None) subdomain = access_info.get("domain") or access_info.get("external_ip") template_name = (pod_status or {}).get("template") if not template_name and db_agent and db_agent.template: template_name = db_agent.template.name framework = (pod_status or {}).get("framework") if not framework and db_agent and db_agent.agent_framework: framework = db_agent.agent_framework.upper() return { "runtime_id": resolved_name, "agent_id": resolved_name, "id": resolved_name, "name": resolved_name, "namespace": resolved_namespace, "status": runtime_status, "runtime_status": runtime_status, "state": runtime_status, "framework": framework, "template": template_name, "subdomain": subdomain, "access_token": None, "access_info": access_info or None, } # ==================== 请求/响应模型 ==================== # Template Management Models class CreateTemplateRequest(BaseModel): """创建模板请求""" name: str = Field(..., min_length=1, max_length=100) display_name: str description: Optional[str] = None agent_type: str = Field(..., description="platform or custom") agent_framework: str = Field(default="langchain", description="langchain, mcp, or a2a") image: str port: Optional[int] = None env_requirements: Optional[Dict] = Field(default_factory=dict) tools_config: Optional[Dict] = Field(default_factory=dict, description="Tools configuration JSON") default_model_provider: Optional[str] = Field(None, description="Default model provider") default_model_name: Optional[str] = Field(None, description="Default model name") cpu_request: Optional[str] = None cpu_limit: Optional[str] = None memory_request: Optional[str] = None memory_limit: Optional[str] = None min_replicas: int = 1 max_replicas: int = 3 target_cpu_utilization: int = 80 class UpdateTemplateRequest(BaseModel): """更新模板请求""" display_name: Optional[str] = None description: Optional[str] = None image: Optional[str] = None port: Optional[int] = None env_requirements: Optional[Dict] = None cpu_request: Optional[str] = None cpu_limit: Optional[str] = None memory_request: Optional[str] = None memory_limit: Optional[str] = None min_replicas: Optional[int] = None max_replicas: Optional[int] = None target_cpu_utilization: Optional[int] = None is_active: Optional[bool] = None class TemplateResponse(BaseModel): """模板响应""" id: int name: str display_name: str description: Optional[str] agent_type: str image: str port: Optional[int] env_requirements: Dict cpu_request: Optional[str] cpu_limit: Optional[str] memory_request: Optional[str] memory_limit: Optional[str] min_replicas: int max_replicas: int target_cpu_utilization: int is_active: bool created_at: datetime class Config: from_attributes = True # Platform Agent Models class CreatePlatformAgentRequest(BaseModel): """创建平台Agent请求""" name: str = Field(..., min_length=1, max_length=63) template_name: str owner_id: str channel_id: Optional[str] = None tenant_id: Optional[str] = None namespace: Optional[str] = Field(default="ai-agents", description="Kubernetes namespace") query_params: Optional[Dict] = Field(default_factory=dict) # NEW: Framework-specific configurations agent_framework: Optional[str] = Field(None, description="Override template framework") tools_config: Optional[Dict] = Field(default_factory=dict, description="Tools configuration") tool_endpoint: Optional[str] = Field(None, description="External tool endpoint") tool_api_key: Optional[str] = Field(None, description="Tool API key") model_provider: Optional[str] = Field(None, description="Model provider") model_name: Optional[str] = Field(None, description="Model name") model_endpoint: Optional[str] = Field(None, description="Model endpoint") model_api_key: Optional[str] = Field(None, description="Model API key") storage_connection_string: Optional[str] = Field(None, description="Storage connection string") storage_account_name: Optional[str] = Field(None, description="Storage account name") # Custom Agent Models class ScalingConfig(BaseModel): """弹性伸缩配置""" min_replicas: int = Field(1, ge=0) max_replicas: int = Field(3, ge=1) target_cpu_utilization: int = Field(80, ge=1, le=100) class CreateCustomAgentRequest(BaseModel): """创建自定义Agent请求""" name: str = Field(..., min_length=1, max_length=63) template_name: str owner_id: str channel_id: Optional[str] = None tenant_id: Optional[str] = None namespace: Optional[str] = Field(default="ai-agents", description="Kubernetes namespace") environment_vars: Dict[str, str] # NEW: Framework-specific configurations agent_framework: Optional[str] = Field(None, description="Override template framework") tools_config: Optional[Dict] = Field(default_factory=dict, description="Tools configuration") tool_endpoint: Optional[str] = Field(None, description="External tool endpoint") tool_api_key: Optional[str] = Field(None, description="Tool API key") model_provider: Optional[str] = Field(None, description="Model provider") model_name: Optional[str] = Field(None, description="Model name") model_endpoint: Optional[str] = Field(None, description="Model endpoint") model_api_key: Optional[str] = Field(None, description="Model API key") storage_connection_string: Optional[str] = Field(None, description="Storage connection string") storage_account_name: Optional[str] = Field(None, description="Storage account name") # Resource configuration cpu_request: Optional[str] = None cpu_limit: Optional[str] = None memory_request: Optional[str] = None memory_limit: Optional[str] = None scaling_config: Optional[ScalingConfig] = None class UpdateAgentEnvRequest(BaseModel): """更新Agent环境变量请求""" environment_vars: Dict[str, str] class UpdateScalingRequest(BaseModel): """更新伸缩配置请求""" min_replicas: Optional[int] = None max_replicas: Optional[int] = None target_cpu_utilization: Optional[int] = None # Unified Agent Response class AgentResponseNew(BaseModel): """Agent响应(新)""" id: int name: str display_name: Optional[str] template_name: str agent_type: str status: str owner_id: str channel_id: Optional[str] tenant_id: Optional[str] service_url: Optional[str] current_replicas: int min_replicas: int max_replicas: int created_at: datetime last_accessed_at: Optional[datetime] class Config: from_attributes = True # Legacy Models (for backward compatibility) class CreateAgentRequest(BaseModel): """创建Agent请求(旧版)""" name: str = Field(..., description="Agent名称", min_length=1, max_length=63) template: str = Field(..., description="模板类型") framework: Optional[str] = Field(default="API", description="Agent框架类型: MCP, A2A, API") config: Dict = Field(default_factory=dict, description="配置信息") env: Optional[Dict[str, str]] = Field(default_factory=dict, description="环境变量") namespace: Optional[str] = Field(default=None, description="Kubernetes命名空间,默认使用环境变量NAMESPACE的值") tool_refs: Optional[List[str]] = Field(default=None, description="外部数据工具标识列表(符合MCP-Server规范)") class AgentResponse(BaseModel): """Agent响应""" name: str runtime_id: Optional[str] = None agent_id: Optional[str] = None id: Optional[str] = None displayName: Optional[str] = None description: Optional[str] = None namespace: str status: str runtime_status: Optional[str] = None state: Optional[str] = None framework: Optional[str] = None created_at: Optional[str] = None template: Optional[str] = None service_port: Optional[int] = None subdomain: Optional[str] = None access_token: Optional[str] = None access_info: Optional[Dict] = None pod_id: Optional[str] = None pod_ip: Optional[str] = None host_ip: Optional[str] = None node_name: Optional[str] = None owner_info: Optional[Dict] = None tools_attached: Optional[int] = Field(default=0, description="附加的外部工具数量") class ResourceUsage(BaseModel): """资源使用情况""" cpu: Optional[str] = None memory: Optional[str] = None available: Optional[bool] = None reason: Optional[str] = None class ResourceInfo(BaseModel): """资源信息(配额和使用情况)""" requests: Optional[Dict] = None limits: Optional[Dict] = None usage: Optional[ResourceUsage] = None class PodStatusResponse(BaseModel): """Pod状态响应""" name: str displayName: Optional[str] = None description: Optional[str] = None namespace: str status: str health_status: Optional[str] = None # 新增:健康状态 (healthy, unhealthy, degraded) template: Optional[str] = None framework: Optional[str] = None # Agent框架类型 (MCP, A2A, API) created_at: Optional[str] = None node: Optional[str] = None pod_ip: Optional[str] = None containers: Optional[List[Dict]] = None # 新增:容器详细信息 resources: Optional[ResourceInfo] = None service_port: Optional[int] = None access_url: Optional[str] = None endpoints: Optional[Dict] = None conditions: Optional[List[Dict]] = None # 访问信息(从数据库读取) access_info: Optional[Dict] = None # 包含 external_ip, domain, URLs 等 class PodMetricsResponse(BaseModel): """Pod资源使用响应""" name: str namespace: Optional[str] = None requests: Dict limits: Dict usage: Optional[Dict] = None # 实时使用情况(需要 metrics-server) timestamp: Optional[str] = None # metrics 时间戳 metrics_available: Optional[bool] = None # metrics-server 是否可用 class MessageResponse(BaseModel): """通用消息响应""" status: str message: str class AgentLifecycleResponse(BaseModel): """模板 Agent 生命周期兼容响应""" runtime_id: str agent_id: str id: str name: str namespace: Optional[str] = None status: str runtime_status: str state: str framework: Optional[str] = None template: Optional[str] = None subdomain: Optional[str] = None access_token: Optional[str] = None access_info: Optional[Dict] = None @app.get("/") async def root(): """健康检查""" return { "service": "AI Agent Manager", "status": "running", "namespace": NAMESPACE } @app.post("/agents", response_model=AgentResponse) async def create_agent(request: CreateAgentRequest, db: Session = Depends(get_db)): """ 创建AI Agent Pod(在独立命名空间中,并创建Service和Ingress) Args: request: 创建请求(name, template, config, namespace可选, user_id可选) db: 数据库会话 Returns: 创建的Agent信息包括pod_id、service和ingress信息 """ try: logger.info(f"收到创建Agent请求: {request.name}, 模板: {request.template}") # 验证模板类型(从数据库动态获取) valid_templates = template_manager.get_template_names() if request.template not in valid_templates: raise HTTPException( status_code=400, detail=f"无效的模板类型。支持的模板: {', '.join(valid_templates)}" ) # 验证 framework 类型 valid_frameworks = ["MCP", "A2A", "API"] framework = (request.framework or "API").upper() if framework not in valid_frameworks: raise HTTPException( status_code=400, detail=f"无效的框架类型。支持的框架: {', '.join(valid_frameworks)}" ) # DNS-1035 名称合规化:确保名称可以用作 K8s 资源名称 original_name = request.name request.name = sanitize_k8s_name(request.name) if original_name != request.name: logger.info(f"🔄 Agent 名称已合规化: '{original_name}' -> '{request.name}' (DNS-1035)") # 合并环境变量到config config_data = request.config.copy() config_data["agent_framework"] = framework # 添加框架类型到配置 if request.env: config_data["env"] = request.env logger.info(f"环境变量: {list(request.env.keys())}") # 处理 tool_refs(符合 MCP-Server 规范) attached_tools = [] if request.tool_refs: logger.info(f"📦 处理外部工具引用: {request.tool_refs}") # 验证所有工具存在 missing_tools = [] for ref in request.tool_refs: tool = tool_storage.get_tool(ref) if tool: attached_tools.append(tool) else: missing_tools.append(ref) if missing_tools: raise HTTPException( status_code=404, detail=f"以下工具不存在: {', '.join(missing_tools)}" ) # 将工具配置添加到 config config_data["tool_refs"] = request.tool_refs config_data["tools_count"] = len(attached_tools) # 标记工具被使用 for ref in request.tool_refs: tool_storage.mark_tool_in_use(ref, request.name) logger.info(f"✅ 已附加 {len(attached_tools)} 个外部工具") # 添加 user_id 标签 user_id = config_data.get("user_id", "default") if "labels" not in config_data: config_data["labels"] = {} config_data["labels"]["user-id"] = user_id config_data["labels"]["managed-by"] = "agent-manager" config_data["labels"]["app"] = request.name config_data["labels"]["framework"] = framework.lower() # 添加框架标签 # 步骤1: 为每个Agent创建独立的命名空间 agent_namespace = k8s_manager.create_agent_namespace( agent_name=request.name, owner_id=user_id ) logger.info(f"✅ Agent {request.name} 将部署在独立命名空间: {agent_namespace}") # 步骤2: 创建Pod(在独立命名空间中) temp_manager = K8sManager(namespace=agent_namespace, kubeconfig_path=KUBECONFIG_PATH) result = temp_manager.create_pod( pod_name=request.name, template=request.template, config_data=config_data ) # 步骤3: 获取服务端口 # OpenClaw 使用 18789 端口(沙箱镜像会自动启用) if request.template == "openclaw": service_port = 18789 else: # 其他模板从数据库动态获取 service_port = template_manager.get_port(request.template) # 步骤4: 创建 Service(OpenClaw 已在 create_openclaw_deployment 中自动创建 Service 和 Ingress) service_info = None dns_info = None # OpenClaw 使用 Ingress + HTTPS,不需要 LoadBalancer Service if request.template == "openclaw": logger.info("OpenClaw 使用 Ingress + HTTPS,跳过 LoadBalancer Service 创建") # Service 和 Ingress 已在 create_openclaw_deployment 中自动创建 # 从 result 中获取访问信息 if "access_info" in result: service_info = { "name": f"{request.name}-service", "type": "ClusterIP", "note": "OpenClaw 使用 Ingress + HTTPS 访问" } elif service_port: try: import time time.sleep(2) # 等待Pod启动 # 创建 LoadBalancer Service(非 OpenClaw 模板) service_info = temp_manager.create_service( service_name=f"{request.name}-service", namespace=agent_namespace, pod_selector={"app": request.name}, service_port=80, # 外部访问端口 target_port=service_port # Pod内部端口 ) logger.info(f"✅ LoadBalancer Service 创建成功: {service_info['name']}") # 步骤5: 等待 LoadBalancer IP 分配并创建 DNS try: logger.info("等待 LoadBalancer 外网 IP 分配...") external_ip = temp_manager.wait_for_loadbalancer_ip( service_name=f"{request.name}-service", namespace=agent_namespace, max_wait=120, # 最多等待2分钟 interval=5 ) service_info["external_ip"] = external_ip logger.info(f"✅ LoadBalancer 外网 IP: {external_ip}") # 步骤6: 自动创建 Azure DNS 记录 try: logger.info(f"创建 DNS 记录: {request.name}.taijiagnet.com") dns_info = temp_manager.create_dns_record( subdomain=request.name, ip_address=external_ip ) logger.info(f"✅ DNS 记录: {dns_info['domain']} -> {external_ip}") except Exception as e: logger.warning(f"DNS 记录创建失败: {str(e)}") dns_info = {"status": "failed", "error": str(e)} except Exception as e: logger.warning(f"等待 LoadBalancer IP 超时: {str(e)}") logger.info(" 外网 IP 将在后台继续分配") except Exception as e: logger.warning(f"创建 LoadBalancer Service 失败: {str(e)}") # 获取 Pod 详细信息(包括 pod_id) try: import time time.sleep(1) # 等待 Pod 创建完成 pod = temp_manager.v1.read_namespaced_pod( name=request.name, namespace=agent_namespace ) result["pod_id"] = pod.metadata.uid result["pod_ip"] = pod.status.pod_ip result["host_ip"] = pod.status.host_ip result["node_name"] = pod.spec.node_name result["namespace"] = agent_namespace # 更新为实际使用的命名空间 result["framework"] = framework # 添加框架类型到响应 result["owner_info"] = { "user_id": user_id, "agent_name": request.name, "namespace": agent_namespace, "framework": framework, "labels": pod.metadata.labels } # 添加 Service 信息到响应 # OpenClaw 的访问信息已在 create_openclaw_deployment 中设置 if request.template == "openclaw": # OpenClaw 使用 Ingress + HTTPS,访问信息已在 result 中 if "access_info" in result and result["access_info"].get("https_url"): logger.info(f" - HTTPS 访问: {result['access_info']['https_url']}") logger.info(f" - 注意: 使用自签名证书,浏览器会显示安全警告") elif service_info: result["service_info"] = service_info # 构建访问信息(非 OpenClaw 模板) external_ip = service_info.get("external_ip") if external_ip: # 有外网 IP result["access_info"] = { "external_ip": external_ip, "ip_url": f"http://{external_ip}:80", "service_url": f"http://{service_info['cluster_ip']}:80", "pod_url": f"http://{pod.status.pod_ip}:{service_port}" if pod.status.pod_ip and service_port else None } # 添加 DNS 信息 if dns_info and dns_info.get("status") == "created": result["dns_info"] = dns_info result["access_info"]["domain"] = dns_info["domain"] result["access_info"]["domain_url"] = f"http://{dns_info['domain']}" result["access_info"]["recommended"] = f"http://{dns_info['domain']}" logger.info(f" - 推荐访问: http://{dns_info['domain']}") else: result["access_info"]["recommended"] = f"http://{external_ip}:80" logger.info(f" - 外网访问: http://{external_ip}:80") else: # IP 还在分配中 result["access_info"] = { "status": "pending", "external_ip": None, "note": "LoadBalancer IP 正在分配中,请稍后查询" } logger.info(f" - 外网 IP 正在分配中") logger.info(f"✅ Agent创建成功!") logger.info(f" - Pod ID: {result['pod_id']}") logger.info(f" - Namespace: {agent_namespace}") logger.info(f" - User: {user_id}") except Exception as e: logger.warning(f"获取Pod详细信息失败: {str(e)}") # 保存到数据库 try: # 提取访问信息 access_info = result.get("access_info", {}) external_ip = access_info.get("external_ip") domain = access_info.get("domain") ip_url = access_info.get("ip_url") domain_url = access_info.get("domain_url") recommended_url = access_info.get("recommended", domain_url or ip_url) # 查找模板以关联 template_id db_template = db.query(Template).filter(Template.name == request.template).first() # 创建Agent记录 db_agent = Agent( name=request.name, display_name=db_template.display_name if db_template else request.name, template_id=db_template.id if db_template else None, owner_id=user_id, agent_type=AgentType.PLATFORM, # 默认为平台类型 status=AgentStatus.RUNNING, agent_framework=framework.lower(), namespace=agent_namespace, service_name=service_info.get("name") if service_info else None, external_ip=external_ip, domain=domain, ip_url=ip_url, domain_url=domain_url, recommended_url=recommended_url, min_replicas=1, max_replicas=1, current_replicas=1 ) db.add(db_agent) db.commit() db.refresh(db_agent) logger.info(f"✅ Agent信息已保存到数据库: {db_agent.id}") except Exception as e: logger.error(f"保存Agent到数据库失败: {str(e)}") # 不抛出异常,因为Agent已经在K8s中创建成功 # 确保 result 包含所有必需字段 if "name" not in result: result["name"] = request.name if "namespace" not in result: result["namespace"] = agent_namespace if "status" not in result: result["status"] = "Pending" # 添加 displayName 和 description(从模板获取) tpl_info = template_manager.get_template(request.template) if tpl_info: result["displayName"] = tpl_info.get("display_name", request.template) result["description"] = tpl_info.get("description") # 添加外部工具信息 if attached_tools: result["tools_attached"] = len(attached_tools) result["runtime_id"] = result["name"] result["agent_id"] = result["name"] result["id"] = result["name"] result["runtime_status"] = _normalize_agent_runtime_status(result.get("status")) result["state"] = result["runtime_status"] access_info = result.get("access_info") or {} result["subdomain"] = access_info.get("domain") or access_info.get("external_ip") result["access_token"] = None try: return AgentResponse(**result) except Exception as validation_error: logger.error(f"AgentResponse 验证失败: {validation_error}") logger.error(f"result 数据: {result}") raise except Exception as e: import traceback error_msg = str(e) if str(e) else repr(e) logger.error(f"创建Agent失败: {error_msg}") logger.error(f"异常详情: {traceback.format_exc()}") raise HTTPException(status_code=500, detail=error_msg) @app.get("/agents/{agent_name}", response_model=AgentLifecycleResponse) async def get_agent(agent_name: str, db: Session = Depends(get_db)): """Return a flat lifecycle snapshot for HM's template-agent runtime contract.""" try: agent_name = sanitize_k8s_name(agent_name) logger.info(f"获取Agent生命周期信息: {agent_name}") db_agent = db.query(Agent).filter(Agent.name == agent_name).first() pod_status = None pod, agent_namespace = _find_agent_pod( agent_name=agent_name, db=db, db_agent=db_agent, cleanup_stale=False, ) if pod and agent_namespace: temp_manager = K8sManager(namespace=agent_namespace, kubeconfig_path=KUBECONFIG_PATH) pod_status = temp_manager.get_pod_status(pod_name=agent_name) if pod_status.get("status") == "not_found": pod_status = None if not db_agent and not pod_status: raise HTTPException(status_code=404, detail=f"Agent {agent_name} 不存在或已被删除") return AgentLifecycleResponse( **_build_agent_lifecycle_response( agent_name=agent_name, db_agent=db_agent, pod_status=pod_status, namespace=agent_namespace, ) ) except HTTPException: raise except Exception as e: logger.error(f"获取Agent生命周期信息失败: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.post("/agents/{agent_name}/stop", response_model=MessageResponse) async def stop_agent(agent_name: str, db: Session = Depends(get_db)): """Idempotently stop a template agent without deleting its runtime metadata.""" try: agent_name = sanitize_k8s_name(agent_name) logger.info(f"收到停止Agent请求: {agent_name}") db_agent = db.query(Agent).filter(Agent.name == agent_name).first() pod, agent_namespace = _find_agent_pod( agent_name=agent_name, db=db, db_agent=db_agent, cleanup_stale=False, ) if not db_agent and not pod: raise HTTPException(status_code=404, detail=f"Agent {agent_name} 不存在或已被删除") if pod and agent_namespace: temp_manager = K8sManager(namespace=agent_namespace, kubeconfig_path=KUBECONFIG_PATH) stop_result = temp_manager.delete_pod(pod_name=agent_name) if stop_result.get("status") not in {"success", "not_found"}: raise HTTPException(status_code=500, detail=stop_result.get("message", "停止 Agent 失败")) if db_agent: db_agent.status = AgentStatus.STOPPED db_agent.current_replicas = 0 db.commit() return MessageResponse( status="success", message=f"Agent {agent_name} 已停止", ) except HTTPException: raise except Exception as e: db.rollback() logger.error(f"停止Agent失败: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.delete("/agents/{agent_name}", response_model=MessageResponse) async def delete_agent(agent_name: str, db: Session = Depends(get_db)): """ 删除AI Agent(包括独立命名空间、LoadBalancer Service、DNS 记录和数据库记录) Args: agent_name: Agent名称 db: 数据库会话 Returns: 删除结果 """ try: logger.info(f"收到删除Agent请求: {agent_name}") # DNS-1035 名称合规化 agent_name = sanitize_k8s_name(agent_name) db_agent = db.query(Agent).filter(Agent.name == agent_name).first() # 保护机制:防止删除 agent-manager 命名空间 computed_namespace = f"agent-{agent_name}"[:63].rstrip('-') if computed_namespace == "agent-manager": logger.error(f"❌ 禁止删除 agent-manager 命名空间!agent_name={agent_name}, computed_namespace={computed_namespace}") raise HTTPException( status_code=400, detail=f"禁止删除 agent-manager 命名空间。这是系统保护命名空间,不能被删除。" ) # region agent log try: import json, time with open("/home/taiji/tools/agent-manager/.cursor/debug.log", "a") as _f: _f.write(json.dumps({ "sessionId": "debug-session", "runId": "pre-fix", "hypothesisId": "H1", "location": "app.py:delete_agent:entry", "message": "delete_agent called", "data": { "agent_name": agent_name, "computed_namespace": f"agent-{agent_name}"[:63].lower().strip('-'), "manager_namespace": "agent-manager" }, "timestamp": int(time.time() * 1000) }) + "\n") except Exception: pass # endregion # 步骤1: 删除 DNS 记录 try: dns_result = k8s_manager.delete_dns_record(subdomain=agent_name) if dns_result.get("status") == "deleted": logger.info(f"✅ DNS 记录已删除: {dns_result.get('domain')}") except Exception as e: logger.warning(f"删除 DNS 记录失败(可忽略): {str(e)}") # 检查是否是 OpenClaw 类型的 Agent is_openclaw = False if db_agent and db_agent.agent_framework: # 从数据库记录判断 pass # 通过命名空间中的资源判断是否为 OpenClaw try: computed_ns = f"agent-{agent_name}".lower().strip('-')[:63] temp_manager = K8sManager(namespace=computed_ns, kubeconfig_path=KUBECONFIG_PATH) # 尝试查找 OpenClaw 特有的资源 try: temp_manager.v1.read_namespaced_config_map(name=f"{agent_name}-config", namespace=computed_ns) is_openclaw = True logger.info(f"检测到 OpenClaw Agent: {agent_name}") except: pass except: pass # 如果是 OpenClaw,先清理其特有资源 if is_openclaw: try: logger.info(f"清理 OpenClaw 专用资源...") temp_manager.delete_openclaw_deployment(agent_name) except Exception as e: logger.warning(f"清理 OpenClaw 资源失败(可忽略): {e}") # 步骤2: 删除 Agent 的独立命名空间(会自动删除Pod、Service等所有资源) result = k8s_manager.delete_agent_namespace(agent_name=agent_name) # region agent log try: import json, time with open("/home/taiji/tools/agent-manager/.cursor/debug.log", "a") as _f: _f.write(json.dumps({ "sessionId": "debug-session", "runId": "pre-fix", "hypothesisId": "H2", "location": "app.py:delete_agent:after_delete_namespace", "message": "delete_agent_namespace result", "data": { "agent_name": agent_name, "result_status": result.get("status"), "result_namespace": result.get("namespace") }, "timestamp": int(time.time() * 1000) }) + "\n") except Exception: pass # endregion if result.get("status") == "not_found": raise HTTPException(status_code=404, detail=result.get("message")) # 步骤3: 从数据库删除Agent记录 try: if db_agent: db.delete(db_agent) db.commit() logger.info(f"✅ Agent数据库记录已删除: {agent_name}") else: logger.warning(f"⚠️ Agent {agent_name} 在数据库中未找到") except Exception as db_error: logger.error(f"删除数据库记录失败: {str(db_error)}") # 不抛出异常,因为K8s资源已经删除 logger.info(f"✅ Agent {agent_name} 及其所有资源删除成功") return MessageResponse(**result) except HTTPException: raise except Exception as e: logger.error(f"删除Agent失败: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.get("/agents/{agent_name}/status", response_model=PodStatusResponse) async def get_agent_status(agent_name: str, db: Session = Depends(get_db)): """ 获取Agent详细状态(包括访问信息) Args: agent_name: Agent名称 db: 数据库会话 Returns: Agent状态信息(包括Pod状态和访问信息) """ try: logger.info(f"获取Agent状态: {agent_name}") db_agent = db.query(Agent).filter(Agent.name == agent_name).first() _, agent_namespace = _find_agent_pod( agent_name=agent_name, db=db, db_agent=db_agent, cleanup_stale=True, ) if not agent_namespace: raise HTTPException(status_code=404, detail=f"Agent {agent_name} 不存在或已被删除") # 使用正确的namespace获取Pod状态 temp_manager = K8sManager(namespace=agent_namespace, kubeconfig_path=KUBECONFIG_PATH) result = temp_manager.get_pod_status(pod_name=agent_name) if result.get("status") == "not_found": _delete_stale_agent_record(db, db_agent, "status_pod_not_found") raise HTTPException(status_code=404, detail=result.get("message")) # 添加数据库中的信息 if db_agent: # 添加框架类型 result["framework"] = db_agent.agent_framework.upper() if db_agent.agent_framework else "API" # 添加 displayName 和 description(从模板获取) if db_agent.template_id and db_agent.template: result["displayName"] = db_agent.template.display_name result["description"] = db_agent.template.description else: tpl_name = result.get("template") if tpl_name: tpl_info = template_manager.get_template(tpl_name) if tpl_info: result["displayName"] = tpl_info.get("display_name") result["description"] = tpl_info.get("description") if not result.get("displayName"): result["displayName"] = db_agent.display_name # 添加访问信息 access_info = {} if db_agent.external_ip: access_info["external_ip"] = db_agent.external_ip if db_agent.ip_url: access_info["ip_url"] = db_agent.ip_url if db_agent.domain: access_info["domain"] = db_agent.domain if db_agent.domain_url: access_info["domain_url"] = db_agent.domain_url if db_agent.recommended_url: access_info["recommended_url"] = db_agent.recommended_url if db_agent.service_name: access_info["service_name"] = db_agent.service_name # 只有当有访问信息时才添加 if access_info: result["access_info"] = access_info logger.info(f"✅ 已添加访问信息: domain={db_agent.domain}, ip={db_agent.external_ip}") else: logger.warning(f"⚠️ Agent {agent_name} 在数据库中没有访问信息") else: logger.warning(f"⚠️ Agent {agent_name} 在数据库中未找到") return PodStatusResponse(**result) except HTTPException: raise except Exception as e: logger.error(f"获取Agent状态失败: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.get("/agents/{agent_name}/metrics", response_model=PodMetricsResponse) async def get_agent_metrics(agent_name: str, db: Session = Depends(get_db)): """ 获取Agent资源使用情况 Args: agent_name: Agent名称 db: 数据库会话 Returns: Agent资源使用信息 """ try: logger.info(f"获取Agent资源信息: {agent_name}") db_agent = db.query(Agent).filter(Agent.name == agent_name).first() pod, agent_namespace = _find_agent_pod( agent_name=agent_name, db=db, db_agent=db_agent, cleanup_stale=True, ) if not pod or not agent_namespace: raise HTTPException(status_code=404, detail=f"Agent {agent_name} 不存在或已被删除") # 使用正确的namespace获取Pod指标 temp_manager = K8sManager(namespace=agent_namespace, kubeconfig_path=KUBECONFIG_PATH) result = temp_manager.get_pod_metrics(pod_name=agent_name) return PodMetricsResponse(**result) except HTTPException: raise except Exception as e: logger.error(f"获取Agent资源信息失败: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.get("/agents") async def list_agents(template: Optional[str] = None, db: Session = Depends(get_db)): """ 列出所有Agent(跨所有命名空间) Args: template: 模板类型过滤(可选) db: 数据库会话 Returns: Agent列表 """ try: logger.info(f"列出Agents, 模板过滤: {template}") all_agents = [] # 方法1: 从数据库获取Agent列表(推荐) try: query = db.query(Agent) db_agents = query.all() for db_agent in db_agents: try: pod, agent_namespace = _find_agent_pod( agent_name=db_agent.name, db=db, db_agent=db_agent, cleanup_stale=True, ) except Exception as e: logger.warning(f"查询 Agent Pod 失败: {db_agent.name}, error={e}") pod = None agent_namespace = db_agent.namespace if not pod: continue pod_status = pod.status.phase pod_ip = pod.status.pod_ip template_name = pod.metadata.labels.get("template") # 获取模板的 displayName 和 description tpl_display_name = db_agent.display_name tpl_description = None tpl_name = template_name or "unknown" if db_agent.template_id and db_agent.template: tpl_display_name = db_agent.template.display_name tpl_description = db_agent.template.description tpl_name = db_agent.template.name elif template_name: tpl_info = template_manager.get_template(template_name) if tpl_info: tpl_display_name = tpl_info.get("display_name", template_name) tpl_description = tpl_info.get("description") tpl_name = template_name agent_info = { "name": db_agent.name, "displayName": tpl_display_name, "description": tpl_description, "namespace": agent_namespace, "status": pod_status, "template": tpl_name, "framework": db_agent.agent_framework or "api", "created_at": db_agent.created_at.isoformat() if db_agent.created_at else None, "pod_ip": pod_ip, "external_ip": db_agent.external_ip, "domain": db_agent.domain, "service_url": db_agent.recommended_url } # 应用模板过滤 if template and agent_info.get("template") != template: continue all_agents.append(agent_info) logger.info(f"从数据库获取到 {len(all_agents)} 个Agent") except Exception as db_error: logger.warning(f"从数据库获取Agent列表失败: {db_error}") # 方法2: 遍历所有以 agent- 开头的命名空间(作为补充) try: namespaces = k8s_manager.v1.list_namespace() agent_namespaces = [ ns.metadata.name for ns in namespaces.items if ns.metadata.name.startswith("agent-") ] # 已经从数据库获取的Agent名称 known_agents = {a["name"] for a in all_agents} for ns_name in agent_namespaces: try: temp_manager = K8sManager(namespace=ns_name, kubeconfig_path=KUBECONFIG_PATH) label_selector = "managed-by=agent-manager" if template: label_selector += f",template={template}" pods = temp_manager.v1.list_namespaced_pod( namespace=ns_name, label_selector=label_selector ) for pod in pods.items: if pod.metadata.name not in known_agents: k8s_tpl_name = pod.metadata.labels.get("template", "unknown") k8s_tpl = template_manager.get_template(k8s_tpl_name) agent_info = { "name": pod.metadata.name, "displayName": k8s_tpl.get("display_name", k8s_tpl_name) if k8s_tpl else k8s_tpl_name, "description": k8s_tpl.get("description") if k8s_tpl else None, "namespace": ns_name, "status": pod.status.phase, "template": k8s_tpl_name, "framework": pod.metadata.labels.get("framework", "api"), "created_at": pod.metadata.creation_timestamp.isoformat() if pod.metadata.creation_timestamp else None, "pod_ip": pod.status.pod_ip } all_agents.append(agent_info) except Exception as e: logger.debug(f"命名空间 {ns_name} 查询失败: {e}") except Exception as ns_error: logger.warning(f"遍历命名空间失败: {ns_error}") return {"agents": all_agents, "count": len(all_agents)} except Exception as e: logger.error(f"列出Agents失败: {str(e)}") raise HTTPException(status_code=500, detail=str(e)) @app.get("/templates") async def list_templates(): """ 列出所有可用的Agent模板及其所需参数 Returns: 模板列表及其配置信息(从数据库动态加载) """ templates = template_manager.get_all_templates() return { "templates": templates, "count": len(templates) } @app.get("/templates/platform") async def list_platform_templates(): """ 获取平台 Agent 镜像列表 Returns: 平台提供的Agent模板列表(从数据库动态加载) """ templates = template_manager.get_all_templates() platform_templates = [t for t in templates if t.get("agent_type") == "platform"] for t in platform_templates: t["type"] = "platform" return { "templates": platform_templates, "count": len(platform_templates), "type": "platform" } @app.get("/templates/custom") async def list_custom_templates(): """ 获取自定义 Agent 镜像列表 Returns: 用户自定义的Agent模板列表(从数据库动态加载) """ templates = template_manager.get_all_templates() custom_templates = [t for t in templates if t.get("agent_type") == "custom"] for t in custom_templates: t["type"] = "custom" return { "templates": custom_templates, "count": len(custom_templates), "type": "custom" } @app.get("/templates/{template_name}") async def get_template_info_endpoint(template_name: str): """ 获取指定模板的详细信息 Args: template_name: 模板名称 Returns: 模板详细信息(端口、所需环境变量等) """ template = template_manager.get_template(template_name) if not template: valid_templates = template_manager.get_template_names() raise HTTPException( status_code=404, detail=f"模板 {template_name} 不存在。可用模板: {', '.join(valid_templates)}" ) return template # ==================== 模板管理 CRUD API ==================== @app.post("/templates/create") async def create_template(request: CreateTemplateRequest, db: Session = Depends(get_db)): """ 创建新的 Agent 模板 通过 API 动态添加模板,无需修改代码或重新部署 Args: request: 模板创建请求 Returns: 创建的模板信息 Example: POST /templates/create { "name": "my_custom_agent", "display_name": "My Custom Agent", "description": "自定义 Agent 描述", "image": "agnettaiji.azurecr.io/ai-agents/my-agent:latest", "port": 8000, "agent_framework": "api", "env_requirements": { "required": {"MY_API_KEY": "API 密钥"} } } """ try: template = template_manager.create_template( name=request.name, image=request.image, display_name=request.display_name, description=request.description, port=request.port or 8000, agent_framework=request.agent_framework or "api", env_requirements=request.env_requirements, tools_config=request.tools_config, cpu_request=request.cpu_request, cpu_limit=request.cpu_limit, memory_request=request.memory_request, memory_limit=request.memory_limit, created_by="api", ) logger.info(f"✅ 模板创建成功: {request.name}") return { "status": "success", "message": f"模板 {request.name} 创建成功", "template": template } except ValueError as e: raise HTTPException(status_code=400, detail=str(e)) except Exception as e: logger.error(f"创建模板失败: {e}") raise HTTPException(status_code=500, detail=str(e)) @app.put("/templates/{template_name}") async def update_template(template_name: str, request: UpdateTemplateRequest, db: Session = Depends(get_db)): """ 更新现有模板 Args: template_name: 模板名称 request: 更新请求(只更新提供的字段) Returns: 更新后的模板信息 """ try: # 检查模板是否存在 if not template_manager.template_exists(template_name): raise HTTPException( status_code=404, detail=f"模板 {template_name} 不存在" ) template = template_manager.update_template( name=template_name, image=request.image, display_name=request.display_name, description=request.description, port=request.port, env_requirements=request.env_requirements, cpu_request=request.cpu_request, cpu_limit=request.cpu_limit, memory_request=request.memory_request, memory_limit=request.memory_limit, is_active=request.is_active, ) logger.info(f"✅ 模板更新成功: {template_name}") return { "status": "success", "message": f"模板 {template_name} 更新成功", "template": template } except ValueError as e: raise HTTPException(status_code=400, detail=str(e)) except Exception as e: logger.error(f"更新模板失败: {e}") raise HTTPException(status_code=500, detail=str(e)) @app.delete("/templates/{template_name}") async def delete_template(template_name: str, force: bool = False, db: Session = Depends(get_db)): """ 删除模板 Args: template_name: 模板名称 force: 是否强制删除(默认软删除/禁用) Returns: 删除结果 """ try: # 检查模板是否存在 if not template_manager.template_exists(template_name): raise HTTPException( status_code=404, detail=f"模板 {template_name} 不存在" ) template_manager.delete_template(template_name, force=force) action = "永久删除" if force else "禁用" logger.info(f"✅ 模板{action}成功: {template_name}") return { "status": "success", "message": f"模板 {template_name} 已{action}" } except ValueError as e: raise HTTPException(status_code=400, detail=str(e)) except Exception as e: logger.error(f"删除模板失败: {e}") raise HTTPException(status_code=500, detail=str(e)) @app.get("/templates/search/{keyword}") async def search_templates(keyword: str): """ 搜索模板 Args: keyword: 搜索关键字(匹配名称、显示名、描述) Returns: 匹配的模板列表 """ templates = template_manager.search_templates(keyword) return { "templates": templates, "count": len(templates), "keyword": keyword } @app.post("/templates/refresh-cache") async def refresh_template_cache(): """ 刷新模板缓存 手动刷新内存中的模板缓存,立即生效新的模板配置 """ template_manager.invalidate_cache() templates = template_manager.get_all_templates() return { "status": "success", "message": "模板缓存已刷新", "count": len(templates) } if __name__ == "__main__": import uvicorn host = os.getenv("SERVICE_HOST", "0.0.0.0") port = int(os.getenv("SERVICE_PORT", "8000")) logger.info(f"启动AI Agent Manager服务: {host}:{port}") uvicorn.run(app, host=host, port=port)