Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion proto/Juniper/juniper_gnmi.proto
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ extend google.protobuf.FileOptions {

// gNMI_service is the current version of the gNMI service, returned through
// the Capabilities RPC.
option (gnmi_service) = "0.8.0";
option (gnmi_service) = "0.7.0";

service gNMI {
// Capabilities allows the client to retrieve the set of capabilities that
Expand Down
6 changes: 6 additions & 0 deletions proto/Juniper/juniper_telemetry.proto
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,9 @@ message TelemetryFieldOptions {
optional bool is_timestamp = 2;
optional bool is_counter = 3;
optional bool is_gauge = 4;
optional bool is_instance = 5;
optional bool is_gauge_min = 6;
optional bool is_gauge_max = 7;
}

message TelemetryStream {
Expand Down Expand Up @@ -72,6 +75,9 @@ message TelemetryStream {
// minor version
optional uint32 version_minor = 8;

// end-of-message marker, set to true when the end of wrap is reached
optional bool eom = 9;

optional IETFSensors ietf = 100;

optional EnterpriseSensors enterprise = 101;
Expand Down
20 changes: 17 additions & 3 deletions proto/Juniper/juniper_telemetry_header_extension.proto
Original file line number Diff line number Diff line change
@@ -1,5 +1,12 @@
syntax = "proto3";

enum StreamType {
INITIAL_SYNC = 0;
ONCHANGE = 1;
PERIODIC = 2;
}


// Present as first gNMI update in all packets
message GnmiJuniperTelemetryHeaderExtension {
// router name:export IP address
Expand Down Expand Up @@ -30,15 +37,22 @@ message GnmiJuniperTelemetryHeaderExtension {
// Stream creation timestamp in milliseconds
int64 stream_creation_timestamp = 10;

// Event timestamp in milliseconds
int64 event_timestamp = 11;
// [Deprecated] Event timestamp in milliseconds
int64 event_timestamp = 11 [deprecated=true];

// Export timestamp in milliseconds
int64 export_timestamp = 12;
int64 export_timestamp = 12;

// Subsequence number
uint64 sub_sequence_number = 13;

// End of marker
bool eom = 14;

// Event publish timestamp in milliseconds
int64 event_publish_timestamp = 15;

// Stream type of packet
StreamType stream_id = 16;

}
193 changes: 150 additions & 43 deletions src/dataManipulation/data_manipulation.cc
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,11 @@ bool DataManipulation::BuildEnvelope(
std::string &json_str_out)
{
Json::Value root;
Json::Value telemetry_json;
JSONCPP_STRING error;
Json::CharReaderBuilder builder_r;
const std::unique_ptr<Json::CharReader> reader(builder_r.newCharReader());

set_sequence_number();
root["event_type"] = "gRPC";
root["serialization"] = "json_string";
Expand All @@ -23,7 +28,18 @@ bool DataManipulation::BuildEnvelope(
root["writer_id"] = writer_id;
root["telemetry_node"] = peer_ip;
root["telemetry_port"] = static_cast<uint16_t>(std::stoi(peer_port));
root["telemetry_data"] = telemetry_body;

/*
* Optimization: Attempt to parse the telemetry_body into a valid JSON object.
* If successful, append it as a nested JSON object under "telemetry_data".
* Otherwise, fallback to appending it as a raw string.
* This prevents payload data from being treated merely as plain text.
*/
if (reader->parse(telemetry_body.c_str(), telemetry_body.c_str() + telemetry_body.length(), &telemetry_json, &error)) {
root["telemetry_data"] = telemetry_json;
} else {
root["telemetry_data"] = telemetry_body;
}

if (label_map != nullptr) {
Json::Value jlabel_map;
Expand Down Expand Up @@ -74,19 +90,31 @@ bool DataManipulation::MetaData(std::string &json_str,
"conversion to JSON failure, {}", error);
return false;
} else {
root.clear();
set_sequence_number();
root["event_type"] = "gRPC";
root["serialization"] = "json_string";
root["seq"] = static_cast<uint64_t>(get_sequence_number());
root["timestamp"] = timestamp;
root["writer_id"] = main_cfg_parameters.at("writer_id");
root["telemetry_node"] = peer_ip;
root["telemetry_port"] = static_cast<uint16_t>(std::stoi(peer_port));
root["telemetry_data"] = json_str;
spdlog::get("multi-logger")->info("[MetaData] data-manipulation: "
"{} meta-data added successfully", peer_ip);
}
/*
* Optimization: Temporarily store the parsed JSON into `telemetry_json`.
* Clear the `root` object to rebuild the standard envelope structure.
*/
Json::Value telemetry_json = root;
root.clear();

set_sequence_number();
root["event_type"] = "gRPC";
root["serialization"] = "json_string";
root["seq"] = static_cast<uint64_t>(get_sequence_number());
root["timestamp"] = timestamp;
root["writer_id"] = main_cfg_parameters.at("writer_id");

/*
* Reattach the preserved `telemetry_json` under the "telemetry_data" key
* to maintain a consistent JSON schema.
*/
root["telemetry_data"] = telemetry_json;
root["telemetry_node"] = peer_ip;
root["telemetry_port"] = static_cast<uint16_t>(std::stoi(peer_port));

spdlog::get("multi-logger")->info("[MetaData] data-manipulation: "
"{} meta-data added successfully", peer_ip);
}

json_str_out = Json::writeString(builder_w, root);

Expand Down Expand Up @@ -272,15 +300,40 @@ bool DataManipulation::JuniperExtension(
if (!juniper_tlm_header_ext.system_id().empty()) {
root["system_id"] = juniper_tlm_header_ext.system_id();
}

stream_data_in.clear();
google::protobuf::util::JsonPrintOptions opt;
opt.add_whitespace = false;
google::protobuf::util::MessageToJsonString(

/*
* Optimization: Added error checking for Protobuf to JSON string conversion
* to prevent unexpected empty string parsing errors.
*/
auto status = google::protobuf::util::MessageToJsonString(
juniper_tlm_header_ext,
&stream_data_in,
opt);
root["extension"] = stream_data_in;

if (!status.ok()) {
logger->error("[JuniperExtension] MessageToJsonString failed: {}", status.ToString());
return false;
}

/*
* Optimization: Introduced a JSON parser to parse the extension header string
* back into a genuine JSON object, eliminating raw text formats containing
* escape characters (e.g., '\').
*/
Json::Value ext_root;
JSONCPP_STRING ext_error;
Json::CharReaderBuilder ext_builder;
const std::unique_ptr<Json::CharReader> ext_reader(ext_builder.newCharReader());

if (ext_reader->parse(stream_data_in.c_str(), stream_data_in.c_str() + stream_data_in.length(), &ext_root, &ext_error)) {
root["extension"] = ext_root;
} else {
root["extension"] = stream_data_in;
}
} else {
return false;
}
Expand All @@ -295,33 +348,32 @@ bool DataManipulation::JuniperUpdate(juniper_gnmi::SubscribeResponse &juniper_st
Json::Value &root)
{
auto logger = spdlog::get("multi-logger");
if (logger->should_log(spdlog::level::debug)) {
std::string raw_data;
google::protobuf::util::JsonPrintOptions opt;
opt.add_whitespace = false;
auto status = google::protobuf::util::MessageToJsonString(juniper_stream, &raw_data, opt);
if (!status.ok()) {
logger->error("[JuniperDebug] Failed to convert protobuf to JSON: {}", status.ToString());
}
logger->debug("[JuniperDebug] pre-JuniperUpdate data: {}", raw_data);
}

if (juniper_stream.has_update()) {
const auto &jup = juniper_stream.update();
std::uint64_t notification_timestamp = jup.timestamp();

int path_idx = 0;
/*
* Optimization: Introduced dedicated "tags" and "metrics" JSON objects.
* "tags" stores data classification labels, and "metrics" stores actual values.
* This structure is highly compatible with Time-Series Databases (TSDB).
*/
Json::Value tags(Json::objectValue);
Json::Value metrics(Json::objectValue);
Json::Value sensor_path(Json::arrayValue);

int path_idx = 0;
while (path_idx < jup.prefix().elem_size()) {
Json::Value path_element;
path_element["name"] = jup.prefix().elem().at(path_idx).name();
std::string elem_name = jup.prefix().elem().at(path_idx).name();
path_element["name"] = elem_name;

if (jup.prefix().elem().at(path_idx).key_size() > 0) {
Json::Value filters;
for (const auto &[key, value] :
jup.prefix().elem().at(path_idx).key()) {
for (const auto &[key, value] : jup.prefix().elem().at(path_idx).key()) {
filters[key] = value;
// Dynamically extract all filters to be used as unified tags
tags[elem_name + "_" + key] = value;
}
path_element["filters"] = filters;
}
Expand All @@ -332,27 +384,82 @@ bool DataManipulation::JuniperUpdate(juniper_gnmi::SubscribeResponse &juniper_st
root["sensor_path"] = sensor_path;
root["notification_timestamp"] = notification_timestamp;

std::string path;
Json::Value value;
// Optimization: Moved the JSON Reader initialization outside the loop to significantly reduce performance overhead
Json::CharReaderBuilder rbuilder;
const std::unique_ptr<Json::CharReader> reader(rbuilder.newCharReader());
JSONCPP_STRING err;

for (const auto &_jup : jup.update()) {
int path_idx = 0;
path.clear();
while (path_idx < _jup.path().elem_size()) {
path.append("/");
path.append(_jup.path().elem().at(path_idx).name());
path_idx++;
Json::Value value;

// Optimization: Comprehensively cover all possible gNMI TypedValue formats
if (_jup.val().has_int_val()) {
value = (Json::Int64) _jup.val().int_val();
} else if (_jup.val().has_uint_val()) {
value = (Json::UInt64) _jup.val().uint_val();
} else if (_jup.val().has_string_val()) {
value = _jup.val().string_val();
} else if (_jup.val().has_bool_val()) {
value = _jup.val().bool_val();
} else if (_jup.val().has_float_val()) {
value = _jup.val().float_val();
} else if (_jup.val().has_json_val() || _jup.val().has_json_ietf_val()) {
// Optimization: Merge handling for json_val and OpenConfig's json_ietf_val, attempting to parse them into actual JSON objects
std::string raw_json = _jup.val().has_json_val() ?
_jup.val().json_val() :
_jup.val().json_ietf_val();

if (!reader->parse(raw_json.c_str(), raw_json.c_str() + raw_json.length(), &value, &err)) {
value = raw_json;
}
} else if (_jup.val().has_bytes_val()) {
value = _jup.val().bytes_val();
} else if (_jup.val().has_ascii_val()) {
// Optimization: Added support for legacy ASCII text format
value = _jup.val().ascii_val();
} else {
// Fallback mechanism: Prevent crashes on encountering unknown protocol types and log a warning
value = "unsupported_type";
logger->warn("[JuniperUpdate] Encountered unknown or unsupported gNMI TypedValue from device.");
}

value = _jup.val().json_val();
root[path] = value;
/*
* Core Optimization: Dynamically construct a nested JSON structure (handling infinite depths).
* Transforms flat paths (e.g., /interface/eth0/stats) into a native JSON Tree representation.
*/
Json::Value* current_node = &metrics;
for (int i = 0; i < _jup.path().elem_size(); ++i) {
std::string node_name = _jup.path().elem().at(i).name();

if (_jup.path().elem().at(i).key_size() > 0) {
for (const auto &[k, v] : _jup.path().elem().at(i).key()) {
node_name += "[" + k + "=" + v + "]";
}
}

if (i == _jup.path().elem_size() - 1) {
(*current_node)[node_name] = value;
} else {
if (!current_node->isMember(node_name) || !(*current_node)[node_name].isObject()) {
(*current_node)[node_name] = Json::Value(Json::objectValue);
}
current_node = &((*current_node)[node_name]);
}
}
}

if (!tags.empty()) {
root["tags"] = tags;
}
if (!metrics.empty()) {
root["metrics"] = metrics;
}
}

Json::StreamWriterBuilder builder_w;
builder_w["emitUTF8"] = true;
builder_w["indentation"] = "";
const std::unique_ptr<Json::StreamWriter> writer(
builder_w.newStreamWriter());
const std::unique_ptr<Json::StreamWriter> writer(builder_w.newStreamWriter());
json_str_out = Json::writeString(builder_w, root);

return true;
Expand Down
Loading