Files
ucvl-home-vision/app.py
T

789 lines
36 KiB
Python
Raw Permalink Normal View History

"""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 shlex
import shutil
import socket
import sqlite3
import sys
import threading
import time
import urllib.error
import urllib.parse as url
import urllib.request
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
import onvif
import playback
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:
rows = [dict(json.loads(r['body']), id=r['id']) for r in DB.execute('SELECT * FROM ' + table)]
return sorted(rows, key=lambda row: (row.get('name', ''), row['id']))
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('录像存储来源不正确')
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
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',
}
for quality in ('compat', 'hd'):
paths[stream_name(c, quality)] = {
'source': 'publisher', 'record': False,
'runOnDemand': shlex.join([sys.executable, str(ROOT / 'preview.py'), str(DATA)]) if c['enabled'] else '',
'runOnDemandRestart': True,
'runOnDemandStartTimeout': '45s', 'runOnDemandCloseAfter': '15s',
}
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')]
+ [{'action': 'publish', 'path': '~^cam_[a-f0-9]{16}_(compat|hd)$'}]}],
'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 urllib.error.HTTPError as exc:
if playback and path.startswith('/list?') and exc.code == 404:
try:
detail = json.loads(exc.read(4096))
if detail.get('error') == 'no recording segments found':
return []
except (ValueError, AttributeError):
pass
raise Problem('视频服务暂不可用,或该时段尚无可回放录像', 503)
except Exception:
raise Problem('视频服务暂不可用,或该时段尚无可回放录像', 503)
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)
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=1.5) as s:
s.settimeout(3)
s.sendall(('OPTIONS rtsp://' + str(ip) + ':554/ RTSP/1.0\r\nCSeq: 1\r\nUser-Agent: ZhaoVision\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))
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)
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, download=False):
conn = http.client.HTTPConnection('127.0.0.1', port, timeout=55)
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)
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|compat|hd)/[a-zA-Z0-9_.-]+', target.path)):
raise Problem('视频服务返回了无效的播放跳转', 502)
location = '/media/live' + target.path + ('?' + target.query if target.query else '')
self.send_response(r.status)
if download:
self.send_header('Content-Disposition', 'attachment; filename="recording.mp4"')
if location:
self.send_header('Location', location)
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 compatible_playback(self, path, duration):
started = False
old_timeout = self.connection.gettimeout()
try:
with playback.stream('http://127.0.0.1:' + str(PLAYBACK_PORT) + path, MEDIA_AUTH, duration) as output:
first = output.read1(65536)
if not first:
raise Problem('录像转换失败,请检查录像是否完整及服务器视频驱动', 503)
self.connection.settimeout(15)
self.send_response(200)
self.send_header('Content-Type', 'video/mp4')
self.send_header('Cache-Control', 'no-store')
self.send_header('Accept-Ranges', 'none')
self.send_header('X-Content-Type-Options', 'nosniff')
self.send_header('Connection', 'close')
self.end_headers()
started = True
self.close_connection = True
self.wfile.write(first)
while chunk := output.read1(65536):
self.wfile.write(chunk)
except playback.PlaybackError as exc:
if not started:
raise Problem(str(exc), 503)
except OSError:
if not started:
raise Problem('录像转换连接失败', 503)
finally:
self.connection.settimeout(old_timeout)
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|compat|hd)/[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('单次回放最多一小时')
download = query.get('download') == '1'
q = url.urlencode({'path': stream_name(c), 'start': start.isoformat(), 'duration': duration,
'format': 'mp4' if download else 'fmp4'})
if download:
self.proxy(PLAYBACK_PORT, '/get?' + q, download=True)
else:
self.compatible_playback('/get?' + q, duration)
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()
threading.Thread(target=media_config_watch, daemon=True).start()
bindings = [h.strip() for h in os.environ.get('VISION_BIND', '127.0.0.1').split(',') if h.strip()]
port = int(os.environ.get('VISION_PORT', '8790'))
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()