Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

chore: restore prometheus metrics in v2 #728

Merged
merged 2 commits into from
Dec 30, 2023
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
208 changes: 146 additions & 62 deletions Cargo.lock

Large diffs are not rendered by default.

6 changes: 4 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -47,14 +47,15 @@ serde = { version = "1.0.152", features = ["derive"] }
serde_json = { version = "1.0.104", features = ["arbitrary_precision"] }
strum = "0.24"
strum_macros = "0.25"
prometheus_exporter = { version = "0.8.5", default-features = false }
prometheus_exporter_base = { version = "1.4.0", features = ["hyper", "hyper_server"] }
unicode-truncate = "0.2.0"
thiserror = "1.0.39"
indicatif = "0.17.3"
lazy_static = "1.4.0"
tracing = "0.1.37"
tracing-subscriber = "0.3.17"
tokio = { version = "1", features = ["rt"] }
anyhow = "1.0.77"
tokio = { version = "1", features = ["rt", "rt-multi-thread"] }
async-trait = "0.1.68"
elasticsearch = { version = "8.5.0-alpha.1", optional = true }
murmur3 = { version = "0.5.2", optional = true }
Expand All @@ -72,6 +73,7 @@ file-rotate = { version = "0.7.5", optional = true }
tonic = { version = "0.9.2", features = ["tls", "tls-roots"], optional = true }
futures = { version = "0.3.28", optional = true }


# aws
aws-config = { version = "^1.1", optional = true }
aws-types = { version = "^1.1", optional = true }
Expand Down
1 change: 1 addition & 0 deletions examples/metrics/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
*.json
12 changes: 12 additions & 0 deletions examples/metrics/daemon.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
[source]
type = "N2N"
peers = ["relays-new.cardano-mainnet.iohk.io:3001"]

[intersect]
type = "Point"
value = [4493860, "ce7f821d2140419fea1a7900cf71b0c0a0e94afbb1f814a6717cff071c3b6afc"]

[sink]
type = "Noop"

[metrics]
54 changes: 42 additions & 12 deletions src/bin/oura/daemon.rs
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
use gasket::runtime::Tether;
use oura::{cursor, filters, framework::*, sinks, sources};
use serde::Deserialize;
use std::time::Duration;
use std::{fmt::Debug, sync::Arc, time::Duration};
use tracing::{info, warn};

use crate::console;
use crate::{console, prometheus};

#[derive(Deserialize)]
struct ConfigRoot {
Expand All @@ -16,6 +16,7 @@ struct ConfigRoot {
chain: Option<ChainConfig>,
retries: Option<gasket::retries::Policy>,
cursor: Option<cursor::Config>,
metrics: Option<prometheus::Config>,
}

impl ConfigRoot {
Expand All @@ -40,15 +41,21 @@ impl ConfigRoot {
}
}

struct Runtime {
source: Tether,
filters: Vec<Tether>,
sink: Tether,
cursor: Tether,
pub struct Runtime {
pub source: Tether,
pub filters: Vec<Tether>,
pub sink: Tether,
pub cursor: Tether,
}

impl Debug for Runtime {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Runtime").finish()
}
}

impl Runtime {
fn all_tethers(&self) -> impl Iterator<Item = &Tether> {
pub fn all_tethers(&self) -> impl Iterator<Item = &Tether> {
std::iter::once(&self.source)
.chain(self.filters.iter())
.chain(std::iter::once(&self.sink))
Expand Down Expand Up @@ -132,6 +139,24 @@ fn bootstrap(
Ok(runtime)
}

async fn monitor_loop(
runtime: &Arc<Runtime>,
metrics: Option<prometheus::Config>,
console: Option<super::console::Mode>,
) {
if let Some(metrics) = metrics {
info!("starting metrics exporter");
let runtime = runtime.clone();
tokio::spawn(async { prometheus::initialize(metrics, runtime).await });
}

while !runtime.should_stop() {
// TODO: move console refresh to it's own tokio thread
console::refresh(&console, runtime.all_tethers());
tokio::time::sleep(Duration::from_millis(1500)).await;
}
}

pub fn run(args: &Args) -> Result<(), Error> {
console::initialize(&args.console);

Expand Down Expand Up @@ -167,15 +192,20 @@ pub fn run(args: &Args) -> Result<(), Error> {

let retries = define_gasket_policy(config.retries.as_ref());
let runtime = bootstrap(source, filters, sink, cursor, retries)?;
let runtime = Arc::new(runtime);

let tokio_rt = tokio::runtime::Builder::new_multi_thread()
.enable_io()
.enable_time()
.build()
.unwrap();

info!("oura is running...");

while !runtime.should_stop() {
console::refresh(&args.console, runtime.all_tethers());
std::thread::sleep(Duration::from_millis(1500));
}
tokio_rt.block_on(async { monitor_loop(&runtime, config.metrics, args.console.clone()).await });

info!("Oura is stopping...");

runtime.shutdown();

Ok(())
Expand Down
1 change: 1 addition & 0 deletions src/bin/oura/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ use std::process;

mod console;
mod daemon;
mod prometheus;

#[derive(Parser)]
#[clap(name = "Oura")]
Expand Down
103 changes: 103 additions & 0 deletions src/bin/oura/prometheus.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
//! An utility to keep track of the progress of the pipeline as a whole

use anyhow::{Context, Result};
use gasket::{metrics::Reading, runtime::Tether};
use prometheus_exporter_base::{
prelude::{Authorization, ServerOptions},
render_prometheus, MetricType, PrometheusInstance, PrometheusMetric,
};
use serde::{Deserialize, Serialize};
use std::{
io::{BufWriter, Write},
net::SocketAddr,
sync::Arc,
};
use tracing::warn;

use crate::daemon::Runtime;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Config {
pub address: Option<String>,
}

pub fn render_tether(tether: &Tether, render: &mut impl Write) -> Result<()> {
let readings = tether.read_metrics()?;

let stage = tether.name();

let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)?
.as_millis();

for (name, value) in readings {
let full_name = format!("{stage}_{name}");

let metric_type = match value {
Reading::Count(_) => MetricType::Counter,
Reading::Gauge(_) => MetricType::Gauge,
Reading::Message(_) => MetricType::Summary,
};

let mut pc = PrometheusMetric::build()
.with_name(&full_name)
.with_help("some specific help")
.with_metric_type(metric_type)
.build();

match value {
Reading::Count(x) => {
pc.render_and_append_instance(
&PrometheusInstance::new().with_value(x).with_timestamp(ts),
);
}
Reading::Gauge(x) => {
pc.render_and_append_instance(
&PrometheusInstance::new().with_value(x).with_timestamp(ts),
);
}
Reading::Message(msg) => {
warn!(msg, "can't render message metrics to prometheous");
}
};

writeln!(render, "{}", pc.render())?;
}

Ok(())
}

fn render(runtime: &Runtime) -> Result<String> {
warn!("rendering");
let mut buf = BufWriter::new(Vec::new());

runtime
.all_tethers()
.try_for_each(|x| render_tether(x, &mut buf))?;

let bytes = buf.into_inner()?;
let string = String::from_utf8(bytes)?;

Ok(string)
}

pub async fn initialize(config: Config, runtime: Arc<Runtime>) -> Result<()> {
let addr: SocketAddr = config
.address
.as_deref()
.unwrap_or("0.0.0.0:9186")
.parse()
.context("parsing binding config")?;

let server_options = ServerOptions {
addr,
authorization: Authorization::None,
};

render_prometheus(server_options, runtime, |_, options| async move {
Ok(render(&options).unwrap())
})
.await;

Ok(())
}
15 changes: 14 additions & 1 deletion src/sources/n2c.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@

#[derive(Stage)]
#[stage(
name = "source-n2c",
name = "source",
unit = "NextResponse<BlockContent>",
worker = "Worker"
)]
Expand All @@ -33,6 +33,12 @@

#[metric]
chain_tip: gasket::metrics::Gauge,

#[metric]
current_slot: gasket::metrics::Gauge,

#[metric]
rollback_count: gasket::metrics::Counter,
}

async fn intersect_from_config(
Expand Down Expand Up @@ -104,6 +110,8 @@
stage.breadcrumbs.track(point.clone());

stage.chain_tip.set(tip.0.slot_or_default() as i64);
stage.current_slot.set(slot as i64);
stage.ops_count.inc(1);

Ok(())
}
Expand All @@ -122,6 +130,9 @@
stage.breadcrumbs.track(point.clone());

stage.chain_tip.set(tip.0.slot_or_default() as i64);
stage.current_slot.set(point.slot_or_default() as i64);
stage.rollback_count.inc(1);
stage.ops_count.inc(1);

Ok(())
}
Expand All @@ -138,7 +149,7 @@
async fn bootstrap(stage: &Stage) -> Result<Self, WorkerError> {
debug!("connecting");

let mut peer_session = NodeClient::connect(&stage.config.socket_path, stage.chain.magic)

Check failure on line 152 in src/sources/n2c.rs

View workflow job for this annotation

GitHub Actions / Check (windows-latest, stable)

no function or associated item named `connect` found for struct `NodeClient` in the current scope
.await
.or_retry()?;

Expand Down Expand Up @@ -203,6 +214,8 @@
output: Default::default(),
ops_count: Default::default(),
chain_tip: Default::default(),
current_slot: Default::default(),
rollback_count: Default::default(),
};

Ok(stage)
Expand Down
15 changes: 14 additions & 1 deletion src/sources/n2n.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ use crate::framework::*;

#[derive(Stage)]
#[stage(
name = "source-n2n",
name = "source",
unit = "NextResponse<HeaderContent>",
worker = "Worker"
)]
Expand All @@ -31,6 +31,12 @@ pub struct Stage {

#[metric]
chain_tip: gasket::metrics::Gauge,

#[metric]
current_slot: gasket::metrics::Gauge,

#[metric]
rollback_count: gasket::metrics::Counter,
}

fn to_traverse(header: &HeaderContent) -> Result<MultiEraHeader<'_>, WorkerError> {
Expand Down Expand Up @@ -121,6 +127,8 @@ impl Worker {
stage.breadcrumbs.track(point);

stage.chain_tip.set(tip.0.slot_or_default() as i64);
stage.current_slot.set(slot as i64);
stage.ops_count.inc(1);

Ok(())
}
Expand All @@ -139,6 +147,9 @@ impl Worker {
stage.breadcrumbs.track(point.clone());

stage.chain_tip.set(tip.0.slot_or_default() as i64);
stage.current_slot.set(point.slot_or_default() as i64);
stage.ops_count.inc(1);
stage.rollback_count.inc(1);

Ok(())
}
Expand Down Expand Up @@ -227,7 +238,9 @@ impl Config {
intersect: ctx.intersect.clone(),
output: Default::default(),
ops_count: Default::default(),
rollback_count: Default::default(),
chain_tip: Default::default(),
current_slot: Default::default(),
};

Ok(stage)
Expand Down
7 changes: 7 additions & 0 deletions src/sources/utxorpc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -230,11 +230,17 @@ pub struct Stage {
config: Config,
breadcrumbs: Breadcrumbs,
intersect: IntersectConfig,

pub output: SourceOutputPort,

#[metric]
ops_count: gasket::metrics::Counter,

#[metric]
chain_tip: gasket::metrics::Gauge,

#[metric]
current_slot: gasket::metrics::Gauge,
}

#[derive(Deserialize)]
Expand All @@ -252,6 +258,7 @@ impl Config {
output: Default::default(),
ops_count: Default::default(),
chain_tip: Default::default(),
current_slot: Default::default(),
};

Ok(stage)
Expand Down
Loading