Fix onboard WebKit preview and archive board telemetry locally

This commit is contained in:
DCCONSTRUCTIONS
2026-09-08 01:37:52 +03:00
parent d8afb61229
commit f0802d2713
46 changed files with 2623 additions and 40 deletions
@@ -382,6 +382,15 @@ def create_app(repository_root: Path):
peers = NodeMediaPeers(bridge.rerun, bridge.service.camera_preview)
sensor = NodeK1Sensor(bridge, peers)
from fastapi.responses import Response
from k1link.viewer.node_local_media import start_local, read_local
async def preview_request(request):
data = await request.body()
if len(data) > 4096:
raise ValueError("Preview request exceeds bound")
return await request.json()
@asynccontextmanager
async def lifespan(_app):
yield
@@ -398,6 +407,30 @@ def create_app(repository_root: Path):
async def inventory(request: Request):
return await sensor.inventory(request.headers["X-Node-Id"])
@app.post("/local-preview")
async def preview(request: Request):
try:
command = await preview_request(request)
if datetime.fromisoformat(command["deadline_at"]) <= datetime.now(UTC):
raise ValueError("Preview request expired")
await sensor.admit_preview(command, request.headers["X-Node-Id"])
identifier = await start_local(peers, command['parameters'])
data = await read_local(peers, identifier, command['parameters'].get('after', 0))
return Response(data, media_type="application/vnd.missioncore.preview",
headers={"Cache-Control": "no-store"})
except (ValueError, KeyError, TypeError):
return JSONResponse({"error": "Preview unavailable"}, status_code=409)
@app.post("/local-preview/read")
async def preview_read(request: Request):
try:
data = await preview_request(request)
payload = await read_local(peers, data['peer_id'], data['after'])
return Response(payload, media_type="application/vnd.missioncore.preview",
headers={"Cache-Control": "no-store"})
except (ValueError, KeyError, TypeError):
return JSONResponse({"error": "Preview expired"}, status_code=409)
@app.get("/prepare-safe")
async def prepare_safe():
state = await sensor.raw_state()
@@ -132,6 +132,17 @@ class NodeK1Sensor:
"expected_state_revision": control["state_revision"],
}
async def admit_preview(self, command, node_id):
item = project_sensor(await self.raw_state(), node_id)
if (not item or command.get("action_id") != "offer"
or command.get("session", {}).get("device_id") != item["id"]
or command["session"].get("session_id") != item["snapshot"]["context"]["session_id"]
or item["snapshot"]["acquisition"] != "streaming"
or not item["control"]["acquisition_id"]
or command.get("parameters", {}).get("acquisition_id") != item["control"]["acquisition_id"]):
raise ValueError("Acquisition is not active or changed")
return item
async def execute(self, command, node_id):
state = await self.raw_state()
item = project_sensor(state, node_id)
+390
View File
@@ -0,0 +1,390 @@
"""Idempotent local replica of a board-owned Timescale history, behind fleet mTLS."""
import json
import math
import re
import sqlite3
import threading
import time
from uuid import UUID
SCHEMA = "missioncore.node-system-monitor/v1"
KEY = re.compile(r"^[A-Za-z0-9_.:-]{1,180}$")
def checked_sample(value):
if not isinstance(value, dict) or set(value) != {
"seq",
"at",
"boot_id",
"uptime",
"values",
"events",
}:
raise ValueError("Invalid sample")
if type(value["seq"]) is not int or not 0 < value["seq"] < 2**53:
raise ValueError("Invalid sequence")
if str(UUID(value["boot_id"])) != value["boot_id"]:
raise ValueError("Invalid boot")
for name in ["at", "uptime"]:
if (
type(value[name]) not in (int, float)
or not math.isfinite(value[name])
or value[name] < 0
):
raise ValueError("Invalid time")
if not isinstance(value["values"], dict) or len(value["values"]) > 384:
raise ValueError("Metric budget")
for key, number in value["values"].items():
if not KEY.fullmatch(key) or (
number is not None and (type(number) not in (int, float) or not math.isfinite(number))
):
raise ValueError("Invalid metric")
if not isinstance(value["events"], list) or len(value["events"]) > 64:
raise ValueError("Event budget")
for event in value["events"]:
if (
not isinstance(event, dict)
or set(event) - {"code", "kind", "locations", "at"}
or event.get("code")
not in {
"ui-error",
"ui-rejection",
"ui-render-error",
"system-oom",
"gpu-reset",
"disk-error",
"usb-disconnected",
"usb-error",
"service-failed",
"ui-process-exit",
"ui-load-failed",
}
or not re.fullmatch("[A-Za-z]{1,48}", event.get("kind", ""))
or not isinstance(event.get("locations"), list)
or len(event["locations"]) > 8
or any(
not isinstance(v, str)
or not re.fullmatch(r"[A-Za-z0-9_-]{1,96}\.js:[0-9]{1,9}:[0-9]{1,9}", v)
for v in event["locations"]
)
):
raise ValueError("Invalid event")
if (
type(event.get("at")) not in (int, float)
or not math.isfinite(event["at"])
or event["at"] < 0
):
raise ValueError("Invalid event time")
return json.dumps(value, allow_nan=False, separators=(",", ":"), sort_keys=True)
class MonitorReplica:
def __init__(self, root):
path = root / "monitor.sqlite3"
if path.is_symlink():
raise ValueError("Invalid monitor store")
self.db = sqlite3.connect(path, check_same_thread=False, timeout=2)
try:
path.chmod(0o600)
self.lock = threading.RLock()
self.maintenance = 0
self.db.execute("PRAGMA journal_mode=WAL")
self.db.execute("PRAGMA synchronous=FULL")
self.db.execute("PRAGMA max_page_count=524288")
self.db.execute(
"CREATE TABLE IF NOT EXISTS samples(node TEXT,source TEXT,seq INTEGER,at REAL,"
"body TEXT,PRIMARY KEY(node,source,seq))"
)
self.db.execute("CREATE INDEX IF NOT EXISTS samples_time ON samples(node,at)")
self.db.execute(
"CREATE TABLE IF NOT EXISTS minutes(node TEXT,at INTEGER,body TEXT,"
"PRIMARY KEY(node,at))"
)
self.db.execute(
"CREATE TABLE IF NOT EXISTS events(node TEXT,source TEXT,seq INTEGER,at REAL,"
"body TEXT,PRIMARY KEY(node,source,seq))"
)
self.db.execute("CREATE INDEX IF NOT EXISTS events_time ON events(node,at)")
self.db.execute(
"CREATE TABLE IF NOT EXISTS states(node TEXT PRIMARY KEY,source TEXT,"
"body TEXT,received REAL,ack INTEGER)"
)
self.db.commit()
except Exception:
self.db.close()
raise
def ingest(self, node, value):
if not isinstance(value, dict) or value.get("schema") != SCHEMA:
return None
source = value.get("source_id")
if not isinstance(source, str) or str(UUID(source)) != source:
raise ValueError("Invalid source")
definitions = value.get("definitions")
if not isinstance(definitions, dict) or len(definitions) > 384:
raise ValueError("Definition budget")
for key, definition in definitions.items():
if (
not KEY.fullmatch(key)
or not isinstance(definition, dict)
or set(definition) != {"label", "group", "resource", "unit", "reason"}
):
raise ValueError("Invalid definition")
if any(
v is not None and (not isinstance(v, str) or len(v) > 256)
for v in definition.values()
):
raise ValueError("Invalid definition")
if value.get("latest") is not None:
checked_sample(value["latest"])
batch = value.get("batch") or {}
rows = batch.get("samples", [])
if (
not isinstance(rows, list)
or len(rows) > 32
or (rows and (batch.get("source_id") != source or batch.get("schema") != SCHEMA))
):
raise ValueError("Invalid batch")
serialized = [(row, checked_sample(row)) for row in rows]
if sum(len(text.encode()) for _, text in serialized) > 196608:
raise ValueError("Batch budget")
if any(
a[0]["seq"] >= b[0]["seq"] for a, b in zip(serialized, serialized[1:], strict=False)
):
raise ValueError("Unordered batch")
snapshot = {
key: item
for key, item in value.items()
if key
in {
"schema",
"source_id",
"latest",
"definitions",
"storage",
"database_bytes",
"sample_interval_seconds",
"retention_days",
"database_budget_bytes",
"free_reserve_bytes",
}
}
snapshot.update(
retained_since=batch.get("retained_since"), first_seq=batch.get("first_seq")
)
with self.lock, self.db:
previous = self.db.execute(
"SELECT source,ack FROM states WHERE node=?", (node,)
).fetchone()
ack = previous[1] if previous and previous[0] == source else 0
minutes = {}
for row, text in serialized:
old = self.db.execute(
"SELECT body FROM samples WHERE node=? AND source=? AND seq=?",
(node, source, row["seq"]),
).fetchone()
if old and old[0] != text:
raise ValueError("Conflicting history sequence")
self.db.execute(
"INSERT OR IGNORE INTO samples VALUES(?,?,?,?,?)",
(node, source, row["seq"], row["at"], text),
)
if not old:
minute = int(row["at"] // 60) * 60
if minute not in minutes:
saved = self.db.execute(
"SELECT body FROM minutes WHERE node=? AND at=?", (node, minute)
).fetchone()
minutes[minute] = json.loads(saved[0]) if saved else {}
for key, number in row["values"].items():
if number is not None:
aggregate = minutes[minute].setdefault(key, [number, number, 0, 0])
aggregate[0] = min(aggregate[0], number)
aggregate[1] = max(aggregate[1], number)
aggregate[2] += number
aggregate[3] += 1
if row["events"]:
self.db.execute(
"INSERT INTO events VALUES(?,?,?,?,?)",
(node, source, row["seq"], row["at"], json.dumps(row["events"])),
)
ack = max(ack, row["seq"])
for minute, aggregate in minutes.items():
self.db.execute(
"INSERT INTO minutes VALUES(?,?,?) ON CONFLICT(node,at) "
"DO UPDATE SET body=excluded.body",
(node, minute, json.dumps(aggregate)),
)
self.db.execute(
"INSERT INTO states VALUES(?,?,?,?,?) ON CONFLICT(node) DO UPDATE SET "
"source=excluded.source,body=excluded.body,received=excluded.received,ack=excluded.ack",
(node, source, json.dumps(snapshot), time.time(), ack),
)
if time.time() - self.maintenance > 60:
self.db.execute("DELETE FROM samples WHERE at<?", (time.time() - 7 * 86400,))
self.db.execute("DELETE FROM minutes WHERE at<?", (time.time() - 7 * 86400,))
self.db.execute("DELETE FROM events WHERE at<?", (time.time() - 7 * 86400,))
self.maintenance = time.time()
return dict(source_id=source, after=ack)
def query(self, node, metric="cpu.usage", window=900, end=None):
if not KEY.fullmatch(metric) or window not in (900, 3600, 21600, 86400, 604800):
raise ValueError("Invalid range")
end = time.time() if end is None else end
if not math.isfinite(end) or end < 0:
raise ValueError("Invalid time")
coarse = window >= 21600
width = math.ceil(window / 300 / 60) * 60 if coarse else max(1, window // 300)
if coarse:
end = math.floor(end / 60) * 60
start = end - window
path = '$.values."' + metric + '"'
with self.lock:
row = self.db.execute(
"SELECT body,received,ack FROM states WHERE node=?", (node,)
).fetchone()
if not row:
return dict(available=False, series=[], events=[])
snapshot = json.loads(row[0])
if coarse:
aggregate = '$."' + metric + '"'
points = self.db.execute(
"SELECT CAST(at/? AS INTEGER),min(json_extract(body,?)), "
"max(json_extract(body,?)),sum(json_extract(body,?))*1.0/"
"nullif(sum(json_extract(body,?)),0),sum(json_extract(body,?)) "
"FROM minutes WHERE node=? AND at>=? AND at<? "
"GROUP BY CAST(at/? AS INTEGER) ORDER BY 1",
(
width,
aggregate + "[0]",
aggregate + "[1]",
aggregate + "[2]",
aggregate + "[3]",
aggregate + "[3]",
node,
start,
end,
width,
),
).fetchall()
else:
points = self.db.execute(
"SELECT CAST(at/? AS INTEGER),min(json_extract(body,?)), "
"max(json_extract(body,?)),avg(json_extract(body,?)),"
"count(json_extract(body,?)) FROM samples "
"WHERE node=? AND at>=? AND at<=? GROUP BY CAST(at/? AS INTEGER) ORDER BY 1",
(width, path, path, path, path, node, start, end, width),
).fetchall()
event_rows = self.db.execute(
"SELECT body FROM events WHERE node=? AND at>=? AND at<=? "
"ORDER BY at DESC LIMIT 64",
(node, start, end),
).fetchall()
earliest = self.db.execute(
"SELECT min(at) FROM samples WHERE node=?", (node,)
).fetchone()[0]
latest = snapshot.get("latest")
fresh = bool(latest and -5 <= time.time() - latest["at"] < 10 and time.time() - row[1] < 15)
return dict(
available=True,
**snapshot,
fresh=fresh,
received_at=row[1],
replicated_after=row[2],
archive_since=earliest,
backlog=max(0, (latest or {}).get("seq", 0) - row[2]),
start=start,
end=end,
bucket_seconds=width,
series=[
dict(at=bucket * width, min=minimum, max=maximum, mean=mean, count=count)
for bucket, minimum, maximum, mean, count in points
],
events=[event for record in event_rows for event in json.loads(record[0])][:64],
)
def close(self):
with self.lock:
self.db.close()
class MonitorReceiver:
"""Replica I/O never runs under the fleet/control lock, including startup.
Heartbeats may repeat a batch until this worker commits it. Replacing a
queued retry is safe: the sender advances only from the returned commit ACK.
"""
def __init__(self, root):
self.root = root
self.lock = threading.Lock()
self.pending = {}
self.acks = {}
self.archive = None
self.worker = None
self.wake = threading.Event()
self.stop = threading.Event()
def submit(self, node, value):
if not isinstance(value, dict) or value.get("schema") != SCHEMA:
return None
with self.lock:
if node not in self.pending and len(self.pending) >= 32:
return None
self.pending[node] = value
if self.worker is None:
self.worker = threading.Thread(
target=self.run, name="fleet-monitor-replica", daemon=True
)
self.worker.start()
self.wake.set()
ack = self.acks.get(node)
return ack if ack and ack["source_id"] == value.get("source_id") else None
def run(self):
try:
while not self.stop.is_set():
self.wake.wait(1)
with self.lock:
pending, self.pending = self.pending, {}
self.wake.clear()
if self.archive is None:
try:
self.archive = MonitorReplica(self.root)
except (OSError, ValueError, sqlite3.Error):
continue
for node, value in pending.items():
try:
if self.archive is None:
self.archive = MonitorReplica(self.root)
ack = self.archive.ingest(node, value)
with self.lock:
self.acks[node] = ack
except (OSError, ValueError, TypeError, KeyError, sqlite3.Error):
# No ACK: the board retains and retries its own archive.
continue
finally:
if self.archive is not None:
self.archive.close()
def query(self, *args, **kwargs):
with self.lock:
if self.worker is None:
self.worker = threading.Thread(
target=self.run, name="fleet-monitor-replica", daemon=True
)
self.worker.start()
self.wake.set()
try:
if self.archive is not None:
return self.archive.query(*args, **kwargs)
except sqlite3.Error:
pass
return dict(available=False, series=[], events=[])
def close(self):
self.stop.set()
self.wake.set()
if self.worker is not None:
self.worker.join(5)
+11 -1
View File
@@ -33,6 +33,9 @@ class FleetRegistry:
from .device_enrollment import DeviceEnrollment
self.device_enrollment = DeviceEnrollment()
from .monitor import MonitorReceiver
self.monitor = MonitorReceiver(root)
self.trust = CoreTrust(root)
path = root / "fleet.sqlite3"
if path.is_symlink():
@@ -93,6 +96,7 @@ class FleetRegistry:
server.server_close()
with self.lock:
self.db.close()
self.monitor.close()
def listen(self, address):
from .transport import NodeChannelServer
@@ -407,4 +411,10 @@ class FleetRegistry:
}
row["binding"]["client_pem"] = self.trust.leaf(node_id, key)
self.save(row)
return 200, {"ok": True, "client_pem": row["binding"]["client_pem"], **sensor_response, **enrollment_response}
monitor_ack = None
try:
monitor_ack = self.monitor.submit(node_id, value.get("monitor"))
except (ValueError, TypeError, KeyError, sqlite3.Error):
# Telemetry persistence failure must not block device control.
pass
return 200, {"monitor_ack": monitor_ack, "ok": True, "client_pem": row["binding"]["client_pem"], **sensor_response, **enrollment_response}
+137
View File
@@ -0,0 +1,137 @@
"""Local HTTP carrier for the same bounded, acknowledged Node preview delivery.
The Unix-only plugin endpoint is proxied by the authenticated Node UI. It never
starts acquisition, rewrites network settings, or creates a second recorder.
"""
import asyncio
import json
import time
from uuid import UUID, uuid4
from .node_media import MEDIA_PROTOCOL
class LocalConnection:
def __init__(self):
self.queue = asyncio.Queue(maxsize=40)
self.buffered = 0
self.closed = False
def push(self, kind, payload):
if self.closed:
return
if isinstance(payload, str):
payload = payload.encode()
if len(payload) > 32768:
raise ValueError("Local preview message exceeds bound")
frame = bytes([kind]) + len(payload).to_bytes(4, "big") + payload
self.queue.put_nowait(frame)
self.buffered += len(frame)
async def close(self):
self.closed = True
class LocalChannel:
def __init__(self, connection, label):
self.connection, self.label = connection, label
self.closed = False
@property
def readyState(self):
return "closed" if self.closed or self.connection.closed else "open"
@property
def bufferedAmount(self):
return self.connection.buffered
def send(self, payload):
kind = (1 if self.label == "rrd" else 3) + (not isinstance(payload, str))
self.connection.push(kind, payload)
def close(self):
if not self.closed:
self.closed = True
self.connection.push(5 if self.label == "camera" else 6, b"")
def admit_local(parameters):
view_id, after = parameters.get("view_id"), parameters.get("after", 0)
if not isinstance(view_id, str) or str(UUID(view_id)) != view_id:
raise ValueError("Preview view identifier required")
if type(after) is not int or not 0 <= after <= 2**53 - 1:
raise ValueError("Invalid preview cursor")
return view_id, after
async def start_local(peers, parameters):
view_id, after = admit_local(parameters)
for identifier, entry in list(peers.items.items()):
if entry.get("local") and entry["view_id"] == view_id:
await peers.close(identifier)
if len(peers.items) >= 2:
raise ValueError("Close another live viewer")
connection = LocalConnection()
identifier = "peer_" + uuid4().hex
entry = {
"pc": connection,
"seen": time.monotonic(),
"tasks": [],
"labels": {"rrd", "camera"},
"view_id": view_id,
"after": after,
"subscriber": None,
"local": True,
"reading": False,
}
peers.items[identifier] = entry
connection.push(0, json.dumps({"peer_id": identifier, "media_protocol": MEDIA_PROTOCOL}))
for label in ("rrd", "camera"):
entry["tasks"].append(
asyncio.create_task(peers.deliver(identifier, LocalChannel(connection, label)))
)
async def expiry():
while identifier in peers.items:
await asyncio.sleep(5)
if time.monotonic() - entry["seen"] > 30:
await peers.close(identifier)
entry["tasks"].append(asyncio.create_task(expiry()))
return identifier
async def read_local(peers, identifier, sequence):
acknowledge(peers, identifier, sequence)
entry = peers.items[identifier]
if entry["reading"]:
raise ValueError("Concurrent preview read")
entry["reading"] = True
connection = entry["pc"]
try:
try:
first = await asyncio.wait_for(connection.queue.get(), 1)
except TimeoutError:
return b""
frames = [first]
size = len(first)
while size < 131072 and not connection.queue.empty():
frame = connection.queue.get_nowait()
frames.append(frame)
size += len(frame)
connection.buffered -= size
return b"".join(frames)
finally:
entry["reading"] = False
def acknowledge(peers, identifier, sequence):
entry = peers.items.get(identifier)
if not entry or not entry.get("local"):
raise ValueError("Local preview expired")
if type(sequence) is not int or not 0 <= sequence <= 2**53 - 1:
raise ValueError("Invalid preview acknowledgement")
entry["seen"] = time.monotonic()
if entry["subscriber"] is not None:
entry["subscriber"].acknowledge(sequence)
+5 -3
View File
@@ -252,9 +252,11 @@ class NodeMediaPeers:
async def close(self, identifier):
entry = self.items.pop(identifier, None)
if entry:
for task in entry["tasks"]:
if task is not asyncio.current_task():
task.cancel()
tasks = [task for task in entry["tasks"] if task is not asyncio.current_task()]
for task in tasks:
task.cancel()
# A reconnect must not race the previous delivery lease's finally.
await asyncio.gather(*tasks, return_exceptions=True)
with suppress(Exception):
await entry["pc"].close()
+15
View File
@@ -51,6 +51,21 @@ class AddRequest(BaseModel):
router = APIRouter(prefix="/api/v1/fleet", tags=["fleet"])
@router.get("/{vehicle_id}/monitor")
def board_monitor(vehicle_id: str, response: Response,
fleet: Annotated[FleetRegistry, Depends(local_operator)],
metric: str = "cpu.usage", window: int = 900, end: float | None = None):
response.headers["Cache-Control"] = "no-store"
try:
with fleet.lock:
node_id = fleet.find(vehicle_id)["node_id"]
return fleet.monitor.query(node_id, metric, window, end)
except PairingError as error:
raise HTTPException(404, str(error)) from None
except ValueError:
raise HTTPException(400, "Выберите доступный период и показатель.") from None
@router.get("")
def fleet_list(response: Response, fleet: Annotated[FleetRegistry, Depends(local_operator)]):
response.headers["Cache-Control"] = "no-store"