Files

1590 lines
59 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
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)