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
9 changes: 5 additions & 4 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,24 +25,22 @@ chrono = "~0.4"
config = "~0.13"
evalexpr = "~9.0"
glob = "~0.3"
jwt = "~0.16"
itertools = "~0.10"
hmac = "~0.12"
http = "~0.2"
new_string_template = "~1.4"
regex = "~1.8"
reqwest = { version = "~0.11", default-features = false, features = ["rustls-tls", "json"] }
reqwest = { version = "~0.12", default-features = false, features = ["rustls-tls", "json"] }
schemars = "~0.8"
serde = { version = "~1.0", features = ["derive"] }
serde_json = "~1.0"
serde_yaml = "~0.9"
sha2 = "~0.10"
tokio = { version = "~1.42", features = ["full"] }
tower = { version = "~0.4" }
tower-http = { version = "~0.4", features = ["trace", "request-id", "util"] }
tracing = "~0.1"
tracing-subscriber = { version = "~0.3", features = ["env-filter"] }
uuid = { version = "~1.3", features = ["v4", "fast-rng"] }
zitadel = { version = "~5.7", default-features = false, features = ["credentials"] }

[dev-dependencies]
mockito = "~1.0"
Expand All @@ -51,6 +49,9 @@ tempfile = "~3.5"
tokio-test = "*"
tower = { version = "0.4", features = ["util"] }
hyper = { version = "0.14", features = ["full"] }
jsonwebtoken = { version = "~11.1", default-features = false, features = ["use_pem", "aws_lc_rs"] }
rand = "~0.8"
rsa = "~0.9"


[target.'cfg(all(target_env = "musl", target_pointer_width = "64"))'.dependencies.jemallocator]
Expand Down
157 changes: 99 additions & 58 deletions src/bin/reporter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@

extern crate anyhow;

use anyhow::Context;
use cloudmon_metrics::config::OidcIdentity;
use cloudmon_metrics::sd::{
build_auth_headers, build_component_id_cache, build_incident_data, create_incident,
fetch_components, find_component_id, Component, ComponentAttribute,
Expand All @@ -17,24 +19,14 @@ use reqwest::ClientBuilder;
use tokio::signal;
use tokio::time::{sleep, Duration};

use serde::{Deserialize, Serialize};

use std::collections::HashMap;

use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};

const CLIENT_TIMEOUT_SECS: u64 = 2;

/// Component status for V1 API (legacy, will be replaced)
#[derive(Deserialize, Serialize, Debug)]
pub struct ComponentStatus {
pub name: String,
pub impact: u8,
pub attributes: Vec<ComponentAttribute>,
}

#[tokio::main]
async fn main() {
async fn main() -> anyhow::Result<()> {
//Enable logging
tracing_subscriber::registry()
.with(tracing_subscriber::EnvFilter::new(
Expand All @@ -46,7 +38,15 @@ async fn main() {
tracing::info!("starting cloudmon-metrics-reporter");

// Parse config
let config = Config::new("config.yaml").unwrap();
let config = Config::new("config.yaml").context("Failed to load config.yaml")?;

// Fail closed: the service identity comes from the key file, which is read here once at startup
let oidc_identity = config
.status_dashboard
.as_ref()
.map(|sdb_config| sdb_config.oidc_identity())
.transpose()
.context("Invalid status_dashboard OIDC configuration")?;

// Set up CTRL+C handlers
let ctrl_c = async {
Expand All @@ -68,15 +68,26 @@ async fn main() {

// Execute metric_watcher unless need to stop
tokio::select! {
_ = metric_watcher(&config) => {},
res = metric_watcher(&config, oidc_identity.as_ref()) => {
if let Err(e) = &res {
tracing::error!(error = %e, "metric reporter stopped with an error");
}
// Fatal reporter errors must terminate the process with a non-zero exit code
res.context("metric reporter failed")?;
},
_ = ctrl_c => {},
_ = terminate => {},
}

tracing::info!("stopped cloudmon-metrics-reporter");

Ok(())
}

async fn metric_watcher(config: &Config) {
async fn metric_watcher(
config: &Config,
oidc_identity: Option<&OidcIdentity>,
) -> anyhow::Result<()> {
tracing::info!("starting metric reporter thread");
// Init reqwest client
let req_client: reqwest::Client = ClientBuilder::new()
Expand Down Expand Up @@ -118,10 +129,8 @@ async fn metric_watcher(config: &Config) {
.status_dashboard
.as_ref()
.expect("Status dashboard section is missing");

// Build authorization headers using status_dashboard module (T021, T022, T023 - US3)
// VERIFIED: Existing HMAC-JWT mechanism works unchanged with V2 endpoints
let headers = build_auth_headers(sdb_config.secret.as_deref());
let oidc_identity =
oidc_identity.context("Status Dashboard OIDC service identity is missing")?;

// Initialize component ID cache at startup with retry logic (T024, T025, T026, T027)
// Per FR-006: 3 retry attempts with 60-second delays
Expand All @@ -136,6 +145,11 @@ async fn metric_watcher(config: &Config) {
"attempting to fetch components from Status Dashboard"
);

// D6: a fresh token is requested per authenticated call, it is never reused or cached
let headers = build_auth_headers(oidc_identity)
.await
.context("Failed to obtain Status Dashboard authorization headers")?;

match fetch_components(&req_client, &sdb_config.url, &headers).await {
Ok(components) => {
tracing::info!(
Expand Down Expand Up @@ -175,8 +189,7 @@ async fn metric_watcher(config: &Config) {
let mut component_cache = match component_cache {
Some(cache) => cache,
None => {
tracing::error!("component cache initialization failed, exiting metric_watcher");
return;
anyhow::bail!("component cache initialization failed, exiting metric_watcher");
}
};

Expand Down Expand Up @@ -235,31 +248,45 @@ async fn metric_watcher(config: &Config) {
"component not found in cache, attempting cache refresh"
);

match fetch_components(
&req_client,
&sdb_config.url,
&headers,
)
.await
{
Ok(components) => {
tracing::info!(
component_count = components.len(),
"cache refreshed"
);
component_cache =
build_component_id_cache(components);
// Retry lookup after refresh
component_id = find_component_id(
&component_cache,
comp,
);
match build_auth_headers(oidc_identity).await {
Ok(headers) => {
match fetch_components(
&req_client,
&sdb_config.url,
&headers,
)
.await
{
Ok(components) => {
tracing::info!(
component_count =
components.len(),
"cache refreshed"
);
component_cache =
build_component_id_cache(
components,
);
component_id = find_component_id(
&component_cache,
comp,
);
}
Err(e) => {
tracing::warn!(
error = %e,
component_name =
comp.name.as_str(),
"failed to refresh component cache"
);
}
}
}
Err(e) => {
tracing::warn!(
tracing::error!(
error = %e,
component_name = comp.name.as_str(),
"failed to refresh component cache"
"failed to obtain authorization headers, skipping component cache refresh"
);
}
}
Expand Down Expand Up @@ -302,30 +329,44 @@ async fn metric_watcher(config: &Config) {
"creating incident: health metric indicates service degradation"
);

// Create incident via V2 API
match create_incident(
&req_client,
&sdb_config.url,
&headers,
&incident_data,
)
.await
{
Ok(_) => {
tracing::info!(
component_id = id,
impact = impact,
"incident created successfully"
);
match build_auth_headers(oidc_identity).await {
Ok(headers) => {
match create_incident(
&req_client,
&sdb_config.url,
&headers,
&incident_data,
)
.await
{
Ok(_) => {
tracing::info!(
component_id = id,
impact = impact,
"incident created successfully"
);
}
Err(e) => {
tracing::error!(
error = %e,
component_id = id,
service = component.0
.as_str(),
environment = env
.name
.as_str(),
"failed to create incident"
);
}
}
}
Err(e) => {
// Error logging with details (FR-015)
tracing::error!(
error = %e,
component_id = id,
service = component.0.as_str(),
environment = env.name.as_str(),
"failed to create incident"
"failed to obtain authorization headers, skipping incident creation"
);
}
}
Expand Down
Loading
Loading