修复阶段4之前的问题
This commit is contained in:
parent
36dc4911f1
commit
c36563bc87
@ -39,6 +39,7 @@ private:
|
||||
std::string name_;
|
||||
std::vector<NodeEntry> nodes_;
|
||||
std::atomic<bool> running_{false};
|
||||
std::atomic<bool> stop_requested_{false};
|
||||
};
|
||||
|
||||
class GraphManager {
|
||||
|
||||
@ -21,6 +21,7 @@ public:
|
||||
|
||||
bool Push(T item) {
|
||||
std::unique_lock<std::mutex> lock(mu_);
|
||||
if (stop_) return false;
|
||||
if (queue_.size() >= capacity_) {
|
||||
if (strategy_ == QueueDropStrategy::DropOldest) {
|
||||
queue_.pop_front();
|
||||
@ -55,6 +56,11 @@ public:
|
||||
space_cv_.notify_all();
|
||||
}
|
||||
|
||||
bool IsStopped() const {
|
||||
std::lock_guard<std::mutex> lock(mu_);
|
||||
return stop_;
|
||||
}
|
||||
|
||||
size_t Size() const {
|
||||
std::lock_guard<std::mutex> lock(mu_);
|
||||
return queue_.size();
|
||||
|
||||
@ -96,7 +96,7 @@ void ClipAction::PushPostEventFrame(std::shared_ptr<Frame> frame) {
|
||||
void ClipAction::Drain() {
|
||||
std::unique_lock<std::mutex> lock(queue_mutex_);
|
||||
queue_cv_.wait_for(lock, std::chrono::seconds(30),
|
||||
[this] { return task_queue_.empty(); });
|
||||
[this] { return task_queue_.empty() && in_flight_ == 0; });
|
||||
}
|
||||
|
||||
void ClipAction::WorkerLoop() {
|
||||
@ -110,6 +110,7 @@ void ClipAction::WorkerLoop() {
|
||||
if (task_queue_.empty()) continue;
|
||||
task = std::move(task_queue_.front());
|
||||
task_queue_.pop();
|
||||
++in_flight_;
|
||||
}
|
||||
|
||||
// Wait a bit for post-event frames if needed
|
||||
@ -128,6 +129,12 @@ void ClipAction::WorkerLoop() {
|
||||
if (!url.empty()) {
|
||||
std::cout << "[ClipAction] clip uploaded: " << url << "\n";
|
||||
}
|
||||
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(queue_mutex_);
|
||||
if (in_flight_ > 0) --in_flight_;
|
||||
}
|
||||
queue_cv_.notify_all();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@ -50,6 +50,7 @@ private:
|
||||
std::mutex queue_mutex_;
|
||||
std::condition_variable queue_cv_;
|
||||
std::queue<ClipTask> task_queue_;
|
||||
size_t in_flight_ = 0;
|
||||
|
||||
// Post-event collection state
|
||||
std::mutex post_mutex_;
|
||||
|
||||
@ -84,10 +84,10 @@ void HttpAction::Execute(AlarmEvent& event, std::shared_ptr<Frame> /*frame*/) {
|
||||
}
|
||||
|
||||
void HttpAction::Drain() {
|
||||
// Wait for queue to empty
|
||||
// Wait for queue to be fully drained (including in-flight request).
|
||||
std::unique_lock<std::mutex> lock(queue_mutex_);
|
||||
queue_cv_.wait_for(lock, std::chrono::seconds(5),
|
||||
[this] { return request_queue_.empty(); });
|
||||
[this] { return request_queue_.empty() && in_flight_ == 0; });
|
||||
}
|
||||
|
||||
void HttpAction::WorkerLoop() {
|
||||
@ -101,11 +101,18 @@ void HttpAction::WorkerLoop() {
|
||||
if (request_queue_.empty()) continue;
|
||||
body = std::move(request_queue_.front());
|
||||
request_queue_.pop();
|
||||
++in_flight_;
|
||||
}
|
||||
|
||||
if (!body.empty()) {
|
||||
SendRequest(body);
|
||||
}
|
||||
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(queue_mutex_);
|
||||
if (in_flight_ > 0) --in_flight_;
|
||||
}
|
||||
queue_cv_.notify_all();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@ -32,6 +32,7 @@ private:
|
||||
std::mutex queue_mutex_;
|
||||
std::condition_variable queue_cv_;
|
||||
std::queue<std::string> request_queue_;
|
||||
size_t in_flight_ = 0;
|
||||
};
|
||||
|
||||
} // namespace rk3588
|
||||
|
||||
@ -46,12 +46,12 @@ public:
|
||||
std::cerr << "[input_rtsp] no downstream queue configured for node " << id_ << "\n";
|
||||
return false;
|
||||
}
|
||||
out_queue_ = ctx.output_queues[0];
|
||||
out_queues_ = ctx.output_queues;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool Start() override {
|
||||
if (!out_queue_) return false;
|
||||
if (out_queues_.empty()) return false;
|
||||
running_.store(true);
|
||||
#if defined(RK3588_ENABLE_FFMPEG) && defined(RK3588_ENABLE_MPP)
|
||||
if (use_mpp_) {
|
||||
@ -89,7 +89,7 @@ public:
|
||||
|
||||
void Stop() override {
|
||||
running_.store(false);
|
||||
if (out_queue_) out_queue_->Stop();
|
||||
for (auto& q : out_queues_) q->Stop();
|
||||
if (worker_.joinable()) worker_.join();
|
||||
}
|
||||
|
||||
@ -98,6 +98,12 @@ public:
|
||||
}
|
||||
|
||||
private:
|
||||
void PushToDownstream(FramePtr frame) {
|
||||
for (auto& q : out_queues_) {
|
||||
q->Push(frame);
|
||||
}
|
||||
}
|
||||
|
||||
void LoopStub() {
|
||||
using namespace std::chrono;
|
||||
auto frame_interval = fps_ > 0 ? milliseconds(1000 / fps_) : milliseconds(40);
|
||||
@ -111,11 +117,13 @@ private:
|
||||
frame->planes[1] = {nullptr, width_, width_ * height_ / 2, width_ * height_};
|
||||
frame->frame_id = ++frame_id_;
|
||||
frame->pts = duration_cast<milliseconds>(steady_clock::now().time_since_epoch()).count();
|
||||
out_queue_->Push(frame);
|
||||
PushToDownstream(frame);
|
||||
|
||||
if (frame_id_ % 100 == 0) {
|
||||
std::cout << "[input_rtsp] generated frame " << frame_id_ << " queue="
|
||||
<< out_queue_->Size() << " drops=" << out_queue_->DroppedCount() << "\n";
|
||||
<< (out_queues_.empty() ? 0 : out_queues_[0]->Size())
|
||||
<< " drops=" << (out_queues_.empty() ? 0 : out_queues_[0]->DroppedCount())
|
||||
<< "\n";
|
||||
}
|
||||
std::this_thread::sleep_for(frame_interval);
|
||||
}
|
||||
@ -246,13 +254,12 @@ private:
|
||||
frame->planes[2] = {buffer->data() + y_size + u_size, uv_stride, u_size, y_size + u_size};
|
||||
}
|
||||
|
||||
if (!out_queue_->Push(frame)) {
|
||||
// dropped by queue policy
|
||||
}
|
||||
PushToDownstream(frame);
|
||||
if (frame_id_ % 100 == 0) {
|
||||
std::cout << "[input_rtsp] recv frame " << frame->frame_id
|
||||
<< " queue=" << out_queue_->Size()
|
||||
<< " drops=" << out_queue_->DroppedCount() << "\n";
|
||||
<< " queue=" << (out_queues_.empty() ? 0 : out_queues_[0]->Size())
|
||||
<< " drops=" << (out_queues_.empty() ? 0 : out_queues_[0]->DroppedCount())
|
||||
<< "\n";
|
||||
}
|
||||
}
|
||||
|
||||
@ -511,12 +518,13 @@ private:
|
||||
frame->planes[2] = {frame->data + y_size + u_size, uv_stride, u_size, y_size + u_size};
|
||||
}
|
||||
|
||||
out_queue_->Push(frame);
|
||||
PushToDownstream(frame);
|
||||
|
||||
if (frame_id_ % 100 == 0) {
|
||||
std::cout << "[input_rtsp] mpp frame " << frame->frame_id
|
||||
<< " queue=" << out_queue_->Size()
|
||||
<< " drops=" << out_queue_->DroppedCount() << "\n";
|
||||
<< " queue=" << (out_queues_.empty() ? 0 : out_queues_[0]->Size())
|
||||
<< " drops=" << (out_queues_.empty() ? 0 : out_queues_[0]->DroppedCount())
|
||||
<< "\n";
|
||||
}
|
||||
}
|
||||
#endif
|
||||
@ -527,7 +535,7 @@ private:
|
||||
int width_ = 1920;
|
||||
int height_ = 1080;
|
||||
std::atomic<bool> running_{false};
|
||||
std::shared_ptr<SpscQueue<FramePtr>> out_queue_;
|
||||
std::vector<std::shared_ptr<SpscQueue<FramePtr>>> out_queues_;
|
||||
std::thread worker_;
|
||||
uint64_t frame_id_ = 0;
|
||||
bool use_ffmpeg_ = false;
|
||||
|
||||
@ -153,6 +153,8 @@ bool Graph::Start() {
|
||||
return true; // Already running
|
||||
}
|
||||
|
||||
stop_requested_.store(false);
|
||||
|
||||
for (auto& entry : nodes_) {
|
||||
if (!entry.enabled || !entry.node) continue;
|
||||
|
||||
@ -166,10 +168,17 @@ bool Graph::Start() {
|
||||
if (entry.context.input_queue) {
|
||||
entry.worker = std::thread([this, &entry]() {
|
||||
FramePtr frame;
|
||||
while (running_) {
|
||||
while (true) {
|
||||
if (entry.context.input_queue->Pop(frame, std::chrono::milliseconds(100))) {
|
||||
if (frame) {
|
||||
entry.node->Process(frame);
|
||||
if (frame) entry.node->Process(frame);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (stop_requested_.load()) {
|
||||
if (entry.context.input_queue->IsStopped() &&
|
||||
entry.context.input_queue->Size() == 0) {
|
||||
for (auto& q : entry.context.output_queues) q->Stop();
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
@ -182,40 +191,46 @@ bool Graph::Start() {
|
||||
}
|
||||
|
||||
void Graph::Stop() {
|
||||
bool expected = true;
|
||||
if (!running_.compare_exchange_strong(expected, false)) {
|
||||
// Already stopped or stopping, but we need to ensure threads are joined if called from destructor
|
||||
// If we are in destructor, running_ might be false but threads joined?
|
||||
// We should just proceed to cleanup to be safe.
|
||||
if (!running_.load()) {
|
||||
return;
|
||||
}
|
||||
|
||||
// 1. Drain
|
||||
for (auto& n : nodes_) {
|
||||
if (n.node) n.node->Drain();
|
||||
}
|
||||
stop_requested_.store(true);
|
||||
|
||||
// 2. Wait for data to flush
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(200));
|
||||
auto is_source_like = [](const NodeEntry& n) {
|
||||
return n.role == "source" || !n.context.input_queue;
|
||||
};
|
||||
|
||||
// 3. Stop queues
|
||||
// 1) Stop sources first to ensure no new data is produced.
|
||||
for (auto& n : nodes_) {
|
||||
if (n.context.input_queue) n.context.input_queue->Stop();
|
||||
if (!n.enabled || !n.node) continue;
|
||||
if (!is_source_like(n)) continue;
|
||||
|
||||
n.node->Drain();
|
||||
n.node->Stop();
|
||||
for (auto& q : n.context.output_queues) q->Stop();
|
||||
}
|
||||
|
||||
// 4. Join threads
|
||||
// 2) Join framework workers after upstream queues are stopped and drained.
|
||||
for (auto& n : nodes_) {
|
||||
if (n.worker.joinable()) {
|
||||
n.worker.join();
|
||||
}
|
||||
if (n.worker.joinable()) n.worker.join();
|
||||
}
|
||||
|
||||
// 5. Stop nodes
|
||||
// 3) Drain non-source nodes (no concurrent Process() at this point).
|
||||
for (auto& n : nodes_) {
|
||||
if (!n.enabled || !n.node) continue;
|
||||
if (is_source_like(n)) continue;
|
||||
n.node->Drain();
|
||||
}
|
||||
|
||||
// 4) Stop remaining nodes (reverse order).
|
||||
for (auto it = nodes_.rbegin(); it != nodes_.rend(); ++it) {
|
||||
if (it->node) {
|
||||
if (!it->enabled || !it->node) continue;
|
||||
if (is_source_like(*it)) continue;
|
||||
it->node->Stop();
|
||||
}
|
||||
}
|
||||
|
||||
running_.store(false);
|
||||
}
|
||||
|
||||
GraphManager::GraphManager(std::string plugin_dir)
|
||||
|
||||
Loading…
Reference in New Issue
Block a user