1pub mod admin;
6pub mod config;
7pub mod consumer;
8pub mod handlers;
9pub mod histogram_relay;
10pub mod metrics;
11pub mod middleware;
12pub mod peers;
13pub mod prom_to_mimir;
14pub mod remote_write;
15
16#[macro_export]
21macro_rules! var {
22 ($key:expr) => {
23 match std::env::var($key) {
24 Ok(val) => val,
25 Err(_) => "".into(),
26 }
27 };
28 ($key:expr, $default:expr) => {
29 match std::env::var($key) {
30 Ok(val) => val.parse::<_>().unwrap(),
31 Err(_) => $default,
32 }
33 };
34}
35
36#[cfg(test)]
37mod tests {
38 use std::{net::TcpListener, time::Duration};
39
40 use axum::{Router, http::StatusCode, routing::post};
41 use iota_tls::{ClientCertVerifier, TlsAcceptor};
42 use prometheus_filtered::{Encoder, PROTOBUF_FORMAT};
43
44 use super::*;
45 use crate::{
46 admin::{CertKeyPair, Labels},
47 config::RemoteWriteConfig,
48 histogram_relay::HistogramRelay,
49 peers::IotaNodeProvider,
50 prom_to_mimir::tests::*,
51 };
52
53 async fn run_dummy_remote_write(listener: TcpListener) {
54 async fn handler() -> StatusCode {
56 StatusCode::OK
57 }
58
59 let app = Router::new().route("/v1/push", post(handler));
61
62 listener.set_nonblocking(true).unwrap();
64 let listener = tokio::net::TcpListener::from_std(listener).unwrap();
65 axum::serve(listener, app).await.unwrap();
66 }
67
68 async fn run_dummy_remote_write_very_slow(listener: TcpListener) {
69 async fn handler() -> StatusCode {
74 tokio::time::sleep(Duration::from_secs(60)).await; StatusCode::OK
78 }
79
80 let app = Router::new().route("/v1/push", post(handler));
82
83 listener.set_nonblocking(true).unwrap();
85 let listener = tokio::net::TcpListener::from_std(listener).unwrap();
86 axum::serve(listener, app).await.unwrap();
87 }
88
89 #[tokio::test]
96 async fn test_axum_acceptor() {
97 let CertKeyPair(client_priv_cert, client_pub_key) =
99 admin::generate_self_cert("iota".into());
100 let CertKeyPair(server_priv_cert, _) = admin::generate_self_cert("localhost".into());
101
102 let dummy_remote_write_listener = std::net::TcpListener::bind("localhost:0").unwrap();
104 let dummy_remote_write_address = dummy_remote_write_listener.local_addr().unwrap();
105 let dummy_remote_write_url = format!(
106 "http://localhost:{}/v1/push",
107 dummy_remote_write_address.port()
108 );
109
110 let _dummy_remote_write =
111 tokio::spawn(async move { run_dummy_remote_write(dummy_remote_write_listener).await });
112
113 let mut allower = IotaNodeProvider::new("".into(), Duration::from_secs(30), vec![]);
115 let tls_config = ClientCertVerifier::new(
116 allower.clone(),
117 iota_tls::IOTA_VALIDATOR_SERVER_NAME.to_string(),
118 )
119 .rustls_server_config(
120 vec![server_priv_cert.rustls_certificate()],
121 server_priv_cert.rustls_private_key(),
122 )
123 .unwrap();
124
125 let client = admin::make_reqwest_client(
126 RemoteWriteConfig {
127 url: dummy_remote_write_url.to_owned(),
128 username: "bar".into(),
129 password: "foo".into(),
130 ..Default::default()
131 },
132 "dummy user agent",
133 );
134
135 let app = admin::app(
136 Labels {
137 network: "unittest-network".into(),
138 },
139 client,
140 HistogramRelay::new(),
141 Some(allower.clone()),
142 None,
143 );
144
145 let listener = std::net::TcpListener::bind("localhost:0").unwrap();
146 let server_address = listener.local_addr().unwrap();
147 let server_url = format!(
148 "https://localhost:{}/publish/metrics",
149 server_address.port()
150 );
151
152 let acceptor = TlsAcceptor::new(tls_config);
153 let _server = tokio::spawn(async move {
154 admin::server(listener, app, Some(acceptor)).await.unwrap();
155 });
156
157 let _ = rustls::crypto::ring::default_provider().install_default();
159 let client = reqwest::Client::builder()
160 .tls_certs_only([server_priv_cert.reqwest_certificate()])
161 .identity(client_priv_cert.reqwest_identity())
162 .https_only(true)
163 .build()
164 .unwrap();
165
166 client.get(&server_url).send().await.unwrap_err();
168
169 allower.get_mut().write().unwrap().insert(
172 client_pub_key.to_owned(),
173 peers::AllowedPeer {
174 name: "some-node".into(),
175 public_key: client_pub_key.to_owned(),
176 },
177 );
178
179 let mf = create_metric_family(
180 "foo_metric",
181 "some help this is",
182 None,
183 vec![create_metric_counter(
184 create_labels(vec![("some", "label")]),
185 create_counter(2046.0),
186 )],
187 );
188
189 let mut buf = vec![];
190 let encoder = prometheus_filtered::ProtobufEncoder::new();
191 encoder.encode(&[mf], &mut buf).unwrap();
192
193 let res = client
194 .post(&server_url)
195 .header(reqwest::header::CONTENT_TYPE, PROTOBUF_FORMAT)
196 .body(buf)
197 .send()
198 .await
199 .expect("expected a successful post with a self-signed certificate");
200 let status = res.status();
201 let body = res.text().await.unwrap();
202 assert_eq!("created", body);
203 assert_eq!(status, reqwest::StatusCode::CREATED);
204 }
205
206 #[tokio::test]
208 async fn test_client_timeout() {
209 let CertKeyPair(client_priv_cert, client_pub_key) =
211 admin::generate_self_cert("iota".into());
212 let CertKeyPair(server_priv_cert, _) = admin::generate_self_cert("localhost".into());
213
214 let dummy_remote_write_listener = std::net::TcpListener::bind("localhost:0").unwrap();
216 let dummy_remote_write_address = dummy_remote_write_listener.local_addr().unwrap();
217 let dummy_remote_write_url = format!(
218 "http://localhost:{}/v1/push",
219 dummy_remote_write_address.port()
220 );
221
222 let _dummy_remote_write = tokio::spawn(async move {
223 run_dummy_remote_write_very_slow(dummy_remote_write_listener).await
224 });
225
226 let mut allower = IotaNodeProvider::new("".into(), Duration::from_secs(30), vec![]);
228 let tls_config = ClientCertVerifier::new(
229 allower.clone(),
230 iota_tls::IOTA_VALIDATOR_SERVER_NAME.to_string(),
231 )
232 .rustls_server_config(
233 vec![server_priv_cert.rustls_certificate()],
234 server_priv_cert.rustls_private_key(),
235 )
236 .unwrap();
237
238 let client = admin::make_reqwest_client(
239 RemoteWriteConfig {
240 url: dummy_remote_write_url.to_owned(),
241 username: "bar".into(),
242 password: "foo".into(),
243 ..Default::default()
244 },
245 "dummy user agent",
246 );
247
248 let timeout_secs = Some(2u64);
249
250 let app = admin::app(
251 Labels {
252 network: "unittest-network".into(),
253 },
254 client,
255 HistogramRelay::new(),
256 Some(allower.clone()),
257 timeout_secs,
258 );
259
260 let listener = std::net::TcpListener::bind("localhost:0").unwrap();
261 let server_address = listener.local_addr().unwrap();
262 let server_url = format!(
263 "https://localhost:{}/publish/metrics",
264 server_address.port()
265 );
266
267 let acceptor = TlsAcceptor::new(tls_config);
268 let _server = tokio::spawn(async move {
269 admin::server(listener, app, Some(acceptor)).await.unwrap();
270 });
271
272 let _ = rustls::crypto::ring::default_provider().install_default();
274 let client = reqwest::Client::builder()
275 .tls_certs_only([server_priv_cert.reqwest_certificate()])
276 .identity(client_priv_cert.reqwest_identity())
277 .https_only(true)
278 .build()
279 .unwrap();
280
281 client.get(&server_url).send().await.unwrap_err();
283
284 allower.get_mut().write().unwrap().insert(
287 client_pub_key.to_owned(),
288 peers::AllowedPeer {
289 name: "some-node".into(),
290 public_key: client_pub_key.to_owned(),
291 },
292 );
293
294 let mf = create_metric_family(
295 "foo_metric",
296 "some help this is",
297 None,
298 vec![create_metric_counter(
299 create_labels(vec![("some", "label")]),
300 create_counter(2046.0),
301 )],
302 );
303
304 let mut buf = vec![];
305 let encoder = prometheus_filtered::ProtobufEncoder::new();
306 encoder.encode(&[mf], &mut buf).unwrap();
307
308 let res = client
309 .post(&server_url)
310 .header(reqwest::header::CONTENT_TYPE, PROTOBUF_FORMAT)
311 .body(buf)
312 .send()
313 .await
314 .expect("expected a successful post with a self-signed certificate");
315 let status = res.status();
316 assert_eq!(status, StatusCode::REQUEST_TIMEOUT);
317 }
318}