#include "FeedbackStore.h" #include "sqlite3.h" #include "json/json.h" #include #include #include using namespace std; // 创建路径的父目录(递归,类似 mkdir -p) static void ensureParentDir(const string& filePath) { size_t pos = filePath.find_last_of('/'); if (pos == string::npos) return; // 无目录部分 string dir = filePath.substr(0, pos); if (dir.empty()) return; string cur; size_t start = 0; while (true) { size_t slash = dir.find('/', start); string part = (slash == string::npos) ? dir.substr(start) : dir.substr(start, slash - start); if (!part.empty()) { if (!cur.empty()) cur += '/'; cur += part; mkdir(cur.c_str(), 0755); } if (slash == string::npos) break; start = slash + 1; } } FeedbackStore::FeedbackStore() : m_db(nullptr), m_appendStmt(nullptr) { } FeedbackStore::~FeedbackStore() { close(); } static int64_t nowMs() { using namespace std::chrono; return duration_cast(system_clock::now().time_since_epoch()).count(); } bool FeedbackStore::open(const string& path) { lock_guard lock(m_mutex); if (m_db) return true; ensureParentDir(path); if (sqlite3_open(path.c_str(), &m_db) != SQLITE_OK) { m_db = nullptr; return false; } const char* schema = "CREATE TABLE IF NOT EXISTS fb_log (" " id INTEGER PRIMARY KEY AUTOINCREMENT," " ts INTEGER NOT NULL," " msg TEXT NOT NULL," " payload TEXT NOT NULL" ");" "CREATE INDEX IF NOT EXISTS idx_fb_log_msg_ts ON fb_log(msg, ts);"; char* err = nullptr; if (sqlite3_exec(m_db, schema, nullptr, nullptr, &err) != SQLITE_OK) { sqlite3_free(err); sqlite3_close(m_db); m_db = nullptr; return false; } return true; } void FeedbackStore::close() { lock_guard lock(m_mutex); if (m_appendStmt) { sqlite3_finalize(m_appendStmt); m_appendStmt = nullptr; } if (m_db) { sqlite3_close(m_db); m_db = nullptr; } } bool FeedbackStore::begin() { lock_guard lock(m_mutex); if (!m_db) return false; char* err = nullptr; int rc = sqlite3_exec(m_db, "BEGIN;", nullptr, nullptr, &err); if (err) sqlite3_free(err); return rc == SQLITE_OK; } bool FeedbackStore::append(int64_t ts, const string& msg, const string& payload) { lock_guard lock(m_mutex); if (!m_db) return false; if (!m_appendStmt) { const char* sql = "INSERT INTO fb_log (ts, msg, payload) VALUES (?, ?, ?);"; if (sqlite3_prepare_v2(m_db, sql, -1, &m_appendStmt, nullptr) != SQLITE_OK) { return false; } } sqlite3_bind_int64(m_appendStmt, 1, ts); sqlite3_bind_text(m_appendStmt, 2, msg.c_str(), -1, SQLITE_TRANSIENT); sqlite3_bind_text(m_appendStmt, 3, payload.c_str(), -1, SQLITE_TRANSIENT); int rc = sqlite3_step(m_appendStmt); sqlite3_reset(m_appendStmt); sqlite3_clear_bindings(m_appendStmt); return rc == SQLITE_DONE; } bool FeedbackStore::commit() { lock_guard lock(m_mutex); if (!m_db) return false; char* err = nullptr; int rc = sqlite3_exec(m_db, "COMMIT;", nullptr, nullptr, &err); if (err) sqlite3_free(err); return rc == SQLITE_OK; } void FeedbackStore::prune(int64_t keepSeconds) { lock_guard lock(m_mutex); if (!m_db || keepSeconds <= 0) return; int64_t cutoff = nowMs() - keepSeconds * 1000; const char* sql = "DELETE FROM fb_log WHERE ts < ?;"; sqlite3_stmt* stmt = nullptr; if (sqlite3_prepare_v2(m_db, sql, -1, &stmt, nullptr) == SQLITE_OK) { sqlite3_bind_int64(stmt, 1, cutoff); sqlite3_step(stmt); sqlite3_finalize(stmt); } } string FeedbackStore::querySeries(const string& msg, const string& field, int64_t fromTs, int64_t toTs, int limit) { Json::Value result(Json::arrayValue); lock_guard lock(m_mutex); if (!m_db || msg.empty() || field.empty()) { return "[]"; } if (limit <= 0) limit = 2000; if (toTs <= 0) toTs = nowMs(); const char* sql = "SELECT ts, payload FROM fb_log " "WHERE msg = ? AND ts >= ? AND ts <= ? ORDER BY ts ASC LIMIT ?;"; sqlite3_stmt* stmt = nullptr; if (sqlite3_prepare_v2(m_db, sql, -1, &stmt, nullptr) != SQLITE_OK) { return "[]"; } sqlite3_bind_text(stmt, 1, msg.c_str(), -1, SQLITE_TRANSIENT); sqlite3_bind_int64(stmt, 2, fromTs); sqlite3_bind_int64(stmt, 3, toTs); sqlite3_bind_int(stmt, 4, limit); Json::CharReaderBuilder builder; while (sqlite3_step(stmt) == SQLITE_ROW) { int64_t ts = sqlite3_column_int64(stmt, 0); const char* payload = (const char*)sqlite3_column_text(stmt, 1); if (!payload) continue; Json::Value root; JSONCPP_STRING errs; istringstream stream(payload); if (!Json::parseFromStream(builder, stream, &root, &errs)) continue; if (!root.isMember(field) || !root[field].isNumeric()) continue; Json::Value pt(Json::arrayValue); pt.append((Json::Int64)ts); pt.append(root[field].asDouble()); result.append(pt); } sqlite3_finalize(stmt); Json::StreamWriterBuilder wb; wb.settings_["indentation"] = ""; return Json::writeString(wb, result); } string FeedbackStore::queryRaw(const string& msg, int64_t fromTs, int64_t toTs, int limit) { Json::Value result(Json::objectValue); Json::Value rows(Json::arrayValue); lock_guard lock(m_mutex); if (m_db && !msg.empty()) { if (limit <= 0) limit = 500; if (toTs <= 0) toTs = nowMs(); const char* sql = "SELECT ts, payload FROM fb_log " "WHERE msg = ? AND ts >= ? AND ts <= ? ORDER BY ts ASC LIMIT ?;"; sqlite3_stmt* stmt = nullptr; if (sqlite3_prepare_v2(m_db, sql, -1, &stmt, nullptr) == SQLITE_OK) { sqlite3_bind_text(stmt, 1, msg.c_str(), -1, SQLITE_TRANSIENT); sqlite3_bind_int64(stmt, 2, fromTs); sqlite3_bind_int64(stmt, 3, toTs); sqlite3_bind_int(stmt, 4, limit); while (sqlite3_step(stmt) == SQLITE_ROW) { Json::Value row(Json::objectValue); row["ts"] = (Json::Int64)sqlite3_column_int64(stmt, 0); const char* p = (const char*)sqlite3_column_text(stmt, 1); row["payload"] = p ? p : ""; rows.append(row); } sqlite3_finalize(stmt); } } result["rows"] = rows; Json::StreamWriterBuilder wb; wb.settings_["indentation"] = ""; return Json::writeString(wb, result); }