21 KiB
Redis SSOT + 事件重放 + 三层调度器 Phase 2 实施计划
For agentic workers: REQUIRED: Use superpowers:subagent-driven-development (if subagents available) or superpowers:executing-plans to implement this plan. Steps use checkbox (
- [ ]) syntax for tracking.
Goal: 将 Redis SSOT、事件重放、三层调度器从当前"已实现但未接通"状态推进到设计文档要求的完整落地——Redis 成为运行时唯一真相源,调度器通过 ZPOPMAX 事件驱动,断电恢复可安全重放。
Architecture: 分三步走:(1) 修复关键断路,让 RedisBus/StateManager 完整运行;(2) 建立运行时写穿通道,Processor 事件处理同步更新 Redis;(3) 接入 ZPOPMAX 调度循环和原子抢占 Lua 脚本。所有变更加 SSOT 开关隔离,不破坏现有执行内核。
Tech Stack: Go 1.23.4, go-zero, ent v0.14.5, Redis (go-redis/v9), PostgreSQL, testify
文件结构
必改文件
- Modify:
internal/svc/service_context.go— 调用 RedisBus.Start()、注入 RecoveryCoordinator - Modify:
internal/state/manager.go— 补全 GetTaskState 接口方法 - Modify:
internal/state/redis_manager.go— 修 ApplyJobTransitionScript 版本号、设备心跳 - Modify:
internal/state/memory_manager.go— 完整实现 ApplyTaskState/GetTaskState - Modify:
internal/state/keys.go— 新增 sched:ready_jobs、equipment heartbeat key - Modify:
internal/eventbus/bus.go— Bus 接口增加 Start/Stop 方法(或保障 RedisBus 启动) - Modify:
internal/processor/job_processor.go— 写穿 StateManager、接入 sched:ready_jobs - Modify:
internal/processor/dispatcher.go— Task 超时检测 - Modify:
internal/processor/scheduler_adapter.go— ZPOPMAX 调度循环适配 - Modify:
internal/scheduler/scheduler.go— 接收 ZPOPMAX 单工件输入 - Modify:
internal/recovery/replayer.go— 增量重放、ReplaySince - Modify:
internal/recovery/assessor.go— 接入 restoreActiveOrders - Modify:
internal/robot/common.go— 心跳信号
新建文件
- Create:
internal/processor/recovery_coordinator.go— RecoveryCoordinator 实现 - Create:
internal/state/state_write_through.go— 写穿辅助函数(Processor→Redis) - Create:
internal/scheduler/zpop.go— ZPOPMAX 调度循环 - Create:
internal/state/redis_manager_test.go— RedisStateManager 集成测试 - Create:
internal/processor/ssot_integration_test.go— SSOT 写穿集成测试
现有测试要改
- Modify:
internal/state/manager_test.go— 新增 GetTaskState 测试 - Modify:
internal/scheduler/scheduler_test.go— 新增 ZPOPMAX 输入路径测试
Chunk 1: 修复关键断路
Task 1: 调用 RedisBus.Start() + Bus 接口补全
Files:
- Modify:
internal/eventbus/bus.go - Modify:
internal/svc/service_context.go
Depends on: 无
- Step 1: Bus 接口增加 Start/Stop 方法
在 bus.go 的 Bus interface 追加:
type Bus interface {
Publish(ctx context.Context, event Event) error
Subscribe(eventType EventType, handler EventHandler) (subID string, err error)
SubscribeAll(handler EventHandler) (subID string, err error)
Unsubscribe(subID string) error
Close() error
Start(ctx context.Context) error // 新增
Stop() error // 新增
}
LocalBus 的 Start/Stop 为空操作。
- Step 2: 在 ServiceContext 中调用 bus.Start()
在 service_context.go 的 NewServiceContext 中,bus 创建之后、订阅 EventLogWriter 之前:
if err := bus.Start(context.Background()); err != nil {
slog.Error("svc: bus start failed", "error", err)
}
- Step 3: 验证编译
Run: rtk go build ./...
Expected: PASS
- Step 4: 验证 LocalBus 不受影响
Run: rtk go test ./internal/eventbus/... ./internal/processor/...
Expected: 现有测试全部 PASS
Task 2: 修复 ApplyJobTransitionScript 版本号不一致
Files:
- Modify:
internal/state/redis_manager.go
Depends on: Task 1
- Step 1: 修改 Lua 脚本使版本递增
当前 applyJobTransitionScript 直接使用 ARGV[5] 作为新版本:
-- 当前(有问题)
redis.call('HSET', jobKey, 'status', newStatus, 'version', newVersion)
-- 修复为递增
local currentVer = tonumber(redis.call('HGET', jobKey, 'version') or '0')
local expectedVer = tonumber(ARGV[4])
if currentVer ~= expectedVer then return 0 end
redis.call('HSET', jobKey, 'status', newStatus, 'version', currentVer + 1)
- Step 2: 运行 StateManager 测试
Run: rtk go test ./internal/state/...
Expected: 现有测试 PASS(MemoryStateManager 行为不变)
- Step 3: 记录检查点
Task 3: 补全 StateManager 接口
Files:
- Modify:
internal/state/manager.go - Modify:
internal/state/memory_manager.go - Modify:
internal/state/manager_test.go
Depends on: Task 1
- Step 1: 在接口中增加 GetTaskState
type StateManager interface {
// ... 现有方法 ...
GetTaskState(ctx context.Context, taskID string) (*TaskState, error) // 新增
}
同时定义 TaskState 结构体。
- Step 2: MemoryStateManager 完整实现 ApplyTaskState + GetTaskState
用内存 map 存储 task state,与 job state 同模式。
- Step 3: RedisStateManager 实现 GetTaskState
func (m *RedisStateManager) GetTaskState(ctx context.Context, taskID string) (*TaskState, error) {
key := taskKey(taskID)
fields, err := m.client.HGetAll(ctx, key).Result()
// ...
}
- Step 4: 写测试并验证
func TestGetTaskState(t *testing.T) {
mgr := NewMemoryStateManager()
err := mgr.ApplyTaskState(ctx, "task-1", &TaskState{Status: "RUNNING"})
require.NoError(t, err)
ts, err := mgr.GetTaskState(ctx, "task-1")
require.NoError(t, err)
assert.Equal(t, "RUNNING", ts.Status)
}
Run: rtk go test ./internal/state/...
Expected: PASS
Chunk 2: 运行时 SSOT 写穿通道
Task 4: 创建写穿辅助层
Files:
- Create:
internal/state/state_write_through.go
Depends on: Chunk 1 完成
- Step 1: 创建 WriteThroughHelper
package state
// WriteThroughHelper 封装 Processor→Redis 的写穿逻辑。
// 当 SSOT 禁用时,降级为仅写 MemoryStateManager。
type WriteThroughHelper struct {
mgr StateManager
enabled bool
}
func NewWriteThroughHelper(mgr StateManager, enabled bool) *WriteThroughHelper {
return &WriteThroughHelper{mgr: mgr, enabled: enabled}
}
// SyncJobState 写穿:将 JobRuntime 状态同步到 StateManager。
// enabled=false 时仅更新内存,不写 Redis。
func (h *WriteThroughHelper) SyncJobState(ctx context.Context, jobID int, status string, positionType, positionRefID string, version int) error {
return h.mgr.ApplyJobState(ctx, jobID, map[string]interface{}{
"status": status,
"position_type": positionType,
"position_ref": positionRefID,
"version": version,
})
}
- Step 2: 验证编译
Run: rtk go build ./internal/state/...
Expected: PASS
Task 5: JobProcessor 事件处理器写穿 Redis
Files:
- Modify:
internal/processor/job_processor.go
Depends on: Task 4
- Step 1: 在 JobProcessor 中注入 WriteThroughHelper
type JobProcessor struct {
// ... 现有字段 ...
ssotWriter *state.WriteThroughHelper
}
func (jp *JobProcessor) SetSSOTWriter(w *state.WriteThroughHelper) {
jp.ssotWriter = w
}
- Step 2: 在关键状态变更点调用写穿
在 onJobCompleted、onJobError、handleMachineTaskComplete 等处理器中,状态变更后调用:
if jp.ssotWriter != nil {
jp.ssotWriter.SyncJobState(ctx, jobID,
string(domainStatus),
string(jr.PositionType),
jr.PositionRefID,
jr.version,
)
}
- Step 3: ServiceContext 接线
在 service_context.go 中:
if stateMgr != nil {
writer := state.NewWriteThroughHelper(stateMgr, c.SSOT.Enabled)
orderProcessor.SetSSOTWriter(writer)
}
- Step 4: 验证编译 + 现有测试
Run: rtk go build ./... && rtk go test ./internal/processor/...
Expected: build PASS, tests PASS(旧路径 SSOT.Enabled=false 不受影响)
- Step 5: 记录检查点
Task 6: 写穿集成测试
Files:
- Create:
internal/processor/ssot_integration_test.go
Depends on: Task 5
- Step 1: 写测试 — 验证 Job 状态变更后 Redis 可读
func TestSSOTWriteThrough_JobCompleted(t *testing.T) {
mgr := state.NewMemoryStateManager()
writer := state.NewWriteThroughHelper(mgr, true)
// 模拟状态变更
err := writer.SyncJobState(ctx, 1, "COMPLETED", "ON_DOCK", "", 2)
require.NoError(t, err)
js, err := mgr.GetJobState(ctx, 1)
require.NoError(t, err)
assert.Equal(t, "COMPLETED", js["status"])
}
- Step 2: 运行测试
Run: rtk go test ./internal/processor/... -run SSOT
Expected: PASS
- Step 3: 写测试 — 验证 SSOT.Enabled=false 时 MemoryStateManager 仍工作
func TestSSOTWriteThrough_Disabled(t *testing.T) {
mgr := state.NewMemoryStateManager()
writer := state.NewWriteThroughHelper(mgr, false) // 禁用
err := writer.SyncJobState(ctx, 1, "COMPLETED", "ON_DOCK", "", 2)
require.NoError(t, err)
// MemoryStateManager 仍能读到(本地写仍执行)
}
Run: rtk go test ./internal/processor/... -run SSOT
Expected: PASS
Chunk 3: 事件驱动调度器(sched:ready_jobs + ZPOPMAX)
Task 7: 实现 sched:ready_jobs 推入逻辑
Files:
- Modify:
internal/state/keys.go - Create:
internal/scheduler/zpop.go
Depends on: Chunk 2 完成
- Step 1: 新增 Redis key 定义
在 keys.go 追加:
const (
SchedReadyJobsKey = "sched:ready_jobs" // Sorted Set: score=priority, member=jobID
)
- Step 2: 在写穿层增加 PushReadyJob
在 state_write_through.go:
func (h *WriteThroughHelper) PushReadyJob(ctx context.Context, jobID int, priority float64) error {
if !h.enabled || h.mgr == nil {
return nil
}
if rm, ok := h.mgr.(*RedisStateManager); ok {
return rm.client.ZAdd(ctx, SchedReadyJobsKey, redis.Z{
Score: priority,
Member: strconv.Itoa(jobID),
}).Err()
}
return nil
}
- Step 3: 在 JobProcessor 状态变更点调用 PushReadyJob
当 Job 进入 ON_BUFFER 或 WAITING_UNLOAD 时推送:
if jp.ssotWriter != nil && (newStatus == constants.JobStatus_OnBuffer || newStatus == constants.JobStatus_WaitingUnload) {
jp.ssotWriter.PushReadyJob(ctx, jobID, float64(jr.Priority))
}
- Step 4: 验证编译
Run: rtk go build ./...
Expected: PASS
Task 8: 实现 ZPOPMAX 调度循环
Files:
- Create:
internal/scheduler/zpop.go(如已创建则修改)
Depends on: Task 7
- Step 1: 实现 ZPOPMAX 阻塞循环
package scheduler
// ZPopLoop ZPOPMAX 阻塞调度循环。
// 设计文档 8.0 节:调度引擎阻塞在 ZPOPMAX 上,仅处理变动工件。
func ZPopLoop(ctx context.Context, rdb *redis.Client, handler func(jobID int) error) {
for {
select {
case <-ctx.Done():
return
default:
}
// BZPOPMAX 阻塞等待,超时 5s
results, err := rdb.BZPopMax(ctx, 5*time.Second, state.SchedReadyJobsKey).Result()
if err == redis.Nil {
continue
}
if err != nil {
slog.Error("zpop: failed", "error", err)
time.Sleep(time.Second)
continue
}
jobID, _ := strconv.Atoi(results.Member.(string))
if err := handler(jobID); err != nil {
// 处理失败放回队列(稍后重试)
rdb.ZAdd(ctx, state.SchedReadyJobsKey, redis.Z{
Score: results.Score - 0.1, // 降优先级
Member: results.Member,
})
}
}
}
- Step 2: 在 ServiceContext 启动 ZPOP 循环
if c.SSOT.Enabled && goRedisClient != nil && sched != nil {
go scheduler.ZPopLoop(context.Background(), goRedisClient, func(jobID int) error {
return orderProcessor.ScheduleSingleJob(ctx, jobID)
})
}
- Step 3: 在 JobProcessor 增加 ScheduleSingleJob
func (jp *JobProcessor) ScheduleSingleJob(ctx context.Context, jobID int) error {
jp.jobsMu.RLock()
jr, ok := jp.jobs[jobID]
jp.jobsMu.RUnlock()
if !ok {
return fmt.Errorf("job %d not found", jobID)
}
systemState := BuildSystemState(jp)
candidates := jp.sched.Schedule(jr, systemState)
for _, c := range candidates {
task := CandidateToRobotTask(c)
jp.dispatcher.Enqueue(task)
}
return nil
}
- Step 4: 验证编译
Run: rtk go build ./...
Expected: PASS
Task 9: 整合新旧两条调度路径
Files:
- Modify:
internal/processor/job_processor.go - Modify:
internal/processor/scheduler_adapter.go
Depends on: Task 8
- Step 1: ReadyQueue 仅用于非 SSOT 路径
当 SSOT.Enabled=true 时,readyQueue 返回 nil,notifyReady 返回 nil。Job 仅通过 ZPOPMAX 驱动调度。
当 SSOT.Enabled=false 时,保持现有 ReadyQueue→ScheduleAndSubmit 路径。
- Step 2: 验证两条路径不冲突
Run: rtk go test ./internal/processor/... ./internal/scheduler/...
Expected: PASS
- Step 3: 记录检查点
Chunk 4: 原子抢占 + 恢复增强 + 心跳超时
Task 10: 实现原子抢占 Lua 脚本
Files:
- Modify:
internal/state/redis_manager.go
Depends on: Chunk 3 完成
- Step 1: 新增原子抢占 Lua 脚本
设计文档 8.4 节:原子地将 ON_BUFFER → IN_HANDLING。
-- claimJobScript: 原子抢占 Job
local jobKey = KEYS[1]
local expectedStatus = ARGV[1]
local newStatus = ARGV[2]
local expectedVer = tonumber(ARGV[3])
local currentStatus = redis.call('HGET', jobKey, 'status')
local currentVer = tonumber(redis.call('HGET', jobKey, 'version') or '0')
if currentStatus ~= expectedStatus then return {0, 'status_mismatch'} end
if currentVer ~= expectedVer then return {0, 'version_conflict'} end
redis.call('HSET', jobKey, 'status', newStatus, 'version', currentVer + 1)
return {1, 'ok'}
- Step 2: 暴露 ClaimJob 方法
func (m *RedisStateManager) ClaimJob(ctx context.Context, jobID int, expectedStatus, newStatus string, expectedVer int) (bool, error) {
key := jobKey(jobID)
result, err := m.client.Eval(ctx, claimJobScript, []string{key},
expectedStatus, newStatus, expectedVer).Result()
// 返回 (抢占成功, error)
}
将此方法加入 StateManager 接口。
- Step 3: 调度器在派发前原子抢占
在 scheduler_adapter.go 中,CandidateToRobotTask 前调用 ClaimJob。
- Step 4: 测试
Run: rtk go test ./internal/state/...
Expected: PASS
Task 11: 创建 RecoveryCoordinator 并接入恢复流程
Files:
- Create:
internal/processor/recovery_coordinator.go - Modify:
internal/recovery/assessor.go - Modify:
internal/svc/service_context.go
Depends on: Chunk 2 完成
- Step 1: 定义 RecoveryCoordinator
package processor
type RecoveryCoordinator struct {
entClient *ent.Client
stateMgr state.StateManager
eventBus eventbus.Bus
}
func (rc *RecoveryCoordinator) AssessAndRecover(ctx context.Context, job *ent.Job) (*recovery.RecoveryGrade, error) {
grade, manualRequired, message := recovery.AssessOrderGrade([]*ent.Job{job})
switch {
case grade == "L1":
// 自动恢复:直接走 RestoreOrder
return grade, nil
case grade == "L2":
// 半自动:PL 状态确认后恢复
return grade, nil
case grade == "L3" || manualRequired:
// 创建 ManualAction 等待人工
// ...
}
return grade, nil
}
- Step 2: 在 restoreActiveOrders 中调用 RecoveryCoordinator
修改 service_context.go 的 restoreActiveOrders(),通过 RecoveryCoordinator.AssessAndRecover 替代直接 RestoreOrder。
- Step 3: 验证编译
Run: rtk go build ./...
Expected: PASS
Task 12: 设备心跳与 Task 超时
Files:
- Modify:
internal/state/redis_manager.go - Modify:
internal/state/keys.go - Modify:
internal/processor/dispatcher.go
Depends on: Chunk 3 完成
- Step 1: 设备心跳 — Redis TTL
在 redis_manager.go:
func (m *RedisStateManager) UpdateEquipmentHeartbeat(ctx context.Context, equipID string) error {
return m.client.Set(ctx, equipHeartbeatKey(equipID), time.Now().Unix(), 5*time.Second).Err()
}
func (m *RedisStateManager) IsEquipmentOnline(ctx context.Context, equipID string) bool {
_, err := m.client.Get(ctx, equipHeartbeatKey(equipID)).Result()
return err == nil
}
- Step 2: 心跳触发
在 PLC 信号处理循环中,每次收到设备信号时调用 UpdateEquipmentHeartbeat。
- Step 3: Task 超时检测
在 dispatcher.go 中增加超时检测 goroutine:
func (d *Dispatcher) startTimeoutWatchdog(ctx context.Context) {
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
d.checkTaskTimeouts(ctx)
}
}
}
func (d *Dispatcher) checkTaskTimeouts(ctx context.Context) {
d.activeTasksMu.Lock()
defer d.activeTasksMu.Unlock()
for taskID, task := range d.activeTasks {
if time.Since(task.EnqueueTime) > d.maxExecutionTime {
slog.Warn("task timeout", "taskID", taskID, "jobID", task.JobID)
// 标记超时,触发回退
}
}
}
- Step 4: 验证编译
Run: rtk go build ./...
Expected: PASS
Chunk 5: 集成验证与回归
Task 13: SSOT 端到端集成测试
Files:
- Create:
internal/processor/ssot_integration_test.go(追加)
Depends on: Chunk 2-4 完成
- Step 1: 测试 — 完整 Job 生命周期 SSOT 同步
func TestSSOTFullJobLifecycle(t *testing.T) {
mgr := state.NewMemoryStateManager()
writer := state.NewWriteThroughHelper(mgr, true)
// Job 创建 → ON_BUFFER → PROCESSING → WAITING_UNLOAD → COMPLETED
steps := []struct {
status string
posType constants.PositionType
posRef string
}{
{"ON_BUFFER", constants.PositionType_OnBuffer, "5"},
{"PROCESSING", constants.PositionType_OnEquipment, "3:0"},
{"WAITING_UNLOAD", constants.PositionType_OnEquipment, "3:0"},
{"COMPLETED", constants.PositionType_OnDock, "2:1"},
}
for i, s := range steps {
err := writer.SyncJobState(ctx, 1, s.status, string(s.posType), s.posRef, i+1)
require.NoError(t, err)
js, err := mgr.GetJobState(ctx, 1)
require.NoError(t, err)
assert.Equal(t, s.status, js["status"])
}
}
- Step 2: 运行
Run: rtk go test ./internal/processor/... -run SSOT -v
Expected: PASS
Task 14: 全量回归验证
Depends on: Task 13
- Step 1: 全量编译
Run: rtk go build ./...
Expected: PASS
- Step 2: 核心包测试
Run: rtk go test ./internal/state/... ./internal/processor/... ./internal/station/... ./internal/scheduler/... ./internal/eventbus/... ./internal/recovery/...
Expected: PASS
- Step 3: 全量测试(排除已知失败)
Run: rtk go test ./internal/...
Expected: 无新增失败(排除 camera/db/preload/robot 已有问题)
- Step 4: 覆盖率检查
Run: rtk go test -cover ./internal/state/... ./internal/processor/... ./internal/scheduler/...
Expected: 新增代码路径有测试覆盖
- Step 5: 人工验收清单
- SSOT.Enabled=true + EventBus.Enabled=true 时服务正常启动
- Redis 中可查询到 Job 实时状态
- 设备心跳 TTL 正常刷新
- 调度器可从 sched:ready_jobs ZPOPMAX
- SSOT.Enabled=false 退化为旧路径,行为不变
风险点
- RedisBus.Start() 调用时机 — 必须在订阅者注册之后调用,否则事件丢失。处理:先 SubscribeAll,再 Start。
- 写穿通道影响旧路径性能 — 每次状态变更多一层 Redis 调用。处理:异步写(goroutine + channel),不阻塞主调度路径。
- ZPOPMAX 与 ReadyQueue 双路径冲突 — 同一 Job 可能被两次调度。处理:SSOT 启用时完全禁用 ReadyQueue 路径,单一入口。
- 原子抢占 Lua 脚本与现有 ApplyJobState 版本号竞争 — 两个脚本操作同一 Hash。处理:统一版本递增逻辑,所有写入走 Lua。
- 增量重放幂等 — 已处理事件不可重复应用。处理:last_event_id 去重,版本号比较防御。
验证点
rtk go build ./...通过rtk go test ./internal/state/...PASSrtk go test ./internal/processor/...PASSrtk go test ./internal/scheduler/...PASSrtk go test ./internal/eventbus/...PASSrtk go test ./internal/recovery/...PASS- SSOT.Enabled=false 旧路径回归通过
- Redis 中
job:{id}Hash 实时反映工件状态 sched:ready_jobsSorted Set 正确推送和消费
执行说明
- 当前目录不是 git 仓库,跳过 commit 步骤,改为每个任务结束记录检查点。
- 严格遵守项目规则:只改
schema/*.go后再生成 ent;不要手改ent/。 - 使用
rtk前缀运行所有命令。 - 所有新增和修改的代码通过
rtk go build ./...后交付。