|
| 1 | +/** |
| 2 | + * @file |
| 3 | + * @author Jaroslav Pesek <jaroslav.pesek@fit.cvut.cz> |
| 4 | + * @author Karel Hynek <hynekkar@cesnet.cz> |
| 5 | + * @author Pavel Siska <siska@cesnet.cz> |
| 6 | + * @brief Sampling Module: Sample flowdata |
| 7 | + * |
| 8 | + * This module distributes 1 unirec interface data to `n` trap outputs based on rules. |
| 9 | + * |
| 10 | + * SPDX-License-Identifier: BSD-3-Clause |
| 11 | + */ |
| 12 | + |
| 13 | +#include "logger/logger.hpp" |
| 14 | +#include "flowScatter.hpp" |
| 15 | +#include "unirec/unirec-telemetry.hpp" |
| 16 | +#include <libtrap/trap.h> |
| 17 | + |
| 18 | +#include <appFs.hpp> |
| 19 | +#include <argparse/argparse.hpp> |
| 20 | +#include <atomic> |
| 21 | +#include <csignal> |
| 22 | +#include <iostream> |
| 23 | +#include <stdexcept> |
| 24 | +#include <telemetry.hpp> |
| 25 | +#include <unirec++/unirec.hpp> |
| 26 | +#include <vector> |
| 27 | +#include <algorithm> |
| 28 | + |
| 29 | +using namespace Nemea; |
| 30 | + |
| 31 | +std::atomic<bool> g_stopFlag(false); |
| 32 | + |
| 33 | +void signalHandler(int signum) |
| 34 | +{ |
| 35 | + Nm::loggerGet("signalHandler")->info("Interrupt signal {} received", signum); |
| 36 | + g_stopFlag.store(true); |
| 37 | +} |
| 38 | + |
| 39 | +/** |
| 40 | + * @brief Handle a format change exception by adjusting the template. |
| 41 | + * |
| 42 | + * This function is called when a `FormatChangeException` is caught in the main loop. |
| 43 | + * |
| 44 | + * @param inputInterface Input interface for Unirec communication. |
| 45 | + * @param outputInterfaces Output interfaces for Unirec communication. |
| 46 | + */ |
| 47 | +void handleFormatChange(UnirecInputInterface& inputInterface, |
| 48 | + std::vector<UnirecOutputInterface>& outputInterfaces) |
| 49 | +{ |
| 50 | + inputInterface.changeTemplate(); |
| 51 | + uint8_t dataType; const char* spec = nullptr; |
| 52 | + if (trap_get_data_fmt(TRAPIFC_INPUT, 0, &dataType, &spec) != TRAP_E_OK) { |
| 53 | + throw std::runtime_error("Failed to get updated format from TRAP"); |
| 54 | + } |
| 55 | + for (auto& outIfc : outputInterfaces) { |
| 56 | + outIfc.changeTemplate(spec); |
| 57 | + } |
| 58 | +} |
| 59 | + |
| 60 | +/** |
| 61 | + * @brief Process the next Unirec record and sample them. |
| 62 | + * |
| 63 | + * This function receives the next Unirec record through the bidirectional interface |
| 64 | + * and performs sampling. |
| 65 | + * |
| 66 | + * @param inputInterface Input interface for Unirec communication. |
| 67 | + * @param outputInterfaces Output interfaces for Unirec communication. |
| 68 | + * @param scatter Sampler class for sampling. |
| 69 | + */ |
| 70 | + |
| 71 | +void processNextRecord(UnirecInputInterface& inputInterface, |
| 72 | + std::vector<UnirecOutputInterface>& outputInterfaces, |
| 73 | + Fs::FlowScatter& scatter) |
| 74 | +{ |
| 75 | + std::optional<UnirecRecordView> unirecRecord = inputInterface.receive(); |
| 76 | + if (!unirecRecord) { |
| 77 | + return; |
| 78 | + } |
| 79 | + size_t index = scatter.outputIndex(*unirecRecord); |
| 80 | + outputInterfaces[index].send(*unirecRecord); |
| 81 | +} |
| 82 | + |
| 83 | +/** |
| 84 | + * @brief Process Unirec records. |
| 85 | + * |
| 86 | + * The `processUnirecRecords` function continuously receives Unirec records through the provided |
| 87 | + * bidirectional interface (`biInterface`) and performs sampling. The loop runs indefinitely until |
| 88 | + * an end-of-file condition is encountered. |
| 89 | + * |
| 90 | + * @param inputInterface Input interface for Unirec communication. |
| 91 | + * @param outputInterfaces Output interfaces for Unirec communication. |
| 92 | + * @param scatter Sampler class for sampling. |
| 93 | + */ |
| 94 | +void processUnirecRecords(UnirecInputInterface& inputInterface, |
| 95 | + std::vector<UnirecOutputInterface>& outputInterfaces, |
| 96 | + Fs::FlowScatter& scatter) |
| 97 | +{ |
| 98 | + while (!g_stopFlag.load()) { |
| 99 | + try { |
| 100 | + processNextRecord(inputInterface, outputInterfaces, scatter); |
| 101 | + } catch (FormatChangeException& ex) { |
| 102 | + handleFormatChange(inputInterface, outputInterfaces); |
| 103 | + } catch (EoFException& ex) { |
| 104 | + break; |
| 105 | + } catch (std::exception& ex) { |
| 106 | + throw; |
| 107 | + } |
| 108 | + } |
| 109 | +} |
| 110 | + |
| 111 | +telemetry::Content getScatterTelemetry(const Fs::FlowScatter& scatter) |
| 112 | +{ |
| 113 | + auto stats = scatter.getStats(); |
| 114 | + |
| 115 | + telemetry::Dict dict; |
| 116 | + dict["totalRecords"] = stats.totalRecords; |
| 117 | + // dict["sampledRecords"] = stats.sampledRecords; |
| 118 | + return dict; |
| 119 | +} |
| 120 | + |
| 121 | +int main(int argc, char** argv) |
| 122 | +{ |
| 123 | + argparse::ArgumentParser program("Unirec Flow Scatter"); |
| 124 | + |
| 125 | + Nm::loggerInit(); |
| 126 | + auto logger = Nm::loggerGet("main"); |
| 127 | + |
| 128 | + signal(SIGINT, signalHandler); |
| 129 | + |
| 130 | + try { |
| 131 | + program.add_argument("-r", "--rule") |
| 132 | + .required() |
| 133 | + .help( |
| 134 | + "Specify the rule set.") |
| 135 | + .default_value(std::string("<>:(SRC_IP)")); |
| 136 | + program.add_argument("-c", "--count") |
| 137 | + .required() |
| 138 | + .help("Specify the number of output interfaces.") |
| 139 | + .scan<'i', int>() |
| 140 | + .default_value(5); |
| 141 | + program.add_argument("-m", "--appfs-mountpoint") |
| 142 | + .required() |
| 143 | + .help("path where the appFs directory will be mounted") |
| 144 | + .default_value(std::string("")); |
| 145 | + } catch (const std::exception& ex) { |
| 146 | + logger->error(ex.what()); |
| 147 | + return EXIT_FAILURE; |
| 148 | + } |
| 149 | + try { |
| 150 | + program.parse_known_args(argc, argv); |
| 151 | + } catch (const std::exception& ex) { |
| 152 | + logger->error(ex.what()); |
| 153 | + return EXIT_FAILURE; |
| 154 | + } |
| 155 | + |
| 156 | + size_t outputCount = 0; |
| 157 | + try { |
| 158 | + outputCount = static_cast<size_t>(program.get<int>("--count")); |
| 159 | + if (outputCount < 1 || outputCount > Fs::MAX_OUTPUTS) { |
| 160 | + throw std::runtime_error("Invalid number of output interfaces: " + std::to_string(outputCount) |
| 161 | + + ". Must be in range 1 to " + std::to_string(Fs::MAX_OUTPUTS)); |
| 162 | + } |
| 163 | + } catch (const std::exception& ex) { |
| 164 | + logger->error("Error parsing output count: {}", ex.what()); |
| 165 | + return EXIT_FAILURE; |
| 166 | + } |
| 167 | + |
| 168 | + std::shared_ptr<telemetry::Directory> telemetryRootDirectory; |
| 169 | + telemetryRootDirectory = telemetry::Directory::create(); |
| 170 | + |
| 171 | + std::unique_ptr<telemetry::appFs::AppFsFuse> appFs; |
| 172 | + |
| 173 | + try { |
| 174 | + auto mountPoint = program.get<std::string>("--appfs-mountpoint"); |
| 175 | + if (!mountPoint.empty()) { |
| 176 | + const bool tryToUnmountOnStart = true; |
| 177 | + const bool createMountPoint = true; |
| 178 | + appFs = std::make_unique<telemetry::appFs::AppFsFuse>( |
| 179 | + telemetryRootDirectory, |
| 180 | + mountPoint, |
| 181 | + tryToUnmountOnStart, |
| 182 | + createMountPoint); |
| 183 | + appFs->start(); |
| 184 | + } |
| 185 | + } catch (std::exception& ex) { |
| 186 | + logger->error(ex.what()); |
| 187 | + return EXIT_FAILURE; |
| 188 | + } |
| 189 | + |
| 190 | + try { |
| 191 | + const std::string rule = program.get<std::string>("--rule"); |
| 192 | + |
| 193 | + Unirec unirec({1, static_cast<int>(outputCount), "flowscatter", "Unirec flow scatter module"}); |
| 194 | + |
| 195 | + try { |
| 196 | + unirec.init(argc, argv); |
| 197 | + } catch (HelpException& ex) { |
| 198 | + std::cerr << program; |
| 199 | + return EXIT_SUCCESS; |
| 200 | + } catch (std::exception& ex) { |
| 201 | + logger->error(ex.what()); |
| 202 | + return EXIT_FAILURE; |
| 203 | + } |
| 204 | + |
| 205 | + UnirecInputInterface inputInterface = unirec.buildInputInterface(); |
| 206 | + std::vector<UnirecOutputInterface> outputInterfaces; |
| 207 | + outputInterfaces.reserve(outputCount); |
| 208 | + |
| 209 | + for (size_t i = 0; i < outputCount; ++i) { |
| 210 | + outputInterfaces.emplace_back(unirec.buildOutputInterface()); |
| 211 | + } |
| 212 | + |
| 213 | + Fs::FlowScatter scatter(outputCount, rule); |
| 214 | + |
| 215 | + auto telemetryInputDirectory = telemetryRootDirectory->addDir("input"); |
| 216 | + const telemetry::FileOps inputFileOps |
| 217 | + = {[&inputInterface]() { return Nm::getInterfaceTelemetry(inputInterface); }, nullptr}; |
| 218 | + const auto inputFile = telemetryInputDirectory->addFile("stats", inputFileOps); |
| 219 | + auto telemetryScatterDirectory = telemetryRootDirectory->addDir("flowscatter"); |
| 220 | + const telemetry::FileOps samplerFileOps |
| 221 | + = {[&scatter]() { return getScatterTelemetry(scatter); }, nullptr}; |
| 222 | + const auto samplerFile = telemetryScatterDirectory->addFile("stats", samplerFileOps); |
| 223 | + processUnirecRecords(inputInterface, outputInterfaces, scatter); |
| 224 | + } catch (std::exception& ex) { |
| 225 | + logger->error(ex.what()); |
| 226 | + return EXIT_FAILURE; |
| 227 | + } |
| 228 | + |
| 229 | + return EXIT_SUCCESS; |
| 230 | +} |
0 commit comments