新增 pCanBridge:USBCAN-8E-U CAN->MOOSDB 透传桥 + SQLite 落库

新应用 src/pCanBridge(独立 MOOS App,单通道):
- CanEndpoint:TCP Client 连接 CANET/USBCAN-8E-U 工作端口
  (目标 IP/端口可配置,默认 192.168.0.222:4001=CAN0),
  接收线程流式缓冲按 13 字节切帧,过滤无效 DLC(>8),
  断线自动重连(非阻塞 connect+超时 / SO_RCVTIMEO 读超时)。
- 帧协议(CANET-8E-U 手册 8.1):byte0=帧信息(bit7 FF/bit6 RTR/
  bit3~0 DLC),byte1~4=ID 大端,byte5~12=数据。
- MOOS 映射:变量名 CAN_0x%08X=CAN ID(不缩减),m_sVal=二进制
  data(逐帧透传不限流),m_sSrcAux=通道名(can_channel_name),
  m_dfVal2=原始帧信息字节(用于区分标准/扩展、数据/远程帧)。
- CanDbStore:每帧落库 SQLite(can_frame 表:time/channel/
  frame_info/can_id/dlc/hex,含 time、can_id 索引),参考 pCCU
  DbStore:预处理语句 + WAL + synchronous=NORMAL,路径由 dbpath
  配置;复用 pPowerManger 的 sqlite3 amalgamation。
- AppCast 输出连接状态/收帧/发布/落库统计。

配套:
- missions/h100.moos 新增 pCanBridge 配置块(板卡 dbpath 指向
  data 目录)。
- build-board.sh / deploy.sh 支持 pCanBridge:产物部署 +
  pCanBridge.service 安装,服务列表加入第 5 个服务。
- .gitignore 忽略 bin/pCanBridge 与本机默认数据库。

已验证(板卡 192.168.0.223 实测):服务 active,连接设备正常,
~60 帧/s 透传,落库 9000+ 行 integrity ok、无丢帧。
This commit is contained in:
zjk
2026-08-31 19:14:59 +08:00
parent 7e4d08808a
commit b547772709
16 changed files with 1048 additions and 17 deletions
+2
View File
@@ -282,6 +282,8 @@ src/pPowerMangerHost/webassets_gen.h
# ============================================================
bin/pCCU
bin/pccuTest
bin/pCanBridge
pCanBridge_data.db
# ============================================================
# MS Office 临时锁文件 / 编辑器、备份残留
+18 -2
View File
@@ -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
}
+27 -3
View File
@@ -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" <<EOF
[Unit]
Description=pCanBridge for H100 Power Manager
After=moosdb.service
Requires=moosdb.service
[Service]
Type=simple
WorkingDirectory=${BOARD_BIN}
ExecStart=${BOARD_BIN}/pCanBridge --alias=pCanBridge /root/work/h100/missions/h100.moos
Restart=on-failure
RestartSec=5
[Install]
WantedBy=multi-user.target
EOF
${SSH} "${SSH_TARGET}" "systemctl daemon-reload && systemctl reset-failed pCanBridge 2>/dev/null || true"
ok "pCanBridge.service 已安装"
echo ""
ok "产物: ${BOARD_BIN}/pPowerManger (aarch64)"
${SSH} "${SSH_TARGET}" "file ${BOARD_BIN}/pPowerManger | cut -c1-60"
+14 -12
View File
@@ -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 "停止服务失败"
}
+1
View File
@@ -21,6 +21,7 @@ ADD_SUBDIRECTORY(pPowerManger)
add_subdirectory(pPMtest)
add_subdirectory(pPowerMangerHost)
add_subdirectory(pCCU)
add_subdirectory(pCanBridge)
##############################################################################
# END of CMakeLists.txt
##############################################################################
+43
View File
@@ -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}
)
+169
View File
@@ -0,0 +1,169 @@
#include "CanBridge.h"
#include "CanBridge_Info.h"
#include "MBUtils.h"
#include "MOOS/libMOOS/Comms/MOOSMsg.h"
#include <cstdio>
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<unsigned int>(dlc),
const_cast<uint8_t*>(data));
msg.SetSourceAux(m_channelName); // 通道 -> m_sSrcAux
msg.SetDoubleAux(static_cast<double>(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
+62
View File
@@ -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 <atomic>
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<unsigned long> m_notifyCount{0};
};
} // namespace canbridge
#endif // PCANBRIDGE_CAN_BRIDGE_H
+106
View File
@@ -0,0 +1,106 @@
/****************************************************************/
/* NAME: CanBridge_Info */
/* FILE: CanBridge_Info.cpp */
/****************************************************************/
#include <cstdlib>
#include <iostream>
#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","=<ProcessName> ");
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);
}
+15
View File
@@ -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
+158
View File
@@ -0,0 +1,158 @@
#include "CanDbStore.h"
#include "sqlite3.h"
#include "MOOS/libMOOS/Utils/MOOSUtilityFunctions.h"
#include <cstdio>
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<size_t>(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<std::mutex> 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<std::mutex> 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<sqlite3_int64>(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<CanFrameRow> CanDbStore::queryRecent(uint32_t canId, int limit) {
std::lock_guard<std::mutex> lock(m_mutex);
std::vector<CanFrameRow> 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<const char*>(sqlite3_column_text(stmt, 2)) : "";
r.frameInfo = sqlite3_column_int(stmt, 3);
r.canId = static_cast<uint32_t>(sqlite3_column_int64(stmt, 4));
r.dlc = sqlite3_column_int(stmt, 5);
r.hex = sqlite3_column_text(stmt, 6)
? reinterpret_cast<const char*>(sqlite3_column_text(stmt, 6)) : "";
rows.push_back(std::move(r));
}
sqlite3_finalize(stmt);
return rows;
}
long long CanDbStore::count() const {
std::lock_guard<std::mutex> 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
+67
View File
@@ -0,0 +1,67 @@
#ifndef PCANBRIDGE_CAN_DB_STORE_H
#define PCANBRIDGE_CAN_DB_STORE_H
#include <cstdint>
#include <string>
#include <vector>
#include <mutex>
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<CanFrameRow> 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
+200
View File
@@ -0,0 +1,200 @@
#include "CanEndpoint.h"
#include <sys/socket.h>
#include <netinet/in.h>
#include <netinet/tcp.h>
#include <arpa/inet.h>
#include <poll.h>
#include <unistd.h>
#include <fcntl.h>
#include <cerrno>
#include <cstring>
#include <iostream>
#include <chrono>
#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<uint16_t>(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<struct sockaddr*>(&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<uint32_t>(f[1]) << 24) |
(static_cast<uint32_t>(f[2]) << 16) |
(static_cast<uint32_t>(f[3]) << 8) |
static_cast<uint32_t>(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<int>(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
+87
View File
@@ -0,0 +1,87 @@
#ifndef PCANBRIDGE_CAN_ENDPOINT_H
#define PCANBRIDGE_CAN_ENDPOINT_H
#include <cstdint>
#include <string>
#include <vector>
#include <atomic>
#include <thread>
#include <functional>
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<void(uint8_t frameInfo, uint32_t id,
const uint8_t* data, int dlc)>;
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<bool> m_running{false};
std::atomic<bool> m_connected{false};
std::thread m_thread;
CanFrameCallback m_onFrame;
// TCP 流式接收缓冲(跨 recv 保持半包/粘包状态)
std::vector<uint8_t> m_buf;
std::atomic<unsigned long> m_frameCount{0};
std::atomic<unsigned long> m_errorCount{0};
std::atomic<unsigned long> m_reconnectCount{0};
std::atomic<double> m_lastFrameTime{0.0};
};
} // namespace canbridge
#endif // PCANBRIDGE_CAN_ENDPOINT_H
+49
View File
@@ -0,0 +1,49 @@
/************************************************************/
/* NAME: pCanBridge */
/* FILE: main.cpp */
/************************************************************/
#include <string>
#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;
}
+30
View File
@@ -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
}