Skip to content

Commit 3ea0d58

Browse files
zzylolclaude
andcommitted
style: apply cargo fmt
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent 8df4927 commit 3ea0d58

3 files changed

Lines changed: 25 additions & 28 deletions

File tree

‎asap-query-engine/src/bin/test_e2e_precompute.rs‎

Lines changed: 8 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -353,7 +353,10 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
353353
}
354354
}
355355
let batch_body = build_remote_write_body(batch_timeseries);
356-
println!(" Payload size: {} bytes (snappy-compressed)", batch_body.len());
356+
println!(
357+
" Payload size: {} bytes (snappy-compressed)",
358+
batch_body.len()
359+
);
357360

358361
let t0 = std::time::Instant::now();
359362
let resp = client
@@ -374,12 +377,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
374377
tokio::time::sleep(tokio::time::Duration::from_millis(500)).await;
375378

376379
// Verify samples landed in the store
377-
let batch_results = store.query_precomputed_output(
378-
"fake_metric",
379-
raw_agg_id,
380-
200_000,
381-
210_000,
382-
)?;
380+
let batch_results =
381+
store.query_precomputed_output("fake_metric", raw_agg_id, 200_000, 210_000)?;
383382
let batch_buckets: usize = batch_results.values().map(|v| v.len()).sum();
384383
println!(" Stored {batch_buckets} buckets from batch (expected 1000)");
385384
assert!(
@@ -438,12 +437,8 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
438437
let drain_deadline = std::time::Instant::now() + std::time::Duration::from_secs(60);
439438
let mut tp_buckets: usize;
440439
loop {
441-
let tp_results = store.query_precomputed_output(
442-
"fake_metric",
443-
raw_agg_id,
444-
300_000,
445-
max_ts,
446-
)?;
440+
let tp_results =
441+
store.query_precomputed_output("fake_metric", raw_agg_id, 300_000, max_ts)?;
447442
tp_buckets = tp_results.values().map(|v| v.len()).sum();
448443
if tp_buckets as u64 >= total_samples || std::time::Instant::now() > drain_deadline {
449444
break;

‎asap-query-engine/src/precompute_engine/mod.rs‎

Lines changed: 9 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -15,9 +15,9 @@ use crate::precompute_engine::worker::Worker;
1515
use axum::{body::Bytes, extract::State, http::StatusCode, routing::post, Router};
1616
use std::collections::HashMap;
1717
use std::sync::Arc;
18+
use std::time::Instant;
1819
use tokio::net::TcpListener;
1920
use tokio::sync::mpsc;
20-
use std::time::Instant;
2121
use tracing::{debug_span, info, warn, Instrument};
2222

2323
/// Shared state for the ingest HTTP handler.
@@ -67,10 +67,8 @@ impl PrecomputeEngine {
6767
let router = SeriesRouter::new(senders);
6868

6969
// Build aggregation config map from streaming config
70-
let agg_configs: HashMap<u64, _> = self
71-
.streaming_config
72-
.get_all_aggregation_configs()
73-
.clone();
70+
let agg_configs: HashMap<u64, _> =
71+
self.streaming_config.get_all_aggregation_configs().clone();
7472

7573
// Spawn workers
7674
let mut worker_handles = Vec::with_capacity(num_workers);
@@ -138,10 +136,7 @@ impl PrecomputeEngine {
138136
}
139137

140138
/// Axum handler for Prometheus remote write.
141-
async fn handle_ingest(
142-
State(state): State<Arc<IngestState>>,
143-
body: Bytes,
144-
) -> StatusCode {
139+
async fn handle_ingest(State(state): State<Arc<IngestState>>, body: Bytes) -> StatusCode {
145140
let ingest_span = debug_span!("ingest", body_len = body.len());
146141
let ingest_received_at = Instant::now();
147142

@@ -179,7 +174,11 @@ async fn handle_ingest(
179174
.collect();
180175

181176
// Route all series to workers concurrently
182-
if let Err(e) = state.router.route_batch(by_series_owned, ingest_received_at).await {
177+
if let Err(e) = state
178+
.router
179+
.route_batch(by_series_owned, ingest_received_at)
180+
.await
181+
{
183182
warn!("Batch routing error: {}", e);
184183
return StatusCode::INTERNAL_SERVER_ERROR;
185184
}

‎asap-query-engine/src/precompute_engine/series_router.rs‎

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
1+
use futures::future::try_join_all;
12
use std::collections::HashMap;
23
use std::time::Instant;
3-
use futures::future::try_join_all;
44
use tokio::sync::mpsc;
55
use xxhash_rust::xxh64::xxh64;
66

@@ -82,15 +82,18 @@ impl SeriesRouter {
8282
let sender = &self.senders[worker_idx];
8383
async move {
8484
for msg in messages {
85-
sender.send(msg).await.map_err(|e| {
86-
format!("Failed to send to worker {}: {}", worker_idx, e)
87-
})?;
85+
sender
86+
.send(msg)
87+
.await
88+
.map_err(|e| format!("Failed to send to worker {}: {}", worker_idx, e))?;
8889
}
8990
Ok::<(), String>(())
9091
}
9192
}))
9293
.await
93-
.map_err(|e| -> Box<dyn std::error::Error + Send + Sync> { Box::new(std::io::Error::new(std::io::ErrorKind::Other, e)) })?;
94+
.map_err(|e| -> Box<dyn std::error::Error + Send + Sync> {
95+
Box::new(std::io::Error::new(std::io::ErrorKind::Other, e))
96+
})?;
9497

9598
Ok(())
9699
}

0 commit comments

Comments
 (0)