diff --git a/api_stats.go b/api_stats.go index abdb75b..a8a7b0a 100644 --- a/api_stats.go +++ b/api_stats.go @@ -3,6 +3,7 @@ package workers import ( "encoding/json" "net/http" + "time" ) func (s *apiServer) Stats(w http.ResponseWriter, req *http.Request) { @@ -24,14 +25,15 @@ func (s *apiServer) Stats(w http.ResponseWriter, req *http.Request) { enc.Encode(allStats) } -// Stats containts current stats for a manager +// Stats contains current stats for a manager type Stats struct { - Name string `json:"manager_name"` - Processed int64 `json:"processed"` - Failed int64 `json:"failed"` - Jobs map[string][]JobStatus `json:"jobs"` - Enqueued map[string]int64 `json:"enqueued"` - RetryCount int64 `json:"retry_count"` + Name string `json:"manager_name"` + Processed int64 `json:"processed"` + Failed int64 `json:"failed"` + Jobs map[string][]JobStatus `json:"jobs"` + Enqueued map[string]int64 `json:"enqueued"` + RetryCount int64 `json:"retry_count"` + HeartbeatLastPushedAt time.Time `json:"heartbeat_last_pushed_at"` } // JobStatus contains the status and data for active jobs of a manager diff --git a/heartbeat.go b/heartbeat.go index 8d7334a..dbb5b29 100644 --- a/heartbeat.go +++ b/heartbeat.go @@ -53,6 +53,7 @@ func GenerateProcessNonce() (string, error) { func (m *Manager) buildHeartbeat(heartbeatTime time.Time, ttl time.Duration) (*storage.Heartbeat, error) { queues := []string{} + msgs := map[string]string{} concurrency := 0 busy := 0 @@ -67,6 +68,41 @@ func (m *Manager) buildHeartbeat(heartbeatTime time.Time, ttl time.Duration) (*s w.runnersLock.Lock() for _, r := range w.runners { + + msg := r.inProgressMessage() + if msg == nil { + continue + } + + workerMsg := &HeartbeatWorkerMsg{ + Retry: 1, + Queue: w.queue, + Backtrace: false, + Class: msg.Class(), + Args: msg.Args(), + Jid: msg.Jid(), + CreatedAt: msg.startedAt, // not actually started at + EnqueuedAt: time.Now().UTC().Unix(), + } + + jsonMsg, err := json.Marshal(workerMsg) + if err != nil { + return nil, err + } + + wrapper := &HeartbeatWorkerMsgWrapper{ + Queue: w.queue, + Payload: string(jsonMsg), + RunAt: msg.startedAt, + } + + jsonWrapper, err := json.Marshal(wrapper) + if err != nil { + return nil, err + } + + msgs[r.tid] = string(jsonWrapper) + workerHeartbeat := storage.WorkerHeartbeat{ Pid: pid, Tid: r.tid, @@ -125,6 +161,7 @@ func (m *Manager) buildHeartbeat(heartbeatTime time.Time, ttl time.Duration) (*s ActiveManager: m.IsActive(), WorkerHeartbeats: workerHeartbeats, Ttl: ttl, + WorkerMessages: msgs, } if m.opts.Heartbeat != nil && m.opts.Heartbeat.PrioritizedManager != nil { heartbeat.ManagerPriority = m.opts.Heartbeat.PrioritizedManager.ManagerPriority diff --git a/manager.go b/manager.go index 1b29a61..603a94c 100644 --- a/manager.go +++ b/manager.go @@ -15,19 +15,20 @@ import ( // Manager coordinates work, workers, and signaling needed for job processing type Manager struct { - uuid string - opts Options - schedule *scheduledWorker - workers []*worker - lock sync.Mutex - signal chan os.Signal - running bool - stop chan bool - active bool - logger *log.Logger - startedAt time.Time - processNonce string - heartbeatChannel chan bool + uuid string + opts Options + schedule *scheduledWorker + workers []*worker + lock sync.Mutex + signal chan os.Signal + running bool + stop chan bool + active bool + logger *log.Logger + startedAt time.Time + processNonce string + heartbeatChannel chan bool + heartbeatLastPushedAt time.Time beforeStartHooks []func() duringDrainHooks []func() @@ -245,9 +246,10 @@ func (m *Manager) Producer() *Producer { // GetStats returns the set of stats for the manager func (m *Manager) GetStats() (Stats, error) { stats := Stats{ - Jobs: map[string][]JobStatus{}, - Enqueued: map[string]int64{}, - Name: m.opts.ManagerDisplayName, + Jobs: map[string][]JobStatus{}, + Enqueued: map[string]int64{}, + Name: m.opts.ManagerDisplayName, + HeartbeatLastPushedAt: m.heartbeatLastPushedAt, } var q []string @@ -321,24 +323,25 @@ func (m *Manager) startHeartbeat() error { m.logger.Println("ERR: Failed to get heartbeat time", err) return err } - heartbeat, err := m.sendHeartbeat(heartbeatTime) + _, err = m.sendHeartbeat(heartbeatTime) if err != nil { - m.logger.Println("ERR: Failed to send heartbeat", err) return err } - expireTS := heartbeatTime.Add(-m.opts.Heartbeat.HeartbeatTTL).Unix() - staleMessageUpdates, err := m.handleAllExpiredHeartbeats(context.Background(), expireTS) - if err != nil { - m.logger.Println("ERR: error expiring heartbeat identities", err) - return err - } - for _, afterHeartbeatHook := range m.afterHeartbeatHooks { - err := afterHeartbeatHook(heartbeat, m, staleMessageUpdates) - if err != nil { - m.logger.Println("ERR: Failed to execute after heartbeat hook", err) - return err - } - } + + //expireTS := heartbeatTime.Add(-m.opts.Heartbeat.HeartbeatTTL).Unix() + //staleMessageUpdates, err := m.handleAllExpiredHeartbeats(context.Background(), expireTS) + //if err != nil { + // m.logger.Println("ERR: error expiring heartbeat identities", err) + // return err + //} + //for _, afterHeartbeatHook := range m.afterHeartbeatHooks { + // err := afterHeartbeatHook(heartbeat, m, staleMessageUpdates) + // if err != nil { + // m.logger.Println("ERR: Failed to execute after heartbeat hook", err) + // return err + // } + //} + m.heartbeatLastPushedAt = time.Now() case <-m.heartbeatChannel: return nil } @@ -418,10 +421,14 @@ func (m *Manager) stopHeartbeat() { func (m *Manager) sendHeartbeat(heartbeatTime time.Time) (*storage.Heartbeat, error) { heartbeat, err := m.buildHeartbeat(heartbeatTime, m.opts.Heartbeat.HeartbeatTTL) if err != nil { + m.logger.Println("ERR: Failed to build heartbeat", err) return heartbeat, err } err = m.opts.store.SendHeartbeat(context.Background(), heartbeat) + if err != nil { + m.logger.Println("ERR: Failed to send heartbeat", err) + } return heartbeat, err } diff --git a/storage/redis.go b/storage/redis.go index d84c5b4..50a7a74 100644 --- a/storage/redis.go +++ b/storage/redis.go @@ -160,6 +160,22 @@ func (r *redisStore) SendHeartbeat(ctx context.Context, heartbeat *Heartbeat) er "active_manager", heartbeat.ActiveManager, "worker_heartbeats", workerHeartbeats) + // ensure the heartbeat is automatically cleaned up + pipe.Expire(ctx, managerKey, heartbeat.Ttl) + + // delete the worker key just in case our set is empty + pipe.Del(ctx, GetWorkersKey(managerKey)) + + // send all job message heartbeats + for tid, msg := range heartbeat.WorkerMessages { + // fake the sidekiq thread id + fakeThreadId := fmt.Sprintf("%d-%s", heartbeat.Pid, tid) + pipe.HSet(ctx, GetWorkersKey(managerKey), fakeThreadId, msg) + } + + // make sure the worker is cleaned up + pipe.Expire(ctx, GetWorkersKey(managerKey), heartbeat.Ttl) + _, err = pipe.Exec(ctx) if err != nil && err != redis.Nil { return err diff --git a/storage/storage.go b/storage/storage.go index 8d4ac7e..6618b96 100644 --- a/storage/storage.go +++ b/storage/storage.go @@ -39,15 +39,16 @@ type Retries struct { type Heartbeat struct { Identity string `json:"identity"` - Beat int64 `json:"beat,string"` - Quiet bool `json:"quiet,string"` - Busy int `json:"busy,string"` - RttUS int `json:"rtt_us,string"` - RSS int64 `json:"rss,string"` - Info string `json:"info"` - Pid int `json:"pid,string"` - ManagerPriority int `json:"manager_priority,string"` - ActiveManager bool `json:"active_manager,string"` + Beat int64 `json:"beat,string"` + Quiet bool `json:"quiet,string"` + Busy int `json:"busy,string"` + RttUS int `json:"rtt_us,string"` + RSS int64 `json:"rss,string"` + Info string `json:"info"` + Pid int `json:"pid,string"` + ManagerPriority int `json:"manager_priority,string"` + ActiveManager bool `json:"active_manager,string"` + WorkerMessages map[string]string `json:"worker_messages"` Ttl time.Duration