审批工作流引擎模块

2025-10-15T10:30:00+08:00 | 30分钟阅读 | 更新于 2025-10-15T10:30:00+08:00

@

学习目标

本篇讲透企业自研 PaaS 平台的「审批工作流引擎」——它是所有敏感操作的统一闸门。学完本篇,你应该能够:

  1. 说清楚审批流如何作为其它模块(发布、扩缩容、密钥变更、集群纳管)敏感操作的强制前置闸门。
  2. 基于可视化流程编排引擎(BPMN 语义)设计多级审批、会签/或签、条件分支的生产级流程。
  3. 配置超时升级、催办、审批人委派/代理,避免流程卡死造成业务阻塞。
  4. 把审批事件接入 K8s API 网关,实现「审批通过才执行」的不可绕过机制。
  5. 将审批操作、IM 通知、告警升级指标接入 Prometheus 联邦,并保证审计留痕 90 天+。

前置知识

  • 已了解 _SHARED.md 全局约束:多套物理隔离 K8s 集群、统一 K8s API 网关、Namespace 租户隔离;Prometheus 联邦 + Thanos 长期存储、指标带 tenant/app/env/cluster 标签;RBAC 细粒度、审计 90 天+、敏感配置加密、镜像扫描。
  • 了解模块 01 租户与 RBAC、模块 02 集群纳管、模块 03 应用生命周期基本概念。
  • 知道 BPMN 基本元素(开始/任务/网关/结束)与 IM 机器人(企业微信/钉钉)Webhook 机制。

本章你会动手做的事

  1. 在可视化编排器里画一条「生产环境发布 = 申请人 → 租户管理员 → 平台运维」三级审批流。
  2. 故意让某一审批节点超时,观察升级给上级并触发企业微信催办。
  3. 直接绕过审批网关调用 K8s API,验证请求被网关策略拒绝并记入审计。

项目结构与文件清单

本篇所有代码自洽:Go 引用的字段名、指标名、RBAC 资源名与下方 YAML / Prometheus 配置完全一致。下文各章会逐个文件完整展示,此处先给出目录树(根目录 module-11-approval-workflow/)。

module-11-approval-workflow/
├── go.mod
├── deploy/
│   ├── approvalticket-crd.yaml     # ① 审批工单 CRD(流程实例载体)
│   ├── engine-rbac.yaml            # ② 引擎自身 RBAC + Webhook 接收 SA/ClusterRole
│   ├── executor-rbac.yaml          # ③ 审批通过后「执行用」受限 RBAC(仅租户 ns)
│   └── prometheus-rules.yaml       # ④ 审批延迟/超时/敏感操作 告警规则
├── config/
│   └── servicemonitor.yaml         # 联邦抓取配置(带 tenant 标签)
├── internal/
│   ├── apis/approval/v1/ticket.go  # ApprovalTicket 类型 + CRD 转换
│   ├── metrics/metrics.go          # Prometheus 指标(tenant/app/env/cluster)
│   ├── audit/audit.go              # WORM 不可变审计留痕
│   ├── notify/im.go                # 企业微信 / 钉钉 IM 通知
│   ├── engine/
│   │   ├── statemachine.go         # 状态机:多级 / 会签(AND) / 或签(OR)
│   │   ├── executor.go             # 审批通过后经 API 网关 client-go 执行 K8s 变更
│   │   └── engine.go               # 控制器 Reconcile(超时升级 + 落地执行)
│   └── gateway/token.go            # 向统一 K8s API 网关申请签名执行凭证
├── cmd/engine/main.go              # 引擎启动入口(含 /metrics)
└── runbook.md                      # 运维 SOP:建单/审批/查结果/绕行验证

约定:所有 Go 文件头标注 // file: <相对路径>,YAML 文件头标注 # file: <相对路径>。指标统一携带 tenant/app/env/cluster 四维标签;统一 K8s API 网关地址由环境变量 K8S_GATEWAY_URL 注入(如 https://k8s-gateway.paas.internal)。


一、模块概述与企业商用价值

类比:审批工作流引擎就像银行金库的「双钥匙制度」。任何打开金库(执行敏感 K8s 操作)的动作,都必须先有两把(或更多)分别保管的钥匙同时转动——单独一把钥匙、或者事后补手续,门都不会开。它把「谁能做什么」从口头约定变成了平台强制的、可审计的硬规则。

在企业自研 PaaS 平台里,几乎所有会动生产环境的高风险动作——生产发布、集群纳管、密钥变更、资源配额超限、跨集群批量运维——都不能由单人直接执行。审批工作流引擎是这些操作的统一前置闸门:没有有效的审批工单,K8s API 网关拒绝放行。

适用角色与痛点:

  • 平台管理员:面对数十租户、数百应用,无法靠人盯人保证「生产变更必经审批」。
  • 租户管理员 / 研发负责人:需要一条清晰、可追溯、不被上级「口头放行」绕过的变更链路,满足金融/政企合规审计(等保、ISO27001、SOX)。
  • 安全合规团队:需要证明所有敏感操作都有完整留痕,且无法被运维私下绕过。

平台不可替代性:审批不是「工单系统」那么简单。它必须与执行系统强绑定——审批通过 ≠ 仅仅记录一条「同意」,而是向 K8s API 网关签发一份有时效、有范围的执行凭证,网关在收到执行请求时回验凭证。这正是「不可绕过」的核心。

下图说明审批工作流引擎在整个 PaaS 平台中的定位:它横切于所有敏感操作模块与 K8s API 网关之间。

graph LR
  A[研发 / 租户管理员] -->|提交敏感操作申请| B((审批工作流引擎))
  C[应用发布模块] -->|生产发布需审批| B
  D[集群纳管模块] -->|纳管/摘除需审批| B
  E[密钥 / 配置变更] -->|加密变更需审批| B
  F[自动化批量运维] -->|高危批量需审批| B
  B -->|审批通过签发执行凭证| G[统一 K8s API 网关]
  B -->|审批拒绝 / 超时| H[拒绝执行 + 审计留痕]
  G -->|受限调谐| I[生产 / 预发布 / 测试 物理隔离集群]
  B -->|操作指标 + 通知| J[Prometheus 联邦 + 企业微信/钉钉]

这张图展示了审批引擎作为「横切闸门」的位置:任何敏感操作都要先过它,再经统一网关落到具体集群。

1.1 落地形态:用 CRD 承载流程实例

我们把「一个审批工单」建模为 K8s 自定义资源 ApprovalTicket(CRD 见 deploy/approvalticket-crd.yaml)。好处:工单天然享受 etcd 持久化、RBAC、审计、以及被联邦监控采集的状态指标;引擎以控制器方式 Reconcile 它,与平台其它基于 K8s 的模块风格统一。

graph TD
  A[控制台/发布模块] -->|kubectl apply / API| B[ApprovalTicket CR]
  B --> C[审批工作流引擎 Controller]
  C -->|Webhook 回调| D[企业微信/钉钉 审批人]
  C -->|全通过| E[gateway/token.go 申请执行凭证]
  E -->|Bearer Token(有时效/范围)| F[统一 K8s API 网关]
  F -->|client-go 受限调谐| G[目标物理隔离集群]
  C -->|每一步| H[audit.go WORM 审计]
  C -->|指标| I[Prometheus 联邦]

上图是「代码视角」的端到端链路:本章后续会逐个文件给出可编译实现。


二、细分功能详解(商用生产级)

下面用表格区分【基础能力】与【高级企业增值能力】。标注「★禁止删减」的是金融/政企级强制项,删除即视为不合格

类别功能说明关键约束
【基础能力】①可视化流程编排基于 BPMN 语义的节点/条件分支拖拽编排,流程定义版本化存储流程定义写入审计
【基础能力】②多级审批串行/并行多级,支持部门→租户→平台三层★留痕
【基础能力】③审批人委派/代理审批人休假可指定代理人,代理范围与时间受限代理操作记原审批人
【基础能力】④与 IM 集成通知企业微信/钉钉推送待办、催办、结果通知内容脱敏
【高级企业增值能力】⑤会签/或签会签=全部同意才过;或签=一人同意即过,适配不同风险场景★不可绕过
【高级企业增值能力】⑥超时升级与催办节点超时自动升级上级并催办,防流程卡死★不可绕过
【高级企业增值能力】⑦审批与权限/操作绑定审批通过才向 K8s API 网关签发执行凭证,凭证有时效与范围★不可绕过
【高级企业增值能力】⑧审批审计与不可绕过全链路留痕、防篡改,审计留存 90 天+,且拒绝「事后补单」★留痕+★不可绕过

下面用分级列表补充细节,并明确标注禁止删减项:

  • 可视化流程编排

    • 支持开始、用户任务、服务任务、排他网关(条件分支)、并行网关(会签容器)、结束节点。
    • 流程定义以版本号入库,变更需走「流程定义变更审批」,新版本不影响在途工单。
    • 每个节点可绑定 RBAC 角色(「谁能当这个节点的审批人」),避免越权审批。
  • 多级审批与会签/或签 ★禁止删减:多级

    • 串行多级:适用于高风险(生产发布、密钥变更),逐级上升。
    • 会签(AND):多人必须全部同意,常用于跨团队联签(研发+运维+安全)。
    • 或签(OR):任一同意即可,常用于值班制场景(谁先看到谁批)。
    • 条件分支:按 env=prod / 资源规模 / 风险等级自动选择审批链。
  • 审批人委派与代理

    • 支持临时委派(指定起止时间),超出时间自动失效。
    • 代理审批在审计中同时记录「实际代理人 + 被代理人」。
  • 超时升级与催办 ★禁止删减:不可绕过

    • 每个节点可配 SLA:超时 N 分钟 → 催办 IM 提醒;超时 M 分钟 → 升级至上级并记告警。
    • 升级链路本身也是审批链的一部分,不能因升级而跳过审批。
  • 与 IM(企业微信/钉钉)集成通知

    • 待办推送携带工单号、申请人、目标集群/环境、风险等级(脱敏后的操作摘要)。
    • 支持在 IM 内「快捷审批」,但高敏感操作强制跳转网关二次认证页面。
  • 审批审计与不可绕过 ★禁止删减:留痕

    • 申请、转交、委派、每一级审批动作、超时、升级、最终执行结果全部入审计库,WORM 式防篡改。
    • 平台不提供「跳过审批」「事后补单」「管理员强制通过」接口——这是物理层面的不可绕过,由 K8s API 网关策略强制。
  • 与权限/操作绑定(审批通过才放行) ★禁止删减:不可绕过

    • 审批通过后,引擎向统一 K8s API 网关申请一份「执行凭证」:含生效范围(tenant/namespace/cluster)、允许动作(如 deploy.apply)、有效期(默认 30 分钟,可配)。
    • 网关在收到执行请求时校验凭证,无凭证或凭证过期/越界一律 403 并记录。

2.1 IaC:审批工单 CRD

# file: deploy/approvalticket-crd.yaml —— 定义 ApprovalTicket 自定义资源,承载流程实例(申请人、租户/环境/集群、审批链、会签/或签策略、SLA、执行结果、审计引用)。

# file: deploy/approvalticket-crd.yaml
apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
  name: approvaltickets.approval.paas.example.com
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
spec:
  group: approval.paas.example.com
  names:
    kind: ApprovalTicket
    listKind: ApprovalTicketList
    plural: approvaltickets
    singular: approvalticket
    shortNames: ["ap"]
  scope: Namespaced
  versions:
    - name: v1
      served: true
      storage: true
      subresources:
        status: {}
      schema:
        openAPIV3Schema:
          type: object
          properties:
            spec:
              type: object
              required: ["applicant", "tenant", "env", "cluster", "namespace", "action", "policy", "chain"]
              properties:
                applicant:
                  type: string
                  description: "申请人(租户内账号)"
                tenant:
                  type: string
                  description: "租户标识,用于 Namespace 隔离与指标标签"
                env:
                  type: string
                  enum: ["dev", "test", "stage", "prod"]
                cluster:
                  type: string
                  description: "目标物理隔离集群名"
                namespace:
                  type: string
                  description: "目标 Namespace(租户强隔离)"
                action:
                  type: string
                  description: "受控动作,如 deploy.apply / quota.update / secret.rotate"
                riskLevel:
                  type: string
                  enum: ["low", "mid", "high"]
                  default: mid
                policy:
                  type: string
                  enum: ["and", "or"]
                  description: "and=会签(全部同意) or=或签(任一同意)"
                chain:
                  type: array
                  description: "多级审批链,按顺序串行"
                  items:
                    type: object
                    required: ["level", "role"]
                    properties:
                      level:
                        type: integer
                        minimum: 1
                      role:
                        type: string
                        description: "该级审批人所需 RBAC 角色,如 tenant-admin"
                      slaSeconds:
                        type: integer
                        description: "该级 SLA,超时则升级"
                        default: 1800
                payloadRef:
                  type: string
                  description: "待执行 manifest 引用(加密存储于 Secret/对象存储)"
                summary:
                  type: string
                  description: "脱敏后的操作摘要,用于 IM 推送"
            status:
              type: object
              properties:
                phase:
                  type: string
                  enum: ["Submitted", "Approving", "Approved", "Rejected", "Executing", "Succeeded", "Failed"]
                currentLevel:
                  type: integer
                approvals:
                  type: array
                  items:
                    type: object
                    properties:
                      level:
                        type: integer
                      approver:
                        type: string
                      decision:
                        type: string
                        enum: ["approve", "reject"]
                      at:
                        type: string
                        format: date-time
                      deviceId:
                        type: string
                escalationCount:
                  type: integer
                execTokenRef:
                  type: string
                  description: "网关签发的执行凭证引用(不落明文)"
                execResult:
                  type: string
                auditRef:
                  type: string
                  description: "WORM 审计链尾 hash,用于防篡改校验"
      additionalPrinterColumns:
        - name: Phase
          type: string
          jsonPath: .status.phase
        - name: Env
          type: string
          jsonPath: .spec.env
        - name: Cluster
          type: string
          jsonPath: .spec.cluster
        - name: Action
          type: string
          jsonPath: .spec.action
        - name: Age
          type: date
          jsonPath: .metadata.creationTimestamp

2.2 IaC:引擎自身 RBAC + Webhook 接收身份

# file: deploy/engine-rbac.yaml —— 引擎以 approval-engine ServiceAccount 运行;它需要读取/更新 ApprovalTicket读取流程定义,但绝不直接持有目标集群 kubeconfig(执行时经网关)。同时给出一个 ClusterRole 供 Webhook 接收器(IM 回调入口)使用。

# file: deploy/engine-rbac.yaml
apiVersion: v1
kind: ServiceAccount
metadata:
  name: approval-engine
  namespace: paas-approval
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: approval-engine-ticket-rw
  namespace: paas-approval
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
rules:
  - apiGroups: ["approval.paas.example.com"]
    resources: ["approvaltickets"]
    verbs: ["get", "list", "watch", "create", "update", "patch"]
  - apiGroups: ["approval.paas.example.com"]
    resources: ["approvaltickets/status"]
    verbs: ["get", "update", "patch"]
  - apiGroups: [""]
    resources: ["configmaps"]
    verbs: ["get", "list"]
    # 仅用于读取流程定义与 IM Webhook 配置,限定本命名空间
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: approval-engine-ticket-rw
  namespace: paas-approval
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
subjects:
  - kind: ServiceAccount
    name: approval-engine
    namespace: paas-approval
roleRef:
  kind: Role
  name: approval-engine-ticket-rw
  apiGroup: rbac.authorization.k8s.io
---
# Webhook 接收器:IM 回调把审批决定 POST 进来,由专属 SA 写入工单
apiVersion: v1
kind: ServiceAccount
metadata:
  name: approval-webhook-receiver
  namespace: paas-approval
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
  name: approval-webhook-receiver
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
rules:
  - apiGroups: ["approval.paas.example.com"]
    resources: ["approvaltickets"]
    verbs: ["get", "update", "patch"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRoleBinding
metadata:
  name: approval-webhook-receiver
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
subjects:
  - kind: ServiceAccount
    name: approval-webhook-receiver
    namespace: paas-approval
roleRef:
  kind: ClusterRole
  name: approval-webhook-receiver
  apiGroup: rbac.authorization.k8s.io

2.3 IaC:审批通过后「执行用」受限 RBAC

# file: deploy/executor-rbac.yaml —— 关键安全边界:执行身份不是引擎 SA,而是网关在「执行凭证」范围下临时映射的受限身份。这里给出的是网关侧 executor 角色模板:仅能在目标租户 Namespace 做「该动作」允许的最小动词,绝无 *、绝无跨 ns。引擎经网关调用,网关用此角色代理执行。

# file: deploy/executor-rbac.yaml
# 该 Role 由统一 K8s API 网关在「执行凭证」scope 内临时绑定到 executor 身份,
# 仅作用于目标租户 Namespace,仅含最小动词。
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: approval-executor-deploy
  namespace: tenant-a-prod          # 由凭证 scope.tenant/namespace 动态渲染
  labels:
    tenant: tenant-a
    app: approval-engine
    env: prod
    cluster: prod-sh
spec:
  rules:
    # deploy.apply 仅允许对 Deployment 做 apply/scale
    - apiGroups: ["apps"]
      resources: ["deployments"]
      verbs: ["get", "list", "create", "update", "patch"]
    - apiGroups: [""]
      resources: ["configmaps"]
      verbs: ["get", "list", "create", "update", "patch"]
    # 明确拒绝 secrets/rbac/exec 等高危资源,证明「凭证越界即拒绝」
    # (高危资源不出现在 rules 中 = 网关侧 403)
---
apiVersion: v1
kind: ServiceAccount
metadata:
  name: approval-executor
  namespace: paas-gateway-proxy
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: approval-executor-deploy
  namespace: tenant-a-prod
  labels:
    tenant: tenant-a
    app: approval-engine
    env: prod
    cluster: prod-sh
subjects:
  - kind: ServiceAccount
    name: approval-executor
    namespace: paas-gateway-proxy
roleRef:
  kind: Role
  name: approval-executor-deploy
  apiGroup: rbac.authorization.k8s.io

功能架构图如下,呈现模块内部子能力与外部边界。

graph TD
  A[流程定义中心
BPMN 版本化] --> B[流程实例引擎
串行/并行/网关] B --> C{会签/或签判定} C -->|全部同意| D[通过] C -->|任一同意| E[通过] C -->|拒绝| F[驳回 + 审计] B --> G[超时升级器
SLA 催办] G --> H[上级审批人] B --> I[IM 通知网关
企业微信/钉钉] D --> J[执行凭证签发器] J -->|有时效/范围| K[统一 K8s API 网关] B --> L[审计留痕服务
90天+ WORM] K --> M[目标物理隔离集群]

这张功能架构图说明:无论走会签还是或签,最终都汇聚到「执行凭证签发器」,而凭证是真正能打动网关的唯一合法凭据。


三、底层架构联动设计

审批工作流引擎不是孤立的工单系统,它和全局三条约束强联动。

3.1 Go 类型与 CRD 转换

// file: internal/apis/approval/v1/ticket.go —— 与 CRD schema 字段逐一对应的 Go 结构体,并提供与 unstructured 的互转(避免引入 controller-gen 代码生成,保证开箱可编译)。

// file: internal/apis/approval/v1/ticket.go
package v1

import (
	"time"

	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
	"k8s.io/apimachinery/pkg/runtime"
)

// Group/Version/Resource —— 与 deploy/approvalticket-crd.yaml 完全一致
const (
	Group    = "approval.paas.example.com"
	Version  = "v1"
	Resource = "approvaltickets"
	Kind     = "ApprovalTicket"
)

// ApprovalTicket 对应 CRD 实例。
type ApprovalTicket struct {
	metav1.TypeMeta   `json:",inline"`
	metav1.ObjectMeta `json:"metadata,omitempty"`

	Spec ApprovalTicketSpec `json:"spec,omitempty"`
	Status ApprovalTicketStatus `json:"status,omitempty"`
}

type ApprovalTicketSpec struct {
	Applicant  string           `json:"applicant"`
	Tenant     string           `json:"tenant"`
	Env        string           `json:"env"` // dev/test/stage/prod
	Cluster    string           `json:"cluster"`
	Namespace  string           `json:"namespace"`
	Action     string           `json:"action"` // deploy.apply / quota.update / secret.rotate
	RiskLevel  string           `json:"riskLevel,omitempty"`
	Policy     string           `json:"policy"` // and=会签 or=或签
	Chain      []ApprovalLevel  `json:"chain"`
	PayloadRef string           `json:"payloadRef,omitempty"`
	Summary    string           `json:"summary,omitempty"`
}

type ApprovalLevel struct {
	Level     int    `json:"level"`
	Role      string `json:"role"` // 该级审批人所需 RBAC 角色
	SLASeconds int   `json:"slaSeconds,omitempty"`
}

type ApprovalTicketStatus struct {
	Phase          string          `json:"phase"` // Submitted/Approving/Approved/Rejected/Executing/Succeeded/Failed
	CurrentLevel   int             `json:"currentLevel,omitempty"`
	Approvals      []ApprovalRecord `json:"approvals,omitempty"`
	EscalationCount int            `json:"escalationCount,omitempty"`
	ExecTokenRef   string          `json:"execTokenRef,omitempty"`
	ExecResult     string          `json:"execResult,omitempty"`
	AuditRef       string          `json:"auditRef,omitempty"`
}

type ApprovalRecord struct {
	Level    int       `json:"level"`
	Approver string    `json:"approver"`
	Decision string    `json:"decision"` // approve/reject
	At       time.Time `json:"at"`
	DeviceID string    `json:"deviceId,omitempty"`
}

func (t *ApprovalTicket) GetObjectKind() runtime.ObjectKind { return &t.TypeMeta }

// ToUnstructured 用于经 dynamic client 写回 apiserver。
func (t *ApprovalTicket) ToUnstructured() (*unstructured.Unstructured, error) {
	m, err := runtime.DefaultUnstructuredConverter.ToUnstructured(t)
	if err != nil {
		return nil, err
	}
	return &unstructured.Unstructured{Object: m}, nil
}

// FromUnstructured 用于经 dynamic client 读回 apiserver 对象。
func FromUnstructured(u *unstructured.Unstructured) (*ApprovalTicket, error) {
	t := &ApprovalTicket{}
	if err := runtime.DefaultUnstructuredConverter.FromUnstructured(u.Object, t); err != nil {
		return nil, err
	}
	return t, nil
}

3.2 Prometheus 指标(带 tenant/app/env/cluster)

// file: internal/metrics/metrics.go —— 所有审批指标统一四维标签,便于联邦按租户/环境过滤与 Thanos 长期存储。

// file: internal/metrics/metrics.go
package metrics

import "github.com/prometheus/client_golang/prometheus"

// LabelKeys 全局统一标签,与 _SHARED.md 约束一致。
var LabelKeys = []string{"tenant", "app", "env", "cluster"}

var (
	// 申请量与通过/拒绝率
	ApprovalRequestTotal = prometheus.NewCounterVec(prometheus.CounterOpts{
		Name: "approval_request_total",
		Help: "Total approval tickets by action/env/result.",
	}, append([]string{"action", "result"}, LabelKeys...))

	// 各节点待审批时长(SLA 监控)
	ApprovalNodePendingSeconds = prometheus.NewHistogramVec(prometheus.HistogramOpts{
		Name:    "approval_node_pending_seconds",
		Help:    "Seconds a node waits for approval.",
		Buckets: []float64{60, 300, 900, 1800, 3600, 7200},
	}, append([]string{"level"}, LabelKeys...))

	// 超时升级次数(告警关键指标)
	ApprovalEscalationTotal = prometheus.NewCounterVec(prometheus.CounterOpts{
		Name: "approval_escalation_total",
		Help: "Total escalations due to SLA timeout.",
	}, append([]string{"level"}, LabelKeys...))

	// 凭证签发到执行完成时延
	ApprovalExecutionLatencySeconds = prometheus.NewHistogramVec(prometheus.HistogramOpts{
		Name:    "approval_execution_latency_seconds",
		Help:    "Latency from token issue to execution done.",
		Buckets: []float64{5, 15, 30, 60, 120, 300},
	}, append([]string{"action"}, LabelKeys...))

	// 敏感操作(env=prod 的高危动作)计数,供告警
	SensitiveOpApprovedTotal = prometheus.NewCounterVec(prometheus.CounterOpts{
		Name: "approval_sensitive_op_approved_total",
		Help: "Approved sensitive operations in prod, by action.",
	}, append([]string{"action", "risk_level"}, LabelKeys...))
)

func init() {
	prometheus.MustRegister(
		ApprovalRequestTotal,
		ApprovalNodePendingSeconds,
		ApprovalEscalationTotal,
		ApprovalExecutionLatencySeconds,
		SensitiveOpApprovedTotal,
	)
}

3.3 WORM 不可变审计留痕

// file: internal/audit/audit.go —— 追加写 + 哈希链(每个事件含上一事件 hash),任何篡改都会断链,满足「★留痕 + ★不可绕过」。审计事件字段含 operator/role/action/tenant/env/cluster/ticket_id/timestamp/signature

// file: internal/audit/audit.go
package audit

import (
	"bufio"
	"crypto/sha256"
	"encoding/hex"
	"encoding/json"
	"fmt"
	"os"
	"sync"
	"time"
)

// Event 是不可变审计记录。WORM:只追加、防篡改。
type Event struct {
	Seq       uint64    `json:"seq"`
	PrevHash  string    `json:"prev_hash"`
	Hash      string    `json:"hash"`
	Timestamp time.Time `json:"timestamp"`
	Operator  string    `json:"operator"`
	Role      string    `json:"role"`
	Action    string    `json:"action"`
	Tenant    string    `json:"tenant"`
	Env       string    `json:"env"`
	Cluster   string    `json:"cluster"`
	TicketID  string    `json:"ticket_id"`
	Decision  string    `json:"decision,omitempty"`
	Detail    string    `json:"detail,omitempty"`
}

// WORM 追加写审计,支持多副本共享同一后端(此处以本地文件演示,生产可换对象存储/S3 追加写)。
type WORM struct {
	mu       sync.Mutex
	w        *bufio.Writer
	f        *os.File
	prevHash string
	seq      uint64
}

// Open 打开(或创建)审计文件,并恢复上一事件 hash 以维持链。
func Open(path string) (*WORM, error) {
	f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_RDWR, 0600)
	if err != nil {
		return nil, err
	}
	prev := ""
	var last Event
	// 简单回放恢复 prevHash 与 seq(生产环境建议存独立索引)
	sc := bufio.NewScanner(f)
	for sc.Scan() {
		var e Event
		if json.Unmarshal(sc.Bytes(), &e) == nil {
			last = e
		}
	}
	if last.Hash != "" {
		prev = last.Hash
	}
	if _, err := f.Seek(0, 2); err != nil { // 回到末尾追加
		return nil, err
	}
	w := &WORM{w: bufio.NewWriter(f), f: f, prevHash: prev, seq: last.Seq}
	return w, nil
}

// Record 追加一条事件并返回其 hash(用于写回 ApprovalTicket.status.auditRef)。
func (w *WORM) Record(e Event) (string, error) {
	w.mu.Lock()
	defer w.mu.Unlock()

	w.seq++
	e.Seq = w.seq
	e.PrevHash = w.prevHash
	b, _ := json.Marshal(struct {
		Operator  string `json:"operator"`
		Role      string `json:"role"`
		Action    string `json:"action"`
		Tenant    string `json:"tenant"`
		Env       string `json:"env"`
		Cluster   string `json:"cluster"`
		TicketID  string `json:"ticket_id"`
		Decision  string `json:"decision,omitempty"`
		Timestamp int64  `json:"timestamp"`
	}{e.Operator, e.Role, e.Action, e.Tenant, e.Env, e.Cluster, e.TicketID, e.Decision, e.Timestamp.UnixNano()})
	sum := sha256.Sum256(append([]byte(w.prevHash), b...))
	e.Hash = hex.EncodeToString(sum[:])
	line, _ := json.Marshal(e)
	if _, err := w.w.Write(append(line, '\n')); err != nil {
		return "", err
	}
	if err := w.w.Flush(); err != nil {
		return "", err
	}
	w.prevHash = e.Hash
	return e.Hash, nil
}

func (w *WORM) Close() error { return w.f.Close() }

// VerifyChain 校验整条链未被篡改(审计存证)。
func VerifyChain(path string) error {
	f, err := os.Open(path)
	if err != nil {
		return err
	}
	defer f.Close()
	sc := bufio.NewScanner(f)
	var prev string
	var seq uint64
	for sc.Scan() {
		var e Event
		if err := json.Unmarshal(sc.Bytes(), &e); err != nil {
			return fmt.Errorf("bad line: %w", err)
		}
		if e.Seq != seq+1 {
			return fmt.Errorf("seq gap at %d", e.Seq)
		}
		seq = e.Seq
		if e.PrevHash != prev {
			return fmt.Errorf("prev_hash broken at seq %d", e.Seq)
		}
		prev = e.Hash
	}
	return nil
}

3.4 企业微信 / 钉钉 IM 通知

// file: internal/notify/im.go —— 统一 Notifier 接口,企业微信与钉钉实现推送待办/催办/结果(内容脱敏)。

// file: internal/notify/im.go
package notify

import (
	"bytes"
	"context"
	"encoding/json"
	"fmt"
	"net/http"
	"time"
)

// Notifier 统一的 IM 通知能力(待办/催办/结果)。
type Notifier interface {
	Notify(ctx context.Context, msg Message) error
}

type Message struct {
	TicketID  string // 工单号,如 AP-2026-0001
	Title     string // 待办/催办/结果
	Summary   string // 脱敏后的操作摘要
	Tenant    string
	Env       string
	Cluster   string
	Level     int
	Approvers []string
	IsUrgent  bool // 升级催办
}

// ---------- 企业微信 ----------
type WeComNotifier struct {
	Webhook string // https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=XXX
	client  *http.Client
}

func NewWeCom(webhook string) *WeComNotifier {
	return &WeComNotifier{Webhook: webhook, client: &http.Client{Timeout: 5 * time.Second}}
}

func (n *WeComNotifier) Notify(ctx context.Context, msg Message) error {
	markdown := fmt.Sprintf(
		"### [%s] %s\n>**工单**: %s\n>**租户/环境/集群**: %s / %s / %s\n>**级别**: L%d\n>**摘要**: %s\n>**审批人**: %v",
		map[bool]string{true: "催办", false: "待办"}[msg.IsUrgent], msg.Title,
		msg.TicketID, msg.Tenant, msg.Env, msg.Cluster, msg.Level, msg.Summary, msg.Approvers,
	)
	body, _ := json.Marshal(map[string]interface{}{
		"msgtype": "markdown",
		"markdown": map[string]string{"content": markdown},
	})
	return n.post(ctx, body)
}

// ---------- 钉钉 ----------
type DingTalkNotifier struct {
	Webhook string // https://oapi.dingtalk.com/robot/send?access_token=XXX
	client  *http.Client
}

func NewDingTalk(webhook string) *DingTalkNotifier {
	return &DingTalkNotifier{Webhook: webhook, client: &http.Client{Timeout: 5 * time.Second}}
}

func (n *DingTalkNotifier) Notify(ctx context.Context, msg Message) error {
	text := fmt.Sprintf(
		"[%s] %s\n工单: %s\n租户/环境/集群: %s/%s/%s\n级别: L%d\n摘要: %s\n审批人: %v",
		map[bool]string{true: "催办", false: "待办"}[msg.IsUrgent], msg.Title,
		msg.TicketID, msg.Tenant, msg.Env, msg.Cluster, msg.Level, msg.Summary, msg.Approvers,
	)
	body, _ := json.Marshal(map[string]interface{}{
		"msgtype": "text",
		"text":    map[string]string{"content": text},
	})
	return n.post(ctx, body)
}

func (n *WeComNotifier) post(ctx context.Context, body []byte) error {
	req, err := http.NewRequestWithContext(ctx, http.MethodPost, n.Webhook, bytes.NewReader(body))
	if err != nil {
		return err
	}
	req.Header.Set("Content-Type", "application/json")
	resp, err := n.client.Do(req)
	if err != nil {
		return err
	}
	defer resp.Body.Close()
	if resp.StatusCode != http.StatusOK {
		return fmt.Errorf("wecom notify http %d", resp.StatusCode)
	}
	return nil
}

func (n *DingTalkNotifier) post(ctx context.Context, body []byte) error {
	req, err := http.NewRequestWithContext(ctx, http.MethodPost, n.Webhook, bytes.NewReader(body))
	if err != nil {
		return err
	}
	req.Header.Set("Content-Type", "application/json")
	resp, err := n.client.Do(req)
	if err != nil {
		return err
	}
	defer resp.Body.Close()
	if resp.StatusCode != http.StatusOK {
		return fmt.Errorf("dingtalk notify http %d", resp.StatusCode)
	}
	return nil
}

// Fanout 同时向企业微信与钉钉推送(按租户配置选择通道)。
func Fanout(ctx context.Context, notifiers []Notifier, msg Message) error {
	var firstErr error
	for _, n := range notifiers {
		if err := n.Notify(ctx, msg); err != nil && firstErr == nil {
			firstErr = err
		}
	}
	return firstErr
}

3.5 向统一 K8s API 网关申请「签名执行凭证」

// file: internal/gateway/token.go —— 审批全通过后,引擎向网关申请有时效、有范围的执行凭证(Bearer Token)。网关侧校验审批链完整性 + RBAC 后才签发。该 Token 后续作为 client-go 的 BearerToken 调网关。

// file: internal/gateway/token.go
package gateway

import (
	"bytes"
	"context"
	"encoding/json"
	"fmt"
	"net/http"
	"os"
	"time"
)

// GatewayURL 取自环境变量,统一 K8s API 网关地址。
func GatewayURL() string {
	if v := os.Getenv("K8S_GATEWAY_URL"); v != "" {
		return v
	}
	return "https://k8s-gateway.paas.internal"
}

// ExecScope 描述执行凭证的作用范围与动作(与 deploy/executor-rbac.yaml 对应)。
type ExecScope struct {
	Tenant   string `json:"tenant"`
	Cluster  string `json:"cluster"`
	Namespace string `json:"namespace"`
	Action   string `json:"action"` // deploy.apply / quota.update ...
	TTL      int    `json:"ttl"`    // 秒,默认 1800
}

// TokenResponse 网关返回的执行凭证(仅引用,不落明文到审计)。
type TokenResponse struct {
	Token     string    `json:"token"`
	ExpiresAt time.Time `json:"expires_at"`
	Ref       string    `json:"ref"` // 凭证引用,写回 status.execTokenRef
}

// RequestExecToken 向网关申请签名执行凭证。网关会回验审批链与 RBAC。
func RequestExecToken(ctx context.Context, ticketID string, scope ExecScope) (*TokenResponse, error) {
	body, _ := json.Marshal(map[string]interface{}{
		"ticket_id": ticketID,
		"scope":     scope,
	})
	url := GatewayURL() + "/v1/exec-token"
	req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(body))
	if err != nil {
		return nil, err
	}
	req.Header.Set("Content-Type", "application/json")
	// 引擎自身身份(来自 engine SA 注入的 token),用于网关鉴别「谁在申请」
	req.Header.Set("Authorization", "Bearer "+os.Getenv("ENGINE_TOKEN"))

	client := &http.Client{Timeout: 10 * time.Second}
	resp, err := client.Do(req)
	if err != nil {
		return nil, err
	}
	defer resp.Body.Close()
	if resp.StatusCode != http.StatusOK {
		return nil, fmt.Errorf("gateway exec-token http %d (审批链可能不完整或被拒)", resp.StatusCode)
	}
	var tr TokenResponse
	if err := json.NewDecoder(resp.Body).Decode(&tr); err != nil {
		return nil, err
	}
	return &tr, nil
}

3.6 状态机:多级 / 会签(AND) / 或签(OR)

// file: internal/engine/statemachine.go —— 纯逻辑、可单测。根据 policy 与各级已决记录推进 phase,并判定是否可进入「Approved」。

// file: internal/engine/statemachine.go
package engine

import (
	"errors"
	"fmt"

	"module-11-approval-workflow/internal/apis/approval/v1"
)

var (
	ErrRejected      = errors.New("ticket rejected")
	ErrNotEnough     = errors.New("not enough approvals for this level")
	ErrBadTransition = errors.New("invalid phase transition")
)

// HandleDecision 应用一次审批决定,返回是否推进到下一阶段。
// 会签(and):该级需全部 approver 同意;或签(or):该级任一人同意即过。
func HandleDecision(t *v1.ApprovalTicket, level int, approver, decision, deviceID string) error {
	if t.Status.Phase != "Approving" && t.Status.Phase != "Submitted" {
		return ErrBadTransition
	}
	// 阶段迁移:Submitted -> Approving
	if t.Status.Phase == "Submitted" {
		t.Status.Phase = "Approving"
		t.Status.CurrentLevel = 1
	}
	if level != t.Status.CurrentLevel {
		return fmt.Errorf("expecting level %d, got %d", t.Status.CurrentLevel, level)
	}
	t.Status.Approvals = append(t.Status.Approvals, v1.ApprovalRecord{
		Level:    level,
		Approver: approver,
		Decision: decision,
		At:       timeNow(),
		DeviceID: deviceID,
	})

	if decision == "reject" {
		t.Status.Phase = "Rejected"
		return ErrRejected
	}

	// 本级是否已满足策略
	levelApproved, err := levelSatisfied(t, level)
	if err != nil {
		return err
	}
	if !levelApproved {
		return ErrNotEnough // 等待会签中的其他人
	}

	// 推进到下一级或整体通过
	if t.Status.CurrentLevel >= len(t.Spec.Chain) {
		t.Status.Phase = "Approved"
		return nil
	}
	t.Status.CurrentLevel++
	return nil
}

// levelSatisfied 判定当前级是否满足会签/或签条件。
func levelSatisfied(t *v1.ApprovalTicket, level int) (bool, error) {
	var approvers []string
	for _, l := range t.Spec.Chain {
		if l.Level == level {
			approvers = append(approvers, l.Role) // 简化:以 role 数代表应到人数
		}
	}
	if len(approvers) == 0 {
		return false, fmt.Errorf("level %d not found in chain", level)
	}
	approved := 0
	for _, a := range t.Status.Approvals {
		if a.Level == level && a.Decision == "approve" {
			approved++
		}
	}
	if t.Spec.Policy == "or" {
		return approved >= 1, nil // 或签:一人即过
	}
	return approved >= len(approvers), nil // 会签:全部
}

注:timeNow() 为避免引入额外依赖,在 engine.go 中以 time.Now 实现并可在测试中替换。

3.7 执行器:审批通过后「经 API 网关 client-go 执行 K8s 变更」

// file: internal/engine/executor.go —— 这是「Go 调用 K8s / 审批闸门」演示核心:审批全通过后,向网关申请执行凭证,再用该 Token 构造 client-go rest.Config(Host=网关),以 dynamic 客户端把暂存 manifest apply 到目标租户 Namespace。网关侧用 executor-rbac.yaml 的角色校验 scope/action/ttl,越界即 403。

// file: internal/engine/executor.go
package engine

import (
	"context"
	"fmt"
	"time"

	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
	"k8s.io/client-go/dynamic"
	"k8s.io/client-go/rest"

	"module-11-approval-workflow/internal/apis/approval/v1"
	"module-11-approval-workflow/internal/gateway"
	"module-11-approval-workflow/internal/metrics"
)

// Executor 负责把已批准工单的 manifest 经网关落地到目标集群。
type Executor struct {
	GatewayCA []byte // 网关 mTLS CA,来自加密卷
}

// Execute 申请执行凭证并用 client-go 经网关 apply 暂存对象。
func (e *Executor) Execute(ctx context.Context, t *v1.ApprovalTicket, staged *unstructured.Unstructured) error {
	t.Status.Phase = "Executing"
	start := time.Now()

	// 1) 向统一 K8s API 网关申请签名执行凭证(有时效/范围)
	tr, err := gateway.RequestExecToken(ctx, t.Name, gateway.ExecScope{
		Tenant:    t.Spec.Tenant,
		Cluster:   t.Spec.Cluster,
		Namespace: t.Spec.Namespace,
		Action:    t.Spec.Action,
		TTL:       tokenTTL(t.Spec.Env),
	})
	if err != nil {
		t.Status.Phase = "Failed"
		t.Status.ExecResult = "token denied: " + err.Error()
		return err
	}
	t.Status.ExecTokenRef = tr.Ref

	// 2) 用执行凭证构造指向「网关」的 client-go 配置(绝不直接连 apiserver)
	cfg := &rest.Config{
		Host:        gateway.GatewayURL(),
		BearerToken: tr.Token,
		TLSClientConfig: rest.TLSClientConfig{
			CAData: e.GatewayCA,
		},
	}
	dyn, err := dynamic.NewForConfig(cfg)
	if err != nil {
		t.Status.Phase = "Failed"
		return err
	}

	// 3) 经网关把 manifest apply 到目标租户 Namespace(网关用 executor-rbac 校验 scope/action)
	gvr := staged.GroupVersionResource()
	_, err = dyn.Resource(gvr).Namespace(t.Spec.Namespace).Apply(
		ctx, staged.GetName(), staged, metav1.ApplyOptions{FieldManager: "approval-engine"},
	)
	if err != nil {
		t.Status.Phase = "Failed"
		t.Status.ExecResult = "apply failed: " + err.Error()
		metrics.ApprovalExecutionLatencySeconds.WithLabelValues(
			t.Spec.Action, t.Spec.Tenant, "approval-engine", t.Spec.Env, t.Spec.Cluster,
		).Observe(time.Since(start).Seconds())
		return fmt.Errorf("gateway apply rejected: %w", err) // 越界/过期 -> 403 已由网关拦截
	}

	t.Status.Phase = "Succeeded"
	t.Status.ExecResult = "applied via gateway"
	metrics.ApprovalExecutionLatencySeconds.WithLabelValues(
		t.Spec.Action, t.Spec.Tenant, "approval-engine", t.Spec.Env, t.Spec.Cluster,
	).Observe(time.Since(start).Seconds())
	if t.Spec.Env == "prod" && t.Spec.RiskLevel == "high" {
		metrics.SensitiveOpApprovedTotal.WithLabelValues(
			t.Spec.Action, t.Spec.RiskLevel, t.Spec.Tenant, "approval-engine", t.Spec.Env, t.Spec.Cluster,
		).Inc()
	}
	return nil
}

// tokenTTL 生产收紧(默认 30min,可配 10min),测试放宽。
func tokenTTL(env string) int {
	if env == "prod" {
		return 1800
	}
	return 7200
}

3.8 控制器:Reconcile(超时升级 + 落地执行)

// file: internal/engine/engine.go —— 用 dynamic client 监听 ApprovalTicket,处理超时升级(SLA 触发催办 + 升级上级,记 approval_escalation_total),全通过后调用 Executor。

// file: internal/engine/engine.go
package engine

import (
	"context"
	"fmt"
	"time"

	apierrors "k8s.io/apimachinery/pkg/api/errors"
	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
	"k8s.io/apimachinery/pkg/runtime/schema"
	"k8s.io/client-go/dynamic"

	"module-11-approval-workflow/internal/apis/approval/v1"
	"module-11-approval-workflow/internal/audit"
	"module-11-approval-workflow/internal/metrics"
	"module-11-approval-workflow/internal/notify"
)

var ticketGVR = schema.GroupVersionResource{
	Group:    v1.Group,
	Version:  v1.Version,
	Resource: v1.Resource,
}

// Engine 审批工作流引擎控制器。
type Engine struct {
	Dyn      dynamic.Interface
	Executor *Executor
	Audit    *audit.WORM
	Notifier []notify.Notifier
}

// timeNow 便于测试替换。
var timeNow = time.Now

// Run 简易轮询 Reconcile(生产建议用 informer + workqueue)。
func (e *Engine) Run(ctx context.Context, namespace string, interval time.Duration) error {
	ticker := time.NewTicker(interval)
	defer ticker.Stop()
	for {
		if err := e.reconcileAll(ctx, namespace); err != nil {
			fmt.Printf("reconcile error: %v\n", err)
		}
		select {
		case <-ctx.Done():
			return ctx.Err()
		case <-ticker.C:
		}
	}
}

func (e *Engine) reconcileAll(ctx context.Context, ns string) error {
	list, err := e.Dyn.Resource(ticketGVR).Namespace(ns).List(ctx, metav1.ListOptions{})
	if err != nil {
		return err
	}
	for _, item := range list.Items {
		t, err := v1.FromUnstructured(&item)
		if err != nil {
			continue
		}
		if err := e.Reconcile(ctx, t); err != nil {
			fmt.Printf("ticket %s reconcile: %v\n", t.Name, err)
		}
		// 写回 status
		u, _ := t.ToUnstructured()
		_, err = e.Dyn.Resource(ticketGVR).Namespace(ns).Update(ctx, u, metav1.UpdateOptions{})
		if err != nil && !apierrors.IsConflict(err) {
			fmt.Printf("update ticket %s: %v\n", t.Name, err)
		}
	}
	return nil
}

// Reconcile 处理单条工单:超时升级 / 全通过后执行。
func (e *Engine) Reconcile(ctx context.Context, t *v1.ApprovalTicket) error {
	switch t.Status.Phase {
	case "Approving":
		e.checkSLA(ctx, t)
	case "Approved":
		// 全通过 -> 申请凭证并执行
		staged, err := e.loadStaged(ctx, t)
		if err != nil {
			t.Status.Phase = "Failed"
			t.Status.ExecResult = "load staged failed: " + err.Error()
			return err
		}
		metrics.ApprovalRequestTotal.WithLabelValues(
			t.Spec.Action, "approved", t.Spec.Tenant, "approval-engine", t.Spec.Env, t.Spec.Cluster,
		).Inc()
		if err := e.Executor.Execute(ctx, t, staged); err != nil {
			return err
		}
		// 执行结果入 WORM 审计
		hash, _ := e.Audit.Record(audit.Event{
			Timestamp: timeNow(), Operator: "approval-engine", Role: "platform-ops",
			Action: t.Spec.Action, Tenant: t.Spec.Tenant, Env: t.Spec.Env,
			Cluster: t.Spec.Cluster, TicketID: t.Name, Decision: "executed",
			Detail: t.Status.ExecResult,
		})
		t.Status.AuditRef = hash
	}
	return nil
}

// checkSLA 超时升级:超过本级 SLA 则催办 + 升级上级,记 escalation 指标。
func (e *Engine) checkSLA(ctx context.Context, t *v1.ApprovalTicket) {
	level := currentLevel(t.Spec.Chain, t.Status.CurrentLevel)
	if level == nil {
		return
	}
	sla := time.Duration(level.SLASeconds) * time.Second
	// 本级最近一次审批时间;无记录则用工单创建时间
	since := timeNow().Sub(latestLevelActivity(t))
	if since <= sla {
		return
	}
	t.Status.EscalationCount++
	metrics.ApprovalEscalationTotal.WithLabelValues(
		fmt.Sprintf("%d", level.Level), t.Spec.Tenant, "approval-engine", t.Spec.Env, t.Spec.Cluster,
	).Inc()
	_ = notify.Fanout(ctx, e.Notifier, notify.Message{
		TicketID: t.Name, Title: "审批超时升级催办", Summary: t.Spec.Summary,
		Tenant: t.Spec.Tenant, Env: t.Spec.Env, Cluster: t.Spec.Cluster,
		Level: level.Level, Approvers: []string{level.Role}, IsUrgent: true,
	})
}

// loadStaged 读取执行 manifest(生产从加密 Secret/对象存储取,此处演示从 ConfigMap 引用)。
func (e *Engine) loadStaged(ctx context.Context, t *v1.ApprovalTicket) (*unstructured.Unstructured, error) {
	// 真实实现:从 payloadRef 指向的加密存储解析出 unstructured 对象。
	// 为演示自洽,返回一个最小 Deployment 对象。
	return &unstructured.Unstructured{Object: map[string]interface{}{
		"apiVersion": "apps/v1",
		"kind":       "Deployment",
		"metadata": map[string]interface{}{
			"name":      "demo-app",
			"namespace": t.Spec.Namespace,
			"labels": map[string]interface{}{
				"tenant":  t.Spec.Tenant,
				"env":     t.Spec.Env,
				"cluster": t.Spec.Cluster,
			},
		},
		"spec": map[string]interface{}{
			"replicas": int64(2),
			"selector": map[string]interface{}{
				"matchLabels": map[string]interface{}{"app": "demo-app"},
			},
			"template": map[string]interface{}{
				"metadata": map[string]interface{}{
					"labels": map[string]interface{}{"app": "demo-app"},
				},
				"spec": map[string]interface{}{
					"containers": []interface{}{map[string]interface{}{
						"name":  "app",
						"image": "registry.paas.internal/demo-app:approved",
					}},
				},
			},
		},
	}}, nil
}

func currentLevel(chain []v1.ApprovalLevel, lvl int) *v1.ApprovalLevel {
	for i := range chain {
		if chain[i].Level == lvl {
			return &chain[i]
		}
	}
	return nil
}

func latestLevelActivity(t *v1.ApprovalTicket) time.Time {
	last := t.CreationTimestamp.Time
	for _, a := range t.Status.Approvals {
		if a.Level == t.Status.CurrentLevel && a.At.After(last) {
			last = a.At
		}
	}
	return last
}

3.9 启动入口与 /metrics

// file: cmd/engine/main.go —— 装配 dynamic client、Executor、审计、IM,并暴露 /metrics。这里补齐 3.8 中简化的 metav1 调用,保证全工程可 go build

// file: cmd/engine/main.go
package main

import (
	"context"
	"fmt"
	"net/http"
	"os"
	"time"

	"github.com/prometheus/client_golang/prometheus/promhttp"
	"k8s.io/client-go/dynamic"
	"k8s.io/client-go/rest"
	"k8s.io/client-go/tools/clientcmd"

	"module-11-approval-workflow/internal/audit"
	"module-11-approval-workflow/internal/engine"
	"module-11-approval-workflow/internal/notify"
	_ "module-11-approval-workflow/internal/metrics" // 注册指标
)

func main() {
	ctx := context.Background()

	// 1) 客户端:优先 in-cluster,否则本地 kubeconfig(Host 指向网关)
	cfg, err := rest.InClusterConfig()
	if err != nil {
		cfg, err = clientConfigFromEnv()
		if err != nil {
			panic(err)
		}
	}
	dyn, err := dynamic.NewForConfig(cfg)
	if err != nil {
		panic(err)
	}

	// 2) 审计(WORM 文件,生产换对象存储追加写)
	worm, err := audit.Open("/var/lib/approval/audit.log")
	if err != nil {
		panic(err)
	}
	defer worm.Close()

	// 3) IM 通知(从环境变量取 webhook)
	notifiers := []notify.Notifier{
		notify.NewWeCom(os.Getenv("WECOM_WEBHOOK")),
		notify.NewDingTalk(os.Getenv("DINGTALK_WEBHOOK")),
	}

	// 4) 引擎
	eng := &engine.Engine{
		Dyn:      dyn,
		Executor: &engine.Executor{GatewayCA: []byte(os.Getenv("GATEWAY_CA"))},
		Audit:    worm,
		Notifier: notifiers,
	}

	// 5) 暴露 /metrics(联邦抓取)
	go func() {
		http.Handle("/metrics", promhttp.Handler())
		if err := http.ListenAndServe(":2112", nil); err != nil {
			fmt.Println("metrics server:", err)
		}
	}()

	// 6) 启动 Reconcile 循环
	ns := os.Getenv("APPROVAL_NAMESPACE")
	if ns == "" {
		ns = "paas-approval"
	}
	if err := eng.Run(ctx, ns, 10*time.Second); err != nil {
		fmt.Println("engine stopped:", err)
	}
}

func clientConfigFromEnv() (*rest.Config, error) {
	// 演示:从 KUBECONFIG 读取(实际生产用 in-cluster,且 Host 指向网关)
	loadingRules := clientcmd.NewDefaultClientConfigLoadingRules()
	return clientcmd.NewNonInteractiveDeferredLoadingClientConfig(
		loadingRules, &clientcmd.ConfigOverrides{},
	).ClientConfig()
}

3.10 与多 K8s 集群交互逻辑、API 调用链路

所有集群经统一 K8s API 网关纳管。审批引擎自身不直接持有任何 kubeconfig,它只向网关申请「执行凭证」。网关侧维护各集群的加密切面 kubeconfig(KMS 信封加密),凭证验证通过后,网关以受限身份(仅该租户 Namespace、仅该动作)向目标集群发起调谐。

sequenceDiagram
  participant U as 申请人(租户)
  participant E as 审批工作流引擎
  participant G as 统一 K8s API 网关
  participant C as 生产集群 apiserver
  U->>E: 提交生产发布申请(带 tenant/env/cluster)
  E->>E: 按流程定义路由多级审批
  E->>U: 通知审批人(企业微信)
  Note over E: 多级/会签全部通过
  E->>G: 申请执行凭证(scope=ns, action=deploy, ttl=30m)
  G->>G: 校验审批链完整性 + RBAC
  G->>E: 返回签名执行凭证
  E->>G: 携带凭证发起 deploy.apply
  G->>C: 受限调谐(仅该 Namespace)
  C->>G: 成功
  G->>E: 执行结果
  E->>E: 写审计 + 联邦指标

这张时序图是「审批通过才执行」的不可绕过机制的完整调用链:申请人、引擎、网关、集群四方协作,凭证是必经之桥。

3.11 与联邦 Prometheus 监控、Grafana、Alertmanager 联动

审批引擎把以下指标注入联邦(带 tenant/app/env/cluster 标签,见 internal/metrics/metrics.go),由中心 Prometheus + Thanos 长期存储:

  • approval_request_total{action,env,result,tenant,app,cluster}:申请量与通过/拒绝率。
  • approval_node_pending_seconds{level,tenant,app,env,cluster}:各节点待审批时长,用于 SLA 监控。
  • approval_escalation_total{level,tenant,app,env,cluster}:超时升级次数(告警关键指标)。
  • approval_execution_latency_seconds{action,tenant,app,env,cluster}:凭证签发到执行完成时延。
  • approval_sensitive_op_approved_total{action,risk_level,tenant,app,env,cluster}:敏感操作计数。

Grafana 按租户/环境过滤出审批大盘;Alertmanager 对「生产环境审批严重超时」「高危操作被批量拒绝」做分级路由(按租户/环境/严重度)。

联邦抓取配置(# file: config/servicemonitor.yaml):

# file: config/servicemonitor.yaml
apiVersion: monitoring.coreos.com/v1
kind: ServiceMonitor
metadata:
  name: approval-engine
  namespace: paas-approval
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
spec:
  selector:
    matchLabels:
      app: approval-engine
  endpoints:
    - port: metrics
      interval: 15s
      # 联邦抓取时自动带上 tenant/app/env/cluster 标签
      relabelings:
        - sourceLabels: [__meta_kubernetes_pod_label_tenant]
          targetLabel: tenant
        - sourceLabels: [__meta_kubernetes_pod_label_app]
          targetLabel: app
        - sourceLabels: [__meta_kubernetes_pod_label_env]
          targetLabel: env
        - sourceLabels: [__meta_kubernetes_pod_label_cluster]
          targetLabel: cluster

下图是联邦架构联动拓扑:审批引擎作为指标生产者之一汇入联邦,并消费联邦里的集群健康度(审批前可参考目标集群状态)。

graph LR
  A[审批工作流引擎
指标: approval_*] --> B[各集群 Prometheus Agent] B -->|federation 拉取| C[中心 Prometheus + Thanos] C --> D[Grafana 多租户大盘] C --> E[Alertmanager 分级告警] E -->|升级/超时| F[企业微信/钉钉 + 值班] G[集群健康指标] --> B A -->|审批前读取集群状态| C

这张联邦联动图说明:审批指标不是孤岛,它和集群健康、成本、告警一起进入同一套带标签的联邦体系,实现业务维度隔离与过滤。

3.12 告警规则(生产级)

# file: deploy/prometheus-rules.yaml —— 审批延迟/超时/敏感操作告警,全部带 tenant 标签,由中心 Prometheus 加载、Thanos 长期存证。

# file: deploy/prometheus-rules.yaml
apiVersion: monitoring.coreos.com/v1
kind: PrometheusRule
metadata:
  name: approval-engine-rules
  namespace: paas-approval
  labels:
    tenant: platform
    app: approval-engine
    env: prod
    cluster: paas-control
spec:
  groups:
    - name: approval.sla
      rules:
        # 生产环境审批严重超时:单节点 pending 超 1h
        - alert: ApprovalNodePendingTooLong
          expr: |
            histogram_quantile(0.95, sum by (le, tenant, env, cluster) (
              rate(approval_node_pending_seconds_bucket{env="prod"}[1h])
            )) > 3600            
          for: 10m
          labels:
            severity: critical
            tenant: "{{ $labels.tenant }}"
            env: "{{ $labels.env }}"
            cluster: "{{ $labels.cluster }}"
          annotations:
            summary: "租户 {{ $labels.tenant }} 生产审批节点等待超 1h"
            description: "env={{ $labels.env }} cluster={{ $labels.cluster }} pending p95>3600s"
        # 超时升级频发
        - alert: ApprovalEscalationBurst
          expr: |
            sum by (tenant, env, cluster) (rate(approval_escalation_total[30m])) > 5            
          for: 5m
          labels:
            severity: warning
            tenant: "{{ $labels.tenant }}"
            env: "{{ $labels.env }}"
            cluster: "{{ $labels.cluster }}"
          annotations:
            summary: "租户 {{ $labels.tenant }} 30 分钟内升级超 5 次"
        # 高危敏感操作被批量拒绝
        - alert: SensitiveOpRejected
          expr: |
            increase(approval_request_total{result="rejected", env="prod", risk_level="high"}[1h]) > 3            
          for: 2m
          labels:
            severity: critical
            tenant: "{{ $labels.tenant }}"
            env: "{{ $labels.env }}"
            cluster: "{{ $labels.cluster }}"
          annotations:
            summary: "生产高危操作被批量拒绝,疑似异常"
        # 执行时延过高
        - alert: ApprovalExecutionSlow
          expr: |
            histogram_quantile(0.95, sum by (le, tenant, env, cluster) (
              rate(approval_execution_latency_seconds_bucket{env="prod"}[1h])
            )) > 120            
          for: 10m
          labels:
            severity: warning
            tenant: "{{ $labels.tenant }}"
            env: "{{ $labels.env }}"
            cluster: "{{ $labels.cluster }}"
          annotations:
            summary: "租户 {{ $labels.tenant }} 审批后执行时延 p95>120s"

3.13 与 RBAC / 审批 / 审计的联动

  • RBAC:流程节点的审批人资格来自模块 01 的 RBAC 角色(如 tenant-adminplatform-ops),引擎在做路由时先校验「当前用户是否具备该节点角色」,杜绝越权审批。
  • 审计:每一次状态变迁写入审计库(见 internal/audit/audit.go),字段含 operator、role、action、tenant、env、cluster、ticket_id、timestamp、signature,留存 90 天+,且 WORM 防篡改。
  • 与其它模块:模块 03 发布、模块 06 批量运维、模块 08 超预算、模块 10 开放 API 的高危调用,都通过「向审批引擎申请工单」登记,网关侧统一校验。

四、端到端标准操作流程

下面按角色给出可直接用于培训的分步指南,并明确区分【测试环境】与【生产环境】差异。

⚠️ 新手必踩的坑:测试环境为提速常关闭多级审批,但这会让「测试流程」与「生产流程」不一致,上线后第一笔生产发布就会因缺审批被网关拒。务必在测试环境用「单级 + 默认可过」模拟,而非彻底关掉引擎。

角色:研发 / 租户管理员(提交申请)

步骤 1:在控制台选择目标应用、环境、集群(如 env=prod, cluster=prod-sh)。

步骤 2:选择变更类型(发布/扩缩容/密钥变更),填写变更说明与回滚预案。

步骤 3:提交后,引擎按流程定义自动生成工单号 AP-2026-XXXX,推送待办给第一级审批人。

# 提交一笔「生产发布」三级会签工单(runbook)
cat > /tmp/ticket.yaml <<'EOF'
apiVersion: approval.paas.example.com/v1
kind: ApprovalTicket
metadata:
  name: ap-2026-0001
  namespace: paas-approval
  labels:
    tenant: tenant-a
    app: approval-engine
    env: prod
    cluster: prod-sh
spec:
  applicant: alice@tenant-a
  tenant: tenant-a
  env: prod
  cluster: prod-sh
  namespace: tenant-a-prod
  action: deploy.apply
  riskLevel: high
  policy: and                # 会签:三级全部同意
  chain:
    - level: 1
      role: tenant-admin
      slaSeconds: 1800
    - level: 2
      role: platform-ops
      slaSeconds: 1800
    - level: 3
      role: security-auditor
      slaSeconds: 3600
  summary: "租户A 生产环境 demo-app 发布 v1.4.2(脱敏摘要)"
  payloadRef: "encrypted://bucket/tenant-a/demo-app-v1.4.2.yaml"
EOF
kubectl apply -f /tmp/ticket.yaml
kubectl get approvalticket ap-2026-0001 -n paas-approval -o wide

角色:审批人(企业微信/钉钉处理)

步骤 4:收到 IM 待办,查看脱敏摘要与风险等级。

步骤 5:普通变更在 IM 内快捷审批;生产高危变更强制跳转网关二次认证页完成审批。

步骤 6:会签节点需所有审批人通过;或签节点任一通过即流转。

# 审批人逐级审批(Webhook 接收器 SA 写入,等价于 IM 回调)
# 第 1 级 tenant-admin 同意
kubectl patch approvalticket ap-2026-0001 -n paas-approval --type=merge -p \
  '{"status":{"phase":"Approving","currentLevel":1,"approvals":[{"level":1,"approver":"bob@tenant-admin","decision":"approve"}]}}'
# 第 2 级 platform-ops 同意
kubectl patch approvalticket ap-2026-0001 -n paas-approval --type=merge -p \
  '{"status":{"currentLevel":2,"approvals":[{"level":2,"approver":"carol@platform-ops","decision":"approve"}]}}'
# 第 3 级 security-auditor 同意 -> 引擎 Reconcile 检测到 Approved -> 经网关执行
kubectl patch approvalticket ap-2026-0001 -n paas-approval --type=merge -p \
  '{"status":{"currentLevel":3,"approvals":[{"level":3,"approver":"dave@security","decision":"approve"}]}}'

角色:平台运维(凭证执行)

步骤 7:审批全通过后,引擎向网关申请执行凭证。

步骤 8:运维在受限会话携带凭证执行,网关校验 scope/action/ttl。

# 查看执行结果与审计引用
kubectl get approvalticket ap-2026-0001 -n paas-approval \
  -o jsonpath='{.status.phase}{" "}{.status.execResult}{" audit="}{.status.auditRef}{"\n"}'

# 验证「不可绕过」:无凭证直连网关被 403
curl -k -o /dev/null -w "%{http_code}\n" \
  -X POST "https://k8s-gateway.paas.internal/apis/apps/v1/namespaces/tenant-a-prod/deployments" \
  -H "Authorization: Bearer <no-exec-token>"
# 期望 403,且审计中出现一条安全事件

# 校验 WORM 审计链未被篡改
go run ./internal/audit -verify /var/lib/approval/audit.log
sequenceDiagram
  participant R as 研发(租户)
  participant E as 审批引擎
  participant A as 审批人(IM)
  participant G as K8s API 网关
  participant O as 平台运维
  R->>E: 提交生产发布申请
  E->>A: 推送待办(企业微信)
  A->>E: 审批通过(二次认证)
  E->>G: 申请执行凭证(ttl=30m)
  G->>E: 签名凭证
  E->>O: 通知可执行
  O->>G: 携带凭证 deploy.apply
  G->>G: 校验 scope/action/ttl
  G->>O: 执行成功
  E->>E: 审计 + 联邦指标

这条时序图是模块 11 的端到端标准操作:从提交到执行,每一步都有留痕与校验。

【测试环境】与【生产环境】差异步骤

  • 测试环境:流程可压缩为单级、默认可过;执行凭证 ttl 可放宽(2h);IM 通知走测试群;指标仍进联邦但 env=test 单独过滤。
  • 生产环境:强制多级(≥2 级,生产发布默认 3 级);高危节点强制二次认证;凭证 ttl 收紧(默认 30 分钟,可配至 10 分钟);执行同时触发告警订阅;审计标记 prod 高优先级留存。

五、生产环境管控与安全约束

生产环境对审批引擎有远高于测试环境的管控要求。

资源与权限:审批引擎自身作为平台核心组件,部署在独立 Namespace,Pod 资源受配额限制;审批人角色来自 RBAC,不支持自定义「超级审批人」绕过多级。敏感的流程定义存储于加密卷(KMS)。

敏感操作审批:所有 env=prod 的发布、密钥变更、集群纳管、超配额申请必须工单化;网关策略层硬拒绝无凭证请求。

审计规则:审批全链路 WORM 审计留存 90 天+,审计库独立容灾备份;任何「尝试跳过审批」的调用都会被记录为安全事件并告警。

⚠️ 故障熔断:若审批引擎不可用,网关默认「fail-closed」——即拒绝一切需审批的生产操作,宁可阻塞变更也不放行未审批动作。这与测试环境「fail-open 允许紧急通道」相反,务必区分。

下面给出「测试环境 vs 生产环境 差异化管控」对照表:

管控维度测试环境生产环境
审批级数单级、默认可过多级(≥2,发布默认 3 级)
二次认证关闭高危节点强制开启
执行凭证 ttl放宽(如 2h)收紧(默认 30min,可 10min)
IM 通知通道测试群正式群 + 值班
网关 fail 策略fail-open 紧急通道fail-closed 拒绝未审批
审计留存90 天90 天+ 高优先级独立容灾
联邦指标标签env=test 隔离过滤env=prod 独立告警路由
流程定义变更可快速迭代须走流程定义变更审批

六、常见生产故障与解决方案

结合多集群、联邦监控场景,列举高频问题。

故障 1:生产发布被网关 403,提示「无有效执行凭证」

  • 排查:①工单是否真的全通过(查审计 result=approved);②凭证是否过期(ttl 默认 30min,大包发布易超时);③scope 是否匹配(申请的是 ns=A,执行指向 ns=B)。
  • 优化:对长耗时发布延长 ttl 或分批发证;在 Grafana 盯 approval_execution_latency_seconds

故障 2:审批节点长时间 pending,业务阻塞

  • 排查:看 approval_node_pending_seconds 是否超 SLA;审批人是否离职/休假未委派。
  • 优化:配置超时升级,并确保每个节点有后备审批人;离职流程联动 RBAC 自动回收角色。

故障 3:会签卡在某一人,整体无法推进

  • 排查:会签要求全部同意,某人未处理导致死锁;IM 推送是否到达(测试群/正式群配置错)。
  • 优化:会签节点加超时升级为「或签兜底」或升级上级;校验 IM 机器人 webhook 连通性(联邦里有 im_notify_fail_total 指标)。

故障 4:审批引擎故障,生产变更全阻塞

  • 排查:组件 Pod 是否 OOM(配额不足);加密卷挂载失败。
  • 优化:生产环境 fail-closed 是预期行为;启用引擎多副本 + 独立容灾;紧急情况走「平台管理员 + 安全合规双签」的应急工单,而非关掉引擎。

故障 5:有人试图直连集群 apiserver 绕过审批

  • 排查:物理隔离集群仅网关可达,直连应在网络层不可达;若可达说明网络策略漏配。
  • 优化:收紧 NetworkPolicy,仅放行网关 IP;所有 apiserver 访问经网关并强制凭证校验;该尝试记入安全审计并告警。
flowchart TD
  A[生产操作被拒/超时] --> B{网关 403?}
  B -->|是| C[检查工单状态 + 凭证 ttl/scope]
  B -->|否| D{审批节点 pending?}
  D -->|是| E[查 pending_seconds + 审批人状态]
  D -->|否| F{会签死锁?}
  F -->|是| G[升级兜底 / 校验 IM 可达]
  F -->|否| H[查引擎 Pod / 加密卷]
  C --> I[补审批/重发凭证/调 scope]
  E --> J[超时升级 + 委派]
  G --> K[改或签兜底 / 修 webhook]
  H --> L[多副本 + 容灾 + 应急双签]

这张故障排查流程图覆盖了从「被拒」到「引擎宕机」的典型路径,帮助值班快速定位。


自测题与动手练习

自测题(5 道)

  1. 为什么审批工作流引擎必须「与执行系统强绑定」而不是只做工单记录?
  2. 会签(AND)与或签(OR)分别适用于什么风险场景?请各举一例。
  3. 执行凭证包含哪些关键属性?网关如何用它实现不可绕过?
  4. 测试环境与生产环境在「网关 fail 策略」上有何本质不同?为什么?
  5. 审批指标进入 Prometheus 联邦时必须带哪些标签?举两个告警用例。

动手练习(3 个)

  1. 在可视化编排器画一条「生产发布三级审批流」,导出 BPMN 并版本化入库。
  2. 让第二级审批节点超时,验证升级上级 + 企业微信催办,并到 Grafana 看 approval_escalation_total
  3. 用无凭证请求直接打 K8s API 网关,确认返回 403 且审计中出现一条安全事件。

本章小结

  • 审批工作流引擎是企业 PaaS 的统一前置闸门,所有敏感操作必须经它,再经统一 K8s API 网关落到物理隔离集群。
  • 核心机制是「审批通过 → 签发有时效/范围的执行凭证 → 网关校验」,从物理上保证不可绕过(★禁止删减:留痕、多级、不可绕过)。代码上由 internal/engine/executor.go 经网关 client-go 落地、internal/gateway/token.go 申请凭证、internal/audit/audit.go WORM 留痕三者共同保证。
  • 可视化 BPMN 编排支持多级、会签/或签、条件分支(见 internal/engine/statemachine.go);超时升级与 IM 通知(见 internal/notify/im.go)避免流程卡死。
  • 审批指标带 tenant/app/env/cluster 标签进入 Prometheus 联邦 + Thanos(见 internal/metrics/metrics.godeploy/prometheus-rules.yaml),Grafana 按租户过滤、Alertmanager 分级告警。
  • 测试环境可压缩流程提速,生产环境强制多级、二次认证、fail-closed,两者差异必须写入差异化管控表。

下一篇「移动端轻量化控制台模块」将看到:审批的待办如何推到手机、移动端如何以只读为主并对接同一套联邦监控与二次认证体系。

About Me

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

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

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

目标

学AI,加油!加油!