141 lines
4.2 KiB
Python
141 lines
4.2 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
MJPEG Video Delay Proxy
|
|
Puffert Frames vom Kamerastream (Port 8081) und liefert sie
|
|
mit einstellbarer Verzögerung auf Port 8083.
|
|
Delay wird aus /tmp/exomy_delay.txt gelesen (Sekunden als float).
|
|
"""
|
|
import collections
|
|
import http.server
|
|
import os
|
|
import socketserver
|
|
import threading
|
|
import time
|
|
import urllib.request
|
|
|
|
DELAY_FILE = '/tmp/exomy_delay.txt'
|
|
SOURCE_URL = 'http://localhost:8081/stream.mjpg'
|
|
PORT = 8083
|
|
BOUNDARY = b'--frame'
|
|
|
|
# Deque: (timestamp_float, jpeg_bytes)
|
|
frame_buffer = collections.deque()
|
|
buffer_lock = threading.Lock()
|
|
latest_frame = None
|
|
latest_lock = threading.Lock()
|
|
|
|
|
|
def read_delay():
|
|
try:
|
|
with open(DELAY_FILE) as f:
|
|
return max(0.0, float(f.read().strip()))
|
|
except Exception:
|
|
return 0.0
|
|
|
|
|
|
def camera_reader():
|
|
global latest_frame
|
|
while True:
|
|
try:
|
|
req = urllib.request.urlopen(SOURCE_URL, timeout=5)
|
|
buf = bytearray()
|
|
while True:
|
|
chunk = req.read(4096)
|
|
if not chunk:
|
|
break
|
|
buf.extend(chunk)
|
|
while True:
|
|
start = buf.find(b'\xff\xd8')
|
|
if start == -1:
|
|
if len(buf) > 1024 * 1024:
|
|
del buf[:-2]
|
|
break
|
|
end = buf.find(b'\xff\xd9', start + 2)
|
|
if end == -1:
|
|
if start > 0:
|
|
del buf[:start]
|
|
break
|
|
jpeg = bytes(buf[start:end + 2])
|
|
del buf[:end + 2]
|
|
ts = time.monotonic()
|
|
with buffer_lock:
|
|
frame_buffer.append((ts, jpeg))
|
|
# Puffer auf 12 Sekunden begrenzen
|
|
cutoff = ts - 12.0
|
|
while frame_buffer and frame_buffer[0][0] < cutoff:
|
|
frame_buffer.popleft()
|
|
with latest_lock:
|
|
latest_frame = jpeg
|
|
except Exception:
|
|
time.sleep(1)
|
|
|
|
|
|
def get_delayed_frame():
|
|
delay = read_delay()
|
|
if delay <= 0:
|
|
with latest_lock:
|
|
return latest_frame
|
|
target_ts = time.monotonic() - delay
|
|
with buffer_lock:
|
|
if not frame_buffer:
|
|
return None
|
|
best = frame_buffer[0][1]
|
|
for ts, jpeg in frame_buffer:
|
|
if ts <= target_ts:
|
|
best = jpeg
|
|
else:
|
|
break
|
|
return best
|
|
|
|
|
|
class ReusableTCPServer(socketserver.ThreadingTCPServer):
|
|
allow_reuse_address = True
|
|
|
|
|
|
class Handler(http.server.BaseHTTPRequestHandler):
|
|
def do_GET(self):
|
|
if self.path != '/stream.mjpg':
|
|
self.send_error(404)
|
|
return
|
|
self.send_response(200)
|
|
self.send_header('Age', '0')
|
|
self.send_header('Cache-Control', 'no-cache, private')
|
|
self.send_header('Pragma', 'no-cache')
|
|
self.send_header('Content-Type', 'multipart/x-mixed-replace; boundary=frame')
|
|
self.end_headers()
|
|
fps = 10
|
|
while True:
|
|
frame = get_delayed_frame()
|
|
if frame is None:
|
|
time.sleep(0.1)
|
|
continue
|
|
try:
|
|
self.wfile.write(BOUNDARY + b'\r\n')
|
|
self.wfile.write(b'Content-Type: image/jpeg\r\n')
|
|
self.wfile.write(f'Content-Length: {len(frame)}\r\n\r\n'.encode())
|
|
self.wfile.write(frame)
|
|
self.wfile.write(b'\r\n')
|
|
self.wfile.flush()
|
|
time.sleep(1.0 / fps)
|
|
except (BrokenPipeError, ConnectionResetError):
|
|
break
|
|
|
|
def log_message(self, fmt, *args):
|
|
return
|
|
|
|
|
|
if __name__ == '__main__':
|
|
# Delay-Datei initialisieren (world-writable damit Admin-API schreiben kann)
|
|
try:
|
|
fd = os.open(DELAY_FILE, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o666)
|
|
os.write(fd, b'0.0')
|
|
os.close(fd)
|
|
except Exception:
|
|
pass
|
|
|
|
threading.Thread(target=camera_reader, daemon=True).start()
|
|
with ReusableTCPServer(('0.0.0.0', PORT), Handler) as server:
|
|
server.daemon_threads = True
|
|
print(f'Video-Delay-Proxy läuft auf Port {PORT}')
|
|
server.serve_forever()
|