lordgamez commented on code in PR #2258:
URL: https://github.com/apache/nifi-minifi-cpp/pull/2258#discussion_r4092720063


##########
minifi_rust/extensions/minifi_tensor/minifi_tensor.md:
##########


Review Comment:
   These could also be added to PROCESSORS.md



##########
minifi_rust/extensions/minifi_tensor/features/resources/.gitignore:
##########


Review Comment:
   Missing license header here and additionally in 
minifi_rust/extensions/minifi_tensor/src/services/tract_model_service.rs, 
minifi_rust/extensions/minifi_tensor/src/utils/dimensions.rs, 
minifi_rust/extensions/minifi_tensor/src/utils/per_channel_f32.rs, 
minifi_rust/extensions/minifi_tensor/src/utils/score_activation.rs, 
minifi_rust/extensions/minifi_tensor/src/utils/tensor_helpers.rs, 
minifi_rust/extensions/minifi_tensor/src/low_level_processors/mod.rs, 
minifi_rust/extensions/minifi_tensor/src/processors/draw_bounding_box.rs, 
minifi_rust/extensions/minifi_tensor/src/services/tract_model_service.rs, 
minifi_rust/extensions/minifi_tensor/src/utils/bounding_box.rs, 
minifi_rust/extensions/minifi_tensor/src/utils/dimensions.rs, 
minifi_rust/extensions/minifi_tensor/src/lib.rs files



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/classify_output.rs:
##########
@@ -0,0 +1,476 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   https://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::utils::score_activation::{ScoreActivation, SoftmaxTerms};
+use crate::utils::tensor_helpers::{deserialize_tensors, tensor_as_f32, 
tensor_shape};
+use classify_output_def::SUCCESS;
+pub(crate) use classify_output_def::{
+    CLASSIFY_OUTPUT_ATTRIBUTES, CONFIDENCE_THRESHOLD, LABEL_INDEX_OFFSET, 
LABELS_FILE_PATH,
+    OUTPUT_ATTRIBUTE_NAME, SCORE_ACTIVATION, SCORE_OUTPUT_INDEX, TOP_K,
+};
+use minifi_native::macros::ComponentIdentifier;
+use minifi_native::{
+    Content, FlowFileTransform, GetAttribute, GetId, GetProperty, InputStream, 
Logger, MinifiError,
+    ProcessError, RouteErrorExt, Schedule, TransformedFlowFile, warn,
+};
+use serde::Serialize;
+use std::path::Path;
+use tract::Tensor;
+
+mod classify_output_def;
+
+#[derive(Serialize, Clone, Debug, PartialEq)]
+struct Prediction {
+    class_id: usize,
+    confidence: f32,
+    #[serde(skip_serializing_if = "Option::is_none")]
+    class_name: Option<String>,
+}
+
+fn load_labels(path: &Path) -> Result<Vec<String>, MinifiError> {
+    let content = std::fs::read_to_string(path).map_err(|e| {
+        MinifiError::custom(format!("Failed to read labels file '{:?}': {}", 
path, e))
+    })?;
+    Ok(content
+        .lines()
+        .map(|line| line.trim().to_string())
+        .collect())
+}
+
+fn top_k(mut scored: Vec<(usize, f32)>, k: usize) -> Vec<(usize, f32)> {
+    scored.sort_by(|&(ai, a), &(bi, b)| b.total_cmp(&a).then(ai.cmp(&bi)));
+    scored.truncate(k);
+    scored
+}
+
+#[derive(ComponentIdentifier)]
+pub(crate) struct ClassifyOutput {
+    top_k: usize,
+    score_output_index: usize,
+    score_activation: ScoreActivation,
+    confidence_threshold: f32,
+    labels: Vec<String>,
+    label_index_offset: usize,
+}
+
+impl Schedule for ClassifyOutput {
+    fn schedule<Ctx: GetProperty, L: Logger>(
+        context: &Ctx,
+        _logger: &L,
+    ) -> Result<Self, MinifiError>
+    where
+        Self: Sized,
+    {
+        let top_k = context.get_property(&TOP_K)?;
+        if top_k == 0 {
+            return Err(MinifiError::validation("Top K must be >= 1"));
+        }
+        let score_output_index = context.get_property(&SCORE_OUTPUT_INDEX)?;
+        let score_activation = context.get_property(&SCORE_ACTIVATION)?;
+        let confidence_threshold = 
context.get_property(&CONFIDENCE_THRESHOLD)?;
+
+        let labels = match context.get_property(&LABELS_FILE_PATH)? {
+            Some(path) => load_labels(&path)?,
+            _ => Vec::new(),
+        };
+        let label_index_offset = context.get_property(&LABEL_INDEX_OFFSET)?;
+        if !labels.is_empty() && label_index_offset >= labels.len() {
+            return Err(MinifiError::validation(format!(
+                "Label index offset ({}) must be smaller than the number of 
labels ({})",
+                label_index_offset,
+                labels.len()
+            )));
+        }
+
+        Ok(Self {
+            top_k,
+            score_output_index,
+            score_activation,
+            confidence_threshold,
+            labels,
+            label_index_offset,
+        })
+    }
+}
+
+impl ClassifyOutput {
+    fn label_for(&self, class_id: usize) -> Option<String> {
+        self.labels
+            .get(class_id.checked_add(self.label_index_offset)?)
+            .cloned()
+    }
+
+    pub(crate) fn classify<'a, Context: GetProperty + GetAttribute + GetId, 
LoggerImpl: Logger>(
+        &self,
+        context: &Context,
+        logger: &LoggerImpl,
+        tensors: Vec<Tensor>,
+    ) -> Result<TransformedFlowFile<'a>, ProcessError> {
+        let score_floats =
+            tensor_as_f32(&tensors, 
self.score_output_index).route_err_to_failure()?;
+        if score_floats.is_empty() {
+            return Err(MinifiError::custom("Score tensor is empty; nothing to 
classify").into());
+        }
+
+        // A classifier head is a single score vector: shape [num_classes] or
+        // [1, .., num_classes]. We rank over the flattened class axis, so any
+        // leading axis > 1 (a real batch) would silently mix rows and yield
+        // class ids past num_classes. Reject it rather than produce garbage.
+        // (`ImageToTensor` emits batch=1 today; this just enforces the 
contract.)
+        let shape = tensor_shape(&tensors, 
self.score_output_index).route_err_to_failure()?;
+        if shape.iter().rev().skip(1).any(|&d| d != 1) {
+            return Err(MinifiError::custom(format!(
+                "ClassifyOutput expects a single score vector (shape 
[num_classes] or \
+                 [1, .., num_classes]); got {shape:?}. A batch dimension > 1 
is not supported."
+            )))
+            .route_err_to_failure();
+        }
+
+        let finite: Vec<(usize, f32)> = score_floats
+            .iter()
+            .copied()
+            .enumerate()
+            .filter(|&(_, s)| s.is_finite())
+            .collect();
+
+        let softmax_terms = SoftmaxTerms::over(finite.iter().map(|&(_, s)| s));
+
+        let predictions: Vec<Prediction> = top_k(finite, self.top_k)
+            .into_iter()
+            .filter_map(|(class_id, raw)| {
+                let confidence = self.score_activation.confidence(raw, 
softmax_terms);
+
+                if confidence >= self.confidence_threshold {
+                    let class_name = self.label_for(class_id);
+                    if class_name.is_none() && !self.labels.is_empty() {
+                        warn!(
+                            logger,
+                            "No label for class id {} (offset {}, {} labels 
loaded); \
+                             the labels file does not match the model's 
classes",
+                            class_id,
+                            self.label_index_offset,
+                            self.labels.len()
+                        );
+                    }
+                    Some(Prediction {
+                        class_id,
+                        confidence,
+                        class_name,
+                    })
+                } else {
+                    None
+                }
+            })
+            .collect();
+
+        let (content, extra_attribute) = match 
context.get_property(&OUTPUT_ATTRIBUTE_NAME)? {
+            None => (
+                Some(Content::Buffer(
+                    serde_json::to_vec(&predictions).route_err_to_failure()?,
+                )),
+                None,
+            ),
+            Some(output_attr) => (
+                None,
+                Some((
+                    output_attr,
+                    
serde_json::to_string(&predictions).route_err_to_failure()?,
+                )),
+            ),
+        };
+
+        let mut transformed = TransformedFlowFile::new(&SUCCESS, content)
+            .with_attribute("mime.type", "application/json")

Review Comment:
   The output attribute description says "Always 'application/json'", but if 
the output is written into the attribute instead of the flow file content, 
should that be the case here?



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/filter_bounding_boxes.rs:
##########
@@ -0,0 +1,593 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   https://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+mod filter_bounding_boxes_def;
+
+use crate::low_level_processors::image_to_tensor::ResizeMode;
+use crate::utils::bounding_box::BoundingBox;
+use crate::utils::dimensions::Dimensions;
+use crate::utils::score_activation::{ScoreActivation, SoftmaxTerms};
+use crate::utils::tensor_helpers::{deserialize_tensors, tensor_as_f32};
+use filter_bounding_boxes_def::SUCCESS;
+pub(crate) use filter_bounding_boxes_def::{
+    BACKGROUND_CLASS_INDEX, BOX_FORMAT, BOX_OUTPUT_INDEX, CLASS_OUTPUT_INDEX, 
CONFIDENCE_THRESHOLD,
+    IOU_THRESHOLD, OUTPUT_ATTRIBUTE_NAME, SCORE_ACTIVATION, SCORE_OUTPUT_INDEX,
+};
+use minifi_native::macros::{ComponentIdentifier, PropertyType};
+use minifi_native::{
+    Content, FlowFileTransform, GetAttribute, GetId, GetProperty, InputStream, 
Logger, MinifiError,
+    ProcessError, RouteErrorExt, Schedule, TransformedFlowFile, debug, trace,
+};
+use strum_macros::{Display, EnumString, IntoStaticStr, VariantNames};
+use tract::Tensor;
+
+#[derive(
+    Debug, Clone, Copy, PartialEq, Display, EnumString, VariantNames, 
IntoStaticStr, PropertyType,
+)]
+#[strum(serialize_all = "PascalCase", const_into_str)]
+pub(crate) enum BoxFormat {
+    /// `[x_min, y_min, x_max, y_max]` — SSD, MobileNet-SSD, most PyTorch 
models.
+    Xyxy,
+    /// `[y_min, x_min, y_max, x_max]` — TensorFlow Object Detection API.
+    Yxyx,
+    /// `[cx, cy, w, h]` — YOLOv3/5/8 raw output (center + size).
+    Cxcywh,
+}
+
+/// Convert the four floats at `box_floats[offset..offset+4]` into a canonical
+/// `(x_min, y_min, x_max, y_max)` tuple, regardless of the source layout.
+fn decode_box(box_floats: &[f32], offset: usize, format: BoxFormat) -> (f32, 
f32, f32, f32) {
+    let a = box_floats[offset];
+    let b = box_floats[offset + 1];
+    let c = box_floats[offset + 2];
+    let d = box_floats[offset + 3];
+    match format {
+        BoxFormat::Xyxy => (a, b, c, d),
+        BoxFormat::Yxyx => (b, a, d, c),
+        BoxFormat::Cxcywh => {
+            let (cx, cy, w, h) = (a, b, c, d);
+            (cx - w / 2.0, cy - h / 2.0, cx + w / 2.0, cy + h / 2.0)
+        }
+    }
+}
+
+struct ScoredClass {
+    class_id: usize,
+    confidence: f32,
+}
+
+fn score_box(
+    logits: &[f32],
+    activation: ScoreActivation,
+    background_class_index: Option<usize>,
+) -> ScoredClass {
+    let num_classes = logits.len();
+
+    let best_valid = logits
+        .iter()
+        .enumerate()
+        .filter(|&(_, &logit)| logit.is_finite())
+        .filter(|&(id, _)| match background_class_index {
+            Some(bg_idx) => !(num_classes > 1 && id == bg_idx),
+            None => true,
+        })
+        .max_by(|a, b| a.1.total_cmp(b.1));
+
+    let (class_id, &best_logit) = match best_valid {
+        Some(val) => val,
+        None => {
+            return ScoredClass {
+                class_id: 0,
+                confidence: f32::NEG_INFINITY,
+            };
+        }
+    };
+
+    let confidence = activation.confidence(best_logit, 
SoftmaxTerms::over(logits.iter().copied()));
+
+    ScoredClass {
+        class_id,
+        confidence,
+    }
+}
+
+#[derive(ComponentIdentifier)]
+pub(crate) struct FilterBoundingBoxes {
+    confidence_threshold: f32,
+    iou_threshold: f32,
+    score_output_index: usize,
+    box_output_index: usize,
+    box_format: BoxFormat,
+    score_activation: ScoreActivation,
+    background_class_index: Option<usize>,
+    class_output_index: Option<usize>,
+}
+
+impl Schedule for FilterBoundingBoxes {
+    fn schedule<Ctx: GetProperty, L: Logger>(
+        context: &Ctx,
+        _logger: &L,
+    ) -> Result<Self, MinifiError> {
+        let confidence_threshold = 
context.get_property(&CONFIDENCE_THRESHOLD)?;
+        let iou_threshold = context.get_property(&IOU_THRESHOLD)?;
+        let score_output_index = context.get_property(&SCORE_OUTPUT_INDEX)?;
+        let box_output_index = context.get_property(&BOX_OUTPUT_INDEX)?;
+        let box_format = context.get_property(&BOX_FORMAT)?;
+        let score_activation = context.get_property(&SCORE_ACTIVATION)?;
+        let background_class_index = 
context.get_property(&BACKGROUND_CLASS_INDEX)?;
+        let class_output_index = context.get_property(&CLASS_OUTPUT_INDEX)?;
+
+        Ok(Self {
+            confidence_threshold,
+            iou_threshold,
+            score_output_index,
+            box_output_index,
+            box_format,
+            score_activation,
+            background_class_index,
+            class_output_index,
+        })
+    }
+}
+
+impl FilterBoundingBoxes {
+    fn result_via_output_attribute<'a, Context: GetProperty>(
+        &self,
+        context: &Context,
+        filtered_boxes: Vec<BoundingBox>,
+    ) -> Result<TransformedFlowFile<'a>, MinifiError> {
+        let output_attr = context.get_property(&OUTPUT_ATTRIBUTE_NAME)?;
+        let content = if output_attr.is_some() {
+            None
+        } else {
+            Some(Content::Buffer(
+                
serde_json::to_vec(&filtered_boxes).map_err(MinifiError::other)?,
+            ))
+        };
+
+        let mut transformed = TransformedFlowFile::new(&SUCCESS, content)
+            .with_attribute("object.count", filtered_boxes.len().to_string());
+        if let Some(attr) = output_attr {
+            transformed = transformed.with_attribute(
+                attr,
+                
serde_json::to_string(&filtered_boxes).map_err(MinifiError::other)?,
+            )
+        } else {
+            transformed = transformed.with_attribute("mime.type", 
"application/json");
+        }
+        Ok(transformed)
+    }
+
+    pub(crate) fn filter<'a, Context: GetProperty, LoggerImpl: Logger>(
+        &self,
+        context: &Context,
+        logger: &LoggerImpl,
+        tensors: Vec<Tensor>,
+        orig_dim: Dimensions,
+        target_dim: Dimensions,
+        resize_mode: ResizeMode,
+    ) -> Result<TransformedFlowFile<'a>, ProcessError> {
+        let score_floats =
+            tensor_as_f32(&tensors, 
self.score_output_index).route_err_to_failure()?;
+        let box_floats = tensor_as_f32(&tensors, 
self.box_output_index).route_err_to_failure()?;
+
+        let (scale_x, scale_y, pad_x, pad_y) = match resize_mode {
+            ResizeMode::Letterbox => {
+                let geometry = orig_dim.letterbox_into(target_dim);
+                (
+                    geometry.scale,
+                    geometry.scale,
+                    geometry.pad_x as f32,
+                    geometry.pad_y as f32,
+                )
+            }
+            ResizeMode::Stretch => (
+                target_dim.width / orig_dim.width,
+                target_dim.height / orig_dim.height,
+                0.0,
+                0.0,
+            ),
+        };
+
+        if !box_floats.len().is_multiple_of(4) {
+            return Err(MinifiError::custom(

Review Comment:
   should we route to failure here instead? Maybe I am misunderstanding, but if 
this returns an error, wouldn't that just rollback the flow file forever as 
this would not change? Same for line 246 and 285



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/image_to_tensor/image_to_tensor_def.rs:
##########
@@ -0,0 +1,206 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   https://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use super::{ColorFormat, ImageToTensor, ResizeFilter, ResizeMode, 
TensorShapeFormat};
+use crate::utils::per_channel_f32::PerChannelF32;
+use minifi_native::{
+    OutputAttribute, ProcessorDefinition, ProcessorInputRequirement, Property, 
PropertyDefinition,
+    Relationship, property_definitions,
+};
+
+pub(crate) const TARGET_WIDTH: Property<u32> = Property::new(
+    "Target width",
+    "Width in pixels the decoded image is resized to before normalisation and \
+                  inference.",
+);
+
+pub(crate) const TARGET_HEIGHT: Property<u32> = Property::new(
+    "Target height",
+    "Height in pixels the decoded image is resized to before normalisation and 
\
+                  inference.",
+);
+
+pub(crate) const RESIZE_FILTER: Property<ResizeFilter> = Property::new(
+    "Resize filter",
+    "Interpolation filter applied when resizing the decoded image. Nearest is 
fastest \
+                  but blocky; Bilinear is a good default; Bicubic and Lanczos3 
are higher-quality \
+                  but slower.",
+)
+.with_default(ResizeFilter::Bilinear.into_str());
+
+pub(crate) const RESIZE_MODE: Property<ResizeMode> = Property::new(
+    "Resize mode",
+    "How the source image is fitted into the target dimensions. 'Stretch' 
scales each \
+                  axis independently, distorting aspect ratio. 'Letterbox' 
preserves aspect ratio \
+                  and pads the remaining border with 'Letterbox pad value' 
(applied in normalised \
+                  output space).",
+)
+.with_default(ResizeMode::Stretch.into_str());
+
+pub(crate) const LETTERBOX_PAD_VALUE: Property<f32> = Property::new(
+    "Letterbox pad value",
+    "Value written for padding pixels when 'Resize mode' is 'Letterbox'. This 
is a \
+                  normalised value (post mean/std), so 0.0 corresponds to a 
neutral input for most \
+                  networks. Ignored when 'Resize mode' is 'Stretch'.",
+)
+.with_default("0.0");
+
+pub(crate) const COLOR_FORMAT: Property<ColorFormat> = Property::new(
+    "Color format",
+    "Colour space of the tensor fed to the model. RGB and BGR produce 
three-channel \
+                  tensors (channel order determined by the format); Grayscale 
produces a \
+                  single-channel luma tensor.",
+)
+.with_default(ColorFormat::Rgb.into_str());
+
+pub(crate) const TENSOR_SHAPE_FORMAT: Property<TensorShapeFormat> = 
Property::new(
+    "Tensor shape format",
+    "Memory layout of the tensor fed to the model. CHW (channels-first) is 
typical \
+                  for PyTorch/ONNX detectors. HWC (channels-last) matches 
TensorFlow/TFLite. \
+                  Ignored for Grayscale (always effectively 1xHxW).",
+)
+.with_default(TensorShapeFormat::Chw.into_str());
+
+pub(crate) const MEAN: Property<PerChannelF32> = Property::new(
+    "Mean",
+    "Mean subtracted from each pixel before dividing by 'Standard Deviation'. 
Accepts \
+                  either a single value (broadcast to all channels) or three 
comma-separated \
+                  values applied per channel in the order dictated by 'Color 
format'. Example: \
+                  '0.485, 0.456, 0.406' for ImageNet-style RGB normalisation.",
+)
+.with_default("0.0");
+
+pub(crate) const STD_DEV: Property<PerChannelF32> = Property::new(
+    "Standard Deviation",
+    "Divisor applied after subtracting 'Mean'. Accepts a single value 
(broadcast) or \
+                  three comma-separated values (per channel). Must be 
non-zero. Example: '255.0' \
+                  to scale u8 pixels into [0.0, 1.0]; '0.229, 0.224, 0.225' 
for ImageNet.",
+)
+.with_default("255.0");
+
+pub(crate) const PIXEL_DIVISOR: Property<f32> = Property::new(
+    "Pixel divisor",
+    "Divisor applied to raw u8 pixel values before subtracting 'Mean' and 
dividing \
+                  by 'Standard Deviation'. Defaults to 1.0 (mean/std 
interpreted in [0, 255] pixel \
+                  space, e.g. UltraFace's mean=127, std=128). Set to 255 to 
bring pixels into \
+                  [0.0, 1.0] first so ImageNet-style mean/std values like 
'0.485, 0.456, 0.406' / \
+                  '0.229, 0.224, 0.225' can be used directly, matching the 
PyTorch / torchvision / \
+                  ONNX MobileNet convention. Must be non-zero.",
+)
+.with_default("1.0");
+
+pub(super) const SUCCESS: Relationship = Relationship {
+    name: "success",
+    description: "The input image was decoded and converted to a tensor.",
+};
+
+pub(super) const FAILURE: Relationship = Relationship {
+    name: "failure",
+    description: "The input flow file could not be decoded as an image.",
+};
+
+pub(super) const TENSORS_LEN_ATTR: OutputAttribute = OutputAttribute {
+    name: "tensors.len",
+    relationships: &["success"],
+    description: "Number of tensors in the output FlowFile. Currently always 
'1'",
+};
+
+pub(super) const TENSOR_BYTES_ATTR: OutputAttribute = OutputAttribute {
+    name: "tensor.0.bytes",
+    relationships: &["success"],
+    description: "Byte length of output tensor.",
+};
+
+pub(super) const TENSOR_SHAPE_ATTR: OutputAttribute = OutputAttribute {
+    name: "tensor.0.shape",
+    relationships: &["success"],
+    description: "Comma-separated dimensions of the output tensor in the 
chosen layout, always \
+                  including a leading batch dimension of 1 (e.g. '1,3,224,224' 
for RGB CHW).",
+};
+
+pub(super) const TENSOR_DTYPE_ATTR: OutputAttribute = OutputAttribute {
+    name: "tensor.0.dtype",
+    relationships: &["success"],
+    description: "Element type of the values in the output tensor. Currently 
always 'F32'.",
+};
+
+pub(super) const IMG_ORG_HEIGHT_ATTR: OutputAttribute = OutputAttribute {
+    name: "image.original.height",
+    relationships: &["success"],
+    description: "The height of the original image before the resizing.",
+};
+
+pub(super) const IMG_ORG_WIDTH_ATTR: OutputAttribute = OutputAttribute {
+    name: "image.original.width",
+    relationships: &["success"],
+    description: "The width of the original image before the resizing.",
+};
+
+pub(super) const IMG_TRG_HEIGHT_ATTR: OutputAttribute = OutputAttribute {
+    name: "image.target.height",
+    relationships: &["success"],
+    description: "The height of the image after the resizing.",
+};
+
+pub(super) const IMG_TRG_WIDTH_ATTR: OutputAttribute = OutputAttribute {
+    name: "image.target.width",
+    relationships: &["success"],
+    description: "The width of the image after the resizing.",
+};
+
+pub(super) const IMG_RESIZE_MODE_ATTR: OutputAttribute = OutputAttribute {
+    name: "image.resize.mode",
+    relationships: &["success"],
+    description: "The resize mode ('Stretch' or 'Letterbox') applied to fit 
the image into the \
+                  target dimensions. Downstream processors such as 
FilterBoundingBoxes use this \
+                  to invert the coordinate mapping correctly.",
+};
+
+impl ProcessorDefinition for ImageToTensor {
+    const DESCRIPTION: &'static str = "Decodes an image from the flow file 
content and converts it into a normalised numeric \
+         tensor suitable for feeding into a downstream inference processor 
such as \
+         InvokeTractModel. Supports RGB / BGR / Grayscale, CHW / HWC layouts, 
stretch or \
+         letterbox resizing, and scalar or per-channel mean/std normalisation. 
The output payload \
+         is the raw little-endian f32 tensor; the 'tensors.len', 
'tensor.{i}.shape' and 'tensor.{i}.dtype' attributes \
+         describe its layout.";
+    const INPUT_REQUIREMENT: ProcessorInputRequirement = 
ProcessorInputRequirement::Required;
+    const SUPPORTS_DYNAMIC_PROPERTIES: bool = false;
+    const SUPPORTS_DYNAMIC_RELATIONSHIPS: bool = false;
+    const OUTPUT_ATTRIBUTES: &'static [OutputAttribute] = &[

Review Comment:
   TENSOR_BYTES_ATTR seems to be missing here



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/classify_output.rs:
##########
@@ -0,0 +1,476 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   https://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use crate::utils::score_activation::{ScoreActivation, SoftmaxTerms};
+use crate::utils::tensor_helpers::{deserialize_tensors, tensor_as_f32, 
tensor_shape};
+use classify_output_def::SUCCESS;
+pub(crate) use classify_output_def::{
+    CLASSIFY_OUTPUT_ATTRIBUTES, CONFIDENCE_THRESHOLD, LABEL_INDEX_OFFSET, 
LABELS_FILE_PATH,
+    OUTPUT_ATTRIBUTE_NAME, SCORE_ACTIVATION, SCORE_OUTPUT_INDEX, TOP_K,
+};
+use minifi_native::macros::ComponentIdentifier;
+use minifi_native::{
+    Content, FlowFileTransform, GetAttribute, GetId, GetProperty, InputStream, 
Logger, MinifiError,
+    ProcessError, RouteErrorExt, Schedule, TransformedFlowFile, warn,
+};
+use serde::Serialize;
+use std::path::Path;
+use tract::Tensor;
+
+mod classify_output_def;
+
+#[derive(Serialize, Clone, Debug, PartialEq)]
+struct Prediction {
+    class_id: usize,
+    confidence: f32,
+    #[serde(skip_serializing_if = "Option::is_none")]
+    class_name: Option<String>,
+}
+
+fn load_labels(path: &Path) -> Result<Vec<String>, MinifiError> {
+    let content = std::fs::read_to_string(path).map_err(|e| {
+        MinifiError::custom(format!("Failed to read labels file '{:?}': {}", 
path, e))
+    })?;
+    Ok(content
+        .lines()
+        .map(|line| line.trim().to_string())
+        .collect())
+}
+
+fn top_k(mut scored: Vec<(usize, f32)>, k: usize) -> Vec<(usize, f32)> {
+    scored.sort_by(|&(ai, a), &(bi, b)| b.total_cmp(&a).then(ai.cmp(&bi)));
+    scored.truncate(k);
+    scored
+}
+
+#[derive(ComponentIdentifier)]
+pub(crate) struct ClassifyOutput {
+    top_k: usize,
+    score_output_index: usize,
+    score_activation: ScoreActivation,
+    confidence_threshold: f32,
+    labels: Vec<String>,
+    label_index_offset: usize,
+}
+
+impl Schedule for ClassifyOutput {
+    fn schedule<Ctx: GetProperty, L: Logger>(
+        context: &Ctx,
+        _logger: &L,
+    ) -> Result<Self, MinifiError>
+    where
+        Self: Sized,
+    {
+        let top_k = context.get_property(&TOP_K)?;
+        if top_k == 0 {
+            return Err(MinifiError::validation("Top K must be >= 1"));
+        }
+        let score_output_index = context.get_property(&SCORE_OUTPUT_INDEX)?;
+        let score_activation = context.get_property(&SCORE_ACTIVATION)?;
+        let confidence_threshold = 
context.get_property(&CONFIDENCE_THRESHOLD)?;
+
+        let labels = match context.get_property(&LABELS_FILE_PATH)? {
+            Some(path) => load_labels(&path)?,
+            _ => Vec::new(),
+        };
+        let label_index_offset = context.get_property(&LABEL_INDEX_OFFSET)?;
+        if !labels.is_empty() && label_index_offset >= labels.len() {
+            return Err(MinifiError::validation(format!(
+                "Label index offset ({}) must be smaller than the number of 
labels ({})",
+                label_index_offset,
+                labels.len()
+            )));
+        }
+
+        Ok(Self {
+            top_k,
+            score_output_index,
+            score_activation,
+            confidence_threshold,
+            labels,
+            label_index_offset,
+        })
+    }
+}
+
+impl ClassifyOutput {
+    fn label_for(&self, class_id: usize) -> Option<String> {
+        self.labels
+            .get(class_id.checked_add(self.label_index_offset)?)
+            .cloned()
+    }
+
+    pub(crate) fn classify<'a, Context: GetProperty + GetAttribute + GetId, 
LoggerImpl: Logger>(
+        &self,
+        context: &Context,
+        logger: &LoggerImpl,
+        tensors: Vec<Tensor>,
+    ) -> Result<TransformedFlowFile<'a>, ProcessError> {
+        let score_floats =
+            tensor_as_f32(&tensors, 
self.score_output_index).route_err_to_failure()?;
+        if score_floats.is_empty() {
+            return Err(MinifiError::custom("Score tensor is empty; nothing to 
classify").into());

Review Comment:
   Same as my other comment in filter_bounding_boxes.rs, would this just 
rollback forever?



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/filter_bounding_boxes.rs:
##########
@@ -0,0 +1,593 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   https://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+mod filter_bounding_boxes_def;
+
+use crate::low_level_processors::image_to_tensor::ResizeMode;
+use crate::utils::bounding_box::BoundingBox;
+use crate::utils::dimensions::Dimensions;
+use crate::utils::score_activation::{ScoreActivation, SoftmaxTerms};
+use crate::utils::tensor_helpers::{deserialize_tensors, tensor_as_f32};
+use filter_bounding_boxes_def::SUCCESS;
+pub(crate) use filter_bounding_boxes_def::{
+    BACKGROUND_CLASS_INDEX, BOX_FORMAT, BOX_OUTPUT_INDEX, CLASS_OUTPUT_INDEX, 
CONFIDENCE_THRESHOLD,
+    IOU_THRESHOLD, OUTPUT_ATTRIBUTE_NAME, SCORE_ACTIVATION, SCORE_OUTPUT_INDEX,
+};
+use minifi_native::macros::{ComponentIdentifier, PropertyType};
+use minifi_native::{
+    Content, FlowFileTransform, GetAttribute, GetId, GetProperty, InputStream, 
Logger, MinifiError,
+    ProcessError, RouteErrorExt, Schedule, TransformedFlowFile, debug, trace,
+};
+use strum_macros::{Display, EnumString, IntoStaticStr, VariantNames};
+use tract::Tensor;
+
+#[derive(
+    Debug, Clone, Copy, PartialEq, Display, EnumString, VariantNames, 
IntoStaticStr, PropertyType,
+)]
+#[strum(serialize_all = "PascalCase", const_into_str)]
+pub(crate) enum BoxFormat {
+    /// `[x_min, y_min, x_max, y_max]` — SSD, MobileNet-SSD, most PyTorch 
models.
+    Xyxy,
+    /// `[y_min, x_min, y_max, x_max]` — TensorFlow Object Detection API.
+    Yxyx,
+    /// `[cx, cy, w, h]` — YOLOv3/5/8 raw output (center + size).
+    Cxcywh,
+}
+
+/// Convert the four floats at `box_floats[offset..offset+4]` into a canonical
+/// `(x_min, y_min, x_max, y_max)` tuple, regardless of the source layout.
+fn decode_box(box_floats: &[f32], offset: usize, format: BoxFormat) -> (f32, 
f32, f32, f32) {
+    let a = box_floats[offset];
+    let b = box_floats[offset + 1];
+    let c = box_floats[offset + 2];
+    let d = box_floats[offset + 3];
+    match format {
+        BoxFormat::Xyxy => (a, b, c, d),
+        BoxFormat::Yxyx => (b, a, d, c),
+        BoxFormat::Cxcywh => {
+            let (cx, cy, w, h) = (a, b, c, d);
+            (cx - w / 2.0, cy - h / 2.0, cx + w / 2.0, cy + h / 2.0)
+        }
+    }
+}
+
+struct ScoredClass {
+    class_id: usize,
+    confidence: f32,
+}
+
+fn score_box(
+    logits: &[f32],
+    activation: ScoreActivation,
+    background_class_index: Option<usize>,
+) -> ScoredClass {
+    let num_classes = logits.len();
+
+    let best_valid = logits
+        .iter()
+        .enumerate()
+        .filter(|&(_, &logit)| logit.is_finite())
+        .filter(|&(id, _)| match background_class_index {
+            Some(bg_idx) => !(num_classes > 1 && id == bg_idx),
+            None => true,
+        })
+        .max_by(|a, b| a.1.total_cmp(b.1));
+
+    let (class_id, &best_logit) = match best_valid {
+        Some(val) => val,
+        None => {
+            return ScoredClass {
+                class_id: 0,
+                confidence: f32::NEG_INFINITY,
+            };
+        }
+    };
+
+    let confidence = activation.confidence(best_logit, 
SoftmaxTerms::over(logits.iter().copied()));
+
+    ScoredClass {
+        class_id,
+        confidence,
+    }
+}
+
+#[derive(ComponentIdentifier)]
+pub(crate) struct FilterBoundingBoxes {
+    confidence_threshold: f32,
+    iou_threshold: f32,
+    score_output_index: usize,
+    box_output_index: usize,
+    box_format: BoxFormat,
+    score_activation: ScoreActivation,
+    background_class_index: Option<usize>,
+    class_output_index: Option<usize>,
+}
+
+impl Schedule for FilterBoundingBoxes {
+    fn schedule<Ctx: GetProperty, L: Logger>(
+        context: &Ctx,
+        _logger: &L,
+    ) -> Result<Self, MinifiError> {
+        let confidence_threshold = 
context.get_property(&CONFIDENCE_THRESHOLD)?;
+        let iou_threshold = context.get_property(&IOU_THRESHOLD)?;
+        let score_output_index = context.get_property(&SCORE_OUTPUT_INDEX)?;
+        let box_output_index = context.get_property(&BOX_OUTPUT_INDEX)?;
+        let box_format = context.get_property(&BOX_FORMAT)?;
+        let score_activation = context.get_property(&SCORE_ACTIVATION)?;
+        let background_class_index = 
context.get_property(&BACKGROUND_CLASS_INDEX)?;
+        let class_output_index = context.get_property(&CLASS_OUTPUT_INDEX)?;
+
+        Ok(Self {
+            confidence_threshold,
+            iou_threshold,
+            score_output_index,
+            box_output_index,
+            box_format,
+            score_activation,
+            background_class_index,
+            class_output_index,
+        })
+    }
+}
+
+impl FilterBoundingBoxes {
+    fn result_via_output_attribute<'a, Context: GetProperty>(
+        &self,
+        context: &Context,
+        filtered_boxes: Vec<BoundingBox>,
+    ) -> Result<TransformedFlowFile<'a>, MinifiError> {
+        let output_attr = context.get_property(&OUTPUT_ATTRIBUTE_NAME)?;
+        let content = if output_attr.is_some() {
+            None
+        } else {
+            Some(Content::Buffer(
+                
serde_json::to_vec(&filtered_boxes).map_err(MinifiError::other)?,
+            ))
+        };
+
+        let mut transformed = TransformedFlowFile::new(&SUCCESS, content)
+            .with_attribute("object.count", filtered_boxes.len().to_string());
+        if let Some(attr) = output_attr {
+            transformed = transformed.with_attribute(
+                attr,
+                
serde_json::to_string(&filtered_boxes).map_err(MinifiError::other)?,
+            )
+        } else {
+            transformed = transformed.with_attribute("mime.type", 
"application/json");
+        }
+        Ok(transformed)
+    }
+
+    pub(crate) fn filter<'a, Context: GetProperty, LoggerImpl: Logger>(
+        &self,
+        context: &Context,
+        logger: &LoggerImpl,
+        tensors: Vec<Tensor>,
+        orig_dim: Dimensions,
+        target_dim: Dimensions,
+        resize_mode: ResizeMode,
+    ) -> Result<TransformedFlowFile<'a>, ProcessError> {
+        let score_floats =
+            tensor_as_f32(&tensors, 
self.score_output_index).route_err_to_failure()?;
+        let box_floats = tensor_as_f32(&tensors, 
self.box_output_index).route_err_to_failure()?;
+
+        let (scale_x, scale_y, pad_x, pad_y) = match resize_mode {
+            ResizeMode::Letterbox => {
+                let geometry = orig_dim.letterbox_into(target_dim);
+                (
+                    geometry.scale,
+                    geometry.scale,
+                    geometry.pad_x as f32,
+                    geometry.pad_y as f32,
+                )
+            }
+            ResizeMode::Stretch => (
+                target_dim.width / orig_dim.width,
+                target_dim.height / orig_dim.height,
+                0.0,
+                0.0,
+            ),
+        };
+
+        if !box_floats.len().is_multiple_of(4) {
+            return Err(MinifiError::custom(
+                "Box tensor byte length is not a multiple of 16 (4 f32 per 
box)",
+            )
+            .into());
+        }
+        let num_boxes = box_floats.len() / 4;
+        if num_boxes == 0 {
+            debug!(logger, "No boxes to filter; emitting empty array");
+            return self
+                .result_via_output_attribute(context, vec![])
+                .route_err_to_failure();
+        }
+
+        let make_box = |i: usize, class_id: usize, confidence: f32| -> 
BoundingBox {
+            let (raw_x_min, raw_y_min, raw_x_max, raw_y_max) =
+                decode_box(&box_floats, i * 4, self.box_format);
+            let true_x_min = (((raw_x_min * target_dim.width) - pad_x) / 
scale_x) / orig_dim.width;
+            let true_y_min =
+                (((raw_y_min * target_dim.height) - pad_y) / scale_y) / 
orig_dim.height;
+            let true_x_max = (((raw_x_max * target_dim.width) - pad_x) / 
scale_x) / orig_dim.width;
+            let true_y_max =
+                (((raw_y_max * target_dim.height) - pad_y) / scale_y) / 
orig_dim.height;
+            BoundingBox {
+                class_id,
+                confidence,
+                x_min: true_x_min.clamp(0.0, 1.0),
+                y_min: true_y_min.clamp(0.0, 1.0),
+                x_max: true_x_max.clamp(0.0, 1.0),
+                y_max: true_y_max.clamp(0.0, 1.0),
+            }
+        };
+
+        let mut valid_boxes = Vec::new();
+
+        match self.class_output_index {
+            // Separate class-id tensor: one score and one class id per box
+            Some(class_index) => {
+                let class_floats = tensor_as_f32(&tensors, 
class_index).route_err_to_failure()?;
+                if score_floats.len() != num_boxes || class_floats.len() != 
num_boxes {
+                    return Err(MinifiError::custom(format!(
+                        "'Class output index' mode expects one score and one 
class id per box \
+                         (num_boxes={}, scores={}, classes={})",
+                        num_boxes,
+                        score_floats.len(),
+                        class_floats.len()
+                    ))
+                    .into());
+                }
+                trace!(
+                    logger,
+                    "Filtering {} boxes with separate class-id tensor 
(activation={:?}, \
+                     box_format={:?})...",
+                    num_boxes,
+                    self.score_activation,
+                    self.box_format
+                );
+                for i in 0..num_boxes {
+                    let confidence = 
self.score_activation.confidence_of_scalar(score_floats[i]);
+                    if confidence < self.confidence_threshold {
+                        continue;
+                    }
+                    let raw_class = class_floats[i];
+                    if raw_class < 0.0 {
+                        continue;
+                    }
+                    if !confidence.is_finite() || !raw_class.is_finite() {
+                        continue;
+                    }
+                    let class_id = raw_class.round() as usize;
+                    if self.background_class_index == Some(class_id) {
+                        continue;
+                    }
+                    valid_boxes.push(make_box(i, class_id, confidence));
+                }
+            }
+            // Per-class score matrix: argmax over classes per box.
+            None => {
+                if !score_floats.len().is_multiple_of(num_boxes) {
+                    return Err(MinifiError::custom(format!(
+                        "Scores length ({}) not divisible by number of boxes 
({})",
+                        score_floats.len(),
+                        num_boxes
+                    ))
+                    .into());
+                }
+                let num_classes = score_floats.len() / num_boxes;
+                trace!(
+                    logger,
+                    "Filtering {} boxes across {} potential classes 
(activation={:?}, \
+                     box_format={:?})...",
+                    num_boxes,
+                    num_classes,
+                    self.score_activation,
+                    self.box_format
+                );
+                for i in 0..num_boxes {
+                    let logits = &score_floats[i * num_classes..(i + 1) * 
num_classes];
+                    let scored =
+                        score_box(logits, self.score_activation, 
self.background_class_index);
+                    if scored.confidence >= self.confidence_threshold {
+                        valid_boxes.push(make_box(i, scored.class_id, 
scored.confidence));
+                    }
+                }
+            }
+        }
+
+        trace!(
+            logger,
+            "Found {} boxes exceeding the {} threshold.",
+            valid_boxes.len(),
+            self.confidence_threshold
+        );
+
+        let filtered_boxes =
+            BoundingBox::apply_non_maximum_suppression(valid_boxes, 
self.iou_threshold);
+
+        self.result_via_output_attribute(context, filtered_boxes)
+            .route_err_to_failure()
+    }
+}
+
+fn resize_mode_from_attributes<Context: GetAttribute>(context: &Context) -> 
ResizeMode {
+    context
+        .get_attribute("image.resize.mode")
+        .ok()
+        .flatten()
+        .and_then(|raw| raw.parse::<ResizeMode>().ok())
+        .unwrap_or(ResizeMode::Letterbox)

Review Comment:
   I am not sure if this is intentional, but in image_to_tensor_def.rs we 
define the default for RESIZE_MODE to ResizeMode::Stretch, but here the default 
is ResizeMode::Letterbox



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to