fix(runner): 稳定异步任务重启恢复
ci / verify (pull_request) Successful in 10m32s
ci / verify (pull_request) Successful in 10m32s
避免恢复任务将已到期的 next_run_at 误插入 River scheduled 状态,并使用专用模型隔离重启集成测试。\n\n验证:Go 全量测试、govulncheck、pnpm lint/test/build/audit、Compose、镜像、迁移、流水线及 SemVer 门禁全部通过。
This commit is contained in:
@@ -156,7 +156,7 @@ func (s *Service) recoverAsyncRiverJobs(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
for _, item := range items {
|
||||
result, err := s.riverClient.Insert(ctx, asyncTaskArgs{TaskID: item.ID}, asyncTaskRecoveryInsertOpts(item))
|
||||
result, err := s.riverClient.Insert(ctx, asyncTaskArgs{TaskID: item.ID}, asyncTaskRecoveryInsertOpts(item, time.Now()))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -196,9 +196,11 @@ func asyncTaskInsertOpts(task store.GatewayTask) *river.InsertOpts {
|
||||
}
|
||||
}
|
||||
|
||||
func asyncTaskRecoveryInsertOpts(item store.AsyncTaskQueueItem) *river.InsertOpts {
|
||||
func asyncTaskRecoveryInsertOpts(item store.AsyncTaskQueueItem, now time.Time) *river.InsertOpts {
|
||||
opts := asyncTaskInsertOpts(store.GatewayTask{ID: item.ID})
|
||||
opts.ScheduledAt = item.NextRunAt
|
||||
if item.NextRunAt.After(now) {
|
||||
opts.ScheduledAt = item.NextRunAt
|
||||
}
|
||||
// A replacement process must not be blocked by a River row that the dead
|
||||
// process left in running state. PostgreSQL execution leases still ensure
|
||||
// that only one recovery job can call the upstream provider.
|
||||
|
||||
Reference in New Issue
Block a user