From 0db74e4453963a51ce01ff909a3708c278121150 Mon Sep 17 00:00:00 2001 From: ktro2828 Date: Sat, 17 Jan 2026 23:00:16 +0900 Subject: [PATCH] refactor: allow to speicfy cudaStream in builder function and constructor Signed-off-by: ktro2828 --- .../builder.hpp | 33 +++++++++++++------ .../jpeg_compressor.hpp | 6 ++-- .../src/builder.cpp | 10 +++--- .../src/jpeg_compressor/jetson.cpp | 17 ++++++---- .../src/jpeg_compressor/nvjpeg.cpp | 20 +++++++---- .../builder.hpp | 15 ++++++--- .../rectifier.hpp | 4 ++- .../src/builder.cpp | 4 +-- .../src/rectifier/npp.cpp | 20 +++++++---- 9 files changed, 86 insertions(+), 43 deletions(-) diff --git a/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/builder.hpp b/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/builder.hpp index 583b728..a186fb7 100644 --- a/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/builder.hpp +++ b/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/builder.hpp @@ -16,6 +16,8 @@ #include +#include + #include #include #include @@ -44,30 +46,35 @@ CompressionType to_compression_type(const std::string & str); * @brief Create a compressor processor. * * @param type Compression type + * @param stream CUDA stream to use for asynchronous operations * @return std::unique_ptr */ -std::unique_ptr create_compressor(CompressionType type); +std::unique_ptr create_compressor(CompressionType type, cudaStream_t stream = nullptr); /** * @brief Create a compressor processor. * * @param type Compression type name in string format + * @param stream CUDA stream to use for asynchronous operations * @return std::unique_ptr */ -std::unique_ptr create_compressor(const std::string & type); +std::unique_ptr create_compressor( + const std::string & type, cudaStream_t stream = nullptr); /** * @brief Create a compressor processor with a free function for the postprocess. * * @param type Compression type * @param fn Free function for the postprocess + * @param stream CUDA stream to use for asynchronous operations * @return std::unique_ptr */ template < typename F, std::enable_if_t, int> = 0> -inline std::unique_ptr create_compressor(CompressionType type, F fn) +inline std::unique_ptr create_compressor( + CompressionType type, F fn, cudaStream_t stream = nullptr) { - auto processor = create_compressor(type); + auto processor = create_compressor(type, stream); auto fp = static_cast(fn); if (fp) processor->register_postprocess(fp); return processor; @@ -78,13 +85,15 @@ inline std::unique_ptr create_compressor(CompressionType type, F fn) * * @param type Compression type name in string format * @param fn Free function for the postprocess + * @param stream CUDA stream to use for asynchronous operations * @return std::unique_ptr */ template < typename F, std::enable_if_t, int> = 0> -inline std::unique_ptr create_compressor(const std::string & type, F fn) +inline std::unique_ptr create_compressor( + const std::string & type, F fn, cudaStream_t stream = nullptr) { - return create_compressor(to_compression_type(type), fn); + return create_compressor(to_compression_type(type), fn, stream); }; /** @@ -92,12 +101,14 @@ inline std::unique_ptr create_compressor(const std::string & type, F * * @param type Compression type * @param obj Object that has a member function for the postprocess + * @param stream CUDA stream to use for asynchronous operations * @return std::unique_ptr */ template -inline std::unique_ptr create_compressor(CompressionType type, Obj * obj) +inline std::unique_ptr create_compressor( + CompressionType type, Obj * obj, cudaStream_t stream = nullptr) { - auto processor = create_compressor(type); + auto processor = create_compressor(type, stream); processor->register_postprocess(obj); return processor; } @@ -107,11 +118,13 @@ inline std::unique_ptr create_compressor(CompressionType type, Obj * * * @param type Compression type name in string format * @param obj Object that has a member function for the postprocess + * @param stream CUDA stream to use for asynchronous operations * @return std::unique_ptr */ template -inline std::unique_ptr create_compressor(const std::string & type, Obj * obj) +inline std::unique_ptr create_compressor( + const std::string & type, Obj * obj, cudaStream_t stream = nullptr) { - return create_compressor(to_compression_type(type), obj); + return create_compressor(to_compression_type(type), obj, stream); } } // namespace accelerated_image_processor::compression diff --git a/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/jpeg_compressor.hpp b/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/jpeg_compressor.hpp index c14a717..793129c 100644 --- a/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/jpeg_compressor.hpp +++ b/src/accelerated_image_processor_compression/include/accelerated_image_processor_compression/jpeg_compressor.hpp @@ -18,6 +18,8 @@ #include #include +#include + #include namespace accelerated_image_processor::compression @@ -57,7 +59,7 @@ class JPEGCompressor : public common::BaseProcessor //!< @brief Factory function to create a CPUJPEGCompressor. std::unique_ptr make_cpujpeg_compressor(); //!< @brief Factory function to create a NvJPEGCompressor. -std::unique_ptr make_nvjpeg_compressor(); +std::unique_ptr make_nvjpeg_compressor(cudaStream_t stream = nullptr); //!< @brief Factory function to create a JetsonJPEGCompressor. -std::unique_ptr make_jetsonjpeg_compressor(); +std::unique_ptr make_jetsonjpeg_compressor(cudaStream_t stream = nullptr); } // namespace accelerated_image_processor::compression diff --git a/src/accelerated_image_processor_compression/src/builder.cpp b/src/accelerated_image_processor_compression/src/builder.cpp index d21bc6f..3b50954 100644 --- a/src/accelerated_image_processor_compression/src/builder.cpp +++ b/src/accelerated_image_processor_compression/src/builder.cpp @@ -59,14 +59,14 @@ CompressionType to_compression_type(const std::string & str) } } -std::unique_ptr create_compressor(CompressionType type) +std::unique_ptr create_compressor(CompressionType type, cudaStream_t stream) { switch (type) { case CompressionType::JPEG: #ifdef JETSON_AVAILABLE - return make_jetsonjpeg_compressor(); + return make_jetsonjpeg_compressor(stream); #elif NVJPEG_AVAILABLE - return make_nvjpeg_compressor(); + return make_nvjpeg_compressor(stream); #elif TURBOJPEG_AVAILABLE return make_cpujpeg_compressor(); #else @@ -79,8 +79,8 @@ std::unique_ptr create_compressor(CompressionType type) } } -std::unique_ptr create_compressor(const std::string & type) +std::unique_ptr create_compressor(const std::string & type, cudaStream_t stream) { - return create_compressor(to_compression_type(type)); + return create_compressor(to_compression_type(type), stream); } } // namespace accelerated_image_processor::compression diff --git a/src/accelerated_image_processor_compression/src/jpeg_compressor/jetson.cpp b/src/accelerated_image_processor_compression/src/jpeg_compressor/jetson.cpp index 787ff51..3c45c5f 100644 --- a/src/accelerated_image_processor_compression/src/jpeg_compressor/jetson.cpp +++ b/src/accelerated_image_processor_compression/src/jpeg_compressor/jetson.cpp @@ -39,9 +39,13 @@ namespace accelerated_image_processor::compression class JetsonJPEGCompressor final : public JPEGCompressor { public: - JetsonJPEGCompressor() : JPEGCompressor(JPEGBackend::JETSON) + JetsonJPEGCompressor(cudaStream_t stream = nullptr) + : JPEGCompressor(JPEGBackend::JETSON), stream_(stream) { - CHECK_CUDA(cudaStreamCreate(&stream_)); + if (stream_ == nullptr) { + CHECK_CUDA(cudaStreamCreateWithFlags(&stream_, cudaStreamNonBlocking)); + own_stream_ = true; + } encoder_ = NvJPEGEncoder::createJPEGEncoder("jpeg_encoder"); } ~JetsonJPEGCompressor() override @@ -60,7 +64,7 @@ class JetsonJPEGCompressor final : public JPEGCompressor delete encoder_; encoder_ = nullptr; } - if (stream_) { + if (stream_ && own_stream_) { cudaStreamDestroy(stream_); stream_ = nullptr; } @@ -168,16 +172,17 @@ class JetsonJPEGCompressor final : public JPEGCompressor std::array yuv_step_bytes_{{0, 0, 0}}; //!< Step sizes in bytes for the YUV data. cudaStream_t stream_{nullptr}; //!< CUDA stream for asynchronous operations. + bool own_stream_{false}; //!< Whether the stream is owned by the compressor. NppStreamContext context_{}; //!< NPP stream context for asynchronous operations. std::optional buffer_; //!< Optional NvBuffer for storing encoded JPEG data. }; -std::unique_ptr make_jetsonjpeg_compressor() +std::unique_ptr make_jetsonjpeg_compressor(cudaStream_t stream) { - return std::make_unique(); + return std::make_unique(stream); } #else -std::unique_ptr make_jetsonjpeg_compressor() +std::unique_ptr make_jetsonjpeg_compressor(cudaStream_t) { return nullptr; } diff --git a/src/accelerated_image_processor_compression/src/jpeg_compressor/nvjpeg.cpp b/src/accelerated_image_processor_compression/src/jpeg_compressor/nvjpeg.cpp index 2241fd7..6d80fe1 100644 --- a/src/accelerated_image_processor_compression/src/jpeg_compressor/nvjpeg.cpp +++ b/src/accelerated_image_processor_compression/src/jpeg_compressor/nvjpeg.cpp @@ -37,9 +37,13 @@ namespace accelerated_image_processor::compression class NvJPEGCompressor final : public JPEGCompressor { public: - NvJPEGCompressor() : JPEGCompressor(JPEGBackend::NVJPEG) + NvJPEGCompressor(cudaStream_t stream = nullptr) + : JPEGCompressor(JPEGBackend::NVJPEG), stream_(stream) { - CHECK_CUDA(cudaStreamCreate(&stream_)); + if (stream_ == nullptr) { + CHECK_CUDA(cudaStreamCreateWithFlags(&stream_, cudaStreamNonBlocking)); + own_stream_ = true; + } CHECK_NVJPEG(nvjpegCreateSimple(&handle_)); CHECK_NVJPEG(nvjpegEncoderStateCreate(handle_, &state_, stream_)); CHECK_NVJPEG(nvjpegEncoderParamsCreate(handle_, ¶ms_, stream_)); @@ -53,7 +57,10 @@ class NvJPEGCompressor final : public JPEGCompressor CHECK_NVJPEG(nvjpegEncoderParamsDestroy(params_)); CHECK_NVJPEG(nvjpegEncoderStateDestroy(state_)); CHECK_NVJPEG(nvjpegDestroy(handle_)); - CHECK_CUDA(cudaStreamDestroy(stream_)); + if (stream_ && own_stream_) { + CHECK_CUDA(cudaStreamDestroy(stream_)); + stream_ = nullptr; + } } private: @@ -114,14 +121,15 @@ class NvJPEGCompressor final : public JPEGCompressor nvjpegEncoderState_t state_; //!< NVJPEG encoder state. nvjpegEncoderParams_t params_; //!< NVJPEG encoder parameters. nvjpegImage_t nv_image_; //!< NVJPEG image buffer. + bool own_stream_{false}; //!< Whether the stream is owned by the compressor. }; -std::unique_ptr make_nvjpeg_compressor() +std::unique_ptr make_nvjpeg_compressor(cudaStream_t stream) { - return std::make_unique(); + return std::make_unique(stream); } #else -std::unique_ptr make_nvjpeg_compressor() +std::unique_ptr make_nvjpeg_compressor(cudaStream_t) { return nullptr; } diff --git a/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/builder.hpp b/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/builder.hpp index 69f7545..225a769 100644 --- a/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/builder.hpp +++ b/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/builder.hpp @@ -18,6 +18,8 @@ #include +#include + #include #include @@ -26,22 +28,24 @@ namespace accelerated_image_processor::pipeline /** * @brief Create a rectifier object. * + * @param stream CUDA stream for asynchronous operations. * @return std::unique_ptr */ -std::unique_ptr create_rectifier(); +std::unique_ptr create_rectifier(cudaStream_t stream = nullptr); /** * @brief Create a rectifier object with a free function for the postprocess. * * @tparam F The type of the callback function. * @param fn Free function for the postprocess. + * @param stream CUDA stream for asynchronous operations. * @return std::unique_ptr */ template < typename F, std::enable_if_t, int> = 0> -inline std::unique_ptr create_rectifier(F fn) +inline std::unique_ptr create_rectifier(F fn, cudaStream_t stream = nullptr) { - auto processor = create_rectifier(); + auto processor = create_rectifier(stream); auto fp = static_cast(fn); if (fp) processor->register_postprocess(fp); return processor; @@ -53,12 +57,13 @@ inline std::unique_ptr create_rectifier(F fn) * @tparam Obj The type of the object. * @tparam Method The type of the member function. * @param obj Pointer to the object. + * @param stream CUDA stream for asynchronous operations. * @return std::unique_ptr */ template -inline std::unique_ptr create_rectifier(Obj * obj) +inline std::unique_ptr create_rectifier(Obj * obj, cudaStream_t stream = nullptr) { - auto processor = create_rectifier(); + auto processor = create_rectifier(stream); processor->register_postprocess(obj); return processor; } diff --git a/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/rectifier.hpp b/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/rectifier.hpp index 1be3ba7..f833f76 100644 --- a/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/rectifier.hpp +++ b/src/accelerated_image_processor_pipeline/include/accelerated_image_processor_pipeline/rectifier.hpp @@ -18,6 +18,8 @@ #include #include +#include + #include #include @@ -89,7 +91,7 @@ class Rectifier : public common::BaseProcessor }; //!< @brief Factory function to create a NppRectifier. -std::unique_ptr make_npp_rectifier(); +std::unique_ptr make_npp_rectifier(cudaStream_t stream = nullptr); //!< @brief Factory function to create a OpenCvCudaRectifier. std::unique_ptr make_opencv_cuda_rectifier(); //!< @brief Factory function to create a CpuRectifier. diff --git a/src/accelerated_image_processor_pipeline/src/builder.cpp b/src/accelerated_image_processor_pipeline/src/builder.cpp index df94d0d..713ec57 100644 --- a/src/accelerated_image_processor_pipeline/src/builder.cpp +++ b/src/accelerated_image_processor_pipeline/src/builder.cpp @@ -21,10 +21,10 @@ namespace accelerated_image_processor::pipeline { -std::unique_ptr create_rectifier() +std::unique_ptr create_rectifier(cudaStream_t stream) { #ifdef NPP_AVAILABLE - return make_npp_rectifier(); + return make_npp_rectifier(stream); #elif OPENCV_CUDA_AVAILABLE return make_opencv_cuda_rectifier(); #else diff --git a/src/accelerated_image_processor_pipeline/src/rectifier/npp.cpp b/src/accelerated_image_processor_pipeline/src/rectifier/npp.cpp index 4e823cb..167712f 100644 --- a/src/accelerated_image_processor_pipeline/src/rectifier/npp.cpp +++ b/src/accelerated_image_processor_pipeline/src/rectifier/npp.cpp @@ -18,6 +18,7 @@ #include #include +#include #ifdef NPP_AVAILABLE #include @@ -35,9 +36,12 @@ namespace accelerated_image_processor::pipeline class NppRectifier final : public Rectifier { public: - NppRectifier() : Rectifier(RectifierBackend::NPP) + NppRectifier(cudaStream_t stream = nullptr) : Rectifier(RectifierBackend::NPP), stream_(stream) { - cudaStreamCreateWithFlags(&stream_, cudaStreamNonBlocking); + if (stream_ == nullptr) { + CHECK_CUDA(cudaStreamCreateWithFlags(&stream_, cudaStreamNonBlocking)); + own_stream_ = true; + } nppSetStream(stream_); } ~NppRectifier() override @@ -58,7 +62,10 @@ class NppRectifier final : public Rectifier nppiFree(dst_); dst_ = nullptr; } - cudaStreamDestroy(stream_); + if (stream_ && own_stream_) { + cudaStreamDestroy(stream_); + stream_ = nullptr; + } } private: @@ -134,14 +141,15 @@ class NppRectifier final : public Rectifier int src_step_{0}; int dst_step_{0}; cudaStream_t stream_; + bool own_stream_{false}; }; -std::unique_ptr make_npp_rectifier() +std::unique_ptr make_npp_rectifier(cudaStream_t stream) { - return std::make_unique(); + return std::make_unique(stream); } #else -std::unique_ptr make_npp_rectifier() +std::unique_ptr make_npp_rectifier(cudaStream_t) { return nullptr; }