Files
2026-09-13 20:24:20 +01:00

1586 lines
58 KiB
Python

#!/usr/bin/env python3
"""
Camera Station — a tiny self-hosted front end for MJPEG and RTSP cameras.
Why a server instead of just opening the HTML file?
1. Browsers "taint" a <canvas> that has drawn an image from another origin,
and a tainted canvas cannot be read back — which would break both photo
capture and recording. Serving the page and the video from one origin
keeps the canvas readable.
2. No browser can play RTSP. RTSP cameras are decoded by ffmpeg here and
handed to the page as MJPEG, so every camera behaves the same way.
Camera credentials stay in this process and are never sent to the page.
Usage:
python3 server.py
python3 server.py --camera "Front door=http://192.168.69.79:8080"
python3 server.py --camera "Drive=rtsp://user:pass@192.168.69.128:554/cam"
python3 server.py --host 127.0.0.1 # only this machine can connect
Then open http://localhost:8000
Standard library only. ffmpeg is needed for RTSP cameras, nothing else.
"""
import argparse
import json
import os
import re
import shutil
import socket
import subprocess
import sys
import threading
import time
import uuid
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from queue import Queue, Empty
from urllib.error import URLError, HTTPError
from urllib.parse import urlparse, parse_qs, quote, unquote
from urllib.request import urlopen
HERE = os.path.dirname(os.path.abspath(__file__))
PAGE = os.path.join(HERE, "index.html")
CONFIG_PAGE = os.path.join(HERE, "config.html")
LIBRARY_PAGE = os.path.join(HERE, "library.html")
# Everything the service produces lives under one directory, so a container can
# mount a single volume and keep recordings, photos and settings together.
DATA_DIR = os.environ.get("CAMERA_STATION_DATA") or os.path.join(HERE, "data")
CONFIG_FILE = os.path.join(DATA_DIR, "cameras.json")
LEGACY_CONFIG = os.path.join(HERE, "cameras.json")
RECORD_DIR = os.path.join(DATA_DIR, "recordings")
PHOTO_DIR = os.path.join(DATA_DIR, "photos")
STORAGE_LIMIT = 10 * 1024 ** 3 # 10 GB of footage, oldest deleted first
SEGMENT_SECONDS = 600 # a new file every 10 minutes, at a keyframe
SWEEP_SECONDS = 20 # how often storage is checked
def ensure_dirs():
for path in (DATA_DIR, RECORD_DIR, PHOTO_DIR):
os.makedirs(path, exist_ok=True)
def slug(text):
out = re.sub(r"[^A-Za-z0-9]+", "-", (text or "").strip()).strip("-")
return out or "camera"
DEFAULT_CAMERAS = [
("Camera 1", "http://192.168.69.79:8080"),
("Camera 2", "rtsp://admin:Admin100pass!@192.168.69.128:554"
"/cam/realmonitor?channel=1&subtype=0"),
]
BOUNDARY = "camerastation"
# --------------------------------------------------------------------------
# helpers
# --------------------------------------------------------------------------
def recording_state(cam):
rec = getattr(cam, "recorder", None)
status = rec.status() if rec else {}
return {"active": bool(getattr(cam, "active", False)),
"recording": bool(status.get("running")),
"since": status.get("since"),
"recordError": status.get("error")}
def safe_host(url):
"""The address with any user:password stripped, for display in the page."""
try:
p = urlparse(url)
host = p.hostname or ""
if p.port:
host += ":%d" % p.port
return host or "camera"
except ValueError:
return "camera"
def ffmpeg_available():
return shutil.which("ffmpeg") is not None
def ffmpeg_timeout_flag():
"""ffmpeg 6 renamed the RTSP socket timeout from -stimeout to -timeout."""
try:
out = subprocess.run(["ffmpeg", "-version"], capture_output=True,
text=True, timeout=5).stdout
m = re.search(r"ffmpeg version n?(\d+)", out)
if m and int(m.group(1)) >= 6:
return "-timeout"
except Exception:
pass
return "-stimeout"
# The config page has to show existing cameras without handing their passwords
# to a browser, so a saved password is displayed as bullets and only travels in
# one direction: you can set a new one, but you can never read the old one back.
KEPT_SECRET = "\u2022" * 8
URL_PARTS = re.compile(r"^([a-zA-Z][a-zA-Z0-9+.\-]*://)([^/@]*@)?(.*)$", re.S)
def redact_url(url):
"""Replace any password in a URL with bullets."""
m = URL_PARTS.match(url or "")
if not m:
return url
scheme, userinfo, rest = m.group(1), m.group(2) or "", m.group(3)
creds = userinfo[:-1]
if ":" not in creds:
return url
return scheme + creds.split(":", 1)[0] + ":" + KEPT_SECRET + "@" + rest
def restore_url(url, previous):
"""Put the stored password back if the page sent the bullets untouched."""
m = URL_PARTS.match(url or "")
if not m:
return url
scheme, userinfo, rest = m.group(1), m.group(2) or "", m.group(3)
creds = userinfo[:-1]
if ":" not in creds:
return url
user, password = creds.split(":", 1)
if password != KEPT_SECRET:
return url # a new password was typed; use it
prev = URL_PARTS.match(previous or "")
if prev and prev.group(2) and ":" in prev.group(2)[:-1]:
old = prev.group(2)[:-1].split(":", 1)[1]
return scheme + user + ":" + old + "@" + rest
return scheme + user + "@" + rest # nothing stored to restore
def valid_url(url):
url = (url or "").strip()
return url.startswith(("http://", "https://", "rtsp://", "rtsps://"))
def load_config():
"""Cameras saved by the config page, or None if there are none yet."""
data = None
for path in (CONFIG_FILE, LEGACY_CONFIG):
try:
with open(path, "r", encoding="utf-8") as fh:
data = json.load(fh)
break
except (OSError, ValueError):
continue
if data is None:
return None
entries = data.get("cameras") if isinstance(data, dict) else data
if not isinstance(entries, list):
return None
out = []
for e in entries:
if isinstance(e, dict) and valid_url(e.get("url")):
out.append({"id": str(e.get("id") or new_id()),
"name": str(e.get("name") or "Camera").strip() or "Camera",
"url": e["url"].strip(),
"active": bool(e.get("active"))})
return out or None
def save_config(entries):
"""Write the camera list, via a temporary file so a crash cannot truncate it."""
ensure_dirs()
tmp = CONFIG_FILE + ".tmp"
with open(tmp, "w", encoding="utf-8") as fh:
json.dump({"cameras": entries}, fh, indent=2)
fh.flush()
os.fsync(fh.fileno())
os.replace(tmp, CONFIG_FILE)
def new_id():
return "cam" + uuid.uuid4().hex[:6]
def jpeg_size(data):
"""Width and height from JPEG headers, without pulling in an image library."""
i, n = 2, len(data)
while i + 9 < n:
if data[i] != 0xFF:
i += 1
continue
marker = data[i + 1]
if marker in (0xD8, 0xD9) or 0xD0 <= marker <= 0xD7:
i += 2
continue
if marker in (0xC0, 0xC1, 0xC2, 0xC3, 0xC5, 0xC6,
0xC7, 0xC9, 0xCA, 0xCB, 0xCD, 0xCE, 0xCF):
return (int.from_bytes(data[i + 7:i + 9], "big"),
int.from_bytes(data[i + 5:i + 7], "big"))
i += 2 + int.from_bytes(data[i + 2:i + 4], "big")
return None
# --------------------------------------------------------------------------
# cameras
# --------------------------------------------------------------------------
class MjpegCamera:
"""An MJPG-Streamer style feed. Relayed byte for byte — nothing re-encoded."""
kind = "mjpeg"
def __init__(self, cid, name, url):
self.id = cid
self.name = name
self.url = url.rstrip("/")
self.available = True
self.reason = None
def describe(self):
return dict(recording_state(self),
id=self.id, name=self.name, kind=self.kind,
host=safe_host(self.url), available=self.available,
reason=self.reason)
def open_upstream(self, action):
return urlopen(self.url + "/?action=" + action, timeout=10)
class RtspCamera:
"""
An RTSP camera, decoded by ffmpeg and re-published as MJPEG.
One ffmpeg process per camera, shared by every viewer, started on the first
subscriber and stopped shortly after the last one leaves. Decoding an H.264
stream twice because two tabs are open would be wasteful, and most cameras
limit how many RTSP sessions they will serve at once.
"""
kind = "rtsp"
IDLE_GRACE = 8.0 # seconds to keep ffmpeg alive after the last viewer
# ffmpeg is chatty about things that are not faults. Two kinds dominate:
#
# * Decoder complaints while it waits for the first keyframe. Joining an
# H.264/HEVC stream mid-GOP means the first frames reference pictures
# that were sent before we connected. This clears itself the moment an
# IDR arrives, usually within a second or two.
# * Broken pipe, on the way down. When a camera is released, its stdout is
# closed, and ffmpeg reports every stage of noticing that.
#
# Neither says anything useful about whether the camera works, and treating
# them as errors puts nonsense like "Error closing file" on screen as the
# reason a camera is unavailable.
DECODER_NOISE = (
"could not find ref with poc",
"error constructing the frame rps",
"missing reference picture",
"reference picture missing",
"co located pocs unavailable",
"decode_slice_header error",
"no frame!",
"mmco: unref short failure",
"illegal short term buffer state",
"error while decoding mb",
"concealing",
"corrupt decoded frame",
"non-existing pps",
"sps unavailable",
"deprecated pixel format",
"last message repeated",
)
SHUTDOWN_NOISE = (
"broken pipe",
"error muxing a packet",
"task finished with error code",
"terminating thread with return code",
"error writing trailer",
"error closing file",
"error submitting a packet to the muxer",
)
@classmethod
def is_noise(cls, text):
low = text.lower()
return any(p in low for p in cls.DECODER_NOISE + cls.SHUTDOWN_NOISE)
def __init__(self, cid, name, url, quality=3, fps=None, size=None, verbose=False):
self.id = cid
self.name = name
self.url = url
self.quality = quality
self.fps = fps
self.size = size
self.verbose = verbose
self.noise_count = 0
self.available = ffmpeg_available()
self.reason = None if self.available else \
"ffmpeg is not installed, so this RTSP camera cannot be decoded."
self._lock = threading.RLock()
self._subs = set()
self._proc = None
self._stopper = None
self._latest = None
self._latest_at = 0
self._first_frame = threading.Event()
self._timeout_flag = ffmpeg_timeout_flag() if self.available else None
self._last_error = None
# ---- description ----
def describe(self):
return dict(recording_state(self),
id=self.id, name=self.name, kind=self.kind,
host=safe_host(self.url), available=self.available,
reason=self.reason or self._last_error)
# ---- process control ----
def input_args(self, timeout_flag):
"""Overridable so tests can substitute a synthetic source."""
return [timeout_flag, "5000000", "-rtsp_transport", "tcp", "-i", self.url]
def _command(self, timeout_flag):
args = ["ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "error",
"-fflags", "nobuffer+discardcorrupt"]
args += self.input_args(timeout_flag)
args += ["-an"]
if self.fps:
args += ["-r", str(self.fps)]
if self.size:
args += ["-s", self.size]
args += ["-f", "mpjpeg", "-q:v", str(self.quality), "-"]
return args
def _spawn(self):
"""Start ffmpeg, retrying once if the timeout flag is the wrong spelling."""
alt = "-stimeout" if self._timeout_flag == "-timeout" else "-timeout"
proc = None
for flag in (self._timeout_flag, alt):
proc = subprocess.Popen(self._command(flag),
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
bufsize=0)
time.sleep(0.3)
if proc.poll() is None:
self._timeout_flag = flag
return proc
err = (proc.stderr.read() or b"").decode("utf-8", "replace")
if "Unrecognized option" not in err and "Option not found" not in err:
lines = [l for l in err.strip().splitlines() if l.strip()]
self._last_error = lines[-1] if lines else "ffmpeg exited immediately."
return proc
return proc
def _ensure_running(self):
with self._lock:
if self._stopper:
self._stopper.cancel()
self._stopper = None
if self._proc and self._proc.poll() is None:
return
self._first_frame.clear()
self._proc = self._spawn()
threading.Thread(target=self._pump, args=(self._proc,), daemon=True).start()
threading.Thread(target=self._drain_errors, args=(self._proc,), daemon=True).start()
def _schedule_stop(self):
with self._lock:
if self._stopper:
self._stopper.cancel()
self._stopper = threading.Timer(self.IDLE_GRACE, self._stop_if_idle)
self._stopper.daemon = True
self._stopper.start()
def _stop_if_idle(self):
with self._lock:
if self._subs:
return
proc, self._proc = self._proc, None
self._latest = None
if proc and proc.poll() is None:
proc.cs_stopping = True
proc.terminate()
try:
proc.wait(timeout=3)
except subprocess.TimeoutExpired:
proc.kill()
def shutdown(self):
"""Stop for good, even if someone is still watching (config changed)."""
with self._lock:
self._subs.clear()
self._stop_if_idle()
def _drain_errors(self, proc):
try:
for line in iter(proc.stderr.readline, b""):
text = line.decode("utf-8", "replace").strip()
if not text:
continue
if getattr(proc, "cs_stopping", False):
# We asked this process to stop. Whatever it says on the way
# out is expected, so it is not worth alarming anyone with.
continue
if self.is_noise(text):
self.noise_count += 1
if self.verbose:
sys.stderr.write(" [%s] %s\n" % (self.id, text))
continue
self._last_error = text
sys.stderr.write(" [%s] %s\n" % (self.id, text))
except Exception:
pass
# ---- mpjpeg parsing ----
def _pump(self, proc):
"""
Read ffmpeg's multipart output and hand whole JPEG frames to viewers.
Frames are split on the Content-length header rather than by hunting for
JPEG end markers, because those bytes can occur inside compressed data.
"""
out = proc.stdout
try:
while True:
line = out.readline()
if not line:
break
low = line.strip().lower()
if not low.startswith(b"content-length"):
continue
try:
size = int(low.split(b":", 1)[1])
except (IndexError, ValueError):
continue
while True: # skip to the blank separator
sep = out.readline()
if not sep or sep in (b"\r\n", b"\n"):
break
frame = out.read(size)
if not frame or len(frame) < size:
break
self._publish(frame)
except Exception as err:
self._last_error = str(err)
finally:
try:
out.close()
except Exception:
pass
def _publish(self, frame):
# Pictures are arriving, so any earlier complaint is history.
if self._last_error is not None:
self._last_error = None
self._latest = frame
self._latest_at = time.time()
self._first_frame.set()
for q in list(self._subs):
if q.full():
try:
q.get_nowait() # drop the stale frame, never the fresh one
except Empty:
pass
try:
q.put_nowait(frame)
except Exception:
pass
# ---- viewer API ----
def subscribe(self):
q = Queue(maxsize=1)
with self._lock:
self._subs.add(q)
self._ensure_running()
return q
def unsubscribe(self, q):
with self._lock:
self._subs.discard(q)
idle = not self._subs
if idle:
self._schedule_stop()
def snapshot(self, wait=10.0):
"""Latest decoded frame, starting the camera briefly if nobody is watching."""
if self._latest is not None and time.time() - self._latest_at < 5:
return self._latest
q = self.subscribe()
try:
self._first_frame.wait(wait)
return self._latest
finally:
self.unsubscribe(q)
def viewers(self):
with self._lock:
return len(self._subs)
def is_running(self):
with self._lock:
return bool(self._proc and self._proc.poll() is None)
# --------------------------------------------------------------------------
# continuous recording
# --------------------------------------------------------------------------
class Recorder:
"""
Writes a camera to disk continuously while it is marked active.
The stream is copied, not re-encoded: what lands on disk is exactly the
bitstream the camera sent, at no quality cost and almost no CPU. That also
means an active camera costs nearly nothing until somebody actually watches
it, because only watching requires decoding.
Footage is cut into timestamped segments so old material can be deleted a
piece at a time. Matroska is used rather than MP4 because a segment that is
still being written stays readable if the service is stopped or the machine
loses power; an unfinalised MP4 would not.
"""
RESTART_DELAY = 5.0
def __init__(self, cam):
self.cam = cam
self.proc = None
self.started_at = 0
self.error = None
self.folder_name = slug(cam.name)
self._lock = threading.RLock()
self._last_attempt = 0
@property
def folder(self):
return os.path.join(RECORD_DIR, self.folder_name)
def is_running(self):
with self._lock:
return bool(self.proc and self.proc.poll() is None)
def command(self):
args = ["ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "error"]
if self.cam.kind == "rtsp":
args += [ffmpeg_timeout_flag(), "5000000",
"-rtsp_transport", "tcp", "-i", self.cam.url]
else:
# An HTTP MJPEG feed carries no timestamps of its own, so take them
# from the clock; without this the file plays back at the wrong speed.
args += ["-use_wallclock_as_timestamps", "1",
"-f", "mjpeg", "-i", self.cam.url + "/?action=stream"]
args += ["-an", "-c", "copy",
"-f", "segment", "-strftime", "1",
"-segment_time", str(SEGMENT_SECONDS),
"-segment_format", "matroska",
"-reset_timestamps", "1",
os.path.join(self.folder, "%Y%m%d-%H%M%S.mkv")]
return args
def start(self):
with self._lock:
if self.is_running():
return
if time.time() - self._last_attempt < self.RESTART_DELAY:
return # a camera that keeps failing is not hammered
self._last_attempt = time.time()
if not ffmpeg_available():
self.error = "ffmpeg is not installed, so nothing can be recorded."
return
os.makedirs(self.folder, exist_ok=True)
try:
self.proc = subprocess.Popen(self.command(), stdout=subprocess.DEVNULL,
stderr=subprocess.PIPE, bufsize=0)
except OSError as err:
self.error = str(err)
return
self.started_at = time.time()
threading.Thread(target=self._watch, args=(self.proc,), daemon=True).start()
def stop(self):
with self._lock:
proc, self.proc = self.proc, None
self.started_at = 0
if proc and proc.poll() is None:
proc.cs_stopping = True
proc.terminate() # ffmpeg closes the current segment cleanly
try:
proc.wait(timeout=8)
except subprocess.TimeoutExpired:
proc.kill()
def rename(self, name):
"""Follow a camera being renamed, carrying its footage across."""
wanted = slug(name)
with self._lock:
if wanted == self.folder_name:
return
running = self.is_running()
if running:
self.stop()
old, new = self.folder, os.path.join(RECORD_DIR, wanted)
try:
if os.path.isdir(old) and not os.path.exists(new):
os.rename(old, new)
except OSError:
pass # leave the old footage where it is
self.folder_name = wanted
if running:
self._last_attempt = 0
self.start()
def _watch(self, proc):
for line in iter(proc.stderr.readline, b""):
text = line.decode("utf-8", "replace").strip()
if not text or getattr(proc, "cs_stopping", False):
continue
if RtspCamera.is_noise(text):
continue
self.error = text
sys.stderr.write(" [rec %s] %s\n" % (self.cam.name, text))
code = proc.wait()
if not getattr(proc, "cs_stopping", False) and code not in (0, None):
sys.stderr.write(" [rec %s] recorder stopped (code %s); will retry\n"
% (self.cam.name, code))
def status(self):
return {"running": self.is_running(),
"since": self.started_at or None,
"folder": self.folder_name,
"error": None if self.is_running() else self.error}
def recording_files():
"""Every stored segment, oldest first."""
found = []
for root, _dirs, names in os.walk(RECORD_DIR):
for name in names:
if not name.endswith(".mkv"):
continue
path = os.path.join(root, name)
try:
st = os.stat(path)
except OSError:
continue
found.append({"path": path,
"folder": os.path.basename(root),
"name": name,
"size": st.st_size,
"mtime": st.st_mtime})
found.sort(key=lambda f: f["mtime"])
return found
def enforce_storage_limit(limit=None):
"""
Delete the oldest footage until the total is back under the limit.
The newest file in each camera's folder is left alone even if that means
briefly exceeding the limit: it is the one being written to right now.
"""
limit = STORAGE_LIMIT if limit is None else limit
files = recording_files()
total = sum(f["size"] for f in files)
if total <= limit:
return {"used": total, "deleted": 0, "freed": 0}
newest_per_folder = {}
for f in files:
newest_per_folder[f["folder"]] = f["path"] # sorted oldest first
protected = set(newest_per_folder.values())
deleted, freed = 0, 0
for f in files:
if total <= limit:
break
if f["path"] in protected:
continue
try:
os.remove(f["path"])
except OSError:
continue
total -= f["size"]
freed += f["size"]
deleted += 1
return {"used": total, "deleted": deleted, "freed": freed}
def storage_summary():
files = recording_files()
folders = {}
for f in files:
entry = folders.setdefault(f["folder"], {"folder": f["folder"], "bytes": 0,
"segments": 0, "oldest": None,
"newest": None})
entry["bytes"] += f["size"]
entry["segments"] += 1
entry["oldest"] = entry["oldest"] or f["mtime"]
entry["newest"] = f["mtime"]
photos = 0
photo_count = 0
try:
for name in os.listdir(PHOTO_DIR):
try:
photos += os.path.getsize(os.path.join(PHOTO_DIR, name))
photo_count += 1
except OSError:
pass
except OSError:
pass
recent = [{"folder": f["folder"], "name": f["name"], "size": f["size"],
"mtime": f["mtime"]} for f in files[-12:]][::-1]
return {"limit": STORAGE_LIMIT,
"used": sum(f["size"] for f in files),
"segments": len(files),
"photoBytes": photos,
"photoCount": photo_count,
"cameras": sorted(folders.values(), key=lambda d: -d["bytes"]),
"recent": recent}
PHOTO_NAME = re.compile(r"^(\d{8})_(\d{6})_(.+?)(?:-(\d+))?\.(png|jpe?g)$", re.I)
SEGMENT_NAME = re.compile(r"^(\d{8})-(\d{6})\.mkv$")
def pretty_time(day, hhmmss):
return "%s-%s-%s %s:%s:%s" % (day[0:4], day[4:6], day[6:8],
hhmmss[0:2], hhmmss[2:4], hhmmss[4:6])
def files_being_written():
"""Segments an active recorder is writing into right now."""
busy = set()
by_folder = {}
for f in recording_files():
by_folder[f["folder"]] = f["path"] # oldest first, so this ends newest
for cam in CAMERAS:
rec = getattr(cam, "recorder", None)
if rec and rec.is_running():
path = by_folder.get(rec.folder_name)
if path:
busy.add(path)
return busy
def recording_entries():
"""Every stored segment, newest first, with the time taken from its name."""
out = []
live = files_being_written()
for f in recording_files():
m = SEGMENT_NAME.match(f["name"])
day = m.group(1) if m else time.strftime("%Y%m%d", time.localtime(f["mtime"]))
clock = m.group(2) if m else time.strftime("%H%M%S", time.localtime(f["mtime"]))
out.append({"kind": "recording", "name": f["name"], "camera": f["folder"],
"day": day, "when": pretty_time(day, clock),
"size": f["size"], "mtime": f["mtime"],
# Still being written to. It can be downloaded, but it is
# only as complete as it was the moment you asked.
"live": f["path"] in live,
"url": "/recordings/%s/%s" % (quote(f["folder"]), quote(f["name"]))})
out.reverse()
return out
def photo_entries():
"""Photos, newest first. The camera and time come out of the filename."""
out = []
try:
names = os.listdir(PHOTO_DIR)
except OSError:
return out
for name in names:
path = os.path.join(PHOTO_DIR, name)
try:
st = os.stat(path)
except OSError:
continue
if not os.path.isfile(path):
continue
m = PHOTO_NAME.match(name)
if m:
day, clock, camera = m.group(1), m.group(2), m.group(3)
when = pretty_time(day, clock)
else:
day = time.strftime("%Y%m%d", time.localtime(st.st_mtime))
clock = time.strftime("%H%M%S", time.localtime(st.st_mtime))
camera, when = "unknown", pretty_time(day, clock)
out.append({"kind": "photo", "name": name, "camera": camera,
"day": day, "when": when, "size": st.st_size,
"mtime": st.st_mtime, "url": "/photos/" + quote(name)})
out.sort(key=lambda e: (e["day"], e["name"]), reverse=True)
return out
def browse(kind, camera=None, day=None, offset=0, limit=300):
entries = photo_entries() if kind == "photos" else recording_entries()
cameras = sorted({e["camera"] for e in entries})
days = sorted({e["day"] for e in entries}, reverse=True)
if camera:
entries = [e for e in entries if e["camera"] == camera]
if day:
entries = [e for e in entries if e["day"] == day]
total = len(entries)
shown = entries[offset:offset + limit]
return {"kind": kind, "cameras": cameras, "days": days,
"total": total, "offset": offset,
"bytes": sum(e["size"] for e in entries),
"files": shown}
def delete_files(kind, names):
"""
Remove stored files.
Paths are rebuilt from the root rather than trusted, and the segment an
active recorder is writing into is refused: deleting it would not stop the
recording, it would just lose the footage being captured.
"""
root = PHOTO_DIR if kind == "photos" else RECORD_DIR
busy = files_being_written() if kind != "photos" else set()
deleted, freed, skipped = 0, 0, []
for rel in names:
target = os.path.normpath(os.path.join(root, rel.lstrip("/")))
if not target.startswith(os.path.normpath(root) + os.sep):
skipped.append(rel)
continue
if target in busy:
skipped.append(rel)
continue
try:
size = os.path.getsize(target)
os.remove(target)
except OSError:
skipped.append(rel)
continue
deleted += 1
freed += size
return {"deleted": deleted, "freed": freed, "skipped": skipped}
def supervise():
"""
Keep active cameras recording and storage inside its limit.
A camera that drops off the network takes its recorder down with it, so the
recorders are restarted from here rather than only when something asks.
"""
while True:
try:
for cam in list(CAMERAS):
rec = getattr(cam, "recorder", None)
if rec is None:
continue
if getattr(cam, "active", False) and cam.available:
if not rec.is_running():
rec.start()
elif rec.is_running():
rec.stop()
enforce_storage_limit()
except Exception as err: # never let the loop die
sys.stderr.write(" supervisor: %s\n" % err)
time.sleep(SWEEP_SECONDS)
# --------------------------------------------------------------------------
# http
# --------------------------------------------------------------------------
CAMERAS = []
LOCKED = False # --lock-config disables editing
class Options:
"""Decoding settings applied to every RTSP camera."""
rtsp_quality = 3
rtsp_fps = None
rtsp_size = None
verbose = False
OPTS = Options()
def make_camera(cid, name, url, opts=None):
# Read settings defensively: this is called with the parsed command line,
# with OPTS, and from tests, and a missing field should fall back to the
# default rather than stop a camera from being created at all.
opts = opts or OPTS
if url.startswith(("rtsp://", "rtsps://")):
cam = RtspCamera(cid, name, url,
quality=getattr(opts, "rtsp_quality", 3),
fps=getattr(opts, "rtsp_fps", None),
size=getattr(opts, "rtsp_size", None),
verbose=getattr(opts, "verbose", False))
else:
cam = MjpegCamera(cid, name, url)
cam.active = False
cam.recorder = Recorder(cam)
return cam
def apply_config(entries):
"""
Swap in a new camera list without disturbing cameras that did not change.
Rebuilding everything on every save would drop the picture on cameras the
edit never touched, and restart an ffmpeg decode for no reason, so cameras
whose address is unchanged are carried over as they are.
"""
previous = {c.id: c for c in CAMERAS}
rebuilt, kept = [], set()
for entry in entries:
cid = entry.get("id") or new_id()
name = entry["name"]
url = entry["url"]
old = previous.get(cid)
# MjpegCamera stores its address with any trailing slash removed, so
# compare against both spellings before deciding it has changed.
unchanged = old is not None and getattr(old, "url", None) in (url, url.rstrip("/"))
if unchanged:
cam = old
if cam.name != name:
cam.name = name
cam.recorder.rename(name) # take the footage folder with it
else:
cam = make_camera(cid, name, url)
cam.active = bool(entry.get("active"))
rebuilt.append(cam)
kept.add(id(cam))
for cam in CAMERAS:
if id(cam) in kept:
continue
if getattr(cam, "recorder", None):
cam.recorder.stop() # dropped from the config, so stop recording it
if isinstance(cam, RtspCamera):
cam.shutdown() # released, so stop decoding it
CAMERAS[:] = rebuilt
return CAMERAS
def config_entries():
"""The current cameras, safe to hand to the config page."""
return [dict(recording_state(c), id=c.id, name=c.name, url=redact_url(c.url))
for c in CAMERAS]
def persist():
"""Write the current cameras back to disk, keeping real passwords."""
save_config([{"id": c.id, "name": c.name, "url": c.url,
"active": bool(getattr(c, "active", False))} for c in CAMERAS])
def probe(url, timeout=15):
"""
Try to pull a single frame, so a camera can be checked before it is saved.
"""
url = (url or "").strip()
if not valid_url(url):
return {"ok": False, "message": "Address must start with http:// or rtsp://"}
if url.startswith(("rtsp://", "rtsps://")):
if not ffmpeg_available():
return {"ok": False, "kind": "rtsp",
"message": "ffmpeg is not installed, so RTSP cameras cannot be used."}
args = ["ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "error",
ffmpeg_timeout_flag(), str(int(timeout * 1000000)),
"-rtsp_transport", "tcp", "-i", url,
"-frames:v", "1", "-f", "mjpeg", "-q:v", "3", "-"]
try:
done = subprocess.run(args, capture_output=True, timeout=timeout + 5)
except subprocess.TimeoutExpired:
return {"ok": False, "kind": "rtsp",
"message": "No response within %ds." % timeout}
if done.stdout[:3] == b"\xff\xd8\xff":
size = jpeg_size(done.stdout)
return {"ok": True, "kind": "rtsp",
"width": size[0] if size else None,
"height": size[1] if size else None}
err = [l for l in done.stderr.decode("utf-8", "replace").splitlines()
if l.strip() and not RtspCamera.is_noise(l)]
return {"ok": False, "kind": "rtsp",
"message": err[-1] if err else "No picture came back."}
base = url.rstrip("/")
for candidate in (base + "/?action=snapshot", base):
try:
with urlopen(candidate, timeout=timeout) as r:
head = r.read(600000)
except (URLError, HTTPError, socket.timeout, OSError) as err:
message = str(getattr(err, "reason", err))
continue
start = head.find(b"\xff\xd8\xff")
if start >= 0:
size = jpeg_size(head[start:])
return {"ok": True, "kind": "mjpeg",
"width": size[0] if size else None,
"height": size[1] if size else None}
message = "Reachable, but no JPEG was found in the reply."
return {"ok": False, "kind": "mjpeg", "message": message}
def find_camera(cid):
for cam in CAMERAS:
if cam.id == cid:
return cam
return None
class Station(ThreadingHTTPServer):
"""
A threading server that stays quiet about viewers leaving.
Long-lived MJPEG responses are aborted every time someone switches camera,
reloads, or closes a tab. That is routine, not a fault, so it should not
print a stack trace. Anything else still gets reported normally.
"""
daemon_threads = True
allow_reuse_address = True
def handle_error(self, request, client_address):
if isinstance(sys.exc_info()[1], (ConnectionError, socket.timeout)):
return
super().handle_error(request, client_address)
class Handler(BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"
server_version = "CameraStation/2.0"
def do_GET(self):
parsed = urlparse(self.path)
path = parsed.path
query = parse_qs(parsed.query)
if path in ("/", "/index.html"):
self.serve_page()
return
if path == "/cameras":
self.serve_json([c.describe() for c in CAMERAS])
return
if path in ("/config", "/config.html"):
self.serve_file(CONFIG_PAGE, "text/html; charset=utf-8")
return
if path == "/api/config":
self.serve_json({"cameras": config_entries(), "locked": LOCKED,
"ffmpeg": ffmpeg_available()})
return
if path == "/api/storage":
self.serve_json(storage_summary())
return
if path in ("/files", "/library", "/files.html"):
self.serve_file(LIBRARY_PAGE, "text/html; charset=utf-8")
return
if path == "/api/files":
kind = (query.get("kind") or ["recordings"])[0]
kind = "photos" if kind == "photos" else "recordings"
def one(key):
value = (query.get(key) or [""])[0].strip()
return value or None
try:
offset = max(0, int((query.get("offset") or ["0"])[0]))
except ValueError:
offset = 0
result = browse(kind, one("camera"), one("day"), offset)
result["limit"] = STORAGE_LIMIT
result["used"] = sum(f["size"] for f in recording_files())
self.serve_json(result)
return
if path.startswith("/recordings/") or path.startswith("/photos/"):
self.serve_stored(path, (query.get("download") or [""])[0] == "1")
return
if path in ("/stream", "/snapshot"):
cid = (query.get("cam") or [None])[0]
cam = find_camera(cid) if cid else (CAMERAS[0] if CAMERAS else None)
if cam is None:
self.fail(404, "No camera named %s" % cid)
return
if path == "/stream":
self.serve_stream(cam)
else:
self.serve_snapshot(cam)
return
self.fail(404, "Nothing is served at %s" % path)
def do_POST(self):
parsed = urlparse(self.path)
path = parsed.path
try:
length = int(self.headers.get("Content-Length") or 0)
except ValueError:
length = 0
if path == "/api/photo":
self.save_photo(parse_qs(parsed.query), length)
return
if length > 1000000:
self.fail(413, "That is too much data.")
return
try:
payload = json.loads(self.rfile.read(length) or b"{}")
except ValueError:
self.fail(400, "Could not read that as JSON.")
return
if path == "/api/test":
self.serve_json(probe(payload.get("url", "")))
return
if path == "/api/active":
cam = find_camera(str(payload.get("id") or ""))
if cam is None:
self.fail(404, "No such camera.")
return
if LOCKED:
self.fail(403, "Recording cannot be changed while --lock-config is set.")
return
cam.active = bool(payload.get("active"))
if cam.active:
cam.recorder.start()
else:
cam.recorder.stop()
try:
persist()
except OSError as err:
self.fail(500, "Could not save that (%s)" % err)
return
self.serve_json(dict(recording_state(cam), id=cam.id))
return
if path == "/api/delete":
if LOCKED:
self.fail(403, "Files cannot be deleted while --lock-config is set.")
return
kind = "photos" if payload.get("kind") == "photos" else "recordings"
names = payload.get("names")
if not isinstance(names, list) or not names:
self.fail(400, "Nothing was selected.")
return
if len(names) > 5000:
self.fail(400, "That is too many files at once.")
return
self.serve_json(delete_files(kind, [str(n) for n in names]))
return
if path == "/api/config":
if LOCKED:
self.fail(403, "Editing is disabled because the server was "
"started with --lock-config.")
return
self.save_cameras(payload)
return
self.fail(404, "Nothing accepts a POST at %s" % path)
def save_photo(self, query, length):
"""
Store a picture taken in the page.
The browser sends the image rather than the server grabbing its own,
because the page may be zoomed in: what gets saved is the exact framing
on screen, cut from the original pixels.
"""
if length <= 0 or length > 80 * 1024 * 1024:
self.fail(400, "No image, or one too large to be a photo.")
return
cam = find_camera((query.get("cam") or [""])[0])
if cam is None:
self.fail(404, "No such camera.")
return
ctype = (self.headers.get("Content-Type") or "").lower()
ext = "jpg" if "jpeg" in ctype else "png"
body = self.rfile.read(length)
if not (body[:3] == b"\xff\xd8\xff" or body[:8] == b"\x89PNG\r\n\x1a\n"):
self.fail(400, "That did not arrive as a PNG or JPEG.")
return
ensure_dirs()
stamp = time.strftime("%Y%m%d_%H%M%S")
name = "%s_%s.%s" % (stamp, slug(cam.name), ext)
target = os.path.join(PHOTO_DIR, name)
# Two photos inside one second should not overwrite each other.
n = 2
while os.path.exists(target):
name = "%s_%s-%d.%s" % (stamp, slug(cam.name), n, ext)
target = os.path.join(PHOTO_DIR, name)
n += 1
try:
with open(target, "wb") as fh:
fh.write(body)
except OSError as err:
self.fail(500, "Could not save the photo (%s)" % err)
return
self.serve_json({"ok": True, "name": name,
"url": "/photos/" + quote(name), "bytes": len(body)})
def save_cameras(self, payload):
incoming = payload.get("cameras")
if not isinstance(incoming, list):
self.fail(400, "Expected a list of cameras.")
return
known = {c.id: c.url for c in CAMERAS}
entries, seen = [], set()
for i, item in enumerate(incoming):
if not isinstance(item, dict):
self.fail(400, "Camera %d is not readable." % (i + 1))
return
name = str(item.get("name") or "").strip() or "Camera %d" % (i + 1)
url = str(item.get("url") or "").strip()
cid = str(item.get("id") or "").strip() or new_id()
# The page never receives passwords, so bullets left in place mean
# "keep the one you already have".
url = restore_url(url, known.get(cid))
if not valid_url(url):
self.fail(400, "%s needs an address starting http:// or rtsp://" % name)
return
if cid in seen:
cid = new_id()
seen.add(cid)
entries.append({"id": cid, "name": name, "url": url,
"active": bool(item.get("active"))})
try:
save_config(entries)
except OSError as err:
self.fail(500, "Could not write cameras.json (%s)" % err)
return
apply_config(entries)
self.serve_json({"ok": True, "cameras": config_entries()})
def do_HEAD(self):
if urlparse(self.path).path in ("/", "/index.html"):
self.send_response(200)
self.send_header("Content-Type", "text/html; charset=utf-8")
self.send_header("Content-Length", "0")
self.end_headers()
else:
self.fail(404, "Nothing is served here")
# ---- handlers ----
def serve_page(self):
self.serve_file(PAGE, "text/html; charset=utf-8")
def serve_file(self, path, ctype):
try:
with open(path, "rb") as fh:
body = fh.read()
except OSError:
self.fail(500, "%s is missing. Keep it next to server.py."
% os.path.basename(path))
return
self.send_response(200)
self.send_header("Content-Type", ctype)
self.send_header("Content-Length", str(len(body)))
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(body)
def serve_stored(self, path, force_download=False):
"""Serve a recording or photo, refusing anything outside the data folder."""
root = RECORD_DIR if path.startswith("/recordings/") else PHOTO_DIR
rel = path.split("/", 2)[2] if path.count("/") >= 2 else ""
target = os.path.normpath(os.path.join(root, unquote(rel)))
if not target.startswith(os.path.normpath(root) + os.sep):
self.fail(403, "That is not inside the recordings folder.")
return
if not os.path.isfile(target):
self.fail(404, "No such file.")
return
ctype = ("video/x-matroska" if target.endswith(".mkv")
else "image/png" if target.endswith(".png")
else "image/jpeg" if target.endswith((".jpg", ".jpeg"))
else "application/octet-stream")
try:
size = os.path.getsize(target)
self.send_response(200)
self.send_header("Content-Type", ctype)
self.send_header("Content-Length", str(size))
# Served inline so a photo can simply be looked at. Saving is driven
# by the link instead, or by asking for it explicitly; forcing every
# request to be a download would stop thumbnails rendering at all.
if force_download:
self.send_header("Content-Disposition",
'attachment; filename="%s"' % os.path.basename(target))
self.end_headers()
with open(target, "rb") as fh:
while True:
chunk = fh.read(65536)
if not chunk:
break
self.wfile.write(chunk)
except ConnectionError:
pass
except OSError as err:
self.fail(500, "Could not read that file (%s)" % err)
def serve_json(self, obj):
body = json.dumps(obj).encode("utf-8")
self.send_response(200)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(body)
def serve_stream(self, cam):
if not cam.available:
self.fail(503, cam.reason or "This camera is unavailable.")
return
if cam.kind == "mjpeg":
self.relay_mjpeg(cam)
else:
self.relay_frames(cam)
def relay_mjpeg(self, cam):
"""Straight passthrough: the camera's own JPEG bytes, untouched."""
try:
upstream = cam.open_upstream("stream")
except (URLError, HTTPError, socket.timeout) as err:
self.fail(502, "Cannot reach %s at %s (%s)" % (cam.name, safe_host(cam.url), err))
return
ctype = upstream.headers.get(
"Content-Type", "multipart/x-mixed-replace; boundary=boundarydonotcross")
self.send_response(200)
self.send_header("Content-Type", ctype)
self.send_header("Cache-Control", "no-store, no-cache")
self.send_header("Pragma", "no-cache")
self.send_header("Connection", "close")
self.end_headers()
self.close_connection = True
try:
while True:
chunk = upstream.read(8192)
if not chunk:
break
self.wfile.write(chunk)
except ConnectionError:
# The viewer closed the tab, switched camera, or reloaded. Windows
# reports this as ConnectionAbortedError and Unix as BrokenPipeError,
# so catch the shared base class rather than one platform's spelling.
pass
finally:
upstream.close()
def relay_frames(self, cam):
"""Re-publish decoded frames as multipart JPEG."""
q = cam.subscribe()
self.send_response(200)
self.send_header("Content-Type",
"multipart/x-mixed-replace; boundary=%s" % BOUNDARY)
self.send_header("Cache-Control", "no-store, no-cache")
self.send_header("Pragma", "no-cache")
self.send_header("Connection", "close")
self.end_headers()
self.close_connection = True
head = ("--%s\r\nContent-Type: image/jpeg\r\nContent-Length: " % BOUNDARY).encode()
try:
while True:
try:
frame = q.get(timeout=15)
except Empty:
break # nothing arriving; let the page reconnect
self.wfile.write(head + str(len(frame)).encode() + b"\r\n\r\n")
self.wfile.write(frame)
self.wfile.write(b"\r\n")
except ConnectionError:
pass # viewer went away; unsubscribing below frees ffmpeg
finally:
cam.unsubscribe(q)
def serve_snapshot(self, cam):
if not cam.available:
self.fail(503, cam.reason or "This camera is unavailable.")
return
if cam.kind == "mjpeg":
try:
upstream = cam.open_upstream("snapshot")
body = upstream.read()
ctype = upstream.headers.get("Content-Type", "image/jpeg")
upstream.close()
except (URLError, HTTPError, socket.timeout) as err:
self.fail(502, "Cannot reach %s at %s (%s)" % (cam.name, safe_host(cam.url), err))
return
else:
body = cam.snapshot()
ctype = "image/jpeg"
if not body:
self.fail(502, "No picture yet from %s. %s"
% (cam.name, cam.describe().get("reason") or ""))
return
self.send_response(200)
self.send_header("Content-Type", ctype)
self.send_header("Content-Length", str(len(body)))
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(body)
# ---- helpers ----
def fail(self, code, message):
body = message.encode("utf-8")
self.send_response(code)
self.send_header("Content-Type", "text/plain; charset=utf-8")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
try:
self.wfile.write(body)
except ConnectionError:
pass
def log_message(self, fmt, *args):
try:
line = fmt % args
except (TypeError, ValueError):
line = fmt
if "/stream" in line:
return
sys.stderr.write(" %s\n" % line)
# --------------------------------------------------------------------------
# startup
# --------------------------------------------------------------------------
def build_camera(index, name, url, args):
return make_camera("cam%d" % (index + 1), name, url, args)
def lan_address():
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
try:
sock.connect(("192.0.2.1", 1)) # never actually sends anything
return sock.getsockname()[0]
except OSError:
return "127.0.0.1"
finally:
sock.close()
def main():
parser = argparse.ArgumentParser(
description="Serve a viewer/recorder page for MJPEG and RTSP cameras.")
parser.add_argument("--camera", action="append", default=[], metavar="NAME=URL",
help="Add a camera. Repeatable. RTSP URLs are decoded with ffmpeg.")
parser.add_argument("--source", help="Shorthand for a single MJPEG camera.")
parser.add_argument("--port", type=int, default=8000, help="Port to listen on")
parser.add_argument("--host", default="0.0.0.0",
help="Address to bind. Use 127.0.0.1 to keep it on this machine only.")
parser.add_argument("--rtsp-quality", type=int, default=3, metavar="1-31",
help="JPEG quality for RTSP cameras, 1 is best (default: %(default)s)")
parser.add_argument("--rtsp-fps", type=float, default=None,
help="Cap the frame rate of RTSP cameras to save CPU")
parser.add_argument("--rtsp-size", default=None, metavar="WxH",
help="Scale RTSP cameras down, e.g. 1280x720, to save CPU")
parser.add_argument("--verbose", action="store_true",
help="Show all ffmpeg output, including routine decoder chatter")
parser.add_argument("--lock-config", action="store_true",
help="Serve the cameras read-only; the config page cannot save")
parser.add_argument("--data-dir", default=None,
help="Where recordings, photos and cameras.json live "
"(default: ./data, or $CAMERA_STATION_DATA)")
parser.add_argument("--storage-limit", default=None, metavar="GB",
help="Total footage to keep before deleting the oldest "
"(default: 10)")
parser.add_argument("--segment-minutes", type=float, default=None,
help="Length of each recording file (default: 10)")
args = parser.parse_args()
global LOCKED, DATA_DIR, CONFIG_FILE, RECORD_DIR, PHOTO_DIR
global STORAGE_LIMIT, SEGMENT_SECONDS
LOCKED = args.lock_config
if args.data_dir:
DATA_DIR = os.path.abspath(args.data_dir)
CONFIG_FILE = os.path.join(DATA_DIR, "cameras.json")
RECORD_DIR = os.path.join(DATA_DIR, "recordings")
PHOTO_DIR = os.path.join(DATA_DIR, "photos")
if args.storage_limit:
try:
STORAGE_LIMIT = int(float(args.storage_limit) * 1024 ** 3)
except ValueError:
parser.error("--storage-limit takes a number of gigabytes")
if args.segment_minutes:
SEGMENT_SECONDS = max(10, int(args.segment_minutes * 60))
ensure_dirs()
OPTS.rtsp_quality = args.rtsp_quality
OPTS.rtsp_fps = args.rtsp_fps
OPTS.rtsp_size = args.rtsp_size
OPTS.verbose = args.verbose
specs = []
if args.source:
specs.append(("Camera 1", args.source))
for entry in args.camera:
name, sep, url = entry.partition("=")
if not sep:
name, url = "Camera %d" % (len(specs) + 1), entry
specs.append((name.strip(), url.strip()))
if specs:
# Cameras named on the command line win for this run and are not saved.
entries = [{"id": "cam%d" % (i + 1), "name": n, "url": u}
for i, (n, u) in enumerate(specs)]
source = "command line"
else:
entries = load_config()
source = "cameras.json"
if entries is None:
entries = [{"id": "cam%d" % (i + 1), "name": n, "url": u}
for i, (n, u) in enumerate(DEFAULT_CAMERAS)]
source = "defaults"
apply_config(entries)
threading.Thread(target=supervise, daemon=True).start()
httpd = Station((args.host, args.port), Handler)
shown = "localhost" if args.host in ("0.0.0.0", "127.0.0.1") else args.host
print("\n Camera Station")
for cam in CAMERAS:
note = "" if cam.available else " (%s)" % cam.reason
mark = "REC " if getattr(cam, "active", False) else " "
print(" %s %-6s %-6s %s%s" % (mark, cam.id, cam.kind, safe_host(cam.url), note))
print(" from %s" % source)
print(" data %s" % DATA_DIR)
print(" keeps %.0f GB of footage, oldest deleted first"
% (STORAGE_LIMIT / 1024 ** 3))
print(" open http://%s:%d" % (shown, args.port))
print(" set up http://%s:%d/config%s"
% (shown, args.port, " (locked)" if LOCKED else ""))
print(" files http://%s:%d/files" % (shown, args.port))
if args.host == "0.0.0.0":
print(" on LAN http://%s:%d" % (lan_address(), args.port))
print(" Ctrl+C to stop\n")
try:
httpd.serve_forever()
except KeyboardInterrupt:
print("\n Stopping.")
for cam in CAMERAS:
if getattr(cam, "recorder", None):
cam.recorder.stop() # closes the segment being written
if isinstance(cam, RtspCamera):
cam._stop_if_idle()
httpd.server_close()
print(" Stopped.\n")
if __name__ == "__main__":
main()