2026-10-03 16:12:25 +08:00
|
|
|
"""Local video management. Python standard library + an external MediaMTX process."""
|
|
|
|
|
import base64
|
|
|
|
|
import concurrent.futures
|
|
|
|
|
import datetime as dt
|
|
|
|
|
import hashlib
|
|
|
|
|
import hmac
|
|
|
|
|
import http.client
|
|
|
|
|
import http.cookies
|
|
|
|
|
import ipaddress
|
|
|
|
|
import json
|
|
|
|
|
import math
|
|
|
|
|
import mimetypes
|
|
|
|
|
import os
|
|
|
|
|
from pathlib import Path
|
|
|
|
|
import re
|
|
|
|
|
import secrets
|
|
|
|
|
import shutil
|
|
|
|
|
import socket
|
|
|
|
|
import sqlite3
|
|
|
|
|
import threading
|
|
|
|
|
import time
|
|
|
|
|
import urllib.error
|
|
|
|
|
import urllib.parse as url
|
|
|
|
|
import urllib.request
|
|
|
|
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
2026-10-03 16:16:23 +08:00
|
|
|
import onvif
|
2026-10-03 16:12:25 +08:00
|
|
|
|
|
|
|
|
ROOT = Path(__file__).resolve().parent
|
|
|
|
|
os.umask(0o077)
|
|
|
|
|
DATA = Path(os.environ.get('VISION_DATA', ROOT / 'data')).resolve()
|
|
|
|
|
DATA.mkdir(mode=0o700, parents=True, exist_ok=True)
|
|
|
|
|
RECORDINGS = Path(os.environ.get('VISION_RECORDINGS', DATA / 'recordings')).resolve()
|
|
|
|
|
RECORDINGS.mkdir(mode=0o700, parents=True, exist_ok=True)
|
|
|
|
|
CONFIG = DATA / 'mediamtx.yml'
|
|
|
|
|
LOCK = threading.RLock()
|
|
|
|
|
DB = sqlite3.connect(DATA / 'vision.db', check_same_thread=False)
|
|
|
|
|
DB.row_factory = sqlite3.Row
|
|
|
|
|
DB.executescript('''
|
|
|
|
|
PRAGMA journal_mode=WAL;
|
|
|
|
|
PRAGMA foreign_keys=ON;
|
|
|
|
|
CREATE TABLE IF NOT EXISTS settings (key TEXT PRIMARY KEY, value TEXT NOT NULL);
|
|
|
|
|
CREATE TABLE IF NOT EXISTS sites (id TEXT PRIMARY KEY, body TEXT NOT NULL);
|
|
|
|
|
CREATE TABLE IF NOT EXISTS cameras (id TEXT PRIMARY KEY, body TEXT NOT NULL);
|
|
|
|
|
CREATE TABLE IF NOT EXISTS assets (id TEXT PRIMARY KEY, body TEXT NOT NULL);
|
|
|
|
|
CREATE TABLE IF NOT EXISTS recorders (id TEXT PRIMARY KEY, body TEXT NOT NULL);
|
|
|
|
|
CREATE TABLE IF NOT EXISTS sessions (hash TEXT PRIMARY KEY, expires REAL NOT NULL);
|
|
|
|
|
CREATE TABLE IF NOT EXISTS audit (id INTEGER PRIMARY KEY, at TEXT, action TEXT, name TEXT);
|
|
|
|
|
''')
|
|
|
|
|
DB.commit()
|
|
|
|
|
os.chmod(DATA / 'vision.db', 0o600)
|
|
|
|
|
STARTED = time.time()
|
|
|
|
|
STORAGE = {'freeGB': None, 'usedGB': None, 'paused': False, 'reason': ''}
|
|
|
|
|
RECENT_RECORDS = {}
|
|
|
|
|
SCAN = {'running': False, 'items': [], 'error': '', 'network': '', 'at': None}
|
|
|
|
|
ATTEMPTS = {}
|
|
|
|
|
API_PORT, HLS_PORT, PLAYBACK_PORT = 19997, 18888, 19996
|
|
|
|
|
MEDIA_AUTH = None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class Problem(Exception):
|
|
|
|
|
def __init__(self, message, status=400):
|
|
|
|
|
super().__init__(message)
|
|
|
|
|
self.status = status
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def setting(key, default=None):
|
|
|
|
|
with LOCK:
|
|
|
|
|
row = DB.execute('SELECT value FROM settings WHERE key=?', (key,)).fetchone()
|
|
|
|
|
return json.loads(row[0]) if row else default
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def set_setting(key, value):
|
|
|
|
|
with LOCK:
|
|
|
|
|
DB.execute('INSERT OR REPLACE INTO settings VALUES (?,?)', (key, json.dumps(value)))
|
|
|
|
|
DB.commit()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def objects(table):
|
|
|
|
|
with LOCK:
|
|
|
|
|
return [dict(json.loads(r['body']), id=r['id']) for r in DB.execute('SELECT * FROM ' + table)]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def get_object(table, key):
|
|
|
|
|
return next((o for o in objects(table) if o['id'] == key), None)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def audit(action, name=''):
|
|
|
|
|
with LOCK:
|
|
|
|
|
DB.execute('INSERT INTO audit (at,action,name) VALUES (?,?,?)',
|
|
|
|
|
(dt.datetime.now(dt.timezone.utc).isoformat(), action, name))
|
|
|
|
|
DB.execute('DELETE FROM audit WHERE id NOT IN (SELECT id FROM audit ORDER BY id DESC LIMIT 500)')
|
|
|
|
|
DB.commit()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def save_object(table, body):
|
|
|
|
|
with LOCK:
|
|
|
|
|
DB.execute('INSERT OR REPLACE INTO ' + table + ' VALUES (?,?)', (body['id'], json.dumps(body)))
|
|
|
|
|
DB.commit()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def clean_text(value, limit=120, required=False):
|
|
|
|
|
if not isinstance(value, str) or len(value) > limit or any(ord(c) < 32 for c in value):
|
|
|
|
|
raise Problem('文本格式或长度不正确')
|
|
|
|
|
value = value.strip()
|
|
|
|
|
if required and not value:
|
|
|
|
|
raise Problem('请填写必填项')
|
|
|
|
|
return value
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def integer(value, lo, hi):
|
|
|
|
|
if isinstance(value, bool) or not isinstance(value, int) or not lo <= value <= hi:
|
|
|
|
|
raise Problem('数值超出允许范围')
|
|
|
|
|
return value
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def host_address(value):
|
|
|
|
|
try:
|
|
|
|
|
ip = ipaddress.ip_address(value)
|
|
|
|
|
if ip.version != 4 or ip.is_global or ip.is_loopback or ip.is_multicast or ip.is_unspecified or ip.is_reserved:
|
|
|
|
|
raise ValueError()
|
|
|
|
|
return str(ip)
|
|
|
|
|
except ValueError:
|
|
|
|
|
raise Problem('请填写局域网或 Tailscale IPv4 地址')
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
SITE_FIELDS = ('name', 'province', 'city', 'district', 'street', 'community', 'owner',
|
|
|
|
|
'building', 'unit', 'floor', 'door', 'address', 'note')
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def save_site(body):
|
|
|
|
|
key = body.get('id') or secrets.token_hex(8)
|
|
|
|
|
if not re.fullmatch('[a-f0-9]{16}', key):
|
|
|
|
|
raise Problem('空间编号不正确')
|
|
|
|
|
item = {k: clean_text(body.get(k, ''), 240 if k in ('address', 'note') else 100,
|
|
|
|
|
required=k == 'name') for k in SITE_FIELDS}
|
|
|
|
|
item.update(id=key, kind=body.get('kind', 'home'))
|
|
|
|
|
if item['kind'] not in ('home', 'company', 'other'):
|
|
|
|
|
raise Problem('请选择家庭、公司或其他用途')
|
|
|
|
|
save_object('sites', item)
|
|
|
|
|
audit('保存空间', item['name'])
|
|
|
|
|
return item
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def save_camera(body):
|
|
|
|
|
key = body.get('id') or secrets.token_hex(8)
|
|
|
|
|
if not re.fullmatch('[a-f0-9]{16}', key):
|
|
|
|
|
raise Problem('设备编号不正确')
|
|
|
|
|
old = get_object('cameras', key) or {}
|
|
|
|
|
item = {k: clean_text(body.get(k, ''), 120, required=k == 'name')
|
|
|
|
|
for k in ('name', 'point', 'username', 'model')}
|
|
|
|
|
asset = get_object('assets', body.get('assetId'))
|
|
|
|
|
if not asset:
|
|
|
|
|
raise Problem('请先建立并选择摄像头对象')
|
|
|
|
|
siblings = [c for c in objects('cameras') if c.get('assetId') == asset['id'] and c['id'] != key]
|
|
|
|
|
if len(siblings) >= asset['lensCount']:
|
|
|
|
|
raise Problem('该摄像头的镜头通道数已达到档案中设置的镜头数')
|
|
|
|
|
site = asset['siteId']
|
|
|
|
|
source_kind = body.get('sourceKind', 'direct')
|
|
|
|
|
recorder = get_object('recorders', body.get('recorderId')) if source_kind == 'recorder' else None
|
|
|
|
|
if source_kind not in ('direct', 'recorder') or (source_kind == 'recorder' and not recorder):
|
|
|
|
|
raise Problem('请选择录像机或直接接入')
|
|
|
|
|
item.update(assetId=asset['id'], sourceKind=source_kind, recorderId=recorder['id'] if recorder else '',
|
|
|
|
|
channelNumber=clean_text(body.get('channelNumber', ''), 30),
|
|
|
|
|
archiveSource=body.get('archiveSource', 'recorder' if recorder else 'local'))
|
|
|
|
|
if item['archiveSource'] not in ('local', 'recorder'):
|
|
|
|
|
raise Problem('录像存储来源不正确')
|
2026-10-03 16:16:23 +08:00
|
|
|
archive_recorder = body.get('archiveRecorderId') or (recorder['id'] if recorder else '')
|
|
|
|
|
if archive_recorder and not get_object('recorders', archive_recorder):
|
|
|
|
|
raise Problem('关联的录像存储设备不存在')
|
|
|
|
|
item['archiveRecorderId'] = archive_recorder
|
2026-10-03 16:12:25 +08:00
|
|
|
item.update(id=key, siteId=site, host=host_address(recorder['host'] if recorder else body.get('host', '')),
|
|
|
|
|
port=integer(recorder['port'] if recorder else body.get('port', 554), 1, 65535))
|
|
|
|
|
for k, default in [('mainPath', '/stream1'), ('subPath', '/stream2')]:
|
|
|
|
|
path = clean_text(body.get(k, default), 500, required=True)
|
|
|
|
|
if not path.startswith('/') or path.startswith('//') or '#' in path:
|
|
|
|
|
raise Problem('码流路径必须以一个 / 开头,不含 #')
|
|
|
|
|
item[k] = path
|
|
|
|
|
password = body.get('password', '')
|
|
|
|
|
if not isinstance(password, str) or len(password) > 256 or any(ord(c) < 32 for c in password):
|
|
|
|
|
raise Problem('设备密码格式不正确')
|
|
|
|
|
item['password'] = '' if recorder else password or old.get('password', '')
|
|
|
|
|
for k in ('enabled', 'record'):
|
|
|
|
|
if not isinstance(body.get(k, False), bool):
|
|
|
|
|
raise Problem('开关值不正确')
|
|
|
|
|
item[k] = body.get(k, False)
|
|
|
|
|
if not old and len(objects('cameras')) >= 32:
|
|
|
|
|
raise Problem('此版本最多配置 32 台设备')
|
|
|
|
|
save_object('cameras', item)
|
|
|
|
|
write_media_config()
|
|
|
|
|
audit('保存摄像头', item['name'])
|
|
|
|
|
return public_camera(item)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def save_asset(body):
|
|
|
|
|
key = body.get('id') or secrets.token_hex(8)
|
|
|
|
|
if not re.fullmatch('[a-f0-9]{16}', key):
|
|
|
|
|
raise Problem('对象编号不正确')
|
|
|
|
|
site = body.get('siteId')
|
|
|
|
|
if not get_object('sites', site):
|
|
|
|
|
raise Problem('请选择所属空间')
|
|
|
|
|
item = {k: clean_text(body.get(k, ''), 120, required=k == 'name')
|
|
|
|
|
for k in ('name', 'brand', 'model', 'serial', 'point', 'note')}
|
|
|
|
|
item.update(id=key, siteId=site, lensCount=integer(body.get('lensCount', 2), 1, 16))
|
|
|
|
|
attached = [c for c in objects('cameras') if c.get('assetId') == key]
|
|
|
|
|
if len(attached) > item['lensCount']:
|
|
|
|
|
raise Problem('镜头数不能小于已经建立的通道数')
|
|
|
|
|
with LOCK:
|
|
|
|
|
save_object('assets', item)
|
|
|
|
|
for c in attached:
|
|
|
|
|
c['siteId'] = site
|
|
|
|
|
save_object('cameras', c)
|
|
|
|
|
audit('保存摄像头对象', item['name'])
|
|
|
|
|
return item
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def save_recorder(body):
|
|
|
|
|
key = body.get('id') or secrets.token_hex(8)
|
|
|
|
|
if not re.fullmatch('[a-f0-9]{16}', key):
|
|
|
|
|
raise Problem('录像机编号不正确')
|
|
|
|
|
old = get_object('recorders', key) or {}
|
|
|
|
|
item = {k: clean_text(body.get(k, ''), 120, required=k == 'name')
|
|
|
|
|
for k in ('name', 'brand', 'model', 'username')}
|
|
|
|
|
site = body.get('siteId')
|
|
|
|
|
if not get_object('sites', site):
|
|
|
|
|
raise Problem('请选择录像机所属空间')
|
|
|
|
|
item.update(id=key, siteId=site, host=host_address(body.get('host', '')),
|
|
|
|
|
port=integer(body.get('port', 554), 1, 65535), driver=body.get('driver', 'generic-rtsp'))
|
|
|
|
|
if item['driver'] not in ('generic-rtsp', 'hikvision-rtsp', 'tplink-rtsp', 'dahua-rtsp'):
|
|
|
|
|
raise Problem('接入驱动不正确')
|
|
|
|
|
password = body.get('password', '')
|
|
|
|
|
if not isinstance(password, str) or len(password) > 256 or any(ord(c) < 32 for c in password):
|
|
|
|
|
raise Problem('密码格式不正确')
|
|
|
|
|
item['password'] = password or old.get('password', '')
|
|
|
|
|
save_object('recorders', item)
|
|
|
|
|
write_media_config()
|
|
|
|
|
audit('保存录像机', item['name'])
|
|
|
|
|
return public_camera(item)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def public_camera(c):
|
|
|
|
|
return {**{k: v for k, v in c.items() if k != 'password'}, 'passwordConfigured': bool(c['password'])}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def stream_name(c, quality='main'):
|
|
|
|
|
return 'cam_' + c['id'] + '_' + quality
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def source(c, quality):
|
|
|
|
|
endpoint = get_object('recorders', c.get('recorderId')) if c.get('sourceKind') == 'recorder' else c
|
|
|
|
|
if not endpoint:
|
|
|
|
|
raise Problem('关联录像机不存在')
|
|
|
|
|
auth = (url.quote(endpoint['username'], safe='') + ':' + url.quote(endpoint['password'], safe='') + '@') if endpoint['username'] else ''
|
|
|
|
|
return 'rtsp://' + auth + endpoint['host'] + ':' + str(endpoint['port']) + c[quality + 'Path']
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def write_media_config():
|
|
|
|
|
with LOCK:
|
|
|
|
|
paths = {}
|
|
|
|
|
policy = setting('storage', {'retentionDays': 7, 'maxGB': 40, 'reserveGB': 8})
|
|
|
|
|
for c in objects('cameras'):
|
|
|
|
|
for quality in ('main', 'sub'):
|
|
|
|
|
paths[stream_name(c, quality)] = {
|
|
|
|
|
'source': source(c, quality) if c['enabled'] else 'publisher',
|
|
|
|
|
'rtspTransport': 'tcp',
|
|
|
|
|
'sourceOnDemand': quality == 'sub' and c['enabled'],
|
|
|
|
|
'record': bool(c['enabled'] and c['record'] and quality == 'main' and not STORAGE['paused']),
|
|
|
|
|
'recordDeleteAfter': str(policy['retentionDays'] * 24) + 'h',
|
|
|
|
|
}
|
|
|
|
|
config = {
|
|
|
|
|
'logLevel': 'warn', 'logDestinations': ['stdout'],
|
|
|
|
|
'api': True, 'apiAddress': '127.0.0.1:' + str(API_PORT),
|
|
|
|
|
'rtsp': True, 'rtspAddress': '127.0.0.1:18554', 'rtspTransports': ['tcp'],
|
|
|
|
|
'rtmp': False, 'srt': False, 'webrtc': False, 'moq': False,
|
|
|
|
|
'hls': True, 'hlsAddress': '127.0.0.1:' + str(HLS_PORT), 'hlsVariant': 'fmp4',
|
|
|
|
|
'hlsSegmentDuration': '1s', 'hlsSegmentCount': 5, 'hlsAllowOrigins': [],
|
|
|
|
|
'playback': True, 'playbackAddress': '127.0.0.1:' + str(PLAYBACK_PORT),
|
|
|
|
|
'playbackAllowOrigins': [],
|
|
|
|
|
'authMethod': 'internal',
|
|
|
|
|
'authInternalUsers': [{'user': 'vision', 'pass': setting('mediaSecret'),
|
|
|
|
|
'ips': ['127.0.0.1', '::1'],
|
|
|
|
|
'permissions': [{'action': a} for a in ('read', 'api', 'playback')]}],
|
|
|
|
|
'pathDefaults': {'recordFormat': 'fmp4', 'recordPath': str(RECORDINGS / '%path' / '%Y-%m-%d_%H-%M-%S-%f'),
|
|
|
|
|
'recordSegmentDuration': '5m', 'recordPartDuration': '1s'},
|
|
|
|
|
'paths': paths,
|
|
|
|
|
}
|
|
|
|
|
raw = json.dumps(config, ensure_ascii=False, indent=2)
|
|
|
|
|
if CONFIG.exists() and CONFIG.read_text(encoding='utf8') == raw:
|
|
|
|
|
return
|
|
|
|
|
temp = DATA / 'mediamtx.tmp'
|
|
|
|
|
temp.write_text(raw, encoding='utf8')
|
|
|
|
|
os.chmod(temp, 0o600)
|
|
|
|
|
os.replace(temp, CONFIG)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def media_json(path, playback=False):
|
|
|
|
|
target = 'http://127.0.0.1:' + str(PLAYBACK_PORT if playback else API_PORT) + path
|
|
|
|
|
try:
|
|
|
|
|
req = urllib.request.Request(target, headers={'Authorization': MEDIA_AUTH})
|
|
|
|
|
with urllib.request.urlopen(req, timeout=5) as r:
|
|
|
|
|
return json.load(r)
|
|
|
|
|
except Exception:
|
|
|
|
|
raise Problem('视频服务暂不可用,或该时段尚无可回放录像', 503)
|
|
|
|
|
|
|
|
|
|
|
2026-10-03 16:26:28 +08:00
|
|
|
def media_config_watch():
|
|
|
|
|
# Reconcile API configuration after atomic file replacements. Multiple rapid
|
|
|
|
|
# writes can be coalesced by a filesystem watcher; the file remains authority.
|
|
|
|
|
applied = None
|
|
|
|
|
while True:
|
|
|
|
|
try:
|
|
|
|
|
raw = CONFIG.read_bytes()
|
|
|
|
|
digest = hashlib.sha256(raw).digest()
|
|
|
|
|
if digest != applied:
|
|
|
|
|
desired = json.loads(raw)['paths']
|
|
|
|
|
existing = {p['name'] for p in media_json('/v3/config/paths/list?itemsPerPage=200').get('items', [])}
|
|
|
|
|
for name, body in desired.items():
|
|
|
|
|
action, method = ('patch', 'PATCH') if name in existing else ('add', 'POST')
|
|
|
|
|
req = urllib.request.Request('http://127.0.0.1:' + str(API_PORT) + '/v3/config/paths/' + action + '/' + name,
|
|
|
|
|
data=json.dumps(body).encode(), method=method,
|
|
|
|
|
headers={'Authorization': MEDIA_AUTH, 'Content-Type': 'application/json'})
|
|
|
|
|
with urllib.request.urlopen(req, timeout=5) as response:
|
|
|
|
|
response.read()
|
|
|
|
|
applied = digest
|
|
|
|
|
except Exception:
|
|
|
|
|
pass # Retry on the next pass, including after a media process restart.
|
|
|
|
|
time.sleep(5)
|
|
|
|
|
|
|
|
|
|
|
2026-10-03 16:12:25 +08:00
|
|
|
def status():
|
|
|
|
|
try:
|
|
|
|
|
data = media_json('/v3/paths/list?itemsPerPage=200')
|
|
|
|
|
ready = {p['name']: {'ready': bool(p.get('ready')), 'tracks': p.get('tracks', []),
|
|
|
|
|
'bytesReceived': p.get('bytesReceived', 0)} for p in data.get('items', [])}
|
|
|
|
|
media = True
|
|
|
|
|
except Problem:
|
|
|
|
|
ready, media = {}, False
|
|
|
|
|
cameras = []
|
|
|
|
|
for c in objects('cameras'):
|
|
|
|
|
if c.get('sourceKind') == 'recorder':
|
|
|
|
|
endpoint = get_object('recorders', c.get('recorderId'))
|
|
|
|
|
if endpoint:
|
|
|
|
|
c = dict(c, host=endpoint['host'], port=endpoint['port'])
|
|
|
|
|
main = ready.get(stream_name(c), {})
|
|
|
|
|
cameras.append({**public_camera(c), 'ready': main.get('ready', False), 'tracks': main.get('tracks', []),
|
|
|
|
|
'recordingRequested': c['enabled'] and c['record'] and not STORAGE['paused'],
|
|
|
|
|
'recordingActive': bool(c['enabled'] and c['record'] and not STORAGE['paused'] and main.get('ready', False)
|
|
|
|
|
and time.time() - RECENT_RECORDS.get(stream_name(c), 0) < 45)})
|
|
|
|
|
with LOCK:
|
|
|
|
|
storage = dict(STORAGE)
|
|
|
|
|
scan = dict(SCAN)
|
|
|
|
|
return {'cameras': cameras, 'assets': objects('assets'), 'recorders': [public_camera(r) for r in objects('recorders')],
|
|
|
|
|
'sites': objects('sites'), 'mediaOnline': media,
|
|
|
|
|
'storage': {**setting('storage'), **storage}, 'scan': scan,
|
|
|
|
|
'version': (ROOT / 'VERSION').read_text().strip(), 'uptime': round(time.time() - STARTED)}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def storage_watch():
|
|
|
|
|
while True:
|
|
|
|
|
try:
|
|
|
|
|
usage = shutil.disk_usage(RECORDINGS)
|
|
|
|
|
used, recent = 0, {}
|
|
|
|
|
for p in RECORDINGS.rglob('*.mp4'):
|
|
|
|
|
try:
|
|
|
|
|
info = p.stat()
|
|
|
|
|
used += info.st_size
|
|
|
|
|
recent[p.parent.name] = max(recent.get(p.parent.name, 0), info.st_mtime)
|
|
|
|
|
except FileNotFoundError:
|
|
|
|
|
pass
|
|
|
|
|
policy = setting('storage')
|
|
|
|
|
free_gb, used_gb = usage.free / 1e9, used / 1e9
|
|
|
|
|
reason = '录像空间已达到配额' if used_gb >= policy['maxGB'] else ('磁盘可用空间不足' if free_gb < policy['reserveGB'] else '')
|
|
|
|
|
with LOCK:
|
|
|
|
|
changed = STORAGE['paused'] != bool(reason)
|
|
|
|
|
STORAGE.update(freeGB=round(free_gb, 2), usedGB=round(used_gb, 2), paused=bool(reason), reason=reason)
|
|
|
|
|
RECENT_RECORDS.clear()
|
|
|
|
|
RECENT_RECORDS.update(recent)
|
|
|
|
|
if changed:
|
|
|
|
|
write_media_config()
|
|
|
|
|
audit('录像空间保护', reason or '空间恢复,继续按设备配置录像')
|
|
|
|
|
except Exception:
|
|
|
|
|
with LOCK:
|
|
|
|
|
STORAGE.update(paused=True, reason='存储状态读取失败')
|
|
|
|
|
try:
|
|
|
|
|
write_media_config()
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
time.sleep(20)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def discover(network):
|
|
|
|
|
def probe(ip):
|
|
|
|
|
try:
|
|
|
|
|
with socket.create_connection((str(ip), 554), timeout=0.7) as s:
|
|
|
|
|
s.settimeout(1)
|
|
|
|
|
s.sendall(('OPTIONS rtsp://' + str(ip) + '/ RTSP/1.0\r\nCSeq: 1\r\n\r\n').encode())
|
|
|
|
|
response = s.recv(1024).decode(errors='replace').split('\r\n')[0]
|
|
|
|
|
if response.startswith('RTSP/'):
|
|
|
|
|
return {'host': str(ip), 'port': 554, 'response': response, 'authenticated': False}
|
|
|
|
|
except OSError:
|
|
|
|
|
return None
|
|
|
|
|
try:
|
|
|
|
|
with concurrent.futures.ThreadPoolExecutor(max_workers=24) as pool:
|
|
|
|
|
items = [p for p in pool.map(probe, network.hosts()) if p]
|
|
|
|
|
with LOCK:
|
|
|
|
|
SCAN.update(items=items, at=dt.datetime.now(dt.timezone.utc).isoformat())
|
|
|
|
|
except Exception:
|
|
|
|
|
with LOCK:
|
|
|
|
|
SCAN['error'] = '探测未完成,请检查主机网络'
|
|
|
|
|
finally:
|
|
|
|
|
with LOCK:
|
|
|
|
|
SCAN['running'] = False
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def parse_time(text):
|
|
|
|
|
try:
|
|
|
|
|
value = dt.datetime.fromisoformat(text.replace('Z', '+00:00'))
|
|
|
|
|
if value.tzinfo is None:
|
|
|
|
|
raise ValueError()
|
|
|
|
|
return value
|
|
|
|
|
except (ValueError, AttributeError):
|
|
|
|
|
raise Problem('时间必须包含时区')
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class Handler(BaseHTTPRequestHandler):
|
|
|
|
|
server_version = 'ZhaoVision'
|
|
|
|
|
|
|
|
|
|
def log_message(self, *_):
|
|
|
|
|
pass # Do not put camera credentials, cookies or playback queries into logs.
|
|
|
|
|
|
|
|
|
|
def answer(self, value, code=200, cookie=None):
|
|
|
|
|
data = json.dumps(value, ensure_ascii=False).encode()
|
|
|
|
|
self.send_response(code)
|
|
|
|
|
self.send_header('Content-Type', 'application/json; charset=utf-8')
|
|
|
|
|
self.send_header('Content-Length', str(len(data)))
|
|
|
|
|
self.send_header('Cache-Control', 'no-store')
|
|
|
|
|
self.send_header('X-Content-Type-Options', 'nosniff')
|
|
|
|
|
if cookie:
|
|
|
|
|
self.send_header('Set-Cookie', cookie)
|
|
|
|
|
self.end_headers()
|
|
|
|
|
self.wfile.write(data)
|
|
|
|
|
|
|
|
|
|
def token(self):
|
|
|
|
|
try:
|
|
|
|
|
c = http.cookies.SimpleCookie(self.headers.get('Cookie', ''))
|
|
|
|
|
return c['vision'].value if 'vision' in c else ''
|
|
|
|
|
except http.cookies.CookieError:
|
|
|
|
|
return ''
|
|
|
|
|
|
|
|
|
|
def authenticated(self):
|
|
|
|
|
token = self.token()
|
|
|
|
|
if not re.fullmatch('[a-f0-9]{64}', token):
|
|
|
|
|
return False
|
|
|
|
|
with LOCK:
|
|
|
|
|
return bool(DB.execute('SELECT 1 FROM sessions WHERE hash=? AND expires>?',
|
|
|
|
|
(hashlib.sha256(token.encode()).hexdigest(), time.time())).fetchone())
|
|
|
|
|
|
|
|
|
|
def require_auth(self):
|
|
|
|
|
if not self.authenticated():
|
|
|
|
|
raise Problem('请先登录', 401)
|
|
|
|
|
|
|
|
|
|
def body(self):
|
|
|
|
|
if self.headers.get_content_type() != 'application/json':
|
|
|
|
|
raise Problem('请求需要 JSON 格式')
|
|
|
|
|
n = int(self.headers.get('Content-Length', '0'))
|
|
|
|
|
if not 0 < n <= 16384:
|
|
|
|
|
raise Problem('请求长度不正确')
|
|
|
|
|
try:
|
|
|
|
|
data = json.loads(self.rfile.read(n))
|
|
|
|
|
if not isinstance(data, dict):
|
|
|
|
|
raise ValueError()
|
|
|
|
|
return data
|
|
|
|
|
except (ValueError, UnicodeError):
|
|
|
|
|
raise Problem('请求内容无法解析')
|
|
|
|
|
|
|
|
|
|
def validate_origin(self):
|
|
|
|
|
origin = self.headers.get('Origin')
|
|
|
|
|
if origin and url.urlsplit(origin).netloc != self.headers.get('Host'):
|
|
|
|
|
raise Problem('请从系统页面提交操作', 403)
|
|
|
|
|
|
|
|
|
|
def login(self, data, setup=False):
|
|
|
|
|
ip = self.client_address[0]
|
|
|
|
|
now = time.time()
|
|
|
|
|
with LOCK:
|
|
|
|
|
for k in list(ATTEMPTS):
|
|
|
|
|
if ATTEMPTS[k][1] < now:
|
|
|
|
|
del ATTEMPTS[k]
|
|
|
|
|
count, expiry = ATTEMPTS.get(ip, (0, now + 900))
|
|
|
|
|
ATTEMPTS[ip] = (count + 1, expiry)
|
|
|
|
|
if count >= 12:
|
|
|
|
|
raise Problem('尝试次数过多,请 15 分钟后重试', 429)
|
|
|
|
|
account = setting('account')
|
|
|
|
|
password = data.get('password', '')
|
|
|
|
|
if not isinstance(password, str) or not 10 <= len(password) <= 128:
|
|
|
|
|
raise Problem('管理密码需要 10 至 128 个字符')
|
|
|
|
|
if setup:
|
|
|
|
|
if account:
|
|
|
|
|
raise Problem('管理员已经初始化', 409)
|
|
|
|
|
code = data.get('code', '')
|
|
|
|
|
expected = (DATA / 'setup-code.txt').read_text().strip()
|
|
|
|
|
if not isinstance(code, str) or not hmac.compare_digest(code, expected):
|
|
|
|
|
raise Problem('初始化码不正确', 403)
|
|
|
|
|
salt = secrets.token_hex(16)
|
|
|
|
|
derived = hashlib.scrypt(password.encode(), salt=salt.encode(), n=16384, r=8, p=1).hex()
|
|
|
|
|
set_setting('account', {'salt': salt, 'hash': derived})
|
|
|
|
|
(DATA / 'setup-code.txt').unlink(missing_ok=True)
|
|
|
|
|
audit('初始化管理员')
|
|
|
|
|
else:
|
|
|
|
|
if not account:
|
|
|
|
|
raise Problem('请先初始化管理员', 409)
|
|
|
|
|
derived = hashlib.scrypt(password.encode(), salt=account['salt'].encode(), n=16384, r=8, p=1).hex()
|
|
|
|
|
if not hmac.compare_digest(derived, account['hash']):
|
|
|
|
|
raise Problem('密码不正确', 403)
|
|
|
|
|
token = secrets.token_hex(32)
|
|
|
|
|
DB.execute('DELETE FROM sessions WHERE expires<?', (now,))
|
|
|
|
|
DB.execute('INSERT INTO sessions VALUES (?,?)', (hashlib.sha256(token.encode()).hexdigest(), now + 86400 * 7))
|
|
|
|
|
DB.commit()
|
|
|
|
|
ATTEMPTS.pop(ip, None)
|
|
|
|
|
secure = '; Secure' if os.environ.get('VISION_SECURE_COOKIE') == '1' else ''
|
|
|
|
|
self.answer({'ok': True}, cookie='vision=' + token + '; HttpOnly; SameSite=Strict; Path=/; Max-Age=604800' + secure)
|
|
|
|
|
|
|
|
|
|
def do_POST(self):
|
|
|
|
|
try:
|
|
|
|
|
self.validate_origin()
|
|
|
|
|
path = url.urlsplit(self.path).path
|
|
|
|
|
data = self.body()
|
|
|
|
|
if path in ('/api/login', '/api/setup'):
|
|
|
|
|
self.login(data, path == '/api/setup')
|
|
|
|
|
return
|
|
|
|
|
self.require_auth()
|
|
|
|
|
if path == '/api/logout':
|
|
|
|
|
with LOCK:
|
|
|
|
|
DB.execute('DELETE FROM sessions WHERE hash=?', (hashlib.sha256(self.token().encode()).hexdigest(),))
|
|
|
|
|
DB.commit()
|
|
|
|
|
self.answer({'ok': True}, cookie='vision=; Path=/; Max-Age=0; HttpOnly; SameSite=Strict')
|
|
|
|
|
elif path == '/api/sites':
|
|
|
|
|
self.answer(save_site(data))
|
|
|
|
|
elif path == '/api/cameras':
|
|
|
|
|
self.answer(save_camera(data))
|
|
|
|
|
elif path == '/api/assets':
|
|
|
|
|
self.answer(save_asset(data))
|
|
|
|
|
elif path == '/api/recorders':
|
|
|
|
|
self.answer(save_recorder(data))
|
2026-10-03 16:16:23 +08:00
|
|
|
elif path == '/api/onvif':
|
|
|
|
|
old = get_object('cameras', data.get('id')) or {}
|
|
|
|
|
host = host_address(data.get('host', ''))
|
|
|
|
|
port = integer(data.get('port', 80), 1, 65535)
|
|
|
|
|
username = clean_text(data.get('username', ''), 120, required=True)
|
|
|
|
|
password = data.get('password') or old.get('password', '')
|
|
|
|
|
if not isinstance(password, str) or len(password) > 256:
|
|
|
|
|
raise Problem('设备密码格式不正确')
|
|
|
|
|
try:
|
|
|
|
|
result = onvif.inspect(host, port, username, password)
|
|
|
|
|
except Exception:
|
|
|
|
|
raise Problem('ONVIF 读取未成功,请检查账号、设备时间、服务端口及 ONVIF 是否启用', 502)
|
|
|
|
|
self.answer(result)
|
2026-10-03 16:12:25 +08:00
|
|
|
elif path == '/api/storage':
|
|
|
|
|
policy = {'retentionDays': integer(data.get('retentionDays'), 1, 365),
|
|
|
|
|
'maxGB': integer(data.get('maxGB'), 5, 100000),
|
|
|
|
|
'reserveGB': integer(data.get('reserveGB'), 2, 1000)}
|
|
|
|
|
set_setting('storage', policy)
|
|
|
|
|
write_media_config()
|
|
|
|
|
audit('更新录像保留策略')
|
|
|
|
|
self.answer({'ok': True})
|
|
|
|
|
elif path == '/api/discover':
|
|
|
|
|
try:
|
|
|
|
|
network = ipaddress.ip_network(data.get('network', ''), strict=False)
|
|
|
|
|
if network.version != 4 or network.prefixlen < 24 or network.prefixlen > 30:
|
|
|
|
|
raise ValueError()
|
|
|
|
|
host_address(str(network.network_address))
|
|
|
|
|
except ValueError:
|
|
|
|
|
raise Problem('请输入最多 254 个地址的局域网网段,例如 192.168.1.0/24')
|
|
|
|
|
with LOCK:
|
|
|
|
|
if SCAN['running']:
|
|
|
|
|
raise Problem('探测正在进行', 409)
|
|
|
|
|
SCAN.update(running=True, items=[], error='', network=str(network))
|
|
|
|
|
threading.Thread(target=discover, args=(network,), daemon=True).start()
|
|
|
|
|
self.answer({'ok': True}, 202)
|
|
|
|
|
else:
|
|
|
|
|
raise Problem('接口不存在', 404)
|
|
|
|
|
except Problem as e:
|
|
|
|
|
self.answer({'error': str(e)}, e.status)
|
|
|
|
|
except (BrokenPipeError, ConnectionResetError):
|
|
|
|
|
pass
|
|
|
|
|
except Exception:
|
|
|
|
|
self.answer({'error': '操作未完成,请检查输入或服务状态'}, 500)
|
|
|
|
|
|
|
|
|
|
def proxy(self, port, path):
|
|
|
|
|
conn = http.client.HTTPConnection('127.0.0.1', port, timeout=25)
|
|
|
|
|
headers = {'Authorization': MEDIA_AUTH}
|
|
|
|
|
if self.headers.get('Range'):
|
|
|
|
|
headers['Range'] = self.headers['Range']
|
|
|
|
|
started = False
|
|
|
|
|
try:
|
|
|
|
|
conn.request('GET', path, headers=headers)
|
|
|
|
|
r = conn.getresponse()
|
|
|
|
|
if r.status >= 400:
|
|
|
|
|
raise Problem('暂时没有可播放画面,请确认设备网络、账号、码流路径和录像时间', 503)
|
2026-10-03 16:45:58 +08:00
|
|
|
location = None
|
|
|
|
|
if 300 <= r.status < 400:
|
|
|
|
|
# MediaMTX starts an HLS session with a relative redirect. The
|
|
|
|
|
# browser must keep that session query on subsequent playlists.
|
|
|
|
|
target = url.urlsplit(url.urljoin(path, r.getheader('Location', '')))
|
|
|
|
|
if (port != HLS_PORT or target.scheme or target.netloc
|
|
|
|
|
or not re.fullmatch(r'/cam_[a-f0-9]{16}_(main|sub)/[a-zA-Z0-9_.-]+', target.path)):
|
|
|
|
|
raise Problem('视频服务返回了无效的播放跳转', 502)
|
|
|
|
|
location = '/media/live' + target.path + ('?' + target.query if target.query else '')
|
2026-10-03 16:12:25 +08:00
|
|
|
self.send_response(r.status)
|
2026-10-03 16:45:58 +08:00
|
|
|
if location:
|
|
|
|
|
self.send_header('Location', location)
|
2026-10-03 16:12:25 +08:00
|
|
|
for k in ('Content-Type', 'Content-Length', 'Content-Range', 'Accept-Ranges'):
|
|
|
|
|
if r.getheader(k):
|
|
|
|
|
self.send_header(k, r.getheader(k))
|
|
|
|
|
self.send_header('Cache-Control', 'no-store')
|
|
|
|
|
self.send_header('X-Content-Type-Options', 'nosniff')
|
|
|
|
|
self.end_headers()
|
|
|
|
|
started = True
|
|
|
|
|
while True:
|
|
|
|
|
chunk = r.read(65536)
|
|
|
|
|
if not chunk:
|
|
|
|
|
break
|
|
|
|
|
self.wfile.write(chunk)
|
|
|
|
|
except (OSError, http.client.HTTPException):
|
|
|
|
|
if not started:
|
|
|
|
|
raise Problem('视频服务连接失败', 503)
|
|
|
|
|
finally:
|
|
|
|
|
conn.close()
|
|
|
|
|
|
|
|
|
|
def do_GET(self):
|
|
|
|
|
try:
|
|
|
|
|
parsed = url.urlsplit(self.path)
|
|
|
|
|
path = parsed.path
|
|
|
|
|
query = {k: v[0] for k, v in url.parse_qs(parsed.query).items()}
|
|
|
|
|
if path == '/api/session':
|
|
|
|
|
self.answer({'authenticated': self.authenticated(), 'setupRequired': not bool(setting('account'))})
|
|
|
|
|
return
|
|
|
|
|
if path == '/healthz':
|
|
|
|
|
self.answer({'status': 'running', 'version': (ROOT / 'VERSION').read_text().strip()})
|
|
|
|
|
return
|
|
|
|
|
if path.startswith(('/api/', '/media/')):
|
|
|
|
|
self.require_auth()
|
|
|
|
|
if path == '/api/state':
|
|
|
|
|
self.answer(status())
|
|
|
|
|
elif path == '/api/audit':
|
|
|
|
|
with LOCK:
|
|
|
|
|
rows = [dict(r) for r in DB.execute('SELECT * FROM audit ORDER BY id DESC LIMIT 50')]
|
|
|
|
|
self.answer(rows)
|
|
|
|
|
elif path == '/api/recordings':
|
|
|
|
|
c = get_object('cameras', query.get('camera'))
|
|
|
|
|
if not c:
|
|
|
|
|
raise Problem('请选择摄像头')
|
|
|
|
|
start, end = parse_time(query.get('start')), parse_time(query.get('end'))
|
|
|
|
|
if not 0 < (end - start).total_seconds() <= 172800:
|
|
|
|
|
raise Problem('每次查询最多两天')
|
|
|
|
|
q = url.urlencode({'path': stream_name(c), 'start': start.isoformat(), 'end': end.isoformat()})
|
|
|
|
|
raw = media_json('/list?' + q, playback=True)
|
|
|
|
|
items = [{'start': r['start'], 'duration': r['duration']} for r in raw]
|
|
|
|
|
self.answer({'items': items})
|
|
|
|
|
elif path.startswith('/media/live/'):
|
|
|
|
|
suffix = path[len('/media/live/'):]
|
|
|
|
|
if not re.fullmatch(r'cam_[a-f0-9]{16}_(main|sub)/[a-zA-Z0-9_.-]+', suffix):
|
|
|
|
|
raise Problem('视频路径不正确')
|
|
|
|
|
self.proxy(HLS_PORT, '/' + suffix + ('?' + parsed.query if parsed.query else ''))
|
|
|
|
|
elif path == '/media/playback':
|
|
|
|
|
c = get_object('cameras', query.get('camera'))
|
|
|
|
|
if not c:
|
|
|
|
|
raise Problem('摄像头不存在', 404)
|
|
|
|
|
start = parse_time(query.get('start'))
|
|
|
|
|
try:
|
|
|
|
|
duration = float(query.get('duration', '300'))
|
|
|
|
|
if not math.isfinite(duration) or not 0 < duration <= 3600:
|
|
|
|
|
raise ValueError()
|
|
|
|
|
except ValueError:
|
|
|
|
|
raise Problem('单次回放最多一小时')
|
|
|
|
|
q = url.urlencode({'path': stream_name(c), 'start': start.isoformat(), 'duration': duration,
|
|
|
|
|
'format': 'mp4' if query.get('download') == '1' else 'fmp4'})
|
|
|
|
|
self.proxy(PLAYBACK_PORT, '/get?' + q)
|
|
|
|
|
else:
|
|
|
|
|
raise Problem('接口不存在', 404)
|
|
|
|
|
else:
|
|
|
|
|
files = {'/': 'index.html', '/app.js': 'app.js', '/style.css': 'style.css',
|
|
|
|
|
'/vendor/hls.min.js': 'vendor/hls.min.js'}
|
|
|
|
|
if path not in files:
|
|
|
|
|
raise Problem('页面不存在', 404)
|
|
|
|
|
file = ROOT / 'web' / files[path]
|
|
|
|
|
raw = file.read_bytes()
|
|
|
|
|
self.send_response(200)
|
|
|
|
|
self.send_header('Content-Type', (mimetypes.guess_type(file.name)[0] or 'application/octet-stream') + '; charset=utf-8')
|
|
|
|
|
self.send_header('Content-Length', str(len(raw)))
|
|
|
|
|
self.send_header('X-Content-Type-Options', 'nosniff')
|
|
|
|
|
self.send_header('X-Frame-Options', 'DENY')
|
|
|
|
|
self.send_header('Referrer-Policy', 'same-origin')
|
|
|
|
|
self.send_header('Content-Security-Policy', "default-src 'self'; script-src 'self'; style-src 'self'; img-src 'self' data:; media-src 'self' blob:; worker-src 'self' blob:; connect-src 'self'; frame-ancestors 'none'; base-uri 'none'; form-action 'self'")
|
|
|
|
|
self.send_header('Cache-Control', 'no-cache')
|
|
|
|
|
self.end_headers()
|
|
|
|
|
self.wfile.write(raw)
|
|
|
|
|
except Problem as e:
|
|
|
|
|
self.answer({'error': str(e)}, e.status)
|
|
|
|
|
except (BrokenPipeError, ConnectionResetError):
|
|
|
|
|
pass
|
|
|
|
|
except Exception:
|
|
|
|
|
self.answer({'error': '请求未完成'}, 500)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def initialize():
|
|
|
|
|
global MEDIA_AUTH
|
|
|
|
|
if not setting('storage'):
|
|
|
|
|
set_setting('storage', {'retentionDays': 7, 'maxGB': 40, 'reserveGB': 8})
|
|
|
|
|
if not setting('mediaSecret'):
|
|
|
|
|
set_setting('mediaSecret', secrets.token_hex(32))
|
|
|
|
|
MEDIA_AUTH = 'Basic ' + base64.b64encode(('vision:' + setting('mediaSecret')).encode()).decode()
|
|
|
|
|
if not setting('account') and not (DATA / 'setup-code.txt').exists():
|
|
|
|
|
(DATA / 'setup-code.txt').write_text(secrets.token_hex(12), encoding='utf8')
|
|
|
|
|
os.chmod(DATA / 'setup-code.txt', 0o600)
|
|
|
|
|
write_media_config()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if __name__ == '__main__':
|
|
|
|
|
initialize()
|
|
|
|
|
threading.Thread(target=storage_watch, daemon=True).start()
|
2026-10-03 16:26:28 +08:00
|
|
|
threading.Thread(target=media_config_watch, daemon=True).start()
|
2026-10-03 16:21:00 +08:00
|
|
|
bindings = [h.strip() for h in os.environ.get('VISION_BIND', '127.0.0.1').split(',') if h.strip()]
|
2026-10-03 16:12:25 +08:00
|
|
|
port = int(os.environ.get('VISION_PORT', '8790'))
|
2026-10-03 16:21:00 +08:00
|
|
|
servers = []
|
|
|
|
|
for bind in dict.fromkeys(bindings):
|
|
|
|
|
server = ThreadingHTTPServer((bind, port), Handler)
|
|
|
|
|
server.daemon_threads = True
|
|
|
|
|
servers.append(server)
|
|
|
|
|
print('ZhaoVision listening on ' + bind + ':' + str(port), flush=True)
|
|
|
|
|
for server in servers[1:]:
|
|
|
|
|
threading.Thread(target=server.serve_forever, daemon=True).start()
|
|
|
|
|
servers[0].serve_forever()
|