()
| 191 | } |
| 192 | |
| 193 | func (r *basePoolManager) runWatcher() { |
| 194 | defer r.consumer.Close() |
| 195 | queued, err := r.store.ListEntityJobsByStatus(r.ctx, r.entity.EntityType, r.entity.ID, params.JobStatusQueued) |
| 196 | if err != nil { |
| 197 | slog.ErrorContext(r.ctx, "failed to list jobs", "error", err) |
| 198 | } |
| 199 | |
| 200 | r.mux.Lock() |
| 201 | for _, job := range queued { |
| 202 | r.jobs[job.ID] = job |
| 203 | } |
| 204 | r.mux.Unlock() |
| 205 | for { |
| 206 | select { |
| 207 | case <-r.quit: |
| 208 | return |
| 209 | case <-r.ctx.Done(): |
| 210 | return |
| 211 | case event, ok := <-r.consumer.Watch(): |
| 212 | if !ok { |
| 213 | return |
| 214 | } |
| 215 | r.handleWatcherEvent(event) |
| 216 | } |
| 217 | } |
| 218 | } |
no test coverage detected