From b62391fa97132a52417106e27d82c9bfebe953d4 Mon Sep 17 00:00:00 2001 From: sladro Date: Mon, 12 Jan 2026 17:53:55 +0800 Subject: [PATCH] =?UTF-8?q?=E6=80=A7=E8=83=BD=E4=BC=98=E5=8C=962=EF=BC=8C?= =?UTF-8?q?=E6=9C=89=E5=BE=85=E6=B5=8B=E8=AF=95=E7=89=88=E6=9C=AC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- include/ai_scheduler.h | 5 + include/utils/spsc_queue.h | 276 ++++++++++-------- plugins/CMakeLists.txt | 10 +- plugins/ai_yolo/ai_yolo_node.cpp | 10 + plugins/osd/osd_node.cpp | 372 ++++++++++++++++++++++++- plugins/preprocess/preprocess_node.cpp | 32 ++- plugins/publish/publish_node.cpp | 128 +++++---- src/ai_scheduler.cpp | 144 ++++++++-- 8 files changed, 758 insertions(+), 219 deletions(-) diff --git a/include/ai_scheduler.h b/include/ai_scheduler.h index 0450537..dc605a3 100644 --- a/include/ai_scheduler.h +++ b/include/ai_scheduler.h @@ -34,6 +34,11 @@ struct InferInput { int height = 0; bool is_nhwc = true; // true: NHWC, false: NCHW + // Optional DMA-BUF input for RKNN zero-copy (best-effort). + // When dma_fd >= 0, AiScheduler will try rknn_create_mem_from_fd + rknn_set_io_mem. + int dma_fd = -1; + int dma_offset = 0; + #if defined(RK3588_ENABLE_RKNN) // The actual data type of `data` passed to RKNN. Default preserves existing behavior. rknn_tensor_type type = RKNN_TENSOR_UINT8; diff --git a/include/utils/spsc_queue.h b/include/utils/spsc_queue.h index 2599bff..8d76072 100644 --- a/include/utils/spsc_queue.h +++ b/include/utils/spsc_queue.h @@ -1,10 +1,13 @@ #pragma once +#include #include #include #include +#include #include #include +#include #include #include @@ -30,106 +33,96 @@ public: SpscQueue(size_t capacity, QueueDropStrategy strategy) : capacity_(capacity == 0 ? 1 : capacity), strategy_(strategy) { - buf_.resize(capacity_); + buffer_.resize(capacity_); + for (size_t i = 0; i < capacity_; ++i) { + buffer_[i].seq.store(i, std::memory_order_relaxed); + } } // Called when the queue transitions from empty -> non-empty (best-effort). // Must be fast and must not call back into this queue. void SetOnDataAvailable(std::function cb) { - std::lock_guard lock(mu_); + std::lock_guard lock(cb_mu_); on_data_available_ = std::move(cb); } bool Push(T item) { - std::function on_data; - bool notify_data = false; - { - std::unique_lock lock(mu_); - if (stop_) return false; - if (size_ >= capacity_) { - if (strategy_ == QueueDropStrategy::DropOldest) { - // Drop the oldest item. - head_ = (head_ + 1) % capacity_; - --size_; - ++dropped_; - } else if (strategy_ == QueueDropStrategy::DropNewest) { - // Drop the incoming item. - ++dropped_; - return true; - } else { - // Block until space is available - space_cv_.wait(lock, [&] { return size_ < capacity_ || stop_; }); - if (stop_) return false; + while (true) { + if (stop_.load(std::memory_order_acquire)) return false; + + if (TryEnqueue(item)) { + pushed_.fetch_add(1, std::memory_order_relaxed); + const size_t prev_size = size_.fetch_add(1, std::memory_order_acq_rel); + if (prev_size == 0) { + std::function cb; + { + std::lock_guard lock(cb_mu_); + cb = on_data_available_; + } + if (cb) cb(); + data_cv_.notify_one(); } + return true; } - const bool was_empty = (size_ == 0); - buf_[tail_] = std::move(item); - tail_ = (tail_ + 1) % capacity_; - ++size_; - ++pushed_; - - // Avoid per-item wakeups when consumer is already running. - if (was_empty) { - notify_data = true; - on_data = on_data_available_; + // Queue full. + if (strategy_ == QueueDropStrategy::DropNewest) { + dropped_.fetch_add(1, std::memory_order_relaxed); + return true; } - } - if (on_data) { - on_data(); + if (strategy_ == QueueDropStrategy::DropOldest) { + if (TryDropOneOldest()) { + continue; + } + std::this_thread::yield(); + continue; + } + + // Block until space is available. + std::unique_lock lock(wait_mu_); + space_cv_.wait(lock, [&] { + return stop_.load(std::memory_order_acquire) || + size_.load(std::memory_order_acquire) < capacity_; + }); } - if (notify_data) { - data_cv_.notify_one(); - } - return true; } bool Pop(T& out, std::chrono::milliseconds timeout) { - std::unique_lock lock(mu_); - if (!data_cv_.wait_for(lock, timeout, [&] { return size_ > 0 || stop_; })) { + if (TryPop(out)) return true; + std::unique_lock lock(wait_mu_); + if (!data_cv_.wait_for(lock, timeout, [&] { + return stop_.load(std::memory_order_acquire) || + size_.load(std::memory_order_acquire) > 0; + })) { return false; } - if (size_ == 0) return false; - - const bool was_full = (size_ >= capacity_); - out = std::move(buf_[head_]); - head_ = (head_ + 1) % capacity_; - --size_; - ++popped_; - if (was_full) { - space_cv_.notify_one(); - } - return true; + return TryPop(out); } // Blocks until an item is available or the queue is stopped. bool Pop(T& out) { - std::unique_lock lock(mu_); - data_cv_.wait(lock, [&] { return size_ > 0 || stop_; }); - if (size_ == 0) return false; - - const bool was_full = (size_ >= capacity_); - out = std::move(buf_[head_]); - head_ = (head_ + 1) % capacity_; - --size_; - ++popped_; - if (was_full) { - space_cv_.notify_one(); + while (true) { + if (TryPop(out)) return true; + std::unique_lock lock(wait_mu_); + data_cv_.wait(lock, [&] { + return stop_.load(std::memory_order_acquire) || + size_.load(std::memory_order_acquire) > 0; + }); + if (stop_.load(std::memory_order_acquire) && + size_.load(std::memory_order_acquire) == 0) { + return false; + } } - return true; } // Non-blocking pop. bool TryPop(T& out) { - std::unique_lock lock(mu_); - if (size_ == 0) return false; - const bool was_full = (size_ >= capacity_); - out = std::move(buf_[head_]); - head_ = (head_ + 1) % capacity_; - --size_; - ++popped_; - if (was_full) { + if (!TryDequeue(out)) return false; + popped_.fetch_add(1, std::memory_order_relaxed); + const size_t prev_size = size_.fetch_sub(1, std::memory_order_acq_rel); + if (prev_size == capacity_) { + // Queue was full before this pop. space_cv_.notify_one(); } return true; @@ -138,85 +131,144 @@ public: // Pop up to max_items currently available (non-blocking). Returns true if any item was popped. bool TryPopBatch(std::vector& out, size_t max_items) { if (max_items == 0) return false; - std::unique_lock lock(mu_); - if (size_ == 0) return false; - const bool was_full = (size_ >= capacity_); - const size_t n = (size_ < max_items) ? size_ : max_items; out.clear(); - out.reserve(n); - for (size_t i = 0; i < n; ++i) { - out.push_back(std::move(buf_[head_])); - head_ = (head_ + 1) % capacity_; + out.reserve(max_items); + for (size_t i = 0; i < max_items; ++i) { + T item; + if (!TryPop(item)) break; + out.push_back(std::move(item)); } - size_ -= n; - popped_ += n; - if (was_full) { - space_cv_.notify_one(); - } - return true; + return !out.empty(); } void Stop() { - std::lock_guard lock(mu_); - stop_ = true; + stop_.store(true, std::memory_order_release); data_cv_.notify_all(); space_cv_.notify_all(); } bool IsStopped() const { - std::lock_guard lock(mu_); - return stop_; + return stop_.load(std::memory_order_acquire); } size_t Size() const { - std::lock_guard lock(mu_); - return size_; + return size_.load(std::memory_order_acquire); } size_t Capacity() const { return capacity_; } size_t DroppedCount() const { - std::lock_guard lock(mu_); - return dropped_; + return dropped_.load(std::memory_order_acquire); } size_t PushedCount() const { - std::lock_guard lock(mu_); - return pushed_; + return pushed_.load(std::memory_order_acquire); } size_t PoppedCount() const { - std::lock_guard lock(mu_); - return popped_; + return popped_.load(std::memory_order_acquire); } Stats GetStats() const { - std::lock_guard lock(mu_); Stats s; - s.size = size_; + s.size = size_.load(std::memory_order_acquire); s.capacity = capacity_; - s.dropped = dropped_; - s.pushed = pushed_; - s.popped = popped_; - s.stopped = stop_; + s.dropped = dropped_.load(std::memory_order_acquire); + s.pushed = pushed_.load(std::memory_order_acquire); + s.popped = popped_.load(std::memory_order_acquire); + s.stopped = stop_.load(std::memory_order_acquire); return s; } private: + struct Cell { + std::atomic seq; + T data; + }; + + // Lock-free bounded MPMC queue (Vyukov). We still expose it as SPSC in this project, + // but DropOldest may dequeue from the producer thread. + bool TryEnqueue(T& item) { + size_t pos = enqueue_pos_.load(std::memory_order_relaxed); + while (true) { + Cell& cell = buffer_[pos % capacity_]; + const size_t seq = cell.seq.load(std::memory_order_acquire); + const intptr_t dif = static_cast(seq) - static_cast(pos); + if (dif == 0) { + if (enqueue_pos_.compare_exchange_weak(pos, pos + 1, std::memory_order_relaxed)) { + cell.data = std::move(item); + cell.seq.store(pos + 1, std::memory_order_release); + return true; + } + } else if (dif < 0) { + return false; // full + } else { + pos = enqueue_pos_.load(std::memory_order_relaxed); + } + } + } + + bool TryDequeue(T& out) { + size_t pos = dequeue_pos_.load(std::memory_order_relaxed); + while (true) { + Cell& cell = buffer_[pos % capacity_]; + const size_t seq = cell.seq.load(std::memory_order_acquire); + const intptr_t dif = static_cast(seq) - static_cast(pos + 1); + if (dif == 0) { + if (dequeue_pos_.compare_exchange_weak(pos, pos + 1, std::memory_order_relaxed)) { + out = std::move(cell.data); + cell.seq.store(pos + capacity_, std::memory_order_release); + return true; + } + } else if (dif < 0) { + return false; // empty + } else { + pos = dequeue_pos_.load(std::memory_order_relaxed); + } + } + } + + bool TryDropOneOldest() { + size_t pos = dequeue_pos_.load(std::memory_order_relaxed); + while (true) { + Cell& cell = buffer_[pos % capacity_]; + const size_t seq = cell.seq.load(std::memory_order_acquire); + const intptr_t dif = static_cast(seq) - static_cast(pos + 1); + if (dif == 0) { + if (dequeue_pos_.compare_exchange_weak(pos, pos + 1, std::memory_order_relaxed)) { + // Drop without moving data out. + cell.seq.store(pos + capacity_, std::memory_order_release); + size_.fetch_sub(1, std::memory_order_acq_rel); + dropped_.fetch_add(1, std::memory_order_relaxed); + return true; + } + } else if (dif < 0) { + return false; + } else { + pos = dequeue_pos_.load(std::memory_order_relaxed); + } + } + } + size_t capacity_ = 0; QueueDropStrategy strategy_ = QueueDropStrategy::DropOldest; - mutable std::mutex mu_; + alignas(64) std::atomic enqueue_pos_{0}; + alignas(64) std::atomic dequeue_pos_{0}; + std::vector buffer_; + + alignas(64) std::atomic size_{0}; + std::atomic stop_{false}; + + std::atomic dropped_{0}; + std::atomic pushed_{0}; + std::atomic popped_{0}; + + mutable std::mutex wait_mu_; std::condition_variable data_cv_; std::condition_variable space_cv_; + + mutable std::mutex cb_mu_; std::function on_data_available_; - std::vector buf_; - size_t head_ = 0; - size_t tail_ = 0; - size_t size_ = 0; - bool stop_ = false; - size_t dropped_ = 0; - size_t pushed_ = 0; - size_t popped_ = 0; }; } // namespace rk3588 diff --git a/plugins/CMakeLists.txt b/plugins/CMakeLists.txt index 182e57b..a9a80cb 100644 --- a/plugins/CMakeLists.txt +++ b/plugins/CMakeLists.txt @@ -290,9 +290,17 @@ set_target_properties(tracker PROPERTIES ) # osd plugin (on-screen display for detection results) -add_library(osd SHARED osd/osd_node.cpp) +add_library(osd SHARED + osd/osd_node.cpp + ${CMAKE_SOURCE_DIR}/src/utils/dma_alloc.cpp +) target_include_directories(osd PRIVATE ${CMAKE_SOURCE_DIR}/include ${CMAKE_SOURCE_DIR}/third_party) target_link_libraries(osd PRIVATE project_options Threads::Threads) +if(RK3588_ENABLE_RGA AND RK_RGA_LIB) + target_compile_definitions(osd PRIVATE RK3588_ENABLE_RGA) + target_include_directories(osd PRIVATE ${RK_RGA_INCLUDE_DIR}) + target_link_libraries(osd PRIVATE ${RK_RGA_LIB}) +endif() set_target_properties(osd PROPERTIES OUTPUT_NAME "osd" LIBRARY_OUTPUT_DIRECTORY ${RK_PLUGIN_OUTPUT_DIR} diff --git a/plugins/ai_yolo/ai_yolo_node.cpp b/plugins/ai_yolo/ai_yolo_node.cpp index fb0871b..a9aa72a 100644 --- a/plugins/ai_yolo/ai_yolo_node.cpp +++ b/plugins/ai_yolo/ai_yolo_node.cpp @@ -2,6 +2,7 @@ #include #include #include +#include #include #include #include @@ -427,6 +428,15 @@ private: if (static_cast(src_stride) == packed_row && frame->data_size >= packed_size) { input.data = src; input.size = packed_size; + + // Best-effort RKNN DMA-BUF zero-copy path. + if (frame->dma_fd >= 0 && frame->data) { + const ptrdiff_t off = src - frame->data; + if (off >= 0 && static_cast(off) + packed_size <= frame->data_size) { + input.dma_fd = frame->dma_fd; + input.dma_offset = static_cast(off); + } + } } else { if (frame->data_size < static_cast(src_stride) * static_cast(h)) { std::cerr << "[ai_yolo] invalid RGB buffer size/stride (data_size=" << frame->data_size diff --git a/plugins/osd/osd_node.cpp b/plugins/osd/osd_node.cpp index 77918b2..450423b 100644 --- a/plugins/osd/osd_node.cpp +++ b/plugins/osd/osd_node.cpp @@ -10,8 +10,15 @@ #include "face/face_result.h" #include "node.h" +#include "utils/dma_alloc.h" #include "utils/logger.h" +#if defined(RK3588_ENABLE_RGA) +#include "im2d.hpp" +#include "im2d_buffer.h" +#include "im2d_type.h" +#endif + namespace rk3588 { namespace { @@ -57,6 +64,62 @@ inline int Clamp(int val, int min_val, int max_val) { return val < min_val ? min_val : (val > max_val ? max_val : val); } +#if defined(RK3588_ENABLE_RGA) +inline uint32_t PackColorArgb(const Color& c) { + return (0xFFu << 24) | (static_cast(c.r) << 16) | + (static_cast(c.g) << 8) | static_cast(c.b); +} + +inline bool WrapFrameAsRgaDst(const Frame& frame, rga_buffer_t& dst, int& rga_stride_bytes) { + if (frame.dma_fd < 0) return false; + + int rga_fmt = 0; + int stride_bytes = 0; + int wstride = 0; + int hstride = frame.height; + + if (frame.format == PixelFormat::RGB) { + rga_fmt = RK_FORMAT_RGB_888; + stride_bytes = frame.planes[0].stride > 0 ? frame.planes[0].stride : (frame.stride > 0 ? frame.stride : frame.width * 3); + wstride = stride_bytes / 3; + if (frame.planes[0].size > 0 && stride_bytes > 0) { + const int hs = frame.planes[0].size / stride_bytes; + if (hs >= frame.height) hstride = hs; + } + } else if (frame.format == PixelFormat::BGR) { + rga_fmt = RK_FORMAT_BGR_888; + stride_bytes = frame.planes[0].stride > 0 ? frame.planes[0].stride : (frame.stride > 0 ? frame.stride : frame.width * 3); + wstride = stride_bytes / 3; + if (frame.planes[0].size > 0 && stride_bytes > 0) { + const int hs = frame.planes[0].size / stride_bytes; + if (hs >= frame.height) hstride = hs; + } + } else if (frame.format == PixelFormat::NV12) { + rga_fmt = RK_FORMAT_YCbCr_420_SP; + stride_bytes = frame.planes[0].stride > 0 ? frame.planes[0].stride : (frame.stride > 0 ? frame.stride : frame.width); + wstride = stride_bytes; + if (frame.planes[0].size > 0 && stride_bytes > 0) { + const int hs = frame.planes[0].size / stride_bytes; + if (hs >= frame.height) hstride = hs; + } + } else if (frame.format == PixelFormat::YUV420) { + rga_fmt = RK_FORMAT_YCbCr_420_P; + stride_bytes = frame.planes[0].stride > 0 ? frame.planes[0].stride : (frame.stride > 0 ? frame.stride : frame.width); + wstride = stride_bytes; + if (frame.planes[0].size > 0 && stride_bytes > 0) { + const int hs = frame.planes[0].size / stride_bytes; + if (hs >= frame.height) hstride = hs; + } + } else { + return false; + } + + rga_stride_bytes = stride_bytes; + dst = wrapbuffer_fd_t(frame.dma_fd, frame.width, frame.height, wstride, hstride, rga_fmt); + return true; +} +#endif + void DrawHLine(uint8_t* data, int w, int h, int stride, PixelFormat fmt, int x1, int x2, int y, int thickness, const Color& color) { x1 = Clamp(x1, 0, w - 1); @@ -290,6 +353,7 @@ public: draw_face_det_ = config.ValueOr("draw_face_det", draw_face_det_); line_width_ = config.ValueOr("line_width", 2); font_scale_ = config.ValueOr("font_scale", 1); + use_rga_bbox_ = config.ValueOr("use_rga_bbox", use_rga_bbox_); if (const SimpleJson* dbg = config.Find("debug"); dbg && dbg->IsObject()) { stats_log_ = dbg->ValueOr("stats", stats_log_); @@ -321,7 +385,13 @@ public: bool Start() override { LogInfo(std::string("[osd] start id=") + id_ + " draw_bbox=" + (draw_bbox_ ? "true" : "false") + " draw_text=" + (draw_text_ ? "true" : "false") + - " draw_face_det=" + (draw_face_det_ ? "true" : "false")); + " draw_face_det=" + (draw_face_det_ ? "true" : "false") + + " use_rga_bbox=" + (use_rga_bbox_ ? "true" : "false")); +#if !defined(RK3588_ENABLE_RGA) + if (use_rga_bbox_ && draw_bbox_) { + LogError("[osd] use_rga_bbox=true but RK3588_ENABLE_RGA is disabled at build; bbox will NOT be drawn id=" + id_); + } +#endif return true; } @@ -355,6 +425,83 @@ private: static FramePtr CloneFrameForWrite(const FramePtr& src) { if (!src || !src->data || src->data_size == 0) return src; + // Prefer DMA clone when src already lives in DMA-BUF. + if (src->dma_fd >= 0) { + if (auto dma = DmaAlloc(src->data_size); dma && dma->valid()) { +#if defined(RK3588_ENABLE_RGA) + // Fast path: RGA copy. + rga_buffer_t rga_src{}; + rga_buffer_t rga_dst{}; + int dummy_stride = 0; + Frame dst_frame = *src; + dst_frame.dma_fd = dma->fd; + dst_frame.data = dma->data(); + dst_frame.data_size = dma->size; + dst_frame.data_owner = dma; + if (WrapFrameAsRgaDst(*src, rga_src, dummy_stride) && + WrapFrameAsRgaDst(dst_frame, rga_dst, dummy_stride)) { + if (imcopy(rga_src, rga_dst, 1, nullptr) == IM_STATUS_SUCCESS) { + auto out = std::make_shared(*src); + out->dma_fd = dma->fd; + out->data = dma->data(); + out->data_size = dma->size; + out->data_owner = dma; + + for (int i = 0; i < 3; ++i) { + FramePlane p = src->planes[i]; + int offset = p.offset; + if (p.data && src->data) { + const ptrdiff_t d = p.data - src->data; + if (d >= 0 && static_cast(d) < src->data_size) { + offset = static_cast(d); + } + } + p.offset = offset; + if (offset >= 0 && static_cast(offset) < out->data_size) { + p.data = out->data + offset; + } else { + p.data = nullptr; + } + out->planes[i] = p; + } + return out; + } + } +#endif + // Fallback: CPU memcpy (requires cache sync on linux). + DmaSyncStartFd(src->dma_fd); + DmaSyncStart(dma.get()); + memcpy(dma->data(), src->data, src->data_size); + DmaSyncEnd(dma.get()); + DmaSyncEndFd(src->dma_fd); + + auto out = std::make_shared(*src); + out->dma_fd = dma->fd; + out->data = dma->data(); + out->data_size = dma->size; + out->data_owner = dma; + + for (int i = 0; i < 3; ++i) { + FramePlane p = src->planes[i]; + int offset = p.offset; + if (p.data && src->data) { + const ptrdiff_t d = p.data - src->data; + if (d >= 0 && static_cast(d) < src->data_size) { + offset = static_cast(d); + } + } + p.offset = offset; + if (offset >= 0 && static_cast(offset) < out->data_size) { + p.data = out->data + offset; + } else { + p.data = nullptr; + } + out->planes[i] = p; + } + return out; + } + } + auto buf = std::make_shared>(src->data_size); memcpy(buf->data(), src->data, src->data_size); @@ -409,6 +556,61 @@ private: uint8_t* data = frame->planes[0].data ? frame->planes[0].data : frame->data; PixelFormat fmt = frame->format; + const bool want_rga_bbox = use_rga_bbox_ && draw_bbox_; + bool rga_bbox_ok = false; + +#if defined(RK3588_ENABLE_RGA) + if (want_rga_bbox) { + if (frame->dma_fd < 0) { + LogRgaBboxErrorThrottled("frame has no dma_fd"); + } else if (!(fmt == PixelFormat::RGB || fmt == PixelFormat::BGR || + fmt == PixelFormat::NV12 || fmt == PixelFormat::YUV420)) { + LogRgaBboxErrorThrottled("unsupported format=" + std::to_string(static_cast(fmt))); + } else { + rga_buffer_t dst{}; + int rga_stride_bytes = 0; + if (!WrapFrameAsRgaDst(*frame, dst, rga_stride_bytes)) { + LogRgaBboxErrorThrottled("WrapFrameAsRgaDst failed"); + } else { + (void)rga_stride_bytes; + im_job_handle_t job = imbeginJob(); + if (!job) { + LogRgaBboxErrorThrottled("imbeginJob failed"); + } else { + bool ok = true; + for (const auto& det : frame->det->items) { + int x = Clamp(static_cast(det.bbox.x), 0, w - 1); + int y = Clamp(static_cast(det.bbox.y), 0, h - 1); + int rw = static_cast(det.bbox.w); + int rh = static_cast(det.bbox.h); + rw = Clamp(rw, 1, w - x); + rh = Clamp(rh, 1, h - y); + im_rect rect{x, y, rw, rh}; + const uint32_t c = PackColorArgb(GetClassColor(det.cls_id)); + if (imrectangleTask(job, dst, rect, c, line_width_) <= 0) { + ok = false; + break; + } + } + if (ok) { + ok = (imendJob(job, IM_SYNC, 0, nullptr) > 0); + } else { + (void)imendJob(job, IM_SYNC, 0, nullptr); + } + rga_bbox_ok = ok; + if (!rga_bbox_ok) { + LogRgaBboxErrorThrottled("imrectangle/imendJob failed"); + } + } + } + } + } +#else + if (want_rga_bbox) { + LogRgaBboxErrorThrottled("RK3588_ENABLE_RGA is disabled at build"); + } +#endif + int stride; if (fmt == PixelFormat::RGB || fmt == PixelFormat::BGR) { stride = frame->planes[0].stride > 0 ? frame->planes[0].stride : (frame->stride > 0 ? frame->stride : w * 3); @@ -418,6 +620,13 @@ private: return; // Unsupported format } + const bool do_cpu_bbox = draw_bbox_ && !use_rga_bbox_; + + const bool do_cpu_text = draw_text_; + if ((do_cpu_bbox || do_cpu_text) && frame->dma_fd >= 0) { + DmaSyncStartFd(frame->dma_fd); + } + for (const auto& det : frame->det->items) { int x1 = static_cast(det.bbox.x); int y1 = static_cast(det.bbox.y); @@ -426,11 +635,11 @@ private: Color color = GetClassColor(det.cls_id); - if (draw_bbox_) { + if (do_cpu_bbox) { DrawRect(data, w, h, stride, fmt, x1, y1, x2, y2, line_width_, color); } - if (draw_text_) { + if (do_cpu_text) { char label[96]; if (det.track_id >= 0) { snprintf(label, sizeof(label), "%s %.0f%% id=%d", GetLabel(det.cls_id), det.score * 100, det.track_id); @@ -442,6 +651,10 @@ private: DrawText(data, w, h, stride, fmt, x1, text_y, label, font_scale_, color); } } + + if ((do_cpu_bbox || do_cpu_text) && frame->dma_fd >= 0) { + DmaSyncEndFd(frame->dma_fd); + } } void DrawFaceRecog(FramePtr frame) { @@ -452,6 +665,61 @@ private: uint8_t* data = frame->planes[0].data ? frame->planes[0].data : frame->data; PixelFormat fmt = frame->format; + const bool want_rga_bbox = use_rga_bbox_ && draw_bbox_; + bool rga_bbox_ok = false; + +#if defined(RK3588_ENABLE_RGA) + if (want_rga_bbox) { + if (frame->dma_fd < 0) { + LogRgaBboxErrorThrottled("frame has no dma_fd"); + } else if (!(fmt == PixelFormat::RGB || fmt == PixelFormat::BGR || + fmt == PixelFormat::NV12 || fmt == PixelFormat::YUV420)) { + LogRgaBboxErrorThrottled("unsupported format=" + std::to_string(static_cast(fmt))); + } else { + rga_buffer_t dst{}; + int rga_stride_bytes = 0; + if (!WrapFrameAsRgaDst(*frame, dst, rga_stride_bytes)) { + LogRgaBboxErrorThrottled("WrapFrameAsRgaDst failed"); + } else { + (void)rga_stride_bytes; + im_job_handle_t job = imbeginJob(); + if (!job) { + LogRgaBboxErrorThrottled("imbeginJob failed"); + } else { + bool ok = true; + for (const auto& it : frame->face_recog->items) { + int x = Clamp(static_cast(it.bbox.x), 0, w - 1); + int y = Clamp(static_cast(it.bbox.y), 0, h - 1); + int rw = static_cast(it.bbox.w); + int rh = static_cast(it.bbox.h); + rw = Clamp(rw, 1, w - x); + rh = Clamp(rh, 1, h - y); + im_rect rect{x, y, rw, rh}; + const Color color = it.unknown ? Color{255, 0, 0} : Color{0, 255, 0}; + if (imrectangleTask(job, dst, rect, PackColorArgb(color), line_width_) <= 0) { + ok = false; + break; + } + } + if (ok) { + ok = (imendJob(job, IM_SYNC, 0, nullptr) > 0); + } else { + (void)imendJob(job, IM_SYNC, 0, nullptr); + } + rga_bbox_ok = ok; + if (!rga_bbox_ok) { + LogRgaBboxErrorThrottled("imrectangle/imendJob failed"); + } + } + } + } + } +#else + if (want_rga_bbox) { + LogRgaBboxErrorThrottled("RK3588_ENABLE_RGA is disabled at build"); + } +#endif + int stride; if (fmt == PixelFormat::RGB || fmt == PixelFormat::BGR) { stride = frame->planes[0].stride > 0 ? frame->planes[0].stride : (frame->stride > 0 ? frame->stride : w * 3); @@ -461,6 +729,12 @@ private: return; } + const bool do_cpu_bbox = draw_bbox_ && !use_rga_bbox_; + const bool do_cpu_text = draw_text_; + if ((do_cpu_bbox || do_cpu_text) && frame->dma_fd >= 0) { + DmaSyncStartFd(frame->dma_fd); + } + for (const auto& it : frame->face_recog->items) { int x1 = static_cast(it.bbox.x); int y1 = static_cast(it.bbox.y); @@ -469,11 +743,11 @@ private: const Color color = it.unknown ? Color{255, 0, 0} : Color{0, 255, 0}; - if (draw_bbox_) { + if (do_cpu_bbox) { DrawRect(data, w, h, stride, fmt, x1, y1, x2, y2, line_width_, color); } - if (draw_text_) { + if (do_cpu_text) { char label[96]; const char* name = it.best_name.empty() ? (it.unknown ? "unknown" : "") : it.best_name.c_str(); snprintf(label, sizeof(label), "%s %.2f", name, it.best_sim); @@ -482,6 +756,10 @@ private: DrawText(data, w, h, stride, fmt, x1, text_y, label, font_scale_, color); } } + + if ((do_cpu_bbox || do_cpu_text) && frame->dma_fd >= 0) { + DmaSyncEndFd(frame->dma_fd); + } } void DrawFaceDet(FramePtr frame) { @@ -492,6 +770,61 @@ private: uint8_t* data = frame->planes[0].data ? frame->planes[0].data : frame->data; PixelFormat fmt = frame->format; + const bool want_rga_bbox = use_rga_bbox_ && draw_bbox_; + bool rga_bbox_ok = false; + +#if defined(RK3588_ENABLE_RGA) + if (want_rga_bbox) { + if (frame->dma_fd < 0) { + LogRgaBboxErrorThrottled("frame has no dma_fd"); + } else if (!(fmt == PixelFormat::RGB || fmt == PixelFormat::BGR || + fmt == PixelFormat::NV12 || fmt == PixelFormat::YUV420)) { + LogRgaBboxErrorThrottled("unsupported format=" + std::to_string(static_cast(fmt))); + } else { + rga_buffer_t dst{}; + int rga_stride_bytes = 0; + if (!WrapFrameAsRgaDst(*frame, dst, rga_stride_bytes)) { + LogRgaBboxErrorThrottled("WrapFrameAsRgaDst failed"); + } else { + (void)rga_stride_bytes; + im_job_handle_t job = imbeginJob(); + if (!job) { + LogRgaBboxErrorThrottled("imbeginJob failed"); + } else { + bool ok = true; + const Color color{0, 255, 255}; + for (const auto& it : frame->face_det->faces) { + int x = Clamp(static_cast(it.bbox.x), 0, w - 1); + int y = Clamp(static_cast(it.bbox.y), 0, h - 1); + int rw = static_cast(it.bbox.w); + int rh = static_cast(it.bbox.h); + rw = Clamp(rw, 1, w - x); + rh = Clamp(rh, 1, h - y); + im_rect rect{x, y, rw, rh}; + if (imrectangleTask(job, dst, rect, PackColorArgb(color), line_width_) <= 0) { + ok = false; + break; + } + } + if (ok) { + ok = (imendJob(job, IM_SYNC, 0, nullptr) > 0); + } else { + (void)imendJob(job, IM_SYNC, 0, nullptr); + } + rga_bbox_ok = ok; + if (!rga_bbox_ok) { + LogRgaBboxErrorThrottled("imrectangle/imendJob failed"); + } + } + } + } + } +#else + if (want_rga_bbox) { + LogRgaBboxErrorThrottled("RK3588_ENABLE_RGA is disabled at build"); + } +#endif + int stride; if (fmt == PixelFormat::RGB || fmt == PixelFormat::BGR) { stride = frame->planes[0].stride > 0 ? frame->planes[0].stride : (frame->stride > 0 ? frame->stride : w * 3); @@ -501,6 +834,12 @@ private: return; } + const bool do_cpu_bbox = draw_bbox_ && !use_rga_bbox_; + const bool do_cpu_text = draw_text_; + if ((do_cpu_bbox || do_cpu_text) && frame->dma_fd >= 0) { + DmaSyncStartFd(frame->dma_fd); + } + const Color color{0, 255, 255}; for (const auto& it : frame->face_det->faces) { int x1 = static_cast(it.bbox.x); @@ -508,10 +847,10 @@ private: int x2 = static_cast(it.bbox.x + it.bbox.w); int y2 = static_cast(it.bbox.y + it.bbox.h); - if (draw_bbox_) { + if (do_cpu_bbox) { DrawRect(data, w, h, stride, fmt, x1, y1, x2, y2, line_width_, color); } - if (draw_text_) { + if (do_cpu_text) { char label[64]; snprintf(label, sizeof(label), "face %.0f%%", it.score * 100); int text_y = y1 - 8 * font_scale_; @@ -519,16 +858,35 @@ private: DrawText(data, w, h, stride, fmt, x1, text_y, label, font_scale_, color); } } + + if ((do_cpu_bbox || do_cpu_text) && frame->dma_fd >= 0) { + DmaSyncEndFd(frame->dma_fd); + } } std::string id_; bool draw_bbox_ = true; bool draw_text_ = true; bool draw_face_det_ = true; + bool use_rga_bbox_ = true; int line_width_ = 2; int font_scale_ = 1; std::vector labels_; + void LogRgaBboxErrorThrottled(const std::string& reason) { + const auto now_ms = std::chrono::duration_cast( + std::chrono::steady_clock::now().time_since_epoch()) + .count(); + if (last_rga_bbox_err_ms_ > 0 && + now_ms - static_cast(last_rga_bbox_err_ms_) < 1000) { + return; + } + last_rga_bbox_err_ms_ = static_cast(now_ms); + LogError("[osd] RGA bbox failed: " + reason + " id=" + id_); + } + + uint64_t last_rga_bbox_err_ms_ = 0; + std::shared_ptr> input_queue_; std::vector>> output_queues_; uint64_t processed_ = 0; diff --git a/plugins/preprocess/preprocess_node.cpp b/plugins/preprocess/preprocess_node.cpp index 55c2094..b24366c 100644 --- a/plugins/preprocess/preprocess_node.cpp +++ b/plugins/preprocess/preprocess_node.cpp @@ -489,15 +489,9 @@ private: rga_buffer_t dst_buf{}; DmaBufferPtr src_dma_buf; // keep alive if we allocate/copy - // Prepare tmp buffer outside the RGA critical section. + // Lazily allocate tmp buffer only when we must fall back to a 2-step pipeline. + // (resize + cvtcolor). Prefer a single improcess() call to reduce scheduling overhead. DmaBufferPtr tmp_dma; - if (need_resize && need_cvt) { - tmp_dma = DmaAlloc(CalcImageSizeStrided(dst_wstride, dst_hstride, frame->format)); - if (!tmp_dma || !tmp_dma->valid()) { - std::cerr << "[preprocess] DMA alloc for tmp failed\n"; - return false; - } - } const bool can_cpu_read_src = (frame->data != nullptr && frame->data_size > 0); auto CopySrcToDmaIfPossible = [&]() -> bool { @@ -572,6 +566,28 @@ private: IM_STATUS status = IM_STATUS_SUCCESS; if (need_resize && need_cvt) { + // Try to fuse resize + CSC in a single call. + if (Check(src_buf, dst_buf)) { + rga_buffer_t pat{}; + const im_rect sr{0, 0, src_buf.width, src_buf.height}; + const im_rect dr{0, 0, dst_buf.width, dst_buf.height}; + const im_rect pr{0, 0, 0, 0}; + status = improcess(src_buf, dst_buf, pat, sr, dr, pr, + 0, nullptr, nullptr, IM_SYNC); + if (status == IM_STATUS_SUCCESS) { + return true; + } + } + + // Fallback: 2-step (resize + cvtcolor). + if (!tmp_dma || !tmp_dma->valid()) { + tmp_dma = DmaAlloc(CalcImageSizeStrided(dst_wstride, dst_hstride, frame->format)); + if (!tmp_dma || !tmp_dma->valid()) { + err = "DMA alloc for tmp failed"; + return false; + } + } + rga_buffer_t tmp = wrapbuffer_fd_t(tmp_dma->fd, out_w, out_h, dst_wstride, dst_hstride, src_fmt_rga); if (!Check(src_buf, tmp) || !Check(tmp, dst_buf)) { diff --git a/plugins/publish/publish_node.cpp b/plugins/publish/publish_node.cpp index 1842ab8..1b179fb 100644 --- a/plugins/publish/publish_node.cpp +++ b/plugins/publish/publish_node.cpp @@ -109,21 +109,68 @@ public: // Global init (safe to call multiple times) avformat_network_init(); + // Async write queue: bound memory and never block the realtime pipeline. + queue_ = std::make_shared>(queue_size_, QueueDropStrategy::DropOldest); + running_ = true; - monitor_thread_ = std::thread(&AvMuxer::MonitorLoop, this); + worker_thread_ = std::thread(&AvMuxer::WorkerLoop, this); LogInfo("[publish] muxer init async for " + url_); return true; } bool WriteFrame(const EncodedPacket& pkt) { - std::lock_guard lock(mutex_); - if (!ready_ || !fmt_ || !stream_) return false; + if (!running_.load(std::memory_order_acquire) || !queue_) return false; if (pkt.data.empty()) return false; + EncodedPacket copy = pkt; + return queue_->Push(std::move(copy)); + } + + void Close() { + running_ = false; + if (queue_) queue_->Stop(); + if (worker_thread_.joinable()) worker_thread_.join(); + } + +private: + void WorkerLoop() { + while (running_.load(std::memory_order_acquire)) { + if (!TryOpen()) { + if (!running_.load(std::memory_order_acquire)) break; + std::this_thread::sleep_for(std::chrono::seconds(1)); + continue; + } + + ready_ = true; + warned_ = false; + last_pts_ = -1; + LogInfo("[publish] muxer server ready: " + url_); + + while (running_.load(std::memory_order_acquire) && ready_) { + EncodedPacket pkt; + if (!queue_ || !queue_->Pop(pkt, std::chrono::milliseconds(100))) { + continue; + } + if (!WriteFrameInternal(pkt)) { + ready_ = false; + } + } + + CleanupContext(); + ready_ = false; + + if (running_.load(std::memory_order_acquire)) { + std::this_thread::sleep_for(std::chrono::seconds(1)); + } + } + CleanupContext(); + } + + bool WriteFrameInternal(const EncodedPacket& pkt) { + if (!fmt_ || !stream_ || pkt.data.empty()) return false; AVPacket* out = av_packet_alloc(); if (!out) return false; - if (av_new_packet(out, static_cast(pkt.data.size())) < 0) { av_packet_free(&out); return false; @@ -131,19 +178,18 @@ public: std::memcpy(out->data, pkt.data.data(), pkt.data.size()); out->stream_index = stream_->index; out->flags = pkt.key ? AV_PKT_FLAG_KEY : 0; - + // Simple PTS mapping int64_t pts = av_rescale_q(pkt.pts_ms, AVRational{1, 1000}, stream_->time_base); if (pts <= last_pts_) pts = last_pts_ + 1; last_pts_ = pts; - + out->pts = pts; out->dts = pts; out->duration = av_rescale_q(1000 / std::max(1, fps_), AVRational{1, 1000}, stream_->time_base); - int ret = av_interleaved_write_frame(fmt_, out); + const int ret = av_interleaved_write_frame(fmt_, out); av_packet_free(&out); - if (ret < 0) { if (!warned_) { char errbuf[128]; @@ -151,58 +197,19 @@ public: std::cerr << "[publish] write failed for " << url_ << ": " << errbuf << ", resetting...\n"; warned_ = true; } - // Signal monitor thread to reset - ready_ = false; - cv_.notify_all(); + return false; } - return ret >= 0; + return true; } - void Close() { - running_ = false; - // Break the av_io_open wait if possible (via CheckInterrupt) - cv_.notify_all(); - if (monitor_thread_.joinable()) monitor_thread_.join(); - } - -private: - void MonitorLoop() { - while (running_) { - if (TryOpen()) { - { - std::lock_guard lock(mutex_); - ready_ = true; - warned_ = false; - LogInfo("[publish] muxer server ready: " + url_); - } - - // Wait until error occurs or stop requested - std::unique_lock lock(mutex_); - cv_.wait(lock, [this] { return !running_ || !ready_; }); - } - - // Cleanup context - { - std::lock_guard lock(mutex_); - if (fmt_) { - // Try to write trailer if logical end (not socket error) - // But usually we are here because of error, so av_write_trailer might hang/fail. - // We skip trailer on error reset to avoid blocking. - if (fmt_->pb) { - avio_closep(&fmt_->pb); - } - avformat_free_context(fmt_); - fmt_ = nullptr; - stream_ = nullptr; - } - ready_ = false; - } - - if (running_) { - // Delay before retry - std::this_thread::sleep_for(std::chrono::seconds(1)); - } + void CleanupContext() { + if (!fmt_) return; + if (fmt_->pb) { + avio_closep(&fmt_->pb); } + avformat_free_context(fmt_); + fmt_ = nullptr; + stream_ = nullptr; } bool TryOpen() { @@ -265,7 +272,6 @@ private: av_dict_free(&opts); { - std::lock_guard lock(mutex_); fmt_ = fmt; stream_ = stream; } @@ -299,11 +305,11 @@ private: std::vector extradata_; std::string proto_; std::string url_; - + std::atomic running_{false}; - std::thread monitor_thread_; - std::mutex mutex_; - std::condition_variable cv_; + std::thread worker_thread_; + std::shared_ptr> queue_; + size_t queue_size_ = 64; bool ready_ = false; bool warned_ = false; diff --git a/src/ai_scheduler.cpp b/src/ai_scheduler.cpp index b017fde..78ba9dc 100644 --- a/src/ai_scheduler.cpp +++ b/src/ai_scheduler.cpp @@ -315,25 +315,68 @@ InferResult AiScheduler::Infer(ModelHandle handle, const InferInput& input) { return result; } - // Setup input - rknn_input inputs[1]; - memset(inputs, 0, sizeof(inputs)); - inputs[0].index = 0; - inputs[0].type = input.type; - inputs[0].size = input.size; - inputs[0].fmt = input.is_nhwc ? RKNN_TENSOR_NHWC : RKNN_TENSOR_NCHW; - inputs[0].buf = const_cast(input.data); - inputs[0].pass_through = 0; + struct InputMemGuard { + rknn_context ctx = 0; + rknn_tensor_mem* mem = nullptr; + ~InputMemGuard() { + if (ctx && mem) { + rknn_destroy_mem(ctx, mem); + mem = nullptr; + } + } + }; + InputMemGuard input_mem{ctx->ctx, nullptr}; + bool used_io_mem = false; - int ret = rknn_inputs_set(ctx->ctx, ctx->n_input, inputs); - if (ret < 0) { - result.error = "rknn_inputs_set failed: " + std::to_string(ret); - total_errors_.fetch_add(1); - return result; + // Best-effort RKNN zero-copy input via DMA-BUF. + if (input.dma_fd >= 0) { + void* virt_base = const_cast(input.data); + const int32_t offset = input.dma_offset; + if (offset != 0) { + virt_base = static_cast(virt_base) - offset; + } + input_mem.mem = rknn_create_mem_from_fd(ctx->ctx, input.dma_fd, + virt_base, + static_cast(input.size), + offset); + if (input_mem.mem) { + rknn_tensor_attr attr = ctx->input_attrs.empty() ? rknn_tensor_attr{} : ctx->input_attrs[0]; + attr.index = 0; + attr.type = input.type; + attr.fmt = input.is_nhwc ? RKNN_TENSOR_NHWC : RKNN_TENSOR_NCHW; + + const int mem_ret = rknn_set_io_mem(ctx->ctx, input_mem.mem, &attr); + if (mem_ret == 0) { + used_io_mem = true; + } else { + // Fallback to normal inputs_set path. + rknn_destroy_mem(ctx->ctx, input_mem.mem); + input_mem.mem = nullptr; + } + } + } + + if (!used_io_mem) { + // Setup input (legacy copy path). + rknn_input inputs[1]; + memset(inputs, 0, sizeof(inputs)); + inputs[0].index = 0; + inputs[0].type = input.type; + inputs[0].size = input.size; + inputs[0].fmt = input.is_nhwc ? RKNN_TENSOR_NHWC : RKNN_TENSOR_NCHW; + inputs[0].buf = const_cast(input.data); + inputs[0].pass_through = 0; + + int ret = rknn_inputs_set(ctx->ctx, ctx->n_input, inputs); + if (ret < 0) { + result.error = "rknn_inputs_set failed: " + std::to_string(ret); + total_errors_.fetch_add(1); + return result; + } } // Run inference - ret = rknn_run(ctx->ctx, nullptr); + int ret = rknn_run(ctx->ctx, nullptr); if (ret < 0) { result.error = "rknn_run failed: " + std::to_string(ret); total_errors_.fetch_add(1); @@ -428,24 +471,65 @@ AiScheduler::BorrowedInferResult AiScheduler::InferBorrowed(ModelHandle handle, result.infer_lock = std::unique_lock(ctx->infer_mutex); result.keepalive = ctx; - // Setup input - rknn_input inputs[1]; - memset(inputs, 0, sizeof(inputs)); - inputs[0].index = 0; - inputs[0].type = input.type; - inputs[0].size = input.size; - inputs[0].fmt = input.is_nhwc ? RKNN_TENSOR_NHWC : RKNN_TENSOR_NCHW; - inputs[0].buf = const_cast(input.data); - inputs[0].pass_through = 0; + struct InputMemGuard { + rknn_context ctx = 0; + rknn_tensor_mem* mem = nullptr; + ~InputMemGuard() { + if (ctx && mem) { + rknn_destroy_mem(ctx, mem); + mem = nullptr; + } + } + }; + InputMemGuard input_mem{ctx->ctx, nullptr}; + bool used_io_mem = false; - int ret = rknn_inputs_set(ctx->ctx, ctx->n_input, inputs); - if (ret < 0) { - result.error = "rknn_inputs_set failed: " + std::to_string(ret); - total_errors_.fetch_add(1); - return result; + // Best-effort RKNN zero-copy input via DMA-BUF. + if (input.dma_fd >= 0) { + void* virt_base = const_cast(input.data); + const int32_t offset = input.dma_offset; + if (offset != 0) { + virt_base = static_cast(virt_base) - offset; + } + input_mem.mem = rknn_create_mem_from_fd(ctx->ctx, input.dma_fd, + virt_base, + static_cast(input.size), + offset); + if (input_mem.mem) { + rknn_tensor_attr attr = ctx->input_attrs.empty() ? rknn_tensor_attr{} : ctx->input_attrs[0]; + attr.index = 0; + attr.type = input.type; + attr.fmt = input.is_nhwc ? RKNN_TENSOR_NHWC : RKNN_TENSOR_NCHW; + const int mem_ret = rknn_set_io_mem(ctx->ctx, input_mem.mem, &attr); + if (mem_ret == 0) { + used_io_mem = true; + } else { + rknn_destroy_mem(ctx->ctx, input_mem.mem); + input_mem.mem = nullptr; + } + } } - ret = rknn_run(ctx->ctx, nullptr); + if (!used_io_mem) { + // Setup input (legacy copy path). + rknn_input inputs[1]; + memset(inputs, 0, sizeof(inputs)); + inputs[0].index = 0; + inputs[0].type = input.type; + inputs[0].size = input.size; + inputs[0].fmt = input.is_nhwc ? RKNN_TENSOR_NHWC : RKNN_TENSOR_NCHW; + inputs[0].buf = const_cast(input.data); + inputs[0].pass_through = 0; + + int ret = rknn_inputs_set(ctx->ctx, ctx->n_input, inputs); + if (ret < 0) { + result.error = "rknn_inputs_set failed: " + std::to_string(ret); + total_errors_.fetch_add(1); + return result; + } + } + + int ret = rknn_run(ctx->ctx, nullptr); if (ret < 0) { result.error = "rknn_run failed: " + std::to_string(ret); total_errors_.fetch_add(1);