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


##########
minifi_rust/extensions/minifi_tensor/src/utils/tensor_helpers.rs:
##########
@@ -0,0 +1,210 @@
+use image::{DynamicImage, ImageResult};
+use minifi_native::{GetAttribute, InputStream, MinifiError};
+use strum_macros::{Display, EnumString};
+use tract::__ndarray_interop::TensorInterface;
+use tract::Tensor;
+use tract::prelude::DatumType;
+
+tract::impl_ndarray_interop!();
+
+fn parse_tensor_shape<Context: GetAttribute>(
+    context: &Context,
+    id: usize,
+) -> Result<Vec<usize>, MinifiError> {
+    let shape_str = context.get_required_attribute(&format!("tensor.{}.shape", 
id))?;
+
+    if shape_str.trim().is_empty() {
+        return Ok(Vec::new());
+    }
+
+    let shape = shape_str
+        .split(',')
+        .map(|s| s.trim().parse::<usize>())
+        .collect::<Result<Vec<usize>, _>>()?;
+
+    Ok(shape)
+}
+
+#[derive(Debug, Clone, Copy, PartialEq, Display, EnumString)]
+#[strum(serialize_all = "PascalCase", const_into_str)]
+pub(crate) enum MinifiDatumType {
+    F32,
+}
+
+impl From<MinifiDatumType> for DatumType {
+    fn from(value: MinifiDatumType) -> Self {
+        match value {
+            MinifiDatumType::F32 => DatumType::F32,
+        }
+    }
+}
+
+fn numeric_datum_type_from_str(s: &str) -> Option<DatumType> {
+    Some(match s {
+        "U8" => DatumType::U8,
+        "U16" => DatumType::U16,
+        "U32" => DatumType::U32,
+        "U64" => DatumType::U64,
+        "I8" => DatumType::I8,
+        "I16" => DatumType::I16,
+        "I32" => DatumType::I32,
+        "I64" => DatumType::I64,
+        "F16" => DatumType::F16,
+        "F32" => DatumType::F32,
+        "F64" => DatumType::F64,
+        _ => return None,
+    })
+}
+
+fn parse_tensor_dtype<Context: GetAttribute>(
+    context: &Context,
+    id: usize,
+) -> Result<DatumType, MinifiError> {
+    let dtype_str = context.get_required_attribute(&format!("tensor.{}.dtype", 
id))?;
+    numeric_datum_type_from_str(&dtype_str).ok_or_else(|| {
+        MinifiError::custom(format!(
+            "Unsupported tensor.{}.dtype '{}': only numeric tensors can be 
read",
+            id, dtype_str
+        ))
+    })
+}
+
+pub(crate) fn deserialize_tensors<Context: GetAttribute>(
+    context: &Context,
+    input_stream: &mut dyn InputStream,
+) -> Result<Vec<Tensor>, MinifiError> {
+    let mut result = vec![];
+
+    let mut flow_file_contents = Vec::new();
+    input_stream.read_to_end(&mut flow_file_contents)?;
+    let number_of_tensors = context
+        .get_required_attribute("tensors.len")?
+        .parse::<usize>()?;
+
+    let mut cursor = 0usize;
+    for i in 0..number_of_tensors {
+        let tensor_len = context
+            .get_required_attribute(&format!("tensor.{}.bytes", i))?
+            .parse::<usize>()?;
+        let tensor_shape = parse_tensor_shape(context, i)?;
+        let tensor_dtype = parse_tensor_dtype(context, i)?;
+        if cursor + tensor_len > flow_file_contents.len() {
+            return Err(MinifiError::custom(
+                "FlowFile contents are not in sync with tensor attributes",
+            ));
+        }
+        let tensor_data = &flow_file_contents[cursor..cursor + tensor_len];
+        result.push(Tensor::from_bytes(
+            tensor_dtype,
+            &tensor_shape,
+            tensor_data,
+        )?);
+        cursor += tensor_len;
+    }
+
+    if cursor != flow_file_contents.len() {
+        Err(MinifiError::custom(
+            "FlowFile contents are not in sync with tensor attributes",
+        ))
+    } else {
+        Ok(result)
+    }
+}
+
+pub(crate) fn tensor_as_f32(tensors: &[Tensor], index: usize) -> 
Result<Vec<f32>, MinifiError> {
+    let tensor = tensors
+        .get(index)
+        .ok_or(MinifiError::custom("Invalid shape of tensors"))?;
+    let casted = tensor.convert_to(DatumType::F32)?;
+    Ok(casted.as_slice::<f32>()?.to_vec())
+}

Review Comment:
   If I understand correctly, all processors can only work with F32. Is it 
feasible to add support for other formats later, possibly to leverage hardware 
support for F16, I8, and other smaller formats, for faster execution?



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/classify_output.rs:
##########
@@ -0,0 +1,429 @@
+// 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;
+use crate::utils::tensor_helpers::{deserialize_tensors, tensor_as_f32, 
tensor_shape};
+use classify_output_def::SUCCESS;
+pub(crate) use classify_output_def::{
+    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,
+};
+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_end().to_string())

Review Comment:
   minor, but do we allow whitespace characters other than a line break as part 
of the label? If yes, then we should only trim the line ending. If not, then we 
should trim the beginning and the end too. 



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/image_to_tensor.rs:
##########
@@ -0,0 +1,625 @@
+// 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.
+pub(crate) mod image_to_tensor_def;
+
+use 
crate::low_level_processors::image_to_tensor::image_to_tensor_def::TENSOR_BYTES_ATTR;
+use crate::utils::dimensions::Dimensions;
+use crate::utils::per_channel_f32::PerChannelF32;
+use crate::utils::tensor_helpers::{MinifiDatumType, load_as_image};
+pub(crate) use image_to_tensor_def::{
+    COLOR_FORMAT, LETTERBOX_PAD_VALUE, MEAN, PIXEL_DIVISOR, RESIZE_FILTER, 
RESIZE_MODE, STD_DEV,
+    TARGET_HEIGHT, TARGET_WIDTH, TENSOR_SHAPE_FORMAT,
+};
+use image_to_tensor_def::{
+    IMG_ORG_HEIGHT_ATTR, IMG_ORG_WIDTH_ATTR, IMG_RESIZE_MODE_ATTR, SUCCESS, 
TENSOR_DTYPE_ATTR,
+    TENSOR_SHAPE_ATTR,
+};
+use image_to_tensor_def::{IMG_TRG_HEIGHT_ATTR, IMG_TRG_WIDTH_ATTR, 
TENSORS_LEN_ATTR};
+use minifi_native::macros::{ComponentIdentifier, PropertyType};
+use minifi_native::{
+    FlowFileTransform, GetAttribute, GetControllerService, GetId, GetProperty, 
InputStream, Logger,
+    MinifiError, ProcessError, RouteErrorExt, Schedule, TransformedFlowFile,
+};
+use strum_macros::{Display, EnumString, IntoStaticStr, VariantNames};
+use tract::Tensor;
+
+tract::impl_ndarray_interop!();
+
+#[derive(
+    Debug, Clone, Copy, PartialEq, Display, EnumString, VariantNames, 
IntoStaticStr, PropertyType,
+)]
+#[strum(serialize_all = "PascalCase", const_into_str)]
+pub(crate) enum ResizeFilter {
+    Nearest,
+    Bilinear,
+    Bicubic,
+    Lanczos3,
+}
+
+impl From<ResizeFilter> for image::imageops::FilterType {
+    fn from(filter: ResizeFilter) -> Self {
+        match filter {
+            ResizeFilter::Nearest => image::imageops::FilterType::Nearest,
+            ResizeFilter::Bilinear => image::imageops::FilterType::Triangle,
+            ResizeFilter::Bicubic => image::imageops::FilterType::CatmullRom,
+            ResizeFilter::Lanczos3 => image::imageops::FilterType::Lanczos3,
+        }
+    }
+}
+
+#[derive(
+    Debug, Clone, Copy, PartialEq, Display, EnumString, VariantNames, 
IntoStaticStr, PropertyType,
+)]
+#[strum(serialize_all = "UPPERCASE", const_into_str)]
+pub(crate) enum ColorFormat {
+    Rgb,
+    Bgr,
+    Grayscale,
+}
+
+#[derive(
+    Debug, Clone, Copy, PartialEq, Display, EnumString, VariantNames, 
IntoStaticStr, PropertyType,
+)]
+#[strum(serialize_all = "UPPERCASE", const_into_str)]
+pub(crate) enum TensorShapeFormat {
+    Chw, // channel, height, width
+    Hwc, // height, width, channel
+}
+
+#[derive(
+    Debug, Clone, Copy, PartialEq, Display, EnumString, VariantNames, 
IntoStaticStr, PropertyType,
+)]
+#[strum(serialize_all = "PascalCase", const_into_str)]
+pub(crate) enum ResizeMode {
+    Stretch,
+    Letterbox,
+}
+
+#[derive(ComponentIdentifier)]
+pub(crate) struct ImageToTensor {
+    target_width: u32,
+    target_height: u32,
+    resize_filter: ResizeFilter,
+    resize_mode: ResizeMode,
+    color_format: ColorFormat,
+    tensor_shape_format: TensorShapeFormat,
+    mean: PerChannelF32,
+    std_dev: PerChannelF32,
+    pixel_divisor: f32,
+    letterbox_pad_value: f32,
+}
+
+impl Schedule for ImageToTensor {
+    fn schedule<Ctx: GetProperty, L: Logger>(
+        context: &Ctx,
+        _logger: &L,
+    ) -> Result<Self, MinifiError>
+    where
+        Self: Sized,
+    {
+        let target_width = context.get_property(&TARGET_WIDTH)?;
+        let target_height = context.get_property(&TARGET_HEIGHT)?;
+        let resize_filter = context.get_property(&RESIZE_FILTER)?;
+        let resize_mode = context.get_property(&RESIZE_MODE)?;
+        let color_format = context.get_property(&COLOR_FORMAT)?;
+        let tensor_shape_format = context.get_property(&TENSOR_SHAPE_FORMAT)?;
+        let mean = context.get_property(&MEAN)?;
+        let std_dev = context.get_property(&STD_DEV)?;
+        if std_dev.contains_zero() {
+            return Err(MinifiError::validation(
+                "Standard Deviation components must be non-zero",
+            ));
+        }
+        let pixel_divisor = context.get_property(&PIXEL_DIVISOR)?;
+        if pixel_divisor == 0.0 {
+            return Err(MinifiError::validation("Pixel divisor must be 
non-zero"));
+        }
+        let letterbox_pad_value = context.get_property(&LETTERBOX_PAD_VALUE)?;
+
+        Ok(Self {
+            target_width,
+            target_height,
+            resize_filter,
+            resize_mode,
+            color_format,
+            tensor_shape_format,
+            mean,
+            std_dev,
+            pixel_divisor,
+            letterbox_pad_value,
+        })
+    }
+}
+
+struct MaskedRgbImage {
+    img: image::RgbImage,
+    mask: Vec<bool>,
+}
+
+impl ImageToTensor {
+    fn stretch_resize(&self, img: image::DynamicImage) -> MaskedRgbImage {
+        let resized = img
+            .resize_exact(
+                self.target_width,
+                self.target_height,
+                self.resize_filter.into(),
+            )
+            .to_rgb8();
+        let mask = vec![true; (self.target_width * self.target_height) as 
usize];
+        MaskedRgbImage { img: resized, mask }
+    }
+
+    fn letterbox_resize(&self, img: image::DynamicImage) -> MaskedRgbImage {
+        let (src_w, src_h) = (img.width() as f32, img.height() as f32);
+        let scale = (self.target_width as f32 / src_w).min(self.target_height 
as f32 / src_h);
+        let new_w = (src_w * scale).round().max(1.0) as u32;
+        let new_h = (src_h * scale).round().max(1.0) as u32;
+        let scaled = img
+            .resize_exact(new_w, new_h, self.resize_filter.into())
+            .to_rgb8();
+
+        let pad_x = (self.target_width - new_w) / 2;
+        let pad_y = (self.target_height - new_h) / 2;
+
+        let mut canvas = image::RgbImage::from_pixel(
+            self.target_width,
+            self.target_height,
+            image::Rgb([0, 0, 0]),
+        );
+        image::imageops::overlay(&mut canvas, &scaled, pad_x as i64, pad_y as 
i64);
+
+        // Mask to track which pixel is part of source and which is padding
+        let mut mask = vec![false; (self.target_width * self.target_height) as 
usize];
+        for y in pad_y..(new_h + pad_y) {
+            for x in pad_x..(new_w + pad_x) {
+                mask[(y * self.target_width + x) as usize] = true;
+            }
+        }
+        MaskedRgbImage { img: canvas, mask }
+    }
+
+    fn resize_rgb(&self, img: image::DynamicImage) -> MaskedRgbImage {
+        match self.resize_mode {
+            ResizeMode::Stretch => self.stretch_resize(img),
+            ResizeMode::Letterbox => self.letterbox_resize(img),
+        }
+    }
+
+    pub fn tensor_bytes(&self, img: image::DynamicImage) -> Vec<u8> {
+        let num_channels: usize = match self.color_format {
+            ColorFormat::Grayscale => 1,
+            _ => 3,
+        };
+        let total_pixels = (self.target_width * self.target_height) as usize;
+        let mut tensor_bytes = Vec::with_capacity(total_pixels * num_channels 
* 4);
+
+        let masked_img = self.resize_rgb(img);

Review Comment:
   Would it make sense to special-case when image dimensions match the target 
dimensions? In my tests, I've just specified the width/height of my test images 
as the target, and it matches what the image classification model expects, so 
no stretching/letterboxing is necessary.



##########
minifi_rust/extensions/minifi_tensor/src/low_level_processors/classify_output.rs:
##########
@@ -0,0 +1,429 @@
+// 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;
+use crate::utils::tensor_helpers::{deserialize_tensors, tensor_as_f32, 
tensor_shape};
+use classify_output_def::SUCCESS;
+pub(crate) use classify_output_def::{
+    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,
+};
+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_end().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
+}

Review Comment:
   possible future optimization idea: partial sort
   ignore for now



-- 
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