Skip to content

Commit d9282fa

Browse files
refactor: added explicit trait for sources in asap-query-engine/precompute_engine (#302)
* refactor: added explicit trait for sources in asap-query-engine/precompute_engine * added missing comments
1 parent b3d98c9 commit d9282fa

13 files changed

Lines changed: 297 additions & 190 deletions

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

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,9 @@ use query_engine_rust::drivers::ingest::prometheus_remote_write::{
1010
};
1111
use query_engine_rust::precompute_engine::config::{LateDataPolicy, PrecomputeEngineConfig};
1212
use query_engine_rust::precompute_engine::output_sink::OutputSink;
13-
use query_engine_rust::precompute_engine::PrecomputeEngine;
13+
use query_engine_rust::precompute_engine::{
14+
HttpIngestConfig, HttpIngestSource, IngestSource, PrecomputeEngine,
15+
};
1416
use query_engine_rust::stores::{SimpleMapStore, Store};
1517
use std::collections::HashMap;
1618
use std::sync::atomic::{AtomicU64, Ordering};
@@ -166,7 +168,6 @@ async fn start_engine(
166168
) {
167169
let config = PrecomputeEngineConfig {
168170
num_workers: workers,
169-
ingest_port: port,
170171
allowed_lateness_ms: 5_000,
171172
max_buffer_per_series: 100_000,
172173
flush_interval_ms: 100,
@@ -175,7 +176,9 @@ async fn start_engine(
175176
raw_mode_aggregation_id: 0,
176177
late_data_policy: LateDataPolicy::Drop,
177178
};
178-
let engine = PrecomputeEngine::new(config, streaming_config, sink);
179+
let sources: Vec<Box<dyn IngestSource>> =
180+
vec![Box::new(HttpIngestSource::new(HttpIngestConfig { port }))];
181+
let engine = PrecomputeEngine::new(config, streaming_config, sink, sources);
179182
tokio::spawn(async move {
180183
if let Err(err) = engine.run().await {
181184
eprintln!("precompute engine on port {port} failed: {err}");

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

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,9 @@ use query_engine_rust::drivers::ingest::prometheus_remote_write::{
1717
};
1818
use query_engine_rust::precompute_engine::config::{LateDataPolicy, PrecomputeEngineConfig};
1919
use query_engine_rust::precompute_engine::output_sink::StoreOutputSink;
20-
use query_engine_rust::precompute_engine::PrecomputeEngine;
20+
use query_engine_rust::precompute_engine::{
21+
HttpIngestConfig, HttpIngestSource, IngestSource, PrecomputeEngine,
22+
};
2123
use query_engine_rust::stores::{SimpleMapStore, Store};
2224
use std::collections::HashMap;
2325
use std::sync::Arc;
@@ -225,7 +227,6 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
225227

226228
let engine_config = PrecomputeEngineConfig {
227229
num_workers: NUM_WORKERS,
228-
ingest_port: INGEST_PORT,
229230
allowed_lateness_ms: 5_000,
230231
max_buffer_per_series: 10_000,
231232
flush_interval_ms: 1_000,
@@ -235,7 +236,11 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
235236
late_data_policy: LateDataPolicy::Drop,
236237
};
237238
let output_sink = Arc::new(StoreOutputSink::new(store.clone()));
238-
let engine = PrecomputeEngine::new(engine_config, streaming_config, output_sink);
239+
let sources: Vec<Box<dyn IngestSource>> =
240+
vec![Box::new(HttpIngestSource::new(HttpIngestConfig {
241+
port: INGEST_PORT,
242+
}))];
243+
let engine = PrecomputeEngine::new(engine_config, streaming_config, output_sink, sources);
239244
tokio::spawn(async move {
240245
if let Err(e) = engine.run().await {
241246
eprintln!("Precompute engine error: {e}");

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

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ use query_engine_rust::drivers::query::adapters::AdapterConfig;
77
use query_engine_rust::engines::SimpleEngine;
88
use query_engine_rust::precompute_engine::config::{LateDataPolicy, PrecomputeEngineConfig};
99
use query_engine_rust::precompute_engine::output_sink::{RawPassthroughSink, StoreOutputSink};
10-
use query_engine_rust::precompute_engine::PrecomputeEngine;
10+
use query_engine_rust::precompute_engine::{HttpIngestConfig, HttpIngestSource, PrecomputeEngine};
1111
use query_engine_rust::stores::SimpleMapStore;
1212
use query_engine_rust::{HttpServer, HttpServerConfig};
1313
use std::sync::Arc;
@@ -128,7 +128,6 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
128128
// Build the precompute engine config
129129
let engine_config = PrecomputeEngineConfig {
130130
num_workers: args.num_workers,
131-
ingest_port: args.ingest_port,
132131
allowed_lateness_ms: args.allowed_lateness_ms,
133132
max_buffer_per_series: args.max_buffer_per_series,
134133
flush_interval_ms: args.flush_interval_ms,
@@ -146,8 +145,13 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
146145
Arc::new(StoreOutputSink::new(store))
147146
};
148147

148+
let sources: Vec<Box<dyn query_engine_rust::precompute_engine::IngestSource>> =
149+
vec![Box::new(HttpIngestSource::new(HttpIngestConfig {
150+
port: args.ingest_port,
151+
}))];
152+
149153
// Build and run the engine
150-
let engine = PrecomputeEngine::new(engine_config, streaming_config, output_sink);
154+
let engine = PrecomputeEngine::new(engine_config, streaming_config, output_sink, sources);
151155

152156
info!("Starting precompute engine...");
153157
engine.run().await?;

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

Lines changed: 26 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,9 @@ use query_engine_rust::precompute_engine::config::{LateDataPolicy, PrecomputeEng
2121
use query_engine_rust::precompute_engine::output_sink::{
2222
NoopOutputSink, RawPassthroughSink, StoreOutputSink,
2323
};
24-
use query_engine_rust::precompute_engine::PrecomputeEngine;
24+
use query_engine_rust::precompute_engine::{
25+
HttpIngestConfig, HttpIngestSource, IngestSource, PrecomputeEngine,
26+
};
2527
use query_engine_rust::stores::SimpleMapStore;
2628
use query_engine_rust::utils::file_io::{read_inference_config, read_streaming_config};
2729
use query_engine_rust::{HttpServer, HttpServerConfig};
@@ -142,7 +144,6 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
142144
// Start precompute engine
143145
let engine_config = PrecomputeEngineConfig {
144146
num_workers: 2,
145-
ingest_port: INGEST_PORT,
146147
allowed_lateness_ms: 5000,
147148
max_buffer_per_series: 10000,
148149
flush_interval_ms: 200,
@@ -152,7 +153,16 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
152153
late_data_policy: LateDataPolicy::Drop,
153154
};
154155
let output_sink = Arc::new(StoreOutputSink::new(store.clone()));
155-
let engine = PrecomputeEngine::new(engine_config, streaming_config.clone(), output_sink);
156+
let sources: Vec<Box<dyn IngestSource>> =
157+
vec![Box::new(HttpIngestSource::new(HttpIngestConfig {
158+
port: INGEST_PORT,
159+
}))];
160+
let engine = PrecomputeEngine::new(
161+
engine_config,
162+
streaming_config.clone(),
163+
output_sink,
164+
sources,
165+
);
156166
tokio::spawn(async move {
157167
if let Err(e) = engine.run().await {
158168
eprintln!("Precompute engine error: {e}");
@@ -282,7 +292,6 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
282292
let raw_agg_id: u64 = 1;
283293
let raw_engine_config = PrecomputeEngineConfig {
284294
num_workers: 4,
285-
ingest_port: RAW_INGEST_PORT,
286295
allowed_lateness_ms: 5000,
287296
max_buffer_per_series: 10000,
288297
flush_interval_ms: 200,
@@ -292,7 +301,16 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
292301
late_data_policy: LateDataPolicy::Drop,
293302
};
294303
let raw_sink = Arc::new(RawPassthroughSink::new(store.clone()));
295-
let raw_engine = PrecomputeEngine::new(raw_engine_config, streaming_config.clone(), raw_sink);
304+
let raw_sources: Vec<Box<dyn IngestSource>> =
305+
vec![Box::new(HttpIngestSource::new(HttpIngestConfig {
306+
port: RAW_INGEST_PORT,
307+
}))];
308+
let raw_engine = PrecomputeEngine::new(
309+
raw_engine_config,
310+
streaming_config.clone(),
311+
raw_sink,
312+
raw_sources,
313+
);
296314
tokio::spawn(async move {
297315
if let Err(e) = raw_engine.run().await {
298316
eprintln!("Raw precompute engine error: {e}");
@@ -630,7 +648,6 @@ async fn run_single_bench(
630648
let noop_sink = Arc::new(NoopOutputSink::new());
631649
let engine_config = PrecomputeEngineConfig {
632650
num_workers,
633-
ingest_port: port,
634651
allowed_lateness_ms: 5000,
635652
max_buffer_per_series: 100_000,
636653
flush_interval_ms: 100,
@@ -639,7 +656,9 @@ async fn run_single_bench(
639656
raw_mode_aggregation_id: 0,
640657
late_data_policy: LateDataPolicy::Drop,
641658
};
642-
let engine = PrecomputeEngine::new(engine_config, streaming_config, noop_sink.clone());
659+
let sources: Vec<Box<dyn IngestSource>> =
660+
vec![Box::new(HttpIngestSource::new(HttpIngestConfig { port }))];
661+
let engine = PrecomputeEngine::new(engine_config, streaming_config, noop_sink.clone(), sources);
643662
tokio::spawn(async move {
644663
if let Err(e) = engine.run().await {
645664
eprintln!("Bench engine error: {e}");

‎asap-query-engine/src/lib.rs‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,10 @@ pub use drivers::{
4747

4848
pub use precompute_engine::config::{LateDataPolicy, PrecomputeEngineConfig};
4949
pub use precompute_engine::output_sink::StoreOutputSink;
50-
pub use precompute_engine::{PrecomputeEngine, PrecomputeEngineHandle};
50+
pub use precompute_engine::{
51+
HttpIngestConfig, HttpIngestSource, IngestContext, IngestSource, PrecomputeEngine,
52+
PrecomputeEngineHandle,
53+
};
5154

5255
pub use query_tracker::{QueryTracker, QueryTrackerConfig};
5356

‎asap-query-engine/src/main.rs‎

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -17,9 +17,10 @@ use query_engine_rust::precompute_engine::PrecomputeWorkerDiagnostics;
1717
use query_engine_rust::utils::file_io::{read_inference_config, read_streaming_config};
1818
use query_engine_rust::InferenceConfig;
1919
use query_engine_rust::{
20-
HttpServer, HttpServerConfig, KafkaConsumer, KafkaConsumerConfig, OtlpReceiver,
21-
OtlpReceiverConfig, PrecomputeEngine, PrecomputeEngineConfig, PrecomputeEngineHandle, Result,
22-
SimpleEngine, SimpleMapStore, StoreOutputSink,
20+
HttpIngestConfig, HttpIngestSource, HttpServer, HttpServerConfig, IngestSource, KafkaConsumer,
21+
KafkaConsumerConfig, OtlpReceiver, OtlpReceiverConfig, PrecomputeEngine,
22+
PrecomputeEngineConfig, PrecomputeEngineHandle, Result, SimpleEngine, SimpleMapStore,
23+
StoreOutputSink,
2324
};
2425

2526
#[derive(Parser, Debug)]
@@ -342,7 +343,6 @@ async fn main() -> Result<()> {
342343
let precompute_handle = if enable_precompute {
343344
let precompute_config = PrecomputeEngineConfig {
344345
num_workers: args.precompute_num_workers,
345-
ingest_port: args.prometheus_remote_write_port,
346346
allowed_lateness_ms: args.precompute_allowed_lateness_ms,
347347
max_buffer_per_series: args.precompute_max_buffer_per_series,
348348
flush_interval_ms: args.precompute_flush_interval_ms,
@@ -352,7 +352,16 @@ async fn main() -> Result<()> {
352352
late_data_policy: LateDataPolicy::Drop,
353353
};
354354
let output_sink = Arc::new(StoreOutputSink::new(store.clone()));
355-
let pe = PrecomputeEngine::new(precompute_config, streaming_config.clone(), output_sink);
355+
let sources: Vec<Box<dyn IngestSource>> =
356+
vec![Box::new(HttpIngestSource::new(HttpIngestConfig {
357+
port: args.prometheus_remote_write_port,
358+
}))];
359+
let pe = PrecomputeEngine::new(
360+
precompute_config,
361+
streaming_config.clone(),
362+
output_sink,
363+
sources,
364+
);
356365
let worker_diagnostics = pe.diagnostics();
357366
// Extract the handle before run() consumes the engine.
358367
pe_engine_handle = Some(pe.handle());

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

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -14,8 +14,6 @@ pub enum LateDataPolicy {
1414
pub struct PrecomputeEngineConfig {
1515
/// Number of worker threads for parallel processing.
1616
pub num_workers: usize,
17-
/// Port for the Prometheus remote write ingest endpoint.
18-
pub ingest_port: u16,
1917
/// Maximum allowed lateness for out-of-order samples (milliseconds).
2018
/// Samples arriving later than this behind the watermark are dropped.
2119
pub allowed_lateness_ms: i64,
@@ -38,7 +36,6 @@ impl Default for PrecomputeEngineConfig {
3836
fn default() -> Self {
3937
Self {
4038
num_workers: 4,
41-
ingest_port: 9090,
4239
allowed_lateness_ms: 5_000,
4340
max_buffer_per_series: 10_000,
4441
flush_interval_ms: 1_000,
@@ -58,7 +55,6 @@ mod tests {
5855
fn test_default_config() {
5956
let config = PrecomputeEngineConfig::default();
6057
assert_eq!(config.num_workers, 4);
61-
assert_eq!(config.ingest_port, 9090);
6258
assert_eq!(config.allowed_lateness_ms, 5_000);
6359
assert_eq!(config.max_buffer_per_series, 10_000);
6460
assert_eq!(config.flush_interval_ms, 1_000);

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

Lines changed: 32 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,14 @@
11
use crate::data_model::StreamingConfig;
22
use crate::precompute_engine::config::PrecomputeEngineConfig;
3-
use crate::precompute_engine::ingest_handler::{
4-
handle_prometheus_ingest, handle_victoriametrics_ingest, IngestState,
5-
};
3+
use crate::precompute_engine::ingest_source::{IngestContext, IngestSource};
64
use crate::precompute_engine::output_sink::OutputSink;
75
use crate::precompute_engine::series_router::{SeriesRouter, WorkerMessage};
86
use crate::precompute_engine::worker::{Worker, WorkerRuntimeConfig};
97
use arc_swap::ArcSwap;
108
use asap_types::aggregation_config::AggregationConfig;
11-
use axum::{routing::post, Router};
129
use std::collections::HashMap;
1310
use std::sync::atomic::{AtomicI64, AtomicUsize};
1411
use std::sync::Arc;
15-
use tokio::net::TcpListener;
1612
use tokio::sync::mpsc;
1713
use tracing::{info, warn};
1814

@@ -64,14 +60,15 @@ impl PrecomputeEngineHandle {
6460

6561
/// The top-level precompute engine orchestrator.
6662
///
67-
/// Creates worker threads, the series router, and the Axum ingest server.
63+
/// Creates worker threads and drives all registered ingest sources.
6864
/// Call `handle()` before `run()` to obtain a `PrecomputeEngineHandle` for
6965
/// applying runtime config updates while the engine is running.
7066
pub struct PrecomputeEngine {
7167
config: PrecomputeEngineConfig,
7268
streaming_config: Arc<StreamingConfig>,
7369
output_sink: Arc<dyn OutputSink>,
7470
diagnostics: Arc<PrecomputeWorkerDiagnostics>,
71+
sources: Vec<Box<dyn IngestSource>>,
7572
/// Channels created at construction so handle() can be extracted before run().
7673
senders: Vec<mpsc::Sender<WorkerMessage>>,
7774
receivers: Option<Vec<mpsc::Receiver<WorkerMessage>>>,
@@ -84,6 +81,7 @@ impl PrecomputeEngine {
8481
config: PrecomputeEngineConfig,
8582
streaming_config: Arc<StreamingConfig>,
8683
output_sink: Arc<dyn OutputSink>,
84+
sources: Vec<Box<dyn IngestSource>>,
8785
) -> Self {
8886
let worker_group_counts = (0..config.num_workers)
8987
.map(|_| Arc::new(AtomicUsize::new(0)))
@@ -96,8 +94,6 @@ impl PrecomputeEngine {
9694
worker_watermarks,
9795
});
9896

99-
// Build channels and initial agg_configs at construction time so that
100-
// handle() can be called before run().
10197
let channel_size = config.channel_buffer_size;
10298
let mut senders = Vec::with_capacity(config.num_workers);
10399
let mut receivers = Vec::with_capacity(config.num_workers);
@@ -119,6 +115,7 @@ impl PrecomputeEngine {
119115
streaming_config,
120116
output_sink,
121117
diagnostics,
118+
sources,
122119
senders,
123120
receivers: Some(receivers),
124121
ingest_agg_configs,
@@ -139,8 +136,8 @@ impl PrecomputeEngine {
139136
}
140137
}
141138

142-
/// Start the precompute engine. This spawns worker tasks and the HTTP
143-
/// ingest server, then blocks until shutdown.
139+
/// Start the precompute engine. This spawns worker tasks and all registered
140+
/// ingest sources, then blocks until shutdown.
144141
pub async fn run(mut self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
145142
let num_workers = self.config.num_workers;
146143

@@ -185,48 +182,49 @@ impl PrecomputeEngine {
185182
worker_handles.push(handle);
186183
}
187184

188-
info!(
189-
"PrecomputeEngine started with {} workers on port {}",
190-
num_workers, self.config.ingest_port
191-
);
185+
info!("PrecomputeEngine started with {} workers", num_workers);
192186

193-
// Build the ingest state, sharing the same Arc<ArcSwap> as the handle so
194-
// that PrecomputeEngineHandle::update_streaming_config swaps are visible here.
195-
let ingest_state = Arc::new(IngestState {
196-
router,
197-
samples_ingested: std::sync::atomic::AtomicU64::new(0),
187+
// Build the ingest context shared by all sources.
188+
// The ArcSwap pointer is the same one held by PrecomputeEngineHandle, so
189+
// handle.update_streaming_config() is immediately visible to every source.
190+
let ctx = IngestContext {
191+
router: router.clone(),
198192
agg_configs: self.ingest_agg_configs.clone(),
199193
pass_raw_samples: self.config.pass_raw_samples,
200-
});
194+
};
201195

202-
// Start flush timer
203-
let flush_state = ingest_state.clone();
196+
// Flush timer: periodically signal workers to close idle windows.
197+
let flush_router = router.clone();
204198
let flush_interval_ms = self.config.flush_interval_ms;
205199
tokio::spawn(async move {
206200
let mut interval =
207201
tokio::time::interval(tokio::time::Duration::from_millis(flush_interval_ms));
208202
loop {
209203
interval.tick().await;
210-
if let Err(e) = flush_state.router.broadcast_flush().await {
204+
if let Err(e) = flush_router.broadcast_flush().await {
211205
warn!("Flush broadcast error: {}", e);
212206
break;
213207
}
214208
}
215209
});
216210

217-
// Start the Axum HTTP server for ingest (Prometheus + VictoriaMetrics)
218-
let app = Router::new()
219-
.route("/api/v1/write", post(handle_prometheus_ingest))
220-
.route("/api/v1/import", post(handle_victoriametrics_ingest))
221-
.with_state(ingest_state);
222-
223-
let addr = format!("0.0.0.0:{}", self.config.ingest_port);
224-
info!("Ingest server listening on {}", addr);
211+
// Spawn each ingest source.
212+
let mut source_handles = Vec::with_capacity(self.sources.len());
213+
for source in self.sources {
214+
let ctx = ctx.clone();
215+
let handle = tokio::spawn(async move {
216+
if let Err(e) = source.run(ctx).await {
217+
warn!("Ingest source error: {}", e);
218+
}
219+
});
220+
source_handles.push(handle);
221+
}
225222

226-
let listener = TcpListener::bind(&addr).await?;
227-
axum::serve(listener, app).await?;
223+
// Block until all sources finish (normally only on shutdown).
224+
for handle in source_handles {
225+
let _ = handle.await;
226+
}
228227

229-
// Wait for workers to finish (this only happens on shutdown)
230228
for handle in worker_handles {
231229
let _ = handle.await;
232230
}

0 commit comments

Comments
 (0)