学习目标
本篇讲透企业自研 PaaS 平台的「审批工作流引擎」——它是所有敏感操作的统一闸门。学完本篇,你应该能够:
- 说清楚审批流如何作为其它模块(发布、扩缩容、密钥变更、集群纳管)敏感操作的强制前置闸门。
- 基于可视化流程编排引擎(BPMN 语义)设计多级审批、会签/或签、条件分支的生产级流程。
- 配置超时升级、催办、审批人委派/代理,避免流程卡死造成业务阻塞。
- 把审批事件接入 K8s API 网关,实现「审批通过才执行」的不可绕过机制。
- 将审批操作、IM 通知、告警升级指标接入 Prometheus 联邦,并保证审计留痕 90 天+。
前置知识:
- 已了解 _SHARED.md 全局约束:多套物理隔离 K8s 集群、统一 K8s API 网关、Namespace 租户隔离;Prometheus 联邦 + Thanos 长期存储、指标带
tenant/app/env/cluster标签;RBAC 细粒度、审计 90 天+、敏感配置加密、镜像扫描。 - 了解模块 01 租户与 RBAC、模块 02 集群纳管、模块 03 应用生命周期基本概念。
- 知道 BPMN 基本元素(开始/任务/网关/结束)与 IM 机器人(企业微信/钉钉)Webhook 机制。
本章你会动手做的事:
- 在可视化编排器里画一条「生产环境发布 = 申请人 → 租户管理员 → 平台运维」三级审批流。
- 故意让某一审批节点超时,观察升级给上级并触发企业微信催办。
- 直接绕过审批网关调用 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 并记录。
- 审批通过后,引擎向统一 K8s API 网关申请一份「执行凭证」:含生效范围(tenant/namespace/cluster)、允许动作(如
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-admin、platform-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 道)
- 为什么审批工作流引擎必须「与执行系统强绑定」而不是只做工单记录?
- 会签(AND)与或签(OR)分别适用于什么风险场景?请各举一例。
- 执行凭证包含哪些关键属性?网关如何用它实现不可绕过?
- 测试环境与生产环境在「网关 fail 策略」上有何本质不同?为什么?
- 审批指标进入 Prometheus 联邦时必须带哪些标签?举两个告警用例。
动手练习(3 个)
- 在可视化编排器画一条「生产发布三级审批流」,导出 BPMN 并版本化入库。
- 让第二级审批节点超时,验证升级上级 + 企业微信催办,并到 Grafana 看
approval_escalation_total。 - 用无凭证请求直接打 K8s API 网关,确认返回 403 且审计中出现一条安全事件。
本章小结
- 审批工作流引擎是企业 PaaS 的统一前置闸门,所有敏感操作必须经它,再经统一 K8s API 网关落到物理隔离集群。
- 核心机制是「审批通过 → 签发有时效/范围的执行凭证 → 网关校验」,从物理上保证不可绕过(★禁止删减:留痕、多级、不可绕过)。代码上由
internal/engine/executor.go经网关 client-go 落地、internal/gateway/token.go申请凭证、internal/audit/audit.goWORM 留痕三者共同保证。 - 可视化 BPMN 编排支持多级、会签/或签、条件分支(见
internal/engine/statemachine.go);超时升级与 IM 通知(见internal/notify/im.go)避免流程卡死。 - 审批指标带
tenant/app/env/cluster标签进入 Prometheus 联邦 + Thanos(见internal/metrics/metrics.go与deploy/prometheus-rules.yaml),Grafana 按租户过滤、Alertmanager 分级告警。 - 测试环境可压缩流程提速,生产环境强制多级、二次认证、fail-closed,两者差异必须写入差异化管控表。
下一篇「移动端轻量化控制台模块」将看到:审批的待办如何推到手机、移动端如何以只读为主并对接同一套联邦监控与二次认证体系。