diff --git a/apps/control-station/src/components/RecordedFmp4Player.tsx b/apps/control-station/src/components/RecordedFmp4Player.tsx index 90db69e..6d48de7 100644 --- a/apps/control-station/src/components/RecordedFmp4Player.tsx +++ b/apps/control-station/src/components/RecordedFmp4Player.tsx @@ -28,7 +28,82 @@ export interface RecordedMediaArchive { export type RecordedMediaPresentationState = "loading" | "ready" | "waiting" | "error"; export const RECORDED_MEDIA_DURATION_TOLERANCE_SECONDS = 1; +const RECORDED_MEDIA_SOURCE_OPEN_TIMEOUT_MS = 10_000; +const RECORDED_MEDIA_TARGET_TIMEOUT_MS = 10_000; +const RECORDED_MEDIA_FRAGMENT_TIMEOUT_MS = 15_000; +const RECORDED_MEDIA_REQUIRED_AHEAD_SEGMENTS = 12; +const RECORDED_MEDIA_SEGMENTS_AHEAD = 36; +const RECORDED_MEDIA_MAX_RESIDENT_SEGMENTS = 160; +const RECORDED_MEDIA_RETAIN_BEHIND_SEGMENTS = 60; let recordedMediaWorkerGeneration = 0; +const recordedMediaSourceOwners = new WeakMap(); + +function claimRecordedMediaSource(video: HTMLVideoElement): object { + const owner = {}; + recordedMediaSourceOwners.set(video, owner); + return owner; +} + +function releaseRecordedMediaSource(video: HTMLVideoElement, owner: object): boolean { + if (recordedMediaSourceOwners.get(video) !== owner) return false; + recordedMediaSourceOwners.delete(video); + return true; +} + +export function recordedMediaDecodeStartSequence( + randomAccessSequences: readonly number[], + targetSequence: number, +): number | null { + if (!Number.isInteger(targetSequence) || targetSequence < 1) return null; + let left = 0; + let right = randomAccessSequences.length - 1; + let selected: number | null = null; + while (left <= right) { + const middle = Math.floor((left + right) / 2); + const sequence = randomAccessSequences[middle]; + if (!Number.isInteger(sequence) || sequence < 1) return null; + if (sequence <= targetSequence) { + selected = sequence; + left = middle + 1; + } else { + right = middle - 1; + } + } + return selected; +} + +export function recordedMediaSegmentAppendOrder( + appended: ReadonlySet, + decodeStartSequence: number, + desiredEndSequence: number, +): number[] { + if ( + !Number.isInteger(decodeStartSequence) + || !Number.isInteger(desiredEndSequence) + || decodeStartSequence < 1 + || desiredEndSequence < decodeStartSequence + ) return []; + const missing: number[] = []; + for (let sequence = decodeStartSequence; sequence <= desiredEndSequence; sequence += 1) { + if (!appended.has(sequence)) missing.push(sequence); + } + return missing; +} + +function recordedMediaTimeRangesContain( + ranges: TimeRanges, + targetSeconds: number, + toleranceSeconds = 0.15, +): boolean { + if (!Number.isFinite(targetSeconds)) return false; + for (let index = 0; index < ranges.length; index += 1) { + if ( + ranges.start(index) - toleranceSeconds <= targetSeconds + && ranges.end(index) + toleranceSeconds >= targetSeconds + ) return true; + } + return false; +} export function recordedMediaPresentationState( state: "loading" | "ready" | "error", @@ -211,13 +286,15 @@ async function mountRecordedEpochStream( // The generation token binds the native media request to the exact manifest. // The browser range-streams the virtual init+fragment file from the server; // no complete camera archive is copied into JavaScript memory. + video.preload = "auto"; + const sourceOwner = claimRecordedMediaSource(video); + video.src = descriptor.streamUrl; const cleanup = () => { + if (!releaseRecordedMediaSource(video, sourceOwner)) return; video.pause(); video.removeAttribute("src"); video.load(); }; - video.preload = "auto"; - video.src = descriptor.streamUrl; video.load(); try { await waitForSeekableArchive( @@ -232,22 +309,36 @@ async function mountRecordedEpochStream( } } +interface RecordedSegmentTarget { + readonly revision: number; + readonly sequence: number; + readonly decodeStart: number; + readonly readyEnd: number; + readonly desiredEnd: number; + readonly targetSeconds: number; + forceReset: boolean; + resetAttempts: number; +} + interface RecordedSegmentStreamRuntime { readonly generation: string; readonly mediaSource: MediaSource; readonly sourceBuffer: SourceBuffer; + readonly video: HTMLVideoElement; readonly objectUrl: string; readonly epoch: ObservationRecordedMediaEpoch; readonly contractGeneration: string; readonly segmentCount: number; + readonly randomAccessSequences: readonly number[]; + readonly segmentEndTimesSeconds: readonly number[]; readonly appended: Set; readonly abort: AbortController; - desiredStart: number; - desiredEnd: number; - target: number; + target: RecordedSegmentTarget | null; + fetchAbort: AbortController | null; + notifiedRevision: number; pumping: boolean; disposed: boolean; - onTargetReady: ((sequence: number) => void) | null; + onTargetBuffered: ((target: RecordedSegmentTarget) => void) | null; } export function recordedMediaFragmentUrl( @@ -269,7 +360,12 @@ function waitForMediaSourceOpen(mediaSource: MediaSource, signal: AbortSignal): if (signal.aborted) return Promise.reject(new DOMException("Aborted", "AbortError")); if (mediaSource.readyState === "open") return Promise.resolve(); return new Promise((resolve, reject) => { + const timeout = globalThis.setTimeout(() => { + cleanup(); + reject(new Error("Recorded fragment MediaSource did not open")); + }, RECORDED_MEDIA_SOURCE_OPEN_TIMEOUT_MS); const cleanup = () => { + globalThis.clearTimeout(timeout); mediaSource.removeEventListener("sourceopen", onOpen); mediaSource.removeEventListener("sourceclose", onClose); signal.removeEventListener("abort", onAbort); @@ -328,15 +424,170 @@ function appendRecordedFragment( }); } -async function fetchRecordedFragment(url: string, signal: AbortSignal): Promise { - const response = await fetch(url, { - method: "GET", - cache: "force-cache", - headers: { Accept: "video/mp4" }, - signal, +function removeRecordedMediaRange( + sourceBuffer: SourceBuffer, + startSeconds: number, + endSeconds: number, + signal: AbortSignal, +): Promise { + if (signal.aborted) return Promise.reject(new DOMException("Aborted", "AbortError")); + if (!Number.isFinite(startSeconds) || !Number.isFinite(endSeconds) || endSeconds <= startSeconds) { + return Promise.resolve(); + } + return new Promise((resolve, reject) => { + const cleanup = () => { + sourceBuffer.removeEventListener("updateend", onUpdateEnd); + sourceBuffer.removeEventListener("error", onError); + signal.removeEventListener("abort", onAbort); + }; + const onUpdateEnd = () => { + cleanup(); + resolve(); + }; + const onError = () => { + cleanup(); + reject(new Error("Recorded fragment eviction failed")); + }; + const onAbort = () => { + cleanup(); + reject(new DOMException("Aborted", "AbortError")); + }; + sourceBuffer.addEventListener("updateend", onUpdateEnd, { once: true }); + sourceBuffer.addEventListener("error", onError, { once: true }); + signal.addEventListener("abort", onAbort, { once: true }); + try { + sourceBuffer.remove(startSeconds, endSeconds); + } catch (error) { + cleanup(); + reject(error); + } }); - if (!response.ok) throw new Error(`Recorded fragment returned HTTP ${response.status}`); - return response.arrayBuffer(); +} + +async function fetchRecordedFragment(url: string, signal: AbortSignal): Promise { + if (signal.aborted) throw new DOMException("Aborted", "AbortError"); + const controller = new AbortController(); + const abort = () => controller.abort(); + const timeout = globalThis.setTimeout( + () => controller.abort(), + RECORDED_MEDIA_FRAGMENT_TIMEOUT_MS, + ); + signal.addEventListener("abort", abort, { once: true }); + try { + const response = await fetch(url, { + method: "GET", + cache: "force-cache", + headers: { Accept: "video/mp4" }, + signal: controller.signal, + }); + if (!response.ok) throw new Error(`Recorded fragment returned HTTP ${response.status}`); + return await response.arrayBuffer(); + } catch (error) { + if (controller.signal.aborted && !signal.aborted) { + throw new Error("Recorded fragment request timed out"); + } + throw error; + } finally { + globalThis.clearTimeout(timeout); + signal.removeEventListener("abort", abort); + } +} + +function recordedSegmentChainComplete( + appended: ReadonlySet, + target: RecordedSegmentTarget, +): boolean { + for (let sequence = target.decodeStart; sequence <= target.readyEnd; sequence += 1) { + if (!appended.has(sequence)) return false; + } + return true; +} + +function recordedSegmentStartSeconds( + segmentEndTimesSeconds: readonly number[], + sequence: number, +): number | null { + if (!Number.isInteger(sequence) || sequence < 1 || sequence > segmentEndTimesSeconds.length) { + return null; + } + return sequence === 1 ? 0 : segmentEndTimesSeconds[sequence - 2] ?? null; +} + +function recordedSegmentTargetBuffered( + runtime: RecordedSegmentStreamRuntime, + target: RecordedSegmentTarget, +): boolean { + const readyEndSeconds = runtime.segmentEndTimesSeconds[target.readyEnd - 1]; + if (!Number.isFinite(readyEndSeconds)) return false; + const readyProbeSeconds = Math.max(target.targetSeconds, readyEndSeconds - 0.001); + return ( + recordedSegmentChainComplete(runtime.appended, target) + && recordedMediaTimeRangesContain(runtime.sourceBuffer.buffered, target.targetSeconds, 0.02) + && recordedMediaTimeRangesContain(runtime.sourceBuffer.buffered, readyProbeSeconds, 0.02) + ); +} + +async function resetRecordedSegmentBuffer(runtime: RecordedSegmentStreamRuntime): Promise { + const duration = runtime.mediaSource.duration; + if (runtime.sourceBuffer.buffered.length && Number.isFinite(duration) && duration > 0) { + await removeRecordedMediaRange(runtime.sourceBuffer, 0, duration, runtime.abort.signal); + } + runtime.appended.clear(); +} + +async function evictRecordedSegmentHistory( + runtime: RecordedSegmentStreamRuntime, + target: RecordedSegmentTarget, +): Promise { + if (runtime.appended.size <= RECORDED_MEDIA_MAX_RESIDENT_SEGMENTS) return; + const retainedTarget = Math.max(1, target.sequence - RECORDED_MEDIA_RETAIN_BEHIND_SEGMENTS); + const retainedStart = recordedMediaDecodeStartSequence( + runtime.randomAccessSequences, + retainedTarget, + ) ?? 1; + const removeEnd = recordedSegmentStartSeconds( + runtime.segmentEndTimesSeconds, + retainedStart, + ) ?? 0; + if (removeEnd > 0) { + await removeRecordedMediaRange(runtime.sourceBuffer, 0, removeEnd, runtime.abort.signal); + } + for (const sequence of runtime.appended) { + if (sequence < retainedStart) runtime.appended.delete(sequence); + } +} + +async function fetchRecordedRuntimeFragment( + runtime: RecordedSegmentStreamRuntime, + targetRevision: number, + sequence: number, +): Promise { + runtime.fetchAbort?.abort(); + const controller = new AbortController(); + runtime.fetchAbort = controller; + const abortFetch = () => controller.abort(); + runtime.abort.signal.addEventListener("abort", abortFetch, { once: true }); + try { + const url = recordedMediaFragmentUrl(runtime.epoch, runtime.contractGeneration, sequence); + const payload = await fetchRecordedFragment(url, controller.signal); + if ( + runtime.disposed + || runtime.abort.signal.aborted + || runtime.target?.revision !== targetRevision + ) return null; + return payload; + } catch (error) { + if ( + controller.signal.aborted + && !runtime.disposed + && !runtime.abort.signal.aborted + && runtime.target?.revision !== targetRevision + ) return null; + throw error; + } finally { + runtime.abort.signal.removeEventListener("abort", abortFetch); + if (runtime.fetchAbort === controller) runtime.fetchAbort = null; + } } async function pumpRecordedSegmentWindow(runtime: RecordedSegmentStreamRuntime): Promise { @@ -344,29 +595,98 @@ async function pumpRecordedSegmentWindow(runtime: RecordedSegmentStreamRuntime): runtime.pumping = true; try { while (!runtime.disposed && !runtime.abort.signal.aborted) { - const targetMissing = !runtime.appended.has(runtime.target) ? runtime.target : null; - let next = targetMissing; - if (next === null) { - for (let sequence = runtime.desiredStart; sequence <= runtime.desiredEnd; sequence += 1) { - if (!runtime.appended.has(sequence)) { - next = sequence; - break; - } - } + const target = runtime.target; + if (!target) break; + if (target.forceReset) { + await resetRecordedSegmentBuffer(runtime); + const latest = runtime.target; + if (!latest) break; + latest.forceReset = false; + continue; } - if (next === null) break; - const url = recordedMediaFragmentUrl(runtime.epoch, runtime.contractGeneration, next); - const payload = await fetchRecordedFragment(url, runtime.abort.signal); - if (runtime.disposed || runtime.abort.signal.aborted) return; - await appendRecordedFragment(runtime.sourceBuffer, payload, runtime.abort.signal); - runtime.appended.add(next); - if (runtime.appended.has(runtime.target)) runtime.onTargetReady?.(runtime.target); + const targetDecoderReady = recordedSegmentTargetBuffered(runtime, target); + if ( + targetDecoderReady + && runtime.notifiedRevision !== target.revision + && runtime.onTargetBuffered !== null + ) { + runtime.notifiedRevision = target.revision; + runtime.onTargetBuffered(target); + } + const next = recordedMediaSegmentAppendOrder( + runtime.appended, + target.decodeStart, + target.desiredEnd, + )[0] ?? null; + if (next !== null) { + const payload = await fetchRecordedRuntimeFragment(runtime, target.revision, next); + if (payload === null) continue; + await appendRecordedFragment(runtime.sourceBuffer, payload, runtime.abort.signal); + runtime.appended.add(next); + continue; + } + if (targetDecoderReady) { + await evictRecordedSegmentHistory(runtime, target); + if (runtime.target === target) break; + continue; + } + if (target.resetAttempts < 1) { + target.resetAttempts += 1; + target.forceReset = true; + continue; + } + throw new Error("Recorded fragment target is not decoder-ready"); } } finally { runtime.pumping = false; } } +function waitForRecordedVideoTarget( + video: HTMLVideoElement, + targetSeconds: number, + signal: AbortSignal, +): Promise { + const targetReady = () => ( + video.readyState >= HTMLMediaElement.HAVE_CURRENT_DATA + && !video.seeking + && Math.abs(video.currentTime - targetSeconds) <= 0.25 + && recordedMediaTimeRangesContain(video.buffered, targetSeconds) + ); + if (signal.aborted) return Promise.reject(new DOMException("Aborted", "AbortError")); + if (targetReady()) return Promise.resolve(); + return new Promise((resolve, reject) => { + const events = ["loadeddata", "canplay", "seeked", "timeupdate"] as const; + const timeout = globalThis.setTimeout(() => { + cleanup(); + reject(new Error("Recorded fragment target did not become decoder-ready")); + }, RECORDED_MEDIA_TARGET_TIMEOUT_MS); + const cleanup = () => { + globalThis.clearTimeout(timeout); + for (const event of events) video.removeEventListener(event, onProgress); + video.removeEventListener("error", onError); + signal.removeEventListener("abort", onAbort); + }; + const onProgress = () => { + if (!targetReady()) return; + cleanup(); + resolve(); + }; + const onError = () => { + cleanup(); + reject(new Error("Recorded fragment target decode failed")); + }; + const onAbort = () => { + cleanup(); + reject(new DOMException("Aborted", "AbortError")); + }; + for (const event of events) video.addEventListener(event, onProgress); + video.addEventListener("error", onError, { once: true }); + signal.addEventListener("abort", onAbort, { once: true }); + onProgress(); + }); +} + export function RecordedFmp4Player({ source, playback, @@ -378,6 +698,7 @@ export function RecordedFmp4Player({ admissionKey = null, onAdmissionChange, onPlaybackChange, + onPlayingRejected, }: { source: ObservationSourceDescriptor; playback?: RecordedObservationPlayback | null; @@ -389,12 +710,15 @@ export function RecordedFmp4Player({ admissionKey?: string | null; onAdmissionChange?: (state: RecordedCameraAdmissionState) => void; onPlaybackChange?: (playback: RecordedObservationPlayback) => void; + onPlayingRejected?: () => void; }) { const videoRef = useRef(null); const onAdmissionChangeRef = useRef(onAdmissionChange); onAdmissionChangeRef.current = onAdmissionChange; const onPlaybackChangeRef = useRef(onPlaybackChange); onPlaybackChangeRef.current = onPlaybackChange; + const onPlayingRejectedRef = useRef(onPlayingRejected); + onPlayingRejectedRef.current = onPlayingRejected; const workerRef = useRef<{ admissionKey: string | null; generation: number } | null>(null); if (!workerRef.current || workerRef.current.admissionKey !== admissionKey) { recordedMediaWorkerGeneration += 1; @@ -431,7 +755,12 @@ export function RecordedFmp4Player({ const [bufferRevision, setBufferRevision] = useState(0); const segmentedRuntimeRef = useRef(null); const [segmentedRuntimeGeneration, setSegmentedRuntimeGeneration] = useState(null); + const targetRevisionRef = useRef(0); + const targetReadyAbortRef = useRef(null); + const playAttemptRevisionRef = useRef(0); const currentSeconds = playback?.currentSeconds ?? contract?.timelineStartSeconds ?? 0; + const playbackPlayingRef = useRef(Boolean(playback?.playing)); + playbackPlayingRef.current = Boolean(playback?.playing); const playbackRate = playback?.rate && Number.isFinite(playback.rate) ? Math.min(4, Math.max(0.25, playback.rate)) : 1; @@ -440,16 +769,17 @@ export function RecordedFmp4Player({ [archive?.manifest.epochs, currentSeconds], ); const segmented = Boolean( - segmentSequence !== null - && Number.isInteger(segmentSequence) - && segmentSequence >= 1 - && segmentCount !== null + segmentCount !== null && Number.isInteger(segmentCount) - && segmentCount >= segmentSequence + && segmentCount >= 1 && typeof MediaSource !== "undefined" && epoch + && epoch.segmentCount === segmentCount + && epoch.randomAccessSequences.length > 0 + && epoch.segmentEndTimesSeconds.length === segmentCount && MediaSource.isTypeSupported(epoch.mediaType), ); + const directPlaybackSeconds = segmented ? null : currentSeconds; const waitingForEpoch = Boolean(archive && !epoch); const selectedGeneration = contract && epoch ? `${contract.manifestGenerationSha256}:${epoch.ordinal}:${epoch.timelineStartSeconds}:${epoch.timelineEndSeconds}` @@ -546,7 +876,15 @@ export function RecordedFmp4Player({ useEffect(() => { const video = videoRef.current; - if (!video || !archive || !contract || !epoch || !segmented || segmentCount === null) return; + const epochSegmentCount = epoch?.segmentCount ?? null; + if ( + !video + || !archive + || !contract + || !epoch + || !segmented + || epochSegmentCount === null + ) return; const abort = new AbortController(); const mediaSource = new MediaSource(); const objectUrl = URL.createObjectURL(mediaSource); @@ -557,11 +895,13 @@ export function RecordedFmp4Player({ setSegmentedRuntimeGeneration(null); setState("loading"); video.pause(); + const sourceOpened = waitForMediaSourceOpen(mediaSource, abort.signal); + const sourceOwner = claimRecordedMediaSource(video); video.src = objectUrl; video.load(); void (async () => { - await waitForMediaSourceOpen(mediaSource, abort.signal); + await sourceOpened; const sourceBuffer = mediaSource.addSourceBuffer(epoch.mediaType); const initUrl = recordedMediaFragmentUrl( epoch, @@ -576,18 +916,21 @@ export function RecordedFmp4Player({ generation, mediaSource, sourceBuffer, + video, objectUrl, epoch, contractGeneration: contract.manifestGenerationSha256, - segmentCount, + segmentCount: epochSegmentCount, + randomAccessSequences: epoch.randomAccessSequences, + segmentEndTimesSeconds: epoch.segmentEndTimesSeconds, appended: new Set(), abort, - desiredStart: 1, - desiredEnd: 1, - target: 1, + target: null, + fetchAbort: null, + notifiedRevision: 0, pumping: false, disposed: false, - onTargetReady: null, + onTargetBuffered: null, }; segmentedRuntimeRef.current = runtime; setSegmentedRuntimeGeneration(generation); @@ -605,15 +948,17 @@ export function RecordedFmp4Player({ }); return () => { - abort.abort(); if (runtime) { runtime.disposed = true; - runtime.onTargetReady = null; + runtime.fetchAbort?.abort(); + runtime.onTargetBuffered = null; } + abort.abort(); + targetReadyAbortRef.current?.abort(); if (segmentedRuntimeRef.current === runtime) segmentedRuntimeRef.current = null; setSegmentedRuntimeGeneration(null); - video.pause(); - if (video.src === objectUrl) { + if (releaseRecordedMediaSource(video, sourceOwner)) { + video.pause(); video.removeAttribute("src"); video.load(); } @@ -631,45 +976,65 @@ export function RecordedFmp4Player({ || !segmented || segmentedRuntimeGeneration !== runtime.generation || segmentSequence === null + || !Number.isInteger(segmentSequence) + || segmentSequence < 1 + || segmentSequence > runtime.segmentCount ) return; const archiveByteLength = archive.byteLength; - runtime.target = segmentSequence; - runtime.desiredStart = Math.max(1, segmentSequence - 2); - runtime.desiredEnd = Math.min(runtime.segmentCount, segmentSequence + 36); - const targetSeconds = recordedMediaLocalTime( - runtime.epoch.timelineStartSeconds, - currentSeconds, - runtime.epoch.timelineEndSeconds - runtime.epoch.timelineStartSeconds, + const decodeStart = recordedMediaDecodeStartSequence( + runtime.randomAccessSequences, + segmentSequence, ); - const markReady = (sequence: number) => { - if ( - runtime.disposed - || segmentedRuntimeRef.current !== runtime - || runtime.target !== sequence - ) return; - try { - if (Math.abs(video.currentTime - targetSeconds) > 0.2) video.currentTime = targetSeconds; - } catch { - setState("error"); - return; - } - setBufferRevision((revision) => revision + 1); - setReadyGeneration(runtime.generation); - setState("ready"); - reportAdmission({ - phase: "ready", - byteLength: archiveByteLength, - message: null, - }); - }; - runtime.onTargetReady = markReady; - if (runtime.appended.has(segmentSequence)) { - markReady(segmentSequence); - } else { + if (decodeStart === null) { setReadyGeneration(null); - setState("loading"); + setState("error"); + reportAdmission({ + phase: "error", + byteLength: archiveByteLength, + message: "Для кадра записанной камеры нет random-access фрагмента.", + }); + return; } - void pumpRecordedSegmentWindow(runtime).catch((error: unknown) => { + const targetSeconds = recordedSegmentStartSeconds( + runtime.segmentEndTimesSeconds, + segmentSequence, + ); + if (targetSeconds === null) { + setReadyGeneration(null); + setState("error"); + reportAdmission({ + phase: "error", + byteLength: archiveByteLength, + message: "Для кадра записанной камеры нет точной media timestamp.", + }); + return; + } + const previousTarget = runtime.target; + const readyEnd = Math.min( + runtime.segmentCount, + segmentSequence + RECORDED_MEDIA_REQUIRED_AHEAD_SEGMENTS, + ); + const desiredEnd = Math.min( + runtime.segmentCount, + segmentSequence + RECORDED_MEDIA_SEGMENTS_AHEAD, + ); + const candidateTarget: RecordedSegmentTarget = { + revision: previousTarget?.revision ?? 0, + sequence: segmentSequence, + decodeStart, + readyEnd, + desiredEnd, + targetSeconds, + forceReset: false, + resetAttempts: 0, + }; + const rollingTarget = Boolean( + playbackPlayingRef.current + && previousTarget + && runtime.notifiedRevision === previousTarget.revision + && recordedSegmentTargetBuffered(runtime, candidateTarget), + ); + const reportPumpError = (error: unknown) => { if ( runtime.disposed || runtime.abort.signal.aborted @@ -682,11 +1047,95 @@ export function RecordedFmp4Player({ byteLength: archiveByteLength, message: "Покадровый фрагмент записанной камеры недоступен.", }); - }); - return () => { - if (runtime.onTargetReady === markReady) runtime.onTargetReady = null; }; - }, [archive?.byteLength, currentSeconds, segmentSequence, segmented, segmentedRuntimeGeneration]); + if (rollingTarget && previousTarget) { + runtime.target = { + ...candidateTarget, + revision: previousTarget.revision, + }; + runtime.onTargetBuffered = null; + void pumpRecordedSegmentWindow(runtime).catch(reportPumpError); + return; + } + + targetRevisionRef.current += 1; + const target: RecordedSegmentTarget = { + ...candidateTarget, + revision: targetRevisionRef.current, + forceReset: Boolean( + runtime.appended.size + && !recordedMediaTimeRangesContain(video.buffered, targetSeconds, 0.02) + ), + }; + const targetReadyAbort = new AbortController(); + targetReadyAbortRef.current?.abort(); + targetReadyAbortRef.current = targetReadyAbort; + runtime.fetchAbort?.abort(); + runtime.target = target; + const alreadyBuffered = recordedSegmentTargetBuffered(runtime, target); + video.pause(); + if (!alreadyBuffered) { + setReadyGeneration(null); + setState("loading"); + } + const markBuffered = (bufferedTarget: RecordedSegmentTarget) => { + if ( + runtime.disposed + || segmentedRuntimeRef.current !== runtime + || runtime.target?.revision !== bufferedTarget.revision + || targetReadyAbort.signal.aborted + ) return; + void (async () => { + try { + if (Math.abs(video.currentTime - bufferedTarget.targetSeconds) > 0.05) { + video.currentTime = bufferedTarget.targetSeconds; + } + await waitForRecordedVideoTarget( + video, + bufferedTarget.targetSeconds, + targetReadyAbort.signal, + ); + if ( + targetReadyAbort.signal.aborted + || runtime.disposed + || segmentedRuntimeRef.current !== runtime + || runtime.target?.revision !== bufferedTarget.revision + ) return; + setBufferRevision((revision) => revision + 1); + setReadyGeneration(runtime.generation); + setState("ready"); + reportAdmission({ + phase: "ready", + byteLength: archiveByteLength, + message: null, + }); + } catch (error) { + if ( + targetReadyAbort.signal.aborted + || (error instanceof DOMException && error.name === "AbortError") + ) return; + setReadyGeneration(null); + setState("error"); + reportAdmission({ + phase: "error", + byteLength: archiveByteLength, + message: "Кадр записанной камеры не стал decoder-ready.", + }); + } + })(); + }; + runtime.onTargetBuffered = markBuffered; + if (alreadyBuffered) { + runtime.notifiedRevision = target.revision; + markBuffered(target); + } + void pumpRecordedSegmentWindow(runtime).catch(reportPumpError); + return () => { + targetReadyAbort.abort(); + if (targetReadyAbortRef.current === targetReadyAbort) targetReadyAbortRef.current = null; + if (runtime.onTargetBuffered === markBuffered) runtime.onTargetBuffered = null; + }; + }, [archive?.byteLength, segmentSequence, segmented, segmentedRuntimeGeneration]); useEffect(() => { const video = videoRef.current; @@ -739,14 +1188,23 @@ export function RecordedFmp4Player({ }, [archive?.byteLength, contract, epoch, segmented]); useEffect(() => { + playAttemptRevisionRef.current += 1; + const playAttemptRevision = playAttemptRevisionRef.current; const video = videoRef.current; if (!video || !epoch || visualState !== "ready") return; - const target = recordedMediaLocalTime( - epoch.timelineStartSeconds, - currentSeconds, - video.duration, - ); - if (Number.isFinite(target) && Math.abs(video.currentTime - target) > 0.35) { + const target = segmented + ? null + : recordedMediaLocalTime( + epoch.timelineStartSeconds, + directPlaybackSeconds ?? epoch.timelineStartSeconds, + video.duration, + ); + if ( + !segmented + && target !== null + && Number.isFinite(target) + && Math.abs(video.currentTime - target) > 0.35 + ) { try { video.currentTime = target; } catch { @@ -762,17 +1220,28 @@ export function RecordedFmp4Player({ } video.playbackRate = playbackRate; if (playback?.playing) { - void video.play().catch(() => undefined); + void video.play().catch(() => { + if (playAttemptRevisionRef.current !== playAttemptRevision) return; + onPlayingRejectedRef.current?.(); + setReadyGeneration(null); + setState("error"); + reportAdmission({ + phase: "error", + byteLength: archive?.byteLength ?? null, + message: "Запуск записанной камеры отклонён браузером.", + }); + }); } else { video.pause(); } }, [ archive?.byteLength, bufferRevision, - currentSeconds, + directPlaybackSeconds, epoch, playback?.playing, playbackRate, + segmented, visualState, ]); @@ -781,15 +1250,6 @@ export function RecordedFmp4Player({ if ((!interactive && onPlaybackChange === undefined) || !video || !epoch || visualState !== "ready") return; let videoFrameRequest: number | null = null; const emitPlayback = () => { - const expectedLocalSeconds = recordedMediaLocalTime( - epoch.timelineStartSeconds, - currentSeconds, - video.duration, - ); - // A controlled seek first moves the shared LAB clock and only then fetches - // its exact fMP4 fragment. Do not let the previously displayed native - // frame race that pending seek and rewind the shared clock. - if (Math.abs(video.currentTime - expectedLocalSeconds) > 0.35) return; onPlaybackChangeRef.current?.({ currentSeconds: epoch.timelineStartSeconds + video.currentTime, playing: !video.paused && !video.ended, @@ -816,7 +1276,7 @@ export function RecordedFmp4Player({ video.cancelVideoFrameCallback(videoFrameRequest); } }; - }, [currentSeconds, epoch, interactive, onPlaybackChange, visualState]); + }, [epoch, interactive, onPlaybackChange, visualState]); return (
void; + onPlayingRejected?: () => void; }) { return (
@@ -50,6 +52,7 @@ export function RecordedEvidenceVideoScene({ segmentSequence={segmentSequence} segmentCount={segmentCount} onPlaybackChange={onPlaybackChange} + onPlayingRejected={onPlayingRejected} /> {semanticOverlay ? ( 0.000001) { + throw new ObservationSessionContractError( + "Временной индекс фрагментов не покрывает codec epoch целиком.", + ); + } + } declaredBytes += byteLength; if (!Number.isSafeInteger(declaredBytes)) { throw new ObservationSessionContractError( @@ -953,6 +1039,9 @@ export function decodeObservationRecordedMediaManifest( mediaType, byteLength, streamUrl, + segmentCount, + randomAccessSequences, + segmentEndTimesSeconds, }; }); if (declaredBytes !== manifestByteLength) { diff --git a/apps/control-station/src/workspaces/laboratory/M4ReplayThreatVisual.tsx b/apps/control-station/src/workspaces/laboratory/M4ReplayThreatVisual.tsx index 964740c..b15d3ef 100644 --- a/apps/control-station/src/workspaces/laboratory/M4ReplayThreatVisual.tsx +++ b/apps/control-station/src/workspaces/laboratory/M4ReplayThreatVisual.tsx @@ -37,7 +37,7 @@ import type { } from "../../core/laboratory/m4ReplayThreat"; import { buildM4LocalSurface } from "../../core/laboratory/m4LocalSurface"; import { recordedObservationSources } from "../../core/observation/recordedObservationSources"; -import { replayObservationSession } from "../../core/observation/sessionArchive"; +import { resolveObservationSessionReplay } from "../../core/observation/useObservationSessions"; import type { ObservationSourceDescriptor } from "../../core/runtime/contracts"; import { useM4ThreatTimelineFrame, @@ -148,14 +148,11 @@ export function M4ReplayThreatVisual({ const controller = new AbortController(); setVideoLoading(true); setVideoError(null); - void replayObservationSession(timeline.recordedSourceSessionId, { + void resolveObservationSessionReplay(timeline.recordedSourceSessionId, { signal: controller.signal, }) - .then((replay) => { - if (replay.kind !== "ready") { - throw new Error("RIGHT-видео RAVNOVES00 ещё готовится к воспроизведению."); - } - const source = recordedObservationSources(replay.launch).find( + .then((launch) => { + const source = recordedObservationSources(launch).find( (candidate) => candidate.modality === "video" && candidate.semanticChannelId === "camera.video.recorded", ); @@ -532,9 +529,14 @@ export function M4ReplayThreatVisual({ semanticOverlay={mediaMode === "video" ? semanticOverlay : undefined} ariaLabel={`M4.6 recorded-realtime frame ${frame?.sequence ?? 0}: ${activeBoxes.length} proposals`} interactive={false} - segmentSequence={frame ? frame.sequence + 1 : undefined} + segmentSequence={ + timelineFrame.activeSequence === null + ? undefined + : timelineFrame.activeSequence + 1 + } segmentCount={timeline.frameCount} onPlaybackChange={playbackController.synchronize} + onPlayingRejected={() => playbackController.setPlaying(false)} /> ) : videoError ? ( diff --git a/apps/control-station/src/workspaces/laboratory/useM4ThreatTimeline.ts b/apps/control-station/src/workspaces/laboratory/useM4ThreatTimeline.ts index 54e14ff..5a1ea07 100644 --- a/apps/control-station/src/workspaces/laboratory/useM4ThreatTimeline.ts +++ b/apps/control-station/src/workspaces/laboratory/useM4ThreatTimeline.ts @@ -29,6 +29,18 @@ export function m4ThreatChunkWindowStarts( ).filter((start) => start >= 0 && start < frameCount); } +export function cancelM4ThreatChunkRequestsOutsideWindow( + inFlight: Map, + desiredStarts: readonly number[], +): void { + const desired = new Set(desiredStarts); + for (const [start, controller] of inFlight) { + if (desired.has(start)) continue; + controller.abort(); + inFlight.delete(start); + } +} + export function useM4ThreatTimelineMetadata(resultId: string) { const [timeline, setTimeline] = useState(null); const [error, setError] = useState(null); @@ -105,6 +117,7 @@ export function useM4ThreatTimelineFrame({ chunkSize, timeline.frameCount, ); + cancelM4ThreatChunkRequestsOutsideWindow(inFlight.current, starts); for (const start of starts) { if (chunksRef.current.has(start) || inFlight.current.has(start)) continue; const controller = new AbortController(); diff --git a/apps/control-station/test/m4ReplayThreat.test.mjs b/apps/control-station/test/m4ReplayThreat.test.mjs index 1cfb57c..2b4c244 100644 --- a/apps/control-station/test/m4ReplayThreat.test.mjs +++ b/apps/control-station/test/m4ReplayThreat.test.mjs @@ -14,6 +14,7 @@ let selectM4ThreatTimelineSequence; let advanceRecordedEvidencePlayback; let synchronizeRecordedEvidencePlayback; let m4ThreatChunkWindowStarts; +let cancelM4ThreatChunkRequestsOutsideWindow; let buildM4LocalSurface; const resultId = `m4-threat-replay-${"a".repeat(64)}`; @@ -38,7 +39,10 @@ before(async () => { } = await server.ssrLoadModule( "/src/components/laboratory/useRecordedEvidencePlayback.ts", )); - ({ m4ThreatChunkWindowStarts } = await server.ssrLoadModule( + ({ + m4ThreatChunkWindowStarts, + cancelM4ThreatChunkRequestsOutsideWindow, + } = await server.ssrLoadModule( "/src/workspaces/laboratory/useM4ThreatTimeline.ts", )); ({ buildM4LocalSurface } = await server.ssrLoadModule( @@ -375,6 +379,23 @@ test("M4.6 spatial buffering keeps previous, active and two future chunks", () = assert.deepEqual(m4ThreatChunkWindowStarts(0, 24, 4489), [0, 24, 48]); }); +test("M4.6 spatial buffering drops stale in-flight windows across rapid jumps", () => { + const aborted = []; + const controller = (start) => ({ abort: () => aborted.push(start) }); + const inFlight = new Map( + m4ThreatChunkWindowStarts(0, 24, 4489).map((start) => [start, controller(start)]), + ); + for (const activeStart of [1488, 4488]) { + const desired = m4ThreatChunkWindowStarts(activeStart, 24, 4489); + cancelM4ThreatChunkRequestsOutsideWindow(inFlight, desired); + for (const start of desired) { + if (!inFlight.has(start)) inFlight.set(start, controller(start)); + } + } + assert.deepEqual([...inFlight.keys()], [4464, 4488]); + assert.deepEqual(aborted, [0, 24, 48, 1464, 1488, 1512, 1536]); +}); + test("recorded evidence clock advances by selected rate and stops at the sealed end", () => { const range = { startSeconds: 10, endSeconds: 20 }; assert.deepEqual( @@ -448,6 +469,7 @@ test("M4.6 viewer keeps media and spatial panes on one playback clock", async () assert.match(visual, /const spatialFrame = frame\?\.spatialAvailable/); assert.match(visual, / playbackController\.setPlaying\(false\)\}/); assert.match(visual, /currentSeconds: playbackController\.playback\.currentSeconds/); - assert.doesNotMatch(visual, /setPlaying\(false\)/); assert.match(visualCss, /m4-replay-threat-visual__deck > \.nodedc-split-pane/); assert.match(visualCss, /m4-replay-threat-visual__pane-toolbar\[data-pane-toolbar="media"\]/); assert.match(visualCss, /m4-replay-threat-visual__pane-toolbar\[data-pane-toolbar="spatial"\]/); diff --git a/apps/control-station/test/observationSessions.test.mjs b/apps/control-station/test/observationSessions.test.mjs index 7dc16c1..ee3124c 100644 --- a/apps/control-station/test/observationSessions.test.mjs +++ b/apps/control-station/test/observationSessions.test.mjs @@ -953,7 +953,7 @@ test("replay decodes opaque recorded cameras and their same-origin fMP4 manifest ); const manifest = decodeObservationRecordedMediaManifest({ - schema_version: "missioncore.observation-recorded-media/v3", + schema_version: "missioncore.observation-recorded-media/v4", source_id: source.id, generation_sha256: "c".repeat(64), byte_length: 3_266, @@ -970,11 +970,17 @@ test("replay decodes opaque recorded cameras and their same-origin fMP4 manifest "/manifest", `/epochs/1/recording.mp4?generation=${"c".repeat(64)}`, ), + segment_count: 2, + random_access_sequences: [1], + segment_end_times_seconds: [9.75, 19.75], }], }, decoded.mediaSources[0]); assert.equal(manifest.epochs[0].byteLength, 3_266); assert.match(manifest.epochs[0].streamUrl, /recording\.mp4\?generation=/); assert.equal(manifest.epochs[0].timelineStartSeconds, 0.25); + assert.equal(manifest.epochs[0].segmentCount, 2); + assert.deepEqual(manifest.epochs[0].randomAccessSequences, [1]); + assert.deepEqual(manifest.epochs[0].segmentEndTimesSeconds, [9.75, 19.75]); assert.throws( () => decodeObservationRecordedMediaManifest({ ...{ @@ -992,7 +998,7 @@ test("replay decodes opaque recorded cameras and their same-origin fMP4 manifest ); assert.throws( () => decodeObservationRecordedMediaManifest({ - schema_version: "missioncore.observation-recorded-media/v3", + schema_version: "missioncore.observation-recorded-media/v4", source_id: source.id, generation_sha256: "c".repeat(64), byte_length: 3_267, @@ -1106,7 +1112,7 @@ test("recorded camera contracts reject foreign origins, path escapes and unknown })); assert.throws( () => decodeObservationRecordedMediaManifest({ - schema_version: "missioncore.observation-recorded-media/v3", + schema_version: "missioncore.observation-recorded-media/v4", source_id: base.id, generation_sha256: "c".repeat(64), byte_length: 384, @@ -1120,6 +1126,9 @@ test("recorded camera contracts reject foreign origins, path escapes and unknown media_type: 'video/mp4; codecs="avc1.640028"', byte_length: 384, stream_url: "file:///private/recording.mp4", + segment_count: 1, + random_access_sequences: [1], + segment_end_times_seconds: [20], }], }, decoded.mediaSources[0]), /небезопасный API URL/, diff --git a/apps/control-station/test/recordedCameraBuffering.test.mjs b/apps/control-station/test/recordedCameraBuffering.test.mjs index 04b7f30..6aac1d2 100644 --- a/apps/control-station/test/recordedCameraBuffering.test.mjs +++ b/apps/control-station/test/recordedCameraBuffering.test.mjs @@ -9,6 +9,8 @@ let fetchRecordedMediaArchive; let recordedMediaPresentationState; let recordedMediaSeekableCoverage; let recordedMediaFragmentUrl; +let recordedMediaDecodeStartSequence; +let recordedMediaSegmentAppendOrder; before(async () => { server = await createServer({ @@ -21,6 +23,8 @@ before(async () => { recordedMediaPresentationState, recordedMediaSeekableCoverage, recordedMediaFragmentUrl, + recordedMediaDecodeStartSequence, + recordedMediaSegmentAppendOrder, } = await server.ssrLoadModule("/src/components/RecordedFmp4Player.tsx")); }); @@ -29,6 +33,7 @@ after(async () => { }); function fixture({ byteLength = 36_000_000_000 } = {}) { + const segmentCount = 1_501; const manifestUrl = "/api/v1/observation-sessions/session-1/media/camera-1/manifest"; const generation = "a".repeat(64); const streamUrl = manifestUrl.replace( @@ -36,7 +41,7 @@ function fixture({ byteLength = 36_000_000_000 } = {}) { `/epochs/1/recording.mp4?generation=${generation}`, ); const manifest = { - schema_version: "missioncore.observation-recorded-media/v3", + schema_version: "missioncore.observation-recorded-media/v4", source_id: "recorded.camera.camera-1", generation_sha256: generation, byte_length: byteLength, @@ -50,6 +55,12 @@ function fixture({ byteLength = 36_000_000_000 } = {}) { media_type: 'video/mp4; codecs="avc1.640028"', byte_length: byteLength, stream_url: streamUrl, + segment_count: segmentCount, + random_access_sequences: [1, 1491, 1501], + segment_end_times_seconds: Array.from( + { length: segmentCount }, + (_value, index) => ((index + 1) * 36_000) / segmentCount, + ), }], }; const source = { @@ -92,6 +103,8 @@ test("multi-hour camera admission fetches only its compact generation-bound mani assert.deepEqual(requested, [source.manifestUrl]); assert.equal(archive.byteLength, 36_000_000_000); assert.equal(archive.manifest.epochs[0].streamUrl, streamUrl); + assert.equal(archive.manifest.epochs[0].segmentCount, 1_501); + assert.deepEqual(archive.manifest.epochs[0].randomAccessSequences, [1, 1491, 1501]); }); test("first camera manifest request rejects a replaced generation", async () => { @@ -154,8 +167,16 @@ test("recorded player keeps full-archive range fallback and uses bounded generat assert.match(source, /video\.src\s*=\s*descriptor\.streamUrl/); assert.doesNotMatch(source, /new Blob\(/); assert.match(source, /new MediaSource\(\)/); - assert.match(source, /segmentSequence \+ 36/); + assert.match(source, /RECORDED_MEDIA_SEGMENTS_AHEAD = 36/); assert.match(source, /pumpRecordedSegmentWindow/); + assert.match(source, /waitForRecordedVideoTarget/); + assert.match(source, /removeRecordedMediaRange/); + + assert.equal(recordedMediaDecodeStartSequence([1, 1491, 1501], 1500), 1491); + assert.deepEqual( + recordedMediaSegmentAppendOrder(new Set([1491, 1492]), 1491, 1500), + [1493, 1494, 1495, 1496, 1497, 1498, 1499, 1500], + ); const { manifest, generation } = fixture(); const epoch = { @@ -165,6 +186,9 @@ test("recorded player keeps full-archive range fallback and uses bounded generat mediaType: manifest.epochs[0].media_type, byteLength: manifest.epochs[0].byte_length, streamUrl: manifest.epochs[0].stream_url, + segmentCount: manifest.epochs[0].segment_count, + randomAccessSequences: manifest.epochs[0].random_access_sequences, + segmentEndTimesSeconds: manifest.epochs[0].segment_end_times_seconds, }; assert.equal( recordedMediaFragmentUrl(epoch, generation, "init"), diff --git a/src/k1link/sessions/camera_frame.py b/src/k1link/sessions/camera_frame.py index 6d069c6..efdfcd0 100644 --- a/src/k1link/sessions/camera_frame.py +++ b/src/k1link/sessions/camera_frame.py @@ -100,7 +100,7 @@ class RecordedCameraFrameService: target = self._inspector.get_segment(manifest, epoch.ordinal, sequence) fragments: list[bytes] = [target.payload] key_sequence = sequence - while not _fragment_is_sync(fragments[0]): + while not epoch.segments[key_sequence - 1].random_access: key_sequence -= 1 if key_sequence < 1 or sequence - key_sequence > _MAX_KEYFRAME_DISTANCE: raise SessionIntegrityError("camera frame has no bounded sync fragment") @@ -171,92 +171,6 @@ def _frame_location( raise SessionIntegrityError("camera frame is outside the recorded media manifest") -def _fragment_is_sync(payload: bytes) -> bool: - tfhd_default_flags: int | None = None - sample_flags: int | None = None - sample_count: int | None = None - for box_type, body in _walk_boxes(payload): - if box_type == b"tfhd": - if len(body) < 8: - raise SessionIntegrityError("camera fragment tfhd is truncated") - flags = int.from_bytes(body[1:4], "big") - offset = 8 - for mask, size in ((0x000001, 8), (0x000002, 4), (0x000008, 4), (0x000010, 4)): - if flags & mask: - offset += size - if flags & 0x000020: - if offset + 4 > len(body): - raise SessionIntegrityError("camera fragment default flags are truncated") - tfhd_default_flags = struct.unpack_from(">I", body, offset)[0] - elif box_type == b"trun": - if len(body) < 8: - raise SessionIntegrityError("camera fragment trun is truncated") - flags = int.from_bytes(body[1:4], "big") - sample_count = struct.unpack_from(">I", body, 4)[0] - if sample_count != 1: - raise SessionIntegrityError("camera fragment must contain exactly one sample") - offset = 8 - if flags & 0x000001: - offset += 4 - if flags & 0x000004: - if offset + 4 > len(body): - raise SessionIntegrityError("camera fragment first flags are truncated") - sample_flags = struct.unpack_from(">I", body, offset)[0] - offset += 4 - per_sample_sizes = ( - (0x000100, 4), - (0x000200, 4), - (0x000400, 4), - (0x000800, 4), - ) - for mask, size in per_sample_sizes: - if flags & mask: - if offset + size > len(body): - raise SessionIntegrityError("camera fragment sample data is truncated") - if mask == 0x000400: - sample_flags = struct.unpack_from(">I", body, offset)[0] - offset += size - if sample_count != 1: - raise SessionIntegrityError("camera fragment has no unique video sample") - effective_flags = sample_flags if sample_flags is not None else tfhd_default_flags - if effective_flags is None: - raise SessionIntegrityError("camera fragment sample flags are unavailable") - return (effective_flags & 0x00010000) == 0 - - -def _walk_boxes(payload: bytes): - containers = {b"moof", b"traf"} - pending = [(0, len(payload))] - boxes = 0 - while pending: - start, end = pending.pop() - offset = start - while offset + 8 <= end: - boxes += 1 - if boxes > 64: - raise SessionIntegrityError("camera fragment box budget was exceeded") - size = struct.unpack_from(">I", payload, offset)[0] - box_type = payload[offset + 4 : offset + 8] - header = 8 - if size == 1: - if offset + 16 > end: - raise SessionIntegrityError("camera fragment extended box is truncated") - size = struct.unpack_from(">Q", payload, offset + 8)[0] - header = 16 - elif size == 0: - size = end - offset - if size < header or offset + size > end: - raise SessionIntegrityError("camera fragment box size is invalid") - body_start = offset + header - body_end = offset + size - yield box_type, payload[body_start:body_end] - if box_type in containers: - pending.append((body_start, body_end)) - offset = body_end - if offset != end: - raise SessionIntegrityError("camera fragment box boundary is invalid") - - def _jpeg_dimensions(payload: bytes) -> tuple[int, int]: if len(payload) < 4 or payload[:2] != b"\xff\xd8" or payload[-2:] != b"\xff\xd9": raise SessionIntegrityError("camera frame decoder returned an invalid JPEG") diff --git a/src/k1link/sessions/media.py b/src/k1link/sessions/media.py index 39e5989..26424fb 100644 --- a/src/k1link/sessions/media.py +++ b/src/k1link/sessions/media.py @@ -18,8 +18,8 @@ from .models import RecordedMediaArtifact, ReplayCommand, SessionIntegrityError CAMERA_ARCHIVE_SCHEMA = "missioncore.camera-recording/v1" CAMERA_INDEX_SCHEMA = "missioncore.camera-recording-index/v1" -RECORDED_MEDIA_MANIFEST_SCHEMA = "missioncore.observation-recorded-media/v2" -RECORDED_MEDIA_PREPARATION_SCHEMA = "missioncore.recorded-media-preparation/v1" +RECORDED_MEDIA_MANIFEST_SCHEMA = "missioncore.observation-recorded-media/v4" +RECORDED_MEDIA_PREPARATION_SCHEMA = "missioncore.recorded-media-preparation/v3" MAX_MEDIA_SUMMARY_BYTES = 2 * 1024 * 1024 MAX_MEDIA_INDEX_LINE_BYTES = 64 * 1024 MAX_INIT_BYTES = 8 * 1024 * 1024 @@ -43,6 +43,8 @@ class RecordedMediaSegment: path: Path byte_length: int sha256: str + random_access: bool + end_time_seconds: float @dataclass(frozen=True, slots=True) @@ -91,12 +93,14 @@ class _Mp4VideoTiming: track_id: int timescale: int default_sample_duration: int | None + default_sample_flags: int | None @dataclass(frozen=True, slots=True) class _Mp4VideoFragmentTiming: base_decode_time: int duration_units: int + random_access: bool @property def end_decode_time(self) -> int: @@ -194,9 +198,7 @@ class RecordedMediaInspector: "recorded media preparation could not be inspected" ) from exc if stat.S_ISLNK(metadata.st_mode) or not stat.S_ISREG(metadata.st_mode): - raise SessionIntegrityError( - "recorded media preparation is not a regular file" - ) + raise SessionIntegrityError("recorded media preparation is not a regular file") try: sidecar.unlink() except OSError as exc: @@ -472,6 +474,8 @@ def _sidecar_manifest_document(manifest: RecordedMediaManifest) -> dict[str, Any "sequence": segment.sequence, "byte_length": segment.byte_length, "sha256": segment.sha256, + "random_access": segment.random_access, + "end_time_seconds": segment.end_time_seconds, } for segment in epoch.segments ], @@ -581,33 +585,46 @@ def _epoch_from_sidecar( ): raise SessionIntegrityError("recorded media prepared epoch descriptor is invalid") segment_documents = value.get("segments") - if ( - not isinstance(segment_documents, list) - or not segment_documents - ): + if not isinstance(segment_documents, list) or not segment_documents: raise SessionIntegrityError("recorded media prepared segments are invalid") segments_root = epoch_path / "segments" segments: list[RecordedMediaSegment] = [] + previous_end_time_seconds = 0.0 for sequence, segment in enumerate(segment_documents, start=1): if not isinstance(segment, dict) or segment.get("sequence") != sequence: raise SessionIntegrityError("recorded media prepared segment is inconsistent") byte_length = segment.get("byte_length") sha256 = segment.get("sha256") + random_access = segment.get("random_access") + end_time_seconds = segment.get("end_time_seconds") if ( not _positive_int(byte_length) or int(byte_length) > MAX_MEDIA_SEGMENT_BYTES or not isinstance(sha256, str) or _SHA256_PATTERN.fullmatch(sha256) is None + or not isinstance(random_access, bool) + or not _finite_non_negative_number(end_time_seconds) + or float(end_time_seconds) <= previous_end_time_seconds ): raise SessionIntegrityError("recorded media prepared segment is invalid") + previous_end_time_seconds = float(end_time_seconds) segments.append( RecordedMediaSegment( sequence=sequence, path=segments_root / f"{sequence}.m4s", byte_length=int(byte_length), sha256=sha256, + random_access=random_access, + end_time_seconds=previous_end_time_seconds, ) ) + if not math.isclose( + previous_end_time_seconds, + float(end) - float(start), + rel_tol=0.0, + abs_tol=1e-9, + ): + raise SessionIntegrityError("recorded media prepared segment timeline is inconsistent") return RecordedMediaEpoch( ordinal=ordinal, path=epoch_path, @@ -863,6 +880,8 @@ def _manifest_generation_sha256( "sequence": segment.sequence, "byte_length": segment.byte_length, "sha256": segment.sha256, + "random_access": segment.random_access, + "end_time_seconds": segment.end_time_seconds, } for segment in epoch.segments ], @@ -895,10 +914,7 @@ def _read_epoch( ): raise SessionIntegrityError("recorded media summary is incompatible") segment_count = summary.get("segment_count") - if ( - not _non_negative_int(segment_count) - or not 1 <= int(segment_count) <= MAX_SAFE_INTEGER - ): + if not _non_negative_int(segment_count) or not 1 <= int(segment_count) <= MAX_SAFE_INTEGER: raise SessionIntegrityError("recorded media segment count is invalid") match = _EPOCH_PATTERN.fullmatch(epoch.name) if ( @@ -922,10 +938,7 @@ def _read_epoch( expected_count=int(segment_count), ) expected_index_sha256 = summary.get("index_sha256") - if ( - not isinstance(expected_index_sha256, str) - or index_sha256 != expected_index_sha256 - ): + if not isinstance(expected_index_sha256, str) or index_sha256 != expected_index_sha256: raise SessionIntegrityError("recorded media index digest changed") use_monotonic_clock = all( @@ -957,7 +970,7 @@ def _read_epoch( segments_root = (epoch / "segments").resolve(strict=True) if not segments_root.is_dir() or segments_root.parent != epoch: raise SessionIntegrityError("recorded media segment directory is invalid") - segments: list[RecordedMediaSegment] = [] + segment_descriptors: list[tuple[int, Path, int, str]] = [] for sequence, entry in enumerate(entries, start=1): byte_length = entry.get("length") digest = entry.get("sha256") @@ -977,14 +990,7 @@ def _read_epoch( raise SessionIntegrityError("recorded media segment length is outside bounds") if metadata.st_size != int(byte_length): raise SessionIntegrityError("recorded media segment length changed") - segments.append( - RecordedMediaSegment( - sequence=sequence, - path=path, - byte_length=int(byte_length), - sha256=digest, - ) - ) + segment_descriptors.append((sequence, path, int(byte_length), digest)) init_path = (epoch / "init.mp4").resolve() init_payload = _read_confined_file(init_path, epoch, MAX_INIT_BYTES) @@ -997,28 +1003,31 @@ def _read_epoch( fragment_timings: list[_Mp4VideoFragmentTiming] = [] stream_sha256 = hashlib.sha256(init_payload) valid_bytes = len(init_payload) - for segment in segments: + for _sequence, path, byte_length, digest in segment_descriptors: payload = _read_confined_file( - segment.path, + path, segments_root, - segment.byte_length, + byte_length, ) - if hashlib.sha256(payload).hexdigest() != segment.sha256: + if hashlib.sha256(payload).hexdigest() != digest: raise SessionIntegrityError("recorded media segment digest changed") stream_sha256.update(payload) valid_bytes += len(payload) - fragment_timings.append( - _mp4_video_fragment_timing( - payload, - timing, - _Mp4ParseBudget(), - ) + fragment_timing = _mp4_video_fragment_timing( + payload, + timing, + _Mp4ParseBudget(), ) + fragment_timings.append(fragment_timing) if ( summary.get("valid_bytes") != valid_bytes or summary.get("stream_sha256") != stream_sha256.hexdigest() ): raise SessionIntegrityError("recorded media stream aggregate changed") + if not fragment_timings[0].random_access: + raise SessionIntegrityError( + "recorded media codec epoch does not begin with a random-access fragment" + ) for previous, current in zip( fragment_timings, fragment_timings[1:], @@ -1039,6 +1048,32 @@ def _read_epoch( timeline_end_seconds = timeline_start_seconds + total_duration_seconds if not math.isfinite(timeline_end_seconds) or timeline_end_seconds < timeline_start_seconds: raise SessionIntegrityError("recorded media epoch timeline is invalid") + epoch_duration_seconds = timeline_end_seconds - timeline_start_seconds + segments: list[RecordedMediaSegment] = [] + cumulative_duration_units = 0 + previous_end_time_seconds = 0.0 + for index, (descriptor, fragment_timing) in enumerate( + zip(segment_descriptors, fragment_timings, strict=True), + start=1, + ): + sequence, path, byte_length, digest = descriptor + cumulative_duration_units += fragment_timing.duration_units + end_time_seconds = cumulative_duration_units / timing.timescale + if index == len(fragment_timings): + end_time_seconds = epoch_duration_seconds + if not math.isfinite(end_time_seconds) or end_time_seconds <= previous_end_time_seconds: + raise SessionIntegrityError("recorded media segment timeline is invalid") + previous_end_time_seconds = end_time_seconds + segments.append( + RecordedMediaSegment( + sequence=sequence, + path=path, + byte_length=byte_length, + sha256=digest, + random_access=fragment_timing.random_access, + end_time_seconds=end_time_seconds, + ) + ) return RecordedMediaEpoch( ordinal=ordinal, path=epoch, @@ -1129,17 +1164,17 @@ def _mp4_video_timing(payload: bytes, budget: _Mp4ParseBudget) -> _Mp4VideoTimin raise SessionIntegrityError("recorded media init has no unique moov box") moov = moov_payloads[0] moov_boxes = tuple(_iter_mp4_boxes(moov, budget)) - defaults: dict[int, int] = {} + defaults: dict[int, tuple[int, int]] = {} for box_type, box_payload in moov_boxes: if box_type != b"mvex": continue for child_type, child_payload in _iter_mp4_boxes(box_payload, budget): if child_type != b"trex": continue - track_id, trex_default_duration = _parse_trex(child_payload) + track_id, trex_default_duration, trex_default_flags = _parse_trex(child_payload) if track_id in defaults: raise SessionIntegrityError("recorded media init repeats a trex track") - defaults[track_id] = trex_default_duration + defaults[track_id] = (trex_default_duration, trex_default_flags) video_tracks: list[tuple[int, int]] = [] for box_type, trak_payload in moov_boxes: @@ -1152,13 +1187,15 @@ def _mp4_video_timing(payload: bytes, budget: _Mp4ParseBudget) -> _Mp4VideoTimin if len(video_tracks) != 1: raise SessionIntegrityError("recorded media init has no unique video track") track_id, timescale = video_tracks[0] - default_duration = defaults.get(track_id) + default_sample = defaults.get(track_id) + default_duration = None if default_sample is None else default_sample[0] return _Mp4VideoTiming( track_id=track_id, timescale=timescale, default_sample_duration=( default_duration if default_duration is not None and default_duration > 0 else None ), + default_sample_flags=None if default_sample is None else default_sample[1], ) @@ -1219,7 +1256,7 @@ def _mp4_video_fragment_timing( tfhd_payloads = [box for kind, box in boxes if kind == b"tfhd"] if len(tfhd_payloads) != 1: raise SessionIntegrityError("recorded media fragment tfhd is ambiguous") - track_id, fragment_default_duration = _parse_tfhd(tfhd_payloads[0]) + track_id, fragment_default_duration, fragment_default_flags = _parse_tfhd(tfhd_payloads[0]) if track_id != timing.track_id: continue tfdt_payloads = [box for kind, box in boxes if kind == b"tfdt"] @@ -1230,15 +1267,28 @@ def _mp4_video_fragment_timing( if not trun_payloads: raise SessionIntegrityError("recorded media video fragment has no trun box") default_duration = fragment_default_duration or timing.default_sample_duration - duration = sum( - _parse_trun_duration_units(trun, default_duration, budget) for trun in trun_payloads + default_flags = ( + fragment_default_flags + if fragment_default_flags is not None + else timing.default_sample_flags ) + trun_descriptors = tuple( + _parse_trun_descriptor(trun, default_duration, default_flags, budget) + for trun in trun_payloads + ) + fragment_sample_count = sum(item[2] for item in trun_descriptors) + if fragment_sample_count != 1: + raise SessionIntegrityError( + "recorded media video fragment must contain exactly one sample" + ) + duration = sum(item[0] for item in trun_descriptors) if base_decode_time > MAX_SAFE_INTEGER - duration: raise SessionIntegrityError("recorded media fragment decode time is outside bounds") matching_timings.append( _Mp4VideoFragmentTiming( base_decode_time=base_decode_time, duration_units=duration, + random_access=(trun_descriptors[0][1] & 0x00010000) == 0, ) ) if len(matching_timings) != 1: @@ -1292,15 +1342,16 @@ def _parse_hdlr_type(payload: bytes) -> bytes: return payload[8:12] -def _parse_trex(payload: bytes) -> tuple[int, int]: +def _parse_trex(payload: bytes) -> tuple[int, int, int]: _full_box_version(payload) return ( _read_u32(payload, 4, "trex track id"), _read_u32(payload, 12, "trex default sample duration"), + _read_u32(payload, 20, "trex default sample flags"), ) -def _parse_tfhd(payload: bytes) -> tuple[int, int | None]: +def _parse_tfhd(payload: bytes) -> tuple[int, int | None, int | None]: flags = _full_box_flags(payload) track_id = _read_u32(payload, 4, "tfhd track id") cursor = 8 @@ -1311,10 +1362,16 @@ def _parse_tfhd(payload: bytes) -> tuple[int, int | None]: if flags & 0x000008: default_duration = _read_u32(payload, cursor, "tfhd default sample duration") cursor += 4 - for flag in (0x000010, 0x000020): - if flags & flag: - cursor = _advance_box_cursor(payload, cursor, 4, "tfhd optional field") - return track_id, default_duration if default_duration and default_duration > 0 else None + if flags & 0x000010: + cursor = _advance_box_cursor(payload, cursor, 4, "tfhd default sample size") + default_flags: int | None = None + if flags & 0x000020: + default_flags = _read_u32(payload, cursor, "tfhd default sample flags") + return ( + track_id, + default_duration if default_duration and default_duration > 0 else None, + default_flags, + ) def _parse_tfdt(payload: bytes) -> int: @@ -1328,11 +1385,12 @@ def _parse_tfdt(payload: bytes) -> int: raise SessionIntegrityError("recorded media tfdt version is unsupported") -def _parse_trun_duration_units( +def _parse_trun_descriptor( payload: bytes, default_duration: int | None, + default_sample_flags: int | None, budget: _Mp4ParseBudget, -) -> int: +) -> tuple[int, int, int]: flags = _full_box_flags(payload) sample_count = _read_u32(payload, 4, "trun sample count") if sample_count < 1: @@ -1341,30 +1399,51 @@ def _parse_trun_duration_units( cursor = 8 if flags & 0x000001: cursor = _advance_box_cursor(payload, cursor, 4, "trun data offset") + first_sample_flags: int | None = None if flags & 0x000004: - cursor = _advance_box_cursor(payload, cursor, 4, "trun first sample flags") - per_sample_width = sum(4 for flag in (0x000100, 0x000200, 0x000400, 0x000800) if flags & flag) - if per_sample_width and sample_count > (len(payload) - cursor) // per_sample_width: - raise SessionIntegrityError("recorded media trun samples are truncated") - if flags & 0x000100: - duration = 0 - for _ in range(sample_count): + first_sample_flags = _read_u32(payload, cursor, "trun first sample flags") + cursor += 4 + if first_sample_flags is not None and flags & 0x000400: + raise SessionIntegrityError("recorded media trun sample flags are ambiguous") + + duration = 0 + first_per_sample_flags: int | None = None + for sample_index in range(sample_count): + if flags & 0x000100: sample_duration = _read_u32(payload, cursor, "trun sample duration") if sample_duration <= 0: raise SessionIntegrityError("recorded media sample duration is invalid") duration += sample_duration - cursor += per_sample_width - return duration - if default_duration is None or default_duration <= 0: - raise SessionIntegrityError("recorded media sample duration is unavailable") - if per_sample_width: - _advance_box_cursor( - payload, - cursor, - sample_count * per_sample_width, - "trun sample table", - ) - return sample_count * default_duration + cursor += 4 + elif default_duration is None or default_duration <= 0: + raise SessionIntegrityError("recorded media sample duration is unavailable") + else: + duration += default_duration + if flags & 0x000200: + cursor = _advance_box_cursor(payload, cursor, 4, "trun sample size") + if flags & 0x000400: + sample_flags = _read_u32(payload, cursor, "trun sample flags") + if sample_index == 0: + first_per_sample_flags = sample_flags + cursor += 4 + if flags & 0x000800: + cursor = _advance_box_cursor( + payload, + cursor, + 4, + "trun sample composition time offset", + ) + + effective_flags = ( + first_per_sample_flags + if first_per_sample_flags is not None + else first_sample_flags + if first_sample_flags is not None + else default_sample_flags + ) + if effective_flags is None: + raise SessionIntegrityError("recorded media sample flags are unavailable") + return duration, effective_flags, sample_count def _full_box_version(payload: bytes) -> int: @@ -1614,9 +1693,11 @@ def _read_media_index( "recorded media index length does not match its summary" ) after = os.fstat(stream.fileno()) - if ( - (after.st_dev, after.st_ino, after.st_size, after.st_mtime_ns) - != (before.st_dev, before.st_ino, before.st_size, before.st_mtime_ns) + if (after.st_dev, after.st_ino, after.st_size, after.st_mtime_ns) != ( + before.st_dev, + before.st_ino, + before.st_size, + before.st_mtime_ns, ): raise SessionIntegrityError("recorded media index changed during validation") return entries, digest.hexdigest() diff --git a/src/k1link/web/session_api.py b/src/k1link/web/session_api.py index 39f64d0..76b2f87 100644 --- a/src/k1link/web/session_api.py +++ b/src/k1link/web/session_api.py @@ -48,7 +48,8 @@ from k1link.viewer.recorded import ( from k1link.viewer.rerun_bridge import RerunSceneSettings DEFAULT_REPLAY_ACTION_ID = "stream.start-replay" -RECORDED_MEDIA_STREAM_MANIFEST_SCHEMA = "missioncore.observation-recorded-media/v3" +RECORDED_MEDIA_STREAM_MANIFEST_SCHEMA = "missioncore.observation-recorded-media/v4" +RECORDED_PERCEPTION_STREAM_MANIFEST_SCHEMA = "missioncore.observation-recorded-media/v3" SAFE_SOURCE_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:-]{0,255}$") SAFE_SHA256 = re.compile(r"^[a-f0-9]{64}$") MAX_SAFE_INTEGER = 9_007_199_254_740_991 @@ -1706,7 +1707,7 @@ def _recorded_perception_manifest_document( encoded_result_id = quote(video.result_id, safe="") base = f"/api/v1/observation-sessions/{encoded_session_id}/perception-media/{encoded_result_id}" return { - "schema_version": RECORDED_MEDIA_STREAM_MANIFEST_SCHEMA, + "schema_version": RECORDED_PERCEPTION_STREAM_MANIFEST_SCHEMA, "source_id": video.public_source_id, "generation_sha256": video.sha256, "byte_length": video.byte_length, @@ -1795,6 +1796,13 @@ def _recorded_media_manifest_document( "media_type": epoch.media_type, "byte_length": epoch.init_byte_length + sum(segment.byte_length for segment in epoch.segments), + "segment_count": len(epoch.segments), + "random_access_sequences": sorted( + segment.sequence for segment in epoch.segments if segment.random_access + ), + "segment_end_times_seconds": [ + segment.end_time_seconds for segment in epoch.segments + ], "stream_url": ( f"{base}/epochs/{epoch.ordinal}/recording.mp4" f"?generation={manifest.generation_sha256}" diff --git a/tests/test_session_api.py b/tests/test_session_api.py index b0cbdff..b48f7e3 100644 --- a/tests/test_session_api.py +++ b/tests/test_session_api.py @@ -41,6 +41,7 @@ from k1link.web.session_api import ( RecordedPerceptionRequest, RecordedPointColorsRequest, ReplayRequest, + _recorded_media_manifest_document, build_session_router, ) @@ -128,6 +129,11 @@ def make_recorded_h264_fixture( timescale: int = 1_000, sample_duration: int = 500, base_decode_time: int = 0, + trex_default_sample_flags: int = 0, + tfhd_default_sample_flags: int | None = None, + first_sample_flags: int | None = None, + sample_flags: int | None = None, + fragment_sample_count: int = 1, ) -> tuple[bytes, bytes]: def box(box_type: bytes, payload: bytes = b"") -> bytes: return (8 + len(payload)).to_bytes(4, "big") + box_type + payload @@ -156,14 +162,32 @@ def make_recorded_h264_fixture( track_id.to_bytes(4, "big") + (1).to_bytes(4, "big") + sample_duration.to_bytes(4, "big") - + b"\x00" * 8, + + b"\x00" * 4 + + trex_default_sample_flags.to_bytes(4, "big"), ) avcc = box(b"avcC", b"\x01\x64\x00\x28") init = box(b"ftyp", b"isom") + box(b"moov", trak + box(b"mvex", trex) + avcc) - tfhd = full_box(b"tfhd", track_id.to_bytes(4, "big"), flags=0x020000) + tfhd_flags = 0x020000 + tfhd_payload = track_id.to_bytes(4, "big") + if tfhd_default_sample_flags is not None: + tfhd_flags |= 0x000020 + tfhd_payload += tfhd_default_sample_flags.to_bytes(4, "big") + tfhd = full_box(b"tfhd", tfhd_payload, flags=tfhd_flags) tfdt = full_box(b"tfdt", base_decode_time.to_bytes(4, "big")) - trun = full_box(b"trun", (1).to_bytes(4, "big")) + trun_flags = 0 + trun_payload = fragment_sample_count.to_bytes(4, "big") + if first_sample_flags is not None: + trun_flags |= 0x000004 + trun_payload += first_sample_flags.to_bytes(4, "big") + if sample_flags is not None: + trun_flags |= 0x000400 + trun_payload += sample_flags.to_bytes(4, "big") * fragment_sample_count + trun = full_box( + b"trun", + trun_payload, + flags=trun_flags, + ) fragment = box(b"moof", box(b"traf", tfhd + tfdt + trun)) + box(b"mdat", b"frame") return init, fragment @@ -1644,7 +1668,7 @@ def test_session_router_exposes_opaque_recorded_media_manifest_and_ranges( if_match=f'"sha256:{source["manifest_generation_sha256"]}"', ) manifest = json.loads(manifest_response.body) - assert manifest["schema_version"] == "missioncore.observation-recorded-media/v3" + assert manifest["schema_version"] == "missioncore.observation-recorded-media/v4" assert manifest["source_id"] == source["id"] assert manifest["generation_sha256"] == source["manifest_generation_sha256"] assert manifest["byte_length"] == source["byte_length"] @@ -1659,6 +1683,9 @@ def test_session_router_exposes_opaque_recorded_media_manifest_and_ranges( "timeline_end_seconds": 0.5, "media_type": 'video/mp4; codecs="avc1.640028"', "byte_length": len(init) + len(segment), + "segment_count": 1, + "random_access_sequences": [1], + "segment_end_times_seconds": [0.5], "stream_url": ( f"{source['manifest_url'].removesuffix('/manifest')}/epochs/1/recording.mp4" f"?generation={generation}" @@ -1815,6 +1842,136 @@ def test_session_router_exposes_opaque_recorded_media_manifest_and_ranges( assert confined.value.status_code == 409 +def test_recorded_media_persists_and_exposes_random_access_sequences( + tmp_path: Path, +) -> None: + repository = tmp_path / "repo" + sessions = repository / "sessions" + session = make_legacy_session(sessions, "20260716T205632Z_viewer_live") + init, first = make_recorded_h264_fixture() + _, second = make_recorded_h264_fixture( + base_decode_time=500, + tfhd_default_sample_flags=0x00010000, + ) + _, third = make_recorded_h264_fixture( + base_decode_time=1_000, + tfhd_default_sample_flags=0x00010000, + first_sample_flags=0, + ) + _, fourth = make_recorded_h264_fixture( + base_decode_time=1_500, + sample_flags=0x00010000, + ) + writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1) + writer.append("init", init) + for sequence, fragment in enumerate((first, second, third, fourth), start=1): + writer.append( + "media", + fragment, + host_epoch_ns=1_000_000_000 + sequence * 500_000_000, + host_monotonic_ns=2_000_000_000 + sequence * 500_000_000, + ) + writer.close() + + store = SessionStore(repository, data_dir=tmp_path / "data") + store.reconcile_archive(xgrids_k1_archive_source(sessions)) + command = store.prepare_replay(session.name) + artifact = store.list_recorded_media(session.name)[0] + cache_root = tmp_path / "prepared-media" + manifest = RecordedMediaInspector(cache_root).inspect(artifact, command) + + assert [segment.random_access for segment in manifest.epochs[0].segments] == [ + True, + False, + True, + False, + ] + document = _recorded_media_manifest_document(manifest) + assert document["schema_version"] == "missioncore.observation-recorded-media/v4" + assert document["epochs"][0]["segment_count"] == 4 + assert document["epochs"][0]["random_access_sequences"] == [1, 3] + assert document["epochs"][0]["segment_end_times_seconds"] == [0.5, 1.0, 1.5, 2.0] + assert document["epochs"][0]["segment_end_times_seconds"][-1] == ( + manifest.epochs[0].timeline_end_seconds - manifest.epochs[0].timeline_start_seconds + ) + + sidecar = json.loads(next(cache_root.glob("*.json")).read_text(encoding="utf-8")) + assert sidecar["schema_version"] == "missioncore.recorded-media-preparation/v3" + assert sidecar["manifest"]["schema_version"] == ("missioncore.observation-recorded-media/v4") + assert [ + segment["random_access"] for segment in sidecar["manifest"]["epochs"][0]["segments"] + ] == [True, False, True, False] + assert [ + segment["end_time_seconds"] for segment in sidecar["manifest"]["epochs"][0]["segments"] + ] == [0.5, 1.0, 1.5, 2.0] + restarted = RecordedMediaInspector(cache_root).inspect(artifact, command) + assert [segment.end_time_seconds for segment in restarted.epochs[0].segments] == [ + 0.5, + 1.0, + 1.5, + 2.0, + ] + + +def test_recorded_media_rejects_ambiguous_trun_sample_flags() -> None: + init, fragment = make_recorded_h264_fixture( + first_sample_flags=0, + sample_flags=0, + ) + timing = recorded_media_module._mp4_video_timing( + init, + recorded_media_module._Mp4ParseBudget(), + ) + + with pytest.raises(SessionIntegrityError, match="sample flags are ambiguous"): + recorded_media_module._mp4_video_fragment_timing( + fragment, + timing, + recorded_media_module._Mp4ParseBudget(), + ) + + +def test_recorded_media_rejects_multi_sample_video_fragment() -> None: + init, fragment = make_recorded_h264_fixture(fragment_sample_count=2) + timing = recorded_media_module._mp4_video_timing( + init, + recorded_media_module._Mp4ParseBudget(), + ) + + with pytest.raises(SessionIntegrityError, match="exactly one sample"): + recorded_media_module._mp4_video_fragment_timing( + fragment, + timing, + recorded_media_module._Mp4ParseBudget(), + ) + + +def test_recorded_media_rejects_codec_epoch_without_initial_random_access( + tmp_path: Path, +) -> None: + repository = tmp_path / "repo" + sessions = repository / "sessions" + session = make_legacy_session(sessions, "20260716T205632Z_viewer_live") + init, non_random_access = make_recorded_h264_fixture(sample_flags=0x00010000) + writer = CameraArchiveWriter(session, "sensor.camera.private-left", 1) + writer.append("init", init) + writer.append( + "media", + non_random_access, + host_epoch_ns=1_500_000_000, + host_monotonic_ns=2_500_000_000, + ) + writer.close() + store = SessionStore(repository, data_dir=tmp_path / "data") + store.reconcile_archive(xgrids_k1_archive_source(sessions)) + + with pytest.raises(SessionIntegrityError, match="does not begin with a random-access"): + RecordedMediaInspector().inspect( + store.list_recorded_media(session.name)[0], + store.prepare_replay(session.name), + ) + + def test_replay_exposes_generation_bound_perception_video_with_native_ranges( tmp_path: Path, ) -> None: @@ -2373,6 +2530,10 @@ def test_recorded_media_preparation_sidecar_reuses_and_rebuilds_generation( sidecar = sidecars[0] assert stat.S_IMODE(sidecar.stat().st_mode) == 0o600 assert str(session) not in sidecar.read_text(encoding="utf-8") + prepared = json.loads(sidecar.read_text(encoding="utf-8")) + assert prepared["schema_version"] == "missioncore.recorded-media-preparation/v3" + assert prepared["manifest"]["epochs"][0]["segments"][0]["random_access"] is True + assert prepared["manifest"]["epochs"][0]["segments"][0]["end_time_seconds"] == 0.5 def forbidden_reparse(*_args: object, **_kwargs: object) -> None: raise AssertionError("restart must reuse the prepared media sidecar") @@ -2383,6 +2544,7 @@ def test_recorded_media_preparation_sidecar_reuses_and_rebuilds_generation( restarted = RecordedMediaInspector(cache_root).inspect(artifact, command) assert restarted.generation_sha256 == initial.generation_sha256 assert restarted.timeline_end_seconds == initial.timeline_end_seconds + assert [segment.end_time_seconds for segment in restarted.epochs[0].segments] == [0.5] original_read_manifest = recorded_media_module._read_manifest reparses = 0 @@ -2392,6 +2554,28 @@ def test_recorded_media_preparation_sidecar_reuses_and_rebuilds_generation( reparses += 1 return original_read_manifest(*args, **kwargs) + old_body = json.loads(json.dumps(prepared)) + old_body.pop("checksum_sha256") + old_body["schema_version"] = "missioncore.recorded-media-preparation/v2" + old_body["manifest"]["schema_version"] = "missioncore.observation-recorded-media/v3" + for epoch in old_body["manifest"]["epochs"]: + for old_segment in epoch["segments"]: + old_segment.pop("end_time_seconds") + old_checksum = hashlib.sha256(recorded_media_module._canonical_json(old_body)).hexdigest() + sidecar.write_bytes( + recorded_media_module._canonical_json({**old_body, "checksum_sha256": old_checksum}) + ) + with monkeypatch.context() as context: + context.setattr(recorded_media_module, "_read_manifest", counted_reparse) + migrated = RecordedMediaInspector(cache_root).inspect(artifact, command) + assert migrated.generation_sha256 == initial.generation_sha256 + assert reparses == 1 + migrated_sidecar = json.loads(sidecar.read_text(encoding="utf-8")) + assert migrated_sidecar["schema_version"] == "missioncore.recorded-media-preparation/v3" + assert migrated_sidecar["manifest"]["epochs"][0]["segments"][0]["random_access"] is True + assert migrated_sidecar["manifest"]["epochs"][0]["segments"][0]["end_time_seconds"] == 0.5 + + reparses = 0 sidecar.write_bytes(b"{corrupt") with monkeypatch.context() as context: context.setattr(recorded_media_module, "_read_manifest", counted_reparse)