| 163 | } |
| 164 | |
| 165 | func (c *Controller) Stop() error { |
| 166 | slog.DebugContext(c.ctx, "stopping entity controller", "entity", c.consumerID) |
| 167 | c.mux.Lock() |
| 168 | defer c.mux.Unlock() |
| 169 | if !c.running { |
| 170 | return nil |
| 171 | } |
| 172 | slog.DebugContext(c.ctx, "stopping entity controller") |
| 173 | |
| 174 | c.Entities.Range(func(key, value any) bool { |
| 175 | entityID := key.(string) |
| 176 | worker := value.(*Worker) |
| 177 | if err := worker.Stop(); err != nil { |
| 178 | slog.ErrorContext(c.ctx, "stopping worker for entity", "entity_id", entityID, "error", err) |
| 179 | } |
| 180 | return true |
| 181 | }) |
| 182 | |
| 183 | c.running = false |
| 184 | close(c.quit) |
| 185 | c.consumer.Close() |
| 186 | slog.DebugContext(c.ctx, "stopped entity controller", "entity", c.consumerID) |
| 187 | return nil |
| 188 | } |
| 189 | |
| 190 | func (c *Controller) retryFailedWorkers() { |
| 191 | c.mux.Lock() |