safesight-edge/plugins/alarm/actions/http_action.cpp
sladro 644fb590fa
Some checks are pending
CI / host-build (push) Waiting to run
CI / rk3588-cross-build (push) Waiting to run
修改bug
2026-01-05 16:02:45 +08:00

175 lines
5.3 KiB
C++

#include "http_action.h"
#include <algorithm>
#include <chrono>
#include <ctime>
#include <iostream>
#include <sstream>
#if __has_include(<curl/curl.h>)
#include <curl/curl.h>
#define HAS_CURL 1
#else
#define HAS_CURL 0
#endif
namespace rk3588 {
HttpAction::~HttpAction() {
Drain();
Stop();
}
bool HttpAction::Init(const SimpleJson& config) {
url_ = config.ValueOr<std::string>("url", "");
timeout_ms_ = config.ValueOr<int>("timeout_ms", 3000);
include_media_url_ = config.ValueOr<bool>("include_media_url", true);
method_ = config.ValueOr<std::string>("method", "POST");
max_queue_size_ = static_cast<size_t>(std::max(0, config.ValueOr<int>("max_queue_size", static_cast<int>(max_queue_size_))));
const std::string policy = config.ValueOr<std::string>("queue_policy", "drop_oldest");
if (policy == "drop_newest") queue_policy_ = QueuePolicy::DropNewest;
else queue_policy_ = QueuePolicy::DropOldest;
if (url_.empty()) {
std::cerr << "[HttpAction] url is required\n";
return false;
}
#if HAS_CURL
curl_global_init(CURL_GLOBAL_DEFAULT);
#endif
running_.store(true);
worker_ = std::thread(&HttpAction::WorkerLoop, this);
std::cout << "[HttpAction] initialized, url=" << url_ << "\n";
return true;
}
void HttpAction::Execute(AlarmEvent& event, std::shared_ptr<Frame> /*frame*/) {
// Build JSON body
std::ostringstream oss;
oss << "{";
oss << "\"node_id\":\"" << event.node_id << "\",";
oss << "\"rule_name\":\"" << event.rule_name << "\",";
oss << "\"timestamp\":" << event.timestamp_ms << ",";
oss << "\"frame_id\":" << event.frame_id << ",";
if (include_media_url_) {
if (!event.snapshot_url.empty()) {
oss << "\"snapshot_url\":\"" << event.snapshot_url << "\",";
}
if (!event.clip_url.empty()) {
oss << "\"clip_url\":\"" << event.clip_url << "\",";
}
}
oss << "\"detections\":[";
for (size_t i = 0; i < event.detections.size(); ++i) {
const auto& det = event.detections[i];
if (i > 0) oss << ",";
oss << "{\"cls_id\":" << det.cls_id
<< ",\"score\":" << det.score
<< ",\"bbox\":{\"x\":" << det.bbox.x
<< ",\"y\":" << det.bbox.y
<< ",\"w\":" << det.bbox.w
<< ",\"h\":" << det.bbox.h << "}"
<< ",\"track_id\":" << det.track_id << "}";
}
oss << "]}";
std::string json_body = oss.str();
{
std::lock_guard<std::mutex> lock(queue_mutex_);
if (max_queue_size_ > 0 && request_queue_.size() >= max_queue_size_) {
++dropped_total_;
if (queue_policy_ == QueuePolicy::DropOldest) {
request_queue_.pop();
request_queue_.push(json_body);
} else {
// DropNewest: drop this request
}
} else {
request_queue_.push(json_body);
}
}
queue_cv_.notify_one();
}
void HttpAction::Drain() {
if (!running_.load()) return;
// 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() && in_flight_ == 0; });
}
void HttpAction::Stop() {
bool was_running = running_.exchange(false);
(void)was_running;
queue_cv_.notify_all();
if (worker_.joinable()) worker_.join();
}
void HttpAction::WorkerLoop() {
while (running_.load()) {
std::string body;
{
std::unique_lock<std::mutex> lock(queue_mutex_);
queue_cv_.wait_for(lock, std::chrono::milliseconds(100),
[this] { return !request_queue_.empty() || !running_.load(); });
if (!running_.load() && request_queue_.empty()) break;
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();
}
}
bool HttpAction::SendRequest(const std::string& json_body) {
#if HAS_CURL
CURL* curl = curl_easy_init();
if (!curl) {
std::cerr << "[HttpAction] curl_easy_init failed\n";
return false;
}
struct curl_slist* headers = nullptr;
headers = curl_slist_append(headers, "Content-Type: application/json");
curl_easy_setopt(curl, CURLOPT_URL, url_.c_str());
curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
curl_easy_setopt(curl, CURLOPT_POSTFIELDS, json_body.c_str());
curl_easy_setopt(curl, CURLOPT_TIMEOUT_MS, timeout_ms_);
curl_easy_setopt(curl, CURLOPT_CONNECTTIMEOUT_MS, timeout_ms_ / 2);
CURLcode res = curl_easy_perform(curl);
if (res != CURLE_OK) {
std::cerr << "[HttpAction] request failed: " << curl_easy_strerror(res) << "\n";
}
curl_slist_free_all(headers);
curl_easy_cleanup(curl);
return res == CURLE_OK;
#else
// Stub implementation for Windows without curl
std::cout << "[HttpAction] HTTP POST to " << url_ << "\n"
<< " Body: " << json_body << "\n";
return true;
#endif
}
} // namespace rk3588