1use std::sync::atomic::{AtomicU64, Ordering};
8use std::sync::Arc;
9use std::time::Duration;
10
11use arrow_array::RecordBatch;
12use arrow_schema::SchemaRef;
13use async_trait::async_trait;
14use crossfire::{mpsc, AsyncRx, TryRecvError};
15use tokio::net::TcpListener;
16use tokio::sync::{watch, Notify};
17use tokio::task::JoinHandle;
18use tonic::transport::server::TcpIncoming;
19
20use opentelemetry_proto::tonic::collector::logs::v1::logs_service_server::LogsServiceServer;
21use opentelemetry_proto::tonic::collector::metrics::v1::metrics_service_server::MetricsServiceServer;
22use opentelemetry_proto::tonic::collector::trace::v1::trace_service_server::TraceServiceServer;
23
24use crate::checkpoint::SourceCheckpoint;
25use crate::config::{ConnectorConfig, ConnectorState};
26use crate::connector::{
27 ConnectorTaskOwner, ConnectorTaskTracker, SourceBatch, SourceConnector, SourceConsistency,
28 SourceContract, SourceInputMode, SourcePosition, SourceStart, SourceTopology,
29};
30use crate::error::ConnectorError;
31
32use super::config::{OtelSignal, OtelSourceConfig};
33use super::schema::{logs_schema, metrics_schema, traces_schema};
34use super::server::OtelReceiver;
35
36const SERVER_CLOSE_TIMEOUT: Duration = Duration::from_secs(5);
37
38struct TrackedServerTask {
39 handle: Option<JoinHandle<Result<(), ConnectorError>>>,
40}
41
42struct TaskExitNotify(Arc<Notify>);
43
44impl Drop for TaskExitNotify {
45 fn drop(&mut self) {
46 self.0.notify_one();
47 }
48}
49
50enum ServerWait {
51 Completed(Result<(), ConnectorError>),
52 TimedOut,
53}
54
55impl TrackedServerTask {
56 fn spawn(
57 owner: &ConnectorTaskOwner,
58 exit_notify: Arc<Notify>,
59 future: impl std::future::Future<Output = Result<(), ConnectorError>> + Send + 'static,
60 ) -> Result<Self, ConnectorError> {
61 let task_guard = owner.track().ok_or_else(|| {
62 ConnectorError::Internal("OTel source task generation is already retired".into())
63 })?;
64 let handle = tokio::spawn(async move {
65 let _task_guard = task_guard;
66 let _exit_notify = TaskExitNotify(exit_notify);
67 future.await
68 });
69 Ok(Self {
70 handle: Some(handle),
71 })
72 }
73
74 async fn wait_until(&mut self, deadline: tokio::time::Instant) -> ServerWait {
75 let Some(handle) = self.handle.as_mut() else {
76 return ServerWait::Completed(Ok(()));
77 };
78 match tokio::time::timeout_at(deadline, handle).await {
79 Ok(_) if tokio::time::Instant::now() >= deadline => {
80 self.handle.take();
81 ServerWait::TimedOut
82 }
83 Ok(Ok(result)) => {
84 self.handle.take();
85 ServerWait::Completed(result)
86 }
87 Ok(Err(error)) => {
88 self.handle.take();
89 ServerWait::Completed(Err(ConnectorError::Internal(format!(
90 "OTel gRPC server task failed: {error}"
91 ))))
92 }
93 Err(_) => ServerWait::TimedOut,
94 }
95 }
96
97 async fn take_finished(&mut self) -> Option<Result<(), ConnectorError>> {
98 let handle = self.handle.as_ref()?;
99 if !handle.is_finished() {
100 return None;
101 }
102 let result = self.handle.take()?.await;
103 Some(result.unwrap_or_else(|error| {
104 Err(ConnectorError::Internal(format!(
105 "OTel gRPC server task failed: {error}"
106 )))
107 }))
108 }
109
110 fn abort(&mut self) {
111 if let Some(handle) = self.handle.take() {
112 handle.abort();
113 }
114 }
115}
116
117pub struct OtelSource {
123 config: OtelSourceConfig,
124 schema: SchemaRef,
125 state: ConnectorState,
126 batch_rx: Option<AsyncRx<mpsc::Array<RecordBatch>>>,
127 data_ready: Arc<Notify>,
128 server_task: Option<TrackedServerTask>,
129 shutdown_tx: Option<watch::Sender<bool>>,
130 records_received: Arc<AtomicU64>,
132 requests_received: Arc<AtomicU64>,
133 checkpoint_seq: u64,
134 server_failure: Option<String>,
135 task_owner: ConnectorTaskOwner,
136 task_tracker: ConnectorTaskTracker,
137}
138
139impl OtelSource {
140 #[must_use]
142 pub fn new(schema: SchemaRef, _registry: Option<&prometheus::Registry>) -> Self {
143 let (task_owner, task_tracker) = ConnectorTaskOwner::new();
144 Self {
145 config: OtelSourceConfig::default(),
146 schema,
147 state: ConnectorState::Created,
148 batch_rx: None,
149 data_ready: Arc::new(Notify::new()),
150 server_task: None,
151 shutdown_tx: None,
152 records_received: Arc::new(AtomicU64::new(0)),
153 requests_received: Arc::new(AtomicU64::new(0)),
154 checkpoint_seq: 0,
155 server_failure: None,
156 task_owner,
157 task_tracker,
158 }
159 }
160
161 fn request_shutdown(&mut self) {
162 if let Some(shutdown) = self.shutdown_tx.take() {
163 shutdown.send_replace(true);
164 }
165 self.batch_rx.take();
166 }
167
168 async fn observe_server_exit(&mut self) {
169 if self.server_failure.is_some() || self.state != ConnectorState::Running {
170 return;
171 }
172 let Some(server) = self.server_task.as_mut() else {
173 return;
174 };
175 let Some(result) = server.take_finished().await else {
176 return;
177 };
178 self.server_task.take();
179 self.server_failure = Some(match result {
180 Ok(()) => "OTel gRPC server stopped unexpectedly".into(),
181 Err(ConnectorError::ConnectionFailed(message) | ConnectorError::Internal(message)) => {
182 message
183 }
184 Err(error) => error.to_string(),
185 });
186 self.state = ConnectorState::Failed;
187 }
188
189 fn terminal_server_error(message: &str) -> ConnectorError {
190 ConnectorError::InvalidState {
191 expected: "live OTLP gRPC server".into(),
192 actual: format!("server generation terminated: {message}"),
193 }
194 }
195}
196
197impl Drop for OtelSource {
198 fn drop(&mut self) {
199 self.request_shutdown();
200 if let Some(server) = self.server_task.as_mut() {
201 server.abort();
202 }
203 }
204}
205
206#[async_trait]
207impl SourceConnector for OtelSource {
208 fn terminal_task_tracker(&self) -> Option<ConnectorTaskTracker> {
209 Some(self.task_tracker.clone())
210 }
211
212 async fn start(&mut self, request: SourceStart) -> Result<(), ConnectorError> {
213 let (config, position, _) = request.into_parts();
214 if let SourcePosition::Resume { attempt, .. } = position {
215 return Err(ConnectorError::ConfigurationError(format!(
216 "OTLP is an ephemeral source and cannot resume checkpoint attempt {attempt:?}"
217 )));
218 }
219 if !matches!(self.state, ConnectorState::Created | ConnectorState::Closed)
220 || self.server_task.is_some()
221 {
222 return Err(ConnectorError::InvalidState {
223 expected: "Created or fully closed".into(),
224 actual: format!("{}", self.state),
225 });
226 }
227
228 let candidate_config = OtelSourceConfig::from_config(&config)?;
229
230 let candidate_schema = match candidate_config.signals {
231 OtelSignal::Traces => traces_schema(),
232 OtelSignal::Metrics => metrics_schema(),
233 OtelSignal::Logs => logs_schema(),
234 };
235
236 let (batch_tx, batch_rx) =
237 mpsc::bounded_async::<RecordBatch>(candidate_config.channel_capacity);
238
239 let addr = candidate_config.socket_addr();
240
241 let listener = TcpListener::bind(&addr)
244 .await
245 .map_err(|e| ConnectorError::ConnectionFailed(format!("failed to bind {addr}: {e}")))?;
246 let incoming = TcpIncoming::from(listener).with_nodelay(Some(true));
247
248 let (shutdown_tx, shutdown_rx) = watch::channel(false);
249 let service_guard = self.task_owner.track().ok_or_else(|| {
250 ConnectorError::Internal("OTel source task generation is already retired".into())
251 })?;
252
253 let receiver = OtelReceiver::new(
254 batch_tx,
255 Arc::clone(&candidate_schema),
256 Arc::clone(&self.data_ready),
257 Arc::clone(&self.records_received),
258 Arc::clone(&self.requests_received),
259 candidate_config.batch_size,
260 service_guard,
261 );
262
263 let server_task = match candidate_config.signals {
265 OtelSignal::Traces => spawn_grpc_server(
266 &self.task_owner,
267 TraceServiceServer::new(receiver),
268 incoming,
269 shutdown_rx,
270 Arc::clone(&self.data_ready),
271 ),
272 OtelSignal::Metrics => spawn_grpc_server(
273 &self.task_owner,
274 MetricsServiceServer::new(receiver),
275 incoming,
276 shutdown_rx,
277 Arc::clone(&self.data_ready),
278 ),
279 OtelSignal::Logs => spawn_grpc_server(
280 &self.task_owner,
281 LogsServiceServer::new(receiver),
282 incoming,
283 shutdown_rx,
284 Arc::clone(&self.data_ready),
285 ),
286 }?;
287
288 self.config = candidate_config;
289 self.schema = candidate_schema;
290 self.batch_rx = Some(batch_rx);
291 self.shutdown_tx = Some(shutdown_tx);
292 self.server_task = Some(server_task);
293 self.server_failure = None;
294 self.state = ConnectorState::Running;
295
296 tracing::info!(
297 %addr,
298 signals = ?self.config.signals,
299 batch_size = self.config.batch_size,
300 "OTel source connector started"
301 );
302
303 Ok(())
304 }
305
306 async fn poll_batch(
307 &mut self,
308 max_records: usize,
309 ) -> Result<Option<SourceBatch>, ConnectorError> {
310 if let Some(error) = &self.server_failure {
311 return Err(Self::terminal_server_error(error));
312 }
313 let rx = self.batch_rx.as_ref().ok_or(ConnectorError::InvalidState {
314 expected: "Running".into(),
315 actual: format!("{}", self.state),
316 })?;
317
318 let mut total_rows = 0usize;
319 let mut batches: Vec<RecordBatch> = Vec::new();
320 let mut disconnected = false;
321
322 loop {
323 match rx.try_recv() {
324 Ok(batch) => {
325 total_rows += batch.num_rows();
326 batches.push(batch);
327 if total_rows >= max_records {
328 break;
329 }
330 }
331 Err(TryRecvError::Empty) => break,
332 Err(TryRecvError::Disconnected) => {
333 disconnected = true;
334 break;
335 }
336 }
337 }
338
339 self.observe_server_exit().await;
340
341 if batches.is_empty() {
342 return if let Some(error) = &self.server_failure {
343 Err(Self::terminal_server_error(error))
344 } else if disconnected {
345 self.state = ConnectorState::Closed;
346 Err(ConnectorError::Closed)
347 } else {
348 Ok(None)
349 };
350 }
351
352 self.checkpoint_seq += 1;
353
354 if batches.len() == 1 {
355 return Ok(Some(SourceBatch::new(batches.into_iter().next().unwrap())));
356 }
357
358 let schema = batches[0].schema();
359 let combined =
360 arrow_select::concat::concat_batches(&schema, batches.iter()).map_err(|e| {
361 ConnectorError::ReadError(format!("failed to concatenate OTel batches: {e}"))
362 })?;
363
364 Ok(Some(SourceBatch::new(combined)))
365 }
366
367 async fn discover_schema(
368 &mut self,
369 properties: &std::collections::HashMap<String, String>,
370 ) -> Result<(), ConnectorError> {
371 let Some(sig) = properties
372 .get("signals")
373 .or_else(|| properties.get("signal"))
374 else {
375 return Ok(());
376 };
377 let signal = OtelSignal::parse(sig).map_err(|e| {
378 ConnectorError::ConfigurationError(format!("invalid OTel signal '{sig}': {e}"))
379 })?;
380 self.schema = match signal {
381 OtelSignal::Traces => traces_schema(),
382 OtelSignal::Metrics => metrics_schema(),
383 OtelSignal::Logs => logs_schema(),
384 };
385 Ok(())
386 }
387
388 fn schema(&self) -> SchemaRef {
389 Arc::clone(&self.schema)
390 }
391
392 fn checkpoint(&self) -> SourceCheckpoint {
393 let mut cp = SourceCheckpoint::new();
394 cp.set_offset("batch_sequence", self.checkpoint_seq.to_string());
395 cp.set_offset(
396 "records_received",
397 self.records_received.load(Ordering::Relaxed).to_string(),
398 );
399 cp.set_offset(
400 "requests_received",
401 self.requests_received.load(Ordering::Relaxed).to_string(),
402 );
403 cp.set_metadata("connector", "otel");
404 cp.set_metadata("signals", format!("{:?}", self.config.signals));
405 cp
406 }
407
408 async fn close(&mut self) -> Result<(), ConnectorError> {
409 tracing::info!("OTel source connector shutting down");
410
411 self.request_shutdown();
412
413 let mut completed = false;
414 let mut close_error = self
415 .server_failure
416 .as_ref()
417 .map(|error| Self::terminal_server_error(error));
418 if let Some(server) = self.server_task.as_mut() {
419 let deadline = tokio::time::Instant::now() + SERVER_CLOSE_TIMEOUT;
420 match server.wait_until(deadline).await {
421 ServerWait::Completed(Ok(())) => completed = true,
422 ServerWait::Completed(Err(error)) => {
423 completed = true;
424 tracing::warn!(%error, "OTel gRPC server task failed while closing");
425 close_error = Some(error);
426 }
427 ServerWait::TimedOut => {
428 server.abort();
429 completed = true;
430 close_error = Some(ConnectorError::Internal(
431 "OTel gRPC server exceeded its close deadline; connector generation retired"
432 .into(),
433 ));
434 tracing::warn!("OTel gRPC server exceeded its close deadline and was aborted");
435 }
436 }
437 }
438 if completed {
439 self.server_task.take();
440 }
441
442 if let Some(error) = close_error {
443 self.state = ConnectorState::Failed;
444 Err(error)
445 } else {
446 self.state = ConnectorState::Closed;
447 Ok(())
448 }
449 }
450
451 fn data_ready_notify(&self) -> Option<Arc<Notify>> {
452 Some(Arc::clone(&self.data_ready))
453 }
454
455 fn contract(&self, _config: &ConnectorConfig) -> Result<SourceContract, ConnectorError> {
456 Ok(SourceContract::new(
457 SourceConsistency::Ephemeral,
458 SourceTopology::NodeLocalIngress,
459 SourceInputMode::AppendOnly,
460 ))
461 }
462}
463
464fn spawn_grpc_server<S>(
466 owner: &ConnectorTaskOwner,
467 svc: S,
468 incoming: TcpIncoming,
469 mut shutdown_rx: watch::Receiver<bool>,
470 data_ready: Arc<Notify>,
471) -> Result<TrackedServerTask, ConnectorError>
472where
473 S: tonic::codegen::Service<
474 tonic::codegen::http::Request<tonic::body::Body>,
475 Response = tonic::codegen::http::Response<tonic::body::Body>,
476 Error = std::convert::Infallible,
477 > + tonic::server::NamedService
478 + Clone
479 + Send
480 + Sync
481 + 'static,
482 S::Future: Send + 'static,
483{
484 TrackedServerTask::spawn(owner, data_ready, async move {
485 tonic::transport::Server::builder()
486 .add_service(svc)
487 .serve_with_incoming_shutdown(incoming, async move {
488 let _ = shutdown_rx.wait_for(|&v| v).await;
489 })
490 .await
491 .map_err(|error| {
492 tracing::error!(%error, "OTel gRPC server exited with error");
493 ConnectorError::ConnectionFailed(format!(
494 "OTel gRPC server exited with error: {error}"
495 ))
496 })
497 })
498}
499
500impl std::fmt::Debug for OtelSource {
501 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
502 f.debug_struct("OtelSource")
503 .field("state", &self.state)
504 .field("config", &self.config)
505 .field(
506 "records_received",
507 &self.records_received.load(Ordering::Relaxed),
508 )
509 .field(
510 "requests_received",
511 &self.requests_received.load(Ordering::Relaxed),
512 )
513 .field("checkpoint_seq", &self.checkpoint_seq)
514 .field(
515 "server_running",
516 &self
517 .server_task
518 .as_ref()
519 .and_then(|task| task.handle.as_ref().map(|handle| !handle.is_finished())),
520 )
521 .field("has_shutdown_tx", &self.shutdown_tx.is_some())
522 .finish_non_exhaustive()
523 }
524}
525
526#[cfg(test)]
527mod tests;