File sources now keep STATE_FILE on graceful shutdown (EOF or SIGTERM): the final batch snapshot is written instead of deleted, so the dashboard's /api/current-batch still shows the recount result after the process exits. On startup, STATE_FILE is cleared unless a live instance is detected via a /proc cmdline scan (rtsp/http source in argv or config). Live sources keep the old resume-at-start / delete-on-finish behavior. SPEC: V19, T21, B2. AGENTS: STATE_FILE lifecycle section.
482 lines
17 KiB
C++
482 lines
17 KiB
C++
#include "batch_store.hpp"
|
|
#include <filesystem>
|
|
#include <fstream>
|
|
#include <chrono>
|
|
#include <ctime>
|
|
#include <sstream>
|
|
#include <iomanip>
|
|
#include <algorithm>
|
|
|
|
static std::string now_iso() {
|
|
auto now = std::chrono::system_clock::now();
|
|
auto t = std::chrono::system_clock::to_time_t(now);
|
|
std::ostringstream ss;
|
|
ss << std::put_time(std::localtime(&t), "%Y-%m-%dT%H:%M:%S");
|
|
return ss.str();
|
|
}
|
|
|
|
BatchStore::BatchStore(const std::string& db_path,
|
|
const std::string& state_file,
|
|
const std::string& camera_name,
|
|
const std::string& object_label,
|
|
const std::string& cutoff_time,
|
|
float batch_timeout,
|
|
float ignore_batch_label_timeout,
|
|
int min_object_per_batch,
|
|
int min_duration_per_batch,
|
|
std::function<void(const std::string&)> logger)
|
|
: db_path_(db_path), state_file_(state_file), camera_name_(camera_name),
|
|
object_label_(object_label), cutoff_time_str_(cutoff_time),
|
|
batch_timeout_(batch_timeout),
|
|
ignore_batch_label_timeout_(ignore_batch_label_timeout),
|
|
min_object_per_batch_(min_object_per_batch),
|
|
min_duration_per_batch_(min_duration_per_batch),
|
|
logger_(logger) {
|
|
|
|
std::filesystem::path dbp(db_path_);
|
|
if (dbp.has_parent_path())
|
|
std::filesystem::create_directories(dbp.parent_path());
|
|
|
|
std::filesystem::path sfp(state_file_);
|
|
if (sfp.has_parent_path())
|
|
std::filesystem::create_directories(sfp.parent_path());
|
|
|
|
int rc = sqlite3_open(db_path_.c_str(), &db_);
|
|
if (rc != SQLITE_OK) throw std::runtime_error("Failed to open database: " + std::string(sqlite3_errmsg(db_)));
|
|
|
|
init_db();
|
|
current_state_ = load_state();
|
|
previous_state_ = current_state_;
|
|
}
|
|
|
|
BatchStore::~BatchStore() {
|
|
if (db_) sqlite3_close(db_);
|
|
}
|
|
|
|
void BatchStore::init_db() {
|
|
const char* sql1 = R"(
|
|
CREATE TABLE IF NOT EXISTS batches (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
counting_date TEXT NOT NULL,
|
|
batch_number INTEGER NOT NULL,
|
|
camera_name TEXT NOT NULL,
|
|
object_label TEXT NOT NULL,
|
|
count INTEGER NOT NULL,
|
|
start_time TEXT NOT NULL,
|
|
end_time TEXT NOT NULL,
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
UNIQUE(counting_date, batch_number, camera_name, object_label)
|
|
)
|
|
)";
|
|
const char* sql2 = R"(
|
|
CREATE TABLE IF NOT EXISTS daily_summaries (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
counting_date TEXT NOT NULL,
|
|
camera_name TEXT NOT NULL,
|
|
object_label TEXT NOT NULL,
|
|
total_count INTEGER NOT NULL DEFAULT 0,
|
|
total_batches INTEGER NOT NULL DEFAULT 0,
|
|
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
UNIQUE(counting_date, camera_name, object_label)
|
|
)
|
|
)";
|
|
char* err = nullptr;
|
|
sqlite3_exec(db_, sql1, nullptr, nullptr, &err);
|
|
if (err) { logger_("DB init error: " + std::string(err)); sqlite3_free(err); }
|
|
err = nullptr;
|
|
sqlite3_exec(db_, sql2, nullptr, nullptr, &err);
|
|
if (err) { logger_("DB init error: " + std::string(err)); sqlite3_free(err); }
|
|
}
|
|
|
|
std::string BatchStore::get_counting_date() const {
|
|
auto now = std::chrono::system_clock::now();
|
|
auto t = std::chrono::system_clock::to_time_t(now);
|
|
auto* tm = std::localtime(&t);
|
|
|
|
int h, m;
|
|
sscanf(cutoff_time_str_.c_str(), "%d:%d", &h, &m);
|
|
|
|
if (tm->tm_hour < h || (tm->tm_hour == h && tm->tm_min < m)) {
|
|
std::ostringstream ss;
|
|
ss << std::put_time(tm, "%Y-%m-%d");
|
|
return ss.str();
|
|
} else {
|
|
auto tomorrow = now + std::chrono::hours(24);
|
|
t = std::chrono::system_clock::to_time_t(tomorrow);
|
|
tm = std::localtime(&t);
|
|
std::ostringstream ss;
|
|
ss << std::put_time(tm, "%Y-%m-%d");
|
|
return ss.str();
|
|
}
|
|
}
|
|
|
|
json BatchStore::load_state() {
|
|
std::ifstream f(state_file_);
|
|
if (!f.good()) return nullptr;
|
|
|
|
try {
|
|
json state = json::parse(f);
|
|
std::string current_date = get_counting_date();
|
|
if (state.value("counting_date", "") != current_date) {
|
|
logger_("State file belongs to previous counting day. Finalizing.");
|
|
insert_batch(state["counting_date"], state["batch_number"],
|
|
state["count"], state["start_time"], now_iso());
|
|
std::filesystem::remove(state_file_);
|
|
return nullptr;
|
|
}
|
|
std::ostringstream ss;
|
|
ss << "Resumed batch #" << state["batch_number"] << " with count=" << state["count"];
|
|
logger_(ss.str());
|
|
reset_batch_timer();
|
|
return state;
|
|
} catch (const std::exception& e) {
|
|
logger_("Failed to load state: " + std::string(e.what()));
|
|
return nullptr;
|
|
}
|
|
}
|
|
|
|
void BatchStore::save_state() {
|
|
if (current_state_.is_null()) {
|
|
std::filesystem::remove(state_file_);
|
|
return;
|
|
}
|
|
std::ofstream f(state_file_);
|
|
f << current_state_.dump(2);
|
|
}
|
|
|
|
int BatchStore::get_next_batch_number(const std::string& counting_date) {
|
|
const char* sql = R"(
|
|
SELECT COALESCE(MAX(batch_number), 0)
|
|
FROM batches
|
|
WHERE counting_date = ? AND camera_name = ? AND object_label = ?
|
|
)";
|
|
sqlite3_stmt* stmt;
|
|
sqlite3_prepare_v2(db_, sql, -1, &stmt, nullptr);
|
|
sqlite3_bind_text(stmt, 1, counting_date.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 2, camera_name_.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 3, object_label_.c_str(), -1, SQLITE_TRANSIENT);
|
|
int result = 0;
|
|
if (sqlite3_step(stmt) == SQLITE_ROW)
|
|
result = sqlite3_column_int(stmt, 0) + 1;
|
|
sqlite3_finalize(stmt);
|
|
return result;
|
|
}
|
|
|
|
void BatchStore::start_new_batch(const std::string& counting_date) {
|
|
int batch_number = get_next_batch_number(counting_date);
|
|
std::string now = now_iso();
|
|
|
|
json counted_ids = json::array();
|
|
if (!previous_state_.is_null() && previous_state_.contains("counted_event_ids")) {
|
|
auto& ids = previous_state_["counted_event_ids"];
|
|
int start = std::max(0, static_cast<int>(ids.size()) - carry_ids_);
|
|
for (int i = start; i < (int)ids.size(); ++i)
|
|
counted_ids.push_back(ids[i]);
|
|
}
|
|
|
|
current_state_ = {
|
|
{"counting_date", counting_date},
|
|
{"batch_number", batch_number},
|
|
{"count", 0},
|
|
{"start_time", now},
|
|
{"last_detection_time", now},
|
|
{"counted_event_ids", counted_ids}
|
|
};
|
|
save_state();
|
|
|
|
std::ostringstream ss;
|
|
ss << "Started batch #" << batch_number << " for " << counting_date;
|
|
logger_(ss.str());
|
|
}
|
|
|
|
void BatchStore::reset_batch_timer() {
|
|
batch_timer_running_ = true;
|
|
uint64_t gen = ++batch_timer_gen_;
|
|
if (batch_timer_thread_.joinable()) batch_timer_thread_.detach();
|
|
batch_timer_thread_ = std::thread([this, gen]() {
|
|
auto start = std::chrono::steady_clock::now();
|
|
while (batch_timer_gen_ == gen && batch_timer_running_) {
|
|
auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(
|
|
std::chrono::steady_clock::now() - start).count();
|
|
if (elapsed >= batch_timeout_) {
|
|
on_batch_timeout();
|
|
break;
|
|
}
|
|
std::this_thread::sleep_for(std::chrono::seconds(1));
|
|
}
|
|
});
|
|
}
|
|
|
|
void BatchStore::on_batch_timeout() {
|
|
batch_timer_running_ = false;
|
|
logger_("Batch inactivity timeout reached");
|
|
end_batch("timeout");
|
|
}
|
|
|
|
void BatchStore::ignore_batch_label_fn() {
|
|
if (!ignore_label_active_) {
|
|
ignore_label_active_ = true;
|
|
ignore_batch_label_ = true;
|
|
if (ignore_label_thread_.joinable()) ignore_label_thread_.detach();
|
|
ignore_label_thread_ = std::thread([this]() {
|
|
auto start = std::chrono::steady_clock::now();
|
|
while (ignore_label_active_) {
|
|
auto elapsed = std::chrono::duration_cast<std::chrono::seconds>(
|
|
std::chrono::steady_clock::now() - start).count();
|
|
if (elapsed >= ignore_batch_label_timeout_) {
|
|
on_ignore_batch_label_timeout();
|
|
break;
|
|
}
|
|
std::this_thread::sleep_for(std::chrono::seconds(1));
|
|
}
|
|
});
|
|
}
|
|
}
|
|
|
|
void BatchStore::on_ignore_batch_label_timeout() {
|
|
ignore_label_active_ = false;
|
|
ignore_batch_label_ = false;
|
|
}
|
|
|
|
std::pair<int, bool> BatchStore::record_ayam_crossing(int track_id) {
|
|
std::lock_guard<std::mutex> lock(state_mutex_);
|
|
std::string counting_date = get_counting_date();
|
|
bool started_new = false;
|
|
|
|
if (current_state_.is_null()) {
|
|
start_new_batch(counting_date);
|
|
started_new = true;
|
|
} else if (current_state_["counting_date"].get<std::string>() != counting_date) {
|
|
end_batch_locked("cutoff");
|
|
start_new_batch(counting_date);
|
|
started_new = true;
|
|
cutoff_triggered_ = true;
|
|
}
|
|
|
|
std::string event_key = std::to_string(track_id);
|
|
auto& ids = current_state_["counted_event_ids"];
|
|
bool found = false;
|
|
for (const auto& id : ids) {
|
|
if (id.get<std::string>() == event_key) { found = true; break; }
|
|
}
|
|
if (!found) {
|
|
current_state_["count"] = current_state_["count"].get<int>() + 1;
|
|
ids.push_back(event_key);
|
|
}
|
|
|
|
current_state_["last_detection_time"] = now_iso();
|
|
save_state();
|
|
reset_batch_timer();
|
|
return {current_state_["count"].get<int>(), started_new};
|
|
}
|
|
|
|
bool BatchStore::record_talenan_crossing(int track_id, float confidence) {
|
|
if (ignore_batch_label_) return false;
|
|
std::lock_guard<std::mutex> lock(state_mutex_);
|
|
{
|
|
int count = current_state_.is_null() ? 0 : current_state_["count"].get<int>();
|
|
std::string start_time = current_state_.is_null() ? "" : current_state_["start_time"];
|
|
std::string end_time = now_iso();
|
|
auto duration_sec = 0LL;
|
|
if (!start_time.empty()) {
|
|
std::tm tm = {};
|
|
std::istringstream ss(start_time);
|
|
ss >> std::get_time(&tm, "%Y-%m-%dT%H:%M:%S");
|
|
auto start_tp = std::chrono::system_clock::from_time_t(std::mktime(&tm));
|
|
auto end_tp = std::chrono::system_clock::now();
|
|
duration_sec = std::chrono::duration_cast<std::chrono::seconds>(end_tp - start_tp).count();
|
|
}
|
|
std::ostringstream msg;
|
|
int bn = current_state_["batch_number"].get<int>();
|
|
double cps = duration_sec > 0 ? static_cast<double>(count) / duration_sec : 0.0;
|
|
msg << "Batch #" << bn << " ended"
|
|
<< " | count=" << count
|
|
<< " | duration=" << duration_sec << "s"
|
|
<< " | cps=" << std::fixed << std::setprecision(2) << cps
|
|
<< " | start=" << start_time
|
|
<< " | end=" << end_time
|
|
<< " | conf=" << std::fixed << std::setprecision(3) << confidence;
|
|
logger_(msg.str());
|
|
}
|
|
ignore_batch_label_fn();
|
|
end_batch_locked("talenan");
|
|
batch_timer_running_ = false;
|
|
return true;
|
|
}
|
|
|
|
void BatchStore::end_batch(const std::string& closed_by) {
|
|
std::lock_guard<std::mutex> lock(state_mutex_);
|
|
end_batch_locked(closed_by);
|
|
}
|
|
|
|
bool BatchStore::end_batch_locked(const std::string& closed_by) {
|
|
if (current_state_.is_null()) return false;
|
|
|
|
previous_state_ = current_state_;
|
|
auto state = current_state_;
|
|
|
|
std::string start_time = state["start_time"];
|
|
std::string end_time = now_iso();
|
|
|
|
// parse duration
|
|
std::tm tm = {};
|
|
std::istringstream ss(start_time);
|
|
ss >> std::get_time(&tm, "%Y-%m-%dT%H:%M:%S");
|
|
auto start_tp = std::chrono::system_clock::from_time_t(std::mktime(&tm));
|
|
auto end_tp = std::chrono::system_clock::now();
|
|
auto duration_sec = std::chrono::duration_cast<std::chrono::seconds>(end_tp - start_tp).count();
|
|
|
|
int count = state["count"].get<int>();
|
|
int batch_number = state["batch_number"].get<int>();
|
|
std::string counting_date = state.value("counting_date", "");
|
|
|
|
if (count < min_object_per_batch_ || duration_sec < min_duration_per_batch_) {
|
|
current_state_ = nullptr;
|
|
if (preserve_state_file_ && closed_by == "shutdown") {
|
|
std::ofstream f(state_file_);
|
|
f << state.dump(2);
|
|
} else {
|
|
save_state();
|
|
}
|
|
batch_timer_running_ = false;
|
|
if (batch_closed_cb_) batch_closed_cb_(batch_number, counting_date);
|
|
return false;
|
|
}
|
|
|
|
try {
|
|
insert_batch(state["counting_date"], batch_number,
|
|
count, state["start_time"], end_time);
|
|
} catch (const std::exception& e) {
|
|
logger_("Failed to persist batch: " + std::string(e.what()));
|
|
return false;
|
|
}
|
|
|
|
current_state_ = nullptr;
|
|
if (preserve_state_file_ && closed_by == "shutdown") {
|
|
std::ofstream f(state_file_);
|
|
f << state.dump(2);
|
|
} else {
|
|
save_state();
|
|
}
|
|
batch_timer_running_ = false;
|
|
if (batch_closed_cb_) batch_closed_cb_(batch_number, counting_date);
|
|
return true;
|
|
}
|
|
|
|
void BatchStore::insert_batch(const std::string& counting_date, int batch_number,
|
|
int count, const std::string& start_time, const std::string& end_time) {
|
|
const char* sql1 = R"(
|
|
INSERT INTO batches (counting_date, batch_number, camera_name, object_label, count, start_time, end_time)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?)
|
|
)";
|
|
sqlite3_stmt* stmt;
|
|
sqlite3_prepare_v2(db_, sql1, -1, &stmt, nullptr);
|
|
sqlite3_bind_text(stmt, 1, counting_date.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_int(stmt, 2, batch_number);
|
|
sqlite3_bind_text(stmt, 3, camera_name_.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 4, object_label_.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_int(stmt, 5, count);
|
|
sqlite3_bind_text(stmt, 6, start_time.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 7, end_time.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_step(stmt);
|
|
sqlite3_finalize(stmt);
|
|
|
|
const char* sql2 = R"(
|
|
INSERT INTO daily_summaries (counting_date, camera_name, object_label, total_count, total_batches)
|
|
VALUES (?, ?, ?, ?, 1)
|
|
ON CONFLICT(counting_date, camera_name, object_label)
|
|
DO UPDATE SET total_count = total_count + excluded.total_count,
|
|
total_batches = total_batches + excluded.total_batches,
|
|
updated_at = CURRENT_TIMESTAMP
|
|
)";
|
|
sqlite3_prepare_v2(db_, sql2, -1, &stmt, nullptr);
|
|
sqlite3_bind_text(stmt, 1, counting_date.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 2, camera_name_.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 3, object_label_.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_int(stmt, 4, count);
|
|
sqlite3_step(stmt);
|
|
sqlite3_finalize(stmt);
|
|
}
|
|
|
|
int BatchStore::get_closed_total_for_day(const std::string& counting_date) {
|
|
const char* sql = R"(
|
|
SELECT COALESCE(total_count, 0)
|
|
FROM daily_summaries
|
|
WHERE counting_date = ? AND camera_name = ? AND object_label = ?
|
|
)";
|
|
sqlite3_stmt* stmt;
|
|
sqlite3_prepare_v2(db_, sql, -1, &stmt, nullptr);
|
|
sqlite3_bind_text(stmt, 1, counting_date.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 2, camera_name_.c_str(), -1, SQLITE_TRANSIENT);
|
|
sqlite3_bind_text(stmt, 3, object_label_.c_str(), -1, SQLITE_TRANSIENT);
|
|
int result = 0;
|
|
if (sqlite3_step(stmt) == SQLITE_ROW)
|
|
result = sqlite3_column_int(stmt, 0);
|
|
sqlite3_finalize(stmt);
|
|
return result;
|
|
}
|
|
|
|
int BatchStore::display_total() {
|
|
return get_closed_total_for_day(get_counting_date()) + current_batch_count();
|
|
}
|
|
|
|
void BatchStore::set_batch_closed_callback(std::function<void(int, const std::string&)> cb) {
|
|
batch_closed_cb_ = std::move(cb);
|
|
}
|
|
|
|
void BatchStore::reset() {
|
|
std::lock_guard<std::mutex> lock(state_mutex_);
|
|
if (!current_state_.is_null()) {
|
|
end_batch_locked("reset");
|
|
}
|
|
current_state_ = nullptr;
|
|
batch_timer_running_ = false;
|
|
std::error_code ec;
|
|
std::filesystem::remove(state_file_, ec);
|
|
}
|
|
|
|
bool BatchStore::check_cutoff_reset() {
|
|
return cutoff_triggered_.exchange(false);
|
|
}
|
|
|
|
int BatchStore::current_batch_number() const {
|
|
if (current_state_.is_null()) return 0;
|
|
return current_state_.value("batch_number", 0);
|
|
}
|
|
|
|
int BatchStore::current_batch_count() const {
|
|
if (current_state_.is_null()) return 0;
|
|
return current_state_.value("count", 0);
|
|
}
|
|
|
|
void BatchStore::cutoff_watcher_loop() {
|
|
while (!shutdown_flag_) {
|
|
std::this_thread::sleep_for(std::chrono::seconds(60));
|
|
std::lock_guard<std::mutex> lock(state_mutex_);
|
|
if (current_state_.is_null()) continue;
|
|
if (current_state_["counting_date"].get<std::string>() != get_counting_date()) {
|
|
logger_("Daily cutoff reached - finalizing batch");
|
|
end_batch_locked("cutoff");
|
|
cutoff_triggered_ = true;
|
|
}
|
|
}
|
|
}
|
|
|
|
void BatchStore::start_cutoff_watcher() {
|
|
std::thread t(&BatchStore::cutoff_watcher_loop, this);
|
|
t.detach();
|
|
}
|
|
|
|
void BatchStore::ping_activity() {
|
|
std::lock_guard<std::mutex> lock(state_mutex_);
|
|
if (!current_state_.is_null()) {
|
|
reset_batch_timer();
|
|
}
|
|
}
|
|
|
|
void BatchStore::shutdown() {
|
|
shutdown_flag_ = true;
|
|
batch_timer_running_ = false;
|
|
end_batch("shutdown");
|
|
}
|