490 lines
18 KiB
C++
490 lines
18 KiB
C++
#include "Data_Source.h"
|
||
#include <codecvt>
|
||
#include <string_view>
|
||
#include "../server/Global.h"
|
||
#include "Database.h"
|
||
Psc::serial::Serial* create_serial(std::string_view serial_name, Baud_Rate_Type baud_rate) {
|
||
auto serial = new Psc::serial::Serial;
|
||
serial->set_serial_name(serial_name);
|
||
serial->set_baud_rate(baud_rate);
|
||
serial->set_parity(Psc::serial::Parity::NoParity);
|
||
serial->set_data_bits(Psc::serial::DataBits::Data8);
|
||
serial->set_stop_bits(Psc::serial::StopBits::OneStop);
|
||
serial->set_flow_control(Psc::serial::FlowControl::HardwareControl);
|
||
serial->set_buffer_byte_size(10 * 1024);
|
||
if (!serial->open()) {
|
||
std::cerr << "createSerial " + std::string(serial_name) + ":" + std::to_string(baud_rate) + " 打开串口失败!\n";
|
||
Psc::fail_fast();
|
||
return nullptr;
|
||
}
|
||
else {
|
||
// std::cerr << "createSerial " + serial_name + ":" +
|
||
// std::to_string(baud_rate) + " 打开串口成功!\n";
|
||
}
|
||
return serial;
|
||
}
|
||
std::shared_ptr<Data_Source> create_from_json(const JSON* that_json, bool not_exist_use_default_value) {
|
||
std::string type = that_json->get_string("type");
|
||
std::shared_ptr<Data_Source> ret = create_data_source_from_type(type);
|
||
ret->from_json(that_json, not_exist_use_default_value);
|
||
return ret;
|
||
}
|
||
std::shared_ptr<Data_Source> Data_Source::that() {
|
||
return shared_from_this();
|
||
}
|
||
bool Data_Source::registered() const {
|
||
auto all_source = Global::instance()->mode_acs.data_source_config.map.list();
|
||
for (auto& item : all_source) {
|
||
if (item.get() == this) {
|
||
return true;
|
||
}
|
||
}
|
||
return false;
|
||
}
|
||
Psc::JSON Data_Source::get_state() {
|
||
Psc::JSON ret = Psc::JSON::object();
|
||
ret.append({"state", Psc::to_string(state)});
|
||
ret.append_list(get_custom_state_json().children);
|
||
ret.append({"ds: read_speed(byte)", read_speed});
|
||
ret.append({"ds: value_statistics(byte)", value_statistics});
|
||
return ret;
|
||
}
|
||
Data_Source::Data_Source() {
|
||
// parse_format = std::make_shared<Input_Format>();
|
||
// 写入到配置文件是懒加载 其他保存时他跟着保存
|
||
base_station.handle_when_updated = [this](SSR::HULC_Status_Message msg) {
|
||
if (!settings.member<&Data_Source_Data::update_form_gps>().read([](const auto& value) { return value; }))
|
||
return;
|
||
settings.write([&](auto& value) {
|
||
value.lat = msg.get_latitude();
|
||
value.lon = msg.get_longitude();
|
||
value.alt = msg.Alt;
|
||
});
|
||
};
|
||
}
|
||
std::string Data_Source::thread_key() const {
|
||
// return "Data_Source_Handle_Thread:[" + key + "]";
|
||
return key + "_DS_HT";
|
||
}
|
||
void File_Data_Source::before_handle_msg(std::shared_ptr<SSR::Mode_Msg>& msg) {
|
||
auto cur = std::dynamic_pointer_cast<SSR::Mode_Msg>(msg);
|
||
if (cur != nullptr) {
|
||
auto& pre = last_prase_msg;
|
||
if (!pre) {
|
||
pre = cur;
|
||
return;
|
||
}
|
||
cur->day_num = pre->day_num;
|
||
if (cur->time() < pre->time()) {
|
||
++cur->day_num;
|
||
}
|
||
}
|
||
}
|
||
std::vector<std::string> readLines(std::string_view path) {
|
||
std::vector<std::string> lines;
|
||
auto path_string = std::string(path);
|
||
std::filesystem::path p = std::filesystem::path((const char8_t*)path_string.c_str());
|
||
// std::filesystem::path p = std::filesystem::u8path(path);
|
||
std::ifstream istream(p);
|
||
if (!istream) {
|
||
std::cerr << "读取模拟数据源 readLines 无法打开文件:" << path_string << std::endl;
|
||
Psc::fail_fast();
|
||
}
|
||
// 将整个文件读入一个字符串
|
||
std::stringstream stream;
|
||
stream << istream.rdbuf();
|
||
// 使用 stringstream 按行分割内容
|
||
std::string line;
|
||
char ch;
|
||
while (stream.get(ch)) {
|
||
if (ch == ' ') {
|
||
continue;
|
||
}
|
||
if (ch == '\r' || ch == '\n') {
|
||
if (!line.empty()) {
|
||
lines.push_back(line);
|
||
line.clear();
|
||
}
|
||
// if (ch == '\r') {
|
||
// stream.get(); // 吃掉 \n
|
||
// }
|
||
// if (ch == '\n') {
|
||
// stream.get(); // 吃掉 \n
|
||
// }
|
||
// 处理 \r\n 组合:如果当前是 \r,下一个是 \n,跳过它
|
||
// if (ch == '\r' && stream.peek() == '\n') {
|
||
// stream.get(); // 吃掉 \n
|
||
// }
|
||
}
|
||
else {
|
||
line += ch;
|
||
}
|
||
if (line.size() == 1000) {
|
||
lines.push_back(line);
|
||
line.clear();
|
||
}
|
||
}
|
||
// 最后一行如果没有换行符也处理一下
|
||
if (!line.empty()) {
|
||
lines.push_back(line);
|
||
}
|
||
std::ostringstream oss;
|
||
oss << path_string + " 总行数:" + std::to_string(lines.size()) + " 有效行数:" + std::to_string(lines.size()) +
|
||
"\n" << std::flush;
|
||
std::cout << oss.str() << std::flush;
|
||
return lines;
|
||
}
|
||
std::string extractID(std::string_view logLine) {
|
||
size_t lastSpacePos = logLine.rfind(" "); // 查找最后一个空格
|
||
if (lastSpacePos != std::string::npos) {
|
||
return std::string(logLine.substr(lastSpacePos + 1)); // 从最后一个空格之后提取字符串
|
||
}
|
||
return std::string(logLine);
|
||
}
|
||
time_t convert_to_timestamp(std::string_view str) {
|
||
// 创建一个结构体 tm 来存储解析后的时间
|
||
std::tm timeStruct = {};
|
||
std::istringstream ss{std::string(str)};
|
||
ss >> std::get_time(&timeStruct, "%Y-%m-%d %H:%M:%S");
|
||
if (ss.fail()) {
|
||
std::cerr << "Failed to parse time" << std::endl;
|
||
return -1; // 如果解析失败,返回 -1
|
||
}
|
||
// 将 tm 转换为 time_t(时间戳)
|
||
time_t timestamp = std::mktime(&timeStruct);
|
||
if (timestamp == -1) {
|
||
std::cerr << "Failed to convert to time_t" << std::endl;
|
||
return -1; // 如果转换失败,返回 -1
|
||
}
|
||
return timestamp;
|
||
}
|
||
std::tuple<std::string, time_t> extractID2(std::string_view logLine) {
|
||
size_t lastSpacePos = logLine.rfind(" "); // 查找最后一个空格
|
||
time_t t = convert_to_timestamp(logLine.substr(0, 19));
|
||
if (lastSpacePos != std::string::npos) {
|
||
return {std::string(logLine.substr(lastSpacePos + 1)), t}; // 从最后一个空格之后提取字符串
|
||
}
|
||
return {std::string(logLine), t};
|
||
}
|
||
std::string bin_format(std::string_view hex, time_t t) {
|
||
std::tm* currentTime = std::localtime(&t);
|
||
auto sec = currentTime->tm_hour * 3600 + currentTime->tm_min * 60 + currentTime->tm_sec;
|
||
SSR::MLAT_timestamp a(sec, 0);
|
||
return create_Binary_Format_memory(hex, 0, &a);
|
||
}
|
||
std::optional<std::string> File_Data_Source::get_raw_line(int& ret_index) {
|
||
if (index == 0 && part_infos.empty()) {
|
||
ret_index = -1;
|
||
return std::nullopt;
|
||
}
|
||
std::optional<std::string> log_info = std::nullopt;
|
||
if (index == part_infos.size()) {
|
||
const auto config = specific.read([](const auto& value) { return value; });
|
||
if (config.play_mode == Play_Mode::loop) {
|
||
index = 0;
|
||
}
|
||
else {
|
||
if (!have_report_play_back_all_success) {
|
||
std::ostringstream oss;
|
||
oss << SSR::get_current_date() << " " << SSR::get_current_time()
|
||
<< " [进程:" << std::this_thread::get_id() << "] " << config.file_path + " 全部加载成功!\n"
|
||
<< std::flush;
|
||
std::cout << oss.str() << std::flush;
|
||
have_report_play_back_all_success = true;
|
||
}
|
||
ret_index = -1;
|
||
return std::nullopt;
|
||
}
|
||
}
|
||
ret_index = index + 1;
|
||
return part_infos[index++];
|
||
}
|
||
// 1A 33 1A 1A F1 FB 87 73 7E 7F a8001d81a87543b0a80000
|
||
std::string get_true_from_raw_line(Data_Source* ds, File_Data_Type data_type, int index, std::string_view log_info) {
|
||
std::string ret;
|
||
// 2025-01-06 07:20:17 [info] dc2541e65001401a3119b2dd1c7af67132201
|
||
if (data_type == File_Data_Type::BIN) {
|
||
return std::string(log_info);
|
||
}
|
||
else if (data_type == File_Data_Type::BIN_Text) {
|
||
auto info = extractID(log_info);
|
||
ret = hex2mem(info);
|
||
}
|
||
else if (data_type == File_Data_Type::AVR) {
|
||
auto [info, time] = extractID2(log_info);
|
||
ret = bin_format(info, time);
|
||
}
|
||
// 1A 33 10 01 A5 31 FB 9A 5D 8D 78 0E 47 99 08 8E 32 B0 08 BE E9 8A E6 1A 32
|
||
// 10 01 AA 14 61 7E 5E 02 E6 0F 38 7D 87 08 1A 32 10 01 AD B5 64 F7 63 5D 78
|
||
// 0F 9D BA 93 BA
|
||
else if (data_type == File_Data_Type::BIN_Blank_Text) {
|
||
std::string info(log_info);
|
||
std::string result = ds->last_char;
|
||
ds->last_char = "";
|
||
for (char c : info) {
|
||
if (std::isxdigit(c)) {
|
||
// 会判断字符 c 是否是十六进制数字字符
|
||
result += c;
|
||
}
|
||
else if (std::isspace(c)) {
|
||
//
|
||
}
|
||
else {
|
||
std::ostringstream oss;
|
||
oss << "存在非法字符[" << c << "]: value:" << (int)c << " in:" << log_info << LOG_POS;
|
||
server_logger->c_debug("非法字符:", {}, oss.str());
|
||
return "";
|
||
}
|
||
}
|
||
auto size = result.size();
|
||
if (size % 2 != 0) {
|
||
// 保存最后一个字符
|
||
ds->last_char = result.empty() ? "" : std::string(1, result.back());
|
||
result.pop_back();
|
||
}
|
||
ret = hex2mem(result);
|
||
}
|
||
else if (data_type == File_Data_Type::BIN_Blank_One_Line_With_Escape) {
|
||
auto size = log_info.size();
|
||
if (size % 2 != 0) {
|
||
return hex2mem(std::string(log_info.substr(0, size - 1)));
|
||
}
|
||
return hex2mem(std::string(log_info));
|
||
}
|
||
else if (data_type == File_Data_Type::BIN_Blank_One_Line_No_Escape) {
|
||
auto size = log_info.size();
|
||
std::string data;
|
||
if (size % 2 != 0) {
|
||
data = std::string(log_info.substr(0, size - 1));
|
||
}
|
||
else {
|
||
data = std::string(log_info);
|
||
}
|
||
return SSR::packet_to_escape_format(hex2mem(data));
|
||
}
|
||
else if (data_type == File_Data_Type::SIMPLE_BIN_Blank) {
|
||
std::string info(log_info);
|
||
std::string result;
|
||
for (char c : info) {
|
||
if (c != ' ') {
|
||
result += c;
|
||
}
|
||
}
|
||
// 1 + 1 + 6 + 1 + 2/7/14
|
||
// 1a + type time_stamp signal_level data
|
||
int size = result.size();
|
||
if (size != 2 * 2 && size != 7 * 2 && size != 14 * 2) {
|
||
std::cout << "line_offset:[" << index << "] " << "大小" << size
|
||
<< " 不为偶数,或不为2,7,14 不能转换成二进制格式:" << result << std::endl;
|
||
return "";
|
||
}
|
||
std::string head("\x1a");
|
||
std::string type;
|
||
std::string time_stamp = "\x11\x22\x33\x44\x55\x66";
|
||
std::string signal_level("\xFF");
|
||
if (size == 2 * 2) {
|
||
type = R"(1)";
|
||
}
|
||
else if (size == 7 * 2) {
|
||
type = R"(2)";
|
||
}
|
||
else if (size == 14 * 2) {
|
||
type = R"(3)";
|
||
}
|
||
ret = SSR::packet_to_escape_format(head + type + time_stamp + signal_level + hex2mem(result));
|
||
// std::cout << ret.size() << std::endl;
|
||
}
|
||
else {
|
||
std::cout << "未知的回放类型 " << static_cast<int>(data_type) << std::endl;
|
||
Psc::fail_fast();
|
||
}
|
||
return ret;
|
||
}
|
||
asio::awaitable<void> File_Data_Source::_close() {
|
||
std::lock_guard g(mtx);
|
||
index = 0;
|
||
part_infos.clear();
|
||
co_return;
|
||
}
|
||
asio::awaitable<void> File_Data_Source::_open() {
|
||
std::lock_guard g(mtx);
|
||
const auto config = specific.read([](const auto& value) { return value; });
|
||
auto path = Psc::get_abs_path(config.file_path);
|
||
namespace fs = std::filesystem;
|
||
if (!fs::exists(path)) {
|
||
state = "文件不存在";
|
||
co_return;
|
||
}
|
||
if (config.data_type == File_Data_Type::BIN) {
|
||
part_infos = readBinaryFileAsString(path, 1000);
|
||
}
|
||
else {
|
||
part_infos = readLines(path);
|
||
}
|
||
state = "已加载" + to_string(part_infos.size()) + "长度数据!";
|
||
co_return;
|
||
}
|
||
std::vector<std::string> File_Data_Source::readBinaryFileAsString(std::string_view filepath, size_t part_size) {
|
||
// 打开文件(以二进制模式)
|
||
std::ifstream file(filepath.data(), std::ios::binary);
|
||
if (!file) {
|
||
std::cerr << "无法打开文件: " << filepath << std::endl;
|
||
return {}; // 文件打开失败,返回空vector
|
||
}
|
||
// 获取文件大小
|
||
file.seekg(0, std::ios::end);
|
||
size_t fileSize = file.tellg();
|
||
file.seekg(0, std::ios::beg);
|
||
// 存储读取的部分
|
||
std::vector<std::string> parts;
|
||
// 读取文件的每一部分
|
||
size_t bytesRead = 0;
|
||
while (bytesRead < fileSize) {
|
||
// 计算每个部分的长度(最后一部分可能小于part_size)
|
||
size_t remaining = fileSize - bytesRead;
|
||
size_t currentPartSize = (remaining < part_size) ? remaining : part_size;
|
||
// 创建buffer来存储当前部分
|
||
std::string buffer(currentPartSize, '\0');
|
||
// 读取当前部分
|
||
file.read(&buffer[0], currentPartSize);
|
||
// 将读取的部分添加到vector
|
||
parts.push_back(std::move(buffer));
|
||
// 更新已读取字节数
|
||
bytesRead += currentPartSize;
|
||
}
|
||
// 关闭文件
|
||
file.close();
|
||
return parts;
|
||
}
|
||
void File_Data_Source::origin_data_transform_mode_data(std::string& mode_data) {
|
||
std::lock_guard g(mtx);
|
||
assert(mode_data.size() == 0);
|
||
int ret_index = 0;
|
||
const auto config = specific.read([](const auto& value) { return value; });
|
||
if (config.play_mode != Play_Mode::analysis) {
|
||
std::optional<std::string> log_info = get_raw_line(ret_index);
|
||
mode_data = log_info.has_value() ? get_true_from_raw_line(this, config.data_type, ret_index, log_info.value()) : "";
|
||
}
|
||
else {
|
||
int num = 0;
|
||
std::string ret;
|
||
std::optional<std::string> log_info = get_raw_line(ret_index);
|
||
while (log_info.has_value()) {
|
||
num++;
|
||
auto data = log_info.value();
|
||
auto cur = get_true_from_raw_line(this, config.data_type, ret_index, data);
|
||
ret += cur;
|
||
// 一定要在这个位置,否则会导致遗漏报文
|
||
if (num == 1000) {
|
||
break;
|
||
}
|
||
log_info = get_raw_line(ret_index);
|
||
}
|
||
// std::cout << "size == " << t.size() << std::endl;
|
||
mode_data = ret;
|
||
}
|
||
}
|
||
void File_Data_Source::handle_mode_s(std::shared_ptr<SSR::Mode_S_Msg>& msg) {
|
||
if (!msg)
|
||
return;
|
||
cache_list.push_back(std::make_shared<Cache_Msg>(std::move(msg)));
|
||
if (cache_list.empty())
|
||
return;
|
||
if (pre_cache == nullptr) {
|
||
pre_cache = cache_list.front();
|
||
cache_list.pop_front();
|
||
SSR::Play_Back_Time_Point& start = player_clock.start_time;
|
||
start.day_num = pre_cache->day_num;
|
||
start.day_sec = pre_cache->msg->mlat_timestamp.daysec;
|
||
}
|
||
auto sys_time = player_clock.get_cur_time_point();
|
||
while (true) {
|
||
if (cache_list.empty())
|
||
break;
|
||
auto pre = pre_cache;
|
||
auto cur = cache_list.front();
|
||
cur->day_num = pre->day_num;
|
||
if (cur->time() < pre->time()) {
|
||
cur->day_num++;
|
||
}
|
||
cache_list.pop_front();
|
||
pre_cache = cur;
|
||
Data_Source::handle_mode_s(cur->msg);
|
||
if (cur->time() > sys_time) {
|
||
break;
|
||
}
|
||
}
|
||
}
|
||
void Dll_Data_Source::origin_data_transform_mode_data(std::string& mode_data) {
|
||
auto start = std::chrono::steady_clock::now();
|
||
if (read_func_ptr == nullptr)
|
||
return;
|
||
const auto config = specific.read([](const auto& value) { return value; });
|
||
auto buf = reinterpret_cast<char*>(buffer.data());
|
||
auto len = read_func_ptr(buf, config.buffer_size);
|
||
vs.update(len);
|
||
std::string data(buf, len);
|
||
if (data.empty()) {
|
||
auto end = std::chrono::steady_clock::now();
|
||
auto cost = std::chrono::duration_cast<std::chrono::microseconds>(end - start).count();
|
||
// std::cout << "Dll_Data_Source::read cost: " << cost << " us\n";
|
||
return;
|
||
}
|
||
auto ret = get_true_from_raw_line(this, config.data_type, -1, data);
|
||
auto end = std::chrono::steady_clock::now();
|
||
auto cost = std::chrono::duration_cast<std::chrono::microseconds>(end - start).count();
|
||
// std::cout << "Dll_Data_Source::read cost: " << cost << " us\n";
|
||
mode_data = ret;
|
||
}
|
||
asio::awaitable<void> Shared_Memory_Data_Source::_open() {
|
||
const auto config = specific.read([](const auto& value) { return value; });
|
||
sm = std::make_unique<Psc::SM_RingBuffer>();
|
||
sm->init(config.shared_memory_name, config.shared_memory_size);
|
||
co_return;
|
||
}
|
||
asio::awaitable<void> Shared_Memory_Data_Source::_close() {
|
||
sm.reset();
|
||
co_return;
|
||
}
|
||
void Shared_Memory_Data_Source::origin_data_transform_mode_data(std::string& mode_data) {
|
||
auto size = sm->shm.size();
|
||
size_t length = size;
|
||
std::string data;
|
||
data.resize(size);
|
||
bool ok = sm->read((uint8_t*)data.data(), length);
|
||
if (!ok) {
|
||
mode_data = "";
|
||
}
|
||
if (length != size) {
|
||
data.resize(length);
|
||
}
|
||
if (Global::instance()->console_config.mode_s_console) {
|
||
std::cout << "read:" << data << std::endl;
|
||
}
|
||
const auto config = specific.read([](const auto& value) { return value; });
|
||
auto ret = get_true_from_raw_line(this, config.data_type, -1, data);
|
||
mode_data = ret;
|
||
}
|
||
void Data_Source_Config::server(Global* g) {
|
||
auto& svr = g->svr;
|
||
auto& api = g->api;
|
||
svr.Post(api + "get_data_source_connect", [this](HTTP_Param) {
|
||
CHECK_JSON_PARAM
|
||
HTTP_REQUIRE_VALUE(key, params.try_get_string("key"))
|
||
auto value = map.get(key);
|
||
JSON ret{nullptr};
|
||
if (value.has_value())
|
||
ret = value.value()->get_all_connect_feed();
|
||
res->setBody(warp(ret).to_json_string());
|
||
});
|
||
svr.Post(api + "get_data_source_state", [this](HTTP_Param) {
|
||
CHECK_JSON_PARAM
|
||
HTTP_REQUIRE_VALUE(key, params.try_get_string("key"))
|
||
auto value = map.get(key);
|
||
JSON ret{nullptr};
|
||
if (value.has_value())
|
||
ret = value.value()->get_state();
|
||
res->setBody(warp(ret).to_json_string());
|
||
});
|
||
}
|