切换主题
设计原则与模式
问题:如何新增能力而不修改整个运行器
一个文件处理器支持缩略图、压缩包清理和水印。如果通用运行器中到处写 if capability == ...,每个新增能力都可能影响会话、恢复和事件上报。真正的变化点是“如何执行具体处理”,通用生命周期应尽量稳定。
模式为重复出现的设计问题提供结构。它的价值在于划清变化边界。图形长得像某种 UML,不能证明需要那个模式。
原理:把原则落实到具体契约
单一职责是单一变化原因,而不是一个方法。开闭原则鼓励在可预见变化点通过扩展实现新增行为,但不要求所有未来情况都提前抽象。
接口隔离要求调用者不依赖无关能力;依赖反转要求稳定策略不依赖具体技术。里氏替换要求实现遵守契约:如果接口允许重试,实现就不能悄悄在每次调用中产生不可去重副作用。
最小例子:Strategy 与委托
go
// 教学片段:共同契约只覆盖可替换的执行部分。
type ImageProcessor interface {
Process(context.Context, Input) (Output, error)
}
type Runner struct { processor ImageProcessor }
func (r *Runner) Run(ctx context.Context, in Input) (Output, error) {
return r.processor.Process(ctx, in)
}水印和缩略图可以在满足相同输入输出语义时实现该契约。如果一个能力返回瞬时文件,另一个建立长期流媒体 Session,就要重新确认共同语义,不能只为统一接口而抹掉生命周期差别。
方案比较:常见结构解决什么问题
| 结构 | 解决的问题 | 成本与限制 |
|---|---|---|
| Adapter | 把数据库或远端协议转换为核心契约 | 需要维护映射与错误分类 |
| Strategy | 同一职责的不同算法或执行实现 | 必须拥有可替换语义 |
| 组合与委托 | 由小职责拼成大入口 | 组合层应避免吸收业务逻辑 |
| 构造函数与显式注入 | 明确有效状态和依赖 | 不等同于 GoF Factory Method |
| 直接函数或分支 | 变化少、规模小 | 分支增长后再识别稳定抽象 |
结构体持有接口不必都叫 Strategy,NewX 也不必都叫 Factory。Observer 的订阅语义与可靠消息的持久化、ACK、重投并不等价;Outbox 是可靠性模式,有独立事务要求。
真实案例:适配与执行扩展
PostgreSQL 与 HTTP 客户端实现核心端口,是 Adapter 的实际应用。Worker 的 CapabilityExecutor 把通用执行命令委托给具体能力,可从策略与委托的角度理解。Runtime 管理会话和恢复,不需要自己实现 ZIP 清理算法。
已核对的实现 · 本地代码快照
Capability 可替换执行契约 backend · c6dbe05c
worker/internal/core/execution/ports/outbound/capability_executor.go · 第 9–27 行
符号:CapabilityExecutor · 核对日期 2026-10-02
来源与提交版本一致
// CapabilityExecutor 将通用执行命令委托给具体 Capability。
type CapabilityExecutor interface {
// Execute 根据命令中的能力名称和版本执行,并通过 reporter 上报进度。
Execute(context.Context, domain.Command, ProgressReporter) (*domain.AttemptEvent, error)
}
// Progress 是能力执行过程中上报给控制面的结构化进度。
type Progress struct {
Stage string `json:"stage"` // 当前处理阶段名称。
Status string `json:"stage_status"` // 阶段状态。
Progress int `json:"progress"` // 整体 0 到 100 的进度。
Completed int `json:"completed"` // 当前阶段已完成条目数。
Total int `json:"total"` // 当前阶段总条目数。
Failed int `json:"failed"` // 当前阶段失败条目数。
}
// ProgressReporter 将能力进度转换为本地检查点和控制面尝试心跳。
type ProgressReporter func(context.Context, Progress) error
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
Worker 的依赖与运行状态 backend · c6dbe05c
worker/internal/core/execution/application/runtime.go · 第 48–100 行
符号:Runtime · 核对日期 2026-10-02
来源与提交版本一致
// Runtime 是 Worker 的长期运行应用:注册会话、恢复本地状态、消费命令并可靠上报事件。
type Runtime struct {
registration Registration
control ControlPlane
commands CommandStream
store LocalStore
capabilityExecutor CapabilityExecutor
ids executionports.IDGenerator
observer executionports.Observation
now func() time.Time
retryInitial time.Duration
outboxWake chan struct{} // 合并事件发送唤醒,可靠数据始终保存在本地数据库。
outboxDrain chan chan error // 请求发送器完成一轮发送并返回结果。
admissionMu sync.Mutex // 串行化本地容量检查与尝试接收。
activeMu sync.Mutex // 保护正在执行的协程及其取消函数。
active map[string]context.CancelFunc // 本进程持有的运行中尝试。
canceled map[string]bool // 记录当前会话已收到的取消。
draining bool // 当前进程的排空状态,持久标记用于重启恢复。
resourceInventory func() executiondomain.NodeInventory
resourceUsage func(context.Context) executiondomain.ResourceUsage
}
func New(registration Registration, control ControlPlane, commands CommandStream, store LocalStore, capabilityExecutor CapabilityExecutor, ids executionports.IDGenerator, observer executionports.Observation) (*Runtime, error) {
if registration.InstanceID == "" || registration.Name == "" || registration.NodeID == "" {
return nil, executiondomain.Wrap(
executiondomain.Invalid,
"Worker 注册信息缺少实例、名称或 Node",
errors.New("worker instance ID, name, and node ID are required"),
)
}
if registration.Inventory.CPU.CapacityMillis < 1 {
return nil, executiondomain.Wrap(
executiondomain.Invalid,
"Worker CPU 资源采集不可用",
errors.New("worker CPU inventory is required"),
)
}
if control == nil || commands == nil || store == nil || capabilityExecutor == nil || ids == nil || observer == nil {
return nil, errors.New("control plane, command stream, local store, capability, ID generator, and observer are required")
}
return &Runtime{
registration: registration,
control: control,
commands: commands,
store: store,
capabilityExecutor: capabilityExecutor,
ids: ids,
observer: observer,
now: time.Now,
retryInitial: time.Second,
outboxWake: make(chan struct{}, 1),
outboxDrain: make(chan chan error),
active: make(map[string]context.CancelFunc),片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
核心定义协作契约 backend · c6dbe05c
server/internal/core/execution/ports/outbound/repository.go · 第 9–24 行
符号:Repository · 核对日期 2026-10-02
来源与提交版本一致
// Idempotency 描述一次写操作的调用方、键和请求指纹。
type Idempotency struct {
PrincipalID string // 隔离不同调用方的同名幂等键。
Operation string // 隔离同一调用方的不同业务动作。
Key string // 调用方提供的幂等键。
RequestHash []byte // 用于拒绝复用同一键但内容不同的请求。
}
// Repository 原子持久化 Task、Execution、首个事件和幂等记录。
type Repository interface {
// Create 创建聚合;重复且请求指纹一致时返回原结果。
Create(context.Context, *domain.Task, *domain.Execution, domain.Event, Idempotency) (*domain.Task, *domain.Execution, error)
// Get 在应用范围内读取任务。
Get(context.Context, Scope, domain.ID) (*domain.Task, error)
}
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
Server 的 HTTP 入口按职责组合 Handler,组合根通过显式构造接入依赖。这些是可观察的结构,不因此声称使用了完整的 GoF Composite 或抽象工厂层级。
失败与边界:抽象只对真正的共同语义有效
一个“万能执行”接口如果包含几十个可选字段,每个实现都忽略其中一半,调用者必须猜它支持什么。这是接口形状统一、语义不统一。使用明确的能力契约、版本和输入 Schema,公开差异。
不建议为了“未来可以替换”给所有值对象和纯函数加接口;测试纯函数无需 mock。抽象外部时间、数据库和网络则有明确的边界和测试收益。
迁移练习与参考答案
练习:导出系统新增 CSV 与 PDF,还要上传文件和发送完成通知。策略接口应该包含导出、上传、通知三个方法吗?
参考答案:导出格式是变化点,定义生成文档的策略;上传与通知是其他协作者,由应用用例编排。接口应明确输出如何流式处理、失败是否留下临时文件。PDF 的大内存需求和 CSV 的流式能力可作为独立资源约束,不隐瞒在共同签名后。
继续阅读:Worker 持久化与恢复、测试与故障验证。