Close closes the TaskQueue.
(ctx context.Context)
| 644 | |
| 645 | // Close closes the TaskQueue. |
| 646 | func (q *TaskQueue) Close(ctx context.Context) error { |
| 647 | _, err := q.Redis.Pipelined(ctx, func(p redis.Pipeliner) error { |
| 648 | q.consumerIDs.Range(func(k, v any) bool { |
| 649 | p.XGroupDelConsumer(ctx, InputTaskKey(q.Key), q.Group, k.(string)) |
| 650 | p.XGroupDelConsumer(ctx, ReadyTaskKey(q.Key), q.Group, k.(string)) |
| 651 | return true |
| 652 | }) |
| 653 | return nil |
| 654 | }) |
| 655 | return ConvertError(err) |
| 656 | } |
| 657 | |
| 658 | // Add adds a task s to the queue with a timestamp startAt. |
| 659 | func (q *TaskQueue) Add(ctx context.Context, r redis.Cmdable, s string, startAt time.Time, replace bool) error { |