From ec79efb921d6c70569114336311e59e4805f67b2 Mon Sep 17 00:00:00 2001 From: Kingston Date: Sat, 19 Sep 2026 22:16:47 -0700 Subject: [PATCH] Add reusable video preparation and raw measurement APIs --- docs/README.md | 1 + docs/VIDEO_PRIMITIVES.md | 72 ++++ docs/how-to/sample-source-video.md | 19 + pyproject.toml | 2 + .../_video_measurements/_frame_statistics.py | 15 +- src/hflow/blur.py | 105 +++++ src/hflow/camera_motion.py | 109 +++++- src/hflow/importers/video.py | 63 ++- src/hflow/mcap_video.py | 363 ++++++++++++++++++ src/hflow/source_sampling.py | 190 ++++++++- src/hflow/video_statistics.py | 60 +++ tests/test_blur.py | 74 ++++ tests/test_camera_shake_summary.py | 117 ++++++ tests/test_mcap_video.py | 246 ++++++++++++ tests/test_source_sampling.py | 70 ++++ tests/test_video_import.py | 53 +++ tests/test_video_statistics.py | 87 +++++ uv.lock | 34 +- 18 files changed, 1662 insertions(+), 18 deletions(-) create mode 100644 docs/VIDEO_PRIMITIVES.md create mode 100644 src/hflow/blur.py create mode 100644 src/hflow/mcap_video.py create mode 100644 src/hflow/video_statistics.py create mode 100644 tests/test_blur.py create mode 100644 tests/test_camera_shake_summary.py create mode 100644 tests/test_mcap_video.py create mode 100644 tests/test_video_statistics.py diff --git a/docs/README.md b/docs/README.md index 07a2acd2..7834b32b 100644 --- a/docs/README.md +++ b/docs/README.md @@ -82,6 +82,7 @@ Reference pages define stable inputs, outputs, configuration, and stored-data contracts. - [Embedded integration boundaries](./EMBEDDED_BOUNDARIES.md) +- [Local video preparation and measurements](./VIDEO_PRIMITIVES.md) - [Canonical episode format](./FORMAT.md) - [Catalog tables and curation API](./CATALOG.md) - [Environment variables](./ENVIRONMENT.md) diff --git a/docs/VIDEO_PRIMITIVES.md b/docs/VIDEO_PRIMITIVES.md new file mode 100644 index 00000000..cf90434e --- /dev/null +++ b/docs/VIDEO_PRIMITIVES.md @@ -0,0 +1,72 @@ +# Local video preparation and measurements + +These file-level APIs do not require a catalog, scheduler, model service, or +persistent workspace. Callers retain source identity and own output lifetimes. +FFmpeg-based APIs use HFlow's binary policy and may download managed binaries +unless the caller configures an installed toolchain. + +## MCAP camera export + +`hflow.mcap_video.export_mcap_camera(source, output, camera_topic=None, limits=...)` +requires the optional `hflow[video]` dependency. It streams one camera from an MCAP +into MP4, returning `PreparedMcapVideo` with the output path, selected topic, +original first log timestamp, and duration in milliseconds. Multiple cameras +require an exact topic. Existing destinations are never overwritten. + +Frame times are MCAP log times relative to the first selected frame, rounded to +microseconds. Each frame lasts until the next; the final frame repeats the prior +positive interval. At least two frames are required. Duplicate rounded timestamps, +unsupported formats, changing dimensions, B-frames, missing H.264 codec headers, +and gaps beyond the MP4 duration representation are rejected. H.264 access units +are remuxed with lossless AUD repair; supported JPEG/PNG and 8-bit raw images are +encoded as lossless RGB H.264. Raw recordings remain unchanged. + +This is synchronous native decoding. A caller needing a hard execution deadline +must isolate it in a process and reap that process before deleting temporary media. +The existing canonical `camera_video` enrichment keeps its constant-rate contract. + +## Direct model-video preparation + +`hflow.importers.video.prepare_model_video(source, output, config, limits=..., +transform_config=...)` uses `VideoImportConfig` and `TransformConfig` to produce +canonical model-input pixels without first writing an MCAP. It shares the video +importer's fixed-rate JPEG rendering and canonical H.264 encoding, including the +single-frame cadence. Output is an atomically published caller-owned MP4. + +Expected media failures return `UnreadableVideo` or `UnsupportedVideo`; operational +failures raise. Existing destinations raise `FileExistsError`. The JPEG intermediate +is intentional: bypassing it would change model-input pixels. This helper does not +change canonical transformation defaults or identities. + +## Frame statistics + +`hflow.video_statistics.measure_video_frame_statistics(video, settings=..., +toolchain=None, instrument_output_cache_path=None)` exposes file-level statistics, +settings, provenance, and errors as a public API. No persistent cache is created +unless explicitly requested. + +`FrameStatisticsSettings.luma_range` accepts `LumaRangePolicy.PRESERVE` (the existing +instrument behavior) or `FULL` (convert the declared input range to full-range luma +inside the measurement filter graph). The graph and settings are recorded in +provenance. Measure unpadded source framing; black model-input borders would affect +pixel statistics. Full-range measurement needs no intermediate video encode. + +## Raw blur and shake summaries + +`hflow.blur.measure_video_blur(video, timeout_seconds=120, executable=None)` uses +FFmpeg's default `blurdetect` settings at the supplied frame cadence and resolution. +`BlurSummary` contains observed/scored frame counts and the mean finite raw score. +Nonfinite scores are unassessed, zero assessed frames produce a null mean, and +negative scores fail. `summarize_blur_scores` applies the same accounting to a +caller-supplied score stream. + +`hflow.camera_motion.summarize_camera_shake(observations)` consumes filtered motion +observations in constant summary memory. Mean and RMS weight adjacent-pair duration; +maximum, assessment duration, and missing-context counts remain explicit. The final +frame has no following pair and contributes no extra duration. Missing observations +never become zero shake. Field of view is caller configuration, not calibration. + +These measurements do not classify footage as blurry or unstable and are not +accuracy estimates. See [streaming camera motion](how-to/stream-camera-motion.md) +for extraction and filtering contracts, and [source sampling](how-to/sample-source-video.md) +for original-frame evidence selection. diff --git a/docs/how-to/sample-source-video.md b/docs/how-to/sample-source-video.md index 950cc235..5cbca0ef 100644 --- a/docs/how-to/sample-source-video.md +++ b/docs/how-to/sample-source-video.md @@ -96,3 +96,22 @@ and selected FFmpeg build in their version or input contract. Changing selection can change measured results. This is a new original-source API; existing `Episode.frames()` and Build AI `FrameSampling` keep their declared frame-rate behavior. No canonical transform version changes are needed to use these helpers. + +## Nearest keyframes and unpadded evidence + +`SourceSamplingMode.NEAREST_KEYFRAMES` selects unique codec keyframes nearest +caller-supplied `keyframe_positions`, expressed as fractions of the requested +window. Ties choose the earlier frame; results are chronological. Empty windows +remain empty and this mode has no uniform fallback. Position count must fit +`maximum_frames`. + +Import `SourceFrameResize` from `hflow.source_sampling`. Set `resize=SourceFrameResize.FIT` +to fit within `width` and `height` without adding padding; the default `PAD` retains +the existing padded canvas. `scaling_algorithm` accepts `lanczos` (default) or +`bicubic`, and `jpeg_quality` accepts FFmpeg quality values 1 through 31 (default 5). +For example, a caller can choose positions `(0.15, 0.5, 0.85)`, a 960×960 fit box, +and JPEG quality 2. These are evidence settings, not a classification policy. + +Nearest-keyframe probing shares the extraction deadline and is bounded by +`maximum_probe_bytes`. Returned timestamps retain the source's playback clock and +rational time base; subtract the window start explicitly for relative timestamps. diff --git a/pyproject.toml b/pyproject.toml index 304a4059..9e0f6c66 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -54,6 +54,7 @@ mediapipe = ["mediapipe>=1.0.1"] # dependency is real, and confined to the one check that needs it. Headless # because a data pipeline runs where there is no display. motion = ["opencv-python-headless>=5.0.0.93"] +video = ["av>=18.1.0"] native-build = ["Cython>=3.3.0", "setuptools>=84.0.0"] openai = ["openai>=3.3.1"] @@ -67,6 +68,7 @@ hflow = "hflow.cli:main" [dependency-groups] dev = [ + "av>=18.1.0", # the video extra, exercised by timestamp-preserving camera export tests "Cython>=3.3.0", # the native-build extra, present so its real build path is tested "hflow-server", # the workspace API package, present in dev so the suite tests it "obstore>=0.11.1", # the bucket extra, present in dev so the suite tests it diff --git a/src/hflow/_video_measurements/_frame_statistics.py b/src/hflow/_video_measurements/_frame_statistics.py index 871d747e..b7a76821 100644 --- a/src/hflow/_video_measurements/_frame_statistics.py +++ b/src/hflow/_video_measurements/_frame_statistics.py @@ -41,6 +41,13 @@ class UnsupportedVideoMeasurementToolchainError(RuntimeError): """The supplied FFmpeg build lacks a filter required by the measurement.""" +class LumaRangePolicy(StrEnum): + """Whether measurements retain decoded luma or normalize its declared range.""" + + PRESERVE = "preserve" + FULL = "full" + + class LumaRangeEvidence(StrEnum): """What decoded luma samples show about the nominal limited range.""" @@ -76,8 +83,11 @@ class FrameStatisticsSettings: freeze_noise_tolerance_decibels: float = -60.0 freeze_minimum_duration_seconds: float = 2.0 overexposed_average_luma_threshold: float = 235.0 + luma_range: LumaRangePolicy = LumaRangePolicy.PRESERVE def __post_init__(self) -> None: + if not isinstance(self.luma_range, LumaRangePolicy): + raise ValueError("luma_range must be a LumaRangePolicy") require_int( self.black_frame_minimum_pixel_share_percent, "black_frame_minimum_pixel_share_percent", @@ -595,7 +605,10 @@ def _temporary_instrument_cache_output( def frame_statistics_filter_graph(settings: FrameStatisticsSettings) -> str: """Return the effective single-pass FFmpeg measurement graph.""" - return ( + normalization = ( + "scale=in_range=auto:out_range=full," if settings.luma_range is LumaRangePolicy.FULL else "" + ) + return normalization + ( "format=pix_fmts=yuv420p," f"blackframe=amount=0:threshold={settings.black_pixel_luma_threshold}," "freezedetect=" diff --git a/src/hflow/blur.py b/src/hflow/blur.py new file mode 100644 index 00000000..65326adf --- /dev/null +++ b/src/hflow/blur.py @@ -0,0 +1,105 @@ +"""Raw FFmpeg blur scores with explicit finite-score coverage, not quality labels.""" + +from __future__ import annotations + +import math +import subprocess +import tempfile +from collections.abc import Iterable +from dataclasses import dataclass +from pathlib import Path + +from hflow.ffmpeg import ffmpeg_path + + +@dataclass(frozen=True, slots=True) +class BlurSummary: + frame_count: int + scored_frame_count: int + mean_blur_score: float | None + + +def summarize_blur_scores(scores: Iterable[float]) -> BlurSummary: + """Average finite raw frame scores; unavailable scores never become zero.""" + frame_count = 0 + scored_frame_count = 0 + mean_blur_score = 0.0 + for score in scores: + frame_count += 1 + # blurdetect emits NaN when it cannot measure edges, e.g. a flat image. + if not math.isfinite(score): + continue + if score < 0: + raise RuntimeError("Blur analysis returned an invalid score") + scored_frame_count += 1 + mean_blur_score += (score - mean_blur_score) / scored_frame_count + return BlurSummary( + frame_count=frame_count, + scored_frame_count=scored_frame_count, + mean_blur_score=mean_blur_score if scored_frame_count else None, + ) + + +def measure_video_blur( + video_path: Path, *, timeout_seconds: float = 120.0, executable: Path | None = None +) -> BlurSummary: + """Score a local video window at its supplied resolution and frame cadence. + + Uses FFmpeg's default blurdetect settings on the first video stream. The + result is the frame-weighted mean of finite raw scores, not a percentage or + probability. No resizing, frame sampling, classification, or model call is + performed. Results contain raw scores and coverage, without affected-time percentages. + """ + if not math.isfinite(timeout_seconds) or timeout_seconds <= 0: + raise ValueError("timeout_seconds must be positive and finite") + try: + with tempfile.TemporaryFile() as metadata_output: + subprocess.run( + [ + str(executable or ffmpeg_path()), + "-hide_banner", + "-loglevel", + "error", + "-nostdin", + "-xerror", + "-protocol_whitelist", + "file", + "-noautorotate", + "-threads", + "1", + "-i", + str(video_path.resolve()), + "-map", + "0:v:0", + "-an", + "-sn", + "-dn", + "-filter_threads", + "1", + "-vf", + "blurdetect,metadata=mode=print:key=lavfi.blur:file=-", + "-fps_mode", + "passthrough", + "-f", + "null", + "-", + ], + stdin=subprocess.DEVNULL, + stdout=metadata_output, + stderr=subprocess.PIPE, + timeout=timeout_seconds, + check=True, + ) + metadata_output.seek(0) + summary = summarize_blur_scores( + float(line.removeprefix(b"lavfi.blur=")) + for line in metadata_output + if line.startswith(b"lavfi.blur=") + ) + metadata_output.seek(0) + emitted_frame_count = sum(line.startswith(b"frame:") for line in metadata_output) + if summary.frame_count != emitted_frame_count: + raise RuntimeError("Blur analysis returned incomplete frame scores") + return summary + except Exception as error: + raise RuntimeError("Blur analysis failed") from error diff --git a/src/hflow/camera_motion.py b/src/hflow/camera_motion.py index 33d528d4..1d79e1b3 100644 --- a/src/hflow/camera_motion.py +++ b/src/hflow/camera_motion.py @@ -5,9 +5,12 @@ episodes, catalogs, or orchestration. See docs/how-to/stream-camera-motion.md. """ -from collections.abc import Iterator +import math +from collections.abc import Iterable, Iterator from contextlib import contextmanager +from dataclasses import dataclass from pathlib import Path +from typing import assert_never from hflow._video_measurement_toolchain import resolved_video_measurement_toolchain from hflow._video_measurements._motion_fit import ( @@ -41,6 +44,7 @@ "CameraMotionTransform", "CameraShakeObservation", "CameraShakeSettings", + "CameraShakeSummary", "MeasuredCameraMotion", "MeasuredCameraShake", "MotionFitEvidence", @@ -50,6 +54,7 @@ "filter_camera_shake", "iter_frame_motion", "stream_camera_motion", + "summarize_camera_shake", ] @@ -76,3 +81,105 @@ def stream_camera_motion( yield observations finally: observations.close() + + +@dataclass(frozen=True, slots=True) +class CameraShakeSummary: + """Continuous residual rates and their observed frame-pair coverage. + + Angular rates use the caller's approximate field of view. A fitted motion + estimate can still be unreliable; coverage does not establish accuracy. + The final frame's display duration is outside the observed pair intervals. + """ + + pair_count: int + measured_motion_pair_count: int + measured_shake_pair_count: int + insufficient_context_pair_count: int + unmeasured_context_pair_count: int + observed_seconds: float + assessed_seconds: float + mean_shake_degrees_per_second: float | None + rms_shake_degrees_per_second: float | None + maximum_shake_degrees_per_second: float | None + + @property + def unassessed_seconds(self) -> float: + return self.observed_seconds - self.assessed_seconds + + @property + def assessed_fraction(self) -> float | None: + if self.observed_seconds == 0: + return None + return self.assessed_seconds / self.observed_seconds + + +def summarize_camera_shake(observations: Iterable[CameraShakeObservation]) -> CameraShakeSummary: + """Reduce a complete filtered motion stream with constant summary memory. + + Missing motion and filter context stay unassessed, never zero shake. Rates + are weighted by the adjacent-pair durations; no video frames or rate history + are retained. Ordering and cadence belong to HFlow's filter contract. + """ + pair_count = 0 + measured_motion_pair_count = 0 + measured_shake_pair_count = 0 + insufficient_context_pair_count = 0 + unmeasured_context_pair_count = 0 + observed_seconds = 0.0 + assessed_seconds = 0.0 + mean_shake_degrees_per_second = 0.0 + rms_shake_degrees_per_second = 0.0 + maximum_shake_degrees_per_second = 0.0 + + for observation in observations: + pair_count += 1 + pair_seconds = observation.motion.end_seconds - observation.motion.start_seconds + observed_seconds += pair_seconds + if isinstance(observation.motion.measurement, MeasuredCameraMotion): + measured_motion_pair_count += 1 + + shake = observation.shake + if isinstance(shake, MeasuredCameraShake): + measured_shake_pair_count += 1 + assessed_seconds += pair_seconds + shake_degrees_per_second = shake.residual.magnitude_degrees_per_second + duration_share = pair_seconds / assessed_seconds + mean_shake_degrees_per_second += duration_share * ( + shake_degrees_per_second - mean_shake_degrees_per_second + ) + # The weighted Euclidean norm avoids squaring large finite rates. + rms_shake_degrees_per_second = math.hypot( + rms_shake_degrees_per_second * math.sqrt(1 - duration_share), + shake_degrees_per_second * math.sqrt(duration_share), + ) + maximum_shake_degrees_per_second = max( + maximum_shake_degrees_per_second, shake_degrees_per_second + ) + else: + match shake.reason: + case "insufficient_context": + insufficient_context_pair_count += 1 + case "unmeasured_context": + unmeasured_context_pair_count += 1 + case unknown_reason: + assert_never(unknown_reason) + + return CameraShakeSummary( + pair_count=pair_count, + measured_motion_pair_count=measured_motion_pair_count, + measured_shake_pair_count=measured_shake_pair_count, + insufficient_context_pair_count=insufficient_context_pair_count, + unmeasured_context_pair_count=unmeasured_context_pair_count, + observed_seconds=observed_seconds, + assessed_seconds=assessed_seconds, + mean_shake_degrees_per_second=( + mean_shake_degrees_per_second if measured_shake_pair_count else None + ), + rms_shake_degrees_per_second=( + rms_shake_degrees_per_second if measured_shake_pair_count else None + ), + maximum_shake_degrees_per_second=( + maximum_shake_degrees_per_second if measured_shake_pair_count else None + ), + ) diff --git a/src/hflow/importers/video.py b/src/hflow/importers/video.py index c1103151..22e68add 100644 --- a/src/hflow/importers/video.py +++ b/src/hflow/importers/video.py @@ -20,8 +20,10 @@ from hflow._pinned_asset import sha256_hex_of_file from hflow.ffmpeg import ffmpeg_path, ffmpeg_version from hflow.ffmpeg._process import media_input_was_rejected, run_media_command -from hflow.format import METADATA_RECORD_EPISODE, NANOSECONDS_PER_SECOND +from hflow.format import GOP_SECONDS, METADATA_RECORD_EPISODE, NANOSECONDS_PER_SECOND from hflow.media import UnreadableVideo, UnsupportedVideo, VideoLimits, VideoProperties, probe_video +from hflow.transform import TransformConfig +from hflow.video import encode_images_to_h264, write_access_units_to_mp4 _IMPORT_METADATA_RECORD = "video_import/v1" _MAXIMUM_TIMESTAMP_NS = (1 << 64) - 1 @@ -335,3 +337,62 @@ def prepare_video_episode( return UnreadableVideo() except _UnsupportedExcerpt: return UnsupportedVideo() + + +def prepare_model_video( + source_video: Path, + output: Path, + config: VideoImportConfig, + *, + limits: VideoLimits = VideoLimits(), + transform_config: TransformConfig = TransformConfig(), +) -> Path | UnreadableVideo | UnsupportedVideo: + """Prepare canonical model-input pixels directly, without an intermediate MCAP. + + Shares the importer's fixed-rate JPEG recipe, then the canonical H.264 encoder. + This preserves the intentional JPEG compression and single-frame cadence. + No catalog, episode or persistent workspace is created. Output is published + atomically without replacing any existing path; the caller owns its lifetime. + """ + source_video = source_video.resolve(strict=True) + if output.exists() or output.is_symlink(): + raise FileExistsError(output) + inspection = probe_video(source_video, limits=limits) + if not isinstance(inspection, VideoProperties): + return inspection + if ( + config.image_width * config.image_height > limits.maximum_frame_pixels + or config.image_hz > limits.maximum_frames_per_second + or config.duration_s > limits.maximum_duration_seconds + ): + return UnsupportedVideo() + try: + _require_excerpt_duration(inspection, config) + except _UnsupportedExcerpt: + return UnsupportedVideo() + output.parent.mkdir(parents=True, exist_ok=True) + with tempfile.TemporaryDirectory(dir=output.parent, prefix=".model-video-") as directory: + frame_directory = Path(directory) + try: + _render_frames(source_video, config, frame_directory, limits) + except _UnreadableImport: + return UnreadableVideo() + frame_paths = tuple( + frame_directory / f"frame_{index + 1:010d}.jpg" for index in range(config.frame_count) + ) + if not all(frame_path.is_file() for frame_path in frame_paths): + return UnreadableVideo() + frames_per_second = config.image_hz if len(frame_paths) >= 2 else 1.0 + access_units = encode_images_to_h264( + [frame_path.read_bytes() for frame_path in frame_paths], + fps=frames_per_second, + gop_frames=max(1, round(GOP_SECONDS[transform_config.gop_preset] * frames_per_second)), + crf=transform_config.crf, + ) + staged_video = write_access_units_to_mp4( + (access_unit.data for access_unit in access_units), + fps=frames_per_second, + output=frame_directory / "video.mp4", + ) + os.link(staged_video, output) + return output diff --git a/src/hflow/mcap_video.py b/src/hflow/mcap_video.py new file mode 100644 index 00000000..d2b2b7d5 --- /dev/null +++ b/src/hflow/mcap_video.py @@ -0,0 +1,363 @@ +"""Stream one MCAP camera to MP4 while preserving irregular log-time intervals. + +Requires the optional ``video`` extra. Decoding is synchronous; callers own +process isolation when a hard native-decoder deadline is required. +""" + +from __future__ import annotations + +import base64 +import math +import os +import tempfile +from collections.abc import Iterator, Mapping +from dataclasses import dataclass, field, replace +from fractions import Fraction +from itertools import chain +from pathlib import Path + +import av +import numpy as np +from av.container import OutputContainer +from av.stream import Stream +from av.video.codeccontext import VideoCodecContext +from av.video.stream import VideoStream + +from hflow import Episode, TopicInfo +from hflow.format import PASSTHROUGH_VIDEO_SCHEMA_NAMES +from hflow.media import VideoLimits +from hflow.video import ( + _annex_b_nal_offsets_and_types, + ensure_access_unit_delimiter, + scan_picture_coding_types, + split_annex_b_stream, +) + +MAXIMUM_FRAME_BYTES = 32 * 1024 * 1024 +NANOSECONDS_PER_SECOND = 1_000_000_000 +VIDEO_TICKS_PER_SECOND = 1_000_000 +NANOSECONDS_PER_VIDEO_TICK = NANOSECONDS_PER_SECOND // VIDEO_TICKS_PER_SECOND +MAXIMUM_VIDEO_INTERVAL_TICKS = 2**31 - 1 +VIDEO_TIME_BASE = Fraction(1, VIDEO_TICKS_PER_SECOND) +COMPRESSED_IMAGE_SCHEMAS = frozenset( + {"foxglove.CompressedImage", "sensor_msgs/msg/CompressedImage"} +) +RAW_IMAGE_SCHEMAS = frozenset({"foxglove.RawImage", "sensor_msgs/msg/Image"}) +RAW_PIXEL_FORMATS = { + "rgb8": ("rgb24", 3), + "bgr8": ("bgr24", 3), + "rgba8": ("rgba", 4), + "bgra8": ("bgra", 4), + "mono8": ("gray", 1), +} + + +class McapMediaError(RuntimeError): + """The selected camera could not be prepared within the supported profile.""" + + +@dataclass(frozen=True, slots=True) +class PreparedMcapVideo: + video_path: Path = field(repr=False) + duration_millis: int + camera_topic: str = field(repr=False) + start_timestamp_ns: int + + +@dataclass(frozen=True, slots=True) +class CameraFrameInterval: + timestamp_ns: int + duration_ns: int + message: object = field(repr=False) + + +def _field(message: object, name: str) -> object: + return message.get(name) if isinstance(message, Mapping) else getattr(message, name, None) + + +def _frame_bytes(message: object) -> bytes: + payload = _field(message, "data") + if isinstance(payload, str): + if len(payload) > MAXIMUM_FRAME_BYTES * 2: + raise McapMediaError() + payload = base64.b64decode(payload, validate=True) + if not isinstance(payload, bytes | bytearray | memoryview | list | np.ndarray): + raise McapMediaError() + if not 0 < len(payload) <= MAXIMUM_FRAME_BYTES: + raise McapMediaError() + return bytes(payload) + + +def _select_camera(episode: Episode, camera_topic: str | None) -> TopicInfo: + cameras = episode.cameras + if camera_topic is None: + if len(cameras) != 1: + raise McapMediaError() + camera_topic = cameras[0] + if camera_topic not in cameras: + raise McapMediaError() + channels = [channel for channel in episode.channels.values() if channel.topic == camera_topic] + if len(channels) != 1: + raise McapMediaError() + return channels[0] + + +def _camera_messages( + episode: Episode, channel: TopicInfo, limits: VideoLimits +) -> Iterator[tuple[int, object]]: + maximum_frames = math.ceil(limits.maximum_duration_seconds * limits.maximum_frames_per_second) + frame_count = 0 + for batch in episode.iter_decoded_batches( + channel_ids=[channel.channel_id], batch_max_messages=32, batch_max_bytes=8 * 1024 * 1024 + ): + for timestamp, message in zip(batch.log_times, batch.messages, strict=True): + frame_count += 1 + if frame_count > maximum_frames: + raise McapMediaError() + yield int(timestamp), message + + +def _frame_intervals( + messages: Iterator[tuple[int, object]], limits: VideoLimits +) -> Iterator[CameraFrameInterval]: + previous_timestamp, previous_message = next(messages) + start_timestamp = previous_timestamp + previous_interval = 0 + frame_count = 1 + maximum_duration_ns = int(limits.maximum_duration_seconds * NANOSECONDS_PER_SECOND) + for timestamp, message in messages: + interval = timestamp - previous_timestamp + if interval <= 0 or timestamp - start_timestamp > maximum_duration_ns: + raise McapMediaError() + frame_count += 1 + yield CameraFrameInterval(previous_timestamp - start_timestamp, interval, previous_message) + previous_timestamp, previous_message, previous_interval = timestamp, message, interval + final_duration = previous_timestamp - start_timestamp + previous_interval + if previous_interval == 0 or not ( + limits.minimum_duration_seconds * NANOSECONDS_PER_SECOND + <= final_duration + <= maximum_duration_ns + ): + raise McapMediaError() + # Log timestamps measure arrival time and can bunch together during capture. + # Bound sustained density while preserving those positive timestamp gaps. + if frame_count * NANOSECONDS_PER_SECOND > limits.maximum_frames_per_second * final_duration: + raise McapMediaError() + yield CameraFrameInterval( + previous_timestamp - start_timestamp, previous_interval, previous_message + ) + + +def _require_dimensions(width: int, height: int, limits: VideoLimits) -> None: + if width <= 0 or height <= 0 or width * height > limits.maximum_frame_pixels: + raise McapMediaError() + + +def _decode_image(message: object, schema_name: str, limits: VideoLimits) -> av.VideoFrame: + payload = _frame_bytes(message) + if schema_name in COMPRESSED_IMAGE_SCHEMAS: + encoding = _field(message, "format") + if not isinstance(encoding, str): + raise McapMediaError() + if "jpeg" in encoding.lower() or "jpg" in encoding.lower(): + codec_name = "mjpeg" + elif "png" in encoding.lower(): + codec_name = "png" + else: + raise McapMediaError() + decoder = av.CodecContext.create(codec_name, "r") + decoder.thread_count = 1 + frames = decoder.decode(av.Packet(payload)) + if len(frames) != 1 or not isinstance(frames[0], av.VideoFrame): + raise McapMediaError() + frame = frames[0] + _require_dimensions(frame.width, frame.height, limits) + return frame.reformat(format="rgb24") + if schema_name not in RAW_IMAGE_SCHEMAS: + raise McapMediaError() + encoding = _field(message, "encoding") + width, height, step = (_field(message, name) for name in ("width", "height", "step")) + if ( + not isinstance(encoding, str) + or encoding not in RAW_PIXEL_FORMATS + or type(width) is not int + or type(height) is not int + or type(step) is not int + ): + raise McapMediaError() + _require_dimensions(width, height, limits) + pixel_format, channels = RAW_PIXEL_FORMATS[encoding] + if step < width * channels or len(payload) != step * height: + raise McapMediaError() + pixels = np.frombuffer(payload, dtype=np.uint8).reshape(height, step)[:, : width * channels] + shape = (height, width) if channels == 1 else (height, width, channels) + return av.VideoFrame.from_ndarray( + np.ascontiguousarray(pixels.reshape(shape)), format=pixel_format + ).reformat(format="rgb24") + + +def _h264_packet( + message: object, limits: VideoLimits, parameter_sets: dict[int, bytes] +) -> tuple[av.Packet, tuple[int, int] | None]: + if _field(message, "format") != "h264": + raise McapMediaError() + payload = ensure_access_unit_delimiter(_frame_bytes(message)) + coding = scan_picture_coding_types(payload) + units = split_annex_b_stream(payload) + if coding.picture_count != 1 or coding.b_picture_count or len(units) != 1: + raise McapMediaError() + # Codec configuration persists across access units. Later IDRs may omit + # SPS/PPS, and an SPS-only update must still pass our dimension limits. + parameter_sets_changed = False + nal_offsets = _annex_b_nal_offsets_and_types(payload) + for index, (start_offset, nal_type) in enumerate(nal_offsets): + if nal_type in (7, 8): + end_offset = nal_offsets[index + 1][0] if index + 1 < len(nal_offsets) else len(payload) + parameter_sets[nal_type] = payload[start_offset:end_offset] + parameter_sets_changed = True + if set(parameter_sets) != {7, 8}: + raise McapMediaError() + dimensions = None + if parameter_sets_changed: + parser = av.CodecContext.create("h264", "r") + parser.parse(parameter_sets[7] + parameter_sets[8] + payload) + parser.parse(None) + _require_dimensions(parser.width, parser.height, limits) + dimensions = (parser.width, parser.height) + packet = av.Packet(payload) + packet.is_keyframe = units[0].is_keyframe + return packet, dimensions + + +def _write_camera_video( + output: OutputContainer, + intervals: Iterator[CameraFrameInterval], + channel: TopicInfo, + limits: VideoLimits, +) -> int: + stream: Stream | None = None + parameter_sets: dict[int, bytes] = {} + source_dimensions: tuple[int, int] | None = None + duration_ns = 0 + timestamp_ticks = 0 + for interval in intervals: + duration_ns = interval.timestamp_ns + interval.duration_ns + # Round each absolute endpoint once, then share it with the next PTS. + # Accumulating rounded intervals would drift on irregular camera timing. + end_timestamp_ticks = ( + duration_ns + NANOSECONDS_PER_VIDEO_TICK // 2 + ) // NANOSECONDS_PER_VIDEO_TICK + duration_ticks = end_timestamp_ticks - timestamp_ticks + if not 0 < duration_ticks <= MAXIMUM_VIDEO_INTERVAL_TICKS: + raise McapMediaError() + if channel.schema_name in PASSTHROUGH_VIDEO_SCHEMA_NAMES: + packet, dimensions = _h264_packet(interval.message, limits, parameter_sets) + if stream is None: + if not packet.is_keyframe or dimensions is None: + raise McapMediaError() + source_dimensions = dimensions + stream = output.add_mux_stream( + "h264", + width=dimensions[0], + height=dimensions[1], + time_base=VIDEO_TIME_BASE, + ) + elif dimensions is not None and dimensions != source_dimensions: + raise McapMediaError() + packet.pts = packet.dts = timestamp_ticks + packet.time_base = VIDEO_TIME_BASE + packet.stream = stream + else: + frame = _decode_image(interval.message, channel.schema_name, limits) + if stream is None: + stream = output.add_stream( + "libx264rgb", + rate=30, + width=frame.width, + height=frame.height, + pix_fmt="rgb24", + time_base=VIDEO_TIME_BASE, + options={ + "crf": "0", + "preset": "ultrafast", + "tune": "zerolatency", + "x264-params": "bframes=0:repeat-headers=1:aud=1", + }, + ) + stream.codec_context.thread_count = 1 + if not isinstance(stream, VideoStream) or (frame.width, frame.height) != ( + stream.width, + stream.height, + ): + raise McapMediaError() + frame.pts, frame.time_base = timestamp_ticks, VIDEO_TIME_BASE + packets = stream.encode(frame) + if len(packets) != 1: + raise McapMediaError() + packet = packets[0] + packet.duration = duration_ticks + output.mux(packet) + timestamp_ticks = end_timestamp_ticks + if ( + isinstance(stream, VideoStream) + and isinstance(stream.codec_context, VideoCodecContext) + and stream.encode(None) + ): + raise McapMediaError() + return (duration_ns + 999_999) // 1_000_000 + + +def _export_camera( + source_path: Path, output_path: Path, camera_topic: str | None, limits: VideoLimits +) -> PreparedMcapVideo: + with Episode(source_path) as episode: + channel = _select_camera(episode, camera_topic) + messages = _camera_messages(episode, channel, limits) + first_message = next(messages) + with av.open( + output_path, + "w", + format="mp4", + options={"video_track_timescale": str(VIDEO_TICKS_PER_SECOND)}, + ) as output: + duration_millis = _write_camera_video( + output, + _frame_intervals(chain((first_message,), messages), limits), + channel, + limits, + ) + return PreparedMcapVideo(output_path, duration_millis, channel.topic, first_message[0]) + + +def export_mcap_camera( + source_path: Path, + output_path: Path, + *, + camera_topic: str | None = None, + limits: VideoLimits = VideoLimits(), +) -> PreparedMcapVideo: + """Export a selected camera without modifying source or existing destination. + + PTS are MCAP log times relative to the first selected frame, rounded to the + nearest microsecond. The last frame repeats the preceding positive interval. + Requires at least two frames with distinct rounded times and gaps no larger + than 2,147.483647 seconds. H.264 without B-frames is remuxed with lossless AUD + repair. JPEG/PNG and supported 8-bit raw images are encoded as lossless RGB + H.264. Metadata includes the selected topic and original first log timestamp. + """ + source_path = source_path.resolve(strict=True) + if not source_path.is_file() or source_path.stat().st_size == 0: + raise McapMediaError("MCAP source must be a nonempty local file") + if output_path.exists() or output_path.is_symlink(): + raise FileExistsError(output_path) + output_path.parent.mkdir(parents=True, exist_ok=True) + with tempfile.TemporaryDirectory(dir=output_path.parent, prefix=".camera-export-") as directory: + staged_path = Path(directory) / "video.mp4" + try: + result = _export_camera(source_path, staged_path, camera_topic, limits) + except (StopIteration, RuntimeError, ValueError) as error: + raise McapMediaError( + "MCAP camera cannot be exported with the requested profile" + ) from error + os.link(staged_path, output_path) + return replace(result, video_path=output_path) diff --git a/src/hflow/source_sampling.py b/src/hflow/source_sampling.py index 833b9ee9..688299f2 100644 --- a/src/hflow/source_sampling.py +++ b/src/hflow/source_sampling.py @@ -2,6 +2,7 @@ from __future__ import annotations +import json import math import re import shutil @@ -29,6 +30,12 @@ class SourceSamplingMode(StrEnum): UNIFORM = "uniform" KEYFRAMES = "keyframes" KEYFRAMES_FIRST = "keyframes_first" + NEAREST_KEYFRAMES = "nearest_keyframes" + + +class SourceFrameResize(StrEnum): + PAD = "pad" + FIT = "fit" class KeyframeFallbackReason(StrEnum): @@ -47,7 +54,9 @@ class SourceFrameSampling: keyframes were selected or their span is less than half the window. All sizes are output bounds; callers own source download/decoder memory - limits. The canvas preserves aspect ratio with black padding. Increasing + limits. The default canvas preserves aspect ratio with black padding; FIT omits + padding. NEAREST_KEYFRAMES selects unique keyframes closest to the configured + relative positions, breaking ties toward the earlier frame without fallback. Increasing ``maximum_window_millis`` explicitly permits longer source excerpts. """ @@ -60,6 +69,11 @@ class SourceFrameSampling: timeout_seconds: float = 120.0 maximum_frame_bytes: int = 2 * 1024 * 1024 maximum_log_bytes: int = 2 * 1024 * 1024 + maximum_probe_bytes: int = 8 * 1024 * 1024 + keyframe_positions: tuple[float, ...] = () + resize: SourceFrameResize = SourceFrameResize.PAD + scaling_algorithm: Literal["lanczos", "bicubic"] = "lanczos" + jpeg_quality: int = 5 def __post_init__(self) -> None: if not isinstance(self.mode, SourceSamplingMode): @@ -72,9 +86,29 @@ def __post_init__(self) -> None: "height", "maximum_frame_bytes", "maximum_log_bytes", + "maximum_probe_bytes", ): require_positive_int(getattr(self, field_name), field_name) require_positive_float(self.timeout_seconds, "timeout_seconds") + if not isinstance(self.resize, SourceFrameResize): + raise ValueError("resize must be a SourceFrameResize") + if self.scaling_algorithm not in ("lanczos", "bicubic"): + raise ValueError("unsupported scaling_algorithm") + if type(self.jpeg_quality) is not int or not 1 <= self.jpeg_quality <= 31: + raise ValueError("jpeg_quality must be an integer between 1 and 31") + if not isinstance(self.keyframe_positions, tuple) or any( + isinstance(position, bool) + or not isinstance(position, int | float) + or not math.isfinite(position) + or not 0 <= position <= 1 + for position in self.keyframe_positions + ): + raise ValueError("keyframe_positions must be finite relative positions between 0 and 1") + if self.mode is SourceSamplingMode.NEAREST_KEYFRAMES: + if not 0 < len(self.keyframe_positions) <= self.maximum_frames: + raise ValueError("nearest keyframes require bounded keyframe_positions") + elif self.keyframe_positions: + raise ValueError("keyframe_positions require nearest-keyframe sampling") if self.width % 2 or self.height % 2: raise ValueError("width and height must be even for JPEG chroma subsampling") @@ -97,7 +131,11 @@ class SourceFrameSamples: window: SourceWindow frames: tuple[SampledSourceFrame, ...] requested_mode: SourceSamplingMode - actual_mode: Literal[SourceSamplingMode.UNIFORM, SourceSamplingMode.KEYFRAMES] + actual_mode: Literal[ + SourceSamplingMode.UNIFORM, + SourceSamplingMode.KEYFRAMES, + SourceSamplingMode.NEAREST_KEYFRAMES, + ] fallback_reason: KeyframeFallbackReason | None @@ -149,7 +187,14 @@ def _source_time_base(source_path: Path, executable: Path, *, deadline: float) - maximum_log_bytes=65_536, ) try: - time_base = Fraction(completed.stdout.decode("ascii").strip()) + # MPEG-TS may repeat the selected stream inside a program section. + # Repeated identical values still describe one unambiguous time base. + time_bases = { + Fraction(value) for value in completed.stdout.decode("ascii").splitlines() if value + } + if len(time_bases) != 1: + raise ValueError("ambiguous time base") + time_base = time_bases.pop() except (ValueError, ZeroDivisionError): raise SourceSamplingError("source video has no valid presentation time base") from None if time_base <= 0: @@ -157,16 +202,110 @@ def _source_time_base(source_path: Path, executable: Path, *, deadline: float) - return time_base +def _nearest_keyframe_times( + source_path: Path, + window: SourceWindow, + settings: SourceFrameSampling, + *, + executable: Path, + deadline: float, + time_base: Fraction, +) -> tuple[Fraction, ...]: + origin_result = _run_sampling_command( + [ + str(executable), + "-v", + "error", + "-protocol_whitelist", + "file", + "-show_entries", + "format=start_time", + "-of", + "default=noprint_wrappers=1:nokey=1", + str(source_path), + ], + deadline=deadline, + maximum_log_bytes=65_536, + ) + try: + origin = Fraction(origin_result.stdout.decode("ascii").strip()) + except (ValueError, ZeroDivisionError) as error: + raise SourceSamplingError("source video has no valid playback origin") from error + start = Fraction(window.start_millis, 1000) + end = Fraction(window.end_millis, 1000) + packet_result = _run_sampling_command( + [ + str(executable), + "-v", + "error", + "-protocol_whitelist", + "file", + "-select_streams", + "v:0", + "-read_intervals", + f"%{float(origin + end):.9f}", + "-show_packets", + "-show_entries", + "packet=pts,flags", + "-of", + "json", + str(source_path), + ], + deadline=deadline, + maximum_log_bytes=settings.maximum_probe_bytes, + ) + try: + document = json.loads(packet_result.stdout) + if not isinstance(document, dict) or not isinstance(document.get("packets"), list): + raise ValueError("invalid packet document") + available_times: set[Fraction] = set() + for packet in document["packets"]: + if not isinstance(packet, dict) or not isinstance(packet.get("flags"), str): + raise ValueError("invalid packet") + if "K" not in packet["flags"]: + continue + if type(packet.get("pts")) is not int: + raise ValueError("keyframe has no presentation timestamp") + timestamp = packet["pts"] * time_base - origin + if start <= timestamp < end: + if (timestamp / time_base).denominator != 1: + raise ValueError("playback origin is not aligned to the source time base") + available_times.add(timestamp) + except (ValueError, TypeError, KeyError) as error: + raise SourceSamplingError("source keyframes have invalid presentation times") from error + if not available_times: + return () + return tuple( + sorted( + { + min( + available_times, + key=lambda timestamp: ( + abs(timestamp - (start + (end - start) * Fraction(str(position)))), + timestamp, + ), + ) + for position in settings.keyframe_positions + } + ) + ) + + def _extract_frames( source_path: Path, output_directory: Path, window: SourceWindow, settings: SourceFrameSampling, - mode: Literal[SourceSamplingMode.UNIFORM, SourceSamplingMode.KEYFRAMES], + mode: Literal[ + SourceSamplingMode.UNIFORM, + SourceSamplingMode.KEYFRAMES, + SourceSamplingMode.NEAREST_KEYFRAMES, + ], *, executable: Path, deadline: float, time_base: Fraction, + probe_executable: Path, ) -> tuple[SampledSourceFrame, ...]: duration_seconds = window.duration_millis / 1000 start_seconds = window.start_millis / 1000 @@ -205,14 +344,27 @@ def _extract_frames( "(isnan(prev_selected_t)+" f"gt({current_bin}\\,{previous_bin}))" ) - video_filter = ",".join( - ( - selection_filter, - f"scale={settings.width}:{settings.height}:force_original_aspect_ratio=decrease:flags=lanczos", - f"pad={settings.width}:{settings.height}:(ow-iw)/2:(oh-ih)/2:black", - "showinfo", + selected_times: tuple[Fraction, ...] = () + if mode is SourceSamplingMode.NEAREST_KEYFRAMES: + selected_times = _nearest_keyframe_times( + source_path, + window, + settings, + executable=probe_executable, + deadline=deadline, + time_base=time_base, ) - ) + if not selected_times: + return () + selection_filter = "select=" + "+".join( + f"eq(pts\\,{int(timestamp / time_base)})" for timestamp in selected_times + ) + resize_filters = [ + f"scale={settings.width}:{settings.height}:force_original_aspect_ratio=decrease:flags={settings.scaling_algorithm}", + ] + if settings.resize is SourceFrameResize.PAD: + resize_filters.append(f"pad={settings.width}:{settings.height}:(ow-iw)/2:(oh-ih)/2:black") + video_filter = ",".join((selection_filter, *resize_filters, "showinfo")) arguments = [ str(executable), "-hide_banner", @@ -226,7 +378,7 @@ def _extract_frames( "-threads", "2", ] - if mode is SourceSamplingMode.KEYFRAMES: + if mode in (SourceSamplingMode.KEYFRAMES, SourceSamplingMode.NEAREST_KEYFRAMES): arguments.extend(("-skip_frame", "nokey")) # Preserve the original playback PTS through a seek. Adding a seek offset # back to rebased PTS would introduce rounding at non-tick-aligned starts. @@ -260,7 +412,7 @@ def _extract_frames( "-c:v", "mjpeg", "-q:v", - "5", + str(settings.jpeg_quality), "-threads", "2", "-f", @@ -284,6 +436,8 @@ def _extract_frames( int(match.group(1)) * Fraction(numerator, denominator) for match in _TIMESTAMP_PATTERN.finditer(diagnostics) ) + if mode is SourceSamplingMode.NEAREST_KEYFRAMES and timestamps != selected_times: + raise SourceSamplingError("extraction did not reproduce the selected keyframes") frame_paths = tuple(sorted(output_directory.glob("frame_*.jpg"))) if len(frame_paths) != len(timestamps) or len(frame_paths) > settings.maximum_frames: raise SourceSamplingError("source frame extraction produced inconsistent samples") @@ -336,9 +490,15 @@ def sample_source_frames( probe_executable = ffprobe_path() deadline = time.monotonic() + settings.timeout_seconds output_directory.mkdir(parents=True, exist_ok=False) - actual_mode: Literal[SourceSamplingMode.UNIFORM, SourceSamplingMode.KEYFRAMES] = ( + actual_mode: Literal[ + SourceSamplingMode.UNIFORM, + SourceSamplingMode.KEYFRAMES, + SourceSamplingMode.NEAREST_KEYFRAMES, + ] = ( SourceSamplingMode.UNIFORM if settings.mode is SourceSamplingMode.UNIFORM + else SourceSamplingMode.NEAREST_KEYFRAMES + if settings.mode is SourceSamplingMode.NEAREST_KEYFRAMES else SourceSamplingMode.KEYFRAMES ) fallback_reason = None @@ -357,6 +517,7 @@ def sample_source_frames( executable=executable, deadline=deadline, time_base=time_base, + probe_executable=probe_executable, ) if settings.mode is SourceSamplingMode.KEYFRAMES_FIRST: if len(frames) < 2: @@ -378,6 +539,7 @@ def sample_source_frames( executable=executable, deadline=deadline, time_base=time_base, + probe_executable=probe_executable, ) published_frames = tuple( SampledSourceFrame( diff --git a/src/hflow/video_statistics.py b/src/hflow/video_statistics.py new file mode 100644 index 00000000..eaae36d8 --- /dev/null +++ b/src/hflow/video_statistics.py @@ -0,0 +1,60 @@ +"""Public file-level frame statistics with explicit range and cache policy. + +Measurements use original framing: callers should not add model-input padding. +Defaults retain the existing instrument's range handling; FULL explicitly converts +declared input range to full-range luma before measurement. Neither mode changes +the source file. +""" + +from pathlib import Path + +from hflow._video_measurement_toolchain import resolved_video_measurement_toolchain +from hflow._video_measurements._frame_statistics import ( + FRAME_STATISTICS_DEFINITION_VERSION, + FrameStatisticsExecutionError, + FrameStatisticsParseError, + FrameStatisticsProvenance, + FrameStatisticsSettings, + LumaRangeEvidence, + LumaRangePolicy, + VideoFrameStatistics, + VideoTimeInterval, +) +from hflow._video_measurements._frame_statistics import ( + measure_video_frame_statistics as _measure_video_frame_statistics, +) +from hflow._video_measurements._toolchain import VideoMeasurementToolchain + +__all__ = [ + "FRAME_STATISTICS_DEFINITION_VERSION", + "FrameStatisticsExecutionError", + "FrameStatisticsParseError", + "FrameStatisticsProvenance", + "FrameStatisticsSettings", + "LumaRangeEvidence", + "LumaRangePolicy", + "VideoFrameStatistics", + "VideoMeasurementToolchain", + "VideoTimeInterval", + "measure_video_frame_statistics", +] + + +def measure_video_frame_statistics( + video: Path, + *, + settings: FrameStatisticsSettings = FrameStatisticsSettings(), + toolchain: VideoMeasurementToolchain | None = None, + instrument_output_cache_path: Path | None = None, +) -> VideoFrameStatistics: + """Measure without persistent output unless a cache path is explicitly supplied. + + Toolchain resolution follows HFlow's binary policy and may download binaries. + The result records effective settings, filter graph and binary version. + """ + return _measure_video_frame_statistics( + video, + settings=settings, + toolchain=toolchain if toolchain is not None else resolved_video_measurement_toolchain(), + instrument_output_cache_path=instrument_output_cache_path, + ) diff --git a/tests/test_blur.py b/tests/test_blur.py new file mode 100644 index 00000000..7de54bc8 --- /dev/null +++ b/tests/test_blur.py @@ -0,0 +1,74 @@ +"""Raw-score and missing-evidence outcomes for the CPU blur adapter.""" + +from __future__ import annotations + +import math +from pathlib import Path + +import cv2 +import numpy as np +import pytest +from numpy.typing import NDArray + +from hflow.blur import measure_video_blur, summarize_blur_scores + + +def test_summary_preserves_raw_scores_and_excludes_unavailable_frames() -> None: + summary = summarize_blur_scores(iter((120.0, math.nan, 240.0, math.inf, -math.inf))) + + assert summary.frame_count == 5 + assert summary.scored_frame_count == 2 + assert summary.mean_blur_score == pytest.approx(180.0) + + +@pytest.mark.parametrize("scores", [(), (math.nan, math.inf, -math.inf)]) +def test_missing_blur_evidence_has_no_score(scores: tuple[float, ...]) -> None: + summary = summarize_blur_scores(iter(scores)) + + assert summary.frame_count == len(scores) + assert summary.scored_frame_count == 0 + assert summary.mean_blur_score is None + + +def write_static_video(video_path: Path, frame: NDArray[np.uint8]) -> None: + frame_height, frame_width = frame.shape + writer = cv2.VideoWriter( + str(video_path), + cv2.VideoWriter.fourcc(*"FFV1"), + 5.0, + (frame_width, frame_height), + isColor=False, + ) + if not writer.isOpened(): + raise RuntimeError("Could not generate blur test video") + try: + for _ in range(3): + writer.write(frame) + finally: + writer.release() + + +def test_video_adapter_distinguishes_defocus_from_missing_detail(tmp_path: Path) -> None: + row_indices, column_indices = np.indices((128, 192)) + sharp_frame = (((row_indices // 16 + column_indices // 16) % 2) * 255).astype(np.uint8) + defocused_frame = cv2.GaussianBlur(sharp_frame, (0, 0), sigmaX=3.0).astype(np.uint8) + featureless_frame = np.full_like(sharp_frame, 128) + sharp_path = tmp_path / "sharp.mkv" + defocused_path = tmp_path / "defocused.mkv" + featureless_path = tmp_path / "featureless.mkv" + write_static_video(sharp_path, sharp_frame) + write_static_video(defocused_path, defocused_frame) + write_static_video(featureless_path, featureless_frame) + + sharp_summary = measure_video_blur(sharp_path) + defocused_summary = measure_video_blur(defocused_path) + featureless_summary = measure_video_blur(featureless_path) + + assert sharp_summary.frame_count == sharp_summary.scored_frame_count == 3 + assert defocused_summary.frame_count == defocused_summary.scored_frame_count == 3 + assert sharp_summary.mean_blur_score is not None + assert defocused_summary.mean_blur_score is not None + assert defocused_summary.mean_blur_score > sharp_summary.mean_blur_score + assert featureless_summary.frame_count == 3 + assert featureless_summary.scored_frame_count == 0 + assert featureless_summary.mean_blur_score is None diff --git a/tests/test_camera_shake_summary.py b/tests/test_camera_shake_summary.py new file mode 100644 index 00000000..ac534fe2 --- /dev/null +++ b/tests/test_camera_shake_summary.py @@ -0,0 +1,117 @@ +import math + +import pytest + +from hflow.camera_motion import ( + AngularMotionRate, + CameraMotionObservation, + CameraMotionStreamSettings, + CameraMotionTransform, + CameraShakeObservation, + CameraShakeSettings, + MeasuredCameraMotion, + MeasuredCameraShake, + MotionFitEvidence, + UnavailableCameraShake, + UnmeasuredCameraMotion, + summarize_camera_shake, +) + + +def shake_observation( + pair_index: int, + *, + duration_seconds: float, + shake: MeasuredCameraShake | UnavailableCameraShake, + motion_measured: bool = True, +) -> CameraShakeObservation: + motion_settings = CameraMotionStreamSettings(frames_per_second=1 / duration_seconds) + motion = ( + MeasuredCameraMotion( + CameraMotionTransform(0.0, 0.0, 0.0, 1.0), + MotionFitEvidence(20, 20, 20, 0.0), + ) + if motion_measured + else UnmeasuredCameraMotion("insufficient_tracks", MotionFitEvidence(20, 0, 0, None)) + ) + return CameraShakeObservation( + motion=CameraMotionObservation( + pair_index=pair_index, + start_seconds=pair_index / motion_settings.frames_per_second, + end_seconds=(pair_index + 1) / motion_settings.frames_per_second, + frame_width_pixels=320, + measurement=motion, + settings=motion_settings, + ), + shake=shake, + settings=CameraShakeSettings(), + ) + + +def measured_shake(rate: float) -> MeasuredCameraShake: + return MeasuredCameraShake( + residual=AngularMotionRate(0.0, rate, 0.0), + smoothed_motion=AngularMotionRate(0.0, 0.0, 0.0), + ) + + +def test_summary_weights_shake_and_coverage_by_observed_duration() -> None: + # The reducer can consume multiple clips with different fixed frame cadences. + observations = ( + shake_observation(0, duration_seconds=0.25, shake=measured_shake(2.0)), + shake_observation( + 1, duration_seconds=0.25, shake=UnavailableCameraShake("insufficient_context") + ), + shake_observation(0, duration_seconds=0.5, shake=measured_shake(4.0)), + shake_observation( + 1, + duration_seconds=0.5, + shake=UnavailableCameraShake("unmeasured_context"), + motion_measured=False, + ), + ) + + summary = summarize_camera_shake(iter(observations)) + + assert summary.pair_count == 4 + assert summary.measured_motion_pair_count == 3 + assert summary.measured_shake_pair_count == 2 + assert summary.insufficient_context_pair_count == 1 + assert summary.unmeasured_context_pair_count == 1 + assert summary.observed_seconds == pytest.approx(1.5) + assert summary.assessed_seconds == pytest.approx(0.75) + assert summary.unassessed_seconds == pytest.approx(0.75) + assert summary.assessed_fraction == pytest.approx(0.5) + assert summary.mean_shake_degrees_per_second == pytest.approx(10 / 3) + assert summary.rms_shake_degrees_per_second == pytest.approx(math.sqrt(12)) + assert summary.maximum_shake_degrees_per_second == pytest.approx(4.0) + + +@pytest.mark.parametrize("empty", [True, False]) +def test_no_assessed_motion_is_missing_instead_of_zero(empty: bool) -> None: + observations = ( + () + if empty + else ( + shake_observation( + 0, duration_seconds=0.5, shake=UnavailableCameraShake("insufficient_context") + ), + shake_observation( + 1, + duration_seconds=0.5, + shake=UnavailableCameraShake("unmeasured_context"), + motion_measured=False, + ), + ) + ) + + summary = summarize_camera_shake(iter(observations)) + + assert summary.pair_count == (0 if empty else 2) + assert summary.observed_seconds == (0 if empty else 1) + assert summary.assessed_seconds == 0 + assert summary.unassessed_seconds == summary.observed_seconds + assert summary.assessed_fraction == (None if empty else 0) + assert summary.mean_shake_degrees_per_second is None + assert summary.rms_shake_degrees_per_second is None + assert summary.maximum_shake_degrees_per_second is None diff --git a/tests/test_mcap_video.py b/tests/test_mcap_video.py new file mode 100644 index 00000000..1c9fb294 --- /dev/null +++ b/tests/test_mcap_video.py @@ -0,0 +1,246 @@ +from __future__ import annotations + +import hashlib +import json +import re +from dataclasses import replace +from fractions import Fraction +from pathlib import Path + +import av +import numpy as np +import pytest +from foxglove_schemas_protobuf.CompressedVideo_pb2 import CompressedVideo +from mcap.writer import Writer +from mcap_protobuf.schema import build_file_descriptor_set + +from hflow.mcap_video import McapMediaError, export_mcap_camera +from hflow.media import VideoLimits + +CAMERA_START_NS = 1_700_000_000_000_000_000 +# Real MCAP capture can deliver adjacent frames in bursts without exceeding +# the supported average frame rate. +FRAME_TIMES_NS = (0, 147_079, 5_000_147_079, 5_500_147_079) +SOURCE_DURATION_MILLIS = 6001 +LIMITS = VideoLimits(maximum_frame_pixels=256, maximum_frames_per_second=60, timeout_seconds=20) + + +def _h264_packets(width: int = 16) -> list[bytes]: + encoder = av.CodecContext.create("libx264", "w") + encoder.width = encoder.height = width + encoder.pix_fmt = "yuv420p" + encoder.time_base = Fraction(1, 30) + encoder.thread_count = 1 + encoder.options = { + "preset": "ultrafast", + "tune": "zerolatency", + "x264-params": "keyint=2:min-keyint=2:scenecut=0:repeat-headers=1:aud=1:bframes=0", + } + packets = [] + for index in range(len(FRAME_TIMES_NS)): + frame = av.VideoFrame.from_ndarray(np.full((width, width, 3), index * 40, dtype=np.uint8)) + frame.pts = index + packets.extend(bytes(packet) for packet in encoder.encode(frame)) + packets.extend(bytes(packet) for packet in encoder.encode(None)) + return packets + + +def _omit_h264_headers(packet: bytes, nal_types: frozenset[int]) -> bytes: + return b"".join( + b"\x00\x00\x00\x01" + nal + for nal in re.split(b"\x00\x00(?:\x00)?\x01", packet) + if nal and nal[0] & 31 not in nal_types + ) + + +def _write_h264_mcap( + path: Path, + topics: tuple[str, ...] = ("/camera/front",), + packets: list[bytes] | None = None, + frame_times_ns: tuple[int, ...] = FRAME_TIMES_NS, +) -> list[bytes]: + packets = _h264_packets() if packets is None else packets + with path.open("wb") as output: + writer = Writer(output) + writer.start() + schema_id = writer.register_schema( + name="foxglove.CompressedVideo", + encoding="protobuf", + data=build_file_descriptor_set(CompressedVideo).SerializeToString(), + ) + for topic in topics: + channel_id = writer.register_channel(topic, "protobuf", schema_id) + for timestamp, packet in zip(frame_times_ns, packets, strict=True): + message = CompressedVideo(format="h264", data=packet) + message.timestamp.FromNanoseconds(CAMERA_START_NS + timestamp) + writer.add_message( + channel_id, + CAMERA_START_NS + timestamp, + message.SerializeToString(), + CAMERA_START_NS + timestamp, + ) + writer.finish() + return packets + + +@pytest.mark.parametrize("repeat_parameter_sets", [True, False]) +def test_mcap_h264_preserves_frames_gaps_timestamps_and_original_bytes( + tmp_path: Path, repeat_parameter_sets: bool +) -> None: + source = tmp_path / "source.mcap" + packets = _h264_packets() + if not repeat_parameter_sets: + packets = [ + packets[0], + *(_omit_h264_headers(packet, frozenset({7, 8})) for packet in packets[1:]), + ] + _write_h264_mcap(source, packets=packets) + original_digest = hashlib.sha256(source.read_bytes()).hexdigest() + result = export_mcap_camera(source, tmp_path / "camera.mp4", limits=LIMITS) + + assert result.camera_topic == "/camera/front" + assert result.start_timestamp_ns == CAMERA_START_NS + assert result.duration_millis == SOURCE_DURATION_MILLIS + assert hashlib.sha256(source.read_bytes()).hexdigest() == original_digest + with av.open(result.video_path) as video: + frames = list(video.decode(video=0)) + for frame, timestamp in zip(frames, FRAME_TIMES_NS, strict=True): + assert frame.pts is not None and frame.time_base is not None + assert frame.pts * frame.time_base == Fraction((timestamp + 500) // 1000, 1_000_000) + decoder = av.CodecContext.create("h264", "r") + expected_frames = [frame for packet in packets for frame in decoder.decode(av.Packet(packet))] + for actual, expected in zip(frames, expected_frames, strict=True): + np.testing.assert_array_equal( + actual.to_ndarray(format="rgb24"), expected.to_ndarray(format="rgb24") + ) + assert not list(tmp_path.glob(".camera-export-*")) + + +def test_mcap_requires_initial_h264_codec_configuration(tmp_path: Path) -> None: + source = tmp_path / "source.mcap" + packets = _h264_packets() + packets[0] = _omit_h264_headers(packets[0], frozenset({7, 8})) + _write_h264_mcap(source, packets=packets) + output = tmp_path / "camera.mp4" + + with pytest.raises(McapMediaError): + export_mcap_camera(source, output, limits=LIMITS) + + assert not output.exists() + + +def test_multicamera_mcap_requires_an_exact_selection(tmp_path: Path) -> None: + source = tmp_path / "source.mcap" + _write_h264_mcap(source, ("/camera/front", "/camera/back")) + output = tmp_path / "camera.mp4" + for camera_topic in (None, "front", "/unknown"): + with pytest.raises(McapMediaError): + export_mcap_camera(source, output, camera_topic=camera_topic, limits=LIMITS) + assert not output.exists() + result = export_mcap_camera(source, output, camera_topic="/camera/back", limits=LIMITS) + assert result.camera_topic == "/camera/back" + assert result.duration_millis == SOURCE_DURATION_MILLIS + + +@pytest.mark.parametrize("repeat_picture_parameter_set", [True, False]) +def test_mcap_rejects_a_camera_resolution_change( + tmp_path: Path, repeat_picture_parameter_set: bool +) -> None: + source = tmp_path / "source.mcap" + changed_packets = _h264_packets(32)[:2] + if not repeat_picture_parameter_set: + changed_packets[0] = _omit_h264_headers(changed_packets[0], frozenset({8})) + _write_h264_mcap(source, packets=_h264_packets()[:2] + changed_packets) + destination = tmp_path / "camera.mp4" + + with pytest.raises(McapMediaError): + export_mcap_camera( + source, destination, limits=replace(LIMITS, maximum_frame_pixels=32 * 32) + ) + + assert not destination.exists() + + +@pytest.mark.parametrize("compressed", [False, True]) +def test_image_mcap_preserves_pixels_and_ignores_unselected_cameras( + tmp_path: Path, compressed: bool +) -> None: + source = tmp_path / "raw.mcap" + pixels = np.arange(16 * 16 * 3, dtype=np.uint8).reshape(16, 16, 3) + payload: dict[str, object] = { + "width": 16, + "height": 16, + "step": 48, + "encoding": "rgb8", + "data": list(pixels.tobytes()), + } + schema_name = "foxglove.RawImage" + if compressed: + encoder = av.CodecContext.create("png", "w") + encoder.width = encoder.height = 16 + encoder.pix_fmt = "rgb24" + encoded = encoder.encode(av.VideoFrame.from_ndarray(pixels))[0] + payload = {"format": "png", "data": list(bytes(encoded))} + schema_name = "foxglove.CompressedImage" + with source.open("wb") as output: + writer = Writer(output) + writer.start() + schema_id = writer.register_schema(schema_name, "jsonschema", b"{}") + selected_channel = writer.register_channel("/rgb", "json", schema_id) + other_channel = writer.register_channel("/unused", "json", schema_id) + for timestamp in FRAME_TIMES_NS: + writer.add_message( + selected_channel, + CAMERA_START_NS + timestamp, + json.dumps(payload).encode(), + CAMERA_START_NS + timestamp, + ) + writer.add_message( + other_channel, CAMERA_START_NS, b'{"data":"unsupported"}', CAMERA_START_NS + ) + writer.finish() + result = export_mcap_camera(source, tmp_path / "raw.mp4", camera_topic="/rgb", limits=LIMITS) + assert result.duration_millis == SOURCE_DURATION_MILLIS + assert result.start_timestamp_ns == CAMERA_START_NS + with av.open(result.video_path) as video: + frames = list(video.decode(video=0)) + assert len(frames) == len(FRAME_TIMES_NS) + for frame, timestamp in zip(frames, FRAME_TIMES_NS, strict=True): + assert frame.pts is not None and frame.time_base is not None + assert frame.pts * frame.time_base == Fraction((timestamp + 500) // 1000, 1_000_000) + np.testing.assert_array_equal(frame.to_ndarray(format="rgb24"), pixels) + + +@pytest.mark.parametrize( + ("limits", "frame_times_ns"), + ( + (VideoLimits(maximum_frame_pixels=1, timeout_seconds=20), FRAME_TIMES_NS), + (VideoLimits(maximum_duration_seconds=2, timeout_seconds=20), FRAME_TIMES_NS), + (VideoLimits(maximum_frames_per_second=0.5, timeout_seconds=20), FRAME_TIMES_NS), + (LIMITS, (0, 100, 1_000_000_000, 2_000_000_000)), + (LIMITS, (0, 2_147_483_648_000, 2_148_483_648_000, 2_149_483_648_000)), + ), +) +def test_rejected_mcap_leaves_original_without_publishing_output( + tmp_path: Path, limits: VideoLimits, frame_times_ns: tuple[int, ...] +) -> None: + source = tmp_path / "source.mcap" + _write_h264_mcap(source, frame_times_ns=frame_times_ns) + original_digest = hashlib.sha256(source.read_bytes()).hexdigest() + destination = tmp_path / "existing.mp4" + with pytest.raises(McapMediaError): + export_mcap_camera(source, destination, limits=limits) + + assert not destination.exists() + assert hashlib.sha256(source.read_bytes()).hexdigest() == original_digest + assert not list(tmp_path.glob(".camera-export-*")) + + +def test_camera_export_preserves_an_existing_destination(tmp_path: Path) -> None: + source = tmp_path / "source.mcap" + _write_h264_mcap(source) + destination = tmp_path / "existing.mp4" + destination.write_bytes(b"existing result") + with pytest.raises(FileExistsError): + export_mcap_camera(source, destination, limits=LIMITS) + assert destination.read_bytes() == b"existing result" diff --git a/tests/test_source_sampling.py b/tests/test_source_sampling.py index c5756d47..4e502b12 100644 --- a/tests/test_source_sampling.py +++ b/tests/test_source_sampling.py @@ -481,3 +481,73 @@ def test_example_emits_complete_windows_with_readable_frame_paths( frames = [frame for record in records for frame in record["frames"]] assert [frame["timestamp_seconds"] for frame in frames] == list(range(6)) assert all(cv2.imread(frame["path"]) is not None for frame in frames) + + +@pytest.mark.parametrize("timestamp_offset", [0, 5, -1]) +def test_nearest_keyframes_preserve_ties_pixels_and_playback_origin( + color_video: Path, tmp_path: Path, timestamp_offset: int +) -> None: + from hflow.source_sampling import SourceFrameResize + + shifted_source = tmp_path / ("shifted.ts" if timestamp_offset < 0 else "shifted.mp4") + subprocess.run( + [ + str(ffmpeg_path()), + "-v", + "error", + "-i", + str(color_video), + "-map", + "0:v:0", + "-c:v", + "copy", + "-output_ts_offset", + str(timestamp_offset), + "-avoid_negative_ts", + "disabled", + str(shifted_source), + ], + check=True, + capture_output=True, + ) + samples = sample_source_frames( + shifted_source, + tmp_path / "nearest", + window=SourceWindow(0, 6000), + settings=SourceFrameSampling( + mode=SourceSamplingMode.NEAREST_KEYFRAMES, + keyframe_positions=(0.15, 0.5, 0.85), + resize=SourceFrameResize.FIT, + width=960, + height=960, + jpeg_quality=2, + scaling_algorithm="bicubic", + ), + ) + assert [frame.timestamp_seconds for frame in samples.frames] == [0, 2, 4] + assert samples.actual_mode is SourceSamplingMode.NEAREST_KEYFRAMES + assert samples.fallback_reason is None + for frame, dominant_channel in zip(samples.frames, (2, 1, 0), strict=True): + pixels = cv2.imread(str(frame.path)) + assert pixels is not None + assert pixels.shape == (640, 960, 3) + assert int(pixels[320, 480, dominant_channel]) > 100 + assert np.all(pixels.max(axis=2) > 30) + + +def test_nearest_keyframes_deduplicate_and_leave_empty_windows_empty( + color_video: Path, tmp_path: Path +) -> None: + settings = SourceFrameSampling( + mode=SourceSamplingMode.NEAREST_KEYFRAMES, + keyframe_positions=(0.2, 0.5, 0.8), + ) + sparse = sample_source_frames( + color_video, tmp_path / "sparse", window=SourceWindow(0, 1900), settings=settings + ) + empty = sample_source_frames( + color_video, tmp_path / "empty-nearest", window=SourceWindow(250, 1250), settings=settings + ) + assert [frame.timestamp_seconds for frame in sparse.frames] == [0] + assert empty.frames == () + assert empty.fallback_reason is None diff --git a/tests/test_video_import.py b/tests/test_video_import.py index 836cec54..93267f2b 100644 --- a/tests/test_video_import.py +++ b/tests/test_video_import.py @@ -433,3 +433,56 @@ def test_start_time_upper_bound_uses_field_guard() -> None: with pytest.raises(ValueError) as exc_info: replace(VideoImportConfig(duration_s=1), start_time_ns=value) assert str(exc_info.value) == (f"start_time_ns must be in [0, {value - 1}], got {value}") + + +@pytest.mark.parametrize("duration_s,image_hz", [(1.0, 4.0), (0.1, 1.0)]) +def test_direct_model_video_matches_canonical_decoded_pixels( + source_video: Path, tmp_path: Path, duration_s: float, image_hz: float +) -> None: + from hflow.importers.video import prepare_model_video + + configuration = VideoImportConfig( + duration_s=duration_s, image_hz=image_hz, image_width=80, image_height=80 + ) + imported = import_video_episode(source_video, tmp_path / "import.mcap", configuration) + application = hflow.App("model-parity", data_root=tmp_path / "workspace", default_checks=()) + report = asyncio.run(application.process(imported, record=False, stages={hflow.Stage.SYNC})) + assert not report.has_errors, report.summary() + output = tmp_path / "direct.mp4" + assert prepare_model_video(source_video, output, configuration) == output + from hflow.video import write_access_units_to_mp4 + + with report.canonical_path.open("rb") as canonical_stream: + canonical_messages = list( + make_reader( + canonical_stream, decoder_factories=[DecoderFactory()] + ).iter_decoded_messages() + ) + reference_video = write_access_units_to_mp4( + (decoded.data for _schema, _channel, _message, decoded in canonical_messages), + fps=image_hz if configuration.frame_count > 1 else 1.0, + output=tmp_path / "reference.mp4", + ) + fingerprints = [ + subprocess.check_output( + [ + str(ffmpeg_path()), + "-v", + "error", + "-i", + str(video), + "-map", + "0:v:0", + "-f", + "framemd5", + "-", + ], + timeout=30, + ) + for video in (reference_video, output) + ] + assert fingerprints[0] == fingerprints[1] + original_output = output.read_bytes() + with pytest.raises(FileExistsError): + prepare_model_video(source_video, output, configuration) + assert output.read_bytes() == original_output diff --git a/tests/test_video_statistics.py b/tests/test_video_statistics.py new file mode 100644 index 00000000..a367b6a2 --- /dev/null +++ b/tests/test_video_statistics.py @@ -0,0 +1,87 @@ +"""Public full-range measurements match explicit lossless normalization.""" + +import hashlib +import subprocess +from pathlib import Path + +import pytest + +from hflow.ffmpeg import ffmpeg_path +from hflow.video_statistics import ( + FrameStatisticsSettings, + LumaRangePolicy, + measure_video_frame_statistics, +) + + +@pytest.mark.parametrize("color", ["black", "gray", "white"]) +def test_full_range_statistics_match_explicit_normalization_without_persistent_outputs( + tmp_path: Path, color: str +) -> None: + source = tmp_path / "source.mp4" + normalized = tmp_path / "normalized.mp4" + subprocess.run( + [ + str(ffmpeg_path()), + "-v", + "error", + "-f", + "lavfi", + "-i", + f"color={color}:s=96x64:r=8:d=3", + "-c:v", + "libx264", + "-threads", + "1", + "-pix_fmt", + "yuv420p", + str(source), + ], + check=True, + capture_output=True, + ) + source_digest = hashlib.sha256(source.read_bytes()).digest() + subprocess.run( + [ + str(ffmpeg_path()), + "-v", + "error", + "-i", + str(source), + "-vf", + "scale=in_range=auto:out_range=full", + "-c:v", + "libx264", + "-crf", + "0", + "-threads", + "1", + "-pix_fmt", + "yuv420p", + "-color_range", + "pc", + str(normalized), + ], + check=True, + capture_output=True, + ) + reference = measure_video_frame_statistics( + normalized, settings=FrameStatisticsSettings(luma_range=LumaRangePolicy.FULL) + ) + observed = measure_video_frame_statistics( + source, settings=FrameStatisticsSettings(luma_range=LumaRangePolicy.FULL) + ) + for field_name in ( + "black_frame_percent", + "clipped_highlight_frame_percent", + "crushed_shadow_frame_percent", + "freeze_total_seconds", + "average_luma_mean", + ): + assert getattr(observed, field_name) == pytest.approx(getattr(reference, field_name)) + assert observed.average_luma_mean == pytest.approx( + {"black": 0, "gray": 128, "white": 255}[color] + ) + assert observed.provenance.settings.luma_range is LumaRangePolicy.FULL + assert hashlib.sha256(source.read_bytes()).digest() == source_digest + assert set(tmp_path.iterdir()) == {source, normalized} diff --git a/uv.lock b/uv.lock index 42a89d2e..eec02d09 100644 --- a/uv.lock +++ b/uv.lock @@ -273,6 +273,32 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/64/b4/17d4b0b2a2dc85a6df63d1157e028ed19f90d4cd97c36717afef2bc2f395/attrs-26.1.0-py3-none-any.whl", hash = "sha256:c647aa4a12dfbad9333ca4e71fe62ddc36f4e63b2d260a37a8b83d2f043ac309", size = 67548, upload-time = "2026-03-19T14:22:23.645Z" }, ] +[[package]] +name = "av" +version = "18.1.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/8d/f4/f22114d30d3435e38c6af2b4870f37b864403dca6ae7af747a289ce0a18e/av-18.1.0.tar.gz", hash = "sha256:47bfc286e1bc9de7ab4681fc2b575cd2460a66919d31ffe1bd5aa54fae531a28", size = 4451061, upload-time = "2026-08-12T22:28:18.761Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/05/d4/d7cdc8bff143c17a6d35924375ae28dd692cacde38700a7d419fde54f44a/av-18.1.0-cp311-abi3-macosx_11_0_x86_64.whl", hash = "sha256:ae75d8bb6467895ed1f8572ededf7ffa49eac07f6e483222f5d7d62a41d12f04", size = 22546147, upload-time = "2026-08-12T22:27:11.851Z" }, + { url = "https://files.pythonhosted.org/packages/3f/c9/37a619297492256b77d5ed906e7d8166c10a26ed251dccf1ae03ab19bff6/av-18.1.0-cp311-abi3-macosx_14_0_arm64.whl", hash = "sha256:b30a4e8d934558e19602b68998a4d9ac9f250fa0dacef216f7e8e40153b13316", size = 18217603, upload-time = "2026-08-12T22:27:14.713Z" }, + { url = "https://files.pythonhosted.org/packages/d9/84/2464ffb64c08c5ce8b522c8e74594714414e3b0575267652c5c51c0574b9/av-18.1.0-cp311-abi3-manylinux_2_28_aarch64.whl", hash = "sha256:6fc837cc51adf80331ac850779cd53b5d4c4460b0ebe9057a02a921c6736f19d", size = 33640142, upload-time = "2026-08-12T22:27:17.835Z" }, + { url = "https://files.pythonhosted.org/packages/27/3a/204dbfc3e08eb4cdc6e6ff57be02150bc44523ebdb50182d10025792ebd9/av-18.1.0-cp311-abi3-manylinux_2_28_x86_64.whl", hash = "sha256:8a032e8d8ebc73dec079364b9b4a6837638a2d106e8472314e685ffbf163e700", size = 35786210, upload-time = "2026-08-12T22:27:20.984Z" }, + { url = "https://files.pythonhosted.org/packages/e1/99/b0d04ec553ff9a7e00455458dfa3a39c8a8f627b273056b4e5fe57d590de/av-18.1.0-cp311-abi3-manylinux_2_31_armv7l.whl", hash = "sha256:3c8b1f8b46f99d52e2d8b0ed5d0cdadf172d24794d46e2077b16e44ed08e26ff", size = 39379798, upload-time = "2026-08-12T22:27:24.432Z" }, + { url = "https://files.pythonhosted.org/packages/56/b1/e00d4feae59160149df6126585e726fdc6300798fd40c5dd324879e81f68/av-18.1.0-cp311-abi3-musllinux_1_2_aarch64.whl", hash = "sha256:ab5ac081bc9eaf54109120d4e56284674fecfbe520d9aa1707c7fa911ec5f4d2", size = 34690321, upload-time = "2026-08-12T22:27:27.769Z" }, + { url = "https://files.pythonhosted.org/packages/dc/94/836fa987e3084d11a21489f11357fb24843ef3aa8faf74ddddfc603d5062/av-18.1.0-cp311-abi3-musllinux_1_2_x86_64.whl", hash = "sha256:191224788d87af06c31784a395bb73f14b72f33d7f4871ace0157de2abdc6276", size = 36859932, upload-time = "2026-08-12T22:27:31.403Z" }, + { url = "https://files.pythonhosted.org/packages/33/b4/76ba21e46704f632004276b85289a1582e95f5eff760436d6149875a1881/av-18.1.0-cp311-abi3-win_amd64.whl", hash = "sha256:ea1480b7a8d5405cb5f382b344731bf125fd2c1c6fae3964f6c48595628387ff", size = 27595679, upload-time = "2026-08-12T22:27:35.177Z" }, + { url = "https://files.pythonhosted.org/packages/4f/ad/a3135884c5753b09773176b97201ae602f67ad14206c395ff838d66bf9b0/av-18.1.0-cp311-abi3-win_arm64.whl", hash = "sha256:5509ec12aaa19fd6601de13cfa6f4cdad450da07982118510592875d970454d6", size = 20257584, upload-time = "2026-08-12T22:27:38.472Z" }, + { url = "https://files.pythonhosted.org/packages/4f/5b/4a756265d7fb164336c8d377bca21c39cfa2c178be23cedee840a69b59c5/av-18.1.0-cp314-cp314t-macosx_11_0_x86_64.whl", hash = "sha256:b36b0bae9e4c62f9487c99481ec15e4e3870fcc868522cd6d18fc2d6bfa04f01", size = 22795654, upload-time = "2026-08-12T22:27:42.016Z" }, + { url = "https://files.pythonhosted.org/packages/d5/cc/1bc841462114a1adf4f7d87456ab78a6972e23271e71865fcd2bbd0e7360/av-18.1.0-cp314-cp314t-macosx_14_0_arm64.whl", hash = "sha256:025f84494cb23278498f03b0d8117d3e47a1cbc9c44b97eb31875cf02251e46b", size = 18435735, upload-time = "2026-08-12T22:27:45.787Z" }, + { url = "https://files.pythonhosted.org/packages/b8/20/005500ed17a2e62a5e4bb94aa3786942560ec2f55ec1895ebf174c87abef/av-18.1.0-cp314-cp314t-manylinux_2_28_aarch64.whl", hash = "sha256:08a9ae288299cfcbf739dba4ad0c53b9b71f45184303dd45947920d022fed695", size = 37090807, upload-time = "2026-08-12T22:27:50.14Z" }, + { url = "https://files.pythonhosted.org/packages/5c/f7/11e7f6d848d3690c31ca4f8578167393e619177f1493ccc93b9400852d4e/av-18.1.0-cp314-cp314t-manylinux_2_28_x86_64.whl", hash = "sha256:cf8a17466bef07765dbdecc9e66ed9b25d20b4e14f654fbf35345a58ac45fa0c", size = 38976836, upload-time = "2026-08-12T22:27:54.565Z" }, + { url = "https://files.pythonhosted.org/packages/c3/63/b271473b24e806062d31191e40c6d65545e9cf59f80f044eba56dcbba0f4/av-18.1.0-cp314-cp314t-manylinux_2_31_armv7l.whl", hash = "sha256:d49a5c542dfdc00f43c6cdb6cc41dac1781ee206fe180b56aa7433dfa816dfae", size = 40896630, upload-time = "2026-08-12T22:27:59.118Z" }, + { url = "https://files.pythonhosted.org/packages/6b/9f/2ab7fa292a947ad3466ed8e655eefa3b82f535d7ea598c297b4471a937c4/av-18.1.0-cp314-cp314t-musllinux_1_2_aarch64.whl", hash = "sha256:5548b79e2bf1f59b3e9aedc918a72d9dc45b9adaac10ff9470d5dbdda0002e47", size = 37895673, upload-time = "2026-08-12T22:28:03.98Z" }, + { url = "https://files.pythonhosted.org/packages/e9/d8/04507c57249b399c3e4f23f01d221532f357338b5316fd2858fbd343127d/av-18.1.0-cp314-cp314t-musllinux_1_2_x86_64.whl", hash = "sha256:e7ea063f6690193ea335a1d592d6e0274350d45e2ed6af83ee107cb90cbfd84f", size = 39992431, upload-time = "2026-08-12T22:28:08.736Z" }, + { url = "https://files.pythonhosted.org/packages/d6/d6/bc4b95bea9c2353a7e4d62a3fcfad9adcf0f881741c6ce01ee179d539ce3/av-18.1.0-cp314-cp314t-win_amd64.whl", hash = "sha256:e4d48b9f12cad009cc72fe4f4099107de5e819c95f82767f4fd01a01481c0661", size = 28497798, upload-time = "2026-08-12T22:28:13.003Z" }, + { url = "https://files.pythonhosted.org/packages/c1/d2/0c277a46f12647c1833f40496e132fb6001e0d19e6144b5ea30896461feb/av-18.1.0-cp314-cp314t-win_arm64.whl", hash = "sha256:5cd9085028902c9880622bd37a12fd4b33060f06a52311f6f4867ca9f29a2c3b", size = 21421979, upload-time = "2026-08-12T22:28:16.48Z" }, +] + [[package]] name = "beautifulsoup4" version = "4.15.0" @@ -1070,9 +1096,13 @@ native-build = [ openai = [ { name = "openai" }, ] +video = [ + { name = "av" }, +] [package.dev-dependencies] dev = [ + { name = "av" }, { name = "cython" }, { name = "hflow-server" }, { name = "obstore" }, @@ -1087,6 +1117,7 @@ dev = [ [package.metadata] requires-dist = [ + { name = "av", marker = "extra == 'video'", specifier = ">=18.1.0" }, { name = "cython", marker = "extra == 'native-build'", specifier = ">=3.3.0" }, { name = "duckdb", specifier = ">=1.5.5" }, { name = "foxglove-schemas-protobuf", specifier = ">=0.4.0" }, @@ -1107,10 +1138,11 @@ requires-dist = [ { name = "tenacity", specifier = ">=9.1.4" }, { name = "zstandard", specifier = ">=0.25.0" }, ] -provides-extras = ["arrow", "bucket", "mediapipe", "motion", "native-build", "openai"] +provides-extras = ["arrow", "bucket", "mediapipe", "motion", "video", "native-build", "openai"] [package.metadata.requires-dev] dev = [ + { name = "av", specifier = ">=18.1.0" }, { name = "cython", specifier = ">=3.3.0" }, { name = "hflow-server", editable = "packages/hflow-server" }, { name = "obstore", specifier = ">=0.11.1" },