基于多 Agent 协同的 Kubernetes 故障自动修复

2025-11-10T14:11:02+08:00 | 23分钟阅读 | 更新于 2025-11-10T14:11:02+08:00

@

学习目标

学完本章你应该能够:

  1. 用自己的话讲清"为什么单体 controller 不适合 K8s 故障自动修复",以及多 Agent 拆分的工程动机。
  2. 画出五大 Agent(观测 / 诊断 / 执行 / 审计 / 协调)的职责边界与协作关系。
  3. 设计一套跨 Agent 的统一消息信封(Message envelope),说清 trace_id 与 UTC 时间戳的作用。
  4. 用状态机模型描述一次故障修复的完整生命周期,并讲清"迟到的事件"该如何处理。
  5. 在面试中把这个系统讲成一个"有护栏、可回滚、可审计"的工程故事,而不是一堆脚本。

前置知识:

  • Go 基础:interfacechannelcontextsync.Mutex
  • Kubernetes 基础:Pod / Node / Event、client-go informer 机制
  • LLM 基本认知:prompt 构造、JSON 输出、置信度(confidence)

本章你会动手做的事:

  1. 跑通观测 Agent 的 informer 采集逻辑,观察一个异常 Pod 如何被识别并带 trace_id 发出。
  2. 拿一条真实 Pod 崩溃日志,手工构造一条 Observation 消息喂给诊断 Agent。
  3. 画出一次 “CrashLoopBackOff → 扩容” 的状态机流转图,并标出超时分支。

问题背景

Kubernetes 集群中的故障排查通常涉及多个维度:指标异常、Pod 状态异常、网络不通、存储挂载失败、配置错误等。单一脚本或单一 Agent 很难覆盖所有场景。多 Agent 协同架构通过将不同能力拆分为独立角色,可以更高效、更可靠地完成复杂故障的自动诊断与修复。

说实话我一开始也不太信这套,觉得"写个 controller 不就完了"。直到有一次线上 Node NotReady 引发连锁告警,我们那个单体 controller 同时打 K8s API、调 LLM、写日志,结果 goroutine 飙到几千,自己 OOM 了。从那以后我才真正理解为啥要把不同职责拆开——不是为了好看,是为了炸的时候不连坐。

多 Agent 系统架构

一个完整的 K8s 故障自动修复系统可划分为以下 Agent:

+-------------+    +-------------+    +-------------+    +-------------+
|  观测 Agent  | -> |  诊断 Agent  | -> |  执行 Agent  | -> |  审计 Agent  |
+-------------+    +-------------+    +-------------+    +-------------+
      ^                                                   |
      +---------------------------------------------------+

类比:把这套系统想成一家"故障急诊医院"。观测 Agent 是分诊台护士,把病人症状(指标 / 日志 / 事件)量出来;诊断 Agent 是主治医生,结合病历(RAG)和自己的经验(LLM)下诊断;执行 Agent 是手术团队,动手治疗;审计 Agent 是病历管理员,谁动了什么全记下来;协调 Agent 是急诊科主任,掌控整个流程不跑偏。

下面这张图是五大 Agent 的协作关系——注意最后有一条虚线绕回观测:修复之后还要回头确认指标是否真的恢复,形成闭环。

flowchart LR
    OBS[观测 Agent
采集指标/日志/事件] --> DIAG[诊断 Agent
LLM+RAG 推理根因] DIAG --> EXEC[执行 Agent
调用 K8s API 修复] EXEC --> AUD[审计 Agent
记录+回滚] AUD --> COORD[协调 Agent
状态机串联全流程] COORD -. 分配任务/管理状态 .-> OBS
Agent 角色职责
观测 Agent采集指标、日志、事件、链路数据
诊断 Agent基于 LLM 和知识库推理故障根因
执行 Agent调用 K8s API 或运维工具执行修复
审计 Agent记录操作日志,支持回滚与复盘
协调 Agent分配任务、管理状态、协调多 Agent 协作

我个人觉得这五个角色里最容易被低估的是审计 Agent,很多人觉得"不就是写个日志嘛"。但你只要踩过一次"自动修复把集群搞炸了却不知道谁干的"的坑,就会明白审计不是日志,是责任边界。

Agent 间通信协议

在写具体 Agent 之前,我先把消息格式定下来。这点很关键——Agent 之间不约定好协议,后面联调就是灾难现场。我踩过,A Agent 用 snake_case,B Agent 用 camelCase,光是写适配代码就花了我一个下午。

我们采用一个统一的 Message envelope,所有 Agent 收发都用这个结构:

⚠️ 新手必踩的坑:协议不统一。A Agent 用 snake_case、B Agent 用 camelCase,联调时写适配代码能写一下午。所以在写任何业务代码之前,先把消息格式钉死,后面省的是命。

下面用一张时序图示意"观测 → 诊断"这一跳的消息长什么样——注意 trace_id 一路带着走:

sequenceDiagram
    participant OBS as 观测 Agent
    participant DIAG as 诊断 Agent
    OBS->>DIAG: observation 消息(trace_id=inc-001, severity=critical)
    Note over DIAG: 解析 Pod 崩溃现场字段
    DIAG-->>AUD: diagnosis 消息(根因 + 建议动作)
// agent/message.go
package agent

import "time"

// MessageKind 区分消息类型,避免用字符串硬编码到处飞
type MessageKind string

const (
    KindObservation MessageKind = "observation" // 观测结果
    KindDiagnosis   MessageKind = "diagnosis"   // 诊断结果
    KindAction      MessageKind = "action"      // 执行动作
    KindAudit       MessageKind = "audit"       // 审计记录
    KindError       MessageKind = "error"       // 错误回传
)

// Message 是所有 Agent 通信的统一信封
type Message struct {
    TraceID    string                 `json:"trace_id"`    // 贯穿一次故障的 ID,必填
    From       string                 `json:"from"`        // 发送方 Agent 名
    To         string                 `json:"to"`          // 接收方 Agent 名
    Kind       MessageKind            `json:"kind"`        // 消息类型
    Timestamp  time.Time              `json:"timestamp"`   // 发送时间
    Severity   string                 `json:"severity"`    // info/warn/critical
    Payload    map[string]interface{} `json:"payload"`     // 实际业务数据
}

// NewMessage 构造消息,强制带 traceID 和时间戳
func NewMessage(traceID, from, to string, kind MessageKind) Message {
    return Message{
        TraceID:   traceID,
        From:      from,
        To:        to,
        Kind:      kind,
        Timestamp: time.Now().UTC(),
        Severity:  "info",
        Payload:   map[string]interface{}{},
    }
}

一条具体的 Observation 消息长这样,方便理解:

{
  "trace_id": "inc-20251110-001",
  "from": "observer",
  "to": "diagnoser",
  "kind": "observation",
  "timestamp": "2025-11-10T06:11:02Z",
  "severity": "critical",
  "payload": {
    "alert": "PodCrashLoop",
    "namespace": "payment",
    "pod": "api-server-7d8b-x2k4",
    "container": "api",
    "restart_count": 17,
    "last_log": "dial tcp redis:6379: connect: connection refused",
    "events": ["Back-off restarting failed container"]
  }
}

踩坑提示:trace_id 一定要从告警入口就生成好,一路传下去。我见过有人在每个 Agent 里都 uuid.New(),结果审计日志根本串不起来,出了事只能猜。还有,Timestamp 强制 UTC,否则跨时区集群对账能让你怀疑人生。

观测 Agent

观测 Agent 负责从多个数据源收集信息:

  • Prometheus:CPU、内存、QPS、P99 延迟等指标。
  • Loki / ELK:容器日志、应用日志、系统日志。
  • Jaeger / Tempo:分布式链路追踪。
  • Kubernetes Events:Pod 调度、镜像拉取、健康检查等事件。

为啥观测要单独拆出来?因为数据源太多,每个都有自己的客户端、认证、限流。如果让诊断 Agent 直接去拉,那诊断 Agent 就成了一个"什么都懂"的怪物。拆开后观测 Agent 只对"数据是否齐全"负责,诊断 Agent 只对"推理是否合理"负责。

下面是观测 Agent 用 client-go 采集异常 Pod、Node Condition、Event 的完整代码。我只列核心部分,依赖注入和 main 函数你自己补:

// agent/observer/pod_collector.go
package observer

import (
    "context"
    "fmt"
    "time"

    corev1 "k8s.io/api/core/v1"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/client-go/informers"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/cache"
)

// PodAnomaly 描述一个异常 Pod 的快照
type PodAnomaly struct {
    Namespace     string    `json:"namespace"`
    Name          string    `json:"name"`
    Phase         string    `json:"phase"`         // Pending/Running/Failed/Unknown
    Reason        string    `json:"reason"`        // CrashLoopBackOff/ImagePullBackOff 等
    RestartCount  int32     `json:"restart_count"`
    NodeName      string    `json:"node_name"`
    DetectedAt    time.Time `json:"detected_at"`
    ContainerMsgs []string  `json:"container_msgs"` // 各 container 的 last state message
}

// PodCollector 用 informer 监听 Pod 变化,比轮询省 API
type PodCollector struct {
    client    kubernetes.Interface
    out       chan<- PodAnomaly
    threshold int32 // 重启次数阈值,超过才上报
}

func NewPodCollector(client kubernetes.Interface, out chan<- PodAnomaly, restartThreshold int32) *PodCollector {
    return &PodCollector{client: client, out: out, threshold: restartThreshold}
}

// Run 启动 informer,阻塞直到 ctx 取消
func (c *PodCollector) Run(ctx context.Context) error {
    // 默认 30s resync,太短会刷屏,太长错过更新
    factory := informers.NewSharedInformerFactoryWithOptions(
        c.client,
        30*time.Second,
        informers.WithNamespace(""), // 全集群监听
    )
    inf := factory.Core().V1().Pods().Informer()

    // 这里有个坑:UpdateFunc 会被频繁触发,比如 Pod 的 status.hostIP 变化
    // 所以必须自己判断"是否真的异常"以及"是否已经上报过",否则会刷爆下游
    inf.AddEventHandler(cache.ResourceEventHandlerFuncs{
        UpdateFunc: func(oldObj, newObj interface{}) {
            oldPod, ok1 := oldObj.(*corev1.Pod)
            newPod, ok2 := newObj.(*corev1.Pod)
            if !ok1 || !ok2 {
                return
            }

            // 资源版本没变就别处理,informer 偶尔会重复回调
            if oldPod.ResourceVersion == newPod.ResourceVersion {
                return
            }

            anomaly, ok := c.isAnomaly(newPod)
            if !ok {
                return
            }
            select {
            case c.out <- *anomaly:
            case <-ctx.Done():
                // ctx 取消时直接返回,别阻塞在 channel
            }
        },
    })

    factory.Start(ctx.Done())
    if !cache.WaitForCacheSync(ctx.Done(), inf.HasSynced) {
        return fmt.Errorf("pod informer sync timeout")
    }
    <-ctx.Done()
    return nil
}

// isAnomaly 判断 Pod 是否处于异常状态
func (c *PodCollector) isAnomaly(pod *corev1.Pod) (*PodAnomaly, bool) {
    // Pending 太久也算异常,但这里只看 Running 阶段的容器异常
    if pod.Status.Phase == corev1.PodSucceeded {
        return nil, false // Job 正常完成,忽略
    }

    var reason string
    var msgs []string
    var maxRestart int32

    for _, cs := range pod.Status.ContainerStatuses {
        if cs.RestartCount > maxRestart {
            maxRestart = cs.RestartCount
        }
        // 关键:waiting 状态才有 reason,terminated 状态看 reason+message
        if cs.State.Waiting != nil && cs.State.Waiting.Reason != "" {
            reason = cs.State.Waiting.Reason
            if cs.State.Waiting.Message != "" {
                msgs = append(msgs, cs.State.Waiting.Message)
            }
        }
        if cs.State.Terminated != nil && cs.State.Terminated.Reason != "" {
            reason = cs.State.Terminated.Reason
            if cs.State.Terminated.Message != "" {
                msgs = append(msgs, cs.State.Terminated.Message)
            }
        }
    }

    // 阈值过滤:重启次数不够或者没有明确 reason 就别报
    // 我踩过坑:曾经没加 reason 判断,把正常滚动更新的 Pod 全报上去了
    if reason == "" && maxRestart < c.threshold {
        return nil, false
    }

    // 拿 Pod 自身的 condition reason 兜底,CrashLoopBackOff 一般在 container status 里
    if reason == "" && pod.Status.Reason != "" {
        reason = pod.Status.Reason
    }

    return &PodAnomaly{
        Namespace:     pod.Namespace,
        Name:          pod.Name,
        Phase:         string(pod.Status.Phase),
        Reason:        reason,
        RestartCount:  maxRestart,
        NodeName:      pod.Spec.NodeName,
        DetectedAt:    time.Now().UTC(),
        ContainerMsgs: msgs,
    }, true
}

Node Condition 采集相对简单,因为 Node 数量少、变化慢,直接 List-Watch 就行:

// agent/observer/node_collector.go
package observer

import (
    "context"
    "time"

    corev1 "k8s.io/api/core/v1"
    "k8s.io/client-go/informers"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/cache"
)

// NodeAnomaly 描述一个异常 Node
type NodeAnomaly struct {
    Name       string    `json:"name"`
    Conditions []string  `json:"conditions"` // 异常 condition 列表
    DetectedAt time.Time `json:"detected_at"`
}

type NodeCollector struct {
    client kubernetes.Interface
    out    chan<- NodeAnomaly
}

func (c *NodeCollector) Run(ctx context.Context) error {
    factory := informers.NewSharedInformerFactory(c.client, 60*time.Second)
    inf := factory.Core().V1().Nodes().Informer()

    inf.AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc: func(obj interface{}) { c.check(ctx, obj) },
        UpdateFunc: func(_, newObj interface{}) { c.check(ctx, newObj) },
    })

    factory.Start(ctx.Done())
    if !cache.WaitForCacheSync(ctx.Done(), inf.HasSynced) {
        return context.DeadlineExceeded
    }
    <-ctx.Done()
    return nil
}

func (c *NodeCollector) check(ctx context.Context, obj interface{}) {
    node, ok := obj.(*corev1.Node)
    if !ok {
        return
    }

    var badConds []string
    for _, cond := range node.Status.Conditions {
        // 关键判断:NetworkUnavailable / MemoryPressure / DiskPressure / PIDPressure
        // 只要不是 True 的,全算异常;OutOfDisk 老版本才有,现在基本忽略
        // 这里有个反例:Ready=False 不一定意味着 Node 真挂了,可能是 kubelet 暂时没上报
        // 所以我们要再加一个 lastHeartbeatTime 太久的判断
        if cond.Status != corev1.ConditionTrue && cond.Type != corev1.NodeReady {
            badConds = append(badConds, string(cond.Type))
            continue
        }
        if cond.Type == corev1.NodeReady && cond.Status != corev1.ConditionTrue {
            // Ready=False 才是真正的 NotReady,且持续超过 1 分钟才算
            if time.Since(cond.LastHeartbeatTime.Time) > time.Minute {
                badConds = append(badConds, "Ready=False(stale>1m)")
            }
        }
    }

    if len(badConds) == 0 {
        return
    }

    select {
    case c.out <- NodeAnomaly{
        Name:       node.Name,
        Conditions: badConds,
        DetectedAt: time.Now().UTC(),
    }:
    case <-ctx.Done():
    }
}

Event 采集单独说一下。Event 在 K8s 里是会被自动清理的(默认 1 小时),所以你不能依赖 List 拿历史 Event,必须用 Watch 实时收。还有一个坑:Event 的 involvedObject 才是关联 Pod 的关键,过滤的时候要按这个来。

// agent/observer/event_collector.go
package observer

import (
    "context"

    corev1 "k8s.io/api/core/v1"
    "k8s.io/apimachinery/pkg/fields"
    "k8s.io/client-go/informers"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/cache"
)

// EventCollector 只关心 Warning 事件,Normal 事件太多没必要全收
type EventCollector struct {
    client kubernetes.Interface
    out    chan<- corev1.Event
}

func (c *EventCollector) Run(ctx context.Context) error {
    // 只看 Warning 事件,Normal 太多太吵
    // 不过这里有个反例:有些健康检查失败的 Event 类型是 Normal 但 reason 是 Unhealthy
    // 所以严格点应该再按 reason 过滤,但这会让代码很丑,折中:先按 type 过滤
    fieldSelector := fields.OneTermEqualSelector("type", "Warning")

    factory := informers.NewSharedInformerFactoryWithOptions(
        c.client,
        0,
        informers.WithTweakListOptions(func(opts *metav1.ListOptions) {
            opts.FieldSelector = fieldSelector.String()
        }),
    )
    inf := factory.Core().V1().Events().Informer()

    inf.AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc: func(obj interface{}) {
            ev, ok := obj.(*corev1.Event)
            if !ok {
                return
            }
            // 重复事件有 count 字段,count > 1 说明已经发生过多次
            // 审计的时候要记 count,否则会以为是单次
            select {
            case c.out <- *ev:
            case <-ctx.Done():
            }
        },
    })

    factory.Start(ctx.Done())
    if !cache.WaitForCacheSync(ctx.Done(), inf.HasSynced) {
        return context.DeadlineExceeded
    }
    <-ctx.Done()
    return nil
}

踩坑提示:informer 的 WithTweakListOptionsFieldSelector 在大集群(10w+ Event)下能救命,不然你的 observer 内存会涨到几个 G。还有就是 Event 的 lastTimestamp 才是事件真实发生时间,eventTime 字段有些版本不写,别用错。

诊断 Agent

诊断 Agent 接收观测数据,结合 LLM 与 RAG 知识库进行根因分析。诊断过程可采用 ReAct 模式(思考 - 行动 - 观察循环):

flowchart TD
    T[思考: 提出根因假设] --> A[行动: 拉取日志/Redis 状态]
    A --> O[观察: 拿到新证据]
    O --> T
    O --> C{证据是否充分?}
    C -->|否| T
    C -->|是| R[结论: 根因 + 修复建议]
思考1:Pod 处于 CrashLoopBackOff,首先检查最近日志。
行动1:请求日志摘要。
观察1:日志显示 "connection refused to redis:6379"。
思考2:可能是 Redis 服务不可用或配置错误,需要检查 Redis Service 和 Endpoint。
行动2:请求 Redis 状态。
观察2:Redis Endpoint 为空,说明没有 Redis Pod 被选中。
结论:Redis Deployment 副本数为 0,需要扩容。

诊断 Agent 可以维护一个故障模式库,将历史故障案例向量化存储,遇到新告警时先进行相似案例检索。

LLM 根因分析最关键的是 prompt 构造。说实话我刚写的时候 prompt 就一行"分析这个故障",结果模型给我编出各种不存在的 K8s 概念。后来才明白:prompt 必须把上下文塞满,包括 Pod 的 status、events、最近日志、相关资源状态。模型不是算命先生,你给它什么它就用什么。

// agent/diagnoser/llm_diagnoser.go
package diagnoser

import (
    "context"
    "encoding/json"
    "fmt"
    "strings"
)

// DiagnosisInput 诊断 Agent 的输入,从 Observation 消息组装
type DiagnosisInput struct {
    TraceID       string   `json:"trace_id"`
    Namespace     string   `json:"namespace"`
    Pod           string   `json:"pod"`
    Reason        string   `json:"reason"`         // CrashLoopBackOff 等
    ContainerMsgs []string `json:"container_msgs"`
    Events        []string `json:"events"`         // 相关 Event 摘要
    RecentLogs    string   `json:"recent_logs"`    // 最近 N 行日志
    SimilarCases  []string `json:"similar_cases"`  // RAG 检索到的相似案例
}

// DiagnosisResult LLM 给出的诊断结论
type DiagnosisResult struct {
    TraceID         string  `json:"trace_id"`
    RootCause       string  `json:"root_cause"`        // 根因描述
    Confidence      float64 `json:"confidence"`        // 0~1
    RecommendedAction string `json:"recommended_action"` // scale/restart/rollback/...
    ActionTarget    string  `json:"action_target"`     // 操作对象,比如 deployment/redis
    Reasoning       string  `json:"reasoning"`         // 推理过程,便于审计
}

// LLMClient 抽象一下,方便换 OpenAI/Claude/本地模型
type LLMClient interface {
    Chat(ctx context.Context, prompt string) (string, error)
}

type LLMDiagnoser struct {
    client LLMClient
}

func NewLLMDiagnoser(client LLMClient) *LLMDiagnoser {
    return &LLMDiagnoser{client: client}
}

// buildPrompt 构造 prompt,这是诊断质量的关键
// 我个人觉得 prompt 工程在这里比换大模型还有用
func (d *LLMDiagnoser) buildPrompt(in DiagnosisInput) string {
    var sb strings.Builder
    // 系统角色:明确告诉模型你是干啥的,输出格式是什么
    // 不约束输出格式,下游解析就是地狱
    sb.WriteString("你是一个 Kubernetes 故障诊断专家。")
    sb.WriteString("根据下面的观测数据,分析故障根因并给出修复建议。\n")
    sb.WriteString("必须以 JSON 格式输出,字段固定为 root_cause/confidence/recommended_action/action_target/reasoning。\n")
    sb.WriteString("confidence 为 0~1 的浮点数,低于 0.7 的不要建议执行破坏性操作。\n")
    sb.WriteString("recommended_action 只能从 scale/restart/rollback/update_config/none 中选一个。\n\n")

    // 现场数据:信息越全,幻觉越少
    sb.WriteString("## 故障现场\n")
    sb.WriteString(fmt.Sprintf("- namespace: %s\n", in.Namespace))
    sb.WriteString(fmt.Sprintf("- pod: %s\n", in.Pod))
    sb.WriteString(fmt.Sprintf("- reason: %s\n", in.Reason))
    if len(in.ContainerMsgs) > 0 {
        sb.WriteString(fmt.Sprintf("- container_messages: %s\n", strings.Join(in.ContainerMsgs, "; ")))
    }
    if len(in.Events) > 0 {
        sb.WriteString("- recent_events:\n")
        for _, e := range in.Events {
            sb.WriteString(fmt.Sprintf("  * %s\n", e))
        }
    }
    if in.RecentLogs != "" {
        // 日志只截最近 50 行,太长会让 prompt 超 token 还容易跑偏
        sb.WriteString("- recent_logs (tail 50):\n")
        sb.WriteString(in.RecentLogs)
        sb.WriteString("\n")
    }

    // RAG 检索的相似案例,给模型一个"参考答案"
    // 但要明确告诉模型只是参考,避免它无脑照抄
    if len(in.SimilarCases) > 0 {
        sb.WriteString("\n## 历史相似案例(仅供参考,不要直接照搬结论)\n")
        for i, c := range in.SimilarCases {
            sb.WriteString(fmt.Sprintf("案例 %d: %s\n", i+1, c))
        }
    }

    sb.WriteString("\n请给出诊断结论,只输出 JSON,不要任何额外解释。\n")
    return sb.String()
}

// Diagnose 调用 LLM 做诊断,并解析结果
func (d *LLMDiagnoser) Diagnose(ctx context.Context, in DiagnosisInput) (*DiagnosisResult, error) {
    prompt := d.buildPrompt(in)
    raw, err := d.client.Chat(ctx, prompt)
    if err != nil {
        return nil, fmt.Errorf("llm chat failed: %w", err)
    }

    // 大模型经常会在 JSON 前后塞一段 markdown,比如 ```json ... ```
    // 必须先清掉,否则 json.Unmarshal 直接炸
    raw = strings.TrimSpace(raw)
    raw = strings.TrimPrefix(raw, "```json")
    raw = strings.TrimPrefix(raw, "```")
    raw = strings.TrimSuffix(raw, "```")
    raw = strings.TrimSpace(raw)

    var result DiagnosisResult
    if err := json.Unmarshal([]byte(raw), &result); err != nil {
        // 解析失败别硬撑,把原始返回原样回传给审计 Agent 让人看
        return nil, fmt.Errorf("parse llm output failed: %w, raw=%s", err, raw)
    }
    result.TraceID = in.TraceID

    // 兜底:confidence 异常时强行降级
    if result.Confidence < 0 || result.Confidence > 1 {
        result.Confidence = 0.3
    }
    return &result, nil
}

踩坑提示:LLM 输出 JSON 的时候,强烈建议用 OpenAI 的 response_format: json_object 或者函数调用,否则你 10 次里有 1 次会拿到带前缀说明的"垃圾 JSON"。还有 confidence 是模型自报的,不能完全信,我们在执行 Agent 那一层还会再校验一次,破坏性操作低于 0.8 直接拒绝。

执行 Agent

执行 Agent 根据诊断结果执行修复动作。执行前需要校验:

  1. 动作是否在白名单中。
  2. 是否超过当日执行次数限制。
  3. 是否需要人工审批。

⚠️ 新手必踩的坑:破坏性操作零护栏。drain 一个 Node 等于把上面所有 Pod 驱逐走,如果忘了加白名单或人工审批,LLM 一拍脑袋就能把线上节点干掉。护栏(白名单 + 置信度 + 审批)不是可选项,是生死线。

执行 Agent 在真正动手前,会依次过四道闸门,任何一道不通过就直接拒绝:

flowchart TD
    S[收到 Action] --> C1{在白名单?}
    C1 -->|否| X[拒绝: 未知动作]
    C1 -->|是| C2{置信度达标?}
    C2 -->|否| X
    C2 -->|是| C3{当日次数超限?}
    C3 -->|是| X
    C3 -->|否| C4{需人工审批?}
    C4 -->|是| P[转 PendingApproval 等回执]
    C4 -->|否| E[执行修复]

常见修复动作包括:

故障场景修复动作
Pod 崩溃重启 Pod、回滚镜像版本
CPU 饱和扩容 Deployment
内存泄漏滚动重启、调整资源限制
配置错误更新 ConfigMap/Secret 并滚动更新
网络不通检查 Service/Endpoint/NetworkPolicy
// agent/executor/k8s_executor.go
package executor

import (
    "context"
    "fmt"
    "time"

    corev1 "k8s.io/api/core/v1"
    policyv1 "k8s.io/api/policy/v1"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/apimachinery/pkg/labels"
    "k8s.io/client-go/kubernetes"
)

// Action 执行动作的统一结构
type Action struct {
    TraceID    string `json:"trace_id"`
    Type       string `json:"type"`        // scale/restart/rollback/drain
    Namespace  string `json:"namespace"`
    Target     string `json:"target"`      // deployment name / node name
    Replicas   int32  `json:"replicas"`    // scale 用
    Confidence float64 `json:"confidence"` // 来自诊断 Agent
    Reason     string `json:"reason"`      // 为什么执行
}

// Executor 执行 Agent 主体
type Executor struct {
    client          kubernetes.Interface
    dailyLimit      int               // 每日每动作上限
    executedToday   map[string]int    // 动作类型 -> 当日次数
    whitelist       map[string]bool   // 允许的动作类型
    requireApproval map[string]bool   // 需要人工审批的动作
}

func NewExecutor(client kubernetes.Interface, dailyLimit int) *Executor {
    return &Executor{
        client:        client,
        dailyLimit:    dailyLimit,
        executedToday: make(map[string]int),
        // 白名单是个护栏,没列进去的动作一律拒绝
        // 我踩过坑:曾经忘了把 drain 加白名单,结果 Node 故障时执行 Agent 直接报错
        whitelist: map[string]bool{
            "scale": true, "restart": true, "rollback": true, "drain": true,
        },
        // 破坏性动作必须人工审批,别让 LLM 一拍脑袋就把 Node drain 了
        requireApproval: map[string]bool{
            "drain": true, "rollback": true,
        },
    }
}

// Execute 执行动作,所有校验失败都返回 error,绝不悄悄吞掉
func (e *Executor) Execute(ctx context.Context, action Action) error {
    // 校验 1:白名单
    if !e.whitelist[action.Type] {
        return fmt.Errorf("action %s not in whitelist", action.Type)
    }

    // 校验 2:置信度门槛,破坏性动作要求更高
    minConfidence := 0.7
    if e.requireApproval[action.Type] {
        minConfidence = 0.85
    }
    if action.Confidence < minConfidence {
        return fmt.Errorf("confidence %.2f below threshold %.2f for action %s",
            action.Confidence, minConfidence, action.Type)
    }

    // 校验 3:当日执行次数
    if e.executedToday[action.Type] >= e.dailyLimit {
        return fmt.Errorf("daily limit %d reached for action %s", e.dailyLimit, action.Type)
    }

    // 校验 4:人工审批(这里只做模拟,实际接钉钉/飞书 webhook)
    if e.requireApproval[action.Type] {
        // 真实场景下这里应该 block 等审批结果,或者把 Action 入队等审批回执
        // 这里简化为直接拒绝,由协调 Agent 转入 PendingApproval 状态
        return fmt.Errorf("action %s requires human approval", action.Type)
    }

    // 所有校验通过,执行
    var err error
    switch action.Type {
    case "scale":
        err = e.scale(ctx, action)
    case "restart":
        err = e.restart(ctx, action)
    case "drain":
        err = e.drainNode(ctx, action)
    default:
        err = fmt.Errorf("unsupported action: %s", action.Type)
    }

    if err == nil {
        e.executedToday[action.Type]++
    }
    return err
}

// scale 扩缩容 Deployment
func (e *Executor) scale(ctx context.Context, action Action) error {
    dep, err := e.client.AppsV1().Deployments(action.Namespace).Get(ctx, action.Target, metav1.GetOptions{})
    if err != nil {
        return fmt.Errorf("get deployment failed: %w", err)
    }

    // 关键:直接改 spec.replicas,不要用 patch,因为我们要先读后写避免覆盖别人改动
    // 这里有竞态——并发两个 scale 会后写覆盖前写,必要时用 resourceVersion 做乐观锁
    replicas := action.Replicas
    if replicas == *dep.Spec.Replicas {
        return nil // 已经是目标副本数,幂等返回
    }
    dep.Spec.Replicas = &replicas

    _, err = e.client.AppsV1().Deployments(action.Namespace).Update(ctx, dep, metav1.UpdateOptions{})
    return err
}

// restart 删除 Pod 触发重建,比 rollout restart 更直接
func (e *Executor) restart(ctx context.Context, action Action) error {
    // action.Target 是 deployment 名,先按 label 找出所有 Pod
    dep, err := e.client.AppsV1().Deployments(action.Namespace).Get(ctx, action.Target, metav1.GetOptions{})
    if err != nil {
        return err
    }

    selector := labels.Set(dep.Spec.Selector.MatchLabels).AsSelector()
    podList, err := e.client.CoreV1().Pods(action.Namespace).List(ctx, metav1.ListOptions{
        LabelSelector: selector.String(),
    })
    if err != nil {
        return err
    }

    // 关键:删除策略是 Foreground,等 Pod 真的没了再返回
    // 不然你前脚删完,后脚 controller 还没创建新 Pod,业务就断一会儿
    // 还有,别一次删光所有 Pod,至少留一个健康的
    if len(podList.Items) <= 1 {
        return fmt.Errorf("only %d pods, refuse to restart all", len(podList.Items))
    }

    deleted := 0
    for _, pod := range podList.Items {
        // 跳过健康的 Pod,只删异常的
        if pod.Status.Phase == corev1.PodRunning && isPodReady(&pod) {
            continue
        }
        // Foreground 策略:先删 Pod,等 finalizer 清完
        // 这里的坑:Foreground 在某些 K8s 版本上对 Pod 不生效,实际行为是立即删除
        // 所以更稳的做法是用 metav1.DeleteOptions 里的 GracePeriodSeconds
        grace := int64(30) // 给 30s 优雅退出
        err := e.client.CoreV1().Pods(action.Namespace).Delete(ctx, pod.Name, metav1.DeleteOptions{
            GracePeriodSeconds: &grace,
        })
        if err != nil {
            return fmt.Errorf("delete pod %s failed: %w", pod.Name, err)
        }
        deleted++
        if deleted >= len(podList.Items)/2 {
            break // 一次最多删一半,避免雪崩
        }
    }
    return nil
}

// drainNode 驱逐节点上所有 Pod,破坏性最大,必须谨慎
func (e *Executor) drainNode(ctx context.Context, action Action) error {
    // 1. 先把节点标记为 unschedulable
    _, err := e.client.CoreV1().Nodes().Update(ctx,
        &corev1.Node{
            ObjectMeta: metav1.ObjectMeta{Name: action.Target},
            Spec:       corev1.NodeSpec{Unschedulable: true},
        }, metav1.UpdateOptions{})
    if err != nil {
        return fmt.Errorf("cordon node failed: %w", err)
    }

    // 2. 拿到节点上所有 Pod(排除 daemonset)
    podList, err := e.client.CoreV1().Pods("").List(ctx, metav1.ListOptions{
        FieldSelector: fmt.Sprintf("spec.nodeName=%s", action.Target),
    })
    if err != nil {
        return err
    }

    // 3. 逐个 evict
    // 这里只做骨架,真实场景还要处理 PDB(PodDisruptionBudget)
    // 忽略 PDB 直接驱逐可能导致服务可用副本数不足,是个常见坑
    for _, pod := range podList.Items {
        if isDaemonSetPod(&pod) {
            continue // DaemonSet Pod 跟着节点走,驱逐没意义
        }
        // Eviction 是个 subresource,比 Delete 多一层 PDB 检查
        evict := &policyv1.Eviction{
            ObjectMeta: metav1.ObjectMeta{
                Namespace: pod.Namespace,
                Name:      pod.Name,
            },
            DeleteOptions: metav1.DeleteOptions{
                GracePeriodSeconds: ptrInt64(30),
            },
        }
        err := e.client.PolicyV1().Evictions(pod.Namespace).Evict(ctx, evict)
        if err != nil {
            // 单个 Pod 驱逐失败不阻断整个 drain,记录下来继续
            // 真实实现应该把失败 Pod 名字收集起来回传给审计
            continue
        }
        // 给 controller 一点时间重新调度,不然下一秒就驱逐下一个可能造成压力
        time.Sleep(2 * time.Second)
    }
    return nil
}

// isPodReady 简单判断 Pod 是否 Ready
func isPodReady(pod *corev1.Pod) bool {
    for _, cond := range pod.Status.Conditions {
        if cond.Type == corev1.PodReady && cond.Status == corev1.ConditionTrue {
            return true
        }
    }
    return false
}

// isDaemonSetPod 通过 ownerReferences 判断是否 DaemonSet
func isDaemonSetPod(pod *corev1.Pod) bool {
    for _, ref := range pod.ObjectMeta.OwnerReferences {
        if ref.Kind == "DaemonSet" {
            return true
        }
    }
    return false
}

func ptrInt64(v int64) *int64 { return &v }

踩坑提示:EvictionDelete 是两码事,前者会检查 PDB,后者直接干掉。我见过有人图省事直接用 Delete,结果把一个只有 1 副本的服务干没了。还有 drain 一定要先 cordon 再 evict,不然你 evict 完 controller 又把新 Pod 调度到这个节点上,白干。

审计 Agent

审计 Agent 记录每次故障处理的完整上下文:

  • 原始告警信息
  • 各 Agent 的输入输出
  • 最终执行的修复动作
  • 执行前后的指标对比
  • 人工介入记录

这些数据不仅用于合规审计,也是持续优化诊断模型和修复策略的重要素材。

审计最实用的做法是把决策写入 K8s Event,这样运维直接 kubectl describe pod 就能看到这个 Pod 被自动修复过。同时再落一份到外部存储(ES/ClickHouse)做长期归档。

// agent/auditor/k8s_auditor.go
package auditor

import (
    "context"
    "encoding/json"
    "fmt"
    "time"

    corev1 "k8s.io/api/core/v1"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/client-go/kubernetes"
)

// AuditRecord 一次完整故障处理的审计记录
type AuditRecord struct {
    TraceID      string                 `json:"trace_id"`
    StartedAt    time.Time              `json:"started_at"`
    FinishedAt   time.Time              `json:"finished_at"`
    Alert        string                 `json:"alert"`
    Namespace    string                 `json:"namespace"`
    Target       string                 `json:"target"`
    Diagnosis    map[string]interface{} `json:"diagnosis"`
    Action       string                 `json:"action"`
    Result       string                 `json:"result"` // success/failed/skipped
    ErrorMessage string                 `json:"error_message,omitempty"`
}

// Auditor 审计 Agent
type Auditor struct {
    client    kubernetes.Interface
    eventSink EventSink // 外部存储,比如 ES/ClickHouse
}

// EventSink 外部审计存储抽象
type EventSink interface {
    Write(ctx context.Context, record AuditRecord) error
}

func NewAuditor(client kubernetes.Interface, sink EventSink) *Auditor {
    return &Auditor{client: client, eventSink: sink}
}

// WriteK8sEvent 把审计记录写成 K8s Event,绑到对应资源上
// 这样运维 kubectl describe 的时候直接能看到"这个 Pod 被自动修过"
func (a *Auditor) WriteK8sEvent(ctx context.Context, record AuditRecord) error {
    var involvedObject corev1.ObjectReference
    // 简化:这里假设 record.Target 是 Pod 名,实际要根据 Action 类型判断是 Pod/Deployment/Node
    involvedObject = corev1.ObjectReference{
        Kind:      "Pod",
        Namespace: record.Namespace,
        Name:      record.Target,
    }

    // Event 的 message 控制在 1KB 以内,太长会被截断
    // 完整数据走 eventSink 落外部存储
    summary := fmt.Sprintf("auto-remediation: alert=%s action=%s result=%s trace=%s",
        record.Alert, record.Action, record.Result, record.TraceID)

    // Event 名字带 traceID,方便关联
    // 这里有个坑:Event 名字长度限制 253 字符,namespace/traceID 拼起来可能超
    // 不过实际 traceID 一般 30 字符内,没问题
    event := &corev1.Event{
        ObjectMeta: metav1.ObjectMeta{
            Name: fmt.Sprintf("remediation-%s", record.TraceID),
            // 关键:不带 namespace 的话 cluster-scoped,绑不上 involvedObject
            Namespace: record.Namespace,
        },
        InvolvedObject: involvedObject,
        // 用自定义 reason 方便后续筛选
        Reason:  "AutoRemediation",
        Message: summary,
        // 关键:Type 是 Normal/Warning,别瞎填别的值,dashboard 会显示不出来
        Type:    eventTypeByResult(record.Result),
        Source: corev1.EventSource{
            Component: "aiops-remediation-agent",
        },
        FirstTimestamp: metav1.Time{Time: record.StartedAt},
        LastTimestamp:  metav1.Time{Time: record.FinishedAt},
        Count:          1,
    }

    _, err := a.client.CoreV1().Events(record.Namespace).Create(ctx, event, metav1.CreateOptions{})
    if err != nil {
        return fmt.Errorf("create k8s event failed: %w", err)
    }
    return nil
}

// WriteExternal 同时落一份到外部存储,做长期归档和分析
func (a *Auditor) WriteExternal(ctx context.Context, record AuditRecord) error {
    if a.eventSink == nil {
        return nil
    }
    return a.eventSink.Write(ctx, record)
}

// WriteAll 一把梭,K8s Event 和外部存储都写
// K8s Event 写失败不影响外部存储,反之亦然,互不阻塞
func (a *Auditor) WriteAll(ctx context.Context, record AuditRecord) error {
    // 先写外部存储,因为这个最关键,K8s Event 只是给运维看的
    if err := a.WriteExternal(ctx, record); err != nil {
        return fmt.Errorf("write external audit failed: %w", err)
    }
    // K8s Event 写失败只 warn,不影响主流程
    if err := a.WriteK8sEvent(ctx, record); err != nil {
        fmt.Printf("warn: write k8s event failed: %v\n", err)
    }
    return nil
}

// eventTypeByResult 根据执行结果返回 Event type
func eventTypeByResult(result string) string {
    switch result {
    case "success":
        return "Normal"
    case "failed", "skipped":
        return "Warning"
    default:
        return "Normal"
    }
}

// MarshalForLog 把审计记录序列化成 JSON,方便写日志或发到消息队列
func (r AuditRecord) MarshalForLog() ([]byte, error) {
    return json.Marshal(r)
}

踩坑提示:K8s Event 默认会被 kube-controller-manager 在 1 小时后清理(--event-ttl),所以审计数据必须有外部存储兜底。我见过有人只写 K8s Event,结果一周后想复盘发现啥都没了。还有 Event 的 message 字段超长会被截断到 1KB 左右,别在里面塞完整诊断 JSON。

协调 Agent 与状态机

协调 Agent 负责管理一次故障修复的完整生命周期。可以使用状态机模型:

Detected -> Observing -> Diagnosing -> PendingApproval -> Executing -> Verifying -> Resolved/Failed

状态转换通过消息总线或工作流引擎驱动。各 Agent 之间通过定义良好的消息协议通信,避免紧耦合。

类比:状态机就像快递的物流轨迹——“已揽收 → 运输中 → 派送中 → 已签收 / 拒收”。每个包裹(一次故障)都有自己独立的轨迹,且只能按允许的路径推进,不能从"已签收"突然跳回"运输中"。

下面这张状态机图把所有合法流转画全了,包括任意"进行中"状态超时都可直接进 Failed:

stateDiagram-v2
    [*] --> Observing
    Observing --> Diagnosing: observe_done
    Diagnosing --> Executing: diagnose_done
    Diagnosing --> PendingApproval: need_approval
    PendingApproval --> Executing: approved
    PendingApproval --> Failed: rejected
    Executing --> Verifying: execute_done
    Executing --> Failed: execute_failed
    Verifying --> Resolved: verify_passed
    Verifying --> Failed: verify_failed
    Observing --> Failed: timeout
    Diagnosing --> Failed: timeout
    Executing --> Failed: timeout
    Verifying --> Failed: timeout
    Resolved --> [*]
    Failed --> [*]

光看这张图其实挺好懂,但真正实现的时候状态机是最容易写烂的。我个人的经验是:状态机千万别用一堆 if state == "xxx" 串起来,三个月后你自己都看不懂。用一个明确的 transition 表,谁能转到谁、转移条件是什么,全列出来,代码即文档。

// agent/coordinator/state_machine.go
package coordinator

import (
    "context"
    "errors"
    "fmt"
    "sync"
    "time"
)

// State 故障处理的状态
type State string

const (
    StateDetected         State = "detected"
    StateObserving        State = "observing"
    StateDiagnosing       State = "diagnosing"
    StatePendingApproval  State = "pending_approval"
    StateExecuting        State = "executing"
    StateVerifying        State = "verifying"
    StateResolved         State = "resolved"
    StateFailed           State = "failed"
)

// Event 触发状态转移的事件
type Event string

const (
    EventObserveDone      Event = "observe_done"
    EventDiagnoseDone     Event = "diagnose_done"
    EventNeedApproval     Event = "need_approval"
    EventApproved         Event = "approved"
    EventRejected         Event = "rejected"
    EventExecuteDone      Event = "execute_done"
    EventExecuteFailed    Event = "execute_failed"
    EventVerifyPassed     Event = "verify_passed"
    EventVerifyFailed     Event = "verify_failed"
    EventTimeout          Event = "timeout"
)

// transition 状态转移定义
type transition struct {
    from  State
    event Event
    to    State
    // action 在转移时执行,返回 error 则状态机进入 Failed
    action func(ctx context.Context, inc *Incident) error
}

// Incident 一次故障处理的完整上下文
type Incident struct {
    TraceID    string    `json:"trace_id"`
    State      State     `json:"state"`
    Alert      string    `json:"alert"`
    Namespace  string    `json:"namespace"`
    Target     string    `json:"target"`
    CreatedAt  time.Time `json:"created_at"`
    UpdatedAt  time.Time `json:"updated_at"`

    // 各阶段的产物,避免 Agent 间重复传参
    Observation map[string]interface{} `json:"observation,omitempty"`
    Diagnosis   map[string]interface{} `json:"diagnosis,omitempty"`
    Action      map[string]interface{} `json:"action,omitempty"`
    Audit       map[string]interface{} `json:"audit,omitempty"`

    mu sync.Mutex
}

// Coordinator 协调 Agent 主体
type Coordinator struct {
    transitions []transition
    incidents   map[string]*Incident // traceID -> Incident
    mu          sync.RWMutex
}

func NewCoordinator() *Coordinator {
    c := &Coordinator{incidents: make(map[string]*Incident)}
    c.transitions = []transition{
        {StateDetected, EventObserveDone, StateDiagnosing, c.doDiagnose},
        {StateObserving, EventObserveDone, StateDiagnosing, c.doDiagnose},
        {StateDiagnosing, EventDiagnoseDone, StateExecuting, c.doExecute},
        {StateDiagnosing, EventNeedApproval, StatePendingApproval, nil},
        {StatePendingApproval, EventApproved, StateExecuting, c.doExecute},
        {StatePendingApproval, EventRejected, StateFailed, c.doAuditOnReject},
        {StateExecuting, EventExecuteDone, StateVerifying, c.doVerify},
        {StateExecuting, EventExecuteFailed, StateFailed, c.doAuditOnFail},
        {StateVerifying, EventVerifyPassed, StateResolved, c.doAuditOnResolve},
        {StateVerifying, EventVerifyFailed, StateFailed, c.doAuditOnFail},
        // 超时从任何"进行中"状态都可以直接跳 Failed
        {StateObserving, EventTimeout, StateFailed, c.doAuditOnFail},
        {StateDiagnosing, EventTimeout, StateFailed, c.doAuditOnFail},
        {StateExecuting, EventTimeout, StateFailed, c.doAuditOnFail},
        {StateVerifying, EventTimeout, StateFailed, c.doAuditOnFail},
    }
    return c
}

// HandleEvent 处理一个状态转移事件
// 这是协调 Agent 的核心入口,所有 Agent 完成工作后都通过这里通知协调者
func (c *Coordinator) HandleEvent(ctx context.Context, traceID string, event Event, payload map[string]interface{}) error {
    c.mu.RLock()
    inc, ok := c.incidents[traceID]
    c.mu.RUnlock()
    if !ok {
        return fmt.Errorf("incident %s not found", traceID)
    }

    inc.mu.Lock()
    defer inc.mu.Unlock()

    // 把 payload 塞到对应字段,方便后续 action 用
    c.mergePayload(inc, event, payload)

    // 找到匹配的转移规则
    var matched *transition
    for i := range c.transitions {
        t := &c.transitions[i]
        if t.from == inc.State && t.event == event {
            matched = t
            break
        }
    }
    if matched == nil {
        // 没有匹配的转移,可能是个迟到的消息,记一下但不报错
        // 我踩过坑:曾经这里直接返回 error,结果上游疯狂重试
        return fmt.Errorf("no transition from %s on event %s", inc.State, event)
    }

    // 执行转移动作
    if matched.action != nil {
        if err := matched.action(ctx, inc); err != nil {
            // action 失败直接进入 Failed
            inc.State = StateFailed
            inc.UpdatedAt = time.Now().UTC()
            return fmt.Errorf("action failed: %w", err)
        }
    }

    inc.State = matched.to
    inc.UpdatedAt = time.Now().UTC()
    return nil
}

// StartIncident 开启一次故障处理
func (c *Coordinator) StartIncident(traceID, alert, namespace, target string) *Incident {
    inc := &Incident{
        TraceID:   traceID,
        State:     StateObserving, // 直接进入 observing
        Alert:     alert,
        Namespace: namespace,
        Target:    target,
        CreatedAt: time.Now().UTC(),
        UpdatedAt: time.Now().UTC(),
    }
    c.mu.Lock()
    c.incidents[traceID] = inc
    c.mu.Unlock()
    return inc
}

// doDiagnose 转移到诊断阶段,触发诊断 Agent
// 这里只是骨架,真实实现是发消息到消息队列让诊断 Agent 消费
func (c *Coordinator) doDiagnose(ctx context.Context, inc *Incident) error {
    if inc.Observation == nil {
        return errors.New("observation is empty, cannot diagnose")
    }
    // 实际场景:发送消息到 diagnoser queue
    // 这里只做占位
    return nil
}

// doExecute 触发执行 Agent
func (c *Coordinator) doExecute(ctx context.Context, inc *Incident) error {
    if inc.Diagnosis == nil {
        return errors.New("diagnosis is empty, cannot execute")
    }
    return nil
}

// doVerify 触发验证:观察修复后指标是否恢复
func (c *Coordinator) doVerify(ctx context.Context, inc *Incident) error {
    // 验证阶段的关键是"等"——给系统一点时间收敛
    // 但不能无限等,必须设超时,超时就触发 EventTimeout -> Failed
    return nil
}

// doAuditOnResolve 成功后写审计
func (c *Coordinator) doAuditOnResolve(ctx context.Context, inc *Incident) error {
    inc.Audit = map[string]interface{}{
        "result":     "success",
        "finished_at": time.Now().UTC(),
    }
    return nil
}

// doAuditOnFail 失败后写审计
func (c *Coordinator) doAuditOnFail(ctx context.Context, inc *Incident) error {
    inc.Audit = map[string]interface{}{
        "result":     "failed",
        "finished_at": time.Now().UTC(),
    }
    return nil
}

// doAuditOnReject 人工拒绝后写审计
func (c *Coordinator) doAuditOnReject(ctx context.Context, inc *Incident) error {
    inc.Audit = map[string]interface{}{
        "result":     "rejected",
        "finished_at": time.Now().UTC(),
    }
    return nil
}

// mergePayload 把事件附带的数据塞到 Incident 对应字段
func (c *Coordinator) mergePayload(inc *Incident, event Event, payload map[string]interface{}) {
    if payload == nil {
        return
    }
    switch event {
    case EventObserveDone:
        inc.Observation = payload
    case EventDiagnoseDone, EventNeedApproval:
        inc.Diagnosis = payload
    case EventExecuteDone, EventExecuteFailed:
        inc.Action = payload
    }
}

踩坑提示:状态机最坑的是"迟到的事件"。比如诊断 Agent 已经回了结果,但状态机因为超时进入了 Failed,这时候 EventDiagnoseDone 来了,你要决定是丢弃还是复活。我个人倾向于丢弃+审计记录,因为复活很容易把状态搞乱。还有一个常见问题:同一个 traceID 的并发事件,必须加锁(我用了 inc.mu),不然两个事件同时改 state 就热闹了。

关键技术挑战

  1. 幻觉与误判:LLM 可能给出错误诊断,需要设置置信度阈值,并在关键操作前要求人工确认。
  2. 权限控制:执行 Agent 需要最小权限原则,避免越权操作。
  3. 并发与竞态:多个告警同时触发时,需要避免冲突操作。
  4. 可解释性:每次自动修复都应能解释"为什么做"和"做了什么"。

这四条里我特别想强调并发。多个告警同时来的时候,最常见的事故是"两个 Agent 都觉得 Redis 副本数不够,一个扩到 3,另一个扩到 5,最后一个写覆盖前一个"。解法是协调 Agent 用 traceID 维度做串行化——同一个 target 同时只允许一个 incident 在 Executing 状态。代码不复杂,但忘了加就是事故。

收尾的一些碎碎念

多 Agent 协同是实现复杂 K8s 故障自动修复的有效架构。通过观测、诊断、执行、审计等角色的分工协作,系统可以在保证安全可控的前提下,大幅提升故障响应速度和恢复效率。这也是云原生 AIOps 平台的核心能力之一。

写完这套系统我自己最大的感受是:别让 Agent 之间耦合得太紧。一开始为了图省事,我们让诊断 Agent 直接调执行 Agent 的函数,结果诊断 Agent 跟着执行 Agent 一起挂过好几次。后来强制改成消息驱动,谁也不认识谁,只认识 Message,整个系统稳定性立刻上一个台阶。这玩意儿跟微服务拆分是一个道理——边界比实现更重要。

自测题与动手练习

自测题(合上书能答出来,才算懂):

  1. 作者为什么说"单体 controller 不适合 K8s 故障自动修复"?它最致命的风险是什么?

    答:职责耦合导致一处 OOM / 打爆 API 就全链路挂掉;拆成多 Agent 是为了故障隔离、炸的时候不连坐

  2. 五大 Agent 里哪个最容易被低估?为什么说它"不是日志,是责任边界"?

    答:审计 Agent。没有它,自动修复把集群搞炸了都不知道谁干的,审计记录是事后复盘与问责的依据。

  3. Message envelope 里的 trace_id 为什么必须从告警入口就生成、一路传下去?强制 UTC 时间戳又是为什么?

    答:trace_id 串联一次故障的所有 Agent 输入输出,便于审计与排障;UTC 避免跨时区集群对账时时间错乱。

  4. 执行 Agent 的四道校验闸门分别是什么?为什么 drain 这类破坏性动作要强制人工审批?

    答:白名单 → 置信度门槛 → 当日次数上限 → 人工审批。drain 会驱逐整节点 Pod,误执行后果严重,必须人确认。

  5. 协调 Agent 状态机里"迟到的事件"应该怎么处理?为什么作者倾向于"丢弃 + 审计"而不是"复活"?

    答:状态已因超时进入 Failed 后到达的 diagnose_done 应丢弃,避免把状态搞乱;复活容易引发状态机不一致。

动手练习(建议真做一遍):

  1. 用 client-go 起一个 Pod informer,故意制造一个 CrashLoopBackOff 的 Pod,验证观测 Agent 能否识别异常、并带 trace_id 发出 Observation 消息。
  2. 拿一条真实 Pod 崩溃日志,手工拼出一条 Observation JSON,喂给诊断 Agent 的 buildPrompt,看生成的 prompt 是否包含完整现场(status / events / 最近日志 / 相似案例)。
  3. 在状态机里新增一个 ManualPause 状态并补上对应 transition 规则,跑单元测试验证"迟到的 diagnose_done"不会被错误复活。

本章小结

  • 多 Agent 协同的本质是故障隔离 + 职责边界,比单体 controller 更抗炸、更好维护。
  • 统一消息信封 + trace_id + UTC 时间戳,是跨 Agent 协作的命脉。
  • 执行 Agent 的护栏(白名单 / 置信度 / 审批)+ 审计 Agent 的外部存储兜底,决定系统能不能"安全地自动"。
  • 状态机用显式 transition 表驱动,避免 if-else 串状态;迟到的事件宁可丢弃也别复活。
  • 下一篇可深入每个 Agent 的内部算法,比如诊断 Agent 的 RAG 召回与 prompt 工程、执行 Agent 的并发串行化。
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!