-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathenqueue.go
More file actions
111 lines (100 loc) · 2.42 KB
/
Copy pathenqueue.go
File metadata and controls
111 lines (100 loc) · 2.42 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
package bgjob
import (
"context"
"database/sql"
"errors"
"fmt"
"strings"
)
type ExecerContext interface {
ExecContext(ctx context.Context, s string, args ...any) (sql.Result, error)
}
func Enqueue(ctx context.Context, e ExecerContext, req EnqueueRequest) error {
return BulkEnqueue(ctx, e, []EnqueueRequest{req})
}
func BulkEnqueue(ctx context.Context, e ExecerContext, list []EnqueueRequest) error {
jobs, err := requestsToJobs(list)
if err != nil {
return err
}
err = bulkInsert(ctx, e, jobs)
if err == ErrJobAlreadyExist {
return err
}
if err != nil {
return fmt.Errorf("bulk insert: %w", err)
}
return nil
}
func requestsToJobs(list []EnqueueRequest) ([]Job, error) {
if len(list) == 0 {
return nil, errors.New("list is empty. at least one job is expected")
}
jobs := make([]Job, 0, len(list))
now := timeNow()
for _, req := range list {
if req.Queue == "" {
return nil, ErrQueueIsRequired
}
if req.Type == "" {
return nil, ErrTypeIsRequired
}
id := req.Id
if id == "" {
generated, err := nextId()
if err != nil {
return nil, fmt.Errorf("generate id: %w", err)
}
id = generated
}
job := Job{
Id: id,
Queue: req.Queue,
Type: req.Type,
Arg: req.Arg,
RequestId: req.RequestId,
Attempt: 0,
LastError: nil,
NextRunAt: now.Add(req.Delay).Unix(),
CreatedAt: now,
UpdatedAt: now,
}
jobs = append(jobs, job)
}
return jobs, nil
}
func bulkInsert(ctx context.Context, e ExecerContext, jobs []Job) error {
valueStrings := make([]string, 0, len(jobs))
valueArgs := make([]any, 0, len(jobs)*9)
placeholderNum := 0
for _, job := range jobs {
placeholders := make([]string, 0)
for i := 0; i < 9; i++ {
placeholderNum++
placeholders = append(placeholders, fmt.Sprintf("$%d", placeholderNum))
}
valueStrings = append(valueStrings, fmt.Sprintf("(%s)", strings.Join(placeholders, ",")))
valueArgs = append(
valueArgs,
job.Id,
job.Queue,
job.Type,
job.Arg,
job.Attempt,
job.NextRunAt,
job.CreatedAt,
job.UpdatedAt,
job.RequestId,
)
}
query := fmt.Sprintf("INSERT INTO bgjob_job (id, queue, type, arg, attempt, next_run_at, created_at, updated_at, request_id) VALUES %s",
strings.Join(valueStrings, ","))
_, err := e.ExecContext(ctx, query, valueArgs...)
if err != nil {
if strings.Contains(err.Error(), "23505") { //unique_violation
return ErrJobAlreadyExist
}
return err
}
return nil
}