From 3c0dc1c68815faf72dc38389ca0e5572c0dfdaca Mon Sep 17 00:00:00 2001 From: zjk <1553836110@qq.com> Date: Wed, 19 Aug 2026 16:42:48 +0800 Subject: [PATCH] =?UTF-8?q?pPowerMangerHost=EF=BC=9A=E5=8F=8D=E9=A6=88?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=20SQLite=20=E8=90=BD=E5=BA=93=EF=BC=8C?= =?UTF-8?q?=E6=96=B0=E5=A2=9E=E5=8E=86=E5=8F=B2=E6=9F=A5=E8=AF=A2=20API=20?= =?UTF-8?q?=E4=B8=8E=E6=9B=B2=E7=BA=BF=E7=BB=98=E5=9B=BE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 FeedbackStore(sqlite3 封装),fb_log 表按消息落盘,环形时间保留 - HostSim 每次收到反馈即批量写库,支持 DB_FILE/FB_LOG_KEEP_SECONDS 等配置 - 新增 /api/history、/api/rawlog 接口;曲线页新增历史数据加载 - 数据库目录不存在时自动创建;新增板卡原生编译脚本 scripts/build-board.sh --- .gitignore | 2 + scripts/build-board.sh | 164 ++++++++++++++ src/pPowerMangerHost/CMakeLists.txt | 6 + src/pPowerMangerHost/FeedbackStore.cpp | 236 +++++++++++++++++++++ src/pPowerMangerHost/FeedbackStore.h | 48 +++++ src/pPowerMangerHost/HostSim.cpp | 138 +++++++++++- src/pPowerMangerHost/HostSim.h | 12 ++ src/pPowerMangerHost/WebServer.cpp | 22 +- src/pPowerMangerHost/WebServer.h | 4 + src/pPowerMangerHost/pPowerMangerHost.moos | 6 + src/pPowerMangerHost/web/index.html | 71 ++++++- 11 files changed, 705 insertions(+), 4 deletions(-) create mode 100755 scripts/build-board.sh create mode 100644 src/pPowerMangerHost/FeedbackStore.cpp create mode 100644 src/pPowerMangerHost/FeedbackStore.h diff --git a/.gitignore b/.gitignore index 89ae626..334c5c4 100644 --- a/.gitignore +++ b/.gitignore @@ -245,6 +245,8 @@ _deps bin/pPowerManger.log bin/ccuState.txt *.db +# pPowerMangerHost 反馈落盘数据库 +feedback_data.db # PowerManager compiled binaries (built into bin/ by CMake) bin/pPowerManger diff --git a/scripts/build-board.sh b/scripts/build-board.sh new file mode 100755 index 0000000..436b733 --- /dev/null +++ b/scripts/build-board.sh @@ -0,0 +1,164 @@ +#!/bin/bash +#======================================================================= +# FILE: scripts/build-board.sh +# DESC: 把源码同步到目标板卡 (RK3588 / Ubuntu 22.04),在板卡原生编译 +# (aarch64 直编,比本地 QEMU 模拟快约 10 倍),并把产物部署到 +# /root/work/h100/bin/,可选重启 systemd 服务。 +# +# 用法: +# ./scripts/build-board.sh # 同步+编译+部署(不重启) +# ./scripts/build-board.sh --start # 编译后重启 3 个服务 +# ./scripts/build-board.sh --clean # 板卡上重新 cmake 再编译 +# ./scripts/build-board.sh --jobs 4 # 指定并行度(默认板卡 nproc) +# +# 参数: +# --host 目标主机(默认 192.168.0.223) +# --user 登录用户(默认 root) +# --board-src 板卡源码目录(默认 /root/work/h100/src) +# --board-bin 板卡部署目录(默认 /root/work/h100/bin) +# --jobs 编译并行度(默认板卡 nproc) +# --start 编译部署后重启服务 +# --clean 先重新 cmake(配置变更时用) +# -h, --help 显示帮助 +# +# 前置: 板卡已配免密 SSH,且有编译环境(cmake/g++/MOOS/jsoncpp), +# 源码已能本地编译(本脚本只负责同步+板卡编译)。 +#======================================================================= + +set -uo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +PROJECT_ROOT="$(dirname "${SCRIPT_DIR}")" + +HOST="192.168.0.223" +USER="root" +BOARD_SRC="/root/work/h100/src" +BOARD_BIN="/root/work/h100/bin" +JOBS="" +DO_START=0 +DO_CLEAN=0 + +SERVICES=(moosdb pPowerManger pPowerMangerHost) + +info() { printf '\033[1;36m[build-board]\033[0m %s\n' "$*"; } +ok() { printf '\033[1;32m[build-board]\033[0m %s\n' "$*"; } +warn() { printf '\033[1;33m[build-board]\033[0m %s\n' "$*"; } +err() { printf '\033[1;31m[build-board]\033[0m %s\n' "$*"; } + +usage() { + sed -n '2,29p' "${BASH_SOURCE[0]}" | sed 's/^# \{0,1\}//' + exit 0 +} + +#------------------------------------------------------------------- +# 参数解析 +#------------------------------------------------------------------- +while [ $# -gt 0 ]; do + case "$1" in + --host) [ $# -ge 2 ] || { err "--host 需要参数"; exit 1; }; HOST="$2"; shift 2 ;; + --user) [ $# -ge 2 ] || { err "--user 需要参数"; exit 1; }; USER="$2"; shift 2 ;; + --board-src) [ $# -ge 2 ] || { err "--board-src 需要参数"; exit 1; }; BOARD_SRC="$2"; shift 2 ;; + --board-bin) [ $# -ge 2 ] || { err "--board-bin 需要参数"; exit 1; }; BOARD_BIN="$2"; shift 2 ;; + --jobs) [ $# -ge 2 ] || { err "--jobs 需要参数"; exit 1; }; JOBS="$2"; shift 2 ;; + --start) DO_START=1; shift ;; + --clean) DO_CLEAN=1; shift ;; + -h|--help) usage ;; + *) err "未知参数: $1(见 --help)"; exit 1 ;; + esac +done + +SSH_TARGET="${USER}@${HOST}" +SSH="ssh -o BatchMode=yes -o ConnectTimeout=5" + +#------------------------------------------------------------------- +# 1. 前置检查 +#------------------------------------------------------------------- +if ! ${SSH} "${SSH_TARGET}" 'true' 2>/dev/null; then + err "无法免密 SSH 到 ${SSH_TARGET}" + err "请先配置免密登录: ssh-copy-id ${SSH_TARGET}" + exit 1 +fi +ok "SSH 连接正常: ${SSH_TARGET}" + +#------------------------------------------------------------------- +# 2. 同步源码到板卡(排除构建产物 / git / 数据) +#------------------------------------------------------------------- +info "同步源码到板卡 ${BOARD_SRC} ..." +rsync -az --delete \ + --exclude='build/' \ + --exclude='build-arm64/' \ + --exclude='build-native/' \ + --exclude='build-nosa-test/' \ + --exclude='bin/' \ + --exclude='.git/' \ + --exclude='.cache/' \ + --exclude='reports/' \ + --exclude='*.log' \ + --exclude='*.db' \ + --exclude='CMakeCache.txt' \ + "${PROJECT_ROOT}/" "${SSH_TARGET}:${BOARD_SRC}/" || { err "rsync 源码失败"; exit 1; } +ok "源码已同步" + +#------------------------------------------------------------------- +# 3. 板卡原生编译 +#------------------------------------------------------------------- +if [ "${DO_CLEAN}" = "1" ]; then + info "重新 cmake(--clean)..." + ${SSH} "${SSH_TARGET}" "rm -rf ${BOARD_SRC}/build-native && mkdir -p ${BOARD_SRC}/build-native && \ + cd ${BOARD_SRC}/build-native && cmake -DCMAKE_BUILD_TYPE=Release .." \ + || { err "cmake 失败"; exit 1; } + ok "cmake 完成" +else + ${SSH} "${SSH_TARGET}" "mkdir -p ${BOARD_SRC}/build-native" + # 若首次未配置过则先 cmake + if ! ${SSH} "${SSH_TARGET}" "[ -f ${BOARD_SRC}/build-native/Makefile ]"; then + info "首次编译,先 cmake ..." + ${SSH} "${SSH_TARGET}" "cd ${BOARD_SRC}/build-native && cmake -DCMAKE_BUILD_TYPE=Release .." \ + || { err "cmake 失败"; exit 1; } + ok "cmake 完成" + fi +fi + +if [ -z "${JOBS}" ]; then + JOBS="$(${SSH} "${SSH_TARGET}" nproc)" +fi + +info "板卡编译 (make -j${JOBS}) ..." +if ! ${SSH} "${SSH_TARGET}" "cd ${BOARD_SRC}/build-native && make -j${JOBS}"; then + err "板卡编译失败(见上方输出)" + exit 1 +fi +ok "板卡编译完成" + +#------------------------------------------------------------------- +# 4. 部署产物到板卡运行目录 +#------------------------------------------------------------------- +info "部署产物到 ${BOARD_BIN}/ ..." +${SSH} "${SSH_TARGET}" "mkdir -p ${BOARD_BIN} && \ + cp -f ${BOARD_SRC}/bin/pPowerManger ${BOARD_BIN}/pPowerManger && \ + cp -f ${BOARD_SRC}/bin/pPowerMangerHost ${BOARD_BIN}/pPowerMangerHost && \ + chmod +x ${BOARD_BIN}/pPowerManger ${BOARD_BIN}/pPowerMangerHost" \ + || { err "部署产物失败"; exit 1; } +ok "产物已部署" + +echo "" +ok "产物: ${BOARD_BIN}/pPowerManger (aarch64)" +${SSH} "${SSH_TARGET}" "file ${BOARD_BIN}/pPowerManger | cut -c1-60" + +#------------------------------------------------------------------- +# 5. 可选重启服务 +#------------------------------------------------------------------- +if [ "${DO_START}" = "1" ]; then + info "重启服务 ..." + ${SSH} "${SSH_TARGET}" "systemctl restart ${SERVICES[*]}" \ + || { err "重启服务失败"; exit 1; } + sleep 3 + ok "服务状态: $(${SSH} "${SSH_TARGET}" "systemctl is-active ${SERVICES[*]}")" +else + echo "" + info "未重启服务。需要重启请加 --start 或手动:" + echo " systemctl restart ${SERVICES[*]}" +fi + +echo "" +ok "完成。" diff --git a/src/pPowerMangerHost/CMakeLists.txt b/src/pPowerMangerHost/CMakeLists.txt index 6f0d0fe..d7527ed 100644 --- a/src/pPowerMangerHost/CMakeLists.txt +++ b/src/pPowerMangerHost/CMakeLists.txt @@ -9,12 +9,17 @@ SET(MONGOOSE_SRC ${CMAKE_SOURCE_DIR}/src/pPowerManger/httpserver/mongoose.c) +SET(SQLITE_SRC + ${CMAKE_SOURCE_DIR}/src/pPowerManger/sqlit3/sqlite3.c) + SET(SRC HostSim.cpp HostSim_Info.cpp WebServer.cpp + FeedbackStore.cpp main.cpp ${MONGOOSE_SRC} + ${SQLITE_SRC} ) #-------------------------------------------------------- @@ -105,6 +110,7 @@ ADD_EXECUTABLE(pPowerMangerHost ${SRC}) TARGET_INCLUDE_DIRECTORIES(pPowerMangerHost PRIVATE ${CMAKE_SOURCE_DIR}/src/pPowerManger/httpserver + ${CMAKE_SOURCE_DIR}/src/pPowerManger/sqlit3 ${CMAKE_CURRENT_BINARY_DIR}) TARGET_LINK_LIBRARIES(pPowerMangerHost diff --git a/src/pPowerMangerHost/FeedbackStore.cpp b/src/pPowerMangerHost/FeedbackStore.cpp new file mode 100644 index 0000000..aa46acc --- /dev/null +++ b/src/pPowerMangerHost/FeedbackStore.cpp @@ -0,0 +1,236 @@ +#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); +} diff --git a/src/pPowerMangerHost/FeedbackStore.h b/src/pPowerMangerHost/FeedbackStore.h new file mode 100644 index 0000000..3c11220 --- /dev/null +++ b/src/pPowerMangerHost/FeedbackStore.h @@ -0,0 +1,48 @@ +#ifndef FEEDBACKSTORE_H +#define FEEDBACKSTORE_H + +#include +#include +#include +#include + +struct sqlite3; +struct sqlite3_stmt; + +// 轻量 SQLite 封装:按时间序列存储上位机反馈消息的原始 JSON。 +// 表结构: +// fb_log(id INTEGER PK AUTOINCREMENT, ts INTEGER, msg TEXT, payload TEXT) +// 索引: (msg, ts) +class FeedbackStore +{ +public: + FeedbackStore(); + ~FeedbackStore(); + + bool open(const std::string& path); + void close(); + bool isOpen() const { return m_db != nullptr; } + + // 批写入:begin -> append*N -> commit + bool begin(); + bool append(int64_t ts, const std::string& msg, const std::string& payload); + bool commit(); + + // 按保留时间清理:删除 ts 早于 (当前时间 - keepSeconds) 的记录 + void prune(int64_t keepSeconds); + + // 查询某消息某数值字段的时间序列,返回 JSON 数组字符串 [[ts,value],...] + std::string querySeries(const std::string& msg, const std::string& field, + int64_t fromTs, int64_t toTs, int limit); + + // 查询某消息的原始记录,返回 JSON {"rows":[{"ts":..,"payload":".."},...]} + std::string queryRaw(const std::string& msg, + int64_t fromTs, int64_t toTs, int limit); + +private: + sqlite3* m_db; + sqlite3_stmt* m_appendStmt; + std::mutex m_mutex; +}; + +#endif // FEEDBACKSTORE_H diff --git a/src/pPowerMangerHost/HostSim.cpp b/src/pPowerMangerHost/HostSim.cpp index c79a874..b34a643 100644 --- a/src/pPowerMangerHost/HostSim.cpp +++ b/src/pPowerMangerHost/HostSim.cpp @@ -7,18 +7,30 @@ #include #include +#include +#include #include "MBUtils.h" #include "HostSim.h" #include "WebServer.h" +#include "FeedbackStore.h" #include "json/json.h" using namespace std; +static int64_t hostNowMs() +{ + using namespace std::chrono; + return duration_cast(system_clock::now().time_since_epoch()).count(); +} + //--------------------------------------------------------- // Constructor HostSim::HostSim() - : m_feedbackDirty(false), m_webServer(nullptr), m_webPort(18080) + : m_feedbackDirty(false), m_webServer(nullptr), m_webPort(18080), + m_feedbackStore(nullptr), m_dbFile("feedback_data.db"), + m_fbLogEnabled(true), m_fbLogKeepSeconds(86400), + m_fbLogPruneIntervalSeconds(60), m_lastPruneTs(0) { } @@ -27,6 +39,11 @@ HostSim::HostSim() HostSim::~HostSim() { + if (m_feedbackStore) { + m_feedbackStore->close(); + delete m_feedbackStore; + m_feedbackStore = nullptr; + } if (m_webServer) { m_webServer->stop(); delete m_webServer; @@ -79,6 +96,7 @@ bool HostSim::Iterate() } if (dirty && m_webServer) { m_webServer->broadcast(buildSnapshot()); + persistFeedback(hostNowMs()); lock_guard lock(m_feedbackMutex); m_feedbackDirty = false; } @@ -98,14 +116,34 @@ bool HostSim::OnStartUp() for (p = sParams.begin(); p != sParams.end(); p++) { string line = *p; string param = tolower(biteStringX(line, '=')); - string value = line; + string value = stripBlankEnds(line); if (param == "web_port") { m_webPort = atoi(value.c_str()); + } else if (param == "db_file") { + m_dbFile = value; + } else if (param == "fb_log_enabled") { + m_fbLogEnabled = (tolower(value) == "true" || value == "1"); + } else if (param == "fb_log_keep_seconds") { + m_fbLogKeepSeconds = atoll(value.c_str()); + } else if (param == "fb_log_prune_interval_seconds") { + m_fbLogPruneIntervalSeconds = atoll(value.c_str()); } } } + // 启动反馈落盘存储(SQLite) + if (m_fbLogEnabled) { + m_feedbackStore = new FeedbackStore(); + if (m_feedbackStore->open(m_dbFile)) { + cout << "pPowerMangerHost feedback logging to " << m_dbFile << endl; + } else { + cerr << "pPowerMangerHost failed to open feedback DB: " << m_dbFile << endl; + delete m_feedbackStore; + m_feedbackStore = nullptr; + } + } + // 启动网页服务器 m_webServer = new WebServer(); m_webServer->setOnMessage([this](const string& msg) { @@ -114,6 +152,9 @@ bool HostSim::OnStartUp() m_webServer->setOnOpen([this]() { return this->buildSnapshot(); }); + m_webServer->setApiHandler([this](const string& uri, const string& query) { + return this->handleApi(uri, query); + }); m_webServer->start(m_webPort); RegisterVariables(); @@ -205,6 +246,99 @@ string HostSim::buildSnapshot() return Json::writeString(wb, root); } +//--------------------------------------------------------- +// Procedure: persistFeedback +// 把当前反馈缓存批量写入 SQLite(每次收到即存) + +void HostSim::persistFeedback(int64_t nowMs) +{ + if (!m_feedbackStore || !m_feedbackStore->isOpen()) return; + + vector > entries; + { + lock_guard lock(m_feedbackMutex); + for (map::iterator it = m_feedback.begin(); + it != m_feedback.end(); ++it) { + entries.push_back(make_pair(it->first, it->second)); + } + } + + if (entries.empty()) return; + + m_feedbackStore->begin(); + for (size_t i = 0; i < entries.size(); i++) { + m_feedbackStore->append(nowMs, entries[i].first, entries[i].second); + } + m_feedbackStore->commit(); + + if (nowMs - m_lastPruneTs >= m_fbLogPruneIntervalSeconds * 1000) { + m_feedbackStore->prune(m_fbLogKeepSeconds); + m_lastPruneTs = nowMs; + } +} + +//--------------------------------------------------------- +// Procedure: handleApi +// 处理 /api/history 与 /api/rawlog 查询 + +static string urlDecode(const string& s) +{ + string out; + for (size_t i = 0; i < s.size(); i++) { + if (s[i] == '%' && i + 2 < s.size()) { + int v = 0; + for (int j = 1; j <= 2; j++) { + char c = s[i + j]; + v <<= 4; + if (c >= '0' && c <= '9') v |= (c - '0'); + else if (c >= 'a' && c <= 'f') v |= (c - 'a' + 10); + else if (c >= 'A' && c <= 'F') v |= (c - 'A' + 10); + } + out += (char)v; + i += 2; + } else if (s[i] == '+') { + out += ' '; + } else { + out += s[i]; + } + } + return out; +} + +static string queryGet(const string& query, const string& key) +{ + size_t pos = 0; + while (pos <= query.size()) { + size_t amp = query.find('&', pos); + string kv = query.substr(pos, amp == string::npos ? string::npos : amp - pos); + size_t eq = kv.find('='); + if (eq != string::npos && kv.substr(0, eq) == key) { + return urlDecode(kv.substr(eq + 1)); + } + if (amp == string::npos) break; + pos = amp + 1; + } + return ""; +} + +string HostSim::handleApi(const string& uri, const string& query) +{ + if (!m_feedbackStore || !m_feedbackStore->isOpen()) return ""; + + string msg = queryGet(query, "msg"); + string field = queryGet(query, "field"); + int64_t fromTs = atoll(queryGet(query, "from").c_str()); + int64_t toTs = atoll(queryGet(query, "to").c_str()); + int limit = atoi(queryGet(query, "limit").c_str()); + + if (uri == "/api/history") { + return m_feedbackStore->querySeries(msg, field, fromTs, toTs, limit); + } else if (uri == "/api/rawlog") { + return m_feedbackStore->queryRaw(msg, fromTs, toTs, limit); + } + return ""; +} + //--------------------------------------------------------- // Procedure: RegisterVariables diff --git a/src/pPowerMangerHost/HostSim.h b/src/pPowerMangerHost/HostSim.h index 0838970..5e59976 100644 --- a/src/pPowerMangerHost/HostSim.h +++ b/src/pPowerMangerHost/HostSim.h @@ -17,6 +17,7 @@ #include class WebServer; +class FeedbackStore; // 模拟 pPowerManger 上位机:网页 <-> MOOSDB 桥 class HostSim : public CMOOSApp @@ -37,11 +38,14 @@ class HostSim : public CMOOSApp void pushCommand(const std::string& key, const std::string& jsonValue); // 取当前所有反馈变量的 JSON 快照 std::string buildSnapshot(); + // 处理 /api/* 请求,返回 JSON 响应体(空串表示 404) + std::string handleApi(const std::string& uri, const std::string& query); private: void registerFeedbackVars(); bool drainCommands(); void handleWebMessage(const std::string& msg); + void persistFeedback(int64_t nowMs); private: // State variables // 反馈缓存:key -> 最新 JSON 字符串 @@ -56,6 +60,14 @@ class HostSim : public CMOOSApp // 网页服务器 WebServer* m_webServer; int m_webPort; + + // 反馈落盘存储(SQLite) + FeedbackStore* m_feedbackStore; + std::string m_dbFile; + bool m_fbLogEnabled; + int64_t m_fbLogKeepSeconds; + int64_t m_fbLogPruneIntervalSeconds; + int64_t m_lastPruneTs; }; #endif diff --git a/src/pPowerMangerHost/WebServer.cpp b/src/pPowerMangerHost/WebServer.cpp index c05dcb2..dbb06af 100644 --- a/src/pPowerMangerHost/WebServer.cpp +++ b/src/pPowerMangerHost/WebServer.cpp @@ -15,8 +15,24 @@ void WebServer::handleEvent(struct mg_connection* c, int ev, void* ev_data) { return; } - // 静态资源路由:在嵌入式资源表 webassets 中按路径查找并返回 + // /api/* 接口:交给宿主程序处理,返回 JSON std::string path(hm->uri.buf, hm->uri.len); + if (path.rfind("/api/", 0) == 0) { + std::string body; + if (m_apiHandler) { + std::string query(hm->query.buf, hm->query.len); + body = m_apiHandler(path, query); + } + if (!body.empty()) { + mg_http_reply(c, 200, "Content-Type: application/json\r\n", + "%.*s", (int)body.size(), body.c_str()); + } else { + mg_http_reply(c, 404, "Content-Type: text/plain\r\n", "Not Found\n"); + } + return; + } + + // 静态资源路由:在嵌入式资源表 webassets 中按路径查找并返回 if (path == "/" || path == "/index.html") { path = "index.html"; } else if (!path.empty() && path[0] == '/') { @@ -76,6 +92,10 @@ void WebServer::setOnOpen(std::function handler) { m_onOpen = handler; } +void WebServer::setApiHandler(std::function handler) { + m_apiHandler = handler; +} + bool WebServer::start(int port) { if (m_running) return true; m_port = port; diff --git a/src/pPowerMangerHost/WebServer.h b/src/pPowerMangerHost/WebServer.h index 0951541..e3e764e 100644 --- a/src/pPowerMangerHost/WebServer.h +++ b/src/pPowerMangerHost/WebServer.h @@ -27,6 +27,9 @@ public: // 新 WebSocket 客户端连接时回调,返回需要立即推送给该客户端的初始文本(如快照) void setOnOpen(std::function handler); + // 处理 /api/* HTTP 请求,返回 JSON 响应体(空串表示 404) + void setApiHandler(std::function handler); + private: static bool serverThreadCB(void* pParam); bool serverThreadFunc(); @@ -41,6 +44,7 @@ private: std::vector m_wsConnections; std::function m_onMessage; std::function m_onOpen; + std::function m_apiHandler; }; #endif // WEBSERVER_H diff --git a/src/pPowerMangerHost/pPowerMangerHost.moos b/src/pPowerMangerHost/pPowerMangerHost.moos index 08de558..81b49c3 100644 --- a/src/pPowerMangerHost/pPowerMangerHost.moos +++ b/src/pPowerMangerHost/pPowerMangerHost.moos @@ -5,4 +5,10 @@ ProcessConfig = pPowerMangerHost AppTick = 4 CommsTick = 4 WEB_PORT = 18080 + + // 反馈落盘(SQLite) + DB_FILE = feedback_data.db + FB_LOG_ENABLED = true + FB_LOG_KEEP_SECONDS = 86400 // 历史保留时长(秒),24h + FB_LOG_PRUNE_INTERVAL_SECONDS = 60 // 清理检查周期(秒) } diff --git a/src/pPowerMangerHost/web/index.html b/src/pPowerMangerHost/web/index.html index 257e00f..70bf976 100644 --- a/src/pPowerMangerHost/web/index.html +++ b/src/pPowerMangerHost/web/index.html @@ -367,6 +367,25 @@
+ +
+

历史数据(本地 SQLite)

+
+ + + + + + + +
+
+
@@ -1097,6 +1116,52 @@ function drawChart() { }); } +// ============================================================ +// 11.5 历史数据(本地 SQLite) +// ============================================================ +function populateHistMsgs() { + el("histMsg").innerHTML = FEEDBACKS.map(([k, t]) => + ``).join(""); + populateHistFields(); +} + +function populateHistFields() { + const msg = el("histMsg").value; + const d = feedbackData[msg]; + const sel = el("histField"); + if (!d || typeof d !== "object") { sel.innerHTML = ''; return; } + const fields = Object.keys(d).filter(k => typeof d[k] === "number"); + sel.innerHTML = fields.length + ? fields.map(f => ``).join("") + : ''; +} + +function loadHistorySeries() { + const msg = el("histMsg").value; + const field = el("histField").value; + if (!msg || !field) { log("请选择消息和字段", "err"); return; } + const rangeSec = Number(el("histRange").value); + const to = Date.now(); + const from = to - rangeSec * 1000; + const url = "/api/history?msg=" + encodeURIComponent(msg) + + "&field=" + encodeURIComponent(field) + + "&from=" + from + "&to=" + to + "&limit=2000"; + fetch(url).then(r => r.json()).then(points => { + if (!Array.isArray(points) || points.length === 0) { + el("histEcho").textContent = "该字段在所选范围内无数据"; + return; + } + const key = "历史 " + msg + "." + field; + history[key] = { color: COLORS[nextColor % COLORS.length], points: points }; + nextColor++; + renderSeriesBox(); + drawChart(); + el("histEcho").textContent = "已加载 " + points.length + " 点 → " + key; + }).catch(err => { + el("histEcho").textContent = "加载失败: " + err; + }); +} + // ============================================================ // 12. Tab 切换 // ============================================================ @@ -1109,7 +1174,6 @@ function switchTab(name) { }); if (name === "chart") { refreshFieldOptions(); drawChart(); } } - // ============================================================ // 13. 初始化 // ============================================================ @@ -1168,6 +1232,11 @@ function init() { el("btnAddSeries").onclick = () => addSeries(el("fieldSel").value); el("btnClearSeries").onclick = () => { history = {}; nextColor = 0; renderSeriesBox(); drawChart(); }; + // 历史数据 + populateHistMsgs(); + el("histMsg").onchange = populateHistFields; + el("btnLoadHistory").onclick = loadHistorySeries; + // 日志 el("btnClearLog").onclick = () => { el("log").innerHTML = ""; };