-
Notifications
You must be signed in to change notification settings - Fork 242
Expand file tree
/
Copy pathworkers.go
More file actions
78 lines (63 loc) · 1.54 KB
/
Copy pathworkers.go
File metadata and controls
78 lines (63 loc) · 1.54 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
package goworker
import (
"encoding/json"
"fmt"
"sync"
)
type workersMutex struct {
sync.RWMutex
workers map[string]workerFunc
}
func (wm *workersMutex) Add(class string, worker workerFunc) {
wm.Lock()
defer wm.Unlock()
wm.workers[class] = worker
}
func (wm *workersMutex) Get(class string) (worker workerFunc, ok bool) {
wm.RLock()
defer wm.RUnlock()
worker, ok = wm.workers[class]
return
}
var workers *workersMutex
func init() {
workers = &workersMutex{
RWMutex: sync.RWMutex{},
workers: make(map[string]workerFunc),
}
}
// Register registers a goworker worker function. Class
// refers to the Ruby name of the class which enqueues the
// job. Worker is a function which accepts a queue and an
// arbitrary array of interfaces as arguments.
func Register(class string, worker workerFunc) {
workers.Add(class, worker)
}
func Enqueue(job *Job) error {
err := Init()
if err != nil {
return err
}
conn, err := GetConn()
if err != nil {
logger.Criticalf("Error on getting connection on enqueue")
return err
}
defer PutConn(conn)
buffer, err := json.Marshal(job.Payload)
if err != nil {
logger.Criticalf("Cant marshal payload on enqueue")
return err
}
err = conn.Send("RPUSH", fmt.Sprintf("%squeue:%s", workerSettings.Namespace, job.Queue), buffer)
if err != nil {
logger.Criticalf("Cant push to queue")
return err
}
err = conn.Send("SADD", fmt.Sprintf("%squeues", workerSettings.Namespace), job.Queue)
if err != nil {
logger.Criticalf("Cant register queue to list of use queues")
return err
}
return conn.Flush()
}