Skip to content

Worker 持久化与恢复 ​

问题:收到命令后进程马上重启 ​

执行器收到启动命令便回 ACK,然后在内存里启动协程。ACK 成功后协调者可能不再重发,但执行器在持久化前崩溃,工作丢失。反过来,执行已经开始但 ACK 丢失,重发又可能启动两份进程。

正确的顺序要围绕可恢复事实设计:先保存命令及身份,再确认接收;重复命令识别为同一动作。执行结果也先保存,再交给独立发送器可靠上报。

原理:本地库是恢复依据 ​

命令 Inbox 记录接收与处理状态,Attempt 记录执行生命周期,Checkpoint 记录能力可继续的位置,事件 Outbox 保存尚未上报的事实。内存中的 goroutine、channel 和 map 只是当前进程的工作状态。

重启扫描时不能默认重跑所有 RUNNING。不同能力的恢复语义不同:可以从文件偏移续传、可以查询远端任务、可以安全重算、或者必须进入待人工确认。

最小例子:两段可靠边界 ​

text
收到 command_id
  → 本地事务保存 Inbox 和待执行 Attempt
  → ACK 接收
  → 调用能力(可能失败或崩溃)
  → 本地事务保存结果/Checkpoint 和 event_id
  → 发送 event_id
  → Server 持久接收后,本地标记已上报

文件系统与 SQLite 也不是同一个事务。如果文件已写、Checkpoint 未更新,要能通过临时路径、校验和或原子 rename 识别有效产物,不能默认全部已经一致。

方案比较:恢复粒度按能力决定 ​

能力恢复策略必须保存的证据
纯计算用稳定输入重新计算输入版本与输出身份
大文件传输续传或安全重新传输偏移、校验与传输身份
远端长期任务先查询再决定远端 job/session ID
不可查询外部副作用幂等键或人工对账请求身份、未知状态

Checkpoint 只说明记录到哪一步,不能证明外部系统没有多做一步。恢复策略必须解释这个窗口。

真实案例:单发送器与重试分类 ​

Worker Runtime 持有 LocalStore、ControlPlane、CommandStream、CapabilityExecutor 等端口。执行协程保存事件并通过 channel 唤醒发送器,可靠数据留在本地库中。

sendEvents 独占当前会话的发送,网络失败退避,延迟上限 30 秒。只有标为远程投递失败且非认证/授权错误的故障可重试;本地数据库错误不会被伪装成网络重试无限循环。

已核对的实现 · 本地代码快照

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

worker/internal/core/execution/application/outbox.go · 第 34–87 行
符号:sendEvents · 核对日期 2026-10-02
来源与提交版本一致

// sendEvents 独占当前会话的事件发送;能力协程只负责持久化并唤醒发送器。
// 临时网络失败保留事件重试,不取消其他运行中的尝试。
func (runtime *Runtime) sendEvents(ctx context.Context, session Session) error {
	delay := time.Second
	timer := time.NewTimer(0)
	defer timer.Stop()
	for {
		var drained chan error
		select {
		case <-ctx.Done():
			return nil
		case <-runtime.outboxWake:
		case drained = <-runtime.outboxDrain:
		case <-timer.C:
		}
		err := runtime.flushEvents(ctx, session)
		if drained != nil {
			drained <- err
		}
		if err == nil {
			delay = time.Second
			timer.Reset(time.Second)
			continue
		}
		if ctx.Err() != nil {
			return nil
		}

		if !retryableDelivery(err) {
			return err
		}
		runtime.observer.Event(ctx, "event delivery deferred", err)
		// 退避期间不响应生产者唤醒;新增事件仍保留在本地数据库中。
		retry := time.NewTimer(delay)
		select {
		case <-ctx.Done():
			retry.Stop()
			return nil
		case <-retry.C:
		}
		timer.Reset(0)
		delay = min(delay*2, 30*time.Second)
	}
}

// 仅标记远程投递失败;本地数据库错误不能被当作网络故障无限重试。
type deliveryError struct{ error }

func (e *deliveryError) Unwrap() error { return e.error }
func retryableDelivery(err error) bool {
	var delivery *deliveryError
	return errors.As(err, &delivery) && domain.Kind(err) != domain.Unauthenticated && domain.Kind(err) != domain.Forbidden
}

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

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

片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。

SSE 连接只负责命令传输,PostgreSQL worker_commands 与 Worker SQLite Inbox 提供事实。删除 SQLite 不等于清缓存,而是放弃身份和未同步恢复状态。

失败与边界:文件、进程和远端服务 ​

ZIP 清理要拒绝路径穿越、限制展开体积并定义产物提交方式;命令执行要保存进程或外部任务身份,不能因为新进程查不到旧 goroutine 就推断没有任务在运行。

人工 drain 与自动恢复是不同状态。容量缩减或能力变化时,应防止新任务进入,同时处理已开始的 Attempt。清理临时文件也要保留尚有恢复价值的数据。

迁移练习与参考答案 ​

练习:执行器调用云渲染 API,创建成功但还未保存 job ID 就崩溃。如何避免重复渲染?

参考答案:创建前持久化稳定请求 ID,远端按该 ID 幂等创建或提供查找接口;重启查询后保存 job ID 再继续。若远端不支持这些能力,结果处于未知状态,需要对账或人工处理,不能仅靠“先写本地 RUNNING”保证不重复。

继续阅读:资源身份与短期授权、端到端案例。

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