124 lines
6.5 KiB
Python
124 lines
6.5 KiB
Python
"""Capability-driven ONVIF pan/tilt: bounded nudges, never latched movement."""
|
|
import hashlib
|
|
import threading
|
|
import time
|
|
import urllib.parse as url
|
|
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 = {}
|
|
|
|
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', '<m:Category>All</m:Category>')
|
|
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 = '<m:StreamSetup><tt:Stream>RTP-Unicast</tt:Stream><tt:Transport><tt:Protocol>RTSP</tt:Protocol></tt:Transport></m:StreamSetup><m:ProfileToken>' + escape(p.get('token', '')) + '</m:ProfileToken>'
|
|
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', '<m:ConfigurationToken>' + escape(conf.get('token')) + '</m:ConfigurationToken>')
|
|
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 = '<m:ProfileToken>' + escape(info['profile']) + '</m:ProfileToken>'
|
|
def stop():
|
|
self.call(c, info['endpoint'], PTZ, 'Stop', body + '<m:PanTilt>true</m:PanTilt>')
|
|
if action == 'stop':
|
|
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'<m:Velocity><tt:PanTilt x="{x}" y="{y}" space="{SPACE}"/></m:Velocity><m:Timeout>PT1S</m:Timeout>')
|
|
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()
|