diff --git a/server/src/services/roomCameraService.js b/server/src/services/roomCameraService.js index 635d6e4a..3c389e4e 100644 --- a/server/src/services/roomCameraService.js +++ b/server/src/services/roomCameraService.js @@ -14,15 +14,15 @@ function normalizeCamera(camera) { logger.warn('Room camera missing id', camera); return null; } - if (!camera.url) { - logger.warn('Room camera missing url', { id, camera }); + if (!camera.url && !camera.streamUrl && !camera.mjpegUrl) { + logger.warn('Room camera missing url/streamUrl', { id, camera }); return null; } return { id: String(id), name: camera.name || camera.id || String(id), description: camera.description || null, - url: camera.url, + url: camera.url || null, streamUrl: camera.streamUrl || camera.mjpegUrl || null, }; } diff --git a/server/src/services/roomCameraSnapshotService.js b/server/src/services/roomCameraSnapshotService.js index 0844543e..24c8ee80 100644 --- a/server/src/services/roomCameraSnapshotService.js +++ b/server/src/services/roomCameraSnapshotService.js @@ -1,13 +1,17 @@ const EventEmitter = require('events'); +const http = require('http'); +const https = require('https'); const logger = require('../globals/logger').child('roomCameraSnapshot'); const { getRoomCameras, roomCameraEvents } = require('./roomCameraService'); const POLL_INTERVAL_MS = 67; const FETCH_TIMEOUT_MS = 2000; +const STREAM_RETRY_MS = 1500; const cameraState = new Map(); // id -> {frame, ts, error, failures, fetching} const events = new EventEmitter(); // frame, status let pollTimer = null; +const streamState = new Map(); // id -> { req, reconnectTimer } function markState(id, updates = {}) { const prev = cameraState.get(id) || {}; @@ -16,6 +20,107 @@ function markState(id, updates = {}) { return next; } +function isMjpegUrl(rawUrl) { + if (!rawUrl) return false; + const lower = String(rawUrl).toLowerCase(); + return lower.endsWith('.mjpg') || lower.endsWith('.mjpeg') || lower.includes('mjpeg'); +} + +function getStreamUrl(camera) { + if (camera.streamUrl) return camera.streamUrl; + if (isMjpegUrl(camera.url)) return camera.url; + return null; +} + +function stopStream(id) { + const entry = streamState.get(id); + if (!entry) return; + if (entry.req) { + entry.req.destroy(); + } + if (entry.reconnectTimer) { + clearTimeout(entry.reconnectTimer); + } + streamState.delete(id); +} + +function scheduleStreamReconnect(camera) { + const id = camera.id; + const entry = streamState.get(id) || {}; + if (entry.reconnectTimer) return; + entry.reconnectTimer = setTimeout(() => { + const current = streamState.get(id) || {}; + current.reconnectTimer = null; + streamState.set(id, current); + startStream(camera); + }, STREAM_RETRY_MS); + streamState.set(id, entry); +} + +function handleStreamError(camera, err) { + const state = cameraState.get(camera.id) || {}; + const failures = (state.failures || 0) + 1; + markState(camera.id, { error: err.message || 'stream error', failures }); + events.emit('status', { id: camera.id, error: err.message || 'stream error' }); + logger.warn('Stream failed', { id: camera.id, err: err.message }); +} + +function startStream(camera) { + const streamUrl = getStreamUrl(camera); + if (!streamUrl) return; + if (streamState.get(camera.id)?.req) return; + const url = new URL(streamUrl); + const client = url.protocol === 'https:' ? https : http; + const req = client.get( + { + hostname: url.hostname, + port: url.port || (url.protocol === 'https:' ? 443 : 80), + path: `${url.pathname}${url.search}`, + headers: { Accept: 'multipart/x-mixed-replace' }, + }, + (res) => { + if (res.statusCode !== 200) { + res.resume(); + handleStreamError(camera, new Error(`HTTP ${res.statusCode}`)); + scheduleStreamReconnect(camera); + return; + } + let buffer = Buffer.alloc(0); + res.on('data', (chunk) => { + buffer = Buffer.concat([buffer, chunk]); + while (true) { + const start = buffer.indexOf(0xffd8); + if (start === -1) { + if (buffer.length > 2 * 1024 * 1024) { + buffer = buffer.slice(-1024 * 1024); + } + break; + } + const end = buffer.indexOf(0xffd9, start + 2); + if (end === -1) break; + const frame = buffer.slice(start, end + 2); + buffer = buffer.slice(end + 2); + const ts = Date.now(); + markState(camera.id, { frame, ts, error: null, failures: 0 }); + events.emit('frame', { id: camera.id, buffer: frame, ts }); + } + }); + res.on('end', () => { + scheduleStreamReconnect(camera); + }); + res.on('error', (err) => { + handleStreamError(camera, err); + scheduleStreamReconnect(camera); + }); + }, + ); + req.on('error', (err) => { + handleStreamError(camera, err); + scheduleStreamReconnect(camera); + }); + streamState.set(camera.id, { req, reconnectTimer: null }); +} + async function fetchSnapshot(camera) { const { id, url } = camera; const state = cameraState.get(id); @@ -49,16 +154,26 @@ function stopAll() { clearInterval(pollTimer); pollTimer = null; } + Array.from(streamState.keys()).forEach((id) => stopStream(id)); cameraState.clear(); } function startAll() { stopAll(); - pollTimer = setInterval(() => { - getRoomCameras().forEach((camera) => fetchSnapshot(camera)); - }, POLL_INTERVAL_MS); - getRoomCameras().forEach((camera) => fetchSnapshot(camera)); - logger.info('Started snapshot polling', { count: getRoomCameras().length }); + const cameras = getRoomCameras(); + const snapshotCameras = cameras.filter((camera) => !getStreamUrl(camera) && camera.url); + cameras.forEach((camera) => startStream(camera)); + if (snapshotCameras.length) { + pollTimer = setInterval(() => { + snapshotCameras.forEach((camera) => fetchSnapshot(camera)); + }, POLL_INTERVAL_MS); + snapshotCameras.forEach((camera) => fetchSnapshot(camera)); + } + logger.info('Started room camera feeds', { + total: cameras.length, + streaming: cameras.length - snapshotCameras.length, + snapshots: snapshotCameras.length, + }); } function getState(id) {