Files
ucvl-home-vision/network_status.py
T

105 lines
5.5 KiB
Python
Raw Permalink Normal View History

"""Read-only, scoped network checks for registered equipment. Never logs credentials."""
import ipaddress
import json
import re
import socket
import subprocess
import threading
import time
from concurrent.futures import ThreadPoolExecutor
def safe_ip(value):
ip = ipaddress.ip_address(value)
return not (ip.is_loopback or ip.is_link_local or ip.is_multicast or ip.is_unspecified or ip.is_reserved)
def probe(host):
"""Connect only to resolved, validated addresses, preventing DNS rebinding."""
try:
addresses = list(dict.fromkeys(x[4][0] for x in socket.getaddrinfo(host, None, type=socket.SOCK_STREAM)))[:3]
if not addresses or not all(safe_ip(ip) for ip in addresses):
return dict(host=host, reachable=False, ports=[], error='地址不能用于设备网络探测')
except (OSError, ValueError):
return dict(host=host, reachable=False, ports=[], error='地址解析失败')
def connect(port):
for ip in addresses:
try:
with socket.create_connection((ip, port), timeout=.55): return port
except OSError: pass
return None
with ThreadPoolExecutor(max_workers=5) as pool:
ports = [port for port in pool.map(connect, (80, 443, 22, 554, 8790)) if port]
return dict(host=host, reachable=bool(ports), ports=ports,
error='' if ports else '常用服务端口未响应,设备可能离线或未开放这些服务')
def tail_peer(host, cidr):
if not host: return None
try:
result = subprocess.run(['tailscale', 'status', '--json'], capture_output=True, text=True, timeout=2, check=True)
status = json.loads(result.stdout)
peer = next((p for p in (status.get('Peer') or {}).values()
if host in p.get('TailscaleIPs', []) or host.lower().rstrip('.') == p.get('DNSName', '').lower().rstrip('.')), None)
if not peer: return dict(available=False, note='本机 Tailscale 未找到此节点')
routes = list(dict.fromkeys((peer.get('PrimaryRoutes') or [])+(peer.get('AllowedIPs') or [])))
# Online is control-plane presence; it does not prove the subnet is reachable.
return dict(available=True, online=bool(peer.get('Online')), active=bool(peer.get('Active')),
subnet=cidr, subnetAvailable=(cidr in routes) if cidr else None)
except (OSError, subprocess.SubprocessError, ValueError, TypeError):
return dict(available=False, note='本机暂时无法读取 Tailscale 状态')
def management_links(results):
links=[]
for item in results:
host=item['host']; authority='['+host+']' if ':' in host else host
try:
if not safe_ip(host):continue
except ValueError:
if not re.fullmatch(r'[A-Za-z0-9.-]{1,253}',host):continue
ports=item.get('ports',[])
port=next((p for p in (443,80,8790) if p in ports),80)
scheme='https' if port==443 else 'http'
links.append(dict(kind=item['kind'],url=scheme+'://'+authority+(':'+str(port) if port==8790 else '')+'/',port=port,verified=port in ports))
return links
class NetworkStatus:
def __init__(self, app):
self.a=app; self.cache={}; self.lock=threading.Lock(); self.gates={}
def read(self, user, query):
a=self.a
if not a.HOUSEHOLDS.manager(user): raise a.Problem('此操作需要家庭管理员权限',403)
table={'device':'devices','camera':'assets','recorder':'recorders'}.get(query.get('source'))
if not table: raise a.Problem('设备类型不正确')
item=next((d for d in a.HOUSEHOLDS.objects(user,table) if d['id']==query.get('id')),None)
if not item: raise a.Problem('设备不存在或未获授权',404)
hosts=[]
direct=item.get('host') if table=='recorders' else item.get('managementHost')
if table=='assets':
direct=next((c.get('host') for c in a.HOUSEHOLDS.objects(user,'cameras') if c.get('assetId')==item['id'] and c.get('sourceKind')!='recorder' and c.get('host')),direct)
for kind,host in [('local',direct),('tailscale',item.get('tailscaleHost'))]:
if host and host not in [h[1] for h in hosts]: hosts.append((kind,host))
key=(a.HOUSEHOLDS.family_of(item),table,item['id'],tuple(hosts),item.get('networkCidr',''))
with self.lock:
if key in self.cache and time.monotonic()-self.cache[key][0]<30:return {**self.cache[key][1],'cached':True}
# Cap concurrency and collapse repeated requests to the same device.
if key in self.gates:raise a.Problem('正在检查该设备,请稍后再试',429)
if len(self.gates)>=4:raise a.Problem('网络检查繁忙,请稍后再试',429)
self.gates[key]=True
try:
with ThreadPoolExecutor(max_workers=3) as pool:
futures=[(kind,pool.submit(probe,host)) for kind,host in hosts]
tail=pool.submit(tail_peer,item.get('tailscaleHost',''),item.get('networkCidr',''))
targets=[dict(f.result(),kind=kind) for kind,f in futures]
result=dict(id=item['id'],source=query['source'],checkedAt=time.time(),targets=targets,tailscale=tail.result(),links=management_links(targets),cached=False)
with self.lock:
self.cache={k:v for k,v in self.cache.items() if time.monotonic()-v[0]<30}
if len(self.cache)>=256:self.cache.pop(next(iter(self.cache)))
self.cache[key]=(time.monotonic(),result)
return result
finally:
with self.lock:self.gates.pop(key,None)