add federated worker registration and heartbeats
This commit is contained in:
@@ -0,0 +1,70 @@
|
||||
package federation
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
var ErrUnknownWorker = errors.New("unknown worker")
|
||||
|
||||
type Worker struct {
|
||||
ID string `json:"id"`
|
||||
Address string `json:"address"`
|
||||
Capacity int `json:"capacity"`
|
||||
LastSeen time.Time `json:"last_seen"`
|
||||
Online bool `json:"online"`
|
||||
}
|
||||
|
||||
type Registry struct {
|
||||
mu sync.Mutex
|
||||
workers map[string]Worker
|
||||
TTL time.Duration
|
||||
}
|
||||
|
||||
func (r *Registry) init() {
|
||||
if r.TTL <= 0 {
|
||||
r.TTL = 90 * time.Second
|
||||
}
|
||||
if r.workers == nil {
|
||||
r.workers = map[string]Worker{}
|
||||
}
|
||||
}
|
||||
func (r *Registry) Register(w Worker) error {
|
||||
if w.ID == "" {
|
||||
return errors.New("worker id required")
|
||||
}
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.init()
|
||||
w.LastSeen = time.Now().UTC()
|
||||
w.Online = true
|
||||
r.workers[w.ID] = w
|
||||
return nil
|
||||
}
|
||||
func (r *Registry) Heartbeat(id string) error {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.init()
|
||||
w, ok := r.workers[id]
|
||||
if !ok {
|
||||
return ErrUnknownWorker
|
||||
}
|
||||
w.LastSeen = time.Now().UTC()
|
||||
w.Online = true
|
||||
r.workers[id] = w
|
||||
return nil
|
||||
}
|
||||
func (r *Registry) Snapshot() []Worker {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.init()
|
||||
now := time.Now()
|
||||
out := make([]Worker, 0, len(r.workers))
|
||||
for id, w := range r.workers {
|
||||
w.Online = now.Sub(w.LastSeen) <= r.TTL
|
||||
r.workers[id] = w
|
||||
out = append(out, w)
|
||||
}
|
||||
return out
|
||||
}
|
||||
Reference in New Issue
Block a user