This is an automated email from the ASF dual-hosted git repository.

martinzink pushed a commit to branch minifi_rust_impr_3
in repository https://gitbox.apache.org/repos/asf/nifi-minifi-cpp.git

commit 8e0aeb0bea97f1e248f565cb762faa2c4b91af1d
Author: Martin Zink <[email protected]>
AuthorDate: Mon Sep 14 16:57:00 2026 +0200

    MINIFICPP-2901 Rust route_to_failure by default, rollback explicit
---
 minifi_rust/CMakeLists.txt                         | 28 +++++--
 .../src/processors/asciify_german.rs               |  5 +-
 .../src/processors/kamikaze_processor.rs           |  6 +-
 .../src/processors/kamikaze_processor/tests.rs     |  2 +-
 .../src/processors/lorem_ipsum_cs_user.rs          | 10 ++-
 .../src/processors/put_file.rs                     |  4 +-
 minifi_rust/minifi_native/src/api/errors.rs        | 79 ++++++++++++++----
 .../api/processor_wrappers/complex_processor.rs    |  4 +-
 .../src/api/processor_wrappers/flow_file_source.rs | 22 ++---
 .../flow_file_stream_transform.rs                  | 63 ++++++++-------
 .../api/processor_wrappers/flow_file_transform.rs  | 94 ++++++++++++----------
 .../src/c_ffi/c_ffi_processor_definition.rs        |  2 +-
 minifi_rust/minifi_native/src/lib.rs               |  2 +-
 13 files changed, 204 insertions(+), 117 deletions(-)

diff --git a/minifi_rust/CMakeLists.txt b/minifi_rust/CMakeLists.txt
index 63b8a407e..0ecc1441d 100644
--- a/minifi_rust/CMakeLists.txt
+++ b/minifi_rust/CMakeLists.txt
@@ -41,8 +41,26 @@ endif()
 
 include(CTest)
 
-add_test(
-        NAME cargo_tests
-        COMMAND cargo test
-        WORKING_DIRECTORY ${CMAKE_CURRENT_SOURCE_DIR}
-)
+find_program(CARGO_NEXTEST_EXECUTABLE cargo-nextest)
+
+if (CARGO_NEXTEST_EXECUTABLE)
+    message(STATUS "Found cargo-nextest, using it for Rust tests.")
+    add_test(
+            NAME cargo_tests
+            COMMAND cargo nextest run
+            WORKING_DIRECTORY ${CMAKE_CURRENT_SOURCE_DIR}
+    )
+    add_test(
+            NAME cargo_doctests
+            COMMAND cargo test --doc
+            WORKING_DIRECTORY ${CMAKE_CURRENT_SOURCE_DIR}
+    )
+else()
+    message(WARNING "cargo-nextest not found. Falling back to standard cargo 
test. "
+            "For faster test execution, install it via: cargo install 
cargo-nextest --locked")
+    add_test(
+            NAME cargo_tests
+            COMMAND cargo test
+            WORKING_DIRECTORY ${CMAKE_CURRENT_SOURCE_DIR}
+    )
+endif()
diff --git 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/asciify_german.rs 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/asciify_german.rs
index 6d67c7d05..4f1f147f0 100644
--- 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/asciify_german.rs
+++ 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/asciify_german.rs
@@ -21,7 +21,7 @@ use crate::processors::asciify_german::relationships::FAILURE;
 use minifi_native::macros::ComponentIdentifier;
 use minifi_native::{
     FlowFileStreamTransform, GetProperty, InputStream, Logger, MinifiError, 
OutputStream,
-    ProcessError, RouteErrorExt, Schedule, TransformStreamResult,
+    ProcessError, Schedule, TransformStreamResult,
 };
 
 mod relationships;
@@ -56,8 +56,7 @@ impl FlowFileStreamTransform for AsciifyGerman {
                 0xC3 => {
                     let mut next = [0u8; 1];
                     if input_stream.read(&mut next)? == 0 {
-                        Err(MinifiError::custom("Truncated multi-byte sequence 
at EOF"))
-                            .route_err_to_failure()?
+                        Err(MinifiError::custom("Truncated multi-byte sequence 
at EOF"))?
                     }
                     match next[0] {
                         0xA4 => output_stream.write_all(b"ae")?, // รค
diff --git 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor.rs
 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor.rs
index 0faa04f18..4c05880f8 100644
--- 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor.rs
+++ 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor.rs
@@ -88,9 +88,9 @@ impl Trigger for KamikazeProcessorRs {
         L: Logger,
     {
         match self.trigger_behaviour {
-            KamikazeBehaviour::ReturnErr => {
-                Err(MinifiError::custom("it was designed to fail in 
trigger").into())
-            }
+            KamikazeBehaviour::ReturnErr => 
Err(ProcessError::Rollback(MinifiError::custom(
+                "it was designed to fail in trigger",
+            ))),
             KamikazeBehaviour::ReturnOk => Ok(OnTriggerResult::Ok),
             KamikazeBehaviour::Panic => {
                 panic!("KamikazeProcessor::trigger panic")
diff --git 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor/tests.rs
 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor/tests.rs
index 2b7f52166..0520d0285 100644
--- 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor/tests.rs
+++ 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/kamikaze_processor/tests.rs
@@ -77,7 +77,7 @@ fn on_trigger_err() {
     let mut session = MockProcessSession::new();
     assert!(matches!(
         processor.trigger(&mut context, &mut session, &MockLogger::new()),
-        Err(ProcessError::Fatal(MinifiError::CustomError(_)))
+        Err(ProcessError::Rollback(MinifiError::CustomError(_)))
     ));
 }
 
diff --git 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/lorem_ipsum_cs_user.rs
 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/lorem_ipsum_cs_user.rs
index 4d32739e2..61b2d4b99 100644
--- 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/lorem_ipsum_cs_user.rs
+++ 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/lorem_ipsum_cs_user.rs
@@ -25,7 +25,7 @@ use 
crate::processors::lorem_ipsum_cs_user::relationships::SUCCESS;
 use minifi_native::macros::{ComponentIdentifier, PropertyType};
 use minifi_native::{
     Content, FlowFileSource, GeneratedFlowFile, GetControllerService, 
GetProperty, Logger,
-    MinifiError, ProcessError, Schedule, trace,
+    MinifiError, ProcessError, ProcessErrorExt, Schedule, trace,
 };
 use strum_macros::{Display, EnumString, IntoStaticStr, VariantNames};
 
@@ -60,12 +60,16 @@ impl FlowFileSource for LoremIpsumCSUser {
         logger: &LoggerImpl,
     ) -> Result<Vec<GeneratedFlowFile<'a>>, ProcessError> {
         trace!(logger, "generate call {:?}", self);
-        let dummy_controller_service = 
context.get_controller_service(&DUMMY_CONTROLLER_SERVICE)?;
+        let dummy_controller_service = context
+            .get_controller_service(&DUMMY_CONTROLLER_SERVICE)
+            .rollback_err()?;
         trace!(
             logger,
             "optional dummy controller service: {:?}", dummy_controller_service
         );
-        let controller_service = 
context.get_controller_service(&CONTROLLER_SERVICE)?;
+        let controller_service = context
+            .get_controller_service(&CONTROLLER_SERVICE)
+            .rollback_err()?;
         match self.write_method {
             WriteMethod::Buffer => {
                 let generated_flow_file = GeneratedFlowFile::new(
diff --git 
a/minifi_rust/extensions/minifi_rs_playground/src/processors/put_file.rs 
b/minifi_rust/extensions/minifi_rs_playground/src/processors/put_file.rs
index 1dd63ec2d..a977d76c5 100644
--- a/minifi_rust/extensions/minifi_rs_playground/src/processors/put_file.rs
+++ b/minifi_rust/extensions/minifi_rs_playground/src/processors/put_file.rs
@@ -22,7 +22,7 @@ use 
crate::processors::put_file::unix_permissions::PutFileUnixPermissions;
 use minifi_native::macros::{ComponentIdentifier, PropertyType};
 use minifi_native::{
     FlowFileTransform, GetAttribute, GetControllerService, GetId, GetProperty, 
InputStream, Logger,
-    MinifiError, ProcessError, RouteErrorExt, Schedule, TransformedFlowFile, 
trace, warn,
+    MinifiError, ProcessError, Schedule, TransformedFlowFile, trace, warn,
 };
 use std::path::{Path, PathBuf};
 use strum_macros::{Display, EnumString, IntoStaticStr, VariantNames};
@@ -174,7 +174,7 @@ impl FlowFileTransform for PutFileRs {
     ) -> Result<TransformedFlowFile<'a>, ProcessError> {
         trace!(logger, "on_trigger: {:?}", self);
 
-        let destination_path = 
Self::get_destination_path(context).route_err_to_failure()?;
+        let destination_path = Self::get_destination_path(context)?;
 
         if self.directory_is_full(&destination_path) {
             warn!(logger, "Directory is full");
diff --git a/minifi_rust/minifi_native/src/api/errors.rs 
b/minifi_rust/minifi_native/src/api/errors.rs
index 16c652a7c..3535ac44d 100644
--- a/minifi_rust/minifi_native/src/api/errors.rs
+++ b/minifi_rust/minifi_native/src/api/errors.rs
@@ -58,7 +58,7 @@ impl Error for RouteError {}
 #[derive(Debug)]
 pub enum ProcessError {
     Route(RouteError),
-    Fatal(MinifiError),
+    Rollback(MinifiError),
 }
 
 impl From<RouteError> for ProcessError {
@@ -69,23 +69,31 @@ impl From<RouteError> for ProcessError {
 
 impl From<MinifiError> for ProcessError {
     fn from(err: MinifiError) -> Self {
-        ProcessError::Fatal(err)
+        ProcessError::Route(RouteError {
+            relationship: "failure",
+            source: Box::new(err),
+            log_level: LogLevel::Warn,
+        })
     }
 }
 
-macro_rules! process_error_from_fatal {
+macro_rules! process_error_route_to_failure {
     ($($t:ty),* $(,)?) => {
         $(
             impl From<$t> for ProcessError {
                 fn from(err: $t) -> Self {
-                    ProcessError::Fatal(MinifiError::from(err))
+                    ProcessError::Route(RouteError {
+                        relationship: "failure",
+                        source: Box::new(err),
+                        log_level: LogLevel::Warn,
+                    })
                 }
             }
         )*
     };
 }
 
-process_error_from_fatal!(
+process_error_route_to_failure!(
     std::io::Error,
     strum::ParseError,
     ParseBoolError,
@@ -101,22 +109,24 @@ impl fmt::Display for ProcessError {
     fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
         match self {
             ProcessError::Route(err) => write!(f, "{}", err),
-            ProcessError::Fatal(err) => write!(f, "{}", err),
+            ProcessError::Rollback(err) => write!(f, "{}", err),
         }
     }
 }
 
 impl Error for ProcessError {}
 
-pub trait RouteErrorExt<T> {
+pub trait ProcessErrorExt<T> {
     fn route_err(self, rel: &Relationship, level: LogLevel) -> Result<T, 
ProcessError>;
 
     fn route_to(self, relationship: &'static str, level: LogLevel) -> 
Result<T, ProcessError>;
 
     fn route_err_to_failure(self) -> Result<T, ProcessError>;
+
+    fn rollback_err(self) -> Result<T, ProcessError>;
 }
 
-impl<T, E> RouteErrorExt<T> for Result<T, E>
+impl<T, E> ProcessErrorExt<T> for Result<T, E>
 where
     E: Into<Box<dyn Error + Send + Sync + 'static>>,
 {
@@ -137,6 +147,18 @@ where
     fn route_err_to_failure(self) -> Result<T, ProcessError> {
         self.route_to("failure", LogLevel::Warn)
     }
+
+    fn rollback_err(self) -> Result<T, ProcessError> {
+        self.map_err(|e| {
+            let boxed: Box<dyn Error + Send + Sync + 'static> = e.into();
+            match boxed.downcast::<MinifiError>() {
+                // The source is already a MinifiError: keep its variant (and 
its
+                // `to_status()` mapping) instead of re-boxing it into `Other`.
+                Ok(minifi_error) => ProcessError::Rollback(*minifi_error),
+                Err(other) => 
ProcessError::Rollback(MinifiError::Other(other)),
+            }
+        })
+    }
 }
 
 #[derive(Debug)]
@@ -304,23 +326,48 @@ mod tests {
     }
 
     #[test]
-    fn minifi_error_converts_to_fatal_via_from() {
+    fn minifi_error_converts_to_route_to_failure_via_from() {
         let pe: ProcessError = MinifiError::custom("nope").into();
-        assert!(matches!(
-            pe,
-            ProcessError::Fatal(MinifiError::CustomError(_))
-        ));
+        match pe {
+            ProcessError::Route(route) => {
+                assert_eq!(route.relationship, "failure");
+                assert_eq!(route.log_level, LogLevel::Warn);
+            }
+            other => panic!("expected a route error, got {other:?}"),
+        }
     }
 
     #[test]
-    fn raw_error_question_mark_becomes_fatal() {
+    fn raw_error_question_mark_routes_to_failure() {
         fn inner() -> Result<(), ProcessError> {
             Err(io_err())?;
             Ok(())
         }
+        match inner() {
+            Err(ProcessError::Route(route)) => {
+                assert_eq!(route.relationship, "failure");
+                assert_eq!(route.log_level, LogLevel::Warn);
+                assert_eq!(route.source.to_string(), "boom");
+            }
+            other => panic!("expected a route error, got {other:?}"),
+        }
+    }
+
+    #[test]
+    fn rollback_err_wraps_foreign_error_as_other() {
+        let res: Result<(), std::io::Error> = Err(io_err());
+        assert!(matches!(
+            res.rollback_err(),
+            Err(ProcessError::Rollback(MinifiError::Other(_)))
+        ));
+    }
+
+    #[test]
+    fn rollback_err_preserves_minifi_error_variant() {
+        let res: Result<(), MinifiError> = Err(MinifiError::validation("bad"));
         assert!(matches!(
-            inner(),
-            Err(ProcessError::Fatal(MinifiError::IoError(_)))
+            res.rollback_err(),
+            Err(ProcessError::Rollback(MinifiError::ValidationError(_)))
         ));
     }
 }
diff --git 
a/minifi_rust/minifi_native/src/api/processor_wrappers/complex_processor.rs 
b/minifi_rust/minifi_native/src/api/processor_wrappers/complex_processor.rs
index 7699f4b33..b2538e3b8 100644
--- a/minifi_rust/minifi_native/src/api/processor_wrappers/complex_processor.rs
+++ b/minifi_rust/minifi_native/src/api/processor_wrappers/complex_processor.rs
@@ -67,7 +67,7 @@ where
         if let Some(ref mut scheduled_impl) = self.scheduled_impl {
             scheduled_impl.trigger(context, session, &self.logger)
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
@@ -90,7 +90,7 @@ where
         if let Some(ref scheduled_impl) = self.scheduled_impl {
             scheduled_impl.trigger(context, session, &self.logger)
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
diff --git 
a/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_source.rs 
b/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_source.rs
index 23587778e..0e4fc6311 100644
--- a/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_source.rs
+++ b/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_source.rs
@@ -20,8 +20,8 @@ use crate::api::raw_processor::{MultiThreadedTrigger, 
SingleThreadedTrigger};
 use crate::{FlowFileAttribute, impl_with_attributes};
 use crate::{
     GetControllerService, GetProperty, Logger, MinifiError, MultiThreaded, 
OnTriggerResult,
-    ProcessContext, ProcessError, ProcessSession, Processor, Relationship, 
Schedule,
-    SingleThreaded,
+    ProcessContext, ProcessError, ProcessErrorExt, ProcessSession, Processor, 
Relationship,
+    Schedule, SingleThreaded,
 };
 
 pub struct GeneratedFlowFile<'a> {
@@ -75,16 +75,20 @@ where
     }
 
     for new_flow_file_data in generated_flow_files {
-        let mut ff = session.create()?;
+        let mut ff = session.create().rollback_err()?;
         match new_flow_file_data.new_content {
             None => {}
-            Some(Content::Buffer(buffer)) => session.write(&ff, &buffer)?,
-            Some(Content::Stream(stream)) => session.write_from_stream(&ff, 
stream)?,
+            Some(Content::Buffer(buffer)) => session.write(&ff, 
&buffer).rollback_err()?,
+            Some(Content::Stream(stream)) => {
+                session.write_from_stream(&ff, stream).rollback_err()?
+            }
         }
         for (k, v) in &new_flow_file_data.attributes_to_add {
-            session.set_attribute(&mut ff, k, v)?;
+            session.set_attribute(&mut ff, k, v).rollback_err()?;
         }
-        session.transfer(ff, 
new_flow_file_data.target_relationship_name.as_ref())?;
+        session
+            .transfer(ff, new_flow_file_data.target_relationship_name.as_ref())
+            .rollback_err()?;
     }
     Ok(OnTriggerResult::Ok)
 }
@@ -110,7 +114,7 @@ where
             let files = scheduled_impl.generate(context, &self.logger)?;
             handle_generated_flow_files::<PC, PS>(session, files)
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
@@ -134,7 +138,7 @@ where
             let files = scheduled_impl.generate(context, &self.logger)?;
             handle_generated_flow_files::<PC, PS>(session, files)
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
diff --git 
a/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_stream_transform.rs
 
b/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_stream_transform.rs
index f300f99b8..502fe6fc7 100644
--- 
a/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_stream_transform.rs
+++ 
b/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_stream_transform.rs
@@ -21,8 +21,8 @@ use crate::api::raw_processor::{MultiThreadedTrigger, 
SingleThreadedTrigger};
 use crate::{FlowFileAttribute, impl_with_attributes};
 use crate::{
     GetAttribute, GetControllerService, GetProperty, InputStream, LogLevel, 
Logger, MinifiError,
-    MultiThreaded, OnTriggerResult, OutputStream, ProcessContext, 
ProcessError, ProcessSession,
-    Processor, Relationship, Schedule, SingleThreaded,
+    MultiThreaded, OnTriggerResult, OutputStream, ProcessContext, 
ProcessError, ProcessErrorExt,
+    ProcessSession, Processor, Relationship, Schedule, SingleThreaded,
 };
 
 #[derive(Debug)]
@@ -112,34 +112,43 @@ where
     if let Some(mut flow_file) = session.get() {
         let simple_context = ContextSessionFlowFileBundle::new(context, 
session, Some(&flow_file));
 
-        let (relationship, attrs) = session.read_stream(&flow_file, 
|input_stream| {
-            session.write_stream(&flow_file, |output_stream| {
-                let transformed = match transform_fn(&simple_context, 
input_stream, output_stream) {
-                    Ok(t) => t,
-                    Err(ProcessError::Route(route)) => {
-                        route.log(logger);
-                        
TransformStreamResult::route_without_changes_by_name(route.relationship)
-                    }
-                    Err(ProcessError::Fatal(e)) => {
-                        return Err(e);
-                    }
-                };
-
-                Ok((
-                    (
-                        transformed.target_relationship_name,
-                        transformed.attributes_to_add,
-                    ),
-                    transformed.write_status,
-                ))
+        let (relationship, attrs) = session
+            .read_stream(&flow_file, |input_stream| {
+                session.write_stream(&flow_file, |output_stream| {
+                    let transformed =
+                        match transform_fn(&simple_context, input_stream, 
output_stream) {
+                            Ok(t) => t,
+                            Err(ProcessError::Route(route)) => {
+                                route.log(logger);
+                                
TransformStreamResult::route_without_changes_by_name(
+                                    route.relationship,
+                                )
+                            }
+                            Err(ProcessError::Rollback(e)) => {
+                                return Err(e);
+                            }
+                        };
+
+                    Ok((
+                        (
+                            transformed.target_relationship_name,
+                            transformed.attributes_to_add,
+                        ),
+                        transformed.write_status,
+                    ))
+                })
             })
-        })?;
+            .rollback_err()?;
 
         for (k, v) in attrs {
-            session.set_attribute(&mut flow_file, &k, &v)?;
+            session
+                .set_attribute(&mut flow_file, &k, &v)
+                .rollback_err()?;
         }
 
-        session.transfer(flow_file, relationship.as_ref())?;
+        session
+            .transfer(flow_file, relationship.as_ref())
+            .rollback_err()?;
 
         Ok(OnTriggerResult::Ok)
     } else {
@@ -168,7 +177,7 @@ where
                 scheduled_impl.transform(ctx, input, output, &self.logger)
             })
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
@@ -193,7 +202,7 @@ where
                 scheduled_impl.transform(ctx, input, output, &self.logger)
             })
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
diff --git 
a/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_transform.rs 
b/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_transform.rs
index 2f470125e..81f0e84c6 100644
--- 
a/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_transform.rs
+++ 
b/minifi_rust/minifi_native/src/api/processor_wrappers/flow_file_transform.rs
@@ -23,8 +23,8 @@ use crate::api::property::{GetControllerService, GetProperty};
 use crate::api::raw_processor::{MultiThreadedTrigger, SingleThreadedTrigger};
 use crate::{
     GetAttribute, LogLevel, Logger, MinifiError, MultiThreaded, 
OnTriggerResult, ProcessContext,
-    ProcessError, ProcessSession, Relationship, Schedule, SingleThreaded, 
impl_with_attributes,
-    info,
+    ProcessError, ProcessErrorExt, ProcessSession, Relationship, Schedule, 
SingleThreaded,
+    impl_with_attributes, info,
 };
 
 use minifi_native::InputStream;
@@ -147,40 +147,46 @@ where
     if let Some(mut flow_file) = session.get() {
         let simple_context = ContextSessionFlowFileBundle::new(context, 
session, Some(&flow_file));
 
-        let (attrs_to_add, relationship) = session.read_stream(&flow_file, 
|input_stream| {
-            let transformed = match transform_fn(&simple_context, 
input_stream) {
-                Ok(transform_success) => transform_success,
-                Err(ProcessError::Route(route)) => {
-                    route.log(logger);
-                    
TransformedFlowFile::route_without_changes_by_name(route.relationship)
-                }
-                Err(ProcessError::Fatal(e)) => {
-                    return Err(e);
-                }
-            };
-
-            info!(logger, "{:?}", transformed);
-            match transformed.new_content {
-                None => {}
-                Some(Content::Buffer(buffer)) => {
-                    session.write(&flow_file, &buffer)?;
-                }
-                Some(Content::Stream(stream)) => {
-                    session.write_from_stream(&flow_file, stream)?;
-                }
-            };
-
-            Ok((
-                transformed.attributes_to_add,
-                transformed.target_relationship_name,
-            ))
-        })?;
+        let (attrs_to_add, relationship) = session
+            .read_stream(&flow_file, |input_stream| {
+                let transformed = match transform_fn(&simple_context, 
input_stream) {
+                    Ok(transform_success) => transform_success,
+                    Err(ProcessError::Route(route)) => {
+                        route.log(logger);
+                        
TransformedFlowFile::route_without_changes_by_name(route.relationship)
+                    }
+                    Err(ProcessError::Rollback(e)) => {
+                        return Err(e);
+                    }
+                };
+
+                info!(logger, "{:?}", transformed);
+                match transformed.new_content {
+                    None => {}
+                    Some(Content::Buffer(buffer)) => {
+                        session.write(&flow_file, &buffer)?;
+                    }
+                    Some(Content::Stream(stream)) => {
+                        session.write_from_stream(&flow_file, stream)?;
+                    }
+                };
+
+                Ok((
+                    transformed.attributes_to_add,
+                    transformed.target_relationship_name,
+                ))
+            })
+            .rollback_err()?;
 
         for (k, v) in attrs_to_add {
-            session.set_attribute(&mut flow_file, &k, &v)?;
+            session
+                .set_attribute(&mut flow_file, &k, &v)
+                .rollback_err()?;
         }
 
-        session.transfer(flow_file, relationship.as_ref())?;
+        session
+            .transfer(flow_file, relationship.as_ref())
+            .rollback_err()?;
         Ok(OnTriggerResult::Ok)
     } else {
         logger.log(LogLevel::Trace, format_args!("No flowfile to transform"));
@@ -208,7 +214,7 @@ where
                 scheduled_impl.transform(ctx, input, &self.logger)
             })
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
@@ -233,7 +239,7 @@ where
                 scheduled_impl.transform(ctx, input, &self.logger)
             })
         } else {
-            Err(MinifiError::UnscheduledProcessor.into())
+            Err(ProcessError::Rollback(MinifiError::UnscheduledProcessor))
         }
     }
 }
@@ -245,7 +251,7 @@ mod tests {
     use crate::api::raw_processor::MultiThreadedTrigger;
     use crate::{
         GetControllerService, GetId, MockFlowFile, MockLogger, 
MockProcessContext,
-        MockProcessSession, ProcessError, RouteErrorExt,
+        MockProcessSession, ProcessError, ProcessErrorExt,
     };
 
     struct RouteToFailure;
@@ -273,13 +279,13 @@ mod tests {
         }
     }
 
-    struct FatalTransform;
-    impl Schedule for FatalTransform {
+    struct RollbackTransform;
+    impl Schedule for RollbackTransform {
         fn schedule<Ctx: GetProperty, L: Logger>(_c: &Ctx, _l: &L) -> 
Result<Self, MinifiError> {
-            Ok(FatalTransform)
+            Ok(RollbackTransform)
         }
     }
-    impl FlowFileTransform for FatalTransform {
+    impl FlowFileTransform for RollbackTransform {
         fn transform<
             'a,
             Context: GetProperty + GetControllerService + GetAttribute + GetId,
@@ -290,7 +296,7 @@ mod tests {
             _input_stream: &'a mut dyn InputStream,
             _logger: &LoggerImpl,
         ) -> Result<TransformedFlowFile<'a>, ProcessError> {
-            Err(ProcessError::Fatal(MinifiError::custom("real error")))
+            Err(ProcessError::Rollback(MinifiError::custom("real error")))
         }
     }
 
@@ -327,14 +333,14 @@ mod tests {
     }
 
     #[test]
-    fn fatal_error_propagates_and_transfers_nothing() {
+    fn rollback_error_propagates_and_transfers_nothing() {
         let mut processor: Processor<
-            FatalTransform,
+            RollbackTransform,
             FlowFileTransformProcessorType,
             MultiThreaded,
             MockLogger,
         > = Processor::new(MockLogger::new());
-        processor.scheduled_impl = Some(FatalTransform);
+        processor.scheduled_impl = Some(RollbackTransform);
 
         let mut context = MockProcessContext::new();
         let mut session = seeded_session();
@@ -343,7 +349,7 @@ mod tests {
 
         assert!(matches!(
             result,
-            Err(ProcessError::Fatal(MinifiError::CustomError(_)))
+            Err(ProcessError::Rollback(MinifiError::CustomError(_)))
         ));
         assert_eq!(session.num_of_transferred_flow_files(), 0);
     }
diff --git a/minifi_rust/minifi_native/src/c_ffi/c_ffi_processor_definition.rs 
b/minifi_rust/minifi_native/src/c_ffi/c_ffi_processor_definition.rs
index 4114b524f..5ae898595 100644
--- a/minifi_rust/minifi_native/src/c_ffi/c_ffi_processor_definition.rs
+++ b/minifi_rust/minifi_native/src/c_ffi/c_ffi_processor_definition.rs
@@ -40,7 +40,7 @@ fn process_error_to_status<P: RawProcessor>(
     match result {
         Ok(OnTriggerResult::Ok) => minifi_status_MINIFI_STATUS_SUCCESS,
         Ok(OnTriggerResult::Yield) => 
minifi_status_MINIFI_STATUS_PROCESSOR_YIELD,
-        Err(ProcessError::Fatal(err)) => {
+        Err(ProcessError::Rollback(err)) => {
             processor.log(
                 LogLevel::Error,
                 format_args!("Error during trigger {}", err),
diff --git a/minifi_rust/minifi_native/src/lib.rs 
b/minifi_rust/minifi_native/src/lib.rs
index b855291a2..4e70b6cc0 100644
--- a/minifi_rust/minifi_native/src/lib.rs
+++ b/minifi_rust/minifi_native/src/lib.rs
@@ -20,7 +20,7 @@ mod api;
 pub mod c_ffi;
 pub mod mock;
 
-pub use api::errors::{MinifiError, ProcessError, RouteError, RouteErrorExt};
+pub use api::errors::{MinifiError, ProcessError, ProcessErrorExt, RouteError};
 
 pub use api::component_definition_traits::{
     ComponentIdentifier, ControllerServiceDefinition, ProcessorDefinition,

Reply via email to