forked from contribsys/faktory
-
Notifications
You must be signed in to change notification settings - Fork 0
/
workers.go
102 lines (87 loc) · 2.03 KB
/
workers.go
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
package server
import (
"encoding/json"
"fmt"
"sync"
"time"
"github.com/contribsys/faktory/util"
)
type ClientWorker struct {
Hostname string `json:"hostname"`
Wid string `json:"wid"`
Pid int `json:"pid"`
Labels []string `json:"labels"`
Salt string `json:"salt"`
PasswordHash string `json:"pwdhash"`
StartedAt time.Time
lastHeartbeat time.Time
signal string
state string
}
func clientWorkerFromAhoy(data string) (*ClientWorker, error) {
var client ClientWorker
err := json.Unmarshal([]byte(data), &client)
if err != nil {
return nil, err
}
if client.Wid == "" {
return nil, fmt.Errorf("Invalid client Wid")
}
return &client, nil
}
func (worker *ClientWorker) IsQuiet() bool {
return worker.state != ""
}
/*
* Send "quiet" or "terminate" to the given client
* worker process. Other signals are undefined.
*/
func (worker *ClientWorker) Signal(sig string) {
worker.signal = sig
worker.state = sig
}
func (worker *ClientWorker) BusyCount() int {
workingMutex.Lock()
defer workingMutex.Unlock()
count := 0
for _, res := range workingMap {
if res.Wid == worker.Wid {
count += 1
}
}
return count
}
/*
* Removes any heartbeat records over 1 minute old.
*/
func reapHeartbeats(heartbeats map[string]*ClientWorker, mu *sync.Mutex) error {
toDelete := []string{}
for k, worker := range heartbeats {
if worker.lastHeartbeat.Before(time.Now().Add(-1 * time.Minute)) {
toDelete = append(toDelete, k)
}
}
if len(toDelete) > 0 {
mu.Lock()
for _, k := range toDelete {
delete(heartbeats, k)
}
mu.Unlock()
util.Debugf("Reaped %d worker heartbeats", len(toDelete))
}
return nil
}
func updateHeartbeat(client *ClientWorker, heartbeats map[string]*ClientWorker, mu *sync.Mutex) {
mu.Lock()
val, ok := heartbeats[client.Wid]
if ok {
val.lastHeartbeat = time.Now()
} else {
client.StartedAt = time.Now()
client.lastHeartbeat = time.Now()
heartbeats[client.Wid] = client
val = client
}
mu.Unlock()
util.Debugf("%+v", val)
}