Add the missing Docker events watcher
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
1 parent
085095c46f
commit
a2c907079a
1 file changed
+33
@@ -170,6 +170,39 @@ func (m *DockerManager) Run(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// watchEvents listens to container events (start, stop, removal…) so that a
|
||||||
|
// new container is followed right away instead of at the next periodic sync.
|
||||||
|
func (m *DockerManager) watchEvents(ctx context.Context) {
|
||||||
|
filters := url.QueryEscape(`{"type":["container"]}`)
|
||||||
|
for {
|
||||||
|
resp, err := m.get(ctx, "/events?filters="+filters)
|
||||||
|
if err != nil {
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
m.setError(err)
|
||||||
|
sleepCtx(ctx, 5*time.Second)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
dec := json.NewDecoder(resp.Body)
|
||||||
|
for {
|
||||||
|
var ev struct{ Action string }
|
||||||
|
if err := dec.Decode(&ev); err != nil {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
switch ev.Action {
|
||||||
|
case "create", "start", "die", "destroy", "rename":
|
||||||
|
m.requestSync()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
resp.Body.Close()
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
sleepCtx(ctx, time.Second) // stream cut (proxy timeout, Docker restart): reconnect
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (m *DockerManager) requestSync() {
|
func (m *DockerManager) requestSync() {
|
||||||
select {
|
select {
|
||||||
case m.syncCh <- struct{}{}:
|
case m.syncCh <- struct{}{}:
|
||||||
|
|||||||
Reference in new issue
Block a user