From 6369d7ce11411f3f585d5e09ac9e1e5158554dd2 Mon Sep 17 00:00:00 2001 From: "Pavel A. Tomskikh" Date: Mon, 13 Jul 2026 12:42:50 +0700 Subject: [PATCH 1/3] Support metadata passing to converters --- samples/output_converter_metadata/README.md | 81 ++++ .../output_converter_metadata/converter.py | 77 +++ .../docker-compose.l4t.yml | 87 ++++ .../docker-compose.x86.yml | 52 ++ samples/output_converter_metadata/module.yml | 93 ++++ .../output_converter_metadata/set-config.sh | 14 + savant/base/converter.py | 17 +- savant/deepstream/nvinfer/processor.py | 456 ++++++++++-------- 8 files changed, 679 insertions(+), 198 deletions(-) create mode 100644 samples/output_converter_metadata/README.md create mode 100644 samples/output_converter_metadata/converter.py create mode 100644 samples/output_converter_metadata/docker-compose.l4t.yml create mode 100644 samples/output_converter_metadata/docker-compose.x86.yml create mode 100644 samples/output_converter_metadata/module.yml create mode 100755 samples/output_converter_metadata/set-config.sh diff --git a/samples/output_converter_metadata/README.md b/samples/output_converter_metadata/README.md new file mode 100644 index 000000000..2f9082aa8 --- /dev/null +++ b/samples/output_converter_metadata/README.md @@ -0,0 +1,81 @@ +# Per-source Converter Configuration from Etcd + +A simple pipeline demonstrates how metadata processing in output converters works in Savant. In the demo, two RTSP streams are ingested in the module and processed with the PeopleNet model. The output converter is configurable via etcd. + +The resulting streams can be accessed via LL-HLS on `http://locahost:888/stream/city-traffic` and `http://locahost:888/stream/town-centre` or via RTSP on `rtsp://127.0.0.1:554/stream/city-traffic` and `rtsp://127.0.0.1:554/stream/town-centre`. + +Two RTSP streams (`city-traffic` and `town-centre`) are ingested by a single module and +processed with a YOLO11n detector. The detector's output converter is a custom subclass +of the built-in YOLO converter that, for every frame, reads the frame's `source_id` from +the converter `metadata` argument and looks up a per-source configuration object in Etcd +(e.g. the detection `confidence_threshold`). Values can be changed **live** with +`etcdctl` — no pipeline restart required. + +This relies on the output-converter `metadata` argument: when a converter's `__call__` +declares a `metadata` parameter it receives the frame's `NvDsFrameMeta` wrapper +(`source_id`, `pts`, `video_frame`, objects, tags). Converters that do not declare it keep +working unchanged. See `samples/output_converter_metadata/converter.py`. + +The resulting streams can be accessed via LL-HLS on +`http://localhost:888/stream/city-traffic` and `http://localhost:888/stream/town-centre`. + +Tested on platforms: + +- Nvidia Ampere + +## Prerequisites + +```bash +git clone https://github.com/insight-platform/Savant.git +cd Savant +git lfs pull +./utils/check-environment-compatible +``` + +**Note**: Ubuntu 22.04 runtime configuration [guide](https://insight-platform.github.io/Savant/develop/getting_started/0_configure_prod_env.html) helps to configure the runtime to run Savant pipelines. + +## Build Engines + +The demo uses models that are compiled into TensorRT engines the first time the demo is run. This takes time. Optionally, you can prepare the engines before running the demo by using the command: + +```bash +# you are expected to be in Savant/ directory + +./scripts/run_module.py --build-engines samples/output_converter_metadata/module.yml +``` + +## Run Demo + +```bash +# you are expected to be in Savant/ directory + +# if x86 +docker compose -f samples/output_converter_metadata/docker-compose.x86.yml up + +# if Jetson +docker compose -f samples/output_converter_metadata/docker-compose.l4t.yml up + +# open 'rtsp://127.0.0.1:554/stream/city-traffic' in your player +# or visit 'http://127.0.0.1:888/stream/city-traffic' (LL-HLS) + +# open 'rtsp://127.0.0.1:554/stream/town-centre' in your player +# or visit 'http://127.0.0.1:888/stream/town-centre' (LL-HLS) + +# Ctrl+C to stop running the compose bundle +``` + +## Per-source Configuration + +The converter reads the Etcd key `savant/source/` as a JSON object. +Supported fields: `confidence_threshold` and `nms_iou_threshold`. Use the helper script to +set or update a source's configuration (the `etcd` service must be running): + +```bash +# you are expected to be in Savant/samples/output_converter_metadata/ directory + +# keep low-confidence detections on city-traffic (more boxes) +./set-config.sh city-traffic '{"confidence_threshold": 0.2}' + +# require high confidence on town-centre (fewer boxes) +./set-config.sh town-centre '{"confidence_threshold": 0.7}' +``` diff --git a/samples/output_converter_metadata/converter.py b/samples/output_converter_metadata/converter.py new file mode 100644 index 000000000..791d9aa00 --- /dev/null +++ b/samples/output_converter_metadata/converter.py @@ -0,0 +1,77 @@ +"""Detector output converter that pulls per-source config from Etcd.""" + +import json +from typing import Optional, Tuple + +import numpy as np +from savant_rs.utils import eval_expr + +from savant.base.model import ObjectModel +from savant.converter.yolo import TensorToBBoxConverter +from savant.deepstream.meta.frame import NvDsFrameMeta + +# how long a fetched Etcd value stays cached locally (seconds) +CONFIG_CACHE_TTL = 5 + + +class EtcdConfigurableConverter(TensorToBBoxConverter): + """YOLO bbox converter whose thresholds are overridden per source_id from Etcd.""" + + def __init__(self, **kwargs): + self._default_confidence_threshold = kwargs.get('confidence_threshold', 0.25) + self._default_nms_iou_threshold = kwargs.get('nms_iou_threshold', 0.0) + self._configs = {} + super().__init__(**kwargs) + + def _load_source_config(self, source_id: str) -> dict: + expr = f'etcd("source/{source_id}", "")' + val, is_cached = eval_expr(expr, ttl=CONFIG_CACHE_TTL, no_gil=True) + if not is_cached: + if val: + try: + self._configs[source_id] = json.loads(val) + except json.JSONDecodeError: + self.logger.warning( + 'Invalid JSON in Etcd config for source %s: %r', source_id, val + ) + self._configs[source_id] = {} + else: + self._configs[source_id] = {} + + return self._configs.get(source_id) + + def __call__( + self, + *output_layers: np.ndarray, + model: ObjectModel, + roi: Tuple[float, float, float, float], + metadata: Optional[NvDsFrameMeta] = None, + ) -> Optional[np.ndarray]: + """Converts detector output layer tensor to bbox tensor. + + :param output_layers: Output layer tensor + :param model: Model definition, required parameters: input tensor shape, + maintain_aspect_ratio + :param roi: [left, top, width, height] of the rectangle + on which the model infers + :param metadata: Frame metadata. + :return: BBox tensor, see the base converter. + """ + + config = {} + if metadata is not None: + config = self._load_source_config(metadata.source_id) + self.logger.debug( + 'Source %s converter config: %s', metadata.source_id, config + ) + + # per-source override with fallback to construction-time defaults + self.confidence_threshold = config.get( + 'confidence_threshold', self._default_confidence_threshold + ) + self.nms_iou_threshold = config.get( + 'nms_iou_threshold', self._default_nms_iou_threshold + ) + + # reuse the parent's YOLO tensor decoding / NMS / coordinate transform + return super().__call__(*output_layers, model=model, roi=roi) diff --git a/samples/output_converter_metadata/docker-compose.l4t.yml b/samples/output_converter_metadata/docker-compose.l4t.yml new file mode 100644 index 000000000..f77a7aefd --- /dev/null +++ b/samples/output_converter_metadata/docker-compose.l4t.yml @@ -0,0 +1,87 @@ +services: + + rtsp-city-traffic: + image: ghcr.io/insight-platform/savant-adapters-gstreamer-l4t:latest + restart: unless-stopped + volumes: + - zmq_sockets:/tmp/zmq-sockets + environment: + - RTSP_URI=rtsp://hello.savant.video:8554/stream/city-traffic + - ZMQ_ENDPOINT=pub+connect:ipc:///tmp/zmq-sockets/input-video.ipc + - SOURCE_ID=city-traffic + entrypoint: /opt/savant/adapters/gst/sources/rtsp.sh + depends_on: + module: + condition: service_healthy + + rtsp-town-centre: + image: ghcr.io/insight-platform/savant-adapters-gstreamer-l4t:latest + restart: unless-stopped + volumes: + - zmq_sockets:/tmp/zmq-sockets + environment: + - RTSP_URI=rtsp://hello.savant.video:8554/stream/town-centre + - ZMQ_ENDPOINT=pub+connect:ipc:///tmp/zmq-sockets/input-video.ipc + - SOURCE_ID=town-centre + entrypoint: /opt/savant/adapters/gst/sources/rtsp.sh + depends_on: + module: + condition: service_healthy + + module: + privileged: true + image: ghcr.io/insight-platform/savant-deepstream-l4t:latest + restart: unless-stopped + volumes: + - zmq_sockets:/tmp/zmq-sockets + - ../../cache:/cache + - ..:/opt/savant/samples + command: samples/output_converter_metadata/module.yml + environment: + - MODEL_PATH=/cache/models/yolo11 + - DOWNLOAD_PATH=/cache/downloads/yolo11 + - ZMQ_SRC_ENDPOINT=sub+bind:ipc:///tmp/zmq-sockets/input-video.ipc + - ZMQ_SINK_ENDPOINT=pub+bind:ipc:///tmp/zmq-sockets/output-video.ipc + - METRICS_FRAME_PERIOD=1000 + - CODEC=jpeg + depends_on: + etcd: + condition: service_healthy + runtime: nvidia + + always-on-sink: + image: ghcr.io/insight-platform/savant-adapters-deepstream-l4t:latest + restart: unless-stopped + ports: + - "554:554" # RTSP + - "1935:1935" # RTMP + - "888:888" # HLS + - "8889:8889" # WebRTC + volumes: + - zmq_sockets:/tmp/zmq-sockets + - ../assets/stub_imgs:/stub_imgs + environment: + - ZMQ_ENDPOINT=sub+connect:ipc:///tmp/zmq-sockets/output-video.ipc + - SOURCE_IDS=city-traffic,town-centre + - FRAMERATE=25/1 + - STUB_FILE_LOCATION=/stub_imgs/smpte100_1280x720.jpeg + - DEV_MODE=True + command: python -m adapters.ds.sinks.always_on_rtsp + + etcd: + container_name: etcd + image: bitnamilegacy/etcd:3.6.4-debian-12-r4 + restart: unless-stopped + environment: + - ALLOW_NONE_AUTHENTICATION=yes + - ETCD_ADVERTISE_CLIENT_URLS=http://etcd:2379 + ports: + - "2379:2379" + healthcheck: + test: [ "CMD", "/opt/bitnami/scripts/etcd/healthcheck.sh" ] + interval: 5s + timeout: 5s + retries: 3 + +volumes: + zmq_sockets: diff --git a/samples/output_converter_metadata/docker-compose.x86.yml b/samples/output_converter_metadata/docker-compose.x86.yml new file mode 100644 index 000000000..aab391aab --- /dev/null +++ b/samples/output_converter_metadata/docker-compose.x86.yml @@ -0,0 +1,52 @@ +services: + + rtsp-city-traffic: + image: ghcr.io/insight-platform/savant-adapters-gstreamer:latest + extends: + file: docker-compose.l4t.yml + service: rtsp-city-traffic + + rtsp-town-centre: + image: ghcr.io/insight-platform/savant-adapters-gstreamer:latest + extends: + file: docker-compose.l4t.yml + service: rtsp-town-centre + + module: + privileged: true + image: ghcr.io/insight-platform/savant-deepstream:latest + extends: + file: docker-compose.l4t.yml + service: module + runtime: runc + environment: + - CODEC=h264 + deploy: + resources: + reservations: + devices: + - driver: nvidia + count: 1 + capabilities: [ gpu ] + + always-on-sink: + privileged: true + image: ghcr.io/insight-platform/savant-adapters-deepstream:latest + extends: + file: docker-compose.l4t.yml + service: always-on-sink + deploy: + resources: + reservations: + devices: + - driver: nvidia + count: 1 + capabilities: [ gpu ] + + etcd: + extends: + file: docker-compose.l4t.yml + service: etcd + +volumes: + zmq_sockets: diff --git a/samples/output_converter_metadata/module.yml b/samples/output_converter_metadata/module.yml new file mode 100644 index 000000000..7224dc389 --- /dev/null +++ b/samples/output_converter_metadata/module.yml @@ -0,0 +1,93 @@ +# module name, required +name: ${oc.env:MODULE_NAME, 'output_converter_metadata'} + +# base module parameters +parameters: + output_frame: + codec: ${oc.env:CODEC, 'raw-rgba'} + # PyFunc for drawing on frames (default implementation) + draw_func: + module: savant.deepstream.drawfunc + class_name: NvDsDrawFunc + rendered_objects: + detector: + person: + bbox: + border_color: 'FFFFFFFF' # Green + background_color: '00000077' # semi-transparent black + thickness: 2 + label: + format: [ '{label}' ] + font_color: 'FFFFFFFF' # White + background_color: '000000AA' # semi-transparent black + border_width: 2 + border_color: 'FFFFFFFF' # white + padding: [10, 0, 10, 0] # left, top, right, bottom + position: + position: TopLeftOutside + margin_x: 0 + margin_y: -10 + + # Etcd storage to manage processing sources + etcd: + # Etcd hosts to connect to + hosts: [etcd:2379] + # Path in Etcd to watch changes + watch_path: savant + + detected_object: + id: 0 + label: person + + # Custom output converter that overrides its thresholds per source_id from Etcd. + # kwargs are the fallback defaults used when a source has no Etcd config yet. + default_yolo_converter: + module: samples.output_converter_metadata.converter + class_name: EtcdConfigurableConverter + kwargs: + confidence_threshold: 0.25 + nms_iou_threshold: 0.45 + top_k: 300 + + default_yolo_selector: + module: savant.selector.detector + class_name: MinMaxSizeBBoxSelector + kwargs: + min_width: 30 + min_height: 30 + + batch_size: 1 + +# pipeline definition +pipeline: + # source definition is skipped, zeromq source is used by default to connect with source adapters + + # define pipeline's main elements + elements: + # primary detector element, inference is provided by the nvinfer Deepstream element + # model type is detector (other available types are: classifier, custom) + - element: nvinfer@detector + name: detector + model: + remote: + url: s3://savant-data/models/yolo11n/yolo11n.zip + checksum_url: s3://savant-data/models/yolo11n/yolo11n.md5 + parameters: + endpoint: https://eu-central-1.linodeobjects.com + format: onnx + model_file: yolo11n.onnx + batch_size: ${parameters.batch_size} + workspace_size: 6144 + input: + shape: [3, 640, 640] + scale_factor: 0.0039215697906911373 + maintain_aspect_ratio: true + symmetric_padding: true + output: + layer_names: [output0] + num_detected_classes: 80 # required for YOLOv11 + converter: ${parameters.default_yolo_converter} + objects: + - class_id: ${parameters.detected_object.id} + label: ${parameters.detected_object.label} + selector: ${parameters.default_yolo_selector} diff --git a/samples/output_converter_metadata/set-config.sh b/samples/output_converter_metadata/set-config.sh new file mode 100755 index 000000000..6091c186e --- /dev/null +++ b/samples/output_converter_metadata/set-config.sh @@ -0,0 +1,14 @@ +#!/usr/bin/env bash +# Set/update the per-source converter configuration in Etcd. +# +# Usage: set-config.sh +# e.g. ./set-config.sh city-traffic '{"confidence_threshold": 0.2}' +# ./set-config.sh town-centre '{"confidence_threshold": 0.7, "nms_iou_threshold": 0.5}' +# +# The real Etcd key is "/source/", +# where watch_path is "savant" (see module.yml). + +source=${1:-city-traffic} +config=${2:-'{"confidence_threshold": 0.25}'} + +docker exec -it etcd etcdctl put "savant/source/$source" "$config" diff --git a/savant/base/converter.py b/savant/base/converter.py index bc3be244f..13326c3aa 100644 --- a/savant/base/converter.py +++ b/savant/base/converter.py @@ -2,7 +2,7 @@ from abc import abstractmethod from enum import Enum -from typing import Any, List, Optional, Tuple, Union +from typing import TYPE_CHECKING, Any, List, Optional, Tuple, Union import cupy as cp import numpy as np @@ -10,6 +10,9 @@ from .model import AttributeModel, ComplexModel, ObjectModel from .pyfunc import BasePyFuncCallableImpl +if TYPE_CHECKING: + from savant.deepstream.meta.frame import NvDsFrameMeta + class TensorFormat(Enum): """Enum of the array module to be used to represent the tensors.""" @@ -36,8 +39,12 @@ def __call__( *output_layers: Union[np.ndarray, cp.ndarray], model: ObjectModel, roi: Tuple[float, float, float, float], + metadata: Optional['NvDsFrameMeta'] = None, ) -> Any: - """Converts raw model output tensors to a model specific representation.""" + """Converts raw model output tensors to a model specific representation. + + :param metadata: Frame metadata. Optional for backward compatibility. + """ class BaseObjectModelOutputConverter(BaseOutputConverter): @@ -49,6 +56,7 @@ def __call__( *output_layers: Union[np.ndarray, cp.ndarray], model: ObjectModel, roi: Tuple[float, float, float, float], + metadata: Optional['NvDsFrameMeta'] = None, ) -> Optional[np.ndarray]: """Converts raw model output tensors to a numpy array that represents a list of detected bboxes in the format ``(class_id, confidence, xc, yc, @@ -60,6 +68,7 @@ def __call__( maintain_aspect_ratio flag :param roi: ``[top, left, width, height]`` of the rectangle on which the model infers + :param metadata: Frame metadata. Optional for backward compatibility. :return: BBox tensor ``(class_id, confidence, xc, yc, width, height, [angle])`` offset by roi upper left and scaled by roi width and height """ @@ -74,6 +83,7 @@ def __call__( *output_layers: Union[np.ndarray, cp.ndarray], model: AttributeModel, roi: Tuple[float, float, float, float], + metadata: Optional['NvDsFrameMeta'] = None, ) -> Optional[List[Tuple[str, Any, float]]]: """Converts raw model output tensors to a list of values in several formats: @@ -90,6 +100,7 @@ def __call__( :param model: Attribute model :param roi: ``[top, left, width, height]`` of the rectangle on which the model infers + :param metadata: Frame metadata. Optional for backward compatibility. :return: list of attributes values with confidences ``(attr_name, value, confidence)`` """ @@ -104,6 +115,7 @@ def __call__( *output_layers: Union[np.ndarray, cp.ndarray], model: ComplexModel, roi: Tuple[float, float, float, float], + metadata: Optional['NvDsFrameMeta'] = None, ) -> Optional[Tuple[np.ndarray, List[List[Tuple[str, Any, float]]]]]: """Converts raw model output tensors to Savant format. @@ -112,6 +124,7 @@ def __call__( maintain_aspect_ratio flag :param roi: ``[top, left, width, height]`` of the rectangle on which the model infers + :param metadata: Frame metadata. Optional for backward compatibility. :return: a combination of :py:class:`.BaseObjectModelOutputConverter` and :py:class:`.BaseAttributeModelOutputConverter` outputs: diff --git a/savant/deepstream/nvinfer/processor.py b/savant/deepstream/nvinfer/processor.py index 17a4dc605..35b524615 100644 --- a/savant/deepstream/nvinfer/processor.py +++ b/savant/deepstream/nvinfer/processor.py @@ -1,5 +1,7 @@ +import inspect import logging -from typing import Callable, List, Optional, Tuple, Union +from contextlib import contextmanager, nullcontext +from typing import Callable, Dict, List, Optional, Tuple, Union import numpy as np import pyds @@ -8,7 +10,9 @@ nvds_frame_meta_get_nvds_savant_frame_meta, ) from savant_rs.pipeline2 import VideoPipeline +from savant_rs.primitives import VideoFrame from savant_rs.primitives.geometry import BBox +from savant_rs.utils import TelemetrySpan from savant_rs.utils.symbol_mapper import ( build_model_object_key, get_model_id, @@ -20,6 +24,7 @@ from savant.base.input_preproc import ObjectsPreprocessing from savant.base.pyfunc import PyFuncNoopCallException from savant.config.schema import FramePadding, ModelElement +from savant.deepstream.meta.frame import NvDsFrameMeta from savant.deepstream.meta.object import _NvDsObjectMetaImpl from savant.deepstream.utils.attribute import ( nvds_add_attr_meta_to_obj, @@ -122,6 +127,11 @@ def no_op(*args): self._restore_object_meta = self._restore_object_meta_ self._restore_frame = self._restore_frame_ + # cache of "does the resolved converter's __call__ accept `metadata`", + # keyed by the converter class (a dev-mode reload yields a new class + # object, which invalidates the entry automatically) + self._converter_accepts_metadata_cache: Dict[type, bool] = {} + if self._model.output.converter: self.postproc = self._process_custom_model_output self._tensor_meta_to_outputs = nvds_infer_tensor_meta_to_outputs @@ -264,219 +274,243 @@ def _process_custom_model_output(self, buffer: Gst.Buffer): self._model: Union[NvInferAttributeModel, NvInferComplexModel] nvds_batch_meta = pyds.gst_buffer_get_nvds_batch_meta(hash(buffer)) for nvds_frame_meta in nvds_frame_meta_iterator(nvds_batch_meta): - source_id, frame_idx = self._get_frame_source_id_and_idx( + video_frame, video_frame_span = self._get_video_frame( buffer, nvds_frame_meta, ) + source_id = video_frame.source_id if video_frame is not None else None source_info = self._sources.get_source(source_id) frame_rect = self._frame_rect(nvds_frame_meta, source_info.padding) - for nvds_obj_meta in nvds_obj_meta_iterator(nvds_frame_meta): - self._restore_object_meta(nvds_obj_meta) - if not self._is_model_input_object(nvds_obj_meta): - continue - parent_nvds_obj_meta = nvds_obj_meta - for tensor_meta in nvds_tensor_output_iterator( - parent_nvds_obj_meta, gie_uid=self._model_uid - ): - if self._logger.isEnabledFor(logging.TRACE): - self._logger.trace( - 'Converting "%s" element tensor output for frame ' - 'with PTS %s.', - self._element_name, - nvds_frame_meta.buf_pts, - ) - # parse and post-process model output - output_layers = self._tensor_meta_to_outputs( - tensor_meta=tensor_meta, - layer_names=self._model.output.layer_names, - ) - try: - outputs = self._model.output.converter( - *output_layers, - model=self._model, - roi=( - parent_nvds_obj_meta.rect_params.left, - parent_nvds_obj_meta.rect_params.top, - parent_nvds_obj_meta.rect_params.width, - parent_nvds_obj_meta.rect_params.height, - ), - ) - except Exception as exc: # pylint: disable=broad-except - if self._model.output.converter.dev_mode: - if not isinstance(exc, PyFuncNoopCallException): - self._logger.exception('Error calling converter') - outputs = None - else: - raise exc - # for object/complex models output - `bbox_tensor` and - # `selected_bboxes` - indices of selected bboxes and meta - # for attribute/complex models output - `values` - bbox_tensor: Optional[np.ndarray] = None - selected_bboxes: Optional[List] = None - values: Optional[List] = None - - if outputs is None: - continue + # Build the frame-meta wrapper lazily, once per frame, and only when + # the converter declares a `metadata` parameter (backward compat). + pass_metadata = self._converter_accepts_metadata() + if ( + pass_metadata + and video_frame is not None + and video_frame_span is not None + ): + frame_meta_cm = _build_frame_meta_cm( + nvds_frame_meta, video_frame, video_frame_span + ) + else: + frame_meta_cm = nullcontext(None) - # complex model - if self._is_complex_model: - # output converter returns tensor and attribute values - bbox_tensor, values = outputs - assert bbox_tensor.shape[0] == len(values), ( - 'Number of detected boxes and attributes do not match.' - ) + with frame_meta_cm as frame_meta: + converter_kwargs = {'metadata': frame_meta} if pass_metadata else {} + self._process_single_frame_output( + nvds_batch_meta, nvds_frame_meta, frame_rect, converter_kwargs + ) + self._restore_frame(buffer) - # object model - elif self._is_object_model: - # output converter returns tensor with - # (class_id, confidence, xc, yc, width, height, [angle]), - # coordinates in roi scale (parent object scale) - bbox_tensor = outputs + def _process_single_frame_output( + self, + nvds_batch_meta, + nvds_frame_meta, + frame_rect, + converter_kwargs: dict, + ): + """Processes custom model output (converter wrapper) for a single frame.""" - # attribute model - else: - # output converter returns attribute values - values = outputs + for nvds_obj_meta in nvds_obj_meta_iterator(nvds_frame_meta): + self._restore_object_meta(nvds_obj_meta) + if not self._is_model_input_object(nvds_obj_meta): + continue - if bbox_tensor is not None and bbox_tensor.shape[0] > 0: - # object or complex model with non-empty output - if bbox_tensor.shape[1] == 6: # no angle - selection_type = ObjectSelectionType.REGULAR_BBOX + parent_nvds_obj_meta = nvds_obj_meta + for tensor_meta in nvds_tensor_output_iterator( + parent_nvds_obj_meta, gie_uid=self._model_uid + ): + if self._logger.isEnabledFor(logging.TRACE): + self._logger.trace( + 'Converting "%s" element tensor output for frame with PTS %s.', + self._element_name, + nvds_frame_meta.buf_pts, + ) + # parse and post-process model output + output_layers = self._tensor_meta_to_outputs( + tensor_meta=tensor_meta, + layer_names=self._model.output.layer_names, + ) + try: + outputs = self._model.output.converter( + *output_layers, + model=self._model, + roi=( + parent_nvds_obj_meta.rect_params.left, + parent_nvds_obj_meta.rect_params.top, + parent_nvds_obj_meta.rect_params.width, + parent_nvds_obj_meta.rect_params.height, + ), + **converter_kwargs, + ) + except Exception as exc: # pylint: disable=broad-except + if self._model.output.converter.dev_mode: + if not isinstance(exc, PyFuncNoopCallException): + self._logger.exception('Error calling converter') + outputs = None + else: + raise exc + # for object/complex models output - `bbox_tensor` and + # `selected_bboxes` - indices of selected bboxes and meta + # for attribute/complex models output - `values` + bbox_tensor: Optional[np.ndarray] = None + selected_bboxes: Optional[List] = None + values: Optional[List] = None + + if outputs is None: + continue - # xc -> left, yc -> top - bbox_tensor[:, 2] -= bbox_tensor[:, 4] / 2 - bbox_tensor[:, 3] -= bbox_tensor[:, 5] / 2 + # complex model + if self._is_complex_model: + # output converter returns tensor and attribute values + bbox_tensor, values = outputs + assert bbox_tensor.shape[0] == len(values), ( + 'Number of detected boxes and attributes do not match.' + ) - # width to right, height to bottom - bbox_tensor[:, 4] += bbox_tensor[:, 2] - bbox_tensor[:, 5] += bbox_tensor[:, 3] + # object model + elif self._is_object_model: + # output converter returns tensor with + # (class_id, confidence, xc, yc, width, height, [angle]), + # coordinates in roi scale (parent object scale) + bbox_tensor = outputs + + # attribute model + else: + # output converter returns attribute values + values = outputs + + if bbox_tensor is not None and bbox_tensor.shape[0] > 0: + # object or complex model with non-empty output + if bbox_tensor.shape[1] == 6: # no angle + selection_type = ObjectSelectionType.REGULAR_BBOX + + # xc -> left, yc -> top + bbox_tensor[:, 2] -= bbox_tensor[:, 4] / 2 + bbox_tensor[:, 3] -= bbox_tensor[:, 5] / 2 + + # width to right, height to bottom + bbox_tensor[:, 4] += bbox_tensor[:, 2] + bbox_tensor[:, 5] += bbox_tensor[:, 3] + + # clip + bbox_tensor[:, 2][bbox_tensor[:, 2] < frame_rect[0]] = ( + frame_rect[0] + ) + bbox_tensor[:, 3][bbox_tensor[:, 3] < frame_rect[1]] = ( + frame_rect[1] + ) + bbox_tensor[:, 4][bbox_tensor[:, 4] > frame_rect[2]] = ( + frame_rect[2] + ) + bbox_tensor[:, 5][bbox_tensor[:, 5] > frame_rect[3]] = ( + frame_rect[3] + ) - # clip - bbox_tensor[:, 2][bbox_tensor[:, 2] < frame_rect[0]] = ( - frame_rect[0] - ) - bbox_tensor[:, 3][bbox_tensor[:, 3] < frame_rect[1]] = ( - frame_rect[1] - ) - bbox_tensor[:, 4][bbox_tensor[:, 4] > frame_rect[2]] = ( - frame_rect[2] - ) - bbox_tensor[:, 5][bbox_tensor[:, 5] > frame_rect[3]] = ( - frame_rect[3] - ) + # right to width, bottom to height + bbox_tensor[:, 4] -= bbox_tensor[:, 2] + bbox_tensor[:, 5] -= bbox_tensor[:, 3] - # right to width, bottom to height - bbox_tensor[:, 4] -= bbox_tensor[:, 2] - bbox_tensor[:, 5] -= bbox_tensor[:, 3] - - # left -> xc , top-> yc - bbox_tensor[:, 2] += bbox_tensor[:, 4] / 2 - bbox_tensor[:, 3] += bbox_tensor[:, 5] / 2 - - # add 0 angle - bbox_tensor = np.concatenate( - [ - bbox_tensor, - np.zeros( - (bbox_tensor.shape[0], 1), dtype=np.float32 - ), - ], - axis=1, - ) - else: - selection_type = ObjectSelectionType.ROTATED_BBOX + # left -> xc , top-> yc + bbox_tensor[:, 2] += bbox_tensor[:, 4] / 2 + bbox_tensor[:, 3] += bbox_tensor[:, 5] / 2 - # add index column to further filter attribute values + # add 0 angle bbox_tensor = np.concatenate( [ bbox_tensor, - np.arange( - bbox_tensor.shape[0], dtype=np.float32 - ).reshape(-1, 1), + np.zeros((bbox_tensor.shape[0], 1), dtype=np.float32), ], axis=1, ) + else: + selection_type = ObjectSelectionType.ROTATED_BBOX + + # add index column to further filter attribute values + bbox_tensor = np.concatenate( + [ + bbox_tensor, + np.arange(bbox_tensor.shape[0], dtype=np.float32).reshape( + -1, 1 + ), + ], + axis=1, + ) - selected_bboxes = [] - for obj in self._model.output.objects: - cls_bbox_tensor = bbox_tensor[ - bbox_tensor[:, 0] == obj.class_id - ] - if cls_bbox_tensor.shape[0] == 0: - continue - if obj.selector: - try: - cls_bbox_tensor = obj.selector(cls_bbox_tensor) - except Exception as exc: # pylint: disable=broad-except - if obj.selector.dev_mode: - if not isinstance(exc, PyFuncNoopCallException): - self._logger.exception( - 'Error calling selector.' - ) - cls_bbox_tensor = np.zeros((0, 8)) - else: - raise exc - - obj_label = build_model_object_key( - self._element_name, obj.label - ) - obj_cls_id = MERGED_CLASSES[self._element_name].get( - obj.class_id - ) - if obj_cls_id is None: - obj_cls_id = obj.class_id - else: - if self._logger.isEnabledFor(logging.TRACE): - self._logger.trace( - 'Updating %s custom objs id %s -> %s, ' - 'label "%s".', - len(cls_bbox_tensor), - obj.class_id, - obj_cls_id, - obj_label, - ) - for bbox in cls_bbox_tensor: - if self._logger.isEnabledFor(logging.TRACE): - self._logger.trace( - 'Adding obj %s into pyds meta for frame ' - 'with PTS %s.', - bbox[2:7], - nvds_frame_meta.buf_pts, - ) - _nvds_obj_meta = nvds_add_obj_meta_to_frame( - nvds_batch_meta, - nvds_frame_meta, - selection_type, + selected_bboxes = [] + for obj in self._model.output.objects: + cls_bbox_tensor = bbox_tensor[bbox_tensor[:, 0] == obj.class_id] + if cls_bbox_tensor.shape[0] == 0: + continue + if obj.selector: + try: + cls_bbox_tensor = obj.selector(cls_bbox_tensor) + except Exception as exc: # pylint: disable=broad-except + if obj.selector.dev_mode: + if not isinstance(exc, PyFuncNoopCallException): + self._logger.exception( + 'Error calling selector.' + ) + cls_bbox_tensor = np.zeros((0, 8)) + else: + raise exc + + obj_label = build_model_object_key( + self._element_name, obj.label + ) + obj_cls_id = MERGED_CLASSES[self._element_name].get( + obj.class_id + ) + if obj_cls_id is None: + obj_cls_id = obj.class_id + else: + if self._logger.isEnabledFor(logging.TRACE): + self._logger.trace( + 'Updating %s custom objs id %s -> %s, label "%s".', + len(cls_bbox_tensor), + obj.class_id, obj_cls_id, - self._model_uid, - bbox[2:7], - bbox[1], obj_label, - parent=parent_nvds_obj_meta, ) - selected_bboxes.append((int(bbox[7]), _nvds_obj_meta)) - - # attribute or complex model - if values: - if self._is_complex_model: - values = [values[i] for i, _ in selected_bboxes] - else: - selected_bboxes = [(0, nvds_obj_meta)] - values = [values] - for (_, _nvds_obj_meta), _values in zip( - selected_bboxes, values - ): - for attr_name, value, confidence in _values: - nvds_add_attr_meta_to_obj( - frame_meta=nvds_frame_meta, - obj_meta=_nvds_obj_meta, - element_name=self._element_name, - name=attr_name, - value=value, - confidence=confidence, + for bbox in cls_bbox_tensor: + if self._logger.isEnabledFor(logging.TRACE): + self._logger.trace( + 'Adding obj %s into pyds meta for frame ' + 'with PTS %s.', + bbox[2:7], + nvds_frame_meta.buf_pts, ) - self._restore_frame(buffer) + _nvds_obj_meta = nvds_add_obj_meta_to_frame( + nvds_batch_meta, + nvds_frame_meta, + selection_type, + obj_cls_id, + self._model_uid, + bbox[2:7], + bbox[1], + obj_label, + parent=parent_nvds_obj_meta, + ) + selected_bboxes.append((int(bbox[7]), _nvds_obj_meta)) + + # attribute or complex model + if values: + if self._is_complex_model: + values = [values[i] for i, _ in selected_bboxes] + else: + selected_bboxes = [(0, nvds_obj_meta)] + values = [values] + for (_, _nvds_obj_meta), _values in zip(selected_bboxes, values): + for attr_name, value, confidence in _values: + nvds_add_attr_meta_to_obj( + frame_meta=nvds_frame_meta, + obj_meta=_nvds_obj_meta, + element_name=self._element_name, + name=attr_name, + value=value, + confidence=confidence, + ) def _process_regular_detector_output(self, buffer: Gst.Buffer): """Processes output of nvinfer detector. @@ -561,11 +595,30 @@ def _is_model_input_object(self, nvds_obj_meta: pyds.NvDsObjectMeta): and nvds_obj_meta.class_id == self._input_object_class_id ) - def _get_frame_source_id_and_idx( + def _converter_accepts_metadata(self) -> bool: + """Whether the resolved output converter's ``__call__`` declares a + ``metadata`` parameter. + + Cached per converter class to avoid calling ``inspect.signature`` in the + hot per-tensor loop; the result is recomputed automatically after a + dev-mode reload because the reload replaces the class object. + """ + instance = self._model.output.converter.instance + converter_cls = type(instance) + cached = self._converter_accepts_metadata_cache.get(converter_cls) + if cached is None: + try: + cached = 'metadata' in inspect.signature(instance.__call__).parameters + except (TypeError, ValueError): + cached = False + self._converter_accepts_metadata_cache[converter_cls] = cached + return cached + + def _get_video_frame( self, buffer: Gst.Buffer, nvds_frame_meta: pyds.NvDsFrameMeta, - ) -> Tuple[Optional[str], Optional[int]]: + ) -> Tuple[Optional[VideoFrame], Optional[TelemetrySpan]]: savant_batch_meta = gst_buffer_get_savant_batch_meta(buffer) if savant_batch_meta is None: return None, None @@ -574,14 +627,10 @@ def _get_frame_source_id_and_idx( if savant_frame_meta is None: return None, None - frame_idx = savant_frame_meta.idx - video_frame, _ = self._video_pipeline.get_batched_frame( + return self._video_pipeline.get_batched_frame( savant_batch_meta.idx, - frame_idx, + savant_frame_meta.idx, ) - source_id = video_frame.source_id - - return source_id, frame_idx def _frame_rect(self, nvds_frame_meta, frame_padding: Optional[FramePadding]): """Frame rect to clip objects. @@ -602,3 +651,18 @@ def _frame_rect(self, nvds_frame_meta, frame_padding: Optional[FramePadding]): nvds_frame_meta.source_frame_width + frame_rect_shift[0] - 1.0, nvds_frame_meta.source_frame_height + frame_rect_shift[1] - 1.0, ) + + +@contextmanager +def _build_frame_meta_cm(nvds_frame_meta, video_frame, video_frame_span): + """Build a savant :class:`NvDsFrameMeta` wrapper for the frame, and open a nested + telemetry span for the output conversion. + """ + + with video_frame_span.nested_span('convert-output') as telemetry_span: + with NvDsFrameMeta( + nvds_frame_meta, + video_frame, + telemetry_span, + ) as frame_meta: + yield frame_meta From 90ef5a693124fc503231af7e4f361a43759005d2 Mon Sep 17 00:00:00 2001 From: Pavel Tomskikh Date: Mon, 13 Jul 2026 18:02:00 +0700 Subject: [PATCH 2/3] Handle the case when a frame doesn't have metadata Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- savant/deepstream/nvinfer/processor.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/savant/deepstream/nvinfer/processor.py b/savant/deepstream/nvinfer/processor.py index 35b524615..9965d12bf 100644 --- a/savant/deepstream/nvinfer/processor.py +++ b/savant/deepstream/nvinfer/processor.py @@ -278,7 +278,10 @@ def _process_custom_model_output(self, buffer: Gst.Buffer): buffer, nvds_frame_meta, ) - source_id = video_frame.source_id if video_frame is not None else None + if video_frame is not None: + source_id = video_frame.source_id + else: + source_id = self._sources.get_id_by_pad_index(nvds_frame_meta.pad_index) source_info = self._sources.get_source(source_id) frame_rect = self._frame_rect(nvds_frame_meta, source_info.padding) From 6cb11c4d06ff6fa803d43a0f32ea572cf25958b2 Mon Sep 17 00:00:00 2001 From: "Pavel A. Tomskikh" Date: Mon, 13 Jul 2026 18:03:40 +0700 Subject: [PATCH 3/3] Fix README and converter config parsing --- samples/mjpeg_usb_cam/README.md | 2 +- samples/multiple_gige/README.md | 2 +- samples/multiple_rtsp/README.md | 2 +- samples/output_converter_metadata/README.md | 4 ++-- samples/output_converter_metadata/converter.py | 7 +++++-- 5 files changed, 10 insertions(+), 7 deletions(-) diff --git a/samples/mjpeg_usb_cam/README.md b/samples/mjpeg_usb_cam/README.md index 746957c07..3491d2483 100644 --- a/samples/mjpeg_usb_cam/README.md +++ b/samples/mjpeg_usb_cam/README.md @@ -2,7 +2,7 @@ A pipeline demonstrating how to capture MJPEG from a USB camera. MJPEG is a common format for USB/MIPI CSI-2 cameras providing compressed, low-latency video streaming. -The resulting stream can be accessed via LL-HLS on `http://locahost:888/stream/video` +The resulting stream can be accessed via LL-HLS on `http://localhost:888/stream/video` Tested on platforms: diff --git a/samples/multiple_gige/README.md b/samples/multiple_gige/README.md index d7edf146f..1c384aafa 100644 --- a/samples/multiple_gige/README.md +++ b/samples/multiple_gige/README.md @@ -2,7 +2,7 @@ A simple pipeline demonstrates how GigE Vision Source Adapter works in Savant. In the demo video from one GigE Vision camera is passed as raw-rgba frames, and another one is passed as HEVC-encoded frames. Both streams are passed to an Always-On-RTSP sink. -The resulting streams can be accessed via LL-HLS on `http://locahost:888/stream/gige-raw` (raw-rgba frames) and `http://locahost:888/stream/gige-encoded` (HEVC-encoded frames). +The resulting streams can be accessed via LL-HLS on `http://localhost:888/stream/gige-raw` (raw-rgba frames) and `http://localhost:888/stream/gige-encoded` (HEVC-encoded frames). Tested on platforms: diff --git a/samples/multiple_rtsp/README.md b/samples/multiple_rtsp/README.md index 3de67cc6e..ed0367314 100644 --- a/samples/multiple_rtsp/README.md +++ b/samples/multiple_rtsp/README.md @@ -2,7 +2,7 @@ A simple pipeline demonstrates how multiplexed processing works in Savant. In the demo, two RTSP streams are ingested in the module and processed with the PeopleNet model. -The resulting streams can be accessed via LL-HLS on `http://locahost:888/stream/city-traffic` and `http://locahost:888/stream/town-centre`. +The resulting streams can be accessed via LL-HLS on `http://localhost:888/stream/city-traffic` and `http://localhost:888/stream/town-centre`. Tested on platforms: diff --git a/samples/output_converter_metadata/README.md b/samples/output_converter_metadata/README.md index 2f9082aa8..66a974ed3 100644 --- a/samples/output_converter_metadata/README.md +++ b/samples/output_converter_metadata/README.md @@ -1,8 +1,8 @@ # Per-source Converter Configuration from Etcd -A simple pipeline demonstrates how metadata processing in output converters works in Savant. In the demo, two RTSP streams are ingested in the module and processed with the PeopleNet model. The output converter is configurable via etcd. +A simple pipeline demonstrates how metadata processing in output converters works in Savant. In the demo, two RTSP streams are ingested in the module and processed with the YOLO11n model. The output converter is configurable via etcd. -The resulting streams can be accessed via LL-HLS on `http://locahost:888/stream/city-traffic` and `http://locahost:888/stream/town-centre` or via RTSP on `rtsp://127.0.0.1:554/stream/city-traffic` and `rtsp://127.0.0.1:554/stream/town-centre`. +The resulting streams can be accessed via LL-HLS on `http://localhost:888/stream/city-traffic` and `http://localhost:888/stream/town-centre` or via RTSP on `rtsp://127.0.0.1:554/stream/city-traffic` and `rtsp://127.0.0.1:554/stream/town-centre`. Two RTSP streams (`city-traffic` and `town-centre`) are ingested by a single module and processed with a YOLO11n detector. The detector's output converter is a custom subclass diff --git a/samples/output_converter_metadata/converter.py b/samples/output_converter_metadata/converter.py index 791d9aa00..89e50ad9e 100644 --- a/samples/output_converter_metadata/converter.py +++ b/samples/output_converter_metadata/converter.py @@ -29,7 +29,10 @@ def _load_source_config(self, source_id: str) -> dict: if not is_cached: if val: try: - self._configs[source_id] = json.loads(val) + parsed_config = json.loads(val) + self._configs[source_id] = ( + parsed_config if isinstance(parsed_config, dict) else {} + ) except json.JSONDecodeError: self.logger.warning( 'Invalid JSON in Etcd config for source %s: %r', source_id, val @@ -38,7 +41,7 @@ def _load_source_config(self, source_id: str) -> dict: else: self._configs[source_id] = {} - return self._configs.get(source_id) + return self._configs.get(source_id, {}) def __call__( self,