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

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

commit 4d396eaf7ac69eaae842b94fe8d4bffec808ca91
Author: Martin Zink <[email protected]>
AuthorDate: Fri Apr 24 11:09:55 2026 +0200

    MINIFICPP-2760 Input/OutputStreamCallback should return std::expected
    
    Signed-off-by: Ferenc Gerlits <[email protected]>
    
    Closes #2148
---
 core-framework/include/http/HTTPCallback.h         |  10 +-
 core-framework/include/io/StreamPipe.h             |  14 +--
 core-framework/include/utils/ByteArrayCallback.h   |  11 ++-
 core-framework/include/utils/JsonCallback.h        |  21 +++--
 .../utils/LineByLineInputOutputStreamCallback.h    |   2 +-
 .../utils/LineByLineInputOutputStreamCallback.cpp  |  19 ++--
 .../cpp-extension-lib/src/core/ProcessSession.cpp  |  44 +++++----
 .../include/serialization/FlowFileSerializer.h     |   4 +-
 .../include/serialization/FlowFileV3Serializer.h   |   2 +-
 .../include/serialization/PayloadSerializer.h      |   2 +-
 .../include/utils/file/FileReaderCallback.h        |   2 +-
 .../include/utils/file/FileWriterCallback.h        |   2 +-
 .../src/serialization/FlowFileV3Serializer.cpp     |  20 ++--
 .../src/serialization/PayloadSerializer.cpp        |   2 +-
 .../src/utils/file/FileReaderCallback.cpp          |   8 +-
 .../src/utils/file/FileWriterCallback.cpp          |   6 +-
 extensions/aws/processors/FetchS3Object.cpp        |   5 +-
 extensions/aws/processors/PutS3Object.cpp          |  13 ++-
 .../azure/processors/FetchAzureBlobStorage.cpp     |   6 +-
 .../azure/processors/FetchAzureDataLakeStorage.cpp |   8 +-
 extensions/azure/processors/PutAzureBlobStorage.h  |   8 +-
 .../azure/processors/PutAzureDataLakeStorage.cpp   |   6 +-
 .../azure/processors/PutAzureDataLakeStorage.h     |   2 +-
 extensions/azure/storage/AzureBlobStorage.cpp      |   2 +-
 extensions/azure/storage/AzureDataLakeStorage.cpp  |   2 +-
 extensions/bustache/ApplyTemplate.cpp              |   4 +-
 .../couchbase/processors/GetCouchbaseKey.cpp       |   6 +-
 extensions/gcp/processors/FetchGCSObject.cpp       |   8 +-
 extensions/gcp/processors/PutGCSObject.cpp         |   6 +-
 extensions/kafka/PublishKafka.cpp                  |  10 +-
 extensions/libarchive/CompressContent.cpp          |  18 ++--
 extensions/libarchive/CompressContent.h            |  12 +--
 extensions/libarchive/FocusArchiveEntry.cpp        |  15 +--
 extensions/libarchive/FocusArchiveEntry.h          |   2 +-
 extensions/libarchive/MergeContent.cpp             |   4 +-
 extensions/libarchive/MergeContent.h               |  42 +++++----
 extensions/libarchive/UnfocusArchiveEntry.cpp      |   4 +-
 extensions/libarchive/UnfocusArchiveEntry.h        |   2 +-
 .../libarchive/tests/CompressContentTests.cpp      |   6 +-
 extensions/libarchive/tests/MergeFileTests.cpp     |  12 +--
 extensions/lua/LuaProcessSession.cpp               |  12 +--
 extensions/mqtt/processors/ConsumeMQTT.cpp         |  10 +-
 extensions/mqtt/processors/ConsumeMQTT.h           |   2 +-
 extensions/mqtt/processors/PublishMQTT.cpp         |   4 +-
 extensions/opencv/CaptureRTSPFrame.cpp             |   4 +-
 extensions/opencv/MotionDetector.cpp               |  10 +-
 extensions/python/types/PyProcessSession.cpp       |  12 +--
 extensions/python/types/PyRecordSetReader.cpp      |   2 +-
 extensions/rocksdb-repos/tests/SwapTests.cpp       |  11 +--
 extensions/sftp/processors/FetchSFTP.cpp           |   4 +-
 extensions/sftp/processors/PutSFTP.cpp             |   2 +-
 .../controllers/JsonRecordSetWriter.cpp            |   8 +-
 .../controllers/XMLRecordSetWriter.cpp             |   4 +-
 .../processors/ConvertRecord.cpp                   |   4 +-
 .../processors/DefragmentText.cpp                  |   6 +-
 .../processors/EvaluateJsonPath.cpp                |   4 +-
 .../standard-processors/processors/ExtractText.cpp |   8 +-
 .../standard-processors/processors/ExtractText.h   |   2 +-
 .../standard-processors/processors/HashContent.cpp |   5 +-
 .../standard-processors/processors/RouteText.cpp   |   8 +-
 .../standard-processors/processors/SplitJson.cpp   |   4 +-
 .../standard-processors/processors/SplitRecord.cpp |   4 +-
 .../standard-processors/processors/SplitText.cpp   |  23 +++--
 .../standard-processors/processors/TailFile.cpp    |   8 +-
 .../tests/unit/XMLRecordSetWriterTests.cpp         |   2 +-
 libminifi/include/c2/PayloadSerializer.h           |   2 +-
 .../include/core/ProcessSessionReadCallback.h      |   2 +-
 libminifi/src/c2/C2Utils.cpp                       |   2 +-
 libminifi/src/core/ProcessSession.cpp              |  41 ++++----
 libminifi/src/core/ProcessSessionReadCallback.cpp  |   8 +-
 libminifi/src/minifi-c.cpp                         |  10 +-
 .../src/sitetosite/CompressionOutputStream.cpp     |  23 +++--
 libminifi/src/sitetosite/SiteToSiteClient.cpp      |   8 +-
 .../libtest/unit/ContentRepositoryDependentTests.h |  18 ++--
 libminifi/test/libtest/unit/TestBase.h             |   4 +-
 libminifi/test/unit/FlowFileSerializationTests.cpp |   4 +-
 libminifi/test/unit/MetricsTests.cpp               |  12 +--
 libminifi/test/unit/SiteToSiteTests.cpp            |   2 +-
 .../common/include/minifi-cpp/io/StreamCallback.h  | 103 +++++++++++++++++++--
 minifi-api/include/minifi-cpp/core/Property.h      |   2 +-
 80 files changed, 433 insertions(+), 343 deletions(-)

diff --git a/core-framework/include/http/HTTPCallback.h 
b/core-framework/include/http/HTTPCallback.h
index 991085baf..31df3598f 100644
--- a/core-framework/include/http/HTTPCallback.h
+++ b/core-framework/include/http/HTTPCallback.h
@@ -65,7 +65,7 @@ class HttpStreamingCallback final : public 
HTTPUploadByteArrayInputCallback {
     seekInner(lock, pos);
   }
 
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream) override {
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) 
override {
     std::vector<std::byte> vec;
 
     if (stream->size() > 0) {
@@ -76,7 +76,7 @@ class HttpStreamingCallback final : public 
HTTPUploadByteArrayInputCallback {
     return processInner(std::move(vec));
   }
 
-  int64_t process(const uint8_t* data, size_t size) {
+  io::IoResult process(const uint8_t* data, size_t size) {
     std::vector<std::byte> vec;
     vec.resize(size);
     memcpy(vec.data(), data, size);
@@ -156,20 +156,20 @@ class HttpStreamingCallback final : public 
HTTPUploadByteArrayInputCallback {
    * @param vec the buffer to be inserted
    * @return the number of bytes processed (the size of vec)
    */
-  int64_t processInner(std::vector<std::byte>&& vec) {
+  io::IoResult processInner(std::vector<std::byte>&& vec) {
     size_t size = vec.size();
 
     logger_->log_trace("processInner() called, vec.data(): {}, vec.size(): 
{}", static_cast<void*>(vec.data()), size);
 
     if (size == 0U) {
-      return 0U;
+      return io::IoResult::zero();
     }
 
     std::unique_lock<std::mutex> lock(mutex_);
     byte_arrays_.emplace_back(std::move(vec));
     cv.notify_all();
 
-    return size;
+    return io::IoResult::from(size);
   }
 
   /**
diff --git a/core-framework/include/io/StreamPipe.h 
b/core-framework/include/io/StreamPipe.h
index ccd7809d2..b5b77dc9e 100644
--- a/core-framework/include/io/StreamPipe.h
+++ b/core-framework/include/io/StreamPipe.h
@@ -31,12 +31,12 @@
 namespace org::apache::nifi::minifi {
 namespace internal {
 
-inline int64_t pipe(io::InputStream& src, io::OutputStream& dst) {
+inline io::IoResult pipe(io::InputStream& src, io::OutputStream& dst) {
   std::array<std::byte, utils::configuration::DEFAULT_BUFFER_SIZE> buffer{};
-  size_t totalTransferred = 0;
+  uint64_t totalTransferred = 0;
   while (true) {
     const auto readRet = src.read(buffer);
-    if (io::isError(readRet)) return -1;
+    if (io::isError(readRet)) return io::IoResult::error();
     if (readRet == 0) break;
     auto remaining = readRet;
     size_t transferred = 0;
@@ -48,14 +48,14 @@ inline int64_t pipe(io::InputStream& src, io::OutputStream& 
dst) {
       //     - the number of bytes read or
       //     - the number of bytes wrote
       if (io::isError(writeRet)) {
-        return -1;
+        return io::IoResult::error();
       }
       transferred += writeRet;
       remaining -= writeRet;
     }
     totalTransferred += transferred;
   }
-  return gsl::narrow<int64_t>(totalTransferred);
+  return io::IoResult::from(totalTransferred);
 }
 
 }  // namespace internal
@@ -64,7 +64,7 @@ class InputStreamPipe {
  public:
   explicit InputStreamPipe(io::OutputStream& output) : output_(&output) {}
 
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream) const {
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) 
const {
     return internal::pipe(*stream, *output_);
   }
 
@@ -76,7 +76,7 @@ class OutputStreamPipe {
  public:
   explicit OutputStreamPipe(io::InputStream& input) : input_(&input) {}
 
-  int64_t operator()(const std::shared_ptr<io::OutputStream>& stream) const {
+  io::IoResult operator()(const std::shared_ptr<io::OutputStream>& stream) 
const {
     return internal::pipe(*input_, *stream);
   }
 
diff --git a/core-framework/include/utils/ByteArrayCallback.h 
b/core-framework/include/utils/ByteArrayCallback.h
index d07b7bf1d..2a430d406 100644
--- a/core-framework/include/utils/ByteArrayCallback.h
+++ b/core-framework/include/utils/ByteArrayCallback.h
@@ -16,16 +16,17 @@
  */
 #pragma once
 
+#include <condition_variable>
 #include <memory>
 #include <string>
-#include <vector>
 #include <utility>
-#include <condition_variable>
+#include <vector>
 
 #include "concurrentqueue.h"
 #include "core/logging/LoggerFactory.h"
-#include "minifi-cpp/utils/gsl.h"
 #include "minifi-cpp/io/InputStream.h"
+#include "minifi-cpp/io/StreamCallback.h"
+#include "minifi-cpp/utils/gsl.h"
 
 namespace org::apache::nifi::minifi::utils {
 
@@ -36,7 +37,7 @@ class ByteInputCallback {
  public:
   virtual ~ByteInputCallback() = default;
 
-  virtual int64_t operator()(const std::shared_ptr<io::InputStream>& stream) {
+  virtual io::IoResult operator()(const std::shared_ptr<io::InputStream>& 
stream) {
     stream->seek(0);
 
     if (stream->size() > 0) {
@@ -44,7 +45,7 @@ class ByteInputCallback {
       stream->read(vec);
     }
 
-    return gsl::narrow<int64_t>(vec.size());
+    return io::IoResult::from(vec.size());
   }
 
   virtual void close() { }
diff --git a/core-framework/include/utils/JsonCallback.h 
b/core-framework/include/utils/JsonCallback.h
index 37d922bd5..61194ed78 100644
--- a/core-framework/include/utils/JsonCallback.h
+++ b/core-framework/include/utils/JsonCallback.h
@@ -32,18 +32,19 @@ namespace org::apache::nifi::minifi::utils {
 class JsonInputCallback {
  public:
   explicit JsonInputCallback(rapidjson::Document& document) : 
document_(document) {}
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream) {
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) {
     std::string content;
     content.resize(stream->size());
-    const auto read_ret = stream->read(as_writable_bytes(std::span(content)));
+    const size_t read_ret = 
stream->read(as_writable_bytes(std::span(content)));
     if (io::isError(read_ret)) {
-      return -1;
+      return io::IoResult::error();
     }
     rapidjson::ParseResult parse_result = 
document_.Parse<rapidjson::kParseStopWhenDoneFlag>(content.data());
-    if (parse_result.IsError())
-      return -1;
+    if (parse_result.IsError()) {
+      return io::IoResult::error();
+    }
 
-    return read_ret;
+    return io::IoResult::from(read_ret);
   }
  private:
   rapidjson::Document& document_;
@@ -54,14 +55,14 @@ class JsonOutputCallback {
   explicit JsonOutputCallback(rapidjson::Document&& root, 
std::optional<uint8_t> decimal_places)
       : root_(std::move(root)), decimal_places_(decimal_places) {}
 
-  int64_t operator()(const std::shared_ptr<io::OutputStream>& stream) const {
+  io::IoResult operator()(const std::shared_ptr<io::OutputStream>& stream) 
const {
     rapidjson::StringBuffer buffer;
     rapidjson::Writer<rapidjson::StringBuffer> writer(buffer);
     if (decimal_places_.has_value())
       writer.SetMaxDecimalPlaces(decimal_places_.value());
     root_.Accept(writer);
     const auto write_return = stream->write(reinterpret_cast<const 
uint8_t*>(buffer.GetString()), buffer.GetSize());
-    return !io::isError(write_return) ? gsl::narrow<int64_t>(write_return) : 
-1;
+    return io::IoResult::from(write_return);
   }
 
  protected:
@@ -74,14 +75,14 @@ class PrettyJsonOutputCallback {
   explicit PrettyJsonOutputCallback(rapidjson::Document&& root, 
std::optional<uint8_t> decimal_places)
       : root_(std::move(root)), decimal_places_(decimal_places) {}
 
-  int64_t operator()(const std::shared_ptr<io::OutputStream>& stream) const {
+  io::IoResult operator()(const std::shared_ptr<io::OutputStream>& stream) 
const {
     rapidjson::StringBuffer buffer;
     rapidjson::PrettyWriter<rapidjson::StringBuffer> writer(buffer);
     if (decimal_places_.has_value())
       writer.SetMaxDecimalPlaces(decimal_places_.value());
     root_.Accept(writer);
     const auto write_return = stream->write(reinterpret_cast<const 
uint8_t*>(buffer.GetString()), buffer.GetSize());
-    return !io::isError(write_return) ? gsl::narrow<int64_t>(write_return) : 
-1;
+    return io::IoResult::from(write_return);
   }
 
  protected:
diff --git a/core-framework/include/utils/LineByLineInputOutputStreamCallback.h 
b/core-framework/include/utils/LineByLineInputOutputStreamCallback.h
index 72fb0bedd..b199219ac 100644
--- a/core-framework/include/utils/LineByLineInputOutputStreamCallback.h
+++ b/core-framework/include/utils/LineByLineInputOutputStreamCallback.h
@@ -33,7 +33,7 @@ class LineByLineInputOutputStreamCallback {
  public:
   using CallbackType = std::function<std::string(const std::string& 
input_line, bool is_first_line, bool is_last_line)>;
   explicit LineByLineInputOutputStreamCallback(CallbackType callback);
-  std::optional<io::ReadWriteResult> operator()(const 
std::shared_ptr<io::InputStream>& input, const 
std::shared_ptr<io::OutputStream>& output);
+  io::ReadWriteResult operator()(const std::shared_ptr<io::InputStream>& 
input, const std::shared_ptr<io::OutputStream>& output);
 
  private:
   int64_t readInput(io::InputStream& stream);
diff --git a/core-framework/src/utils/LineByLineInputOutputStreamCallback.cpp 
b/core-framework/src/utils/LineByLineInputOutputStreamCallback.cpp
index ec59fff53..fce1e89ff 100644
--- a/core-framework/src/utils/LineByLineInputOutputStreamCallback.cpp
+++ b/core-framework/src/utils/LineByLineInputOutputStreamCallback.cpp
@@ -26,22 +26,21 @@ 
LineByLineInputOutputStreamCallback::LineByLineInputOutputStreamCallback(Callbac
   : callback_(std::move(callback)) {
 }
 
-std::optional<io::ReadWriteResult> 
LineByLineInputOutputStreamCallback::operator()(const 
std::shared_ptr<io::InputStream>& input, const 
std::shared_ptr<io::OutputStream>& output) {
+io::ReadWriteResult LineByLineInputOutputStreamCallback::operator()(const 
std::shared_ptr<io::InputStream>& input, const 
std::shared_ptr<io::OutputStream>& output) {
   gsl_Expects(input);
   gsl_Expects(output);
 
-  io::ReadWriteResult result;
 
-  if (int64_t status = readInput(*input); status <= 0) {
+  if (const int64_t status = readInput(*input); status <= 0) {
     if (status < 0) {
-      return std::nullopt;
+      return io::ReadWriteResult::error();
     }
-    return result;
+    return io::ReadWriteResult::zero();
   }
 
-  result.bytes_read = gsl::narrow<int64_t>(input_.size());
+  const uint64_t bytes_read = input_.size();
 
-  std::size_t total_bytes_written = 0;
+  uint64_t total_bytes_written = 0;
   bool is_first_line = true;
   readLine();
   do {
@@ -49,14 +48,14 @@ std::optional<io::ReadWriteResult> 
LineByLineInputOutputStreamCallback::operator
     std::string output_line = callback_(*current_line_, is_first_line, 
isLastLine());
     const auto bytes_written = output->write(reinterpret_cast<const uint8_t 
*>(output_line.data()), output_line.size());
     if (io::isError(bytes_written)) {
-      return std::nullopt;
+      return io::ReadWriteResult::error();
     }
     total_bytes_written += bytes_written;
     is_first_line = false;
   } while (!isLastLine());
 
-  result.bytes_written = gsl::narrow<int64_t>(total_bytes_written);
-  return result;
+
+  return { bytes_read, total_bytes_written };
 }
 
 int64_t LineByLineInputOutputStreamCallback::readInput(io::InputStream& 
stream) {
diff --git a/extension-framework/cpp-extension-lib/src/core/ProcessSession.cpp 
b/extension-framework/cpp-extension-lib/src/core/ProcessSession.cpp
index 6b22c0bbc..cc7ab50e1 100644
--- a/extension-framework/cpp-extension-lib/src/core/ProcessSession.cpp
+++ b/extension-framework/cpp-extension-lib/src/core/ProcessSession.cpp
@@ -17,10 +17,10 @@
 
 #include "api/core/ProcessSession.h"
 
+#include "api/core/FlowFile.h"
+#include "api/utils/minifi-c-utils.h"
 #include "io/InputStream.h"
 #include "io/OutputStream.h"
-#include "api/utils/minifi-c-utils.h"
-#include "api/core/FlowFile.h"
 #include "minifi-cpp/Exception.h"
 
 namespace org::apache::nifi::minifi::api::core {
@@ -91,21 +91,29 @@ void ProcessSession::remove(FlowFile ff) {
 }
 
 void ProcessSession::write(FlowFile& flow_file, const 
io::OutputStreamCallback& callback) {
-  const auto status = MinifiProcessSessionWrite(impl_, flow_file.get(), [] 
(void* data, MinifiOutputStream* output) {
-    return (*static_cast<const 
io::OutputStreamCallback*>(data))(std::make_shared<MinifiOutputStreamWrapper>(output));
-  }, const_cast<io::OutputStreamCallback*>(&callback));
-  if (status != MINIFI_STATUS_SUCCESS) {
-    throw minifi::Exception(minifi::FILE_OPERATION_EXCEPTION, "Failed to 
process flowfile content");
-  }
+  const auto status = MinifiProcessSessionWrite(
+      impl_,
+      flow_file.get(),
+      [](void* data, MinifiOutputStream* output) -> int64_t {
+        const auto result =
+            (*static_cast<const 
io::OutputStreamCallback*>(data))(std::make_shared<MinifiOutputStreamWrapper>(output));
+        return result.toI64();
+      },
+      const_cast<io::OutputStreamCallback*>(&callback));
+  if (status != MINIFI_STATUS_SUCCESS) { throw 
minifi::Exception(minifi::FILE_OPERATION_EXCEPTION, "Failed to process flowfile 
content"); }
 }
 
 void ProcessSession::read(FlowFile& flow_file, const io::InputStreamCallback& 
callback) {
-  const auto status = MinifiProcessSessionRead(impl_, flow_file.get(), [] 
(void* data, MinifiInputStream* input) {
-    return (*static_cast<const 
io::InputStreamCallback*>(data))(std::make_shared<MinifiInputStreamWrapper>(input));
-  }, const_cast<io::InputStreamCallback*>(&callback));
-  if (status != MINIFI_STATUS_SUCCESS) {
-    throw minifi::Exception(minifi::FILE_OPERATION_EXCEPTION, "Failed to 
process flowfile content");
-  }
+  const auto status = MinifiProcessSessionRead(
+      impl_,
+      flow_file.get(),
+      [](void* data, MinifiInputStream* input) -> int64_t {
+        const auto result =
+            (*static_cast<const 
io::InputStreamCallback*>(data))(std::make_shared<MinifiInputStreamWrapper>(input));
+        return result.toI64();
+      },
+      const_cast<io::InputStreamCallback*>(&callback));
+  if (status != MINIFI_STATUS_SUCCESS) { throw 
minifi::Exception(minifi::FILE_OPERATION_EXCEPTION, "Failed to process flowfile 
content"); }
 }
 
 void ProcessSession::setAttribute(FlowFile& ff, const std::string_view key, 
std::string value) {  // NOLINT(performance-unnecessary-value-param)
@@ -142,18 +150,18 @@ void ProcessSession::writeBuffer(FlowFile& flow_file, 
std::span<const char> buff
 }
 
 void ProcessSession::writeBuffer(FlowFile& flow_file, std::span<const 
std::byte> buffer) {
-  write(flow_file, [buffer](const std::shared_ptr<io::OutputStream>& 
output_stream) {
+  write(flow_file, [buffer](const std::shared_ptr<io::OutputStream>& 
output_stream) -> io::IoResult {
     const auto write_status = output_stream->write(buffer);
-    return io::isError(write_status) ? -1 : gsl::narrow<int64_t>(write_status);
+    return io::IoResult::from(write_status);
   });
 }
 
 std::vector<std::byte> ProcessSession::readBuffer(FlowFile& flow_file) {
   std::vector<std::byte> result;
-  read(flow_file, [&result](const std::shared_ptr<io::InputStream>& 
input_stream) {
+  read(flow_file, [&result](const std::shared_ptr<io::InputStream>& 
input_stream) -> io::IoResult {
     result.resize(input_stream->size());
     const auto read_status = input_stream->read(result);
-    return io::isError(read_status) ? -1 : gsl::narrow<int64_t>(read_status);
+    return io::IoResult::from(read_status);
   });
   return result;
 }
diff --git a/extension-framework/include/serialization/FlowFileSerializer.h 
b/extension-framework/include/serialization/FlowFileSerializer.h
index 85362931f..c8168a565 100644
--- a/extension-framework/include/serialization/FlowFileSerializer.h
+++ b/extension-framework/include/serialization/FlowFileSerializer.h
@@ -38,11 +38,11 @@ class FlowFile;
 
 class FlowFileSerializer {
  public:
-  using FlowFileReader = std::function<int64_t(const 
std::shared_ptr<core::FlowFile>&, const io::InputStreamCallback&)>;
+  using FlowFileReader = std::function<io::IoResult(const 
std::shared_ptr<core::FlowFile>&, const io::InputStreamCallback&)>;
 
   explicit FlowFileSerializer(FlowFileReader reader) : 
reader_(std::move(reader)) {}
 
-  virtual int64_t serialize(const std::shared_ptr<core::FlowFile>& flowFile, 
const std::shared_ptr<io::OutputStream>& out) = 0;
+  virtual io::IoResult serialize(const std::shared_ptr<core::FlowFile>& 
flowFile, const std::shared_ptr<io::OutputStream>& out) = 0;
 
   virtual ~FlowFileSerializer() = default;
 
diff --git a/extension-framework/include/serialization/FlowFileV3Serializer.h 
b/extension-framework/include/serialization/FlowFileV3Serializer.h
index 64c9ce691..fc27215a8 100644
--- a/extension-framework/include/serialization/FlowFileV3Serializer.h
+++ b/extension-framework/include/serialization/FlowFileV3Serializer.h
@@ -38,7 +38,7 @@ class FlowFileV3Serializer : public FlowFileSerializer {
  public:
   using FlowFileSerializer::FlowFileSerializer;
 
-  int64_t serialize(const std::shared_ptr<core::FlowFile>& flowFile, const 
std::shared_ptr<io::OutputStream>& out) override;
+  io::IoResult serialize(const std::shared_ptr<core::FlowFile>& flowFile, 
const std::shared_ptr<io::OutputStream>& out) override;
 };
 
 }  // namespace org::apache::nifi::minifi
diff --git a/extension-framework/include/serialization/PayloadSerializer.h 
b/extension-framework/include/serialization/PayloadSerializer.h
index b6f843ca5..f9b14cf4e 100644
--- a/extension-framework/include/serialization/PayloadSerializer.h
+++ b/extension-framework/include/serialization/PayloadSerializer.h
@@ -28,7 +28,7 @@ class PayloadSerializer : public FlowFileSerializer {
  public:
   using FlowFileSerializer::FlowFileSerializer;
 
-  int64_t serialize(const std::shared_ptr<core::FlowFile>& flowFile, const 
std::shared_ptr<io::OutputStream>& out) override;
+  io::IoResult serialize(const std::shared_ptr<core::FlowFile>& flowFile, 
const std::shared_ptr<io::OutputStream>& out) override;
 };
 
 }  // namespace org::apache::nifi::minifi
diff --git a/extension-framework/include/utils/file/FileReaderCallback.h 
b/extension-framework/include/utils/file/FileReaderCallback.h
index 2bf55cbe0..b206d35a1 100644
--- a/extension-framework/include/utils/file/FileReaderCallback.h
+++ b/extension-framework/include/utils/file/FileReaderCallback.h
@@ -33,7 +33,7 @@ namespace org::apache::nifi::minifi::utils {
 class FileReaderCallback {
  public:
   explicit FileReaderCallback(std::filesystem::path file_path, size_t 
buffer_size);
-  int64_t operator()(const std::shared_ptr<io::OutputStream>& output_stream) 
const;
+  io::IoResult operator()(const std::shared_ptr<io::OutputStream>& 
output_stream) const;
 
  private:
   std::filesystem::path file_path_;
diff --git a/extension-framework/include/utils/file/FileWriterCallback.h 
b/extension-framework/include/utils/file/FileWriterCallback.h
index 2cb8729fa..6f3820e46 100644
--- a/extension-framework/include/utils/file/FileWriterCallback.h
+++ b/extension-framework/include/utils/file/FileWriterCallback.h
@@ -29,7 +29,7 @@ class FileWriterCallback {
  public:
   explicit FileWriterCallback(std::filesystem::path dest_path);
   ~FileWriterCallback();
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream);
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream);
   bool commit();
 
 
diff --git a/extension-framework/src/serialization/FlowFileV3Serializer.cpp 
b/extension-framework/src/serialization/FlowFileV3Serializer.cpp
index 972092ce4..471898096 100644
--- a/extension-framework/src/serialization/FlowFileV3Serializer.cpp
+++ b/extension-framework/src/serialization/FlowFileV3Serializer.cpp
@@ -57,43 +57,43 @@ size_t FlowFileV3Serializer::writeString(const std::string 
&str, const std::shar
   return sum;
 }
 
-int64_t FlowFileV3Serializer::serialize(const std::shared_ptr<core::FlowFile>& 
flowFile, const std::shared_ptr<io::OutputStream>& out) {
+io::IoResult FlowFileV3Serializer::serialize(const 
std::shared_ptr<core::FlowFile>& flowFile, const 
std::shared_ptr<io::OutputStream>& out) {
   size_t sum = 0;
   {
     const auto ret = out->write(MAGIC_HEADER, sizeof(MAGIC_HEADER));
-    if (io::isError(ret)) return -1;
-    if (ret != sizeof(MAGIC_HEADER)) return -1;
+    if (io::isError(ret)) return io::IoResult::error();
+    if (ret != sizeof(MAGIC_HEADER)) return io::IoResult::error();
     sum += ret;
   }
   const auto& attributes = flowFile->getAttributes();
   {
     const auto ret = writeLength(attributes.size(), out);
-    if (io::isError(ret)) return -1;
+    if (io::isError(ret)) return io::IoResult::error();
     sum += ret;
   }
   for (const auto& attrIt : attributes) {
     {
       const auto ret = writeString(attrIt.first, out);
-      if (io::isError(ret)) return -1;
+      if (io::isError(ret)) return io::IoResult::error();
       sum += ret;
     }
     {
       const auto ret = writeString(attrIt.second, out);
-      if (io::isError(ret)) return -1;
+      if (io::isError(ret)) return io::IoResult::error();
       sum += ret;
     }
   }
   {
     const auto ret = out->write(flowFile->getSize());
-    if (io::isError(ret)) return -1;
+    if (io::isError(ret)) return io::IoResult::error();
     sum += ret;
   }
   {
     const auto ret = reader_(flowFile, InputStreamPipe{*out});
-    if (ret < 0) return -1;
-    sum += gsl::narrow<size_t>(ret);
+    if (!ret) return ret;
+    sum += gsl::narrow<size_t>(ret.toI64());
   }
-  return gsl::narrow<int64_t>(sum);
+  return io::IoResult::from(sum);
 }
 
 }  // namespace org::apache::nifi::minifi
diff --git a/extension-framework/src/serialization/PayloadSerializer.cpp 
b/extension-framework/src/serialization/PayloadSerializer.cpp
index 7fd22e171..2a3912cc6 100644
--- a/extension-framework/src/serialization/PayloadSerializer.cpp
+++ b/extension-framework/src/serialization/PayloadSerializer.cpp
@@ -21,7 +21,7 @@
 
 namespace org::apache::nifi::minifi {
 
-int64_t PayloadSerializer::serialize(const std::shared_ptr<core::FlowFile>& 
flowFile, const std::shared_ptr<io::OutputStream>& out) {
+io::IoResult PayloadSerializer::serialize(const 
std::shared_ptr<core::FlowFile>& flowFile, const 
std::shared_ptr<io::OutputStream>& out) {
   return reader_(flowFile, InputStreamPipe{*out});
 }
 
diff --git a/extension-framework/src/utils/file/FileReaderCallback.cpp 
b/extension-framework/src/utils/file/FileReaderCallback.cpp
index 9386d6356..63720b56d 100644
--- a/extension-framework/src/utils/file/FileReaderCallback.cpp
+++ b/extension-framework/src/utils/file/FileReaderCallback.cpp
@@ -31,9 +31,9 @@ FileReaderCallback::FileReaderCallback(std::filesystem::path 
file_path, const si
     logger_(core::logging::LoggerFactory<FileReaderCallback>::getLogger()) {
 }
 
-int64_t FileReaderCallback::operator()(const 
std::shared_ptr<io::OutputStream>& output_stream) const {
+io::IoResult FileReaderCallback::operator()(const 
std::shared_ptr<io::OutputStream>& output_stream) const {
   std::vector<char> buffer(buffer_size_);
-  uint64_t num_bytes_written = 0;
+  size_t num_bytes_written = 0;
 
   std::ifstream input_stream{file_path_, std::ifstream::in | 
std::ifstream::binary};
   if (!input_stream.is_open()) {
@@ -47,14 +47,14 @@ int64_t FileReaderCallback::operator()(const 
std::shared_ptr<io::OutputStream>&
     }
     const auto num_bytes_read = input_stream.gcount();
     logger_->log_trace("Read {} bytes of input", 
std::intmax_t{num_bytes_read});
-    const auto len = gsl::narrow<size_t>(num_bytes_read);
+    const auto len = num_bytes_read;
     output_stream->write(reinterpret_cast<uint8_t*>(buffer.data()), len);
     num_bytes_written += len;
   }
   input_stream.close();
 
   logger_->log_debug("Finished reading {} bytes from the file", 
num_bytes_written);
-  return gsl::narrow<int64_t>(num_bytes_written);
+  return io::IoResult::from(num_bytes_written);
 }
 
 }  // namespace org::apache::nifi::minifi::utils
diff --git a/extension-framework/src/utils/file/FileWriterCallback.cpp 
b/extension-framework/src/utils/file/FileWriterCallback.cpp
index 7abfe8792..6e7ba1ddc 100644
--- a/extension-framework/src/utils/file/FileWriterCallback.cpp
+++ b/extension-framework/src/utils/file/FileWriterCallback.cpp
@@ -34,7 +34,7 @@ FileWriterCallback::~FileWriterCallback() {
   std::filesystem::remove(temp_path_, remove_error);
 }
 
-int64_t FileWriterCallback::operator()(const std::shared_ptr<io::InputStream>& 
stream) {
+io::IoResult FileWriterCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) {
   write_succeeded_ = false;
   size_t size = 0;
   std::array<std::byte, 1024> buffer{};
@@ -43,7 +43,7 @@ int64_t FileWriterCallback::operator()(const 
std::shared_ptr<io::InputStream>& s
 
   do {
     const auto read = stream->read(buffer);
-    if (io::isError(read)) return -1;
+    if (io::isError(read)) return io::IoResult::error();
     if (read == 0) break;
     tmp_file_os.write(reinterpret_cast<char *>(buffer.data()), 
gsl::narrow<std::streamsize>(read));
     size += read;
@@ -55,7 +55,7 @@ int64_t FileWriterCallback::operator()(const 
std::shared_ptr<io::InputStream>& s
     write_succeeded_ = true;
   }
 
-  return gsl::narrow<int64_t>(size);
+  return io::IoResult::from(size);
 }
 
 bool FileWriterCallback::commit() {
diff --git a/extensions/aws/processors/FetchS3Object.cpp 
b/extensions/aws/processors/FetchS3Object.cpp
index e2a12ce45..0dd9d24f0 100644
--- a/extensions/aws/processors/FetchS3Object.cpp
+++ b/extensions/aws/processors/FetchS3Object.cpp
@@ -92,9 +92,10 @@ void FetchS3Object::onTrigger(core::ProcessContext& context, 
core::ProcessSessio
   }
 
   std::optional<minifi::aws::s3::GetObjectResult> result;
-  session.write(flow_file, [&get_object_params, &result, this](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
+  session.write(flow_file, [&get_object_params, &result, this](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
     result = s3_wrapper_->getObject(*get_object_params, *stream);
-    return (result | 
minifi::utils::transform(&s3::GetObjectResult::write_size)).value_or(0);
+    const auto ret = (result | 
minifi::utils::transform(&s3::GetObjectResult::write_size)).value_or(0);
+    return io::IoResult::from(ret);
   });
 
   if (result) {
diff --git a/extensions/aws/processors/PutS3Object.cpp 
b/extensions/aws/processors/PutS3Object.cpp
index 5fb521892..b896e2c78 100644
--- a/extensions/aws/processors/PutS3Object.cpp
+++ b/extensions/aws/processors/PutS3Object.cpp
@@ -283,20 +283,19 @@ void PutS3Object::onTrigger(core::ProcessContext& 
context, core::ProcessSession&
   }
 
   std::optional<minifi::aws::s3::PutObjectResult> result;
-  session.read(flow_file, [this, &flow_file, &put_s3_request_params, 
&result](const std::shared_ptr<io::InputStream>& stream) -> int64_t {
+  session.read(flow_file, [this, &flow_file, &put_s3_request_params, 
&result](const std::shared_ptr<io::InputStream>& stream) -> io::IoResult {
     try {
       if (flow_file->getSize() <= multipart_threshold_) {
         logger_->log_info("Uploading S3 Object '{}' in a single upload", 
put_s3_request_params->object_key);
         result = s3_wrapper_->putObject(*put_s3_request_params, stream, 
flow_file->getSize());
-        return gsl::narrow<int64_t>(flow_file->getSize());
-      } else {
-        logger_->log_info("S3 Object '{}' passes the multipart threshold, 
uploading it in multiple parts", put_s3_request_params->object_key);
-        result = s3_wrapper_->putObjectMultipart(*put_s3_request_params, 
stream, flow_file->getSize(), multipart_size_);
-        return gsl::narrow<int64_t>(flow_file->getSize());
+        return io::IoResult::from(flow_file->getSize());
       }
+      logger_->log_info("S3 Object '{}' passes the multipart threshold, 
uploading it in multiple parts", put_s3_request_params->object_key);
+      result = s3_wrapper_->putObjectMultipart(*put_s3_request_params, stream, 
flow_file->getSize(), multipart_size_);
+      return io::IoResult::from(flow_file->getSize());
     } catch(const aws::s3::StreamReadException& ex) {
       logger_->log_error("Error occurred while uploading to S3: {}", 
ex.what());
-      return -1;
+      return io::IoResult::error();
     }
   });
   if (!result.has_value()) {
diff --git a/extensions/azure/processors/FetchAzureBlobStorage.cpp 
b/extensions/azure/processors/FetchAzureBlobStorage.cpp
index b96da5b5c..4b5fec14f 100644
--- a/extensions/azure/processors/FetchAzureBlobStorage.cpp
+++ b/extensions/azure/processors/FetchAzureBlobStorage.cpp
@@ -68,12 +68,12 @@ void FetchAzureBlobStorage::onTrigger(core::ProcessContext& 
context, core::Proce
 
   auto fetched_flow_file = session.create(flow_file.get());
   std::optional<int64_t> result_size;
-  session.write(fetched_flow_file, [&, this](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
+  session.write(fetched_flow_file, [&, this](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
     result_size = azure_blob_storage_.fetchBlob(*params, *stream);
     if (!result_size) {
-      return 0;
+      return io::IoResult::zero();
     }
-    return gsl::narrow<int64_t>(*result_size);
+    return io::IoResult::from(*result_size);
   });
 
   if (result_size == std::nullopt) {
diff --git a/extensions/azure/processors/FetchAzureDataLakeStorage.cpp 
b/extensions/azure/processors/FetchAzureDataLakeStorage.cpp
index fa51c21f4..f5500fc53 100644
--- a/extensions/azure/processors/FetchAzureDataLakeStorage.cpp
+++ b/extensions/azure/processors/FetchAzureDataLakeStorage.cpp
@@ -73,13 +73,9 @@ void 
FetchAzureDataLakeStorage::onTrigger(core::ProcessContext& context, core::P
 
   auto fetched_flow_file = session.create(flow_file.get());
   std::optional<uint64_t> result;
-  session.write(fetched_flow_file, [&, this](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
+  session.write(fetched_flow_file, [&, this](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
     result = azure_data_lake_storage_.fetchFile(*params, *output_stream);
-    if (!result) {
-      return 0;
-    }
-
-    return gsl::narrow<int64_t>(*result);
+    return io::IoResult::from(result.value_or(0));
   });
 
   if (result == std::nullopt) {
diff --git a/extensions/azure/processors/PutAzureBlobStorage.h 
b/extensions/azure/processors/PutAzureBlobStorage.h
index 06d3dfaa5..662247452 100644
--- a/extensions/azure/processors/PutAzureBlobStorage.h
+++ b/extensions/azure/processors/PutAzureBlobStorage.h
@@ -81,16 +81,16 @@ class PutAzureBlobStorage final : public 
AzureBlobStorageSingleBlobProcessorBase
       , params_(params) {
     }
 
-    int64_t operator()(const std::shared_ptr<io::InputStream>& stream) {
+    io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) {
       std::vector<std::byte> buffer;
       buffer.resize(flow_size_);
-      size_t read_ret = stream->read(buffer);
+      const size_t read_ret = stream->read(buffer);
       if (io::isError(read_ret) || read_ret != flow_size_) {
-        return -1;
+        return io::IoResult::error();
       }
 
       result_ = azure_blob_storage_.uploadBlob(params_, buffer);
-      return read_ret;
+      return io::IoResult::from(read_ret);
     }
 
     std::optional<storage::UploadBlobResult> getResult() const {
diff --git a/extensions/azure/processors/PutAzureDataLakeStorage.cpp 
b/extensions/azure/processors/PutAzureDataLakeStorage.cpp
index 95c591d84..44c9ad3d7 100644
--- a/extensions/azure/processors/PutAzureDataLakeStorage.cpp
+++ b/extensions/azure/processors/PutAzureDataLakeStorage.cpp
@@ -103,16 +103,16 @@ PutAzureDataLakeStorage::ReadCallback::ReadCallback(
     logger_(std::move(logger)) {
 }
 
-int64_t PutAzureDataLakeStorage::ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) {
+io::IoResult PutAzureDataLakeStorage::ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) {
   std::vector<std::byte> buffer;
   buffer.resize(flow_size_);
   size_t read_ret = stream->read(buffer);
   if (io::isError(read_ret) || read_ret != flow_size_) {
-    return -1;
+    return io::IoResult::error();
   }
 
   result_ = azure_data_lake_storage_.uploadFile(params_, buffer);
-  return gsl::narrow<int64_t>(read_ret);
+  return io::IoResult::from(read_ret);
 }
 
 REGISTER_RESOURCE(PutAzureDataLakeStorage, Processor);
diff --git a/extensions/azure/processors/PutAzureDataLakeStorage.h 
b/extensions/azure/processors/PutAzureDataLakeStorage.h
index 4cdd012c8..9b8f23c6a 100644
--- a/extensions/azure/processors/PutAzureDataLakeStorage.h
+++ b/extensions/azure/processors/PutAzureDataLakeStorage.h
@@ -81,7 +81,7 @@ class PutAzureDataLakeStorage final : public 
AzureDataLakeStorageFileProcessorBa
   class ReadCallback {
    public:
     ReadCallback(uint64_t flow_size, storage::AzureDataLakeStorage 
&azure_data_lake_storage, const storage::PutAzureDataLakeStorageParameters 
&params, std::shared_ptr<core::logging::Logger> logger);
-    int64_t operator()(const std::shared_ptr<io::InputStream> &stream);
+    io::IoResult operator()(const std::shared_ptr<io::InputStream> &stream);
 
     [[nodiscard]] storage::UploadDataLakeStorageResult getResult() const {
       return result_;
diff --git a/extensions/azure/storage/AzureBlobStorage.cpp 
b/extensions/azure/storage/AzureBlobStorage.cpp
index 94521e9c7..c3449a135 100644
--- a/extensions/azure/storage/AzureBlobStorage.cpp
+++ b/extensions/azure/storage/AzureBlobStorage.cpp
@@ -78,7 +78,7 @@ bool AzureBlobStorage::deleteBlob(const 
DeleteAzureBlobStorageParameters& params
 std::optional<uint64_t> AzureBlobStorage::fetchBlob(const 
FetchAzureBlobStorageParameters& params, io::OutputStream& stream) {
   try {
     auto fetch_res = blob_storage_client_->fetchBlob(params);
-    return internal::pipe(*fetch_res, stream);
+    return internal::pipe(*fetch_res, stream).inner() | utils::toOptional();
   } catch (const std::exception& ex) {
     logger_->log_error("An exception occurred while fetching blob '{}' of 
container '{}': {}", params.blob_name, params.container_name, ex.what());
     return std::nullopt;
diff --git a/extensions/azure/storage/AzureDataLakeStorage.cpp 
b/extensions/azure/storage/AzureDataLakeStorage.cpp
index cb55dc1a8..939083991 100644
--- a/extensions/azure/storage/AzureDataLakeStorage.cpp
+++ b/extensions/azure/storage/AzureDataLakeStorage.cpp
@@ -95,7 +95,7 @@ bool AzureDataLakeStorage::deleteFile(const 
DeleteAzureDataLakeStorageParameters
 std::optional<uint64_t> AzureDataLakeStorage::fetchFile(const 
FetchAzureDataLakeStorageParameters& params, io::OutputStream& stream) {
   try {
     auto result = data_lake_storage_client_->fetchFile(params);
-    return internal::pipe(*result, stream);
+    return internal::pipe(*result, stream).inner() | utils::toOptional();
   } catch (const std::exception& ex) {
     logger_->log_error("An exception occurred while fetching '{}/{}' of 
filesystem '{}': {}", params.directory_name, params.filename, 
params.file_system_name, ex.what());
     return std::nullopt;
diff --git a/extensions/bustache/ApplyTemplate.cpp 
b/extensions/bustache/ApplyTemplate.cpp
index b178a2649..81af9c741 100644
--- a/extensions/bustache/ApplyTemplate.cpp
+++ b/extensions/bustache/ApplyTemplate.cpp
@@ -43,7 +43,7 @@ void ApplyTemplate::onTrigger(core::ProcessContext& context, 
core::ProcessSessio
   }
 
   std::string template_file = context.getProperty(Template.name, 
flow_file.get()).value_or("");
-  session.write(flow_file, [&template_file, &flow_file, this](const auto& 
output_stream) {
+  session.write(flow_file, [&template_file, &flow_file, this](const auto& 
output_stream) -> io::IoResult {
     logger_->log_info("ApplyTemplate reading template file from {}", 
template_file);
     // TODO(szaszm): we might want to return to memory-mapped input files when 
the next todo is done. Until then, the agents stores the whole result in memory 
anyway, so no point in not doing the same
     // with the template file itself
@@ -62,7 +62,7 @@ void ApplyTemplate::onTrigger(core::ProcessContext& context, 
core::ProcessSessio
     // TODO(calebj) write ostream reciever for format() to prevent excessive 
copying
     std::string ostring = bustache::to_string(format(data));
     output_stream->write(gsl::make_span(ostring).as_span<const std::byte>());
-    return gsl::narrow<int64_t>(ostring.length());
+    return io::IoResult::from(ostring.length());
   });
   session.transfer(flow_file, Success);
 }
diff --git a/extensions/couchbase/processors/GetCouchbaseKey.cpp 
b/extensions/couchbase/processors/GetCouchbaseKey.cpp
index 71744c13a..681a0c2b9 100644
--- a/extensions/couchbase/processors/GetCouchbaseKey.cpp
+++ b/extensions/couchbase/processors/GetCouchbaseKey.cpp
@@ -72,15 +72,15 @@ void GetCouchbaseKey::onTrigger(core::ProcessContext& 
context, core::ProcessSess
         session.putAttribute(*flow_file, attribute_to_put_result_to, 
str_value);
       }
     } else {
-      session.write(flow_file, [&, this](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
+      session.write(flow_file, [&, this](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
         if (document_type_ == CouchbaseValueType::String) {
           auto& value = std::get<std::string>(get_result->value);
           stream->write(value);
-          return gsl::narrow<int64_t>(value.size());
+          return io::IoResult::from(value.size());
         } else {
           auto& value = std::get<std::vector<std::byte>>(get_result->value);
           stream->write(value);
-          return gsl::narrow<int64_t>(value.size());
+          return io::IoResult::from(value.size());
         }
       });
     }
diff --git a/extensions/gcp/processors/FetchGCSObject.cpp 
b/extensions/gcp/processors/FetchGCSObject.cpp
index a42799e28..32e1eb4c7 100644
--- a/extensions/gcp/processors/FetchGCSObject.cpp
+++ b/extensions/gcp/processors/FetchGCSObject.cpp
@@ -37,7 +37,7 @@ class FetchFromGCSCallback {
         client_(client) {
   }
 
-  int64_t operator()(const std::shared_ptr<io::OutputStream>& stream) {
+  io::IoResult operator()(const std::shared_ptr<io::OutputStream>& stream) {
     auto reader = client_.ReadObject(bucket_, key_, encryption_key_, 
generation_, gcs::IfGenerationNotMatch(0));
     auto set_members = gsl::finally([&]{
       status_ = reader.status();
@@ -46,11 +46,11 @@ class FetchFromGCSCallback {
       storage_class_ = reader.storage_class();
     });
     if (!reader)
-      return 0;
+      return io::IoResult::zero();
     std::string contents{std::istreambuf_iterator<char>{reader}, {}};
-    const auto write_ret = 
gsl::narrow<int64_t>(stream->write(gsl::make_span(contents).as_span<std::byte>()));
+    const auto write_ret = 
stream->write(gsl::make_span(contents).as_span<std::byte>());
     reader.Close();
-    return write_ret;
+    return io::IoResult::from(write_ret);
   }
 
   [[nodiscard]] auto getStatus() const noexcept { return status_; }
diff --git a/extensions/gcp/processors/PutGCSObject.cpp 
b/extensions/gcp/processors/PutGCSObject.cpp
index ceea47658..a05380f32 100644
--- a/extensions/gcp/processors/PutGCSObject.cpp
+++ b/extensions/gcp/processors/PutGCSObject.cpp
@@ -38,18 +38,18 @@ class UploadToGCSCallback {
         client_(client) {
   }
 
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream) {
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) {
     std::string content;
     content.resize(stream->size());
     const auto read_ret = stream->read(as_writable_bytes(std::span(content)));
     if (io::isError(read_ret)) {
-      return -1;
+      return io::IoResult::error();
     }
     auto writer = client_.WriteObject(bucket_, key_, hash_value_, 
crc32c_checksum_, encryption_key_, content_type_, predefined_acl_, 
if_generation_match_);
     writer << content;
     writer.Close();
     result_ = writer.metadata();
-    return gsl::narrow<int64_t>(read_ret);
+    return io::IoResult::from(read_ret);
   }
 
   [[nodiscard]] const google::cloud::StatusOr<gcs::ObjectMetadata>& 
getResult() const noexcept {
diff --git a/extensions/kafka/PublishKafka.cpp 
b/extensions/kafka/PublishKafka.cpp
index 7c6c0912c..6f04186cd 100644
--- a/extensions/kafka/PublishKafka.cpp
+++ b/extensions/kafka/PublishKafka.cpp
@@ -248,7 +248,7 @@ class ReadCallback {
   ReadCallback& operator=(const ReadCallback&) = delete;
   ~ReadCallback() = default;
 
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream) {
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) {
     std::vector<std::byte> buffer;
 
     buffer.resize(max_seg_size_);
@@ -269,7 +269,7 @@ class ReadCallback {
         status_ = -1;
         error_ = rd_kafka_err2str(err);
       }
-      return 0;
+      return io::IoResult::zero();
     }
 
     for (size_t segment_num = 0; read_size_ < flow_size_; ++segment_num) {
@@ -277,7 +277,7 @@ class ReadCallback {
       if (io::isError(readRet)) {
         status_ = -1;
         error_ = "Failed to read from stream";
-        return read_size_;
+        return io::IoResult::zero();
       }
       if (readRet == 0) { break; }
 
@@ -289,11 +289,11 @@ class ReadCallback {
         });
         status_ = -1;
         error_ = rd_kafka_err2str(err);
-        return read_size_;
+        return io::IoResult::zero();
       }
       read_size_ += gsl::narrow<uint32_t>(readRet);
     }
-    return read_size_;
+    return io::IoResult::from(gsl::narrow<uint64_t>(read_size_));
   }
 
   const uint64_t flow_size_ = 0;
diff --git a/extensions/libarchive/CompressContent.cpp 
b/extensions/libarchive/CompressContent.cpp
index aad392ab1..039fd699a 100644
--- a/extensions/libarchive/CompressContent.cpp
+++ b/extensions/libarchive/CompressContent.cpp
@@ -131,37 +131,37 @@ void CompressContent::processFlowFile(const 
std::shared_ptr<core::FlowFile>& flo
   std::shared_ptr<core::FlowFile> result = session.create(flowFile.get());
   bool success = true;
   if (encapsulateInTar_) {
-    std::function<int64_t(const std::shared_ptr<io::InputStream>&, const 
std::shared_ptr<io::OutputStream>&)> transformer;
+    std::function<io::IoResult(const std::shared_ptr<io::InputStream>&, const 
std::shared_ptr<io::OutputStream>&)> transformer;
 
     if (compressMode_ == compress_content::CompressionMode::compress) {
       std::string filename;
       flowFile->getAttribute(core::SpecialFlowAttribute::FILENAME, filename);
-      transformer = [&, filename] (const std::shared_ptr<io::InputStream>& in, 
const std::shared_ptr<io::OutputStream>& out) -> int64_t {
+      transformer = [&, filename] (const std::shared_ptr<io::InputStream>& in, 
const std::shared_ptr<io::OutputStream>& out) -> io::IoResult {
         io::WriteArchiveStreamImpl compressor(compressLevel_, compressFormat, 
out);
         if (!compressor.newEntry({filename, in->size()})) {
-          return -1;
+          return io::IoResult::error();
         }
         return internal::pipe(*in, compressor);
       };
     } else {
-      transformer = [&] (const std::shared_ptr<io::InputStream>& in, const 
std::shared_ptr<io::OutputStream>& out) -> int64_t {
+      transformer = [&] (const std::shared_ptr<io::InputStream>& in, const 
std::shared_ptr<io::OutputStream>& out) -> io::IoResult {
         io::ReadArchiveStreamImpl decompressor(in);
         if (!decompressor.nextEntry()) {
           success = false;
-          return 0;  // prevents a session rollback
+          return io::IoResult::zero();  // prevents a session rollback
         }
         auto ret = internal::pipe(decompressor, *out);
-        if (ret < 0) {
+        if (!ret) {
           success = false;
-          return 0;  // prevents a session rollback
+          return io::IoResult::zero();  // prevents a session rollback
         }
         return ret;
       };
     }
     session.write(result, [&] (const auto& out) {
-      return session.read(flowFile, [&] (const auto& in) {
+      return io::IoResult::from(session.read(flowFile, [&] (const auto& in) {
         return transformer(in, out);
-      });
+      }));
     });
   } else {
     CompressContent::GzipWriteCallback callback(compressMode_, compressLevel_, 
flowFile, session);
diff --git a/extensions/libarchive/CompressContent.h 
b/extensions/libarchive/CompressContent.h
index 08667be80..59c53e52f 100644
--- a/extensions/libarchive/CompressContent.h
+++ b/extensions/libarchive/CompressContent.h
@@ -168,37 +168,37 @@ class CompressContent : public core::ProcessorImpl {
     core::ProcessSession& session_;
     bool success_{false};
 
-    int64_t operator()(const std::shared_ptr<io::OutputStream>& output_stream) 
{
+    io::IoResult operator()(const std::shared_ptr<io::OutputStream>& 
output_stream) {
       std::shared_ptr<io::ZlibBaseStream> filterStream;
       if (compress_mode_ == compress_content::CompressionMode::compress) {
         filterStream = 
std::make_shared<io::ZlibCompressStream>(gsl::make_not_null(output_stream.get()),
 io::ZlibCompressionFormat::GZIP, compress_level_);
       } else {
         filterStream = 
std::make_shared<io::ZlibDecompressStream>(gsl::make_not_null(output_stream.get()),
 io::ZlibCompressionFormat::GZIP);
       }
-      session_.read(flow_, [this, &filterStream](const 
std::shared_ptr<io::InputStream>& input_stream) -> int64_t {
+      session_.read(flow_, [this, &filterStream](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
         std::vector<std::byte> buffer(16 * 1024U);
         size_t read_size = 0;
         while (read_size < flow_->getSize()) {
           const auto ret = input_stream->read(buffer);
           if (io::isError(ret)) {
-            return -1;
+            return io::IoResult::error();
           } else if (ret == 0) {
             break;
           } else {
             const auto writeret = 
filterStream->write(gsl::make_span(buffer).subspan(0, ret));
             if (io::isError(writeret) || gsl::narrow<size_t>(writeret) != ret) 
{
-              return -1;
+              return io::IoResult::error();
             }
             read_size += ret;
           }
         }
         filterStream->close();
-        return gsl::narrow<int64_t>(read_size);
+        return io::IoResult::from(read_size);
       });
 
       success_ = filterStream->isFinished();
 
-      return gsl::narrow<int64_t>(flow_->getSize());
+      return io::IoResult::from(flow_->getSize());
     }
   };
 
diff --git a/extensions/libarchive/FocusArchiveEntry.cpp 
b/extensions/libarchive/FocusArchiveEntry.cpp
index acc301194..0a0a12387 100644
--- a/extensions/libarchive/FocusArchiveEntry.cpp
+++ b/extensions/libarchive/FocusArchiveEntry.cpp
@@ -154,7 +154,7 @@ la_ssize_t FocusArchiveEntry::ReadCallback::read_cb(struct 
archive * a, void *d,
   return gsl::narrow<la_ssize_t>(read);
 }
 
-int64_t FocusArchiveEntry::ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) const {
+io::IoResult FocusArchiveEntry::ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) const {
   auto input_archive = processors::archive_read_unique_ptr{archive_read_new()};
   struct archive_entry *entry = nullptr;
   int64_t nlen = 0;
@@ -169,7 +169,7 @@ int64_t FocusArchiveEntry::ReadCallback::operator()(const 
std::shared_ptr<io::In
   // Read each item in the archive
   if (archive_read_open(input_archive.get(), &data, ok_cb, read_cb, ok_cb)) {
     logger_->log_error("FocusArchiveEntry can't open due to archive error: 
{}", archive_error_string(input_archive.get()));
-    return nlen;
+    return io::IoResult::zero();
   }
 
   while (context_->isRunning()) {
@@ -179,14 +179,9 @@ int64_t FocusArchiveEntry::ReadCallback::operator()(const 
std::shared_ptr<io::In
       break;
     }
 
-    if (res < ARCHIVE_OK) {
+    if (res < ARCHIVE_OK) {  // TODO(MINIFICPP-2761)
       logger_->log_error("FocusArchiveEntry can't read header due to archive 
error: {}", archive_error_string(input_archive.get()));
-      return nlen;
-    }
-
-    if (res < ARCHIVE_WARN) {
-      logger_->log_warn("FocusArchiveEntry got archive warning while reading 
header: {}", archive_error_string(input_archive.get()));
-      return nlen;
+      return io::IoResult::from(nlen);
     }
 
     auto entryName = archive_entry_pathname(entry);
@@ -232,7 +227,7 @@ int64_t FocusArchiveEntry::ReadCallback::operator()(const 
std::shared_ptr<io::In
     _archiveMetadata->entryMetadata.push_back(metadata);
   }
 
-  return nlen;
+  return io::IoResult::from(nlen);
 }
 
 FocusArchiveEntry::ReadCallback::ReadCallback(core::ProcessContext* context, 
utils::file::FileManager *file_man, ArchiveMetadata *archiveMetadata)
diff --git a/extensions/libarchive/FocusArchiveEntry.h 
b/extensions/libarchive/FocusArchiveEntry.h
index bcf816cb3..a0e67c958 100644
--- a/extensions/libarchive/FocusArchiveEntry.h
+++ b/extensions/libarchive/FocusArchiveEntry.h
@@ -66,7 +66,7 @@ class FocusArchiveEntry : public core::ProcessorImpl {
   class ReadCallback {
    public:
     explicit ReadCallback(core::ProcessContext* context, 
utils::file::FileManager *file_man, ArchiveMetadata *archiveMetadata);
-    int64_t operator()(const std::shared_ptr<io::InputStream>& stream) const;
+    io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) 
const;
 
    private:
     utils::file::FileManager *file_man_;
diff --git a/extensions/libarchive/MergeContent.cpp 
b/extensions/libarchive/MergeContent.cpp
index 90bff039f..428d8c5e4 100644
--- a/extensions/libarchive/MergeContent.cpp
+++ b/extensions/libarchive/MergeContent.cpp
@@ -231,8 +231,8 @@ bool MergeContent::processBin(core::ProcessSession 
&session, std::unique_ptr<Bin
     return false;
   }
 
-  auto flowFileReader = [&] (const std::shared_ptr<core::FlowFile>& ff, const 
io::InputStreamCallback& cb) {
-    return session.read(ff, cb);
+  auto flowFileReader = [&] (const std::shared_ptr<core::FlowFile>& ff, const 
io::InputStreamCallback& cb) -> io::IoResult {
+    return io::IoResult::from(session.read(ff, cb));
   };
 
   const char* mimeType = nullptr;
diff --git a/extensions/libarchive/MergeContent.h 
b/extensions/libarchive/MergeContent.h
index 6b7330d01..cbdfb8682 100644
--- a/extensions/libarchive/MergeContent.h
+++ b/extensions/libarchive/MergeContent.h
@@ -82,35 +82,39 @@ class BinaryConcatenationMerge : public MergeBin {
     std::deque<std::shared_ptr<core::FlowFile>> &flows_;
     FlowFileSerializer& serializer_;
 
-    int64_t operator()(const std::shared_ptr<io::OutputStream>& stream) const {
+    io::IoResult operator()(const std::shared_ptr<io::OutputStream>& stream) 
const {
       size_t write_size_sum = 0;
       if (!header_.empty()) {
         const auto write_ret = stream->write(reinterpret_cast<const 
uint8_t*>(header_.data()), header_.size());
-        if (io::isError(write_ret))
-          return -1;
+        if (io::isError(write_ret)) {
+          return io::IoResult::error();
+        }
         write_size_sum += write_ret;
       }
       bool isFirst = true;
       for (const auto& flow : flows_) {
         if (!isFirst && !demarcator_.empty()) {
-          const auto write_ret = stream->write(reinterpret_cast<const 
uint8_t*>(demarcator_.data()), demarcator_.size());
-          if (io::isError(write_ret))
-            return -1;
-          write_size_sum += write_ret;
+          const auto write_result = stream->write(reinterpret_cast<const 
uint8_t*>(demarcator_.data()), demarcator_.size());
+          if (io::isError(write_result)) {
+            return io::IoResult::error();
+          }
+          write_size_sum += write_result;
+        }
+        const auto serialization_result = serializer_.serialize(flow, stream);
+        if (!serialization_result) {
+          return serialization_result;
         }
-        const auto len = serializer_.serialize(flow, stream);
-        if (len < 0)
-          return len;
-        write_size_sum += gsl::narrow<size_t>(len);
+        write_size_sum += gsl::narrow<size_t>(*serialization_result);
         isFirst = false;
       }
       if (!footer_.empty()) {
-        const auto write_ret = stream->write(reinterpret_cast<const 
uint8_t*>(footer_.data()), footer_.size());
-        if (io::isError(write_ret))
-          return -1;
-        write_size_sum += write_ret;
+        const auto write_result = stream->write(reinterpret_cast<const 
uint8_t*>(footer_.data()), footer_.size());
+        if (io::isError(write_result)) {
+          return io::IoResult::error();
+        }
+        write_size_sum += write_result;
       }
-      return gsl::narrow<int64_t>(write_size_sum);
+      return io::IoResult::from(write_size_sum);
     }
   };
 
@@ -192,7 +196,7 @@ class ArchiveMerge {
       return totalWrote;
     }
 
-    int64_t operator()(const std::shared_ptr<io::OutputStream>& stream) {
+    io::IoResult operator()(const std::shared_ptr<io::OutputStream>& stream) {
       const auto arch = archive_write_unique_ptr{archive_write_new()};
 
       if (merge_type_ == merge_content_options::MERGE_FORMAT_TAR_VALUE) {
@@ -224,12 +228,12 @@ class ArchiveMerge {
           }
         }
         const auto ret = serializer_.serialize(flow, 
std::make_shared<ArchiveWriter>(*arch, *entry));
-        if (ret < 0) {
+        if (!ret) {
           return ret;
         }
       }
 
-      return gsl::narrow<int64_t>(size_);
+      return io::IoResult::from(size_);
     }
   };
 };
diff --git a/extensions/libarchive/UnfocusArchiveEntry.cpp 
b/extensions/libarchive/UnfocusArchiveEntry.cpp
index eb94b693f..b8de55134 100644
--- a/extensions/libarchive/UnfocusArchiveEntry.cpp
+++ b/extensions/libarchive/UnfocusArchiveEntry.cpp
@@ -150,7 +150,7 @@ la_ssize_t 
UnfocusArchiveEntry::WriteCallback::write_cb(struct archive *, void *
   return io::isError(write_ret) ? -1 : gsl::narrow<la_ssize_t>(write_ret);
 }
 
-int64_t UnfocusArchiveEntry::WriteCallback::operator()(const 
std::shared_ptr<io::OutputStream>& stream) const {
+io::IoResult UnfocusArchiveEntry::WriteCallback::operator()(const 
std::shared_ptr<io::OutputStream>& stream) const {
   UnfocusArchiveEntryWriteData data;
   data.stream = stream;
   auto output_archive = archive_write_unique_ptr{archive_write_new()};
@@ -218,7 +218,7 @@ int64_t 
UnfocusArchiveEntry::WriteCallback::operator()(const std::shared_ptr<io:
     archive_entry_clear(entry.get());
   }
 
-  return nlen;
+  return io::IoResult::from(nlen);
 }
 
 REGISTER_RESOURCE(UnfocusArchiveEntry, Processor);
diff --git a/extensions/libarchive/UnfocusArchiveEntry.h 
b/extensions/libarchive/UnfocusArchiveEntry.h
index 5cafd93a5..2423c95ed 100644
--- a/extensions/libarchive/UnfocusArchiveEntry.h
+++ b/extensions/libarchive/UnfocusArchiveEntry.h
@@ -64,7 +64,7 @@ class UnfocusArchiveEntry : public core::ProcessorImpl {
   class WriteCallback {
    public:
     explicit WriteCallback(ArchiveMetadata *archiveMetadata);
-    int64_t operator()(const std::shared_ptr<io::OutputStream>& stream) const;
+    io::IoResult operator()(const std::shared_ptr<io::OutputStream>& stream) 
const;
    private:
     //! Logger
     std::shared_ptr<Logger> logger_ = 
core::logging::LoggerFactory<UnfocusArchiveEntry>::getLogger();
diff --git a/extensions/libarchive/tests/CompressContentTests.cpp 
b/extensions/libarchive/tests/CompressContentTests.cpp
index c2d38415a..3c7317c4c 100644
--- a/extensions/libarchive/tests/CompressContentTests.cpp
+++ b/extensions/libarchive/tests/CompressContentTests.cpp
@@ -56,16 +56,16 @@ class ReadCallback {
   ReadCallback& operator=(ReadCallback&&) = delete;
   ~ReadCallback() = default;
 
-  int64_t operator()(const std::shared_ptr<minifi::io::InputStream>& stream) {
+  minifi::io::IoResult operator()(const 
std::shared_ptr<minifi::io::InputStream>& stream) {
     int64_t total_read = 0;
     do {
       const auto ret = stream->read(std::span(buffer_).subspan(read_size_));
       if (ret == 0) { break; }
-      if (minifi::io::isError(ret)) { return -1; }
+      if (minifi::io::isError(ret)) { return  minifi::io::IoResult::error(); }
       read_size_ += gsl::narrow<size_t>(ret);
       total_read += gsl::narrow<int64_t>(ret);
     } while (buffer_.size() != read_size_);
-    return total_read;
+    return minifi::io::IoResult::from(total_read);
   }
   void archive_read() {
     const auto archive = 
minifi::processors::archive_read_unique_ptr{archive_read_new()};
diff --git a/extensions/libarchive/tests/MergeFileTests.cpp 
b/extensions/libarchive/tests/MergeFileTests.cpp
index 651ac0ca9..2ab3fc53a 100644
--- a/extensions/libarchive/tests/MergeFileTests.cpp
+++ b/extensions/libarchive/tests/MergeFileTests.cpp
@@ -108,8 +108,8 @@ class FixedBuffer {
     return gsl::narrow<int64_t>(total_read);
   }
 
-  int64_t operator()(const std::shared_ptr<minifi::io::InputStream>& stream) {
-    return write(*stream, capacity_);
+  minifi::io::IoResult operator()(const 
std::shared_ptr<minifi::io::InputStream>& stream) {
+    return minifi::io::IoResult::from(write(*stream, capacity_));
   }
 
  private:
@@ -712,11 +712,11 @@ TEST_CASE("FlowFile serialization", 
"[testFlowFileSerialization]") {
 
   core::ProcessSessionImpl session(context);
 
-  minifi::PayloadSerializer payloadSerializer([&] (const 
std::shared_ptr<core::FlowFile>& ff, const minifi::io::InputStreamCallback& cb) 
{
-    return session.read(ff, cb);
+  minifi::PayloadSerializer payloadSerializer([&] (const 
std::shared_ptr<core::FlowFile>& ff, const minifi::io::InputStreamCallback& cb) 
-> minifi::io::IoResult {
+    return minifi::io::IoResult::from(session.read(ff, cb));
   });
-  minifi::FlowFileV3Serializer ffV3Serializer([&] (const 
std::shared_ptr<core::FlowFile>& ff, const minifi::io::InputStreamCallback& cb) 
{
-    return session.read(ff, cb);
+  minifi::FlowFileV3Serializer ffV3Serializer([&] (const 
std::shared_ptr<core::FlowFile>& ff, const minifi::io::InputStreamCallback& cb) 
-> minifi::io::IoResult {
+    return minifi::io::IoResult::from(session.read(ff, cb));
   });
 
   minifi::FlowFileSerializer* usedSerializer = nullptr;
diff --git a/extensions/lua/LuaProcessSession.cpp 
b/extensions/lua/LuaProcessSession.cpp
index 3a7f3a48d..bc9cd4062 100644
--- a/extensions/lua/LuaProcessSession.cpp
+++ b/extensions/lua/LuaProcessSession.cpp
@@ -58,9 +58,9 @@ void LuaProcessSession::read(const 
std::shared_ptr<LuaScriptFlowFile> &script_fl
     throw std::runtime_error("Access of FlowFile after it has been released");
   }
 
-  session_.read(flow_file, [&input_stream_callback](const 
std::shared_ptr<io::InputStream>& input_stream) -> int64_t {
-    sol::function callback = input_stream_callback["process"];
-    return callback(input_stream_callback, 
std::make_shared<LuaInputStream>(input_stream));
+  session_.read(flow_file, [&input_stream_callback](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
+    const int64_t callback_result = 
input_stream_callback["process"](input_stream_callback, 
std::make_shared<LuaInputStream>(input_stream));
+    return io::IoResult::from(callback_result);
   });
 }
 
@@ -72,9 +72,9 @@ void LuaProcessSession::write(const 
std::shared_ptr<LuaScriptFlowFile> &script_f
     throw std::runtime_error("Access of FlowFile after it has been released");
   }
 
-  session_.write(flow_file, [&output_stream_callback](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
-    sol::function callback = output_stream_callback["process"];
-    return callback(output_stream_callback, 
std::make_shared<LuaOutputStream>(output_stream));
+  session_.write(flow_file, [&output_stream_callback](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
+    const int64_t callback_result = 
output_stream_callback["process"](output_stream_callback, 
std::make_shared<LuaOutputStream>(output_stream));
+    return io::IoResult::from(callback_result);
   });
 }
 
diff --git a/extensions/mqtt/processors/ConsumeMQTT.cpp 
b/extensions/mqtt/processors/ConsumeMQTT.cpp
index 5b5405a0a..45bbb8025 100644
--- a/extensions/mqtt/processors/ConsumeMQTT.cpp
+++ b/extensions/mqtt/processors/ConsumeMQTT.cpp
@@ -163,21 +163,21 @@ std::queue<ConsumeMQTT::SmartMessage> 
ConsumeMQTT::getReceivedMqttMessages() {
   return msg_queue;
 }
 
-int64_t ConsumeMQTT::WriteCallback::operator() (const 
std::shared_ptr<io::OutputStream>& stream) {
+io::IoResult ConsumeMQTT::WriteCallback::operator() (const 
std::shared_ptr<io::OutputStream>& stream) {
   if (message_.contents->payloadlen < 0) {
     success_status_ = false;
     logger_->log_error("Payload length of message is negative, value is [{}]", 
message_.contents->payloadlen);
-    return -1;
+    return io::IoResult::error();
   }
 
-  const auto len = 
stream->write(reinterpret_cast<uint8_t*>(message_.contents->payload), 
gsl::narrow<size_t>(message_.contents->payloadlen));
+  const size_t len = 
stream->write(reinterpret_cast<uint8_t*>(message_.contents->payload), 
gsl::narrow<size_t>(message_.contents->payloadlen));
   if (io::isError(len)) {
     success_status_ = false;
     logger_->log_error("Stream writing error when processing message");
-    return -1;
+    return io::IoResult::error();
   }
 
-  return gsl::narrow<int64_t>(len);
+  return io::IoResult::from(len);
 }
 
 void ConsumeMQTT::putUserPropertiesAsAttributes(const SmartMessage& message, 
const std::shared_ptr<core::FlowFile>& flow_file, core::ProcessSession& 
session) const {
diff --git a/extensions/mqtt/processors/ConsumeMQTT.h 
b/extensions/mqtt/processors/ConsumeMQTT.h
index b4434d46e..a172fbe7e 100644
--- a/extensions/mqtt/processors/ConsumeMQTT.h
+++ b/extensions/mqtt/processors/ConsumeMQTT.h
@@ -142,7 +142,7 @@ class ConsumeMQTT : public 
processors::AbstractMQTTProcessor {
       , logger_(std::move(logger)) {
     }
 
-    int64_t operator() (const std::shared_ptr<io::OutputStream>& stream);
+    io::IoResult operator() (const std::shared_ptr<io::OutputStream>& stream);
 
     [[nodiscard]] bool getSuccessStatus() const {
       return success_status_;
diff --git a/extensions/mqtt/processors/PublishMQTT.cpp 
b/extensions/mqtt/processors/PublishMQTT.cpp
index 9b8fe786d..24bc72e90 100644
--- a/extensions/mqtt/processors/PublishMQTT.cpp
+++ b/extensions/mqtt/processors/PublishMQTT.cpp
@@ -66,9 +66,9 @@ void PublishMQTT::onTriggerImpl(core::ProcessContext& 
context, core::ProcessSess
   std::vector<std::shared_ptr<core::FlowFile>> flow_files;
   if (record_converter_) {
     nonstd::expected<core::RecordSet, std::error_code> record_set;
-    session.read(original_flow_file, [this, &record_set](const 
std::shared_ptr<io::InputStream>& input_stream) {
+    session.read(original_flow_file, [this, &record_set](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
       record_set = record_converter_->record_set_reader->read(*input_stream);
-      return gsl::narrow<int64_t>(input_stream->size());
+      return io::IoResult::from(input_stream->size());
     });
 
     if (!record_set) {
diff --git a/extensions/opencv/CaptureRTSPFrame.cpp 
b/extensions/opencv/CaptureRTSPFrame.cpp
index cb165baf7..19aa01acd 100644
--- a/extensions/opencv/CaptureRTSPFrame.cpp
+++ b/extensions/opencv/CaptureRTSPFrame.cpp
@@ -94,11 +94,11 @@ void CaptureRTSPFrame::onTrigger(core::ProcessContext& 
context, core::ProcessSes
       session.putAttribute(*flow_file, "filename", filename);
       session.putAttribute(*flow_file, "video.backend.driver", 
video_backend_driver_);
 
-      session.write(flow_file, [&frame, this](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
+      session.write(flow_file, [&frame, this](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
         std::vector<uchar> image_buf;
         imencode(image_encoding_, frame, image_buf);
         const auto ret = output_stream->write(image_buf.data(), 
image_buf.size());
-        return io::isError(ret) ? -1 : gsl::narrow<int64_t>(ret);
+        return io::IoResult::from(ret);
       });
       session.transfer(flow_file, Success);
       logger_->log_info("A frame is captured");
diff --git a/extensions/opencv/MotionDetector.cpp 
b/extensions/opencv/MotionDetector.cpp
index b6728f455..3724ff11e 100644
--- a/extensions/opencv/MotionDetector.cpp
+++ b/extensions/opencv/MotionDetector.cpp
@@ -109,15 +109,15 @@ void MotionDetector::onTrigger(core::ProcessContext& 
context, core::ProcessSessi
   }
   cv::Mat frame;
 
-  session.read(flow_file, [&frame](const std::shared_ptr<io::InputStream>& 
input_stream) -> int64_t {
+  session.read(flow_file, [&frame](const std::shared_ptr<io::InputStream>& 
input_stream) -> io::IoResult {
     std::vector<uchar> image_buf;
     image_buf.resize(input_stream->size());
-    const auto ret = 
input_stream->read(as_writable_bytes(std::span(image_buf)));
+    const size_t ret = 
input_stream->read(as_writable_bytes(std::span(image_buf)));
     if (io::isError(ret) || ret != input_stream->size()) {
       throw std::runtime_error("ImageReadCallback failed to fully read flow 
file input stream");
     }
     frame = cv::imdecode(image_buf, -1);
-    return gsl::narrow<int64_t>(ret);
+    return io::IoResult::from(ret);
   });
 
   if (frame.empty()) {
@@ -148,11 +148,11 @@ void MotionDetector::onTrigger(core::ProcessContext& 
context, core::ProcessSessi
 
   session.putAttribute(*flow_file, "filename", filename);
 
-  session.write(flow_file, [&frame, this](const auto& output_stream) -> 
int64_t {
+  session.write(flow_file, [&frame, this](const auto& output_stream) -> 
io::IoResult {
     std::vector<uchar> image_buf;
     imencode(image_encoding_, frame, image_buf);
     const auto ret = output_stream->write(image_buf.data(), image_buf.size());
-    return io::isError(ret) ? -1 : gsl::narrow<int64_t>(ret);
+    return io::IoResult::from(ret);
   });
   session.transfer(flow_file, Success);
   logger_->log_trace("Finish motion detecting");
diff --git a/extensions/python/types/PyProcessSession.cpp 
b/extensions/python/types/PyProcessSession.cpp
index 9391c93ff..1f153b00c 100644
--- a/extensions/python/types/PyProcessSession.cpp
+++ b/extensions/python/types/PyProcessSession.cpp
@@ -67,8 +67,8 @@ void PyProcessSession::read(const 
std::shared_ptr<core::FlowFile>& flow_file, Bo
     throw std::runtime_error("Access of FlowFile after it has been released");
   }
 
-  session_.read(flow_file, [&input_stream_callback](const 
std::shared_ptr<io::InputStream>& input_stream) -> int64_t {
-    return 
Long(Callable(input_stream_callback.getAttribute("process"))(std::weak_ptr(input_stream))).asInt64();
+  session_.read(flow_file, [&input_stream_callback](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
+    return 
io::IoResult::from(Long(Callable(input_stream_callback.getAttribute("process"))(std::weak_ptr(input_stream))).asInt64());
   });
 }
 
@@ -77,8 +77,8 @@ void PyProcessSession::write(const 
std::shared_ptr<core::FlowFile>& flow_file, B
     throw std::runtime_error("Access of FlowFile after it has been released");
   }
 
-  session_.write(flow_file, [&output_stream_callback](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
-    return 
Long(Callable(output_stream_callback.getAttribute("process"))(std::weak_ptr(output_stream))).asInt64();
+  session_.write(flow_file, [&output_stream_callback](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
+    return 
io::IoResult::from(Long(Callable(output_stream_callback.getAttribute("process"))(std::weak_ptr(output_stream))).asInt64());
   });
 }
 
@@ -111,9 +111,9 @@ std::string PyProcessSession::getContentsAsString(const 
std::shared_ptr<core::Fl
   }
 
   std::string content;
-  session_.read(flow_file, [&content](const std::shared_ptr<io::InputStream>& 
input_stream) -> int64_t {
+  session_.read(flow_file, [&content](const std::shared_ptr<io::InputStream>& 
input_stream) -> io::IoResult {
     content.resize(input_stream->size());
-    return 
gsl::narrow<int64_t>(input_stream->read(as_writable_bytes(std::span(content))));
+    return 
io::IoResult::from(input_stream->read(as_writable_bytes(std::span(content))));
   });
   return content;
 }
diff --git a/extensions/python/types/PyRecordSetReader.cpp 
b/extensions/python/types/PyRecordSetReader.cpp
index 666843074..815d2d642 100644
--- a/extensions/python/types/PyRecordSetReader.cpp
+++ b/extensions/python/types/PyRecordSetReader.cpp
@@ -89,7 +89,7 @@ PyObject* PyRecordSetReader::read(PyRecordSetReader* self, 
PyObject* args) {
   nonstd::expected<core::RecordSet, std::error_code> read_result;
   process_session->getSession().read(flow_file, [&record_set_reader, 
&read_result](const std::shared_ptr<io::InputStream>& input_stream) {
     read_result = record_set_reader->read(*input_stream);
-    return gsl::narrow<int64_t>(input_stream->size());
+    return io::IoResult::from(input_stream->size());
   });
 
   if  (!read_result) {
diff --git a/extensions/rocksdb-repos/tests/SwapTests.cpp 
b/extensions/rocksdb-repos/tests/SwapTests.cpp
index 128fe23a9..82677c0a2 100644
--- a/extensions/rocksdb-repos/tests/SwapTests.cpp
+++ b/extensions/rocksdb-repos/tests/SwapTests.cpp
@@ -56,14 +56,11 @@ class OutputProcessor : public core::ProcessorImpl {
 
   void onTrigger(core::ProcessContext&, core::ProcessSession& session) 
override {
     auto id = std::to_string(next_id_++);
-    auto ff = session.create();
+    const auto ff = session.create();
     ff->addAttribute("index", id);
-    session.write(ff, [&] (const std::shared_ptr<minifi::io::OutputStream>& 
output) -> int64_t {
-      auto ret = output->write(as_bytes(std::span(id)));
-      if (minifi::io::isError(ret)) {
-        return -1;
-      }
-      return gsl::narrow<int64_t>(ret);
+    session.write(ff, [&] (const std::shared_ptr<minifi::io::OutputStream>& 
output) -> io::IoResult {
+      const size_t ret = output->write(as_bytes(std::span(id)));
+      return io::IoResult::from(ret);
     });
     session.transfer(ff, Success);
     flow_files_.push_back(ff);
diff --git a/extensions/sftp/processors/FetchSFTP.cpp 
b/extensions/sftp/processors/FetchSFTP.cpp
index 4ce4c41d0..9bbcae6b9 100644
--- a/extensions/sftp/processors/FetchSFTP.cpp
+++ b/extensions/sftp/processors/FetchSFTP.cpp
@@ -107,12 +107,12 @@ void FetchSFTP::onTrigger(core::ProcessContext& context, 
core::ProcessSession& s
 
   /* Download file */
   try {
-    session.write(flow_file, [&remote_file, &client](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
+    session.write(flow_file, [&remote_file, &client](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
       auto bytes_read = client->getFile(remote_file.generic_string(), *stream);
       if (!bytes_read) {
         throw utils::SFTPException{client->getLastError()};
       }
-      return gsl::narrow<int64_t>(*bytes_read);
+      return io::IoResult::from(*bytes_read);
     });
   } catch (const utils::SFTPException& ex) {
     logger_->log_debug("{}", ex.what());
diff --git a/extensions/sftp/processors/PutSFTP.cpp 
b/extensions/sftp/processors/PutSFTP.cpp
index b04feb457..e41285141 100644
--- a/extensions/sftp/processors/PutSFTP.cpp
+++ b/extensions/sftp/processors/PutSFTP.cpp
@@ -248,7 +248,7 @@ bool PutSFTP::processOne(core::ProcessContext& context, 
core::ProcessSession& se
           gsl::narrow<int64_t>(stream->size()) /*expected_size*/)) {
         throw utils::SFTPException{client->getLastError()};
       }
-      return gsl::narrow<int64_t>(stream->size());
+      return io::IoResult::from(stream->size());
     });
   } catch (const utils::SFTPException& ex) {
     logger_->log_debug("{}", ex.what());
diff --git a/extensions/standard-processors/controllers/JsonRecordSetWriter.cpp 
b/extensions/standard-processors/controllers/JsonRecordSetWriter.cpp
index efc3f2326..c018da50a 100644
--- a/extensions/standard-processors/controllers/JsonRecordSetWriter.cpp
+++ b/extensions/standard-processors/controllers/JsonRecordSetWriter.cpp
@@ -75,7 +75,7 @@ void JsonRecordSetWriter::onEnable() {
 }
 
 void JsonRecordSetWriter::writePerLine(const core::RecordSet& record_set, 
const std::shared_ptr<core::FlowFile>& flow_file, core::ProcessSession& 
session) {
-  session.write(flow_file, [&record_set](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
+  session.write(flow_file, [&record_set](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
     int64_t write_result = 0;
     for (const auto& record : record_set) {
       auto doc = rapidjson::Document(rapidjson::kObjectType);
@@ -86,12 +86,12 @@ void JsonRecordSetWriter::writePerLine(const 
core::RecordSet& record_set, const
       doc.Accept(writer);
       write_result += 
gsl::narrow<int64_t>(stream->write(gsl::make_span(fmt::format("{}\n", 
buffer.GetString())).as_span<const std::byte>()));
     }
-    return write_result;
+    return io::IoResult::from(write_result);
   });
 }
 
 void JsonRecordSetWriter::writeAsArray(const core::RecordSet& record_set, 
const std::shared_ptr<core::FlowFile>& flow_file, core::ProcessSession& 
session) const {
-  session.write(flow_file, [this, &record_set](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
+  session.write(flow_file, [this, &record_set](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
     auto doc = rapidjson::Document(rapidjson::kArrayType);
     for (const auto& record : record_set) {
       auto& allocator = doc.GetAllocator();
@@ -107,7 +107,7 @@ void JsonRecordSetWriter::writeAsArray(const 
core::RecordSet& record_set, const
       rapidjson::Writer writer(buffer);
       doc.Accept(writer);
     }
-    return gsl::narrow<int64_t>(stream->write(gsl::make_span(fmt::format("{}", 
buffer.GetString())).as_span<const std::byte>()));
+    return io::IoResult::from(stream->write(gsl::make_span(fmt::format("{}", 
buffer.GetString())).as_span<const std::byte>()));
   });
 }
 
diff --git a/extensions/standard-processors/controllers/XMLRecordSetWriter.cpp 
b/extensions/standard-processors/controllers/XMLRecordSetWriter.cpp
index b98cbb67b..4e283e553 100644
--- a/extensions/standard-processors/controllers/XMLRecordSetWriter.cpp
+++ b/extensions/standard-processors/controllers/XMLRecordSetWriter.cpp
@@ -154,9 +154,9 @@ void XMLRecordSetWriter::write(const core::RecordSet& 
record_set, const std::sha
   }
 
   auto xml_content = convertRecordSetToXml(record_set);
-  session.write(flow_file, [&xml_content](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
+  session.write(flow_file, [&xml_content](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
     stream->write(reinterpret_cast<const uint8_t*>(xml_content.data()), 
xml_content.size());
-    return gsl::narrow<int64_t>(xml_content.size());
+    return io::IoResult::from(xml_content.size());
   });
 }
 
diff --git a/extensions/standard-processors/processors/ConvertRecord.cpp 
b/extensions/standard-processors/processors/ConvertRecord.cpp
index cf37cc191..00d864b8c 100644
--- a/extensions/standard-processors/processors/ConvertRecord.cpp
+++ b/extensions/standard-processors/processors/ConvertRecord.cpp
@@ -41,9 +41,9 @@ void ConvertRecord::onTrigger(core::ProcessContext& context, 
core::ProcessSessio
   }
 
   nonstd::expected<core::RecordSet, std::error_code> record_set;
-  session.read(flow_file, [this, &record_set](const 
std::shared_ptr<io::InputStream>& input_stream) {
+  session.read(flow_file, [this, &record_set](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
     record_set = record_converter_->record_set_reader->read(*input_stream);
-    return gsl::narrow<int64_t>(input_stream->size());
+    return io::IoResult::from(input_stream->size());
   });
   if (!record_set) {
     logger_->log_error("Failed to read record set from flow file: {}", 
record_set.error().message());
diff --git a/extensions/standard-processors/processors/DefragmentText.cpp 
b/extensions/standard-processors/processors/DefragmentText.cpp
index 6b051682b..dae7511de 100644
--- a/extensions/standard-processors/processors/DefragmentText.cpp
+++ b/extensions/standard-processors/processors/DefragmentText.cpp
@@ -193,12 +193,12 @@ void DefragmentText::Buffer::append(core::ProcessSession& 
session, const gsl::no
     store(session, flow_file_to_append);
     return;
   }
-  auto flowFileReader = [&] (const std::shared_ptr<core::FlowFile>& ff, const 
io::InputStreamCallback& cb) {
-    return session.read(ff, cb);
+  auto flowFileReader = [&] (const std::shared_ptr<core::FlowFile>& ff, const 
io::InputStreamCallback& cb) -> io::IoResult {
+    return io::IoResult::from(session.read(ff, cb));
   };
   PayloadSerializer serializer(flowFileReader);
   session.add(buffered_flow_file_);
-  session.append(buffered_flow_file_, [&serializer, 
&flow_file_to_append](const auto& output_stream) -> int64_t {
+  session.append(buffered_flow_file_, [&serializer, 
&flow_file_to_append](const auto& output_stream) -> io::IoResult {
     return serializer.serialize(flow_file_to_append, output_stream);
   });
   updateAppendedAttributes(*buffered_flow_file_);
diff --git a/extensions/standard-processors/processors/EvaluateJsonPath.cpp 
b/extensions/standard-processors/processors/EvaluateJsonPath.cpp
index 9273cee6f..09803d229 100644
--- a/extensions/standard-processors/processors/EvaluateJsonPath.cpp
+++ b/extensions/standard-processors/processors/EvaluateJsonPath.cpp
@@ -84,9 +84,9 @@ std::string EvaluateJsonPath::extractQueryResult(const 
jsoncons::json& query_res
 void EvaluateJsonPath::writeQueryResult(core::ProcessSession& session, 
core::FlowFile& flow_file, const jsoncons::json& query_result, const 
std::string& property_name,
     std::unordered_map<std::string, std::string>& attributes_to_set) const {
   if (destination_ == evaluate_json_path::DestinationType::FlowFileContent) {
-    session.write(flow_file, [&query_result, this](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
+    session.write(flow_file, [&query_result, this](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
       auto result_string = extractQueryResult(query_result);
-      return gsl::narrow<int64_t>(output_stream->write(reinterpret_cast<const 
uint8_t*>(result_string.data()), result_string.size()));
+      return io::IoResult::from(output_stream->write(reinterpret_cast<const 
uint8_t*>(result_string.data()), result_string.size()));
     });
   } else {
     attributes_to_set.emplace(property_name, extractQueryResult(query_result));
diff --git a/extensions/standard-processors/processors/ExtractText.cpp 
b/extensions/standard-processors/processors/ExtractText.cpp
index 0cba55b81..a2e612b6c 100644
--- a/extensions/standard-processors/processors/ExtractText.cpp
+++ b/extensions/standard-processors/processors/ExtractText.cpp
@@ -53,7 +53,7 @@ void ExtractText::onTrigger(core::ProcessContext& context, 
core::ProcessSession&
   session.transfer(flowFile, Success);
 }
 
-int64_t ExtractText::ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) const {
+io::IoResult ExtractText::ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) const {
   const auto flow_file_size = gsl::narrow<size_t>(flowFile_->getSize());
   const auto default_buffer_size = 
utils::configuration::getBufferSize(*ctx_->getConfiguration());
   std::vector<std::byte> buffer((std::min)(flow_file_size, 
default_buffer_size));
@@ -74,7 +74,7 @@ int64_t ExtractText::ReadCallback::operator()(const 
std::shared_ptr<io::InputStr
     const auto ret = stream->read(std::span(buffer).subspan(0, length));
 
     if (io::isError(ret)) {
-      return -1;  // Stream error
+      return io::IoResult::error();  // Stream error
     } else if (ret == 0) {
       break;  // End of stream, no more data
     }
@@ -82,7 +82,7 @@ int64_t ExtractText::ReadCallback::operator()(const 
std::shared_ptr<io::InputStr
     contentStream.write(reinterpret_cast<const char*>(buffer.data()), 
gsl::narrow<std::streamsize>(ret));
     read_size += ret;
     if (contentStream.fail()) {
-      return -1;
+      return io::IoResult::error();
     }
   }
 
@@ -137,7 +137,7 @@ int64_t ExtractText::ReadCallback::operator()(const 
std::shared_ptr<io::InputStr
   } else {
     flowFile_->setAttribute(attrKey, contentStream.str());
   }
-  return gsl::narrow<int64_t>(read_size);
+  return io::IoResult::from(read_size);
 }
 
 ExtractText::ReadCallback::ReadCallback(std::shared_ptr<core::FlowFile> 
flowFile, core::ProcessContext& ctx,  std::shared_ptr<core::logging::Logger> 
lgr)
diff --git a/extensions/standard-processors/processors/ExtractText.h 
b/extensions/standard-processors/processors/ExtractText.h
index fdf17e5e1..9041da167 100644
--- a/extensions/standard-processors/processors/ExtractText.h
+++ b/extensions/standard-processors/processors/ExtractText.h
@@ -109,7 +109,7 @@ class ExtractText : public core::ProcessorImpl {
   class ReadCallback {
    public:
     ReadCallback(std::shared_ptr<core::FlowFile> flowFile, 
core::ProcessContext& ctx, std::shared_ptr<core::logging::Logger> lgr);
-    int64_t operator()(const std::shared_ptr<io::InputStream>& stream) const;
+    io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) 
const;
 
    private:
     std::shared_ptr<core::FlowFile> flowFile_;
diff --git a/extensions/standard-processors/processors/HashContent.cpp 
b/extensions/standard-processors/processors/HashContent.cpp
index f48afc0b3..d81adc978 100644
--- a/extensions/standard-processors/processors/HashContent.cpp
+++ b/extensions/standard-processors/processors/HashContent.cpp
@@ -69,12 +69,11 @@ void HashContent::onTrigger(core::ProcessContext&, 
core::ProcessSession& session
   }
 
   logger_->log_trace("attempting read");
-  session.read(flowFile, [&flowFile, this](const 
std::shared_ptr<io::InputStream>& stream) {
+  session.read(flowFile, [&flowFile, this](const 
std::shared_ptr<io::InputStream>& stream) -> io::IoResult {
     const auto& ret_val = algorithm_(stream);
-
     flowFile->setAttribute(attrKey_, ret_val.first);
 
-    return ret_val.second;
+    return io::IoResult::from(ret_val.second);
   });
   session.transfer(flowFile, Success);
 }
diff --git a/extensions/standard-processors/processors/RouteText.cpp 
b/extensions/standard-processors/processors/RouteText.cpp
index fee76d588..4604ab41e 100644
--- a/extensions/standard-processors/processors/RouteText.cpp
+++ b/extensions/standard-processors/processors/RouteText.cpp
@@ -68,12 +68,12 @@ class RouteText::ReadCallback {
   ReadCallback(route_text::Segmentation segmentation, size_t file_size, Fn&& 
fn)
     : segmentation_(segmentation), file_size_(file_size), fn_(std::move(fn)) {}
 
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream) const {
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream) 
const {
     std::vector<std::byte> buffer;
     buffer.resize(file_size_);
     size_t ret = stream->read(buffer);
     if (io::isError(ret)) {
-      return -1;
+      return io::IoResult::error();
     }
     if (ret != file_size_) {
       throw Exception(PROCESS_SESSION_EXCEPTION, "Couldn't read whole flowfile 
content");
@@ -82,7 +82,7 @@ class RouteText::ReadCallback {
     switch (segmentation_) {
       case route_text::Segmentation::FULL_TEXT: {
         fn_({content, 0});
-        return gsl::narrow<int64_t>(content.length());
+        return io::IoResult::from(content.length());
       }
       case route_text::Segmentation::PER_LINE: {
         // 1-based index as in nifi
@@ -102,7 +102,7 @@ class RouteText::ReadCallback {
           curr = next_line;
           ++segment_idx;
         }
-        return gsl::narrow<int64_t>(content.length());
+        return io::IoResult::from(content.length());
       }
     }
     throw Exception(PROCESSOR_EXCEPTION, "Unknown segmentation strategy");
diff --git a/extensions/standard-processors/processors/SplitJson.cpp 
b/extensions/standard-processors/processors/SplitJson.cpp
index 0e3d90cd4..9498de3f8 100644
--- a/extensions/standard-processors/processors/SplitJson.cpp
+++ b/extensions/standard-processors/processors/SplitJson.cpp
@@ -115,9 +115,9 @@ void SplitJson::onTrigger(core::ProcessContext& context, 
core::ProcessSession& s
     child_flow_file->setAttribute(SplitJson::SegmentOriginalFilename.name,  
original_filename ? original_filename.value() : "");
 
     auto& json_value_to_write = result_array[i];
-    session.write(child_flow_file, [this, &json_value_to_write](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
+    session.write(child_flow_file, [this, &json_value_to_write](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
       auto result_string = jsonValueToString(json_value_to_write);
-      return gsl::narrow<int64_t>(output_stream->write(reinterpret_cast<const 
uint8_t*>(result_string.data()), result_string.size()));
+      return io::IoResult::from(output_stream->write(reinterpret_cast<const 
uint8_t*>(result_string.data()), result_string.size()));
     });
 
     session.transfer(child_flow_file, Split);
diff --git a/extensions/standard-processors/processors/SplitRecord.cpp 
b/extensions/standard-processors/processors/SplitRecord.cpp
index fb32a5baa..ba501473a 100644
--- a/extensions/standard-processors/processors/SplitRecord.cpp
+++ b/extensions/standard-processors/processors/SplitRecord.cpp
@@ -55,9 +55,9 @@ void SplitRecord::onTrigger(core::ProcessContext& context, 
core::ProcessSession&
   }
 
   nonstd::expected<core::RecordSet, std::error_code> record_set;
-  session.read(original_flow_file, [this, &record_set](const 
std::shared_ptr<io::InputStream>& input_stream) {
+  session.read(original_flow_file, [this, &record_set](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
     record_set = record_converter_->record_set_reader->read(*input_stream);
-    return gsl::narrow<int64_t>(input_stream->size());
+    return io::IoResult::from(input_stream->size());
   });
   if (!record_set) {
     logger_->log_error("Failed to read record set from flow file: {}", 
record_set.error().message());
diff --git a/extensions/standard-processors/processors/SplitText.cpp 
b/extensions/standard-processors/processors/SplitText.cpp
index 60dc02442..3573169a8 100644
--- a/extensions/standard-processors/processors/SplitText.cpp
+++ b/extensions/standard-processors/processors/SplitText.cpp
@@ -157,7 +157,7 @@ class ReadCallback {
  public:
   ReadCallback(std::shared_ptr<core::FlowFile> flow_file, const 
SplitTextConfiguration& split_text_config,
     core::ProcessSession& session, size_t buffer_size, 
std::shared_ptr<core::logging::Logger> logger);
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream);
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream);
   std::optional<const char*> error;
   std::vector<std::shared_ptr<org::apache::nifi::minifi::core::FlowFile>> 
results;
 
@@ -325,17 +325,17 @@ void ReadCallback::mergeHeaderAndFragmentFlows(const 
std::shared_ptr<core::FlowF
     logger_->log_error("Failed to clone merged fragment flow!");
     return;
   }
-  session_.write(merged_flow, [this, &fragment_flow, &header_flow](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
-    auto header_write_result = session_.read(header_flow, 
[&output_stream](const std::shared_ptr<io::InputStream>& header_input_stream) 
-> int64_t {
+  session_.write(merged_flow, [this, &fragment_flow, &header_flow](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
+    const int64_t header_write_result = session_.read(header_flow, 
[&output_stream](const std::shared_ptr<io::InputStream>& header_input_stream) 
-> io::IoResult {
       return internal::pipe(*header_input_stream, *output_stream);
     });
     if (header_write_result < 0) {
       logger_->log_error("Failed to write header to fragment!");
-      return header_write_result;
+      return io::IoResult::from(header_write_result);
     }
-    return session_.read(fragment_flow, [&output_stream](const 
std::shared_ptr<io::InputStream>& fragment_input_stream) -> int64_t {
+    return io::IoResult::from(session_.read(fragment_flow, 
[&output_stream](const std::shared_ptr<io::InputStream>& fragment_input_stream) 
-> io::IoResult {
       return internal::pipe(*fragment_input_stream, *output_stream);
-    });
+    }));
   });
   logger_->log_debug("Creating fragment with header with fragment index: {} 
fragment size: {}", emitted_fragment_index_, merged_flow->getSize());
   setAttributesOfDoneSegment(*merged_flow, fragment.text_line_count);
@@ -354,7 +354,7 @@ void ReadCallback::createFragmentFlowWithoutHeader(const 
SplitTextFragmentGenera
   results.push_back(fragment_flow);
 }
 
-int64_t ReadCallback::operator()(const std::shared_ptr<io::InputStream>& 
stream) {
+io::IoResult ReadCallback::operator()(const std::shared_ptr<io::InputStream>& 
stream) {
   SplitTextFragmentGenerator fragment_generator(stream, split_text_config_, 
buffer_size_);
   nonstd::expected<SplitTextFragmentGenerator::Fragment, const char*> 
header_fragment;
   std::shared_ptr<core::FlowFile> header_flow;  // cache header flow file to 
avoid cloning it for each fragment
@@ -362,12 +362,12 @@ int64_t ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream)
     header_fragment = fragment_generator.readHeaderFragment();
     if (!header_fragment) {
       error = header_fragment.error();
-      return gsl::narrow<int64_t>(flow_file_->getSize());
+      return io::IoResult::from(flow_file_->getSize());
     }
     header_flow = session_.clone(*flow_file_, 
gsl::narrow<int64_t>(header_fragment->fragment_offset), 
gsl::narrow<int64_t>(header_fragment->fragment_size));
     if (!header_flow) {
       logger_->log_error("Failed to clone header flow!");
-      return -1;
+      return io::IoResult::error();
     }
   }
 
@@ -386,7 +386,10 @@ int64_t ReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream)
   if (header_flow) {
     session_.remove(header_flow);
   }
-  return fragment_generator.getState() == detail::StreamReadState::EndOfStream 
? gsl::narrow<int64_t>(flow_file_->getSize()) : -1;
+  if (fragment_generator.getState() == detail::StreamReadState::EndOfStream) {
+    return io::IoResult::from(flow_file_->getSize());
+  }
+  return io::IoResult::error();
 }
 
 }  // namespace
diff --git a/extensions/standard-processors/processors/TailFile.cpp 
b/extensions/standard-processors/processors/TailFile.cpp
index c9bdd74a7..334c2ec22 100644
--- a/extensions/standard-processors/processors/TailFile.cpp
+++ b/extensions/standard-processors/processors/TailFile.cpp
@@ -111,7 +111,7 @@ class FileReaderCallback {
     openFile(file_path, offset, input_stream_, logger_);
   }
 
-  int64_t operator()(const std::shared_ptr<io::OutputStream>& output_stream) {
+  io::IoResult operator()(const std::shared_ptr<io::OutputStream>& 
output_stream) {
     io::CRCStream<io::OutputStream> 
crc_stream{gsl::make_not_null(output_stream.get()), checksum_};
 
     uint64_t num_bytes_written = 0;
@@ -143,7 +143,7 @@ class FileReaderCallback {
       latest_flow_file_ends_with_delimiter_ = false;
     }
 
-    return gsl::narrow<int64_t>(num_bytes_written);
+    return io::IoResult::from(num_bytes_written);
   }
 
   uint64_t checksum() const {
@@ -187,7 +187,7 @@ class WholeFileReaderCallback {
     return checksum_;
   }
 
-  int64_t operator()(const std::shared_ptr<io::OutputStream>& output_stream) {
+  io::IoResult operator()(const std::shared_ptr<io::OutputStream>& 
output_stream) {
     std::vector<char> buffer(buffer_size_);
 
     io::CRCStream<io::OutputStream> 
crc_stream{gsl::make_not_null(output_stream.get()), checksum_};
@@ -208,7 +208,7 @@ class WholeFileReaderCallback {
 
     checksum_ = crc_stream.getCRC();
 
-    return gsl::narrow<int64_t>(num_bytes_written);
+    return io::IoResult::from(num_bytes_written);
   }
 
  private:
diff --git 
a/extensions/standard-processors/tests/unit/XMLRecordSetWriterTests.cpp 
b/extensions/standard-processors/tests/unit/XMLRecordSetWriterTests.cpp
index ab887af67..2f7900e80 100644
--- a/extensions/standard-processors/tests/unit/XMLRecordSetWriterTests.cpp
+++ b/extensions/standard-processors/tests/unit/XMLRecordSetWriterTests.cpp
@@ -70,7 +70,7 @@ class XMLRecordSetWriterTestFixture {
       std::vector<std::byte> buffer(input_stream->size());
       input_stream->read(buffer);
       xml_content = std::string(reinterpret_cast<const char*>(buffer.data()), 
buffer.size());
-      return gsl::narrow<int64_t>(input_stream->size());
+      return io::IoResult::from(input_stream->size());
     });
     return xml_content;
   }
diff --git a/libminifi/include/c2/PayloadSerializer.h 
b/libminifi/include/c2/PayloadSerializer.h
index 35590c072..1d7125176 100644
--- a/libminifi/include/c2/PayloadSerializer.h
+++ b/libminifi/include/c2/PayloadSerializer.h
@@ -72,7 +72,7 @@ class PayloadSerializer {
     }
   }
   static void serialize(uint16_t op, const C2Payload &payload, 
std::shared_ptr<io::OutputStream> stream) {
-    uint8_t st;
+    uint8_t st = 0;
     uint32_t size = gsl::narrow<uint32_t>(payload.getNestedPayloads().size());
     stream->write(size);
     for (const auto &nested_payload : payload.getNestedPayloads()) {
diff --git a/libminifi/include/core/ProcessSessionReadCallback.h 
b/libminifi/include/core/ProcessSessionReadCallback.h
index 715f4155e..3693b6392 100644
--- a/libminifi/include/core/ProcessSessionReadCallback.h
+++ b/libminifi/include/core/ProcessSessionReadCallback.h
@@ -31,7 +31,7 @@ class ProcessSessionReadCallback {
  public:
   ProcessSessionReadCallback(std::filesystem::path temp_file, 
std::filesystem::path dest_file, std::shared_ptr<logging::Logger> logger);
   ~ProcessSessionReadCallback();
-  int64_t operator()(const std::shared_ptr<io::InputStream>& stream);
+  io::IoResult operator()(const std::shared_ptr<io::InputStream>& stream);
   bool commit();
 
  private:
diff --git a/libminifi/src/c2/C2Utils.cpp b/libminifi/src/c2/C2Utils.cpp
index 18141b861..00a8f6d9e 100644
--- a/libminifi/src/c2/C2Utils.cpp
+++ b/libminifi/src/c2/C2Utils.cpp
@@ -50,7 +50,7 @@ nonstd::expected<std::shared_ptr<io::BufferStream>, 
std::string> createDebugBund
     if (!archiver->newEntry({filename, file_size})) {
       return nonstd::make_unexpected("Couldn't initialize archive entry for '" 
+ filename + "'");
     }
-    if (gsl::narrow<int64_t>(file_size) != minifi::internal::pipe(*stream, 
*archiver)) {
+    if (gsl::narrow<int64_t>(file_size) != minifi::internal::pipe(*stream, 
*archiver).toI64()) {
       // we have touched the input streams, they cannot be reused
       return nonstd::make_unexpected("Error while writing file '" + filename + 
"' into the debug bundle");
     }
diff --git a/libminifi/src/core/ProcessSession.cpp 
b/libminifi/src/core/ProcessSession.cpp
index 23708c6a3..662f2b00e 100644
--- a/libminifi/src/core/ProcessSession.cpp
+++ b/libminifi/src/core/ProcessSession.cpp
@@ -189,11 +189,11 @@ std::shared_ptr<core::FlowFile> 
ProcessSessionImpl::clone(const FlowFile& parent
   if (record) {
     logger_->log_debug("Cloned parent flow files {} to {}, with {}:{}", 
parent.getUUIDStr(), record->getUUIDStr(), offset, size);
     if (parent.getResourceClaim()) {
-      write(record, [&] (const std::shared_ptr<io::OutputStream>& output) -> 
int64_t {
-        return read(parent, [&] (const std::shared_ptr<io::InputStream>& 
input) -> int64_t {
+      write(record, [&] (const std::shared_ptr<io::OutputStream>& output) -> 
io::IoResult {
+        return io::IoResult::from(read(parent, [&] (const 
std::shared_ptr<io::InputStream>& input) -> io::IoResult {
           io::StreamSlice slice(input, offset, size);
-          return minifi::internal::pipe(slice, *output);
-        });
+          return internal::pipe(slice, *output);
+        }));
       });
     }
     provenance_report_->clone(parent, *record);
@@ -262,14 +262,15 @@ void ProcessSessionImpl::write(core::FlowFile &flow, 
const io::OutputStreamCallb
       throw Exception(FILE_OPERATION_EXCEPTION, "Failed to open flowfile 
content for write");
     }
     const auto callback_result = callback(stream);
-    if (callback_result == MinifiIoStatus::MINIFI_IO_CANCEL) {
+
+    if (callback_result.is_cancelled()) {
       stream->close();
       content_session_->remove(claim);
       claim.reset();
       return;
     }
 
-    if (callback_result < 0) {
+    if (!callback_result) {
       throw Exception(FILE_OPERATION_EXCEPTION, "Failed to process flowfile 
content");
     }
 
@@ -300,7 +301,7 @@ void ProcessSessionImpl::writeBuffer(const 
std::shared_ptr<core::FlowFile>& flow
 void ProcessSessionImpl::writeBuffer(const std::shared_ptr<core::FlowFile>& 
flow_file, std::span<const std::byte> buffer) {
   write(flow_file, [buffer](const std::shared_ptr<io::OutputStream>& 
output_stream) {
     const auto write_status = output_stream->write(buffer);
-    return io::isError(write_status) ? -1 : gsl::narrow<int64_t>(write_status);
+    return io::IoResult::from(write_status);
   });
 }
 
@@ -330,7 +331,7 @@ void ProcessSessionImpl::append(const 
std::shared_ptr<core::FlowFile> &flow, con
     // this prevents an issue if we write, above, with zero length.
     if (stream_size_before_callback > 0)
       stream->seek(stream_size_before_callback);
-    if (callback(stream) < 0) {
+    if (!callback(stream)) {
       throw Exception(FILE_OPERATION_EXCEPTION, "Failed to process flowfile 
content");
     }
     flow->setSize(flow_file_size + (stream->size() - 
stream_size_before_callback));
@@ -357,7 +358,7 @@ void ProcessSessionImpl::appendBuffer(const 
std::shared_ptr<core::FlowFile>& flo
   if (buffer.empty()) { return; }
   append(flow_file, [buffer](const std::shared_ptr<io::OutputStream>& 
output_stream) {
     const auto write_status = output_stream->write(buffer);
-    return io::isError(write_status) ? -1 : gsl::narrow<int64_t>(write_status);
+    return io::IoResult::from(write_status);
   });
 }
 
@@ -390,14 +391,14 @@ int64_t ProcessSessionImpl::read(const core::FlowFile& 
flow_file, const io::Inpu
       return 0;
     }
 
-    auto ret = callback(flow_file_stream);
-    if (ret < 0) {
+    const auto callback_result = callback(flow_file_stream);
+    if (!callback_result) {
       throw Exception(FILE_OPERATION_EXCEPTION, "Failed to process flowfile 
content");
     }
     if (metrics_) {
-      metrics_->bytesRead() += ret;
+      metrics_->bytesRead() += *callback_result;
     }
-    return ret;
+    return callback_result.toI64();
   } catch (const std::exception& exception) {
     logger_->log_debug("Caught Exception {}", exception.what());
     throw;
@@ -432,7 +433,7 @@ int64_t ProcessSessionImpl::readWrite(const 
std::shared_ptr<core::FlowFile> &flo
       throw Exception(FILE_OPERATION_EXCEPTION, "Failed to open flowfile 
content for write");
     }
 
-    auto read_write_result = callback(input_stream, output_stream);
+    const auto read_write_result = callback(input_stream, output_stream);
     if (!read_write_result) {
       throw Exception(FILE_OPERATION_EXCEPTION, "Failed to process flowfile 
content");
     }
@@ -440,15 +441,15 @@ int64_t ProcessSessionImpl::readWrite(const 
std::shared_ptr<core::FlowFile> &flo
     input_stream->close();
     output_stream->close();
 
-    flow->setSize(gsl::narrow<uint64_t>(read_write_result->bytes_written));
+    flow->setSize(gsl::narrow<uint64_t>(read_write_result.bytesWritten()));
     flow->setOffset(0);
     flow->setResourceClaim(output_claim);
     if (metrics_) {
-      metrics_->bytesWritten() += read_write_result->bytes_written;
-      metrics_->bytesRead() += read_write_result->bytes_read;
+      metrics_->bytesWritten() += read_write_result.bytesWritten();
+      metrics_->bytesRead() += read_write_result.bytesRead();
     }
 
-    return read_write_result->bytes_written;
+    return gsl::narrow<int64_t>(read_write_result.bytesWritten());
   } catch (const std::exception& exception) {
     logger_->log_debug("Caught exception during process session readWrite, 
type: {}, what: {}", typeid(exception).name(), exception.what());
     throw;
@@ -460,14 +461,14 @@ int64_t ProcessSessionImpl::readWrite(const 
std::shared_ptr<core::FlowFile> &flo
 
 detail::ReadBufferResult ProcessSessionImpl::readBuffer(const 
std::shared_ptr<core::FlowFile>& flow) {
   detail::ReadBufferResult result;
-  result.status = read(flow, [&result, this](const 
std::shared_ptr<io::InputStream>& input_stream) {
+  result.status = read(flow, [&result, this](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
     result.buffer.resize(input_stream->size());
     const auto read_status = input_stream->read(result.buffer);
     if (read_status != result.buffer.size()) {
       logger_->log_error("readBuffer: {} bytes were requested from the stream 
but {} bytes were read. Rolling back.", result.buffer.size(), read_status);
       throw Exception(PROCESSOR_EXCEPTION, "Failed to read the entire 
FlowFile.");
     }
-    return gsl::narrow<int64_t>(read_status);
+    return io::IoResult::from(read_status);
   });
   return result;
 }
diff --git a/libminifi/src/core/ProcessSessionReadCallback.cpp 
b/libminifi/src/core/ProcessSessionReadCallback.cpp
index 139041fc2..7f1ec0846 100644
--- a/libminifi/src/core/ProcessSessionReadCallback.cpp
+++ b/libminifi/src/core/ProcessSessionReadCallback.cpp
@@ -39,22 +39,22 @@ 
ProcessSessionReadCallback::ProcessSessionReadCallback(std::filesystem::path tem
 }
 
 // Copy the entire file contents to the temporary file
-int64_t ProcessSessionReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) {
+io::IoResult ProcessSessionReadCallback::operator()(const 
std::shared_ptr<io::InputStream>& stream) {
   // Copy file contents into tmp file
   write_succeeded_ = false;
   size_t size = 0;
   std::array<std::byte, 8192> buffer{};
   do {
     const auto read = stream->read(buffer);
-    if (io::isError(read)) return -1;
+    if (io::isError(read)) { return io::IoResult::error(); }
     if (read == 0) break;
     if (!tmp_file_os_.write(reinterpret_cast<char*>(buffer.data()), 
gsl::narrow<std::streamsize>(read))) {
-      return -1;
+      return io::IoResult::error();
     }
     size += read;
   } while (size < stream->size());
         write_succeeded_ = true;
-  return gsl::narrow<int64_t>(size);
+  return io::IoResult::from(size);
 }
 
 // Renames tmp file to final destination
diff --git a/libminifi/src/minifi-c.cpp b/libminifi/src/minifi-c.cpp
index 44152eb25..dac7bbe21 100644
--- a/libminifi/src/minifi-c.cpp
+++ b/libminifi/src/minifi-c.cpp
@@ -382,8 +382,9 @@ MinifiStatus MinifiProcessSessionRead(MinifiProcessSession* 
session, MinifiFlowF
   gsl_Assert(session != MINIFI_NULL);
   gsl_Assert(flowfile != MINIFI_NULL);
   try {
-    
reinterpret_cast<minifi::core::ProcessSession*>(session)->read(*reinterpret_cast<std::shared_ptr<minifi::core::FlowFile>*>(flowfile),
 [&] (auto& input_stream) -> int64_t {
-      return cb(user_ctx, 
reinterpret_cast<MinifiInputStream*>(input_stream.get()));
+    
reinterpret_cast<minifi::core::ProcessSession*>(session)->read(*reinterpret_cast<std::shared_ptr<minifi::core::FlowFile>*>(flowfile),
 [&] (auto& input_stream) -> minifi::io::IoResult {
+      const int64_t cb_result = cb(user_ctx, 
reinterpret_cast<MinifiInputStream*>(input_stream.get()));
+      return minifi::io::IoResult::from(cb_result);
     });
     return MINIFI_STATUS_SUCCESS;
   } catch (...) {
@@ -395,8 +396,9 @@ MinifiStatus 
MinifiProcessSessionWrite(MinifiProcessSession* session, MinifiFlow
   gsl_Assert(session != MINIFI_NULL);
   gsl_Assert(ff != MINIFI_NULL);
   try {
-    
reinterpret_cast<minifi::core::ProcessSession*>(session)->write(*reinterpret_cast<std::shared_ptr<minifi::core::FlowFile>*>(ff),
 [&] (auto& output_stream) -> int64_t {
-      return cb(user_ctx, 
reinterpret_cast<MinifiOutputStream*>(output_stream.get()));
+    
reinterpret_cast<minifi::core::ProcessSession*>(session)->write(*reinterpret_cast<std::shared_ptr<minifi::core::FlowFile>*>(ff),
 [&] (auto& output_stream) -> minifi::io::IoResult {
+      const int64_t cb_result = cb(user_ctx, 
reinterpret_cast<MinifiOutputStream*>(output_stream.get()));
+      return minifi::io::IoResult::from(cb_result);
     });
     return MINIFI_STATUS_SUCCESS;
   } catch (...) {
diff --git a/libminifi/src/sitetosite/CompressionOutputStream.cpp 
b/libminifi/src/sitetosite/CompressionOutputStream.cpp
index 31e1e4388..d6631e55a 100644
--- a/libminifi/src/sitetosite/CompressionOutputStream.cpp
+++ b/libminifi/src/sitetosite/CompressionOutputStream.cpp
@@ -64,14 +64,14 @@ size_t CompressionOutputStream::write(const uint8_t *value, 
size_t len) {
 size_t CompressionOutputStream::compressAndWrite() {
   if (was_data_written_) {
     // Write a continue byte to indicate that there is more data to follow
-    auto ret = internal_stream_.write(static_cast<uint8_t>(1));
+    const auto ret = internal_stream_.write(static_cast<uint8_t>(1));
     if (io::isError(ret)) {
       logger_->log_error("Failed to write continue byte before compression: 
{}", ret);
       return ret;
     }
   }
-  auto ret = internal_stream_.write(reinterpret_cast<const uint8_t 
*>(SYNC_BYTES.data()), SYNC_BYTES.size());
-  if (io::isError(ret)) {
+
+  if (const auto ret = internal_stream_.write(reinterpret_cast<const uint8_t 
*>(SYNC_BYTES.data()), SYNC_BYTES.size()); io::isError(ret)) {
     logger_->log_error("Failed to write sync bytes before compression: {}", 
ret);
     return ret;
   }
@@ -79,7 +79,7 @@ size_t CompressionOutputStream::compressAndWrite() {
   io::BufferStream buffer_stream;
   {
     io::ZlibCompressStream zlib_stream{gsl::make_not_null(&buffer_stream), 
io::ZlibCompressionFormat::ZLIB, Z_BEST_SPEED};
-    ret = zlib_stream.write(gsl::make_span(buffer_).subspan(0, 
buffer_offset_));
+    const auto ret = zlib_stream.write(gsl::make_span(buffer_).subspan(0, 
buffer_offset_));
     if (io::isError(ret)) {
       logger_->log_error("Failed to write data to zlib stream: {}", ret);
       return ret;
@@ -90,27 +90,26 @@ size_t CompressionOutputStream::compressAndWrite() {
   }
 
   // Write the original size of the data before compression
-  ret = internal_stream_.write(gsl::narrow<uint32_t>(buffer_offset_));
-  if (io::isError(ret)) {
+
+  if (const auto ret = 
internal_stream_.write(gsl::narrow<uint32_t>(buffer_offset_)); 
io::isError(ret)) {
     logger_->log_error("Failed to write original size before compression: {}", 
ret);
     return ret;
   }
 
   // Write the compressed size of the data
-  ret = internal_stream_.write(gsl::narrow<uint32_t>(buffer_stream.size()));
-  if (io::isError(ret)) {
+  if (const auto ret = 
internal_stream_.write(gsl::narrow<uint32_t>(buffer_stream.size())); 
io::isError(ret)) {
     return ret;
   }
 
   // Write the compressed data
-  ret = internal::pipe(buffer_stream, internal_stream_);
-  if (io::isError(ret)) {
-    return ret;
+  const io::IoResult pipe_res = internal::pipe(buffer_stream, 
internal_stream_);
+  if (!pipe_res) {
+    return io::STREAM_ERROR;
   }
 
   buffer_offset_ = 0;
   was_data_written_ = true;
-  return ret;
+  return gsl::narrow<size_t>(pipe_res.toI64());
 }
 
 void CompressionOutputStream::flush() {
diff --git a/libminifi/src/sitetosite/SiteToSiteClient.cpp 
b/libminifi/src/sitetosite/SiteToSiteClient.cpp
index 290097b40..c5bd28850 100644
--- a/libminifi/src/sitetosite/SiteToSiteClient.cpp
+++ b/libminifi/src/sitetosite/SiteToSiteClient.cpp
@@ -490,7 +490,7 @@ bool SiteToSiteClient::sendFlowFile(const 
std::shared_ptr<Transaction>& transact
       return false;
     }
     if (flow_file.getSize() > 0) {
-      auto read_result = session.read(flow_file, [&stream](const 
std::shared_ptr<io::InputStream>& input_stream) -> int64_t {
+      auto read_result = session.read(flow_file, [&stream](const 
std::shared_ptr<io::InputStream>& input_stream) -> io::IoResult {
         return internal::pipe(*input_stream, stream);
       });
       if (flow_file.getSize() != gsl::narrow<uint64_t>(read_result)) {
@@ -700,19 +700,19 @@ std::pair<uint64_t, uint64_t> 
SiteToSiteClient::readFlowFiles(const std::shared_
     }
 
     if (receive_header_result->flow_file_data_size > 0) {
-      session.write(flow_file, [&receive_header_result, &stream](const 
std::shared_ptr<io::OutputStream>& output_stream) -> int64_t {
+      session.write(flow_file, [&receive_header_result, &stream](const 
std::shared_ptr<io::OutputStream>& output_stream) -> io::IoResult {
         uint64_t len = receive_header_result->flow_file_data_size;
         std::array<std::byte, utils::configuration::DEFAULT_BUFFER_SIZE> 
buffer{};
         while (len > 0) {
           const auto size = std::min(len, 
uint64_t{utils::configuration::DEFAULT_BUFFER_SIZE});
           const auto ret = 
stream.read(std::as_writable_bytes(std::span(buffer).subspan(0, size)));
           if (ret != size) {
-            return -1;
+            return io::IoResult::error();
           }
           output_stream->write(std::span(buffer).subspan(0, size));
           len -= size;
         }
-        return 
gsl::narrow<int64_t>(receive_header_result->flow_file_data_size);
+        return io::IoResult::from(receive_header_result->flow_file_data_size);
       });
       if (flow_file->getSize() != receive_header_result->flow_file_data_size) {
         std::stringstream message;
diff --git a/libminifi/test/libtest/unit/ContentRepositoryDependentTests.h 
b/libminifi/test/libtest/unit/ContentRepositoryDependentTests.h
index b15bc2f47..d59b6d14d 100644
--- a/libminifi/test/libtest/unit/ContentRepositoryDependentTests.h
+++ b/libminifi/test/libtest/unit/ContentRepositoryDependentTests.h
@@ -37,16 +37,16 @@ namespace ContentRepositoryDependentTests {
 struct ReadUntilItCan {
   std::string value_;
 
-  int64_t operator()(const std::shared_ptr<minifi::io::InputStream> &stream) {
+  minifi::io::IoResult operator()(const 
std::shared_ptr<minifi::io::InputStream> &stream) {
     value_.clear();
     std::array<std::byte, 1024> buffer{};
     size_t bytes_read = 0;
     while (true) {
       const size_t read_result = stream->read(buffer);
       if (minifi::io::isError(read_result))
-        return -1;
+        return minifi::io::IoResult::error();
       if (read_result == 0)
-        return gsl::narrow<int64_t>(bytes_read);
+        return minifi::io::IoResult::from(bytes_read);
       bytes_read += read_result;
       const auto char_view = gsl::make_span(buffer).subspan(0, 
read_result).as_span<const char>();
       value_.append(std::begin(char_view), std::end(char_view));
@@ -179,10 +179,10 @@ inline void 
testErrWrite(std::shared_ptr<core::ContentRepository> content_repo)
   fixture.writeToFlowFile(flow_file, "original_content");
 
   REQUIRE_THROWS(
-  process_session.write(flow_file, [](const 
std::shared_ptr<minifi::io::OutputStream>& output_stream) {
+  process_session.write(flow_file, [](const 
std::shared_ptr<minifi::io::OutputStream>& output_stream) -> 
minifi::io::IoResult {
     std::string str = "new_content";
     output_stream->write(as_bytes(std::span(str)));
-    return MinifiIoStatus::MINIFI_IO_ERROR;
+    return minifi::io::IoResult::error();
   }));
 
   fixture.transferAndCommit(flow_file);
@@ -202,9 +202,9 @@ inline void 
testOkWrite(std::shared_ptr<core::ContentRepository> content_repo) {
 
   CHECK(flow_file->getSize() == 16);
 
-  process_session.write(flow_file, [](const 
std::shared_ptr<minifi::io::OutputStream>& output_stream) {
+  process_session.write(flow_file, [](const 
std::shared_ptr<minifi::io::OutputStream>& output_stream) -> 
minifi::io::IoResult {
     std::string str = "new_content";
-    return output_stream->write(as_bytes(std::span(str)));
+    return 
minifi::io::IoResult::from(output_stream->write(as_bytes(std::span(str))));
   });
 
   fixture.transferAndCommit(flow_file);
@@ -222,10 +222,10 @@ inline void 
testCancelWrite(std::shared_ptr<core::ContentRepository> content_rep
   const auto flow_file = process_session.create();
   fixture.writeToFlowFile(flow_file, "original_content");
 
-  process_session.write(flow_file, [](const 
std::shared_ptr<minifi::io::OutputStream>& output_stream) {
+  process_session.write(flow_file, [](const 
std::shared_ptr<minifi::io::OutputStream>& output_stream) -> 
minifi::io::IoResult {
     std::string str = "new_content";
     output_stream->write(as_bytes(std::span(str)));
-    return MinifiIoStatus::MINIFI_IO_CANCEL;
+    return minifi::io::IoResult::cancelled();
   });
 
   fixture.transferAndCommit(flow_file);
diff --git a/libminifi/test/libtest/unit/TestBase.h 
b/libminifi/test/libtest/unit/TestBase.h
index e8f48832d..f3b85eb63 100644
--- a/libminifi/test/libtest/unit/TestBase.h
+++ b/libminifi/test/libtest/unit/TestBase.h
@@ -260,11 +260,11 @@ class TestPlan {
 
   minifi::core::Processor* 
addProcessor(std::unique_ptr<minifi::core::Processor> processor, const 
std::string &name,
       const minifi::core::Relationship& relationship = 
minifi::core::Relationship("success", "description"), bool linkToPrevious = 
false) {
-    return addProcessor(std::move(processor), name, { relationship }, 
linkToPrevious);
+    return addProcessor(std::move(processor), name, 
std::initializer_list<minifi::core::Relationship>{ relationship }, 
linkToPrevious);
   }
   minifi::core::Processor* addProcessor(const std::string &processor_name, 
const std::string &name,
       const minifi::core::Relationship& relationship = 
minifi::core::Relationship("success", "description"), bool linkToPrevious = 
false) {
-    return addProcessor(processor_name, name, { relationship }, 
linkToPrevious);
+    return addProcessor(processor_name, name, 
std::initializer_list<minifi::core::Relationship>{ relationship }, 
linkToPrevious);
   }
   minifi::core::Processor* 
addProcessor(std::unique_ptr<minifi::core::Processor> processor, const 
std::string &name, const std::initializer_list<minifi::core::Relationship>& 
relationships,
     bool linkToPrevious = false);
diff --git a/libminifi/test/unit/FlowFileSerializationTests.cpp 
b/libminifi/test/unit/FlowFileSerializationTests.cpp
index df35a0617..1f2bcf73d 100644
--- a/libminifi/test/unit/FlowFileSerializationTests.cpp
+++ b/libminifi/test/unit/FlowFileSerializationTests.cpp
@@ -47,7 +47,7 @@ TEST_CASE("Payload Serializer", "[testPayload]") {
   flowFile->addAttribute("first", "one");
   flowFile->addAttribute("second", "two");
 
-  minifi::PayloadSerializer serializer([&] (const 
std::shared_ptr<core::FlowFile>&, const minifi::io::InputStreamCallback& cb) {
+  minifi::PayloadSerializer serializer([&] (const 
std::shared_ptr<core::FlowFile>&, const minifi::io::InputStreamCallback& cb) -> 
minifi::io::IoResult {
     return cb(contentStream);
   });
   serializer.serialize(flowFile, result);
@@ -67,7 +67,7 @@ TEST_CASE("FFv3 Serializer", "[testFFv3]") {
   flowFile->addAttribute("first", "one");
   flowFile->addAttribute("second", "two");
 
-  minifi::FlowFileV3Serializer serializer([&] (const 
std::shared_ptr<core::FlowFile>&, const minifi::io::InputStreamCallback& cb) {
+  minifi::FlowFileV3Serializer serializer([&] (const 
std::shared_ptr<core::FlowFile>&, const minifi::io::InputStreamCallback& cb) -> 
minifi::io::IoResult {
     return cb(contentStream);
   });
   serializer.serialize(flowFile, result);
diff --git a/libminifi/test/unit/MetricsTests.cpp 
b/libminifi/test/unit/MetricsTests.cpp
index 4fe02defd..906d50321 100644
--- a/libminifi/test/unit/MetricsTests.cpp
+++ b/libminifi/test/unit/MetricsTests.cpp
@@ -273,15 +273,15 @@ class DuplicateContentProcessor : public 
minifi::core::ProcessorImpl {
 
     auto flow_file_copy = session.create();
     std::vector<std::byte> buffer;
-    session.read(flow_file, [&](const std::shared_ptr<io::InputStream>& 
stream) -> int64_t {
+    session.read(flow_file, [&](const std::shared_ptr<io::InputStream>& 
stream) -> io::IoResult {
       buffer.resize(stream->size());
-      return gsl::narrow<int64_t>(stream->read(buffer));
+      return io::IoResult::from(stream->read(buffer));
     });
-    session.write(flow_file_copy, [&](const std::shared_ptr<io::OutputStream>& 
stream) -> int64_t {
-      return gsl::narrow<int64_t>(stream->write(buffer));
+    session.write(flow_file_copy, [&](const std::shared_ptr<io::OutputStream>& 
stream) -> io::IoResult {
+      return io::IoResult::from(stream->write(buffer));
     });
-    session.append(flow_file_copy, [&](const 
std::shared_ptr<io::OutputStream>& stream) -> int64_t {
-      return gsl::narrow<int64_t>(stream->write(buffer));
+    session.append(flow_file_copy, [&](const 
std::shared_ptr<io::OutputStream>& stream) -> io::IoResult {
+      return io::IoResult::from(stream->write(buffer));
     });
     session.transfer(flow_file_copy, Success);
     session.transfer(flow_file, Original);
diff --git a/libminifi/test/unit/SiteToSiteTests.cpp 
b/libminifi/test/unit/SiteToSiteTests.cpp
index 90c6b740d..5b4d45c0b 100644
--- a/libminifi/test/unit/SiteToSiteTests.cpp
+++ b/libminifi/test/unit/SiteToSiteTests.cpp
@@ -250,7 +250,7 @@ TEST_CASE("TestSiteToSiteVerifySend using flowfile data", 
"[S2S]") {
   session->write(flow_file, [&](const std::shared_ptr<io::OutputStream>& 
output_stream) {
     std::span<const std::byte> span{reinterpret_cast<const 
std::byte*>(payload.data()), payload.size()};
     output_stream->write(span);
-    return payload.size();
+    return io::IoResult::from(span.size());
   });
   flow_file->updateAttribute("filename", "myfile");
   flow_file->updateAttribute("flow.id", "test");
diff --git a/minifi-api/common/include/minifi-cpp/io/StreamCallback.h 
b/minifi-api/common/include/minifi-cpp/io/StreamCallback.h
index 8467c0f96..7063c7146 100644
--- a/minifi-api/common/include/minifi-cpp/io/StreamCallback.h
+++ b/minifi-api/common/include/minifi-cpp/io/StreamCallback.h
@@ -16,24 +16,109 @@
  */
 #pragma once
 
+#include <cinttypes>
 #include <functional>
 #include <memory>
-#include <optional>
+
+#include "../../minifi-api/include/minifi-c/minifi-c.h"
+#include "Stream.h"
+#include "minifi-cpp/utils/gsl.h"
+#include "utils/expected.h"
 
 namespace org::apache::nifi::minifi::io {
 
 class InputStream;
 class OutputStream;
 
-struct ReadWriteResult {
-  int64_t bytes_written = 0;
-  int64_t bytes_read = 0;
+class IoResult {
+ public:
+  IoResult() = delete;
+  IoResult(const IoResult&) = default;
+  IoResult(IoResult&&) noexcept = default;
+  IoResult& operator=(IoResult&&) noexcept = default;
+  IoResult& operator=(const IoResult&) = default;
+
+  virtual ~IoResult() = default;
+
+  static IoResult error() { return 
IoResult(nonstd::make_unexpected(MINIFI_IO_ERROR)); }
+  static IoResult cancelled() { return 
IoResult(nonstd::make_unexpected(MINIFI_IO_CANCEL)); }
+  static IoResult zero() { return IoResult(0U); }
+
+  template <typename T>
+  static IoResult from(T val) {
+    static_assert(std::is_same_v<T, int64_t> || std::is_same_v<T, size_t> || 
std::is_same_v<T, uint64_t>,
+                  "IoResult::from() can only be called with uint64_t, int64_t 
or size_t");
+
+    if constexpr (std::is_same_v<T, int64_t>) {
+      if (val < 0) {
+        return 
IoResult(nonstd::make_unexpected(static_cast<MinifiIoStatus>(val)));
+      }
+    } else if constexpr (std::is_same_v<T, size_t>) {
+      if (isError(val)) {
+        return IoResult::error();
+      }
+    }
+
+    // Common return path for the valid cases of both types
+    return IoResult(gsl::narrow<uint64_t>(val));
+  }
+
+  [[nodiscard]] int64_t toI64() const {
+    if (result_.has_value()) { return gsl::narrow<int64_t>(*result_); }
+    return result_.error();
+  }
+
+  [[nodiscard]] bool is_cancelled() const { return !result_ && result_.error() 
== MINIFI_IO_CANCEL; }
+
+  bool operator()() const { return result_.has_value(); }
+  bool operator!() const { return !result_.has_value(); }
+
+  uint64_t operator*() const { return *result_; }
+
+  nonstd::expected<uint64_t, MinifiIoStatus> inner() const { return result_; }
+
+ private:
+  explicit IoResult(nonstd::expected<uint64_t, MinifiIoStatus> result) : 
result_(std::move(result)) {}
+
+  nonstd::expected<uint64_t, MinifiIoStatus> result_;
+};
+
+class ReadWriteResult {
+ public:
+  ReadWriteResult() = delete;
+  ReadWriteResult(const ReadWriteResult&) = default;
+  ReadWriteResult(ReadWriteResult&&) noexcept = default;
+  ReadWriteResult& operator=(ReadWriteResult&&) noexcept = default;
+  ReadWriteResult& operator=(const ReadWriteResult&) = default;
+
+  ReadWriteResult(const uint64_t bytes_read, const uint64_t bytes_written)
+      : result_(ReadWrite{.bytes_read = bytes_read, .bytes_written = 
bytes_written}) {}
+
+  static ReadWriteResult zero() { return ReadWriteResult(ReadWrite{.bytes_read 
= 0, .bytes_written = 0}); };
+  static ReadWriteResult error() { return 
ReadWriteResult(nonstd::make_unexpected(MINIFI_IO_ERROR)); }
+  static ReadWriteResult cancelled() { return 
ReadWriteResult(nonstd::make_unexpected(MINIFI_IO_CANCEL)); }
+
+  virtual ~ReadWriteResult() = default;
+
+  bool operator()() const { return result_.has_value(); }
+  bool operator!() const { return !result_.has_value(); }
+
+  [[nodiscard]] uint64_t bytesWritten() const { return result_->bytes_written; 
}
+  [[nodiscard]] uint64_t bytesRead() const { return result_->bytes_read; }
+
+ private:
+  struct ReadWrite {
+    uint64_t bytes_read;
+    uint64_t bytes_written;
+  };
+  explicit ReadWriteResult(nonstd::expected<ReadWrite, MinifiIoStatus> result) 
: result_(std::move(result)) {}
+
+  nonstd::expected<ReadWrite, MinifiIoStatus> result_;
 };
 
-// FlowFile IO Callback functions for input and output
-// throw exception for error
-using InputStreamCallback = std::function<int64_t(const 
std::shared_ptr<InputStream>& input_stream)>;
-using OutputStreamCallback = std::function<int64_t(const 
std::shared_ptr<OutputStream>& output_stream)>;
-using InputOutputStreamCallback = 
std::function<std::optional<ReadWriteResult>(const 
std::shared_ptr<InputStream>& input_stream, const 
std::shared_ptr<OutputStream>& output_stream)>;
+using OutputStreamCallback = std::function<IoResult(const 
std::shared_ptr<OutputStream>& output_stream)>;
+using InputStreamCallback = std::function<IoResult(const 
std::shared_ptr<InputStream>& input_stream)>;
+using InputOutputStreamCallback =
+    std::function<ReadWriteResult(const std::shared_ptr<InputStream>& 
input_stream, const std::shared_ptr<OutputStream>& output_stream)>;
 
 }  // namespace org::apache::nifi::minifi::io
diff --git a/minifi-api/include/minifi-cpp/core/Property.h 
b/minifi-api/include/minifi-cpp/core/Property.h
index e4fb1f986..8d1225494 100644
--- a/minifi-api/include/minifi-cpp/core/Property.h
+++ b/minifi-api/include/minifi-cpp/core/Property.h
@@ -48,7 +48,7 @@ class Property final {
 
   Property(const PropertyReference &);
 
-  virtual ~Property() = default;
+  ~Property() = default;
 
   void setTransient() { is_transient_ = true; }
 

Reply via email to