"""Capability-driven ONVIF pan/tilt: bounded nudges, never latched movement.""" import hashlib import json import threading import time import urllib.parse as url import xml.etree.ElementTree as ET from xml.sax.saxutils import escape import onvif PTZ = 'http://www.onvif.org/ver20/ptz/wsdl' SPACE = 'http://www.onvif.org/ver10/tptz/PanTiltSpaces/VelocityGenericSpace' T = '{' + onvif.SCHEMA + '}' DIRECTIONS = {'left': (-0.25, 0), 'right': (0.25, 0), 'up': (0, 0.25), 'down': (0, -0.25)} class Controller: def __init__(self): self.guard = threading.RLock() self.cache = {} self.locks = {} self.cancellations = {} self.preset_versions = {} def node_key(self, c, info): return c['host'], info['endpoint'][0], info['node'] def read_presets(self, c): info = self.discover(c) node = self.call(c, info['endpoint'], PTZ, 'GetNode', '' + escape(info['node']) + '') maximum = int(node.findtext('.//' + T + 'MaximumNumberOfPresets', '0')) if not 0 < maximum <= 1024: raise ValueError('设备未开放预置位功能') root = self.call(c, info['endpoint'], PTZ, 'GetPresets', '' + escape(info['profile']) + '') items = [] for item in root.findall('.//{' + PTZ + '}Preset'): token = item.get('token', '') if not token or len(token) > 256 or any(p['token'] == token for p in items): raise ValueError('设备返回了无效预置位列表') items.append(dict(token=token, name=item.findtext(T + 'Name', '') or '未命名位置')) positions = [ET.tostring(p, encoding='unicode') for p in root.findall('.//{' + PTZ + '}Preset')] with self.guard: version = self.preset_versions.get(self.node_key(c, info), 0) fingerprint = [info['profile'], info['node'], maximum, sorted(items, key=lambda p: p['token']), sorted(positions), version] revision = hashlib.sha256(json.dumps(fingerprint, sort_keys=True).encode()).hexdigest() return dict(items=items, maximum=maximum, revision=revision, note='预置位保存在摄像机中,两个镜头可能共用云台和位置。'), info def preset_command(self, c, data): action = data.get('action') if action not in ('create', 'remove', 'goto'): raise ValueError('请选择保存当前位置、删除或转到预置位') info = self.discover(c) key = self.node_key(c, info) with self.guard: lock = self.locks.setdefault(key, threading.Lock()) if not lock.acquire(blocking=False): raise ValueError('云台正在执行操作,请稍后再点') try: snapshot, info = self.read_presets(c) if data.get('revision') != snapshot['revision']: raise ValueError('预置位已变化,请重新读取后操作') token = data.get('token') if action != 'create' and not any(p['token'] == token for p in snapshot['items']): raise ValueError('预置位不存在,请重新读取') body = '' + escape(info['profile']) + '' if action == 'goto': cancel = threading.Event() with self.guard: self.cancellations[key] = cancel try: # A finite preset move has no ONVIF Timeout parameter. Bound it on # the server and send Stop even after an ambiguous network failure. self.call(c, info['endpoint'], PTZ, 'GotoPreset', body + '' + escape(token) + '') deadline = time.monotonic() + 10 while time.monotonic() < deadline: if cancel.wait(0.5): return dict(ok=True, completed=False, message='本次预置位转动已停止') status = self.call(c, info['endpoint'], PTZ, 'GetStatus', body) if status.findtext('.//' + T + 'MoveStatus/' + T + 'PanTilt') == 'IDLE': return dict(ok=True, completed=True, message='设备报告转动结束,请核对画面位置') return dict(ok=True, completed=False, message='设备未确认转动结束,已发送停止;请核对画面后重试') finally: try: self.call(c, info['endpoint'], PTZ, 'Stop', body + 'truetrue') finally: with self.guard: self.cancellations.pop(key, None) if action == 'create': name = data.get('name') if not isinstance(name, str) or not 1 <= len(name.strip()) <= 40 or any(ord(x) < 32 for x in name): raise ValueError('位置名称需要 1–40 个字且不包含控制字符') name = name.strip() if any(p['name'].casefold() == name.casefold() for p in snapshot['items']): raise ValueError('位置名称已存在,请换一个名称') if len(snapshot['items']) >= snapshot['maximum']: raise ValueError('设备预置位已满') payload = body + '' + escape(name) + '' with self.guard: self.preset_versions[key] = self.preset_versions.get(key, 0) + 1 response = self.call(c, info['endpoint'], PTZ, 'SetPreset', payload) saved_token = response.findtext('.//{' + PTZ + '}PresetToken') else: with self.guard: self.preset_versions[key] = self.preset_versions.get(key, 0) + 1 self.call(c, info['endpoint'], PTZ, 'RemovePreset', body + '' + escape(token) + '') latest, _ = self.read_presets(c) verified = (not any(p['token'] == token for p in latest['items']) if action == 'remove' else any(p['token'] == saved_token and p['name'] == name for p in latest['items'])) return dict(ok=True, verified=verified, presets=latest, message='设备预置位列表已回读确认' if verified else '设备响应与预置位回读不一致,请核对后再操作') finally: lock.release() def call(self, c, endpoint, namespace, operation, body=''): port, path = endpoint return onvif.request(c['host'], port, path, namespace, operation, body, c['username'], c['password']) @staticmethod def endpoint(c, text): parsed = url.urlsplit(text) if parsed.scheme != 'http' or parsed.hostname != c['host'] or parsed.username or parsed.password: raise ValueError('设备返回了不匹配的云台服务地址') return parsed.port or 80, parsed.path + ('?' + parsed.query if parsed.query else '') def discover(self, c): if c.get('sourceKind', 'direct') != 'direct': raise ValueError('请使用摄像头直连通道控制云台') key = (c['host'], c.get('onvifPort', 80), c['username'], hashlib.sha256(c['password'].encode()).hexdigest(), c['mainPath']) with self.guard: old = self.cache.get(key) if old and old[0] > time.monotonic(): return old[1] caps = self.call(c, (c.get('onvifPort', 80), '/onvif/device_service'), onvif.DEVICE, 'GetCapabilities', 'All') def service(name): value = caps.find('.//' + T + name + '/' + T + 'XAddr') if value is None or not value.text: raise ValueError('设备没有提供云台服务') return self.endpoint(c, value.text) endpoint, media = service('PTZ'), service('Media') profiles = self.call(c, media, onvif.MEDIA, 'GetProfiles') available = [p for p in profiles.findall('.//{' + onvif.MEDIA + '}Profiles') if p.find(T + 'PTZConfiguration') is not None] wanted = url.parse_qs(url.urlsplit(c['mainPath']).query).get('profile', [''])[0] selected = next((p for p in available if p.get('token') == wanted), None) if selected is None: for p in available[:32]: body = 'RTP-UnicastRTSP' + escape(p.get('token', '')) + '' root = self.call(c, media, onvif.MEDIA, 'GetStreamUri', body) uri = root.find('.//' + T + 'Uri') actual, wanted_uri = url.urlsplit(uri.text if uri is not None else ''), url.urlsplit(c['mainPath']) if (actual.hostname == c['host'] and actual.path == wanted_uri.path and sorted(url.parse_qsl(actual.query, keep_blank_values=True)) == sorted(url.parse_qsl(wanted_uri.query, keep_blank_values=True))): selected = p break if selected is None: raise ValueError('这个镜头没有匹配的云台配置') conf = selected.find(T + 'PTZConfiguration') options = self.call(c, endpoint, PTZ, 'GetConfigurationOptions', '' + escape(conf.get('token')) + '') spaces = options.findall('.//' + T + 'ContinuousPanTiltVelocitySpace') space = next((s for s in spaces if s.findtext(T + 'URI') == SPACE), None) if space is None: raise ValueError('设备不支持通用方向微调') for axis in ('XRange', 'YRange'): if float(space.findtext(T + axis + '/' + T + 'Min')) > -0.25 or float(space.findtext(T + axis + '/' + T + 'Max')) < 0.25: raise ValueError('设备云台速度范围不兼容') # Do not issue movement unless the device advertises a one-second watchdog. minimum = options.findtext('.//' + T + 'PTZTimeout/' + T + 'Min') maximum = options.findtext('.//' + T + 'PTZTimeout/' + T + 'Max') def seconds(value): import re m = re.fullmatch(r'PT(?:(\d+(?:\.\d+)?)H)?(?:(\d+(?:\.\d+)?)M)?(?:(\d+(?:\.\d+)?)S)?', value or '') return sum(float(v or 0) * scale for v, scale in zip(m.groups(), (3600, 60, 1))) if m else -1 if not 0 <= seconds(minimum) <= 1 <= seconds(maximum): raise ValueError('设备不支持一秒自动停止保护') result = {'endpoint': endpoint, 'profile': selected.get('token'), 'node': conf.findtext(T + 'NodeToken'), 'supported': True} with self.guard: if len(self.cache) > 128: self.cache.clear() self.cache[key] = (time.monotonic() + 300, result) return result def capabilities(self, c): self.discover(c) return {'supported': True, 'directions': list(DIRECTIONS), 'mode': 'nudge', 'note': '点击一次短距离转动并自动停止。同一摄像头的镜头可能共用云台。'} def command(self, c, action): if action not in (*DIRECTIONS, 'stop'): raise ValueError('请选择上、下、左、右或停止') info = self.discover(c) body = '' + escape(info['profile']) + '' def stop(): self.call(c, info['endpoint'], PTZ, 'Stop', body + 'true') if action == 'stop': with self.guard: cancel = self.cancellations.get(self.node_key(c, info)) if cancel: cancel.set() stop() return {'ok': True, 'stopped': True} # Lock by physical host/node, including commands sent through another lens/account. key = (c['host'], info['endpoint'][0], info['node']) with self.guard: lock = self.locks.setdefault(key, threading.Lock()) if not lock.acquire(blocking=False): raise ValueError('云台正在执行操作,请稍后再点') try: x, y = DIRECTIONS[action] try: self.call(c, info['endpoint'], PTZ, 'ContinuousMove', body + f'PT1S') time.sleep(0.3) finally: # Also stop after an ambiguous network failure. Device timeout is a second guard. stop() return {'ok': True, 'stopped': True} finally: lock.release() CONTROLLER = Controller()