Files
agent_management/k8s_manager_new.py
T
2026-01-06 09:26:11 +00:00

627 lines
23 KiB
Python
Raw 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.
"""
Enhanced Kubernetes Manager - 支持Deployment、Service、HPA和Secrets
"""
from kubernetes import client, config
from kubernetes.client.rest import ApiException
from datetime import datetime
import logging
import os
import base64
logger = logging.getLogger(__name__)
class K8sManager:
"""Kubernetes资源管理器 - 增强版"""
def __init__(self, namespace="ai-agents", kubeconfig_path=None):
"""
初始化K8s管理器
Args:
namespace: 命名空间
kubeconfig_path: kubeconfig文件路径(可选,用于本地开发)
"""
self.namespace = namespace
try:
if kubeconfig_path and os.path.exists(kubeconfig_path):
config.load_kube_config(kubeconfig_path)
logger.info(f"使用kubeconfig: {kubeconfig_path}")
else:
config.load_incluster_config()
logger.info("使用集群内ServiceAccount")
except Exception as e:
logger.error(f"K8s配置加载失败: {str(e)}")
raise
self.core_v1 = client.CoreV1Api()
self.apps_v1 = client.AppsV1Api()
self.autoscaling_v2 = client.AutoscalingV2Api()
self._ensure_namespace()
def _ensure_namespace(self):
"""确保命名空间存在"""
try:
self.core_v1.read_namespace(self.namespace)
logger.info(f"命名空间 {self.namespace} 已存在")
except ApiException as e:
if e.status == 404:
namespace = client.V1Namespace(
metadata=client.V1ObjectMeta(name=self.namespace)
)
self.core_v1.create_namespace(namespace)
logger.info(f"创建命名空间: {self.namespace}")
else:
raise
def create_secret(self, name: str, data: dict) -> dict:
"""
创建Kubernetes Secret存储敏感数据
Args:
name: Secret名称
data: 敏感数据字典
Returns:
Secret信息
"""
try:
# 编码数据为base64
encoded_data = {}
for key, value in data.items():
if isinstance(value, str):
encoded_data[key] = base64.b64encode(value.encode()).decode()
else:
encoded_data[key] = base64.b64encode(str(value).encode()).decode()
secret = client.V1Secret(
metadata=client.V1ObjectMeta(
name=name,
namespace=self.namespace,
labels={
"managed-by": "agent-manager",
"type": "agent-secret"
}
),
type="Opaque",
data=encoded_data
)
result = self.core_v1.create_namespaced_secret(self.namespace, secret)
logger.info(f"Created secret: {name}")
return {"name": name, "namespace": self.namespace}
except ApiException as e:
if e.status == 409:
# Secret已存在,更新它
logger.info(f"Secret {name} exists, updating...")
result = self.core_v1.replace_namespaced_secret(name, self.namespace, secret)
return {"name": name, "namespace": self.namespace}
else:
logger.error(f"Failed to create secret: {e}")
raise
def delete_secret(self, name: str):
"""删除Secret"""
try:
self.core_v1.delete_namespaced_secret(name, self.namespace)
logger.info(f"Deleted secret: {name}")
except ApiException as e:
if e.status != 404:
logger.error(f"Failed to delete secret: {e}")
def create_deployment_and_service(self, name: str, template, agent, env_vars: dict) -> dict:
"""
创建Deployment和Service
Args:
name: Agent名称
template: Template数据库对象
agent: Agent数据库对象
env_vars: 环境变量字典
Returns:
部署结果信息
"""
try:
deployment_name = f"{name}-deployment"
service_name = f"{name}-service"
# 1. 如果有敏感环境变量,创建Secret
secret_name = None
if env_vars:
secret_name = f"{name}-secret"
self.create_secret(secret_name, env_vars)
# 2. 创建Deployment
deployment = self._build_deployment(
name=deployment_name,
image=template.image,
port=template.port,
secret_name=secret_name,
agent=agent,
labels={
"app": name,
"managed-by": "agent-manager",
"template": template.name,
"agent-type": agent.agent_type.value,
"owner": agent.owner_id
}
)
self.apps_v1.create_namespaced_deployment(self.namespace, deployment)
logger.info(f"Created deployment: {deployment_name}")
# 3. 创建Service(如果模板定义了端口)
service_url = None
if template.port:
service = self._build_service(
name=service_name,
port=template.port,
selector={"app": name}
)
self.core_v1.create_namespaced_service(self.namespace, service)
logger.info(f"Created service: {service_name}")
# 生成服务URL(集群内访问)
service_url = f"http://{service_name}.{self.namespace}.svc.cluster.local:{template.port}"
# 4. 创建HPA(如果配置了弹性伸缩)
if agent.max_replicas > agent.min_replicas:
self.create_hpa(
name=f"{name}-hpa",
deployment_name=deployment_name,
min_replicas=agent.min_replicas,
max_replicas=agent.max_replicas,
target_cpu_utilization=agent.target_cpu_utilization
)
return {
"deployment_name": deployment_name,
"service_name": service_name,
"service_url": service_url,
"secret_name": secret_name
}
except Exception as e:
logger.error(f"Failed to create deployment and service: {str(e)}")
# 清理已创建的资源
self._cleanup_resources(deployment_name, service_name, secret_name)
raise
def _build_deployment(self, name: str, image: str, port: int, secret_name: str,
agent, labels: dict) -> client.V1Deployment:
"""构建Deployment对象"""
# 环境变量配置
env_vars = []
if secret_name:
# 从Secret引用环境变量
for key in agent.environment_vars.keys():
env_vars.append(client.V1EnvVar(
name=key,
value_from=client.V1EnvVarSource(
secret_key_ref=client.V1SecretKeySelector(
name=secret_name,
key=key
)
)
))
# 容器配置
container = client.V1Container(
name="agent",
image=image,
image_pull_policy="Always",
env=env_vars if env_vars else None,
resources=client.V1ResourceRequirements(
requests={
"cpu": agent.cpu_request or "100m",
"memory": agent.memory_request or "128Mi"
},
limits={
"cpu": agent.cpu_limit or "500m",
"memory": agent.memory_limit or "512Mi"
}
)
)
# 如果有端口,添加端口配置
if port:
container.ports = [client.V1ContainerPort(container_port=port)]
# Pod模板
template = client.V1PodTemplateSpec(
metadata=client.V1ObjectMeta(
labels=labels
),
spec=client.V1PodSpec(
containers=[container],
image_pull_secrets=[client.V1LocalObjectReference(name="acr-secret")]
)
)
# Deployment规格
spec = client.V1DeploymentSpec(
replicas=agent.min_replicas,
selector=client.V1LabelSelector(
match_labels={"app": labels["app"]}
),
template=template
)
# Deployment对象
deployment = client.V1Deployment(
api_version="apps/v1",
kind="Deployment",
metadata=client.V1ObjectMeta(
name=name,
namespace=self.namespace,
labels=labels
),
spec=spec
)
return deployment
def _build_service(self, name: str, port: int, selector: dict) -> client.V1Service:
"""构建Service对象"""
service = client.V1Service(
api_version="v1",
kind="Service",
metadata=client.V1ObjectMeta(
name=name,
namespace=self.namespace,
labels={
"managed-by": "agent-manager"
}
),
spec=client.V1ServiceSpec(
selector=selector,
ports=[client.V1ServicePort(
port=port,
target_port=port,
protocol="TCP"
)],
type="ClusterIP"
)
)
return service
def create_hpa(self, name: str, deployment_name: str, min_replicas: int,
max_replicas: int, target_cpu_utilization: int) -> dict:
"""
创建HorizontalPodAutoscaler
Args:
name: HPA名称
deployment_name: 目标Deployment名称
min_replicas: 最小副本数
max_replicas: 最大副本数
target_cpu_utilization: 目标CPU利用率(百分比)
Returns:
HPA信息
"""
try:
hpa = client.V2HorizontalPodAutoscaler(
api_version="autoscaling/v2",
kind="HorizontalPodAutoscaler",
metadata=client.V1ObjectMeta(
name=name,
namespace=self.namespace
),
spec=client.V2HorizontalPodAutoscalerSpec(
scale_target_ref=client.V2CrossVersionObjectReference(
api_version="apps/v1",
kind="Deployment",
name=deployment_name
),
min_replicas=min_replicas,
max_replicas=max_replicas,
metrics=[
client.V2MetricSpec(
type="Resource",
resource=client.V2ResourceMetricSource(
name="cpu",
target=client.V2MetricTarget(
type="Utilization",
average_utilization=target_cpu_utilization
)
)
)
]
)
)
result = self.autoscaling_v2.create_namespaced_horizontal_pod_autoscaler(
self.namespace, hpa
)
logger.info(f"Created HPA: {name}")
return {"name": name, "namespace": self.namespace}
except ApiException as e:
logger.error(f"Failed to create HPA: {e}")
raise
def delete_hpa(self, name: str):
"""删除HPA"""
try:
self.autoscaling_v2.delete_namespaced_horizontal_pod_autoscaler(
name, self.namespace
)
logger.info(f"Deleted HPA: {name}")
except ApiException as e:
if e.status != 404:
logger.error(f"Failed to delete HPA: {e}")
def update_deployment_env(self, deployment_name: str, env_vars: dict):
"""
更新Deployment的环境变量(通过更新Secret)
Args:
deployment_name: Deployment名称
env_vars: 新的环境变量字典
"""
try:
# 获取Deployment
deployment = self.apps_v1.read_namespaced_deployment(
deployment_name, self.namespace
)
# 查找Secret名称
secret_name = None
for env in deployment.spec.template.spec.containers[0].env or []:
if env.value_from and env.value_from.secret_key_ref:
secret_name = env.value_from.secret_key_ref.name
break
if secret_name:
# 更新Secret
self.create_secret(secret_name, env_vars)
# 触发Pod重启(通过添加annotation)
if not deployment.spec.template.metadata.annotations:
deployment.spec.template.metadata.annotations = {}
deployment.spec.template.metadata.annotations["kubectl.kubernetes.io/restartedAt"] = \
datetime.utcnow().isoformat()
self.apps_v1.replace_namespaced_deployment(
deployment_name, self.namespace, deployment
)
logger.info(f"Updated deployment env: {deployment_name}")
else:
raise ValueError("No secret found in deployment")
except ApiException as e:
logger.error(f"Failed to update deployment env: {e}")
raise
def delete_deployment_and_service(self, deployment_name: str, service_name: str):
"""
删除Deployment、Service和相关资源
Args:
deployment_name: Deployment名称
service_name: Service名称
"""
try:
# 删除Deployment
try:
self.apps_v1.delete_namespaced_deployment(
deployment_name, self.namespace,
propagation_policy='Foreground'
)
logger.info(f"Deleted deployment: {deployment_name}")
except ApiException as e:
if e.status != 404:
logger.error(f"Failed to delete deployment: {e}")
# 删除Service
try:
self.core_v1.delete_namespaced_service(service_name, self.namespace)
logger.info(f"Deleted service: {service_name}")
except ApiException as e:
if e.status != 404:
logger.error(f"Failed to delete service: {e}")
# 删除HPA
hpa_name = deployment_name.replace("-deployment", "-hpa")
self.delete_hpa(hpa_name)
# 删除Secret
secret_name = deployment_name.replace("-deployment", "-secret")
self.delete_secret(secret_name)
except Exception as e:
logger.error(f"Failed to delete resources: {str(e)}")
raise
def _cleanup_resources(self, deployment_name: str, service_name: str, secret_name: str):
"""清理资源(用于错误恢复)"""
if deployment_name:
try:
self.apps_v1.delete_namespaced_deployment(deployment_name, self.namespace)
except:
pass
if service_name:
try:
self.core_v1.delete_namespaced_service(service_name, self.namespace)
except:
pass
if secret_name:
try:
self.delete_secret(secret_name)
except:
pass
def get_deployment_status(self, deployment_name: str) -> dict:
"""获取Deployment状态"""
try:
deployment = self.apps_v1.read_namespaced_deployment(
deployment_name, self.namespace
)
# 获取 Deployment 对应的 Pods 实际状态
label_selector = f"app={deployment_name.replace('-deployment', '')}"
pods = self.core_v1.list_namespaced_pod(
self.namespace,
label_selector=label_selector
)
# 检查 Pod 的健康状态
health_status = "healthy"
pod_details = []
for pod in pods.items:
pod_health = "healthy"
container_statuses = pod.status.container_statuses or []
for container_status in container_statuses:
container_info = {
"name": container_status.name,
"ready": container_status.ready,
"restart_count": container_status.restart_count
}
# 检查容器状态
if container_status.state.waiting:
container_info["state"] = "waiting"
container_info["reason"] = container_status.state.waiting.reason
pod_health = "unhealthy"
elif container_status.state.terminated:
container_info["state"] = "terminated"
container_info["reason"] = container_status.state.terminated.reason
container_info["exit_code"] = container_status.state.terminated.exit_code
pod_health = "unhealthy"
elif container_status.state.running:
container_info["state"] = "running"
# 检查是否就绪
if not container_status.ready:
pod_health = "unhealthy"
# 检查重启次数
if container_status.restart_count > 5:
pod_health = "degraded"
pod_details.append({
"name": pod.metadata.name,
"phase": pod.status.phase,
"health": pod_health,
"containers": [container_info]
})
# 更新整体健康状态
if pod_health == "unhealthy":
health_status = "unhealthy"
elif pod_health == "degraded" and health_status != "unhealthy":
health_status = "degraded"
return {
"name": deployment_name,
"namespace": self.namespace,
"status": "Running" if deployment.status.available_replicas else "Pending",
"health_status": health_status, # 新增:真实健康状态
"replicas": deployment.status.replicas or 0,
"ready_replicas": deployment.status.ready_replicas or 0,
"available_replicas": deployment.status.available_replicas or 0,
"pods": pod_details, # 新增:Pod详细信息
"conditions": [
{
"type": c.type,
"status": c.status,
"reason": c.reason,
"message": c.message
}
for c in (deployment.status.conditions or [])
]
}
except ApiException as e:
if e.status == 404:
return {"status": "not_found", "message": f"Deployment {deployment_name} not found"}
raise
def get_pod_logs(self, deployment_name: str, lines: int = 100) -> str:
"""获取Pod日志"""
try:
# 查找Deployment对应的Pods
label_selector = f"app={deployment_name.replace('-deployment', '')}"
pods = self.core_v1.list_namespaced_pod(
self.namespace,
label_selector=label_selector
)
if not pods.items:
return "No pods found"
# 获取第一个Pod的日志
pod_name = pods.items[0].metadata.name
logs = self.core_v1.read_namespaced_pod_log(
pod_name, self.namespace,
tail_lines=lines
)
return logs
except ApiException as e:
logger.error(f"Failed to get pod logs: {e}")
raise
# ==================== 向后兼容的方法 ====================
def create_pod(self, pod_name: str, template: str, config_data: dict) -> dict:
"""创建Pod(旧方法,保留向后兼容)"""
# 这个方法现在已被create_deployment_and_service替代
# 但为了兼容性保留
raise NotImplementedError("Use create_deployment_and_service instead")
def delete_pod(self, pod_name: str) -> dict:
"""删除Pod(旧方法)"""
raise NotImplementedError("Use delete_deployment_and_service instead")
def get_pod_status(self, pod_name: str) -> dict:
"""获取Pod状态(旧方法)"""
# 尝试查找对应的Deployment
deployment_name = f"{pod_name}-deployment"
return self.get_deployment_status(deployment_name)
def list_pods(self, label_selector: str = None) -> list:
"""列出Pods"""
try:
if label_selector:
deployments = self.apps_v1.list_namespaced_deployment(
self.namespace,
label_selector=label_selector
)
else:
deployments = self.apps_v1.list_namespaced_deployment(self.namespace)
result = []
for deployment in deployments.items:
result.append({
"name": deployment.metadata.name,
"namespace": self.namespace,
"replicas": deployment.status.replicas or 0,
"ready_replicas": deployment.status.ready_replicas or 0,
"labels": deployment.metadata.labels
})
return result
except ApiException as e:
logger.error(f"Failed to list deployments: {e}")
raise