Preempt 获取一个任务,并启动一个协程定期刷新任务的更新时间
(ctx context.Context)
| 29 | |
| 30 | // Preempt 获取一个任务,并启动一个协程定期刷新任务的更新时间 |
| 31 | func (c *cronJobService) Preempt(ctx context.Context) (domain.Job, error) { |
| 32 | j, err := c.repo.Preempt(ctx) |
| 33 | if err != nil { |
| 34 | return domain.Job{}, err |
| 35 | } |
| 36 | |
| 37 | ticker := time.NewTicker(c.refreshInterval) |
| 38 | go func() { |
| 39 | for { |
| 40 | select { |
| 41 | case <-ticker.C: |
| 42 | c.refresh(j.Id) |
| 43 | case <-ctx.Done(): |
| 44 | ticker.Stop() |
| 45 | return |
| 46 | } |
| 47 | } |
| 48 | }() |
| 49 | j.CancelFunc = func() { |
| 50 | ticker.Stop() |
| 51 | ct, cancel := context.WithTimeout(context.Background(), time.Second) |
| 52 | defer cancel() |
| 53 | er := c.repo.Release(ct, j.Id) |
| 54 | if er != nil { |
| 55 | c.l.Error("Failed to release job", zap.Error(er)) |
| 56 | } |
| 57 | } |
| 58 | return j, nil |
| 59 | } |
| 60 | |
| 61 | // ResetNextTime 重置任务的下次执行时间 |
| 62 | func (c *cronJobService) ResetNextTime(ctx context.Context, dj domain.Job) error { |