()
| 188 | } |
| 189 | |
| 190 | func (c *Controller) retryFailedWorkers() { |
| 191 | c.mux.Lock() |
| 192 | defer c.mux.Unlock() |
| 193 | |
| 194 | if !c.running { |
| 195 | return |
| 196 | } |
| 197 | |
| 198 | c.Entities.Range(func(key, value any) bool { |
| 199 | entityID := key.(string) |
| 200 | worker := value.(*Worker) |
| 201 | |
| 202 | if worker.IsRunning() { |
| 203 | return true |
| 204 | } |
| 205 | |
| 206 | if !c.backoff.ShouldRetry(entityID) { |
| 207 | return true |
| 208 | } |
| 209 | |
| 210 | slog.InfoContext(c.ctx, "retrying failed worker", "entity_id", entityID) |
| 211 | if err := worker.Start(); err != nil { |
| 212 | slog.ErrorContext(c.ctx, "retry failed for worker", "entity_id", entityID, "error", err) |
| 213 | worker.addStatusEvent(fmt.Sprintf("retry failed for worker %s (%s): %s", entityID, worker.Entity.ForgeURL(), err.Error()), params.EventError) |
| 214 | c.backoff.RecordFailure(entityID) |
| 215 | return true |
| 216 | } |
| 217 | |
| 218 | slog.InfoContext(c.ctx, "worker successfully started after retry", "entity_id", entityID) |
| 219 | worker.addStatusEvent(fmt.Sprintf("worker successfully started after retry for entity: %s (%s)", entityID, worker.Entity.ForgeURL()), params.EventInfo) |
| 220 | c.backoff.RecordSuccess(entityID) |
| 221 | return true |
| 222 | }) |
| 223 | } |
| 224 | |
| 225 | // loop drains the watcher consumer channel as fast as possible. |
| 226 | // It exits when the context is cancelled or quit is closed, triggering |
nothing calls this directly
no test coverage detected