diff --git a/plugins/xgrids-k1/frontend/src/sensors/K1Detail.tsx b/plugins/xgrids-k1/frontend/src/sensors/K1Detail.tsx
index 31b375e..274dcb1 100644
--- a/plugins/xgrids-k1/frontend/src/sensors/K1Detail.tsx
+++ b/plugins/xgrids-k1/frontend/src/sensors/K1Detail.tsx
@@ -1,4 +1,7 @@
-import {useRef,useState} from 'react';
+import {K1SpatialSession, k1SpatialPhasePresentation} from '../components/K1SpatialSession';
+import {deviceTelemetry} from '../spatialSessionModel';
+import type {AcquisitionState} from '../spatialSessionModel';
+import {useEffect,useRef,useState} from 'react';
import {ActivityIndicator,Button,Icon,IconButton,SettingsCard,StatusBadge} from '@nodedc/ui-react';
import {perform,type Sensor,type SensorTransport} from './runtime';
import {K1LiveView} from './K1LiveView';
@@ -11,6 +14,9 @@ import {pendingPreview} from './useK1Preview';
export function K1Detail({device,transport,enabled,back,refresh,failure,createRerunHost}:{device:Sensor;transport:SensorTransport;enabled:boolean;back:()=>void;refresh:()=>Promise
;failure:(error:unknown)=>void;createRerunHost?:RerunHostFactory}){
const [busy,setBusy]=useState(''),[tool,setTool]=useState(null),[operationError,setOperationError]=useState('');
const [preview,setPreview]=useState(pendingPreview);
+ const [sceneOpen,setSceneOpen]=useState(['preparing','starting','streaming','stopping'].includes(device.snapshot.acquisition));
+ const acquisitionId=device.control?.acquisition_id;
+ useEffect(()=>{if(acquisitionId&&['preparing','starting','streaming','stopping'].includes(device.snapshot.acquisition))setSceneOpen(true);},[acquisitionId]);
const scene=useK1SceneSettings(device,transport,failure);
const running=useRef(false);
const failedAction=useRef('');
@@ -19,6 +25,7 @@ export function K1Detail({device,transport,enabled,back,refresh,failure,createRe
const pending=busy||(['preparing','starting'].includes(device.snapshot.acquisition)?'start':device.snapshot.acquisition==='stopping'?'stop':'');
async function act(action:'start'|'stop'|'verify'){
if(running.current||!enabled)return;
+ if(action==='start')setSceneOpen(true);
running.current=true;setBusy(action);failedAction.current='';setOperationError('');failure(null);
try{await perform(transport,device,action,action==='verify'?{}:{operator_confirmed:true,control_generation:device.control?.generation,acquisition_id:device.control?.acquisition_id??null});await refresh();}
catch{failedAction.current=action;setOperationError(action==='start'?'Результат запуска пока не подтверждён.':action==='stop'?'Результат остановки пока не подтверждён.':'Не удалось подтвердить связь с K1.');await refresh().catch(()=>{});}
@@ -26,6 +33,11 @@ export function K1Detail({device,transport,enabled,back,refresh,failure,createRe
}
const confirmed=(failedAction.current==='start'&&streaming)||(failedAction.current==='stop'&&device.control?.phase==='completed')||(failedAction.current==='verify'&&manual.connected);
const summary=(!confirmed&&operationError)||(pending==='start'?'Запускаем устройство. Ожидаем завершения калибровки.':pending==='stop'?'Останавливаем устройство. Ожидаем подтверждения.':!manual.connected?k1ConnectionNotice(device,enabled):status.label);
+ const phaseName = (device.control?.acquisition_phase || (pending==='start'?'starting':pending==='stop'?'stopping':streaming?'acquiring':'completed')) as AcquisitionState;
+ const phase = k1SpatialPhasePresentation({state:phaseName},true) ?? {label:summary,detail:'',busy:!!pending};
+ const sessionControls =
+ {manual.showStop&&:undefined} onClick={()=>void act('stop')}>{pending==='stop'?'Остановка устройства…':'Остановить устройство и запись'}}
+ ;
return
{status.label}
{streaming?`${summary}. ${preview.lidar}. ${preview.camera}.`:summary}
{manual.connected&&setTool(null)} scene={scene} connected={enabled}/>}
- {streaming&&createRerunHost&&}
+ {!sceneOpen&&manual.showStop&&}
+ {sceneOpen&&createRerunHost&&setSceneOpen(false)} sessionControls={sessionControls}/>}
;
}
diff --git a/plugins/xgrids-k1/frontend/src/sensors/K1LiveView.tsx b/plugins/xgrids-k1/frontend/src/sensors/K1LiveView.tsx
index 5fd58a2..bd7f71b 100644
--- a/plugins/xgrids-k1/frontend/src/sensors/K1LiveView.tsx
+++ b/plugins/xgrids-k1/frontend/src/sensors/K1LiveView.tsx
@@ -1,11 +1,12 @@
+import type {ReactNode} from 'react';
import { useCallback, useEffect, useRef, useState } from 'react';
-import { Icon, IconButton, SettingsCard, StatusBadge } from '@nodedc/ui-react';
-import { FloatingMediaWindow, ObservationSourcePicker, ObservationTimeline, SpatialToolbarActions, type ObservationWindowRect, type RerunHostFactory, type SpatialSourceDescriptor } from '@mission-core/sensor-sdk';
+import { ApplicationPanel, Icon, IconButton, StatusBadge } from '@nodedc/ui-react';
+import { SpatialScene, EmptySpatialStage, FloatingMediaWindow, ObservationSourcePicker, ObservationTimeline, SpatialToolbarActions, type ObservationWindowRect, type RerunHostFactory, type SpatialSourceDescriptor } from '@mission-core/sensor-sdk';
import type { Sensor, SensorTransport } from './runtime';
import type { K1SceneState, SceneTool } from './K1SceneWindows';
import { useK1Preview, pendingPreview, type PreviewStatus } from './useK1Preview';
import './livePreview.css';
-export function K1LiveView({ device, transport, createRerunHost, onStatus, enabled, scene, openTool }: {
+export function K1LiveView({ device, transport, createRerunHost, onStatus, enabled, scene, openTool, close, sessionControls }: {
device: Sensor;
transport: SensorTransport;
createRerunHost: RerunHostFactory;
@@ -13,8 +14,10 @@ export function K1LiveView({ device, transport, createRerunHost, onStatus, enabl
enabled: boolean;
scene: K1SceneState;
openTool: (tool: SceneTool) => void;
+ close: () => void; sessionControls: ReactNode;
}) {
const spatial = useRef(null), viewport = useRef(null), video = useRef(null);
+ const [focused, setFocused] = useState(false);
const [expanded, setExpanded] = useState(false), [cameraMaximized, setCameraMaximized] = useState(false), [cameraRect, setCameraRect] = useState();
const [visible, setVisible] = useState(() => new Set(['lidar', 'camera']));
const [generation, setGeneration] = useState(0), [status, setStatus] = useState(pendingPreview);
@@ -31,6 +34,7 @@ export function K1LiveView({ device, transport, createRerunHost, onStatus, enabl
return; const timer = setTimeout(() => { attempts.current = 0; }, 10000); return () => clearTimeout(timer); }, [status.presented]);
useEffect(() => { const key = (event: KeyboardEvent) => { if (event.key === 'Escape') {
setExpanded(false);
+ setFocused(false);
setCameraMaximized(false);
} }; window.addEventListener('keydown', key); return () => window.removeEventListener('keydown', key); }, []);
const sources: SpatialSourceDescriptor[] = [
@@ -41,29 +45,38 @@ export function K1LiveView({ device, transport, createRerunHost, onStatus, enabl
next.delete(id);
else
next.add(id); return next; });
- return setExpanded(value => !value)}>}>
-
-
openTool('source')} openLayers={() => openTool('layers')} openDisplay={() => openTool('display')}/>
-
-
-
-
- {!cameraMaximized &&
}
-
- ВИЗУАЛЬНЫЙ ДВИЖОК{status.presented ? 'Визуализатор готов' : status.retry ? 'Восстановление связи' : 'Ожидаем свежие данные'}
-
-
-
КАДР/С{status.presented ? (device.frames?.pcl_fps ?? 0).toFixed(1) : '—'}
-
До публикации{status.presented && device.frames?.mqtt_to_publish_ms != null ? device.frames.mqtt_to_publish_ms.toFixed(1) : '—'} мс
-
Точек в кадре{status.presented ? status.points.toLocaleString('ru-RU') : '—'}
-
- {status.presented && !cameraMaximized &&
Колесо · зум к курсоруWASD · свободный проход
}
- {!cameraMaximized &&
scene.stage({ accumulationSeconds })} onAccumulationCommit={scene.flush} className="scene-timeline"/>}
- { }} onClose={() => { setCameraMaximized(false); setVisible(current => { const next = new Set(current); next.delete('camera'); return next; }); }} status={{status.cameraPresented ? 'Эфир' : 'Ожидание'}} footer={Бортовой компьютерЭфир без буфера}>
-
- {!status.cameraPresented && {status.camera}
}
-
-
-
- ;
+ const hasScene = status.hasScene;
+ const focusedAny = focused || cameraMaximized;
+ return
{status.presented?'Эфир':'Ожидание эфира'}}>
+ openTool('source')} openLayers={()=>openTool('layers')} openDisplay={()=>openTool('display')}/>}
+ renderer={<>
+
+
+
+ {!hasScene&&
}
+ >}
+ deviceControls={sessionControls}
+ sourceControls={focused?
setFocused(false)}>:!cameraMaximized&&
+
+ {hasScene&&visible.has('lidar')&&setFocused(true)}>}
+
}
+ status={{tone:status.presented?'success':'neutral',label:status.presented?'Визуализатор готов':status.retry?'Восстановление связи':hasScene?'Нет свежих данных':'Ожидание источника'}}
+ metrics={<>
+
КАДР/С{status.presented?(device.frames?.pcl_fps??device.frames?.frame_rate_hz??0).toFixed(1):'—'}
+
Точек{status.presented?status.points.toLocaleString('ru-RU'):'—'}
+
До публикации{status.presented&&device.frames?.mqtt_to_publish_ms!=null?device.frames.mqtt_to_publish_ms.toFixed(1):'—'} мс
+ >}
+ navigationReady={hasScene}
+ timeline={!focusedAny&&
scene.stage({accumulationSeconds})} onAccumulationCommit={scene.flush} className="scene-timeline"/>}
+ media={{}} onClose={()=>{setCameraMaximized(false);setVisible(current=>{const next=new Set(current);next.delete('camera');return next;});}}
+ status={{status.cameraPresented?'Эфир':'Ожидание'}}
+ footer={Бортовой компьютерЭфир без буфера}>
+
+ {!status.cameraPresented&&{status.camera}
}
+ }
+ />
+ ;
}
diff --git a/plugins/xgrids-k1/frontend/src/sensors/previewFrames.ts b/plugins/xgrids-k1/frontend/src/sensors/previewFrames.ts
index 1f6cca1..d08e45a 100644
--- a/plugins/xgrids-k1/frontend/src/sensors/previewFrames.ts
+++ b/plugins/xgrids-k1/frontend/src/sensors/previewFrames.ts
@@ -1,5 +1,5 @@
/** Ordered data-channel fragments are reassembled before a decoder sees them. */
-export const MEDIA_PROTOCOL='missioncore.node-preview/v2';
+export const MEDIA_PROTOCOL='missioncore.node-preview/v3';
const MAX_PAYLOAD=8*1024*1024,FRAGMENT_BYTES=16384;
export function previewFrames(accept:(payload:Uint8Array)=>void) {
let pending:Uint8Array|null=null,offset=0;
@@ -23,3 +23,17 @@ export function previewFrames(accept:(payload:Uint8Array)=>void) {
export function assertRrd(bytes:Uint8Array){
if(bytes.length<12||bytes[0]!==82||bytes[1]!==82||bytes[2]!==70||bytes[3]!==50)throw new Error('Invalid RRD recording');
}
+
+/** Kept by the acquisition view while disposable media peers reconnect. */
+export function previewRecording(send:(payload:Uint8Array)=>void) {
+ let after=0;
+ return {
+ get after(){return after;},
+ accept(sequence:number,payload:Uint8Array) {
+ assertRrd(payload);
+ if(!Number.isSafeInteger(sequence)||sequence<1||sequence>after+1)throw new Error('Preview cursor gap');
+ if(sequence<=after)return false;
+ send(payload);after=sequence;return true;
+ },
+ };
+}
diff --git a/plugins/xgrids-k1/frontend/src/sensors/useK1Preview.ts b/plugins/xgrids-k1/frontend/src/sensors/useK1Preview.ts
index 3c98348..7f1bb11 100644
--- a/plugins/xgrids-k1/frontend/src/sensors/useK1Preview.ts
+++ b/plugins/xgrids-k1/frontend/src/sensors/useK1Preview.ts
@@ -1,7 +1,7 @@
import { useEffect, useState, type RefObject } from 'react';
import type { RerunHostFactory, LiveRerunViewer } from '@mission-core/sensor-sdk';
import { perform, type Sensor, type SensorTransport } from './runtime';
-import { assertRrd, MEDIA_PROTOCOL, previewFrames } from './previewFrames';
+import { MEDIA_PROTOCOL, previewFrames, previewRecording } from './previewFrames';
import { previewCamera } from './previewCamera';
import { previewFreshness } from './previewFreshness';
export type PreviewStatus = {
@@ -9,10 +9,11 @@ export type PreviewStatus = {
camera: string;
retry: boolean;
presented: boolean;
+ hasScene: boolean;
cameraPresented: boolean;
points: number;
};
-export const pendingPreview: PreviewStatus = { lidar: 'Лидар: ожидаем данные', camera: 'Камера: ожидаем изображение', retry: false, presented: false, cameraPresented: false, points: 0 };
+export const pendingPreview: PreviewStatus = { lidar: 'Лидар: ожидаем данные', camera: 'Камера: ожидаем изображение', retry: false, presented: false, hasScene: false, cameraPresented: false, points: 0 };
function privateCandidate(sdp: string): string {
return sdp.split('\r\n').filter(line => {
if (!line.startsWith('a=candidate:'))
@@ -24,8 +25,7 @@ function privateCandidate(sdp: string): string {
}
export function useK1Preview(device: Sensor, transport: SensorTransport, createRerunHost: RerunHostFactory, spatial: RefObject, video: RefObject, onStatus: (value: PreviewStatus) => void, generation = 0, enabled = true) {
const [native, setNative] = useState(null);
- // The native runtime belongs to the view. Media recovery opens a new recording
- // channel inside it; it never downloads/recreates the WASM iframe on each retry.
+ // The runtime and native channel belong to this acquisition view, not its peers.
useEffect(() => {
let active = true;
const host = createRerunHost(spatial.current!);
@@ -43,13 +43,31 @@ export function useK1Preview(device: Sensor, transport: SensorTransport, createR
onStatus({ ...pendingPreview, lidar: 'Визуализатор не загрузился. Обновите страницу.' }); });
return () => { active = false; setNative(null); host.dispose(); };
}, [createRerunHost, spatial, onStatus]);
+ const [view, setView] = useState<{
+ id: string; session: string; acquisition: string | null | undefined; channel: ReturnType;
+ recording: ReturnType; previousRecording: string | null; hasScene: boolean;
+ } | null>(null);
useEffect(() => {
- onStatus(enabled ? { ...pendingPreview } : { ...pendingPreview, lidar: 'Лидар: нет связи с БК', camera: 'Камера: нет связи с БК' });
- if (!native || !enabled)
+ if (!native) return;
+ const id = crypto.randomUUID();
+ const channel = native.open_channel('live-acquisition:' + device.snapshot.context.session_id + ':' + device.control?.acquisition_id + ':' + id);
+ setView({id, session:device.snapshot.context.session_id, acquisition:device.control?.acquisition_id, channel, recording: previewRecording(bytes => {if (!channel.ready) throw new Error('Viewer unavailable'); channel.send_rrd(bytes);}), previousRecording: native.get_active_recording_id(), hasScene: false});
+ return () => { setView(null); channel.close();
+ void perform(transport, device, 'close-peer', {view_id:id, retire_view:true}).catch(()=>{});
+ };
+ }, [native, device.snapshot.context.session_id, device.control?.acquisition_id, transport]);
+ useEffect(() => {
+ const currentView = view?.session === device.snapshot.context.session_id && view?.acquisition === device.control?.acquisition_id;
+ const paused = device.snapshot.acquisition === 'streaming'
+ ? {lidar:'Лидар: нет связи с БК', camera:'Камера: нет связи с БК'}
+ : ['starting','preparing'].includes(device.snapshot.acquisition)
+ ? {} : {lidar:'Лидар: приём завершён', camera:'Камера: изображение остановлено'};
+ onStatus({...pendingPreview, hasScene:currentView ? view?.hasScene ?? false : false, ...(!enabled ? paused : {})});
+ if (!native || !view || !currentView || !enabled)
return;
let active = true, failed = false, peerID: string | undefined;
let keepalive: ReturnType | undefined, follow: ReturnType | undefined, iceTimeout: ReturnType | undefined;
- let status = { ...pendingPreview };
+ let status = { ...pendingPreview, hasScene: view.hasScene };
const update = (patch: Partial) => { if (active) {
const next = { ...status, ...patch };
if (JSON.stringify(next) !== JSON.stringify(status)) {
@@ -62,16 +80,23 @@ export function useK1Preview(device: Sensor, transport: SensorTransport, createR
const rrd = pc.createDataChannel('rrd', { ordered: true }), camera = pc.createDataChannel('camera', { ordered: true });
rrd.binaryType = 'arraybuffer';
camera.binaryType = 'arraybuffer';
- const previousRecording = viewer.get_active_recording_id();
- const channel = viewer.open_channel('live-acquisition:' + device.snapshot.context.session_id + ':' + device.control?.acquisition_id + ':' + generation);
+ const {channel, previousRecording} = view;
const fail = () => { if (!active || failed)
return; failed = true; update({ lidar: 'Лидар: восстанавливаем связь', camera: 'Камера: ожидаем соединение', retry: true, presented: false, cameraPresented: false }); clearInterval(follow); clearInterval(keepalive); pc.close(); };
let cameraDecoded = false, cameraFailed = false, lastVideoTime = -1, lastVideoProgress = 0, seenLidar = false;
const cameraFail = () => { cameraFailed = true; update({ camera: 'Камера: изображение недоступно', cameraPresented: false }); camera.close(); };
- const decoder = previewCamera(video.current!, () => { cameraDecoded = true; lastVideoProgress = performance.now(); }, cameraFail);
+ let decoder = previewCamera(video.current!, () => { cameraDecoded = true; lastVideoProgress = performance.now(); }, cameraFail);
const cameraFrames = previewFrames(payload => decoder.push(payload));
- const rrdFrames = previewFrames(payload => { assertRrd(payload); if (!channel.ready)
- throw new Error('Viewer unavailable'); channel.send_rrd(payload); });
+ let batch: {sequence: number; lidar: Record; received: number} | null = null;
+ const rrdFrames = previewFrames(payload => {
+ if (!batch || !channel.ready) throw new Error('Viewer unavailable');
+ if (view.recording.accept(batch.sequence, payload)) {
+ const age = batch.lidar.age_ms;
+ freshness.receive({...batch.lidar, age_ms: typeof age === 'number' ? age + performance.now() - batch.received : age});
+ }
+ rrd.send(JSON.stringify({ack: batch.sequence}));
+ batch = null;
+ });
rrd.onclose = () => { if (active && !failed)
fail(); };
camera.onclose = () => { if (active && !failed)
@@ -80,8 +105,17 @@ export function useK1Preview(device: Sensor, transport: SensorTransport, createR
if (!active)
return;
try {
- if (typeof event.data === 'string')
- freshness.receive(JSON.parse(event.data));
+ if (typeof event.data === 'string') {
+ const value = JSON.parse(event.data);
+ if (value.type === 'preview-unavailable' && value.code === 'resume-expired') {
+ failed = true; clearInterval(follow); clearInterval(keepalive);
+ update({retry:false, presented:false, cameraPresented:false, lidar:'Сессия просмотра истекла. Закройте и откройте пространственную сцену.'});
+ pc.close(); return;
+ }
+ if (batch || value.type !== 'rrd-batch' || !Number.isSafeInteger(value.sequence) || value.sequence < 1 || !value.lidar)
+ throw new Error('Invalid RRD batch');
+ batch = {...value, received: performance.now()};
+ }
else if (event.data instanceof ArrayBuffer)
rrdFrames.push(event.data);
}
@@ -97,6 +131,9 @@ export function useK1Preview(device: Sensor, transport: SensorTransport, createR
const metadata = JSON.parse(event.data);
if (metadata.type !== 'camera-ready' || typeof metadata.mime !== 'string')
throw new Error('Invalid camera metadata');
+ decoder.close();
+ cameraDecoded = false; cameraFailed = false; lastVideoTime = -1;
+ decoder = previewCamera(video.current!, () => {cameraDecoded = true; lastVideoProgress = performance.now();}, cameraFail);
decoder.open(metadata.mime);
}
else if (event.data instanceof ArrayBuffer)
@@ -128,13 +165,14 @@ export function useK1Preview(device: Sensor, transport: SensorTransport, createR
}
catch { /* Native recording/timeline creation is asynchronous. */ }
seenLidar ||= presented;
+ view.hasScene ||= presented;
const v = video.current!;
if (cameraDecoded && v.currentTime !== lastVideoTime) {
lastVideoTime = v.currentTime;
lastVideoProgress = performance.now();
}
const cameraPresented = cameraDecoded && !cameraFailed && performance.now() - lastVideoProgress < 2000;
- update({ presented, cameraPresented, points: presented ? freshness.points : 0,
+ update({ presented, hasScene: view.hasScene, cameraPresented, points: presented ? freshness.points : 0,
lidar: presented ? 'Лидар: данные поступают' : seenLidar ? 'Лидар: нет свежих данных' : 'Лидар: ожидаем данные',
camera: cameraFailed ? 'Камера: изображение недоступно' : cameraPresented ? 'Камера: изображение поступает' : cameraDecoded ? 'Камера: нет свежих кадров' : 'Камера: ожидаем изображение' });
}, 250);
@@ -154,7 +192,7 @@ export function useK1Preview(device: Sensor, transport: SensorTransport, createR
sdp: string;
type: 'answer';
media_protocol: string;
- }>(transport, device, 'offer', { sdp: privateCandidate(pc.localDescription!.sdp), acquisition_id: device.control?.acquisition_id });
+ }>(transport, device, 'offer', { sdp: privateCandidate(pc.localDescription!.sdp), acquisition_id: device.control?.acquisition_id, view_id: view.id, after: view.recording.after });
peerID = answer.peer_id;
if (!active) {
void perform(transport, device, 'close-peer', { peer_id: peerID }).catch(() => { });
@@ -190,9 +228,8 @@ export function useK1Preview(device: Sensor, transport: SensorTransport, createR
rrdFrames.close();
cameraFrames.close();
decoder.close();
- channel.close();
if (peerID)
void perform(transport, device, 'close-peer', { peer_id: peerID }).catch(() => { });
};
- }, [native, device.snapshot.context.session_id, device.control?.acquisition_id, transport, onStatus, video, generation, enabled]);
+ }, [native, view, device.snapshot.context.session_id, device.control?.acquisition_id, transport, onStatus, video, generation, enabled]);
}
diff --git a/plugins/xgrids-k1/frontend/src/spatialSessionModel.ts b/plugins/xgrids-k1/frontend/src/spatialSessionModel.ts
new file mode 100644
index 0000000..eab0540
--- /dev/null
+++ b/plugins/xgrids-k1/frontend/src/spatialSessionModel.ts
@@ -0,0 +1,84 @@
+/** Pure K1 telemetry and phase vocabulary, shared by direct and onboard views. */
+export type AcquisitionState =
+ | "preparing"
+ | "prepared"
+ | "awaiting_external_start"
+ | "starting"
+ | "acquiring"
+ | "awaiting_external_stop"
+ | "stopping"
+ | "finalizing"
+ | "completed"
+ | "failed"
+ | "aborted"
+ | "interrupted";
+
+export interface XgridsK1Metrics {
+ pcl_frames?: number | null;
+ pose_frames?: number | null;
+ mqtt_to_decode_ms?: number | null;
+ decode_ms?: number | null;
+ publish_ms?: number | null;
+ pipeline_ms?: number | null;
+ end_to_end_ms?: number | null;
+ frame_rate?: number | null;
+ frame_rate_hz?: number | null;
+ point_count?: number | null;
+ dropped_preview_frames?: number | null;
+ ai_end_to_end_ms?: number | null;
+ ai_end_to_end_p95_ms?: number | null;
+ ai_frame_rate_hz?: number | null;
+ ai_dropped_frames?: number | null;
+ ai_stale_ms?: number | null;
+ device_elapsed_seconds?: number | null;
+ device_route_distance_meters?: number | null;
+ device_speed_meters_per_second?: number | null;
+ device_speed_mps?: number | null;
+ elapsed_seconds?: number | null;
+ route_distance_meters?: number | null;
+ speed_meters_per_second?: number | null;
+ [key: string]: number | null | undefined;
+}
+
+export function finiteMetric(value: number | null | undefined): number | null {
+ return typeof value === "number" && Number.isFinite(value) ? value : null;
+}
+
+function nonNegativeMetric(value: number | null | undefined): number | null {
+ const finite = finiteMetric(value);
+ return finite !== null && finite >= 0 ? finite : null;
+}
+
+export interface K1DeviceTelemetry {
+ elapsedSeconds: number | null;
+ routeDistanceMeters: number | null;
+ speedMetersPerSecond: number | null;
+}
+
+
+export function deviceTelemetry(
+ metrics: XgridsK1Metrics | null | undefined,
+): K1DeviceTelemetry {
+ return {
+ elapsedSeconds: nonNegativeMetric(
+ metrics?.device_elapsed_seconds ?? metrics?.elapsed_seconds,
+ ),
+ routeDistanceMeters: nonNegativeMetric(
+ metrics?.device_route_distance_meters ?? metrics?.route_distance_meters,
+ ),
+ speedMetersPerSecond: nonNegativeMetric(
+ metrics?.device_speed_meters_per_second ??
+ metrics?.device_speed_mps ??
+ metrics?.speed_meters_per_second,
+ ),
+ };
+}
+
+export function formatNumber(value: number | null, digits = 1): string {
+ if (value === null) return "—";
+ return value.toLocaleString("ru-RU", {
+ maximumFractionDigits: digits,
+ minimumFractionDigits: digits,
+ });
+}
+
diff --git a/plugins/xgrids-k1/frontend/src/styles.css b/plugins/xgrids-k1/frontend/src/styles.css
index c0bc9da..b268907 100644
--- a/plugins/xgrids-k1/frontend/src/styles.css
+++ b/plugins/xgrids-k1/frontend/src/styles.css
@@ -984,21 +984,6 @@ box-sizing: border-box;
}
}
-.xgrids-k1-spatial-controls {
- display: flex;
- min-width: min(42rem, 100%);
- max-width: 100%;
- align-items: center;
- gap: 0.85rem;
- border: 1px solid rgb(255 255 255 / 0.1);
- border-radius: 1rem;
- background: rgb(9 10 13 / 0.88);
- padding: 0.55rem 0.65rem 0.55rem 0.75rem;
- color: var(--nodedc-text-primary);
- box-shadow: 0 0.9rem 2.4rem rgb(0 0 0 / 0.3);
- backdrop-filter: blur(20px);
-}
-
.xgrids-k1-spatial-controls--recovery {
display: grid;
min-width: min(46rem, 100%);
@@ -1108,90 +1093,6 @@ box-sizing: border-box;
display: none;
}
-.xgrids-k1-spatial-controls__phase {
- display: flex;
- min-width: 11rem;
- flex: 1 1 15rem;
- align-items: center;
- gap: 0.58rem;
-}
-
-.xgrids-k1-spatial-controls__phase > span:last-child {
- display: grid;
- min-width: 0;
- gap: 0.15rem;
-}
-
-.xgrids-k1-spatial-controls__phase strong,
-.xgrids-k1-spatial-controls__phase small {
- overflow: hidden;
- text-overflow: ellipsis;
- white-space: nowrap;
-}
-
-.xgrids-k1-spatial-controls__phase strong {
- font-size: 0.66rem;
-}
-
-.xgrids-k1-spatial-controls__phase small {
- color: var(--nodedc-text-muted);
- font-size: 0.53rem;
-}
-
-.xgrids-k1-spatial-controls__telemetry {
- display: flex;
- flex: 0 1 auto;
- align-items: center;
- gap: 0.65rem;
-}
-
-.xgrids-k1-spatial-controls__telemetry > span {
- display: grid;
- gap: 0.12rem;
- white-space: nowrap;
-}
-
-.xgrids-k1-spatial-controls__telemetry small {
- color: var(--nodedc-text-muted);
- font-size: 0.48rem;
-}
-
-.xgrids-k1-spatial-controls__telemetry strong {
- font-size: 0.61rem;
-}
-
-.xgrids-k1-spatial-controls__error {
- display: grid;
- max-width: 17rem;
- gap: 0.12rem;
- color: rgb(var(--nodedc-danger-rgb));
-}
-
-.xgrids-k1-spatial-controls__error strong {
- font-size: 0.61rem;
-}
-
-.xgrids-k1-spatial-controls__error small {
- overflow: hidden;
- color: var(--nodedc-text-secondary);
- font-size: 0.51rem;
- text-overflow: ellipsis;
- white-space: nowrap;
-}
-
-.xgrids-k1-spatial-controls__action-label {
- display: block;
- width: 9.75rem;
- font-size: 0.66rem;
- line-height: 1.08;
- text-align: center;
- white-space: normal;
-}
-
-.xgrids-k1-spatial-controls__action-label--local {
- width: 8.75rem;
-}
-
.connection-action-progress {
display: flex;
min-height: 2.75rem;
diff --git a/plugins/xgrids-k1/packaging/build_deb.py b/plugins/xgrids-k1/packaging/build_deb.py
index 54f5d33..ec6af48 100644
--- a/plugins/xgrids-k1/packaging/build_deb.py
+++ b/plugins/xgrids-k1/packaging/build_deb.py
@@ -22,7 +22,7 @@ from credential_install import PROFILE_ID, validate # noqa: E402
from debian import package # noqa: E402
from runtime_payload import files as runtime_files # noqa: E402
-VERSION = "0.1.9"
+VERSION = "0.1.10"
RESOURCES = (
"plugins/xgrids-k1/profile_loader.py",
"plugins/xgrids-k1/plugin.manifest.json",
@@ -135,7 +135,7 @@ Architecture: amd64
Maintainer: NODE.DC local build
Section: admin
Priority: optional
-Depends: mission-core-node (>= 0.8.9), mission-core-node (<< 0.9.0),
+Depends: mission-core-node (>= 0.8.10), mission-core-node (<< 0.9.0),
systemd, python3, adduser, bluez, network-manager, iproute2, ffmpeg
Breaks: mission-core-node (<< 0.8.0)
Replaces: mission-core-node (<< 0.8.0)
diff --git a/src/k1link/device_plugins/xgrids_k1/node_sensor.py b/src/k1link/device_plugins/xgrids_k1/node_sensor.py
index d3015e9..bc1432f 100644
--- a/src/k1link/device_plugins/xgrids_k1/node_sensor.py
+++ b/src/k1link/device_plugins/xgrids_k1/node_sensor.py
@@ -98,6 +98,7 @@ def project_sensor(snapshot, node_id):
"network_applied": connected or attempt.get("phase") == "network_applied",
"reason_code": None if connected else attempt.get("public_error_code"),
"acquisition_id": acquisition.get("acquisition_id"),
+ "acquisition_phase": acquisition.get("state"),
},
"live_settings": snapshot.get("viewer_settings", {}),
"frames": snapshot.get("metrics", {}),
@@ -149,6 +150,8 @@ class NodeK1Sensor:
return item
if action == "close-peer":
await self.peers.close(params.get("peer_id"))
+ if params.get("retire_view") is True:
+ await self.peers.release_view(params.get("view_id"))
return {"ok": True}
if action == "offer":
acquisition_id = item["control"]["acquisition_id"]
diff --git a/src/k1link/viewer/node_media.py b/src/k1link/viewer/node_media.py
index b14f9d3..b97dd7e 100644
--- a/src/k1link/viewer/node_media.py
+++ b/src/k1link/viewer/node_media.py
@@ -8,12 +8,14 @@ import queue
import time
from contextlib import suppress
from itertools import chain
-from uuid import uuid4
+from uuid import UUID, uuid4
import aioice.ice
from aiortc import RTCConfiguration, RTCPeerConnection, RTCSessionDescription
-MEDIA_PROTOCOL = "missioncore.node-preview/v2"
+from .node_rerun import PreviewResumeError
+
+MEDIA_PROTOCOL = "missioncore.node-preview/v3"
MAX_PAYLOAD = 8 * 1024 * 1024
FRAGMENT_BYTES = 16384
logger = logging.getLogger(__name__)
@@ -60,11 +62,17 @@ class NodeMediaPeers:
async def offer(self, parameters):
admit_sdp(parameters.get("sdp"))
+ 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")
if len(self.items) >= 2:
raise ValueError("Close another live viewer")
pc = RTCPeerConnection(RTCConfiguration(iceServers=[]))
identifier = "peer_" + uuid4().hex
- entry = {"pc": pc, "seen": time.monotonic(), "tasks": [], "labels": set()}
+ entry = {"pc": pc, "seen": time.monotonic(), "tasks": [], "labels": set(),
+ "view_id": view_id, "after": after, "subscriber": None}
self.items[identifier] = entry
@pc.on("datachannel")
@@ -78,6 +86,14 @@ class NodeMediaPeers:
def message(value):
if value == "keepalive":
entry["seen"] = time.monotonic()
+ elif channel.label == "rrd" and isinstance(value, str) and len(value) < 80:
+ try:
+ ack = json.loads(value)
+ sequence = ack.get("ack")
+ if type(sequence) is int and entry["subscriber"] is not None:
+ entry["subscriber"].acknowledge(sequence)
+ except (ValueError, AttributeError):
+ pass
entry["tasks"].append(asyncio.create_task(self.deliver(identifier, channel)))
@@ -110,7 +126,7 @@ class NodeMediaPeers:
await self.close(identifier)
raise
- async def send(self, channel, payload):
+ async def send(self, channel, payload, alive=None):
if not 0 < len(payload) <= MAX_PAYLOAD:
raise RuntimeError("Preview fragment exceeds bound")
# Each binary_stream.read() is an independent RRD. SCTP messages are
@@ -118,9 +134,9 @@ class NodeMediaPeers:
parts = (payload[offset:offset + FRAGMENT_BYTES]
for offset in range(0, len(payload), FRAGMENT_BYTES))
for part in chain((b"MCF1" + len(payload).to_bytes(4, "big"),), parts):
- deadline = time.monotonic() + 2
- while channel.readyState != "open" or channel.bufferedAmount > 1024 * 1024:
- if channel.readyState in {"closed", "closing"} or time.monotonic() > deadline:
+ while channel.readyState != "open" or channel.bufferedAmount > 256 * 1024:
+ if (channel.readyState in {"closed", "closing"}
+ or (alive is not None and not alive())):
raise RuntimeError("Preview consumer unavailable")
await asyncio.sleep(0.01)
channel.send(part)
@@ -166,32 +182,58 @@ class NodeMediaPeers:
try:
entry = self.items[identifier]
if channel.label == "rrd":
- subscriber = await asyncio.to_thread(self.hub.subscribe)
+ # Admission only takes a short in-process lock. Keep it on this
+ # task so cancellation cannot orphan an attached subscription.
+ subscriber = self.hub.subscribe(entry["view_id"], entry["after"])
+ entry["subscriber"] = subscriber
else:
lease = await self.camera_delivery(identifier, channel)
if lease is None:
return
- while identifier in self.items and time.monotonic() - entry["seen"] < 30:
+ def alive():
+ return identifier in self.items and time.monotonic() - entry["seen"] < 30
+ while alive():
if subscriber:
- payload = await asyncio.to_thread(subscriber.read)
+ batch = subscriber.next_batch(wait=False)
+ if batch is None:
+ break
+ if not batch:
+ await asyncio.sleep(0.025)
+ continue
+ sequence, payload, frame = batch
+ channel.send(json.dumps({"type": "rrd-batch", "sequence": sequence,
+ "lidar": subscriber.snapshot(frame)}))
else:
try:
segment = await asyncio.to_thread(lease.segments.get, 0.5)
except queue.Empty:
continue
if segment is None:
- break
+ self.camera.release_delivery(lease, client_closed=True)
+ lease = None
+ lease = await self.camera_delivery(identifier, channel)
+ if lease is None:
+ break
+ continue
kind, payload = segment
if kind == "media":
self.camera.mark_streaming(lease)
if payload is None:
break
if payload:
- await self.send(channel, payload)
+ await self.send(channel, payload, alive)
if subscriber:
- # Ordered after native bytes. Age is source arrival age,
- # not time spent replaying an encoded preview backlog.
- channel.send(json.dumps(subscriber.snapshot()))
+ # One in-flight complete native batch. Network stalls
+ # pause this disposable delivery; the recorder keeps running.
+ while alive() and subscriber.pending is not None:
+ await asyncio.sleep(0.01)
+ except PreviewResumeError:
+ if channel.readyState == "open":
+ channel.send(json.dumps({"type": "preview-unavailable", "code": "resume-expired"}))
+ # Let the receiver consume the terminal reason and close its peer.
+ deadline = time.monotonic() + 1
+ while channel.readyState == "open" and time.monotonic() < deadline:
+ await asyncio.sleep(0.025)
except asyncio.CancelledError:
pass
except Exception as error:
@@ -199,7 +241,7 @@ class NodeMediaPeers:
channel.label, type(error).__name__)
finally:
if subscriber:
- subscriber.close()
+ subscriber.release()
if lease:
self.camera.release_delivery(lease, client_closed=True)
if channel.label == "camera":
@@ -216,6 +258,14 @@ class NodeMediaPeers:
with suppress(Exception):
await entry["pc"].close()
+ async def release_view(self, view_id):
+ if not isinstance(view_id, str) or str(UUID(view_id)) != view_id:
+ raise ValueError("Invalid preview view identifier")
+ for identifier, entry in list(self.items.items()):
+ if entry["view_id"] == view_id:
+ await self.close(identifier)
+ self.hub.release_view(view_id)
+
async def close_all(self):
for identifier in list(self.items):
await self.close(identifier)
diff --git a/src/k1link/viewer/node_rerun.py b/src/k1link/viewer/node_rerun.py
index af4fe08..b5caa4e 100644
--- a/src/k1link/viewer/node_rerun.py
+++ b/src/k1link/viewer/node_rerun.py
@@ -1,8 +1,8 @@
"""Bounded live RRD publication for paired Node viewers, without a TCP listener.
-Every viewer receives a fresh native recording including StoreInfo/blueprint.
-Latest-value queues discard decoded preview frames before encoding; encoded
-RRD bytes are never dropped inside a stream. Slow viewers are closed instead.
+Each view owns one recording for the acquisition, including across peer recovery.
+Decoded frames coalesce before encoding; a bounded encoded outbox waits for the
+viewer. Only complete, acknowledged RRD batches leave that outbox.
"""
import logging
@@ -17,13 +17,25 @@ from k1link.viewer.rerun_bridge import RerunBridge
MAX_ENCODED_CHUNK = 8 * 1024 * 1024
LIVE_FRAME_MAX_AGE_SECONDS = 2.0
+RESUME_GRACE_SECONDS = 300.0
logger = logging.getLogger(__name__)
+_CURRENT_FRAME = object()
+
+
+class PreviewResumeError(ValueError):
+ """A view must be explicitly reopened; silently resetting history is forbidden."""
class RrdSubscriber:
- def __init__(self, settings_provider):
+ def __init__(self, settings_provider, path_provider=None):
self.closed = threading.Event()
- self.inputs = queue.Queue(maxsize=2)
+ self.inputs = {}
+ self.input_ready = threading.Condition()
+ self.delivery_lock = threading.Lock()
+ self.pending = None
+ self.batch_sequence = 0
+ self.detached_at = None
+ self.path_provider = path_provider
self.output = queue.Queue(maxsize=2)
self.settings_provider = settings_provider
self.last_frame = None
@@ -33,22 +45,55 @@ class RrdSubscriber:
def offer(self, envelope):
if self.closed.is_set():
return
- with suppress(queue.Full):
- if self.inputs.full():
- with suppress(queue.Empty):
- self.inputs.get_nowait()
- self.inputs.put_nowait(envelope)
+ with self.input_ready:
+ # Keep one latest envelope per modality: a pose burst must not
+ # starve PCL, and no network wait reaches the acquisition producer.
+ self.inputs[type(envelope)] = envelope
+ self.input_ready.notify()
- def read(self):
+ def attach(self, after=0):
+ if self.closed.is_set() or not self.delivery_lock.acquire(blocking=False):
+ raise RuntimeError("Preview is closed or still attached")
+ if after != self.batch_sequence and not (
+ self.pending is not None and after == self.batch_sequence - 1):
+ self.delivery_lock.release()
+ raise PreviewResumeError("Preview cursor does not match recording")
+ self.acknowledge(after)
+ self.detached_at = None
+ return self
+
+ def release(self):
+ self.detached_at = time.monotonic()
+ self.delivery_lock.release()
+
+ def next_batch(self, *, wait=True):
+ if self.pending is None:
+ value = self._read(wait=wait)
+ if not value:
+ return value
+ self.batch_sequence += 1
+ self.pending = (self.batch_sequence, *value)
+ return self.pending
+
+ def acknowledge(self, sequence):
+ if self.pending is not None and sequence == self.pending[0]:
+ self.pending = None
+
+ def _read(self, *, wait=True):
if self.closed.is_set():
return None
try:
- return self.output.get(timeout=0.5)
+ return self.output.get(timeout=0.5) if wait else self.output.get_nowait()
except queue.Empty:
return b""
- def snapshot(self):
- frame = self.last_frame
+ def read(self):
+ # Kept for native sink inspection; media delivery uses acknowledged batches.
+ value = self._read()
+ return value[0] if value else value
+
+ def snapshot(self, frame=_CURRENT_FRAME):
+ frame = self.last_frame if frame is _CURRENT_FRAME else frame
return {"type": "lidar-state", "sequence": frame[0] if frame else 0,
"age_ms": max(0, (time.monotonic_ns() - frame[1]) / 1_000_000) if frame else None,
"points": frame[2] if frame else 0}
@@ -66,7 +111,8 @@ class RrdSubscriber:
return "webrtc+rrd://" + str(uuid4())
try:
- bridge = RerunBridge(settings_provider=self.settings_provider, recording_output=output)
+ bridge = PreviewRerunBridge(settings_provider=self.settings_provider,
+ recording_output=output, path_provider=self.path_provider)
bridge.begin_session()
while not self.closed.is_set():
payload = binary.read()
@@ -75,13 +121,21 @@ class RrdSubscriber:
if payload and len(payload) > MAX_ENCODED_CHUNK:
break
if payload:
- # Bound both bytes and waiting time. The archive/producer
- # never waits for this disposable preview subscription.
- self.output.put(payload, timeout=0.5)
- try:
- envelope = self.inputs.get(timeout=0.1)
- except queue.Empty:
- continue
+ # Backpressure is a pause, never EOF. At most two encoded
+ # batches plus this one and the in-flight batch are retained.
+ # Do not drop encoded bytes or restart a native recording.
+ while not self.closed.is_set():
+ try:
+ self.output.put((payload, self.last_frame), timeout=0.1)
+ break
+ except queue.Full:
+ continue
+ with self.input_ready:
+ if not self.inputs:
+ self.input_ready.wait(timeout=0.1)
+ if not self.inputs:
+ continue
+ envelope = self.inputs.pop(next(iter(self.inputs)))
received = envelope.context.received_monotonic_ns
if (received is not None
and (time.monotonic_ns() - received) / 1_000_000_000
@@ -104,10 +158,28 @@ class RrdSubscriber:
binary.read()
+class PreviewRerunBridge(RerunBridge):
+ def __init__(self, *, path_provider=None, **kwargs):
+ self.path_provider = path_provider
+ super().__init__(**kwargs)
+
+ def _append_trajectory_pose(self, position, source_time_ns):
+ if self.path_provider is None:
+ return super()._append_trajectory_pose(position, source_time_ns)
+ # Route history belongs to the acquisition, including movement during
+ # preview congestion. The primary bridge already bounds its point count.
+ path = self.path_provider()
+ if self._path == path:
+ return False
+ self._path = path
+ return True
+
+
class NodeRerunBridge(RerunBridge):
def __init__(self, **kwargs):
self.lock = threading.Lock()
self.subscribers = []
+ self.views = {}
self.latest = {}
def output(recording):
@@ -124,6 +196,7 @@ class NodeRerunBridge(RerunBridge):
self.binary.read()
with self.lock:
self.latest[type(envelope)] = (time.monotonic(), envelope)
+ self._expire_views()
self.subscribers = [v for v in self.subscribers if not v.closed.is_set()]
for subscriber in self.subscribers:
subscriber.offer(envelope)
@@ -132,24 +205,48 @@ class NodeRerunBridge(RerunBridge):
super().process_perception(frame)
self.binary.read()
- def subscribe(self):
+ def _expire_views(self):
+ for key, value in list(self.views.items()):
+ if (value.closed.is_set() or (value.detached_at is not None
+ and time.monotonic() - value.detached_at > RESUME_GRACE_SECONDS)):
+ value.close()
+ del self.views[key]
+
+ def subscribe(self, view_id=None, after=0):
with self.lock:
+ self._expire_views()
+ if self._closed:
+ raise RuntimeError("Live acquisition is not active")
+ if view_id in self.views:
+ return self.views[view_id].attach(after)
+ if after:
+ raise PreviewResumeError("Preview recording expired; reopen the view")
self.subscribers = [v for v in self.subscribers if not v.closed.is_set()]
- if self._closed or len(self.subscribers) >= 2:
- raise RuntimeError("Live viewer unavailable")
- subscriber = RrdSubscriber(self._settings_provider)
+ if len(self.subscribers) >= 2:
+ raise RuntimeError("Close another live viewer")
+ subscriber = RrdSubscriber(self._settings_provider, lambda: list(self._path))
for observed, envelope in self.latest.values():
if time.monotonic() - observed <= LIVE_FRAME_MAX_AGE_SECONDS:
subscriber.offer(envelope)
self.subscribers.append(subscriber)
+ if view_id is not None:
+ self.views[view_id] = subscriber
+ subscriber.attach()
return subscriber
+ def release_view(self, view_id):
+ with self.lock:
+ subscriber = self.views.pop(view_id, None)
+ if subscriber is not None:
+ subscriber.close()
+
def close(self):
with self.lock:
for subscriber in self.subscribers:
subscriber.close()
self.subscribers.clear()
self.latest.clear()
+ self.views.clear()
super().close()
self.binary.read()
@@ -163,8 +260,12 @@ class NodeRerunHub:
self.bridge = bridge
return bridge
- def subscribe(self):
+ def subscribe(self, view_id=None, after=0):
bridge = self.bridge
if bridge is None:
raise RuntimeError("Live acquisition is not active")
- return bridge.subscribe()
+ return bridge.subscribe(view_id, after)
+
+ def release_view(self, view_id):
+ if self.bridge is not None:
+ self.bridge.release_view(view_id)
diff --git a/tests/test_k1_installer.py b/tests/test_k1_installer.py
index aba89a6..11e275a 100644
--- a/tests/test_k1_installer.py
+++ b/tests/test_k1_installer.py
@@ -158,7 +158,7 @@ def test_private_release_contains_material_only_in_root_private_member(
position += 60 + length + length % 2
with tarfile.open(fileobj=io.BytesIO(members["control.tar.gz"]), mode="r:gz") as archive:
control = archive.extractfile("control").read().decode()
- assert "Depends: mission-core-node (>= 0.8.9)" in control
+ assert "Depends: mission-core-node (>= 0.8.10)" in control
assert "Replaces: mission-core-node (<< 0.8.0)" in control
diff --git a/tests/test_node_media.py b/tests/test_node_media.py
index 6f6a25e..64389cc 100644
--- a/tests/test_node_media.py
+++ b/tests/test_node_media.py
@@ -1,5 +1,6 @@
import asyncio
import queue
+from uuid import uuid4
import pytest
from aiortc import RTCConfiguration, RTCPeerConnection, RTCSessionDescription
@@ -36,7 +37,7 @@ def test_failed_media_offer_retires_peer_without_camera_channel(monkeypatch):
peers = NodeMediaPeers(None, Camera())
try:
with pytest.raises(ValueError):
- await peers.offer({"sdp": SDP_HEADER})
+ await peers.offer({"sdp": SDP_HEADER, "view_id": str(uuid4())})
assert peers.items == {}
finally:
await peers.close_all()
@@ -50,23 +51,28 @@ def test_native_webrtc_roundtrip_and_missing_camera_preserve_rrd(monkeypatch):
class Subscription:
def __init__(self):
+ self.pending = None
self.output = queue.Queue()
self.output.put(payload)
- def read(self):
+ def next_batch(self, *, wait=True):
try:
- return self.output.get(timeout=0.1)
+ self.pending = (1, self.output.get(timeout=0.1), None)
+ return self.pending
except queue.Empty:
return b""
- def snapshot(self):
+ def snapshot(self, frame=None):
return {"type": "lidar-state", "sequence": 1, "age_ms": 0, "points": 5000}
- def close(self):
+ def acknowledge(self, sequence):
+ self.pending = None
+
+ def release(self):
pass
class Hub:
- def subscribe(self):
+ def subscribe(self, view_id, after):
return Subscription()
class Camera:
@@ -98,11 +104,14 @@ def test_native_webrtc_roundtrip_and_missing_camera_preserve_rrd(monkeypatch):
return
payloads.append(data)
if len(payloads) > 1 and sum(map(len, payloads[1:])) == len(payload):
+ channel.send('{"ack":1}')
received.set()
try:
await client.setLocalDescription(await client.createOffer())
- answer = await peers.offer({"sdp": client.localDescription.sdp})
+ answer = await peers.offer({
+ "sdp": client.localDescription.sdp, "view_id": str(uuid4()),
+ })
await client.setRemoteDescription(
RTCSessionDescription(sdp=answer["sdp"], type="answer")
)
@@ -210,14 +219,23 @@ def test_native_rrd_idle_does_not_close_camera_or_peer(monkeypatch):
video = client.createDataChannel("camera", ordered=True)
ready, resumed, camera_frame = asyncio.Event(), asyncio.Event(), asyncio.Event()
+ batch_sequence, remaining = 0, 0
+
@rrd.on("message")
def rrd_message(data):
+ nonlocal batch_sequence, remaining
if isinstance(data, str):
value = json.loads(data)
- if value["sequence"] == 1:
+ batch_sequence = value["sequence"]
+ if value["lidar"]["sequence"] == 1:
resumed.set()
+ elif data.startswith(b"MCF1") and remaining == 0:
+ remaining = int.from_bytes(data[4:], "big")
else:
- ready.set()
+ remaining -= len(data)
+ if remaining == 0:
+ rrd.send(json.dumps({"ack": batch_sequence}))
+ ready.set()
@video.on("message")
def camera_message(data):
@@ -226,7 +244,9 @@ def test_native_rrd_idle_does_not_close_camera_or_peer(monkeypatch):
try:
await client.setLocalDescription(await client.createOffer())
- answer = await peers.offer({"sdp": client.localDescription.sdp})
+ answer = await peers.offer({
+ "sdp": client.localDescription.sdp, "view_id": str(uuid4()),
+ })
await client.setRemoteDescription(
RTCSessionDescription(sdp=answer["sdp"], type="answer")
)
@@ -251,3 +271,33 @@ def test_native_rrd_idle_does_not_close_camera_or_peer(monkeypatch):
assert not peers.items
asyncio.run(run())
+
+
+def test_network_backpressure_longer_than_two_seconds_is_a_pause():
+ """Bounded synthetic scheduler pause; no sockets, scanner or load generation."""
+ import time
+
+ class Channel:
+ readyState = "open"
+ sent = []
+ blocked = True
+
+ @property
+ def bufferedAmount(self):
+ return 1024 * 1024 if self.blocked else 0
+
+ def send(self, value):
+ self.sent.append(value)
+
+ async def run():
+ channel = Channel()
+ peers = object.__new__(NodeMediaPeers)
+ task = asyncio.create_task(peers.send(channel, b"small-synthetic-payload", lambda: True))
+ started = time.monotonic()
+ await asyncio.sleep(2.1)
+ assert not task.done() and channel.sent == []
+ channel.blocked = False
+ await asyncio.wait_for(task, 1)
+ assert time.monotonic() - started >= 2
+ assert channel.sent == [b"MCF1\x00\x00\x00\x17", b"small-synthetic-payload"]
+ asyncio.run(run())
diff --git a/tests/test_node_rerun.py b/tests/test_node_rerun.py
index d687983..2acb4aa 100644
--- a/tests/test_node_rerun.py
+++ b/tests/test_node_rerun.py
@@ -78,3 +78,63 @@ def test_reopened_viewer_does_not_replay_cached_points_after_source_pause(monkey
finally:
bridge.close()
sub.thread.join(timeout=3)
+
+
+def test_slow_consumer_resumes_same_recording_and_replays_unacknowledged_batch():
+ """A >500ms delivery pause used to retire the recording and lose its route."""
+ bridge = NodeRerunBridge()
+ sub = bridge.subscribe("synthetic-view")
+ try:
+ batch = sub.next_batch()
+ assert batch and batch[1].startswith(b"RRF2")
+ sequence = batch[0]
+ # Produce several tiny frames without draining the bounded outbox.
+ for index in range(1, 8):
+ bridge.process(points(index))
+ time.sleep(0.12)
+ assert not sub.closed.is_set()
+ assert sub.output.qsize() <= 2
+ old_thread = sub.thread
+ sub.release()
+ resumed = bridge.subscribe("synthetic-view", sequence - 1)
+ assert resumed is sub and resumed.thread is old_thread
+ assert resumed.next_batch() == batch # ACK lost: resend the exact RRD.
+ resumed.acknowledge(sequence)
+ following = resumed.next_batch()
+ assert following[0] == sequence + 1
+ resumed.release()
+ again = bridge.subscribe("synthetic-view", following[0])
+ assert again is sub and again.pending is None # Delivered ACK lost at sender.
+ again.release()
+ finally:
+ bridge.close()
+ sub.thread.join(timeout=3)
+ assert not sub.thread.is_alive()
+
+
+def test_resumption_cannot_silently_replace_expired_recording():
+ import pytest
+ bridge = NodeRerunBridge()
+ try:
+ with pytest.raises(ValueError, match="expired"):
+ bridge.subscribe("missing-view", 2)
+ assert bridge.subscribers == []
+ finally:
+ bridge.close()
+
+
+def test_closing_view_releases_capacity_without_waiting_for_resume_grace():
+ bridge = NodeRerunBridge()
+ views = []
+ try:
+ for index in range(4):
+ sub = bridge.subscribe(f"view-{index}")
+ views.append(sub)
+ sub.release()
+ bridge.release_view(f"view-{index}")
+ assert sub.closed.is_set()
+ assert bridge.views == {}
+ finally:
+ bridge.close()
+ for sub in views:
+ sub.thread.join(timeout=3)