Skip to content

设计原则与模式 ​

问题:如何新增能力而不修改整个运行器 ​

一个文件处理器支持缩略图、压缩包清理和水印。如果通用运行器中到处写 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 持久化与恢复、测试与故障验证。

理解原理 · 分析取舍 · 用真实代码检验