diff --git a/.gitignore b/.gitignore index 80c0ec9..1117f00 100644 --- a/.gitignore +++ b/.gitignore @@ -282,6 +282,8 @@ src/pPowerMangerHost/webassets_gen.h # ============================================================ bin/pCCU bin/pccuTest +bin/pCanBridge +pCanBridge_data.db # ============================================================ # MS Office 临时锁文件 / 编辑器、备份残留 diff --git a/missions/h100.moos b/missions/h100.moos index e541719..553ee0a 100644 --- a/missions/h100.moos +++ b/missions/h100.moos @@ -1,6 +1,7 @@ // MOOS file // 板卡 (RK3588 / Ubuntu 20.04) 专用 mission 配置 -// 四个 systemd 服务(moosdb / pPowerManger / pPowerMangerHost / pCCU)共用本文件。 +// 五个 systemd 服务(moosdb / pPowerManger / pPowerMangerHost / pCCU / pCanBridge) +// 共用本文件。 // 由 scripts/deploy.sh 推送到板卡 /root/work/h100/missions/h100.moos ServerHost = localhost @@ -61,10 +62,25 @@ ProcessConfig = pCCU ins_bat_remote_port = 7000 // 电池工况设定:10工房调试 / 30试验实航 / 50科研模式 bat_work_condition = 10 - dbpath = /root/work/h100/data/pccu_data.db logpath = /root/work/h100/data/pCCU.log web_port = 8080 web_enable = true } + +ProcessConfig = pCanBridge +{ + AppTick = 4 + CommsTick = 4 + + // CAN->MOOSDB 透传:USBCAN-8E-U/CANET(TCP Server 模式) + // 工作端口:CAN0=4001,CAN1=4002,...,CAN7=4008 + can_host = 192.168.0.222 + can_port = 4001 + // 通道名(写入消息 m_sSrcAux 辅助字段) + can_channel_name = CAN0 + + // CAN 帧 SQLite 落库路径(can_frame 表) + dbpath = /root/work/h100/data/pCanBridge_data.db +} diff --git a/scripts/build-board.sh b/scripts/build-board.sh index c1b4b82..0268902 100755 --- a/scripts/build-board.sh +++ b/scripts/build-board.sh @@ -7,7 +7,7 @@ # # 用法: # ./scripts/build-board.sh # 同步+编译+部署(不重启) -# ./scripts/build-board.sh --start # 编译后重启 3 个服务 +# ./scripts/build-board.sh --start # 编译后重启全部服务(含 pCanBridge) # ./scripts/build-board.sh --clean # 板卡上重新 cmake 再编译 # ./scripts/build-board.sh --jobs 4 # 指定并行度(默认板卡 nproc) # @@ -68,7 +68,7 @@ DO_CLEAN=0 CCU_HOST="" CCU_PORT="" -SERVICES=(moosdb pPowerManger pPowerMangerHost pCCU) +SERVICES=(moosdb pPowerManger pPowerMangerHost pCCU pCanBridge) info() { printf '\033[1;36m[build-board]\033[0m %s\n' "$*"; } ok() { printf '\033[1;32m[build-board]\033[0m %s\n' "$*"; } @@ -210,10 +210,34 @@ ${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 && \ cp -f ${BOARD_SRC}/bin/pCCU ${BOARD_BIN}/pCCU && \ - chmod +x ${BOARD_BIN}/pPowerManger ${BOARD_BIN}/pPowerMangerHost ${BOARD_BIN}/pCCU" \ + cp -f ${BOARD_SRC}/bin/pCanBridge ${BOARD_BIN}/pCanBridge && \ + chmod +x ${BOARD_BIN}/pPowerManger ${BOARD_BIN}/pPowerMangerHost ${BOARD_BIN}/pCCU ${BOARD_BIN}/pCanBridge" \ || { err "部署产物失败"; exit 1; } ok "产物已部署" +#------------------------------------------------------------------- +# 4.5 安装/刷新 pCanBridge systemd unit(幂等,其它 4 个 unit 已装) +#------------------------------------------------------------------- +info "安装 pCanBridge.service ..." +${SSH} "${SSH_TARGET}" "cat > /etc/systemd/system/pCanBridge.service" </dev/null || true" +ok "pCanBridge.service 已安装" + echo "" ok "产物: ${BOARD_BIN}/pPowerManger (aarch64)" ${SSH} "${SSH_TARGET}" "file ${BOARD_BIN}/pPowerManger | cut -c1-60" diff --git a/scripts/deploy.sh b/scripts/deploy.sh index 2962b8b..918ee4e 100755 --- a/scripts/deploy.sh +++ b/scripts/deploy.sh @@ -1,10 +1,10 @@ #!/bin/bash #======================================================================= # FILE: scripts/deploy.sh -# DESC: 部署 pPowerManger / pPowerMangerHost / pCCU(aarch64 二进制 + -# mission)到目标板卡 (RK3588 / Ubuntu 22.04),并用四个独立 -# systemd 服务(moosdb / pPowerManger / pPowerMangerHost / pCCU) -# 管理。只安装服务 + daemon-reload,默认不启动、不自启。 +# DESC: 部署 pPowerManger / pPowerMangerHost / pCCU / pCanBridge(aarch64 二进制 + +# mission)到目标板卡 (RK3588 / Ubuntu 22.04),并用五个独立 +# systemd 服务(moosdb / pPowerManger / pPowerMangerHost / pCCU / +# pCanBridge)管理。只安装服务 + daemon-reload,默认不启动、不自启。 # # 用法: # ./scripts/deploy.sh setup-ssh 一次性:生成 ed25519 key + ssh-copy-id @@ -28,7 +28,8 @@ # /root/work/h100/missions/ 板卡 mission h100.moos # /root/work/h100/data/ 数据库存储目录 (power_data.db 等) # /etc/systemd/system/ moosdb.service / pPowerManger.service / -# pPowerMangerHost.service / pCCU.service +# pPowerMangerHost.service / pCCU.service / +# pCanBridge.service #======================================================================= set -uo pipefail @@ -42,7 +43,7 @@ DO_START=0 CCU_HOST="" CCU_PORT="" -SERVICES=(moosdb pPowerManger pPowerMangerHost pCCU) +SERVICES=(moosdb pPowerManger pPowerMangerHost pCCU pCanBridge) info() { printf '\033[1;36m[deploy]\033[0m %s\n' "$*"; } ok() { printf '\033[1;32m[deploy]\033[0m %s\n' "$*"; } @@ -101,7 +102,7 @@ cmd_setup_ssh() { #------------------------------------------------------------------- check_local() { local missing=0 - for b in pPowerManger pPowerMangerHost pCCU; do + for b in pPowerManger pPowerMangerHost pCCU pCanBridge; do if [ ! -f "${PROJECT_ROOT}/bin/arm64/${b}" ]; then err "缺少本地 aarch64 二进制: bin/arm64/${b}(先运行 ./scripts/build-arm64.sh)"; missing=1 elif ! file "${PROJECT_ROOT}/bin/arm64/${b}" | grep -q "aarch64"; then @@ -204,6 +205,7 @@ cmd_deploy() { ${SCP} -q "${PROJECT_ROOT}/bin/arm64/pPowerManger" \ "${PROJECT_ROOT}/bin/arm64/pPowerMangerHost" \ "${PROJECT_ROOT}/bin/arm64/pCCU" \ + "${PROJECT_ROOT}/bin/arm64/pCanBridge" \ "${SSH_TARGET}:${BOARD_DIR}/bin/" || { err "推送二进制失败"; return 1; } info "推送 mission h100.moos ..." @@ -221,7 +223,7 @@ cmd_deploy() { "${SSH_TARGET}:${BOARD_DIR}/missions/" || { err "推送 mission 失败"; [ -n "${tmp_moos}" ] && rm -f "${tmp_moos}"; return 1; } [ -n "${tmp_moos}" ] && rm -f "${tmp_moos}" - ${SSH} "${SSH_TARGET}" "chmod +x ${BOARD_DIR}/bin/pPowerManger ${BOARD_DIR}/bin/pPowerMangerHost ${BOARD_DIR}/bin/pCCU" + ${SSH} "${SSH_TARGET}" "chmod +x ${BOARD_DIR}/bin/pPowerManger ${BOARD_DIR}/bin/pPowerMangerHost ${BOARD_DIR}/bin/pCCU ${BOARD_DIR}/bin/pCanBridge" ok "二进制与 mission 已推送" # 2. 生成并安装 4 个 systemd unit @@ -248,7 +250,7 @@ cmd_deploy() { else echo "" info "部署完成。手动启动(或加 --start):" - echo " systemctl start moosdb pPowerManger pPowerMangerHost pCCU" + echo " systemctl start moosdb pPowerManger pPowerMangerHost pCCU pCanBridge" echo " 浏览器: http://${HOST}:18080 (上位机模拟)" echo " http://${HOST}:8090 (pPowerManger)" echo " http://${HOST}:8080 (pCCU)" @@ -259,9 +261,9 @@ cmd_deploy() { # 命令: start / status / stop #------------------------------------------------------------------- cmd_start() { - info "启动服务 (moosdb -> pPowerManger / pPowerMangerHost / pCCU) ..." + info "启动服务 (moosdb -> pPowerManger / pPowerMangerHost / pCCU / pCanBridge) ..." ${SSH} "${SSH_TARGET}" "systemctl start moosdb && sleep 1 && \ - systemctl start pPowerManger pPowerMangerHost pCCU" || { err "启动失败"; return 1; } + systemctl start pPowerManger pPowerMangerHost pCCU pCanBridge" || { err "启动失败"; return 1; } ok "服务已启动" cmd_status } @@ -271,7 +273,7 @@ cmd_status() { } cmd_stop() { - ${SSH} "${SSH_TARGET}" "systemctl stop pPowerMangerHost pPowerManger pCCU moosdb" && \ + ${SSH} "${SSH_TARGET}" "systemctl stop pCanBridge pPowerMangerHost pPowerManger pCCU moosdb" && \ ok "服务已停止" || err "停止服务失败" } diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index ea82536..ebb2c8e 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -21,6 +21,7 @@ ADD_SUBDIRECTORY(pPowerManger) add_subdirectory(pPMtest) add_subdirectory(pPowerMangerHost) add_subdirectory(pCCU) +add_subdirectory(pCanBridge) ############################################################################## # END of CMakeLists.txt ############################################################################## diff --git a/src/pCanBridge/CMakeLists.txt b/src/pCanBridge/CMakeLists.txt new file mode 100644 index 0000000..a611fae --- /dev/null +++ b/src/pCanBridge/CMakeLists.txt @@ -0,0 +1,43 @@ +#-------------------------------------------------------- +# The CMakeLists.txt for: pCanBridge +# CAN(USBCAN-8E-U/CANET TCP) -> MOOSDB 透传桥 + SQLite 落库 +#-------------------------------------------------------- + +if (${WIN32}) + SET(SYSTEM_LIBS wsock32) +else (${WIN32}) + SET(SYSTEM_LIBS m pthread) +endif (${WIN32}) + +# 复用仓库内 pPowerManger 的 sqlite3 amalgamation(相对路径引用,避免重复维护) +SET(PM_DIR ${CMAKE_CURRENT_SOURCE_DIR}/../pPowerManger) + +SET(SHARED_SRC + ${PM_DIR}/sqlit3/sqlite3.c +) + +SET(SRC + ${SHARED_SRC} + CanEndpoint.cpp + CanDbStore.cpp + CanBridge.cpp + CanBridge_Info.cpp + main.cpp +) + +ADD_EXECUTABLE(pCanBridge ${SRC}) + +TARGET_INCLUDE_DIRECTORIES(pCanBridge PRIVATE + ${PM_DIR} + ${PM_DIR}/sqlit3 +) + +TARGET_LINK_LIBRARIES(pCanBridge + ${MOOS_LIBRARIES} + apputil + mbutil + m + pthread + dl + ${SYSTEM_LIBS} +) diff --git a/src/pCanBridge/CanBridge.cpp b/src/pCanBridge/CanBridge.cpp new file mode 100644 index 0000000..f11fa3c --- /dev/null +++ b/src/pCanBridge/CanBridge.cpp @@ -0,0 +1,169 @@ +#include "CanBridge.h" +#include "CanBridge_Info.h" +#include "MBUtils.h" +#include "MOOS/libMOOS/Comms/MOOSMsg.h" +#include + +using namespace std; + +namespace canbridge { + +//--------------------------------------------------------- +// Constructor / Destructor + +CanBridge::CanBridge() {} + +CanBridge::~CanBridge() { + if (m_endpoint) { + m_endpoint->stop(); + delete m_endpoint; + m_endpoint = nullptr; + } + if (m_db) { + m_db->close(); + delete m_db; + m_db = nullptr; + } +} + +//--------------------------------------------------------- +// OnStartUp:读取配置并启动 CAN 链路 + +bool CanBridge::OnStartUp() { + AppCastingMOOSApp::OnStartUp(); + + STRING_LIST sParams; + m_MissionReader.EnableVerbatimQuoting(false); + if (!m_MissionReader.GetConfiguration(GetAppName(), sParams)) + reportConfigWarning("No config block found for " + GetAppName()); + + STRING_LIST::iterator p; + for (p = sParams.begin(); p != sParams.end(); p++) { + string orig = *p; + string line = *p; + string param = stripBlankEnds(tolower(biteStringX(line, '='))); + string value = stripBlankEnds(line); + + bool handled = true; + if (param == "can_host") m_canHost = value; + else if (param == "can_port") m_canPort = atol(value.c_str()); + else if (param == "can_channel_name") m_channelName = value; + else if (param == "dbpath") m_dbPath = value; + else handled = false; + + if (!handled) + reportUnhandledConfigWarning(orig); + } + + registerVariables(); + + // 数据库(CAN 帧落库) + m_db = new CanDbStore(m_dbPath); + if (!m_db->open()) { + reportRunWarning("CAN db open failed: " + m_dbPath + " (" + m_db->lastError() + ")"); + delete m_db; + m_db = nullptr; + } + + // CAN 链路(TCP Client,收帧回调在接收线程中触发) + m_endpoint = new CanEndpoint(); + m_endpoint->configure(m_canHost, m_canPort); + m_endpoint->setFrameCallback( + [this](uint8_t fi, uint32_t id, const uint8_t* data, int dlc) { + handleFrame(fi, id, data, dlc); + }); + if (!m_endpoint->start()) { + reportRunWarning("CAN endpoint start failed: " + m_canHost); + } + + cout << "pCanBridge started: " << m_canHost << ":" << m_canPort + << " (" << m_channelName << "), db: " << m_dbPath << endl; + return true; +} + +//--------------------------------------------------------- +// OnConnectToServer + +bool CanBridge::OnConnectToServer() { + registerVariables(); + return true; +} + +//--------------------------------------------------------- +// registerVariables:纯发布应用,无需注册订阅变量 + +void CanBridge::registerVariables() { + AppCastingMOOSApp::RegisterVariables(); +} + +//--------------------------------------------------------- +// OnNewMail:当前无订阅指令,预留 + +bool CanBridge::OnNewMail(MOOSMSG_LIST &NewMail) { + AppCastingMOOSApp::OnNewMail(NewMail); + MOOSMSG_LIST::iterator p; + for (p = NewMail.begin(); p != NewMail.end(); p++) { + CMOOSMsg &msg = *p; + (void)msg; + } + return true; +} + +//--------------------------------------------------------- +// Iterate:周期任务 + +bool CanBridge::Iterate() { + AppCastingMOOSApp::Iterate(); + AppCastingMOOSApp::PostReport(); + return true; +} + +//--------------------------------------------------------- +// handleFrame:CAN 帧 -> CMOOSMsg(二进制) -> MOOSDB +// +// 在 CanEndpoint 接收线程中调用;m_Comms::Post 线程安全。 + +void CanBridge::handleFrame(uint8_t frameInfo, uint32_t id, + const uint8_t* data, int dlc) { + // 落库(WAL + 预处理语句,~60帧/秒无压力) + if (m_db) m_db->onFrame(m_channelName, frameInfo, id, data, dlc); + + char key[32]; + std::snprintf(key, sizeof(key), "CAN_0x%08X", id); + + // 二进制构造:m_cDataType = MOOS_BINARY_STRING + CMOOSMsg msg(MOOS_NOTIFY, key, + static_cast(dlc), + const_cast(data)); + msg.SetSourceAux(m_channelName); // 通道 -> m_sSrcAux + msg.SetDoubleAux(static_cast(frameInfo)); // FF/RTR/DLC -> m_dfVal2 + m_Comms.Post(msg); + ++m_notifyCount; +} + +//--------------------------------------------------------- +// buildReport:AppCasting 报告 + +bool CanBridge::buildReport() { + m_msgs << "============================================" << "\n"; + m_msgs << "pCanBridge CAN->MOOSDB 透传" << "\n"; + m_msgs << "============================================" << "\n"; + if (m_endpoint) { + m_msgs << "目标: " << m_endpoint->host() << ":" + << m_endpoint->port() << " (" << m_channelName << ")" << "\n"; + m_msgs << "连接状态: " << (m_endpoint->isConnected() ? "已连接" : "断开") << "\n"; + m_msgs << "收帧数: " << m_endpoint->frameCount() << "\n"; + m_msgs << "发布数: " << m_notifyCount << "\n"; + m_msgs << "错误数: " << m_endpoint->errorCount() << "\n"; + m_msgs << "重连次数: " << m_endpoint->reconnectCount() << "\n"; + m_msgs << "最近收帧: " << m_endpoint->lastFrameTime() << "\n"; + } + if (m_db) { + m_msgs << "数据库: " << m_dbPath + << (m_db->isOpen() ? " (已打开)" : " (未打开)") << "\n"; + m_msgs << "落库记录数: " << m_db->count() << "\n"; + } + return true; +} + +} // namespace canbridge diff --git a/src/pCanBridge/CanBridge.h b/src/pCanBridge/CanBridge.h new file mode 100644 index 0000000..6ba2c30 --- /dev/null +++ b/src/pCanBridge/CanBridge.h @@ -0,0 +1,62 @@ +#ifndef PCANBRIDGE_CAN_BRIDGE_H +#define PCANBRIDGE_CAN_BRIDGE_H + +#define UNIX +#include "MOOS/libMOOS/Thirdparty/AppCasting/AppCastingMOOSApp.h" +#include "CanEndpoint.h" +#include "CanDbStore.h" +#include + +namespace canbridge { + +//============================================================================ +// CanBridge:CAN -> MOOSDB 透传桥(MOOS 应用外壳)。 +// +// 数据流: +// USBCAN-8E-U (TCP Server, 默认 192.168.0.222:4001 = CAN0) +// -> CanEndpoint(TCP Client,13 字节切帧) +// -> handleFrame() 解析 +// -> Notify 到 MOOSDB +// +// MOOS 消息映射: +// m_sKey = "CAN_0x%08X" CAN ID(完整 4 字节,标准/扩展帧同格式) +// m_sVal = 二进制 data(dlc 字节,MOOS_BINARY_STRING) +// m_sSrcAux = 通道名(默认 "CAN0",配置项 can_channel_name) +// m_dfVal2 = 原始帧信息字节 byte0(bit7 FF / bit6 RTR / bit3~0 DLC) +// m_nID 不使用(MOOS 内部消息序号) +// +// 每帧同时落库 SQLite(can_frame 表,配置项 dbpath)。 +//============================================================================ + +class CanBridge : public AppCastingMOOSApp { +public: + CanBridge(); + ~CanBridge(); + +protected: + bool OnNewMail(MOOSMSG_LIST &NewMail); + bool Iterate(); + bool OnConnectToServer(); + bool OnStartUp(); + bool buildReport(); + void registerVariables(); + +private: + // 收帧回调(CanEndpoint 接收线程调用) + void handleFrame(uint8_t frameInfo, uint32_t id, const uint8_t* data, int dlc); + + // 配置 + std::string m_canHost = "192.168.0.222"; + long m_canPort = 4001; + std::string m_channelName = "CAN0"; // 写入 m_sSrcAux + std::string m_dbPath = "pCanBridge_data.db"; // SQLite 路径(moos 配置项 dbpath) + + // 组件 + CanEndpoint* m_endpoint = nullptr; + CanDbStore* m_db = nullptr; + std::atomic m_notifyCount{0}; +}; + +} // namespace canbridge + +#endif // PCANBRIDGE_CAN_BRIDGE_H diff --git a/src/pCanBridge/CanBridge_Info.cpp b/src/pCanBridge/CanBridge_Info.cpp new file mode 100644 index 0000000..30a1af9 --- /dev/null +++ b/src/pCanBridge/CanBridge_Info.cpp @@ -0,0 +1,106 @@ +/****************************************************************/ +/* NAME: CanBridge_Info */ +/* FILE: CanBridge_Info.cpp */ +/****************************************************************/ + +#include +#include +#include "CanBridge_Info.h" +#include "ColorParse.h" +#include "ReleaseInfo.h" + +using namespace std; + +void showSynopsis() { + blk("SYNOPSIS: "); + blk("------------------------------------ "); + blk(" The pCanBridge application bridges CAN bus frames from a "); + blk(" USBCAN-8E-U / CANET Ethernet-CAN converter (TCP Server mode) "); + blk(" into the MOOSDB. Each received CAN frame is parsed from the "); + blk(" 13-byte CANET wire format and transparently published as a "); + blk(" binary MOOS message. "); + blk(" "); +} + +void showHelpAndExit() { + blk(" "); + blu("=============================================================== "); + blu("Usage: pCanBridge file.moos [OPTIONS] "); + blu("=============================================================== "); + blk(" "); + showSynopsis(); + blk(" "); + blk("Options: "); + mag(" --alias","= "); + blk(" Launch pCanBridge with the given process name "); + blk(" rather than pCanBridge. "); + mag(" --example, -e "); + blk(" Display example MOOS configuration block. "); + mag(" --help, -h "); + blk(" Display this help message. "); + mag(" --interface, -i "); + blk(" Display MOOS publications and subscriptions. "); + mag(" --version,-v "); + blk(" Display the release version of pCanBridge. "); + blk(" "); + blk("Note: If argv[2] does not otherwise match a known option, "); + blk(" then it will be interpreted as a run alias. This is "); + blk(" to support pAntler launching conventions. "); + blk(" "); + exit(0); +} + +void showExampleConfigAndExit() { + blk(" "); + blu("=============================================================== "); + blu("pCanBridge Example MOOS Configuration "); + blu("=============================================================== "); + blk(" "); + blk("ProcessConfig = pCanBridge "); + blk("{ "); + blk(" AppTick = 4 "); + blk(" CommsTick = 4 "); + blk(" "); + blk(" // USBCAN-8E-U/CANET 目标地址(TCP Server 模式) "); + blk(" can_host = 192.168.0.222 // 转换器 IP "); + blk(" can_port = 4001 // 工作端口: CAN0=4001 CAN1=4002.. "); + blk(" "); + blk(" // 通道名(写入消息的 m_sSrcAux 辅助字段) "); + blk(" can_channel_name = CAN0 "); + blk(" "); + blk(" // CAN 帧 SQLite 落库路径(can_frame 表) "); + blk(" dbpath = pCanBridge_data.db "); + blk("} "); + blk(" "); + exit(0); +} + +void showInterfaceAndExit() { + blk(" "); + blu("=============================================================== "); + blu("pCanBridge INTERFACE "); + blu("=============================================================== "); + blk(" "); + showSynopsis(); + blk(" "); + blk("SUBSCRIPTIONS: "); + blk("------------------------------------ "); + blk(" (none: publish-only application) "); + blk(" "); + blk("PUBLICATIONS: "); + blk("------------------------------------ "); + blk(" CAN_0x%08X = CAN 帧(每帧一条二进制消息) "); + blk(" 变量名: CAN ID 完整 4 字节十六进制, 如 CAN_0x10010001 "); + blk(" m_sVal: 二进制 data(dlc 字节, MOOS_BINARY_STRING) "); + blk(" m_sSrcAux: 通道名(默认 CAN0) "); + blk(" m_dfVal2: 原始帧信息字节 byte0 "); + blk(" bit7=FF(1扩展/0标准) bit6=RTR(1远程/0数据) "); + blk(" bit3~0=DLC "); + blk(" "); + exit(0); +} + +void showReleaseInfoAndExit() { + showReleaseInfo("pCanBridge", "gpl"); + exit(0); +} diff --git a/src/pCanBridge/CanBridge_Info.h b/src/pCanBridge/CanBridge_Info.h new file mode 100644 index 0000000..04f1bc0 --- /dev/null +++ b/src/pCanBridge/CanBridge_Info.h @@ -0,0 +1,15 @@ +/****************************************************************/ +/* NAME: CanBridge_Info */ +/* FILE: CanBridge_Info.h */ +/****************************************************************/ + +#ifndef PCANBRIDGE_INFO_HEADER +#define PCANBRIDGE_INFO_HEADER + +void showSynopsis(); +void showHelpAndExit(); +void showExampleConfigAndExit(); +void showInterfaceAndExit(); +void showReleaseInfoAndExit(); + +#endif diff --git a/src/pCanBridge/CanDbStore.cpp b/src/pCanBridge/CanDbStore.cpp new file mode 100644 index 0000000..d8aa401 --- /dev/null +++ b/src/pCanBridge/CanDbStore.cpp @@ -0,0 +1,158 @@ +#include "CanDbStore.h" +#include "sqlite3.h" +#include "MOOS/libMOOS/Utils/MOOSUtilityFunctions.h" +#include + +namespace canbridge { + +namespace { +void toHex(const uint8_t* data, int len, std::string& out) { + static const char* hex = "0123456789ABCDEF"; + out.clear(); + out.reserve(static_cast(len) * 2); + for (int i = 0; i < len; ++i) { + out += hex[(data[i] >> 4) & 0xF]; + out += hex[data[i] & 0xF]; + } +} +} // namespace + +CanDbStore::CanDbStore(const std::string& dbPath) : m_dbPath(dbPath) {} + +CanDbStore::~CanDbStore() { + close(); +} + +bool CanDbStore::exec(const char* sql) { + char* err = nullptr; + int rc = sqlite3_exec(m_db, sql, nullptr, nullptr, &err); + if (rc != SQLITE_OK) { + m_lastError = err ? err : "sqlite error"; + sqlite3_free(err); + return false; + } + return true; +} + +bool CanDbStore::open() { + if (m_db) return true; + if (sqlite3_open(m_dbPath.c_str(), &m_db) != SQLITE_OK) { + m_lastError = m_db ? sqlite3_errmsg(m_db) : "cannot open database"; + sqlite3_close(m_db); + m_db = nullptr; + return false; + } + sqlite3_busy_timeout(m_db, 2000); + exec("PRAGMA journal_mode=WAL;"); + exec("PRAGMA synchronous=NORMAL;"); + + // can_frame:收到的每个 CAN 帧一条记录 + if (!exec( + "CREATE TABLE IF NOT EXISTS can_frame (" + " id INTEGER PRIMARY KEY AUTOINCREMENT," + " time REAL NOT NULL," // unix 时间(含小数秒) + " channel TEXT NOT NULL," // 通道名(CAN0~CAN7) + " frame_info INTEGER NOT NULL," // 原始帧信息字节 byte0 + " can_id INTEGER NOT NULL," // CAN ID + " dlc INTEGER NOT NULL," // 数据长度 0~8 + " hex TEXT NOT NULL" // 数据十六进制 + ");")) return false; + exec("CREATE INDEX IF NOT EXISTS idx_can_frame_time ON can_frame(time);"); + exec("CREATE INDEX IF NOT EXISTS idx_can_frame_id ON can_frame(can_id);"); + + return prepareInsert(); +} + +bool CanDbStore::prepareInsert() { + const char* sql = + "INSERT INTO can_frame (time, channel, frame_info, can_id, dlc, hex) " + "VALUES (?1, ?2, ?3, ?4, ?5, ?6);"; + if (sqlite3_prepare_v2(m_db, sql, -1, &m_stmtInsert, nullptr) != SQLITE_OK) { + m_lastError = sqlite3_errmsg(m_db); + return false; + } + return true; +} + +void CanDbStore::close() { + std::lock_guard lock(m_mutex); + if (m_stmtInsert) { + sqlite3_finalize(m_stmtInsert); + m_stmtInsert = nullptr; + } + if (m_db) { + sqlite3_close(m_db); + m_db = nullptr; + } +} + +void CanDbStore::onFrame(const std::string& channel, uint8_t frameInfo, + uint32_t canId, const uint8_t* data, int dlc) { + std::lock_guard lock(m_mutex); + if (!m_db || !m_stmtInsert) return; + + std::string hex; + if (data && dlc > 0) toHex(data, dlc, hex); + + sqlite3_reset(m_stmtInsert); + sqlite3_clear_bindings(m_stmtInsert); + sqlite3_bind_double(m_stmtInsert, 1, MOOSTime(false)); + sqlite3_bind_text(m_stmtInsert, 2, channel.c_str(), -1, SQLITE_TRANSIENT); + sqlite3_bind_int(m_stmtInsert, 3, frameInfo); + sqlite3_bind_int64(m_stmtInsert, 4, static_cast(canId)); + sqlite3_bind_int(m_stmtInsert, 5, dlc); + sqlite3_bind_text(m_stmtInsert, 6, hex.c_str(), -1, SQLITE_TRANSIENT); + + if (sqlite3_step(m_stmtInsert) != SQLITE_DONE) { + m_lastError = sqlite3_errmsg(m_db); + } +} + +std::vector CanDbStore::queryRecent(uint32_t canId, int limit) { + std::lock_guard lock(m_mutex); + std::vector rows; + if (!m_db) return rows; + if (limit <= 0) limit = 100; + + std::string sql = "SELECT id, time, channel, frame_info, can_id, dlc, hex " + "FROM can_frame"; + if (canId != 0) { + sql += " WHERE can_id=" + std::to_string(canId); + } + sql += " ORDER BY id DESC LIMIT " + std::to_string(limit) + ";"; + + sqlite3_stmt* stmt = nullptr; + if (sqlite3_prepare_v2(m_db, sql.c_str(), -1, &stmt, nullptr) != SQLITE_OK) { + m_lastError = sqlite3_errmsg(m_db); + return rows; + } + while (sqlite3_step(stmt) == SQLITE_ROW) { + CanFrameRow r; + r.id = sqlite3_column_int64(stmt, 0); + r.time = sqlite3_column_double(stmt, 1); + r.channel = sqlite3_column_text(stmt, 2) + ? reinterpret_cast(sqlite3_column_text(stmt, 2)) : ""; + r.frameInfo = sqlite3_column_int(stmt, 3); + r.canId = static_cast(sqlite3_column_int64(stmt, 4)); + r.dlc = sqlite3_column_int(stmt, 5); + r.hex = sqlite3_column_text(stmt, 6) + ? reinterpret_cast(sqlite3_column_text(stmt, 6)) : ""; + rows.push_back(std::move(r)); + } + sqlite3_finalize(stmt); + return rows; +} + +long long CanDbStore::count() const { + std::lock_guard lock(m_mutex); + if (!m_db) return 0; + sqlite3_stmt* stmt = nullptr; + if (sqlite3_prepare_v2(m_db, "SELECT COUNT(*) FROM can_frame;", -1, &stmt, nullptr) != SQLITE_OK) + return 0; + long long n = 0; + if (sqlite3_step(stmt) == SQLITE_ROW) n = sqlite3_column_int64(stmt, 0); + sqlite3_finalize(stmt); + return n; +} + +} // namespace canbridge diff --git a/src/pCanBridge/CanDbStore.h b/src/pCanBridge/CanDbStore.h new file mode 100644 index 0000000..7d53329 --- /dev/null +++ b/src/pCanBridge/CanDbStore.h @@ -0,0 +1,67 @@ +#ifndef PCANBRIDGE_CAN_DB_STORE_H +#define PCANBRIDGE_CAN_DB_STORE_H + +#include +#include +#include +#include + +struct sqlite3; +struct sqlite3_stmt; + +namespace canbridge { + +//============================================================================ +// CanDbStore:CAN 帧 SQLite 存储。 +// +// 设计参考 src/pCCU/store/DbStore(原始帧日志模式): +// - 收到的每个 CAN 帧逐条落库(can_frame 表), +// 存原始帧信息/ID/DLC/十六进制数据,语义解析交给下游。 +// - 预处理语句 + WAL + synchronous=NORMAL,满足 ~60 帧/秒持续写入。 +// - 内部互斥锁:接收线程写入,Iterate 线程查询/统计。 +// 复用仓库内 src/pPowerManger/sqlit3/sqlite3.c 与 sqlite3.h。 +//============================================================================ + +struct CanFrameRow { + long long id; + double time; // unix 时间(含小数秒) + std::string channel; // 通道名(CAN0~CAN7) + int frameInfo; // 原始帧信息字节 byte0(bit7 FF / bit6 RTR / bit3~0 DLC) + uint32_t canId; // CAN ID + int dlc; + std::string hex; // 数据十六进制 +}; + +class CanDbStore { +public: + explicit CanDbStore(const std::string& dbPath); + ~CanDbStore(); + + bool open(); + void close(); + bool isOpen() const { return m_db != nullptr; } + std::string lastError() const { return m_lastError; } + + // 收帧落库(CanEndpoint 接收线程调用,内部加锁) + void onFrame(const std::string& channel, uint8_t frameInfo, + uint32_t canId, const uint8_t* data, int dlc); + + // 查询最近 N 条(canId=0 表示不过滤;limit<=0 表示默认 100) + std::vector queryRecent(uint32_t canId, int limit); + + long long count() const; + +private: + bool exec(const char* sql); + bool prepareInsert(); + + std::string m_dbPath; + sqlite3* m_db = nullptr; + sqlite3_stmt* m_stmtInsert = nullptr; + mutable std::mutex m_mutex; + std::string m_lastError; +}; + +} // namespace canbridge + +#endif // PCANBRIDGE_CAN_DB_STORE_H diff --git a/src/pCanBridge/CanEndpoint.cpp b/src/pCanBridge/CanEndpoint.cpp new file mode 100644 index 0000000..8792fa8 --- /dev/null +++ b/src/pCanBridge/CanEndpoint.cpp @@ -0,0 +1,200 @@ +#include "CanEndpoint.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include "MOOS/libMOOS/Utils/MOOSUtilityFunctions.h" + +namespace canbridge { + +namespace { +// CANET 系列固定帧长 +const size_t kFrameLen = 13; +// 非阻塞 connect 超时(ms) +const int kConnectTOms = 2000; +// recv 读超时(s):保证断线能及时检出、stop() 能及时退出 +const int kRecvTimeoutS = 1; +// 连接失败后的重试间隔(ms) +const int kRetryInterval = 2000; +} // namespace + +CanEndpoint::CanEndpoint() { + m_buf.reserve(4096); +} + +CanEndpoint::~CanEndpoint() { + stop(); +} + +void CanEndpoint::configure(const std::string& host, long port) { + m_host = host; + m_port = port; +} + +void CanEndpoint::setFrameCallback(CanFrameCallback cb) { + m_onFrame = std::move(cb); +} + +bool CanEndpoint::start() { + if (m_running) return true; + if (m_host.empty() || m_port <= 0) { + std::cerr << "[CanEndpoint] invalid target: " << m_host << ":" << m_port << std::endl; + return false; + } + m_running = true; + m_thread = std::thread([this]() { threadFunc(); }); + std::cout << "[CanEndpoint] started, target " << m_host << ":" << m_port << std::endl; + return true; +} + +void CanEndpoint::stop() { + if (!m_running) return; + m_running = false; + if (m_thread.joinable()) m_thread.join(); + closeSocket(); +} + +//---------------------------------------------------------------------- +// 连接一次目标(非阻塞 connect + poll 超时) +//---------------------------------------------------------------------- +bool CanEndpoint::connectOnce() { + m_fd = ::socket(AF_INET, SOCK_STREAM, 0); + if (m_fd < 0) return false; + + // 非阻塞 connect,避免目标不可达时线程长时间卡死 + int flags = ::fcntl(m_fd, F_GETFL, 0); + ::fcntl(m_fd, F_SETFL, flags | O_NONBLOCK); + + struct sockaddr_in addr; + std::memset(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; + addr.sin_port = htons(static_cast(m_port)); + if (::inet_pton(AF_INET, m_host.c_str(), &addr.sin_addr) != 1) { + std::cerr << "[CanEndpoint] bad host: " << m_host << std::endl; + closeSocket(); + return false; + } + + int rc = ::connect(m_fd, reinterpret_cast(&addr), sizeof(addr)); + if (rc < 0 && errno != EINPROGRESS) { + closeSocket(); + return false; + } + if (rc < 0) { + struct pollfd pfd; + pfd.fd = m_fd; + pfd.events = POLLOUT; + int pr = ::poll(&pfd, 1, kConnectTOms); + if (pr <= 0) { + closeSocket(); + return false; + } + int err = 0; + socklen_t elen = sizeof(err); + if (::getsockopt(m_fd, SOL_SOCKET, SO_ERROR, &err, &elen) < 0 || err != 0) { + closeSocket(); + return false; + } + } + + // 恢复阻塞模式 + 读超时 + 禁用 Nagle + ::fcntl(m_fd, F_SETFL, flags); + struct timeval tv; + tv.tv_sec = kRecvTimeoutS; + tv.tv_usec = 0; + ::setsockopt(m_fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv)); + int one = 1; + ::setsockopt(m_fd, IPPROTO_TCP, TCP_NODELAY, &one, sizeof(one)); + + m_connected = true; + ++m_reconnectCount; + std::cout << "[CanEndpoint] connected to " << m_host << ":" << m_port << std::endl; + return true; +} + +void CanEndpoint::closeSocket() { + m_connected = false; + if (m_fd >= 0) { + ::close(m_fd); + m_fd = -1; + } +} + +//---------------------------------------------------------------------- +// 切帧:缓冲 >= 13 字节即解析并回调 +//---------------------------------------------------------------------- +void CanEndpoint::processBuffer() { + while (m_buf.size() >= kFrameLen) { + const uint8_t* f = m_buf.data(); + + uint8_t frameInfo = f[0]; + uint8_t dlc = frameInfo & 0x0F; + uint32_t id = (static_cast(f[1]) << 24) | + (static_cast(f[2]) << 16) | + (static_cast(f[3]) << 8) | + static_cast(f[4]); + + // 无效 DLC(>8)按坏帧丢弃;同时覆盖以太网心跳包(AA 00 ... 55) + if (dlc > 8) { + ++m_errorCount; + m_buf.erase(m_buf.begin(), m_buf.begin() + kFrameLen); + continue; + } + + if (m_onFrame) { + m_onFrame(frameInfo, id, f + 5, static_cast(dlc)); + } + ++m_frameCount; + m_lastFrameTime = MOOSTime(false); + + m_buf.erase(m_buf.begin(), m_buf.begin() + kFrameLen); + } +} + +//---------------------------------------------------------------------- +// 接收线程:连接 -> 读流 -> 切帧 -> 断线重连 +//---------------------------------------------------------------------- +void CanEndpoint::threadFunc() { + uint8_t tmp[4096]; + while (m_running) { + if (!m_connected) { + if (!connectOnce()) { + ++m_errorCount; + // 分片休眠,保证 stop() 能及时退出 + for (int i = 0; i < kRetryInterval / 100 && m_running; ++i) { + struct timespec ts = {0, 100 * 1000 * 1000}; // 100ms + nanosleep(&ts, nullptr); + } + continue; + } + m_buf.clear(); + } + + ssize_t n = ::recv(m_fd, tmp, sizeof(tmp), 0); + if (n > 0) { + m_buf.insert(m_buf.end(), tmp, tmp + n); + processBuffer(); + } else if (n == 0) { + // 对端关闭 + std::cerr << "[CanEndpoint] connection closed by peer" << std::endl; + closeSocket(); + } else { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR) { + // 读超时/中断:回环检查 m_running + continue; + } + std::cerr << "[CanEndpoint] recv error: " << std::strerror(errno) << std::endl; + closeSocket(); + } + } +} + +} // namespace canbridge diff --git a/src/pCanBridge/CanEndpoint.h b/src/pCanBridge/CanEndpoint.h new file mode 100644 index 0000000..f8a3c2c --- /dev/null +++ b/src/pCanBridge/CanEndpoint.h @@ -0,0 +1,87 @@ +#ifndef PCANBRIDGE_CAN_ENDPOINT_H +#define PCANBRIDGE_CAN_ENDPOINT_H + +#include +#include +#include +#include +#include +#include + +namespace canbridge { + +//============================================================================ +// CanEndpoint:USBCAN-8E-U(CANET 系列)TCP 工作端口的客户端封装。 +// +// 协议(CANET-8E-U 用户手册 8.1 节):TCP 为字节流,每个 CAN 帧固定 13 字节: +// byte0 帧信息:bit7=FF(1扩展/0标准) bit6=RTR(1远程/0数据) bit3~0=DLC(0~8) +// byte1~4 帧 ID(大端,标准帧 11 位有效 / 扩展帧 29 位有效) +// byte5~12 数据 8 字节(有效长度由 DLC 决定) +// +// 职责: +// - 以 TCP Client 连接目标 host:port(非阻塞 connect + 超时) +// - 接收线程阻塞读流 -> 缓冲 -> 按 13 字节切帧 -> 回调 +// - 断线(recv 返回 0 / 错误)自动重连 +// +// 说明:采用 POSIX socket 而非 XPCTcpSocket,便于精确控制连接超时与 +// 读超时(SO_RCVTIMEO),并避免 XPC 系列的异常式错误处理。 +//============================================================================ + +// 收帧回调:frameInfo 为原始帧信息字节;data 有效长度为 dlc 字节 +using CanFrameCallback = std::function; + +class CanEndpoint { +public: + CanEndpoint(); + ~CanEndpoint(); + + // 配置目标地址(须在 start 前调用) + void configure(const std::string& host, long port); + + // 收帧回调(须在 start 前调用) + void setFrameCallback(CanFrameCallback cb); + + // 启动接收线程(内部自动连接/重连);返回是否成功启动线程 + bool start(); + void stop(); + + bool isConnected() const { return m_connected; } + + // 统计 + unsigned long frameCount() const { return m_frameCount; } + unsigned long errorCount() const { return m_errorCount; } + unsigned long reconnectCount() const { return m_reconnectCount; } + double lastFrameTime() const { return m_lastFrameTime; } + + const std::string& host() const { return m_host; } + long port() const { return m_port; } + +private: + void threadFunc(); + bool connectOnce(); + void closeSocket(); + void processBuffer(); + + std::string m_host; + long m_port = 4001; + int m_fd = -1; + + std::atomic m_running{false}; + std::atomic m_connected{false}; + std::thread m_thread; + + CanFrameCallback m_onFrame; + + // TCP 流式接收缓冲(跨 recv 保持半包/粘包状态) + std::vector m_buf; + + std::atomic m_frameCount{0}; + std::atomic m_errorCount{0}; + std::atomic m_reconnectCount{0}; + std::atomic m_lastFrameTime{0.0}; +}; + +} // namespace canbridge + +#endif // PCANBRIDGE_CAN_ENDPOINT_H diff --git a/src/pCanBridge/main.cpp b/src/pCanBridge/main.cpp new file mode 100644 index 0000000..e0e9e68 --- /dev/null +++ b/src/pCanBridge/main.cpp @@ -0,0 +1,49 @@ +/************************************************************/ +/* NAME: pCanBridge */ +/* FILE: main.cpp */ +/************************************************************/ + +#include +#include "MBUtils.h" +#include "ColorParse.h" +#include "CanBridge.h" +#include "CanBridge_Info.h" + +using namespace std; + +int main(int argc, char *argv[]) +{ + string mission_file; + string run_command = argv[0]; + // 默认以程序文件名(去掉路径)作为进程名,保证与 .moos 中 ProcessConfig 匹配 + { + size_t slash = run_command.find_last_of('/'); + if (slash != string::npos) run_command = run_command.substr(slash + 1); + } + + for (int i = 1; i < argc; i++) { + string argi = argv[i]; + if ((argi == "-v") || (argi == "--version") || (argi == "-version")) + showReleaseInfoAndExit(); + else if ((argi == "-e") || (argi == "--example") || (argi == "-example")) + showExampleConfigAndExit(); + else if ((argi == "-h") || (argi == "--help") || (argi == "-help")) + showHelpAndExit(); + else if ((argi == "-i") || (argi == "--interface")) + showInterfaceAndExit(); + else if (strEnds(argi, ".moos") || strEnds(argi, ".moos++")) + mission_file = argv[i]; + else if (strBegins(argi, "--alias=")) + run_command = argi.substr(8); + else if (i == 2) + run_command = argi; + } + + if (mission_file == "") + showHelpAndExit(); + + canbridge::CanBridge bridge; + bridge.Run(run_command.c_str(), mission_file.c_str()); + + return 0; +} diff --git a/src/pCanBridge/pCanBridge.moos b/src/pCanBridge/pCanBridge.moos new file mode 100644 index 0000000..4e16671 --- /dev/null +++ b/src/pCanBridge/pCanBridge.moos @@ -0,0 +1,30 @@ +//============================================================================ +// pCanBridge 配置示例 +// +// 拓扑: +// USBCAN-8E-U(CANET 系列,TCP Server 模式) +// 192.168.0.222:4001 = CAN0 工作端口(CAN1=4002 ... CAN7=4008) +// pCanBridge 以 TCP Client 连接目标,收帧后透传 MOOSDB: +// 变量名 CAN_0x%08X = CAN ID;m_sVal = 二进制 data; +// m_sSrcAux = 通道名;m_dfVal2 = 原始帧信息字节(FF/RTR/DLC) +//============================================================================ + +ProcessConfig = pCanBridge +{ + AppTick = 4 + CommsTick = 4 + + //======== 目标 USBCAN-8E-U 地址(TCP Server 模式) ======== + // 转换器 IP + can_host = 192.168.0.222 + // 工作端口:CAN0=4001,CAN1=4002,...,CAN7=4008 + can_port = 4001 + + //======== 通道标识 ======== + // 写入消息 m_sSrcAux 辅助字段,用于下游区分通道 + can_channel_name = CAN0 + + //======== 存储 ======== + // CAN 帧 SQLite 落库路径(can_frame 表) + dbpath = pCanBridge_data.db +}