+102
-19
@@ -1,5 +1,6 @@
|
||||
/**
|
||||
* Docker event stream → peer broadcasts.
|
||||
* Auto-reconnects when the daemon restarts or the stream ends.
|
||||
*/
|
||||
import { docker, extractVolumesList } from './docker.js'
|
||||
import { peers } from '../core/peer-registry.js'
|
||||
@@ -7,38 +8,86 @@ import { Pushes } from '../../shared/protocol.js'
|
||||
import logger from '../utils/logger.js'
|
||||
|
||||
let dockerEventStream = null
|
||||
let reconnectTimer = null
|
||||
let stopped = false
|
||||
let reconnectAttempt = 0
|
||||
|
||||
const BASE_DELAY_MS = 2000
|
||||
const MAX_DELAY_MS = 30_000
|
||||
|
||||
export async function startDockerEventStream() {
|
||||
stopped = false
|
||||
await openEventStream()
|
||||
}
|
||||
|
||||
async function openEventStream() {
|
||||
if (stopped) return
|
||||
if (reconnectTimer) {
|
||||
clearTimeout(reconnectTimer)
|
||||
reconnectTimer = null
|
||||
}
|
||||
|
||||
try {
|
||||
const stream = await new Promise((resolve, reject) => {
|
||||
docker.getEvents({}, (err, s) => (err ? reject(err) : resolve(s)))
|
||||
})
|
||||
dockerEventStream = stream
|
||||
reconnectAttempt = 0
|
||||
logger.info('Docker event stream connected')
|
||||
|
||||
stream.on('data', async (chunk) => {
|
||||
try {
|
||||
const event = JSON.parse(chunk.toString())
|
||||
if (event.status === 'undefined') return
|
||||
logger.info('Docker event', {
|
||||
status: event.status,
|
||||
id: event.id,
|
||||
type: event.Type,
|
||||
})
|
||||
const lines = chunk
|
||||
.toString()
|
||||
.split('\n')
|
||||
.map((l) => l.trim())
|
||||
.filter(Boolean)
|
||||
|
||||
if (event.Type === 'container') {
|
||||
const containers = await docker.listContainers({ all: true })
|
||||
peers.broadcast(Pushes.containers, { type: 'containers', data: containers })
|
||||
}
|
||||
for (const line of lines) {
|
||||
let event
|
||||
try {
|
||||
event = JSON.parse(line)
|
||||
} catch {
|
||||
continue
|
||||
}
|
||||
if (event.status === 'undefined') continue
|
||||
|
||||
if (event.Type === 'volume' && (event.Action === 'create' || event.Action === 'destroy')) {
|
||||
const volumesResult = await docker.listVolumes()
|
||||
const volumesList = extractVolumesList(volumesResult)
|
||||
peers.broadcast(Pushes.volumes, {
|
||||
type: 'volumes',
|
||||
data: volumesList,
|
||||
success: true,
|
||||
volumes: volumesList,
|
||||
logger.info('Docker event', {
|
||||
status: event.status,
|
||||
id: event.id,
|
||||
type: event.Type,
|
||||
action: event.Action,
|
||||
})
|
||||
|
||||
peers.broadcast(Pushes.dockerEvent, {
|
||||
type: 'dockerEvent',
|
||||
data: event,
|
||||
})
|
||||
|
||||
if (event.Type === 'container') {
|
||||
const containers = await docker.listContainers({ all: true })
|
||||
peers.broadcast(Pushes.containers, { type: 'containers', data: containers })
|
||||
}
|
||||
|
||||
if (event.Type === 'image') {
|
||||
try {
|
||||
const images = await docker.listImages({ all: true })
|
||||
peers.broadcast(Pushes.images, { type: 'images', data: images })
|
||||
} catch (e) {
|
||||
logger.debug('image list on event failed', { error: e.message })
|
||||
}
|
||||
}
|
||||
|
||||
if (event.Type === 'volume' && (event.Action === 'create' || event.Action === 'destroy')) {
|
||||
const volumesResult = await docker.listVolumes()
|
||||
const volumesList = extractVolumesList(volumesResult)
|
||||
peers.broadcast(Pushes.volumes, {
|
||||
type: 'volumes',
|
||||
data: volumesList,
|
||||
success: true,
|
||||
volumes: volumesList,
|
||||
})
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
logger.error('Failed to process Docker event', { error: err.message })
|
||||
@@ -47,16 +96,50 @@ export async function startDockerEventStream() {
|
||||
|
||||
stream.on('error', (err) => {
|
||||
logger.error('Docker event stream error', { error: err.message })
|
||||
scheduleReconnect()
|
||||
})
|
||||
stream.on('end', () => {
|
||||
dockerEventStream = null
|
||||
logger.warn('Docker event stream ended')
|
||||
scheduleReconnect()
|
||||
})
|
||||
} catch (err) {
|
||||
logger.error('Failed to start Docker event stream', { error: err.message })
|
||||
scheduleReconnect()
|
||||
}
|
||||
}
|
||||
|
||||
function scheduleReconnect() {
|
||||
if (stopped) return
|
||||
if (reconnectTimer) return
|
||||
|
||||
if (dockerEventStream) {
|
||||
try {
|
||||
dockerEventStream.destroy()
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
dockerEventStream = null
|
||||
}
|
||||
|
||||
const delay = Math.min(MAX_DELAY_MS, BASE_DELAY_MS * 2 ** reconnectAttempt)
|
||||
reconnectAttempt += 1
|
||||
logger.info('Reconnecting Docker event stream', { delayMs: delay, attempt: reconnectAttempt })
|
||||
reconnectTimer = setTimeout(() => {
|
||||
reconnectTimer = null
|
||||
openEventStream().catch((err) => {
|
||||
logger.error('Event stream reconnect failed', { error: err.message })
|
||||
scheduleReconnect()
|
||||
})
|
||||
}, delay)
|
||||
}
|
||||
|
||||
export function stopDockerEventStream() {
|
||||
stopped = true
|
||||
if (reconnectTimer) {
|
||||
clearTimeout(reconnectTimer)
|
||||
reconnectTimer = null
|
||||
}
|
||||
if (dockerEventStream) {
|
||||
try {
|
||||
dockerEventStream.destroy()
|
||||
|
||||
Reference in New Issue
Block a user