Skip to content

Commit 4ab3845

Browse files
committed
refactor: use service-toolkit runtime helpers; drop local dupes
Replace hand-rolled copies with the shared toolkit modules (service-toolkit-rust `runtime` feature): - config: local env_* readers now delegate to service_toolkit_rust::env (names kept for the fallback-chain call sites in load) - main: tracing init -> telemetry::init_tracing - metrics: /metrics server -> metrics::serve (global registry) - wal: timestamp -> time::now_millis_u128 Drops now-unused deps: hyper, hyper-util, http-body-util, tracing-subscriber. Builds + clippy clean (debug & release).
1 parent 4bab33a commit 4ab3845

6 files changed

Lines changed: 35 additions & 118 deletions

File tree

Cargo.lock

Lines changed: 8 additions & 18 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -12,17 +12,13 @@ tokio = { version = "1", features = ["full"] }
1212
reqwest = { version = "0.12", features = ["json"] }
1313
tonic = "0.12"
1414
prost = "0.13"
15-
service-toolkit-rust = { git = "https://github.com/LeavePulse/service-toolkit.git", default-features = false, features = ["grpc"] }
15+
service-toolkit-rust = { git = "https://github.com/LeavePulse/service-toolkit.git", default-features = false, features = ["grpc", "runtime"] }
1616
serde = { version = "1", features = ["derive"] }
1717
serde_json = "1"
1818
redis = { version = "0.27", features = ["tokio-comp"] }
1919
prometheus = { version = "0.13", features = ["process"] }
20-
hyper = { version = "1", features = ["server", "http1"] }
21-
hyper-util = { version = "0.1", features = ["tokio"] }
22-
http-body-util = "0.1"
2320
sha2 = "0.10"
2421
tracing = "0.1"
25-
tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] }
2622
dotenvy = "0.15"
2723
rand = "0.8"
2824
chrono = { version = "0.4", features = ["serde"] }

src/config.rs

Lines changed: 14 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,10 @@
11
//! Application configuration loaded from environment variables + .env file.
2+
//!
3+
//! The typed env readers below are thin aliases over
4+
//! `service_toolkit_rust::env` so the parsing logic lives in one place; the
5+
//! local names are kept for the many fallback-chain call sites in `load`.
26
3-
use std::env;
7+
use service_toolkit_rust::env as cfg;
48

59
/// Server catalog API settings (server-service).
610
#[derive(Debug, Clone)]
@@ -96,47 +100,34 @@ pub struct Settings {
96100
pub redis: RedisSettings,
97101
}
98102

103+
// Thin aliases over the shared toolkit readers — kept so the fallback-chain
104+
// call sites in `load` (`.max(...)`, `or_else`) stay terse.
99105
fn env_str(key: &str, default: &str) -> String {
100-
env::var(key).unwrap_or_else(|_| default.to_string())
106+
cfg::str(key, default)
101107
}
102108

103109
fn env_str_opt(key: &str) -> Option<String> {
104-
env::var(key).ok().filter(|s| !s.is_empty())
110+
cfg::str_opt(key)
105111
}
106112

107113
fn env_u64(key: &str, default: u64) -> u64 {
108-
env::var(key)
109-
.ok()
110-
.and_then(|v| v.parse().ok())
111-
.unwrap_or(default)
114+
cfg::u64(key, default)
112115
}
113116

114117
fn env_usize(key: &str, default: usize) -> usize {
115-
env::var(key)
116-
.ok()
117-
.and_then(|v| v.parse().ok())
118-
.unwrap_or(default)
118+
cfg::usize(key, default)
119119
}
120120

121121
fn env_f64(key: &str, default: f64) -> f64 {
122-
env::var(key)
123-
.ok()
124-
.and_then(|v| v.parse().ok())
125-
.unwrap_or(default)
122+
cfg::f64(key, default)
126123
}
127124

128125
fn env_u16(key: &str, default: u16) -> u16 {
129-
env::var(key)
130-
.ok()
131-
.and_then(|v| v.parse().ok())
132-
.unwrap_or(default)
126+
cfg::u16(key, default)
133127
}
134128

135129
fn env_bool(key: &str, default: bool) -> bool {
136-
env::var(key)
137-
.ok()
138-
.map(|v| matches!(v.to_lowercase().as_str(), "true" | "1" | "yes" | "on"))
139-
.unwrap_or(default)
130+
cfg::bool(key, default)
140131
}
141132

142133
impl Settings {

src/main.rs

Lines changed: 1 addition & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -13,18 +13,12 @@ mod wal;
1313

1414
use config::Settings;
1515
use tracing::info;
16-
use tracing_subscriber::EnvFilter;
1716

1817
#[tokio::main]
1918
async fn main() {
2019
let settings = Settings::load();
2120

22-
// Initialize tracing.
23-
let filter = EnvFilter::try_new(&settings.log_level).unwrap_or_else(|_| EnvFilter::new("info"));
24-
tracing_subscriber::fmt()
25-
.with_env_filter(filter)
26-
.with_target(true)
27-
.init();
21+
service_toolkit_rust::telemetry::init_tracing(&settings.log_level);
2822

2923
info!(
3024
service = %settings.service_name,

src/metrics.rs

Lines changed: 9 additions & 60 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,16 @@
1-
//! Prometheus metrics definitions and HTTP exporter.
1+
//! Prometheus metrics definitions. The `/metrics` HTTP exporter itself is
2+
//! served by the shared `service_toolkit_rust::metrics` helper over the
3+
//! process-global default registry these `register_*!` macros write to.
24
35
use std::net::SocketAddr;
46
use std::sync::LazyLock;
57

6-
use http_body_util::Full;
7-
use hyper::body::Bytes;
8-
use hyper::service::service_fn;
9-
use hyper::{Request, Response, StatusCode};
10-
use hyper_util::rt::TokioIo;
118
use prometheus::{
12-
Encoder, Gauge, GaugeVec, Histogram, HistogramOpts, IntCounter, IntCounterVec, IntGauge, Opts,
13-
TextEncoder, register_gauge, register_gauge_vec, register_histogram, register_int_counter,
9+
Gauge, GaugeVec, Histogram, HistogramOpts, IntCounter, IntCounterVec, IntGauge, Opts,
10+
register_gauge, register_gauge_vec, register_histogram, register_int_counter,
1411
register_int_counter_vec, register_int_gauge,
1512
};
16-
use tokio::net::TcpListener;
17-
use tracing::{error, info};
13+
use tracing::info;
1814

1915
// ---------------------------------------------------------------------------
2016
// Polling
@@ -280,59 +276,12 @@ pub fn init() {
280276
LazyLock::force(&BUFFER_EVICTED);
281277
}
282278

283-
async fn handle_metrics(
284-
_req: Request<hyper::body::Incoming>,
285-
) -> Result<Response<Full<Bytes>>, hyper::Error> {
286-
let encoder = TextEncoder::new();
287-
let metric_families = prometheus::gather();
288-
let mut buffer = Vec::new();
289-
if encoder.encode(&metric_families, &mut buffer).is_err() {
290-
let resp = Response::builder()
291-
.status(StatusCode::INTERNAL_SERVER_ERROR)
292-
.body(Full::new(Bytes::from("Failed to encode metrics")))
293-
.unwrap();
294-
return Ok(resp);
295-
}
296-
let resp = Response::builder()
297-
.status(StatusCode::OK)
298-
.header("Content-Type", encoder.format_type())
299-
.body(Full::new(Bytes::from(buffer)))
300-
.unwrap();
301-
Ok(resp)
302-
}
303-
304-
/// Start the Prometheus metrics HTTP server on the given address.
279+
/// Start the Prometheus metrics HTTP server on the given address. Serves the
280+
/// process-global default registry via the shared toolkit helper.
305281
pub async fn start_metrics_server(host: &str, port: u16) {
306282
let addr: SocketAddr = format!("{host}:{port}")
307283
.parse()
308284
.unwrap_or_else(|_| SocketAddr::from(([0, 0, 0, 0], port)));
309-
310-
let listener = match TcpListener::bind(addr).await {
311-
Ok(l) => l,
312-
Err(e) => {
313-
error!("Failed to bind metrics server on {addr}: {e}");
314-
return;
315-
}
316-
};
317-
318285
info!("Prometheus metrics server listening on {addr}");
319-
320-
loop {
321-
let (stream, _) = match listener.accept().await {
322-
Ok(v) => v,
323-
Err(e) => {
324-
error!("Metrics server accept error: {e}");
325-
continue;
326-
}
327-
};
328-
let io = TokioIo::new(stream);
329-
tokio::spawn(async move {
330-
if let Err(e) = hyper::server::conn::http1::Builder::new()
331-
.serve_connection(io, service_fn(handle_metrics))
332-
.await
333-
{
334-
error!("Metrics connection error: {e}");
335-
}
336-
});
337-
}
286+
service_toolkit_rust::metrics::serve(addr).await;
338287
}

src/wal.rs

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -9,9 +9,9 @@
99
use std::fs::{self, File, OpenOptions};
1010
use std::io::{self, BufReader, Read, Write};
1111
use std::path::{Path, PathBuf};
12-
use std::time::{SystemTime, UNIX_EPOCH};
1312

1413
use serde_json::Value;
14+
use service_toolkit_rust::time::now_millis_u128;
1515
use tracing::{info, warn};
1616

1717
use crate::metrics::{
@@ -162,10 +162,7 @@ impl WalBuffer {
162162
// Close current segment (drop flushes).
163163
self.current_segment = None;
164164

165-
let ts = SystemTime::now()
166-
.duration_since(UNIX_EPOCH)
167-
.unwrap_or_default()
168-
.as_millis();
165+
let ts = now_millis_u128();
169166
let path = self.dir.join(format!("buffer-{ts}.wal"));
170167

171168
let file = OpenOptions::new().create(true).append(true).open(&path)?;

0 commit comments

Comments
 (0)