1use std::fmt::Write as _;
4use std::sync::Arc;
5use std::time::{Duration, Instant};
6
7use arrow_array::{Array, StringArray};
8use arrow_schema::SchemaRef;
9use async_trait::async_trait;
10use rdkafka::error::{KafkaError, RDKafkaErrorCode};
11use rdkafka::message::OwnedHeaders;
12use rdkafka::producer::{DeliveryFuture, FutureProducer, FutureRecord, Producer};
13use rdkafka::ClientConfig;
14use tracing::{debug, info, warn};
15
16use crate::changelog::collapse_changelog;
17use crate::config::{ConnectorConfig, ConnectorState};
18use crate::connector::{
19 ConnectorTaskOwner, ConnectorTaskTracker, SinkConnector, SinkConsistency, SinkContract,
20 SinkInputMode, SinkTopology, WriteResult,
21};
22use crate::error::{ConnectorError, SerdeError};
23use crate::serde::{self, Format, RecordSerializer};
24
25use super::avro_serializer::AvroSerializer;
26use super::metadata_error::{fetch_error, invalid_response, topic_error};
27use super::partitioner::{
28 KafkaPartitioner, KeyHashPartitioner, RoundRobinPartitioner, StickyPartitioner,
29};
30use super::schema_registry::SchemaRegistryClient;
31use super::sink_config::{KafkaSinkConfig, PartitionStrategy, SinkEnvelope};
32use super::sink_metrics::KafkaSinkMetrics;
33
34const QUEUE_RETRY_TIMEOUT: Duration = Duration::from_millis(500);
37const QUEUE_RETRY_INTERVAL: Duration = Duration::from_millis(100);
38const WRITE_TIMEOUT_HEADROOM: Duration = Duration::from_secs(5);
39
40fn queue_retry_delay(deadline: Instant, now: Instant) -> Option<Duration> {
41 let remaining = deadline.saturating_duration_since(now);
42 (!remaining.is_zero()).then_some(remaining.min(QUEUE_RETRY_INTERVAL))
43}
44
45fn delivery_outcome_unknown(
46 operation: &str,
47 detail: impl std::fmt::Display,
48 retryable: bool,
49) -> ConnectorError {
50 ConnectorError::outcome_unknown(
51 format!(
52 "Kafka {operation} was dispatched but its external outcome is not fully known: {detail}"
53 ),
54 retryable,
55 )
56}
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq)]
59enum KafkaFailureCertainty {
60 DefinitelyNotPersisted,
61 OutcomeUnknown,
62}
63
64#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65enum KafkaFailureScope {
66 Record,
67 Infrastructure,
68 Connector,
69}
70
71#[derive(Debug)]
72struct KafkaFailure {
73 certainty: KafkaFailureCertainty,
74 scope: KafkaFailureScope,
75 retryable: bool,
76 detail: String,
77}
78
79impl KafkaFailure {
80 fn enqueue(error: &KafkaError, operation: &str) -> Self {
83 let (scope, retryable) = kafka_error_policy(error);
84 Self {
85 certainty: KafkaFailureCertainty::DefinitelyNotPersisted,
86 scope,
87 retryable,
88 detail: format!("{operation} enqueue failed before dispatch: {error}"),
89 }
90 }
91
92 fn delivery(error: &KafkaError, operation: &str) -> Self {
96 let (scope, retryable) = kafka_error_policy(error);
97 Self {
98 certainty: KafkaFailureCertainty::OutcomeUnknown,
99 scope,
100 retryable,
101 detail: format!("{operation} delivery failed: {error}"),
102 }
103 }
104
105 fn canceled(operation: &str) -> Self {
106 Self {
107 certainty: KafkaFailureCertainty::OutcomeUnknown,
108 scope: KafkaFailureScope::Infrastructure,
109 retryable: true,
110 detail: format!("{operation} delivery canceled because the producer was dropped"),
111 }
112 }
113
114 fn dlq_eligible(&self) -> bool {
115 self.certainty == KafkaFailureCertainty::DefinitelyNotPersisted
116 && self.scope == KafkaFailureScope::Record
117 && !self.retryable
118 }
119}
120
121fn kafka_error_policy(error: &KafkaError) -> (KafkaFailureScope, bool) {
125 match error.rdkafka_error_code() {
126 Some(
127 RDKafkaErrorCode::KeySerialization
128 | RDKafkaErrorCode::ValueSerialization
129 | RDKafkaErrorCode::MessageSizeTooLarge
130 | RDKafkaErrorCode::InvalidTimestamp
131 | RDKafkaErrorCode::InvalidRecord,
132 ) => (KafkaFailureScope::Record, false),
133 Some(
134 RDKafkaErrorCode::BrokerDestroy
135 | RDKafkaErrorCode::BrokerTransportFailure
136 | RDKafkaErrorCode::Resolve
137 | RDKafkaErrorCode::MessageTimedOut
138 | RDKafkaErrorCode::AllBrokersDown
139 | RDKafkaErrorCode::OperationTimedOut
140 | RDKafkaErrorCode::QueueFull
141 | RDKafkaErrorCode::ISRInsufficient
142 | RDKafkaErrorCode::TimedOutQueue
143 | RDKafkaErrorCode::WaitCache
144 | RDKafkaErrorCode::Interrupted
145 | RDKafkaErrorCode::Retry
146 | RDKafkaErrorCode::PurgeQueue
147 | RDKafkaErrorCode::PurgeInflight
148 | RDKafkaErrorCode::DestroyBroker
149 | RDKafkaErrorCode::UnknownTopicOrPartition
150 | RDKafkaErrorCode::LeaderNotAvailable
151 | RDKafkaErrorCode::NotLeaderForPartition
152 | RDKafkaErrorCode::RequestTimedOut
153 | RDKafkaErrorCode::BrokerNotAvailable
154 | RDKafkaErrorCode::ReplicaNotAvailable
155 | RDKafkaErrorCode::NetworkException
156 | RDKafkaErrorCode::NotEnoughReplicas
157 | RDKafkaErrorCode::NotEnoughReplicasAfterAppend
158 | RDKafkaErrorCode::NotController
159 | RDKafkaErrorCode::KafkaStorageError
160 | RDKafkaErrorCode::ReassignmentInProgress
161 | RDKafkaErrorCode::FencedLeaderEpoch
162 | RDKafkaErrorCode::UnknownLeaderEpoch
163 | RDKafkaErrorCode::StaleBrokerEpoch
164 | RDKafkaErrorCode::EligibleLeadersNotAvailable
165 | RDKafkaErrorCode::ThrottlingQuotaExceeded
166 | RDKafkaErrorCode::UnknownTopicId,
167 ) => (KafkaFailureScope::Infrastructure, true),
168 _ => (KafkaFailureScope::Connector, false),
169 }
170}
171
172fn unresolved_delivery_error(
173 operation: &str,
174 total: usize,
175 applied: usize,
176 definitely_not_persisted: usize,
177 ambiguous: usize,
178 first_error: Option<String>,
179 retryable: bool,
180) -> ConnectorError {
181 let detail = format!(
182 "{definitely_not_persisted} definitely not persisted, {ambiguous} outcome unknown, \
183 {applied} already applied out of {total}; first error: {}",
184 first_error.unwrap_or_else(|| "unknown".into())
185 );
186 if ambiguous > 0 || applied > 0 {
187 delivery_outcome_unknown(operation, detail, retryable)
188 } else if retryable {
189 ConnectorError::WriteError(format!("Kafka {operation} failed: {detail}"))
190 } else {
191 ConnectorError::ConfigurationError(format!("Kafka {operation} was rejected: {detail}"))
192 }
193}
194
195fn record_failure(
196 failure: &KafkaFailure,
197 count: usize,
198 definitely_not_persisted: &mut usize,
199 ambiguous: &mut usize,
200 first_error: &mut Option<String>,
201 retryable: &mut bool,
202) {
203 match failure.certainty {
204 KafkaFailureCertainty::DefinitelyNotPersisted => *definitely_not_persisted += count,
205 KafkaFailureCertainty::OutcomeUnknown => *ambiguous += count,
206 }
207 *retryable &= failure.retryable;
208 first_error.get_or_insert_with(|| failure.detail.clone());
209}
210
211fn kafka_write_timeout(delivery_timeout: Duration) -> Duration {
212 delivery_timeout
213 .saturating_add(QUEUE_RETRY_TIMEOUT)
214 .saturating_add(QUEUE_RETRY_TIMEOUT)
215 .saturating_add(WRITE_TIMEOUT_HEADROOM)
216}
217
218fn producer_creation_error(role: &str, error: &KafkaError) -> ConnectorError {
219 ConnectorError::ConfigurationError(format!("failed to create Kafka {role} producer: {error}"))
220}
221
222fn validate_payload_cardinality(expected: usize, actual: usize) -> Result<(), ConnectorError> {
223 if expected == actual {
224 Ok(())
225 } else {
226 Err(ConnectorError::Serde(SerdeError::RecordCountMismatch {
227 expected,
228 got: actual,
229 }))
230 }
231}
232
233struct KeyBuffer {
237 data: Vec<u8>,
238 offsets: Vec<(usize, usize)>,
239}
240
241impl KeyBuffer {
242 fn with_capacity(num_rows: usize, avg_key_len: usize) -> Self {
243 Self {
244 data: Vec::with_capacity(num_rows * avg_key_len),
245 offsets: Vec::with_capacity(num_rows),
246 }
247 }
248
249 fn push(&mut self, key: &[u8]) {
250 let start = self.data.len();
251 self.data.extend_from_slice(key);
252 self.offsets.push((start, key.len()));
253 }
254
255 fn push_empty(&mut self) {
256 self.offsets.push((0, 0));
257 }
258
259 fn key(&self, i: usize) -> &[u8] {
260 let (start, len) = self.offsets[i];
261 &self.data[start..start + len]
262 }
263
264 #[cfg(test)]
265 fn len(&self) -> usize {
266 self.offsets.len()
267 }
268}
269
270impl std::ops::Index<usize> for KeyBuffer {
271 type Output = [u8];
272
273 fn index(&self, i: usize) -> &[u8] {
274 self.key(i)
275 }
276}
277
278pub struct KafkaSink {
294 producer: Option<FutureProducer>,
296 config: KafkaSinkConfig,
298 serializer: Box<dyn RecordSerializer>,
300 partitioner: Box<dyn KafkaPartitioner>,
302 state: ConnectorState,
304 dlq_producer: Option<FutureProducer>,
306 metrics: KafkaSinkMetrics,
308 schema: SchemaRef,
310 schema_registry: Option<Arc<SchemaRegistryClient>>,
312 avro_schema_id: Arc<std::sync::atomic::AtomicU32>,
314 topic_partition_count: Option<i32>,
316 task_owner: ConnectorTaskOwner,
318 task_tracker: ConnectorTaskTracker,
320}
321
322impl KafkaSink {
323 #[must_use]
330 pub fn new(
331 schema: SchemaRef,
332 config: KafkaSinkConfig,
333 registry: Option<&prometheus::Registry>,
334 ) -> Self {
335 let avro_schema_id = Arc::new(std::sync::atomic::AtomicU32::new(0));
336 let serializer =
337 select_serializer(config.format, &schema, Arc::clone(&avro_schema_id), None)
338 .expect("format validated in KafkaSinkConfig::validate()");
339 let partitioner = select_partitioner(config.partitioner);
340 let (task_owner, task_tracker) = ConnectorTaskOwner::new();
341
342 Self {
343 producer: None,
344 config,
345 serializer,
346 partitioner,
347 state: ConnectorState::Created,
348 dlq_producer: None,
349 metrics: KafkaSinkMetrics::new(registry),
350 schema,
351 schema_registry: None,
352 avro_schema_id,
353 topic_partition_count: None,
354 task_owner,
355 task_tracker,
356 }
357 }
358
359 #[must_use]
366 pub fn with_schema_registry(
367 schema: SchemaRef,
368 config: KafkaSinkConfig,
369 sr_client: SchemaRegistryClient,
370 ) -> Self {
371 let sr = Arc::new(sr_client);
372 let avro_schema_id = Arc::new(std::sync::atomic::AtomicU32::new(0));
373 let serializer = select_serializer(
374 config.format,
375 &schema,
376 Arc::clone(&avro_schema_id),
377 Some(Arc::clone(&sr)),
378 )
379 .expect("format validated in KafkaSinkConfig::validate()");
380 let partitioner = select_partitioner(config.partitioner);
381 let (task_owner, task_tracker) = ConnectorTaskOwner::new();
382
383 Self {
384 producer: None,
385 config,
386 serializer,
387 partitioner,
388 state: ConnectorState::Created,
389 dlq_producer: None,
390 metrics: KafkaSinkMetrics::new(None),
391 schema,
392 schema_registry: Some(sr),
393 avro_schema_id,
394 topic_partition_count: None,
395 task_owner,
396 task_tracker,
397 }
398 }
399
400 #[must_use]
402 pub fn state(&self) -> ConnectorState {
403 self.state
404 }
405
406 #[must_use]
408 pub fn has_schema_registry(&self) -> bool {
409 self.schema_registry.is_some()
410 }
411
412 fn retire_producers(&mut self) {
415 if let Some(producer) = self.producer.take() {
416 self.spawn_producer_drop(producer, "main");
417 }
418 if let Some(producer) = self.dlq_producer.take() {
419 self.spawn_producer_drop(producer, "DLQ");
420 }
421 }
422
423 fn spawn_producer_drop(&self, producer: FutureProducer, role: &'static str) {
424 let Some(terminal_guard) = self.task_owner.track() else {
425 tracing::error!(
428 role,
429 "Kafka producer teardown could not enter the terminal task tracker"
430 );
431 return;
432 };
433 let teardown = move || {
434 let _terminal_guard = terminal_guard;
435 drop(producer);
436 };
437 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
438 drop(runtime.spawn_blocking(teardown));
439 } else if let Err(error) = std::thread::Builder::new()
440 .name("laminardb-kafka-producer-drop".into())
441 .spawn(teardown)
442 {
443 tracing::error!(role, %error, "failed to start Kafka producer teardown thread");
446 }
447 }
448
449 async fn ensure_schema_ready(
451 &mut self,
452 batch_schema: &SchemaRef,
453 ) -> Result<(), ConnectorError> {
454 let schema_changed = self.schema != *batch_schema;
455 let needs_registration = self.config.format == Format::Avro
456 && (schema_changed
457 || self
458 .avro_schema_id
459 .load(std::sync::atomic::Ordering::Relaxed)
460 == 0);
461
462 if needs_registration {
466 if let Some(ref sr) = self.schema_registry {
467 let subject = format!("{}-value", self.config.topic);
468 let avro_schema =
469 super::schema_registry::arrow_to_avro_schema(batch_schema, &self.config.topic)
470 .map_err(ConnectorError::Serde)?;
471 let schema_id = sr
472 .register_schema(
473 &subject,
474 &avro_schema,
475 super::schema_registry::SchemaType::Avro,
476 )
477 .await?;
478 #[allow(clippy::cast_sign_loss)]
479 self.avro_schema_id
480 .store(schema_id as u32, std::sync::atomic::Ordering::Relaxed);
481 info!(subject = %subject, schema_id, "registered Avro schema");
482 }
483 }
484
485 if schema_changed {
486 debug!(
487 old = ?self.schema.fields().iter().map(|f| f.name()).collect::<Vec<_>>(),
488 new = ?batch_schema.fields().iter().map(|f| f.name()).collect::<Vec<_>>(),
489 "sink schema updated from incoming batch"
490 );
491 self.schema = batch_schema.clone();
492 self.serializer = select_serializer(
493 self.config.format,
494 &self.schema,
495 Arc::clone(&self.avro_schema_id),
496 self.schema_registry.clone(),
497 )?;
498 }
499
500 Ok(())
501 }
502
503 fn extract_keys(
507 &self,
508 batch: &arrow_array::RecordBatch,
509 ) -> Result<Option<KeyBuffer>, ConnectorError> {
510 let Some(key_col) = &self.config.key_column else {
511 return Ok(None);
512 };
513
514 let col_idx = batch.schema().index_of(key_col).map_err(|_| {
515 ConnectorError::ConfigurationError(format!(
516 "key column '{key_col}' not found in schema"
517 ))
518 })?;
519
520 let array = batch.column(col_idx);
521 let num_rows = batch.num_rows();
522 let mut buf = KeyBuffer::with_capacity(num_rows, 32);
523
524 if let Some(str_array) = array.as_any().downcast_ref::<StringArray>() {
526 for i in 0..num_rows {
527 if str_array.is_null(i) {
528 buf.push_empty();
529 } else {
530 buf.push(str_array.value(i).as_bytes());
531 }
532 }
533 } else {
534 use std::fmt::Write;
535 let formatter = arrow_cast::display::ArrayFormatter::try_new(
536 array,
537 &arrow_cast::display::FormatOptions::default(),
538 )
539 .map_err(|e| {
540 ConnectorError::Internal(format!(
541 "failed to create array formatter for key column: {e}"
542 ))
543 })?;
544 let mut fmt_buf = String::with_capacity(64);
545 for i in 0..num_rows {
546 if array.is_null(i) {
547 buf.push_empty();
548 } else {
549 fmt_buf.clear();
550 let _ = write!(fmt_buf, "{}", formatter.value(i));
551 buf.push(fmt_buf.as_bytes());
552 }
553 }
554 }
555
556 Ok(Some(buf))
557 }
558
559 async fn enqueue_dlq(
563 &self,
564 payload: &[u8],
565 key: Option<&[u8]>,
566 error_msg: &str,
567 queue_deadline: Instant,
568 ) -> Result<DeliveryFuture, KafkaFailure> {
569 let dlq_producer = self.dlq_producer.as_ref().ok_or_else(|| KafkaFailure {
570 certainty: KafkaFailureCertainty::DefinitelyNotPersisted,
571 scope: KafkaFailureScope::Connector,
572 detail: "DLQ topic or producer is not configured".into(),
573 retryable: false,
574 })?;
575 let dlq_topic = self.config.dlq_topic.as_ref().ok_or_else(|| KafkaFailure {
576 certainty: KafkaFailureCertainty::DefinitelyNotPersisted,
577 scope: KafkaFailureScope::Connector,
578 detail: "DLQ topic or producer is not configured".into(),
579 retryable: false,
580 })?;
581
582 let now = std::time::SystemTime::now()
583 .duration_since(std::time::UNIX_EPOCH)
584 .unwrap_or_else(|_| {
585 tracing::warn!("system clock before Unix epoch — using 0 for DLQ timestamp");
586 std::time::Duration::ZERO
587 })
588 .as_millis()
589 .to_string();
590 let headers = OwnedHeaders::new()
591 .insert(rdkafka::message::Header {
592 key: "__dlq.error",
593 value: Some(error_msg.as_bytes()),
594 })
595 .insert(rdkafka::message::Header {
596 key: "__dlq.topic",
597 value: Some(self.config.topic.as_bytes()),
598 })
599 .insert(rdkafka::message::Header {
600 key: "__dlq.timestamp",
601 value: Some(now.as_bytes()),
602 });
603
604 let mut record = FutureRecord::to(dlq_topic)
605 .payload(payload)
606 .headers(headers);
607
608 if let Some(k) = key {
609 record = record.key(k);
610 }
611
612 Self::enqueue_with_queue_retry(dlq_producer, record, queue_deadline)
613 .await
614 .map_err(|error| KafkaFailure::enqueue(&error, "Kafka DLQ"))
615 }
616
617 async fn enqueue_with_queue_retry(
622 producer: &FutureProducer,
623 mut record: FutureRecord<'_, [u8], [u8]>,
624 queue_deadline: Instant,
625 ) -> Result<DeliveryFuture, KafkaError> {
626 loop {
627 match producer.send_result(record) {
629 Ok(fut) => return Ok(fut),
630 Err((error @ KafkaError::MessageProduction(RDKafkaErrorCode::QueueFull), r)) => {
631 let Some(delay) = queue_retry_delay(queue_deadline, Instant::now()) else {
632 return Err(error);
633 };
634 record = r;
635 tokio::time::sleep(delay).await;
636 }
637 Err((error, _)) => return Err(error),
638 }
639 }
640 }
641
642 #[allow(clippy::cast_possible_truncation, clippy::too_many_lines)] async fn write_upsert_batch(
648 &mut self,
649 batch: &arrow_array::RecordBatch,
650 ) -> Result<WriteResult, ConnectorError> {
651 let key_col = self.config.key_column.clone().ok_or_else(|| {
652 ConnectorError::ConfigurationError("envelope = 'upsert' requires 'key.column'".into())
653 })?;
654 let collapsed = collapse_changelog(batch, std::slice::from_ref(&key_col))?;
655 let rows = collapsed.num_rows();
656 if rows == 0 {
657 return Ok(WriteResult::new(0, 0));
658 }
659 let op_idx = collapsed
660 .schema()
661 .index_of("_op")
662 .map_err(|_| ConnectorError::Internal("collapsed changelog missing _op".into()))?;
663 let ops = collapsed
664 .column(op_idx)
665 .as_any()
666 .downcast_ref::<StringArray>()
667 .ok_or_else(|| ConnectorError::Internal("_op column is not Utf8".into()))?
668 .clone();
669
670 let value_idxs: Vec<usize> = (0..collapsed.num_columns())
672 .filter(|&i| i != op_idx)
673 .collect();
674 let value_batch = collapsed
675 .project(&value_idxs)
676 .map_err(|e| ConnectorError::Internal(format!("project value columns: {e}")))?;
677
678 self.ensure_schema_ready(&value_batch.schema()).await?;
679 let payloads = self.serializer.serialize(&value_batch).map_err(|e| {
680 self.metrics.record_serialization_error();
681 ConnectorError::Serde(e)
682 })?;
683 validate_payload_cardinality(rows, payloads.len())?;
684 let keys = self.extract_keys(&collapsed)?;
685 if let Some(kb) = keys.as_ref() {
688 for i in 0..payloads.len() {
689 if kb.key(i).is_empty() {
690 return Err(ConnectorError::WriteError(format!(
691 "upsert envelope: row {i} has an empty/NULL merge key"
692 )));
693 }
694 }
695 }
696
697 let partition_count =
698 self.topic_partition_count
699 .ok_or_else(|| ConnectorError::InvalidState {
700 expected: "broker topic metadata installed".into(),
701 actual: "partition count is unavailable".into(),
702 })?;
703 let producer = self
704 .producer
705 .as_ref()
706 .ok_or_else(|| ConnectorError::InvalidState {
707 expected: "producer initialized".into(),
708 actual: "producer is None".into(),
709 })?;
710
711 let mut delivery_futures = Vec::with_capacity(rows);
712 let mut enqueue_failure: Option<(KafkaFailure, usize)> = None;
713 let queue_deadline = Instant::now() + QUEUE_RETRY_TIMEOUT;
714 for (i, payload) in payloads.iter().enumerate() {
715 let key: Option<&[u8]> = keys.as_ref().map(|kb| kb.key(i));
716 let is_delete = ops.value(i) == "D";
717 let partition = self.partitioner.partition(key, partition_count);
718 let mut record: FutureRecord<'_, [u8], [u8]> = FutureRecord::to(&self.config.topic);
720 if let Some(k) = key {
721 record = record.key(k);
722 }
723 if !is_delete {
724 record = record.payload(payload.as_slice());
725 }
726 if let Some(p) = partition {
727 record = record.partition(p);
728 }
729 match Self::enqueue_with_queue_retry(producer, record, queue_deadline).await {
730 Ok(future) => delivery_futures.push((Instant::now(), future, is_delete, i)),
731 Err(error) => {
732 let mut failure = KafkaFailure::enqueue(&error, "Kafka upsert");
733 let affected = rows - i;
734 if affected > 1 {
735 let _ = write!(
736 failure.detail,
737 "; {} later record(s) were not attempted",
738 affected - 1
739 );
740 }
741 enqueue_failure = Some((failure, affected));
742 break;
743 }
744 }
745 }
746
747 let mut records_written: usize = 0;
748 let mut bytes_written: u64 = 0;
749 let mut definitely_not_persisted: usize = 0;
750 let mut ambiguous: usize = 0;
751 let mut first_error: Option<String> = None;
752 let mut retryable = true;
753 for (send_time, future, is_delete, i) in delivery_futures {
754 match future.await {
755 Ok(Ok(_)) => {
756 self.metrics
757 .record_produce_latency(send_time.elapsed().as_micros() as u64);
758 records_written += 1;
759 if !is_delete {
760 bytes_written += payloads[i].len() as u64;
761 }
762 }
763 Ok(Err((err, _))) => {
764 let failure = KafkaFailure::delivery(&err, "Kafka upsert");
765 record_failure(
766 &failure,
767 1,
768 &mut definitely_not_persisted,
769 &mut ambiguous,
770 &mut first_error,
771 &mut retryable,
772 );
773 }
774 Err(_canceled) => {
775 let failure = KafkaFailure::canceled("Kafka upsert");
776 record_failure(
777 &failure,
778 1,
779 &mut definitely_not_persisted,
780 &mut ambiguous,
781 &mut first_error,
782 &mut retryable,
783 );
784 }
785 }
786 }
787 if let Some((failure, affected)) = enqueue_failure {
788 record_failure(
789 &failure,
790 affected,
791 &mut definitely_not_persisted,
792 &mut ambiguous,
793 &mut first_error,
794 &mut retryable,
795 );
796 }
797 self.metrics
798 .record_write(records_written as u64, bytes_written);
799 if definitely_not_persisted > 0 || ambiguous > 0 {
800 self.metrics.record_error();
801 return Err(unresolved_delivery_error(
802 "upsert produce",
803 rows,
804 records_written,
805 definitely_not_persisted,
806 ambiguous,
807 first_error,
808 retryable,
809 ));
810 }
811 Ok(WriteResult::new(records_written, bytes_written))
812 }
813}
814
815#[async_trait]
816#[allow(clippy::too_many_lines)]
817impl SinkConnector for KafkaSink {
818 fn terminal_task_tracker(&self) -> Option<ConnectorTaskTracker> {
819 Some(self.task_tracker.clone())
820 }
821
822 fn contract(&self, config: &ConnectorConfig) -> Result<SinkContract, ConnectorError> {
823 let cfg = if config.properties().is_empty() {
824 self.config.clone()
825 } else {
826 KafkaSinkConfig::from_config(config)?
827 };
828 let consistency = SinkConsistency::DurableAtLeastOnce;
831 let topology = if cfg.envelope == SinkEnvelope::Upsert {
835 SinkTopology::Singleton
836 } else {
837 SinkTopology::MultiWriter
838 };
839 let input_mode = if cfg.envelope == SinkEnvelope::Upsert {
840 SinkInputMode::FullChangelog
841 } else {
842 SinkInputMode::AppendOnly
843 };
844 Ok(SinkContract::new(consistency, topology, input_mode))
845 }
846
847 async fn open(&mut self, config: &ConnectorConfig) -> Result<(), ConnectorError> {
848 self.state = ConnectorState::Initializing;
849
850 if !config.properties().is_empty() {
851 let parsed = KafkaSinkConfig::from_config(config)?;
852 self.config = parsed;
853 self.serializer = select_serializer(
854 self.config.format,
855 &self.schema,
856 Arc::clone(&self.avro_schema_id),
857 self.schema_registry.clone(),
858 )?;
859 self.partitioner = select_partitioner(self.config.partitioner);
860 }
861 self.config.validate()?;
862 info!(
863 brokers = %self.config.bootstrap_servers,
864 topic = %self.config.topic,
865 format = %self.config.format,
866 "opening Kafka sink connector"
867 );
868
869 if let Some(ref url) = self.config.schema_registry_url {
870 if self.schema_registry.is_none() {
871 let sr = if let Some(ref ca_path) = self.config.schema_registry_ssl_ca_location {
872 SchemaRegistryClient::with_tls(
873 url,
874 self.config.schema_registry_auth.clone(),
875 ca_path,
876 )?
877 } else {
878 SchemaRegistryClient::new(url, self.config.schema_registry_auth.clone())?
879 };
880 self.schema_registry = Some(Arc::new(sr));
881 }
882 }
883
884 if self.config.format == Format::Avro {
888 if let Some(ref sr) = self.schema_registry {
889 if let Some(ref compat) = self.config.schema_compatibility {
890 let subject = format!("{}-value", self.config.topic);
891 sr.set_compatibility_level(&subject, *compat).await?;
892 }
893 }
894 }
895
896 let rdkafka_config: ClientConfig = self.config.to_rdkafka_config();
897 let producer: FutureProducer = rdkafka_config
898 .create()
899 .map_err(|error| producer_creation_error("main", &error))?;
900 self.producer = Some(producer);
901
902 if self.config.dlq_topic.is_some() {
906 let dlq_config = self.config.to_dlq_rdkafka_config();
907 match dlq_config.create::<FutureProducer>() {
908 Ok(dlq_producer) => self.dlq_producer = Some(dlq_producer),
909 Err(error) => {
910 self.retire_producers();
911 return Err(producer_creation_error("DLQ", &error));
912 }
913 }
914 }
915
916 self.topic_partition_count = None;
919 let producer = self
920 .producer
921 .as_ref()
922 .expect("Kafka producer was installed above")
923 .clone();
924 let topic = self.config.topic.clone();
925 let metadata_guard = self
926 .task_owner
927 .track()
928 .expect("live Kafka sink must admit its metadata lookup");
929 let metadata = tokio::task::spawn_blocking(move || -> Result<i32, ConnectorError> {
930 let _metadata_guard = metadata_guard;
931 let metadata = producer
932 .client()
933 .fetch_metadata(Some(&topic), Duration::from_secs(5))
934 .map_err(|error| fetch_error(&topic, &error))?;
935 let topic_metadata = metadata
936 .topics()
937 .iter()
938 .find(|candidate| candidate.name() == topic.as_str())
939 .ok_or_else(|| invalid_response(&topic, "metadata response omitted the topic"))?;
940 if let Some(error) = topic_metadata.error() {
941 return Err(topic_error(&topic, error.into()));
942 }
943 i32::try_from(topic_metadata.partitions().len())
944 .ok()
945 .filter(|count| *count > 0)
946 .ok_or_else(|| invalid_response(&topic, "broker returned no partitions"))
947 })
948 .await;
949 match metadata {
950 Ok(Ok(count)) => {
951 self.topic_partition_count = Some(count);
952 info!(
953 topic = %self.config.topic,
954 partitions = count,
955 "queried topic partition count from broker"
956 );
957 }
958 Ok(Err(error)) => {
959 self.retire_producers();
960 return Err(error);
961 }
962 Err(error) => {
963 self.retire_producers();
964 return Err(ConnectorError::Internal(format!(
965 "Kafka metadata worker for topic '{}' failed: {error}",
966 self.config.topic
967 )));
968 }
969 }
970
971 self.state = ConnectorState::Running;
972 info!("Kafka sink connector opened successfully");
973 Ok(())
974 }
975
976 #[allow(clippy::cast_possible_truncation)] async fn write_batch(
978 &mut self,
979 batch: &arrow_array::RecordBatch,
980 ) -> Result<WriteResult, ConnectorError> {
981 if self.state != ConnectorState::Running {
982 return Err(ConnectorError::InvalidState {
983 expected: "Running".into(),
984 actual: self.state.to_string(),
985 });
986 }
987
988 if self.config.envelope == SinkEnvelope::Upsert {
991 return self.write_upsert_batch(batch).await;
992 }
993
994 self.ensure_schema_ready(&batch.schema()).await?;
995
996 let payloads = self.serializer.serialize(batch).map_err(|e| {
997 self.metrics.record_serialization_error();
998 ConnectorError::Serde(e)
999 })?;
1000 validate_payload_cardinality(batch.num_rows(), payloads.len())?;
1001
1002 let keys = self.extract_keys(batch)?;
1003 let partition_count =
1004 self.topic_partition_count
1005 .ok_or_else(|| ConnectorError::InvalidState {
1006 expected: "broker topic metadata installed".into(),
1007 actual: "partition count is unavailable".into(),
1008 })?;
1009 let producer = self
1010 .producer
1011 .as_ref()
1012 .ok_or_else(|| ConnectorError::InvalidState {
1013 expected: "producer initialized".into(),
1014 actual: "producer is None".into(),
1015 })?;
1016
1017 let mut delivery_futures = Vec::with_capacity(payloads.len());
1022 let mut dlq_candidates = Vec::new();
1023 let mut enqueue_failure: Option<(KafkaFailure, usize)> = None;
1024 let main_queue_deadline = Instant::now() + QUEUE_RETRY_TIMEOUT;
1025 for (i, payload) in payloads.iter().enumerate() {
1026 let key: Option<&[u8]> = keys.as_ref().map(|kb| kb.key(i)).filter(|k| !k.is_empty());
1027 let partition = self.partitioner.partition(key, partition_count);
1028
1029 let mut record = FutureRecord::to(&self.config.topic).payload(payload.as_slice());
1030 if let Some(k) = key {
1031 record = record.key(k);
1032 }
1033 if let Some(p) = partition {
1034 record = record.partition(p);
1035 }
1036
1037 match Self::enqueue_with_queue_retry(producer, record, main_queue_deadline).await {
1038 Ok(future) => delivery_futures.push((i, Instant::now(), future)),
1039 Err(error) => {
1040 self.metrics.record_error();
1041 let mut failure = KafkaFailure::enqueue(&error, "Kafka");
1042 if self.dlq_producer.is_some() && failure.dlq_eligible() {
1043 dlq_candidates.push((i, failure.detail));
1044 continue;
1045 }
1046
1047 let affected = payloads.len() - i;
1048 if affected > 1 {
1049 let _ = write!(
1050 failure.detail,
1051 "; {} later record(s) were not attempted",
1052 affected - 1
1053 );
1054 }
1055 enqueue_failure = Some((failure, affected));
1056 break;
1057 }
1058 }
1059 }
1060
1061 let dlq_candidate_count = dlq_candidates.len();
1066 let mut dlq_delivery_futures = Vec::with_capacity(dlq_candidate_count);
1067 let mut dlq_enqueue_failure: Option<(KafkaFailure, usize)> = None;
1068 let dlq_queue_deadline = Instant::now() + QUEUE_RETRY_TIMEOUT;
1069 for (position, (row, original_error)) in dlq_candidates.into_iter().enumerate() {
1070 let key = keys
1071 .as_ref()
1072 .map(|buffer| buffer.key(row))
1073 .filter(|key| !key.is_empty());
1074 match self
1075 .enqueue_dlq(&payloads[row], key, &original_error, dlq_queue_deadline)
1076 .await
1077 {
1078 Ok(future) => dlq_delivery_futures.push((row, original_error, future)),
1079 Err(mut failure) => {
1080 let affected = dlq_candidate_count - position;
1081 if affected > 1 {
1082 let _ = write!(
1083 failure.detail,
1084 "; {} later DLQ record(s) were not attempted",
1085 affected - 1
1086 );
1087 }
1088 dlq_enqueue_failure = Some((failure, affected));
1089 break;
1090 }
1091 }
1092 }
1093
1094 let mut records_written: usize = 0;
1095 let mut bytes_written: u64 = 0;
1096 let mut definitely_not_persisted: usize = 0;
1097 let mut ambiguous: usize = 0;
1098 let mut dlq_records: usize = 0;
1099 let mut dlq_bytes: u64 = 0;
1100 let mut first_error: Option<String> = None;
1101 let mut retryable = true;
1102 for (row, send_time, future) in delivery_futures {
1103 match future.await {
1104 Ok(Ok(_)) => {
1105 let latency_us = send_time.elapsed().as_micros() as u64;
1106 self.metrics.record_produce_latency(latency_us);
1107 records_written += 1;
1108 bytes_written += payloads[row].len() as u64;
1109 }
1110 Ok(Err((error, _))) => {
1111 self.metrics.record_error();
1112 let failure = KafkaFailure::delivery(&error, "Kafka");
1113 record_failure(
1114 &failure,
1115 1,
1116 &mut definitely_not_persisted,
1117 &mut ambiguous,
1118 &mut first_error,
1119 &mut retryable,
1120 );
1121 }
1122 Err(_) => {
1123 self.metrics.record_error();
1124 let failure = KafkaFailure::canceled("Kafka");
1125 record_failure(
1126 &failure,
1127 1,
1128 &mut definitely_not_persisted,
1129 &mut ambiguous,
1130 &mut first_error,
1131 &mut retryable,
1132 );
1133 }
1134 }
1135 }
1136 for (row, original_error, future) in dlq_delivery_futures {
1137 match future.await {
1138 Ok(Ok(_)) => {
1139 self.metrics.record_dlq();
1140 dlq_records += 1;
1141 dlq_bytes += payloads[row].len() as u64;
1142 }
1143 Ok(Err((error, _))) => {
1144 self.metrics.record_error();
1145 let mut failure = KafkaFailure::delivery(&error, "Kafka DLQ");
1146 failure.detail = format!("original: {original_error}; {}", failure.detail);
1147 warn!(
1148 original_error = %original_error,
1149 dlq_error = %failure.detail,
1150 "failed to route definitely rejected Kafka record to DLQ"
1151 );
1152 record_failure(
1153 &failure,
1154 1,
1155 &mut definitely_not_persisted,
1156 &mut ambiguous,
1157 &mut first_error,
1158 &mut retryable,
1159 );
1160 }
1161 Err(_) => {
1162 self.metrics.record_error();
1163 let mut failure = KafkaFailure::canceled("Kafka DLQ");
1164 failure.detail = format!("original: {original_error}; {}", failure.detail);
1165 record_failure(
1166 &failure,
1167 1,
1168 &mut definitely_not_persisted,
1169 &mut ambiguous,
1170 &mut first_error,
1171 &mut retryable,
1172 );
1173 }
1174 }
1175 }
1176 if let Some((failure, affected)) = dlq_enqueue_failure {
1177 record_failure(
1178 &failure,
1179 affected,
1180 &mut definitely_not_persisted,
1181 &mut ambiguous,
1182 &mut first_error,
1183 &mut retryable,
1184 );
1185 }
1186 if let Some((failure, affected)) = enqueue_failure {
1187 record_failure(
1188 &failure,
1189 affected,
1190 &mut definitely_not_persisted,
1191 &mut ambiguous,
1192 &mut first_error,
1193 &mut retryable,
1194 );
1195 }
1196
1197 self.metrics
1198 .record_write(records_written as u64, bytes_written);
1199
1200 debug!(
1201 records = records_written,
1202 dlq_records,
1203 bytes = bytes_written,
1204 definitely_not_persisted,
1205 ambiguous,
1206 "wrote batch to Kafka"
1207 );
1208
1209 let applied = records_written + dlq_records;
1210 if definitely_not_persisted > 0 || ambiguous > 0 {
1211 return Err(unresolved_delivery_error(
1212 "produce",
1213 payloads.len(),
1214 applied,
1215 definitely_not_persisted,
1216 ambiguous,
1217 first_error,
1218 retryable,
1219 ));
1220 }
1221
1222 Ok(WriteResult::new(applied, bytes_written + dlq_bytes))
1223 }
1224
1225 fn schema(&self) -> SchemaRef {
1226 self.schema.clone()
1227 }
1228
1229 fn suggested_write_timeout(&self) -> Duration {
1230 kafka_write_timeout(self.config.delivery_timeout)
1234 }
1235
1236 async fn flush(&mut self) -> Result<(), ConnectorError> {
1237 Ok(())
1242 }
1243
1244 async fn close(&mut self) -> Result<(), ConnectorError> {
1245 info!("closing Kafka sink connector");
1246
1247 self.retire_producers();
1251 self.state = ConnectorState::Closed;
1252 info!("Kafka sink connector closed");
1253 Ok(())
1254 }
1255}
1256
1257impl Drop for KafkaSink {
1258 fn drop(&mut self) {
1259 self.retire_producers();
1260 }
1261}
1262
1263impl std::fmt::Debug for KafkaSink {
1264 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1265 f.debug_struct("KafkaSink")
1266 .field("state", &self.state)
1267 .field("topic", &self.config.topic)
1268 .field("format", &self.config.format)
1269 .finish_non_exhaustive()
1270 }
1271}
1272
1273fn select_serializer(
1278 format: Format,
1279 schema: &SchemaRef,
1280 schema_id: Arc<std::sync::atomic::AtomicU32>,
1281 registry: Option<Arc<SchemaRegistryClient>>,
1282) -> Result<Box<dyn RecordSerializer>, ConnectorError> {
1283 match format {
1284 Format::Avro => Ok(Box::new(AvroSerializer::with_shared_schema_id(
1285 schema.clone(),
1286 schema_id,
1287 registry,
1288 ))),
1289 other => serde::create_serializer(other).map_err(|e| {
1290 ConnectorError::ConfigurationError(format!("unsupported sink format '{other}': {e}"))
1291 }),
1292 }
1293}
1294
1295fn select_partitioner(strategy: PartitionStrategy) -> Box<dyn KafkaPartitioner> {
1297 match strategy {
1298 PartitionStrategy::KeyHash => Box::new(KeyHashPartitioner::new()),
1299 PartitionStrategy::RoundRobin => Box::new(RoundRobinPartitioner::new()),
1300 PartitionStrategy::Sticky => Box::new(StickyPartitioner::new(100)),
1301 }
1302}
1303
1304#[cfg(test)]
1305mod tests {
1306 use super::*;
1307 use arrow_array::Int64Array;
1308 use arrow_schema::{DataType, Field, Schema};
1309
1310 struct MismatchedSerializer;
1311
1312 impl RecordSerializer for MismatchedSerializer {
1313 fn serialize(&self, _batch: &arrow_array::RecordBatch) -> Result<Vec<Vec<u8>>, SerdeError> {
1314 Ok(vec![b"one-payload".to_vec()])
1315 }
1316
1317 fn format(&self) -> Format {
1318 Format::Json
1319 }
1320 }
1321
1322 fn test_schema() -> SchemaRef {
1323 Arc::new(Schema::new(vec![
1324 Field::new("id", DataType::Int64, false),
1325 Field::new("value", DataType::Utf8, false),
1326 ]))
1327 }
1328
1329 fn test_config() -> KafkaSinkConfig {
1330 let mut cfg = KafkaSinkConfig::default();
1331 cfg.bootstrap_servers = "localhost:9092".into();
1332 cfg.topic = "output-events".into();
1333 cfg
1334 }
1335
1336 fn two_row_batch() -> arrow_array::RecordBatch {
1337 arrow_array::RecordBatch::try_new(
1338 test_schema(),
1339 vec![
1340 Arc::new(Int64Array::from(vec![1, 2])),
1341 Arc::new(StringArray::from(vec!["a", "b"])),
1342 ],
1343 )
1344 .unwrap()
1345 }
1346
1347 #[test]
1348 fn test_new_defaults() {
1349 let sink = KafkaSink::new(test_schema(), test_config(), None);
1350 assert_eq!(sink.state(), ConnectorState::Created);
1351 assert!(sink.producer.is_none());
1352 assert_eq!(sink.topic_partition_count, None);
1353 }
1354
1355 #[test]
1356 fn local_producer_creation_failure_is_terminal_configuration() {
1357 let mut invalid = ClientConfig::new();
1358 invalid.set("laminardb.invalid.kafka.property", "value");
1359 let Err(error) = invalid.create::<FutureProducer>() else {
1360 panic!("an unknown local librdkafka option must fail client creation");
1361 };
1362
1363 let mapped = producer_creation_error("main", &error);
1364 assert!(matches!(mapped, ConnectorError::ConfigurationError(_)));
1365 assert!(!mapped.is_transient());
1366 }
1367
1368 #[test]
1369 fn malformed_broker_metadata_is_terminal_and_has_no_partition_fallback() {
1370 let error = invalid_response("orders", "metadata response omitted the topic");
1371 assert!(matches!(error, ConnectorError::ConfigurationError(_)));
1372 assert!(!error.is_transient());
1373
1374 let sink = KafkaSink::new(test_schema(), test_config(), None);
1375 assert_eq!(sink.topic_partition_count, None);
1376 }
1377
1378 #[tokio::test]
1379 async fn append_cardinality_mismatch_fails_before_producer_access() {
1380 let mut sink = KafkaSink::new(test_schema(), test_config(), None);
1381 sink.state = ConnectorState::Running;
1382 sink.topic_partition_count = Some(3);
1383 sink.serializer = Box::new(MismatchedSerializer);
1384
1385 let error = sink.write_batch(&two_row_batch()).await.unwrap_err();
1386 assert!(matches!(
1387 error,
1388 ConnectorError::Serde(SerdeError::RecordCountMismatch {
1389 expected: 2,
1390 got: 1
1391 })
1392 ));
1393 }
1394
1395 #[tokio::test]
1396 async fn upsert_cardinality_mismatch_fails_before_producer_access() {
1397 let mut config = test_config();
1398 config.envelope = SinkEnvelope::Upsert;
1399 config.key_column = Some("id".into());
1400 let mut sink = KafkaSink::new(test_schema(), config, None);
1401 sink.state = ConnectorState::Running;
1402 sink.topic_partition_count = Some(3);
1403 sink.serializer = Box::new(MismatchedSerializer);
1404
1405 let error = sink.write_batch(&two_row_batch()).await.unwrap_err();
1406 assert!(matches!(
1407 error,
1408 ConnectorError::Serde(SerdeError::RecordCountMismatch {
1409 expected: 2,
1410 got: 1
1411 })
1412 ));
1413 }
1414
1415 #[tokio::test]
1416 async fn schema_registration_preserves_terminal_registry_error() {
1417 use wiremock::matchers::{method, path};
1418 use wiremock::{Mock, MockServer, ResponseTemplate};
1419
1420 let server = MockServer::start().await;
1421 Mock::given(method("POST"))
1422 .and(path("/subjects/output-events-value/versions"))
1423 .respond_with(ResponseTemplate::new(422).set_body_string("invalid schema"))
1424 .mount(&server)
1425 .await;
1426 let mut config = test_config();
1427 config.format = Format::Avro;
1428 config.schema_registry_url = Some(server.uri());
1429 let registry = SchemaRegistryClient::new(server.uri(), None).unwrap();
1430 let mut sink = KafkaSink::with_schema_registry(test_schema(), config, registry);
1431
1432 let error = sink.ensure_schema_ready(&test_schema()).await.unwrap_err();
1433 assert!(matches!(error, ConnectorError::ConfigurationError(_)));
1434 assert!(!error.is_transient());
1435 assert!(error.to_string().contains("output-events-value"));
1436 }
1437
1438 #[tokio::test]
1439 async fn compatibility_put_preserves_terminal_registry_error() {
1440 use wiremock::matchers::{method, path};
1441 use wiremock::{Mock, MockServer, ResponseTemplate};
1442
1443 let server = MockServer::start().await;
1444 Mock::given(method("PUT"))
1445 .and(path("/config/output-events-value"))
1446 .respond_with(ResponseTemplate::new(401).set_body_string("invalid credentials"))
1447 .mount(&server)
1448 .await;
1449 let mut config = test_config();
1450 config.format = Format::Avro;
1451 config.schema_registry_url = Some(server.uri());
1452 config.schema_compatibility = Some(crate::kafka::config::CompatibilityLevel::Backward);
1453 let registry = SchemaRegistryClient::new(server.uri(), None).unwrap();
1454 let mut sink = KafkaSink::with_schema_registry(test_schema(), config, registry);
1455
1456 let error = sink.open(&ConnectorConfig::new("kafka")).await.unwrap_err();
1457 assert!(matches!(error, ConnectorError::ConfigurationError(_)));
1458 assert!(!error.is_transient());
1459 assert!(error.to_string().contains("output-events-value"));
1460 }
1461
1462 #[test]
1463 fn terminal_tracker_seals_when_sink_is_dropped() {
1464 let sink = KafkaSink::new(test_schema(), test_config(), None);
1465 let terminal = sink.terminal_task_tracker().unwrap();
1466 assert!(!terminal.is_terminated());
1467 drop(sink);
1468 assert!(terminal.is_terminated());
1469 }
1470
1471 #[test]
1472 fn test_schema_returned() {
1473 let schema = test_schema();
1474 let sink = KafkaSink::new(schema.clone(), test_config(), None);
1475 assert_eq!(sink.schema(), schema);
1476 }
1477
1478 #[test]
1479 fn contract_is_multi_writer_durable_at_least_once() {
1480 let sink = KafkaSink::new(test_schema(), test_config(), None);
1481 let contract = sink.contract(&ConnectorConfig::new("kafka")).unwrap();
1482 assert_eq!(contract.consistency, SinkConsistency::DurableAtLeastOnce);
1483 assert_eq!(contract.topology, SinkTopology::MultiWriter);
1484 assert_eq!(contract.input_mode, SinkInputMode::AppendOnly);
1485 assert_eq!(sink.suggested_write_timeout(), Duration::from_secs(126));
1486 }
1487
1488 #[test]
1489 fn delivery_error_codes_never_claim_non_persistence_without_native_status() {
1490 for code in [
1491 RDKafkaErrorCode::MessageTimedOut,
1492 RDKafkaErrorCode::TimedOutQueue,
1493 RDKafkaErrorCode::PurgeQueue,
1494 RDKafkaErrorCode::PurgeInflight,
1495 RDKafkaErrorCode::MessageSizeTooLarge,
1496 RDKafkaErrorCode::TopicAuthorizationFailed,
1497 ] {
1498 let error = KafkaError::MessageProduction(code);
1499 assert_eq!(
1500 KafkaFailure::delivery(&error, "test").certainty,
1501 KafkaFailureCertainty::OutcomeUnknown,
1502 "delivery code {code:?}"
1503 );
1504 }
1505 }
1506
1507 #[test]
1508 fn only_terminal_record_local_enqueue_failures_are_dlq_eligible() {
1509 let too_large = KafkaFailure::enqueue(
1510 &KafkaError::MessageProduction(RDKafkaErrorCode::MessageSizeTooLarge),
1511 "test",
1512 );
1513 assert_eq!(
1514 too_large.certainty,
1515 KafkaFailureCertainty::DefinitelyNotPersisted
1516 );
1517 assert_eq!(too_large.scope, KafkaFailureScope::Record);
1518 assert!(!too_large.retryable);
1519 assert!(too_large.dlq_eligible());
1520
1521 let queue_full = KafkaFailure::enqueue(
1522 &KafkaError::MessageProduction(RDKafkaErrorCode::QueueFull),
1523 "test",
1524 );
1525 assert_eq!(queue_full.scope, KafkaFailureScope::Infrastructure);
1526 assert!(queue_full.retryable);
1527 assert!(!queue_full.dlq_eligible());
1528
1529 let unauthorized = KafkaFailure::enqueue(
1530 &KafkaError::MessageProduction(RDKafkaErrorCode::TopicAuthorizationFailed),
1531 "test",
1532 );
1533 assert_eq!(unauthorized.scope, KafkaFailureScope::Connector);
1534 assert!(!unauthorized.retryable);
1535 assert!(!unauthorized.dlq_eligible());
1536 }
1537
1538 #[test]
1539 fn fatal_and_unknown_codes_fail_closed() {
1540 for code in [
1541 RDKafkaErrorCode::Unknown,
1542 RDKafkaErrorCode::Fatal,
1543 RDKafkaErrorCode::ProducerFenced,
1544 ] {
1545 let failure = KafkaFailure::delivery(&KafkaError::MessageProduction(code), "test");
1546 assert_eq!(failure.scope, KafkaFailureScope::Connector);
1547 assert!(!failure.retryable, "delivery code {code:?}");
1548 }
1549 }
1550
1551 #[test]
1552 fn aggregate_retryability_is_the_conjunction_of_every_failure() {
1553 let transient = KafkaFailure::delivery(
1554 &KafkaError::MessageProduction(RDKafkaErrorCode::RequestTimedOut),
1555 "test",
1556 );
1557 let terminal = KafkaFailure::enqueue(
1558 &KafkaError::MessageProduction(RDKafkaErrorCode::MessageSizeTooLarge),
1559 "test",
1560 );
1561 let mut definitely_not_persisted = 0;
1562 let mut ambiguous = 0;
1563 let mut first_error = None;
1564 let mut retryable = true;
1565 record_failure(
1566 &transient,
1567 1,
1568 &mut definitely_not_persisted,
1569 &mut ambiguous,
1570 &mut first_error,
1571 &mut retryable,
1572 );
1573 record_failure(
1574 &terminal,
1575 2,
1576 &mut definitely_not_persisted,
1577 &mut ambiguous,
1578 &mut first_error,
1579 &mut retryable,
1580 );
1581
1582 assert_eq!(definitely_not_persisted, 2);
1583 assert_eq!(ambiguous, 1);
1584 assert!(!retryable);
1585 }
1586
1587 #[test]
1588 fn suggested_timeout_tracks_driver_deadline_with_constant_headroom() {
1589 let mut config = test_config();
1590 config.delivery_timeout = Duration::from_secs(42);
1591 let sink = KafkaSink::new(test_schema(), config, None);
1592 assert_eq!(sink.suggested_write_timeout(), Duration::from_secs(48));
1593 }
1594
1595 #[test]
1596 fn queue_retry_wait_is_bounded_across_records() {
1597 let start = Instant::now();
1598 let deadline = start + QUEUE_RETRY_TIMEOUT;
1599 let mut now = start;
1600 let mut total_wait = Duration::ZERO;
1601
1602 for _ in 0..32 {
1603 if let Some(delay) = queue_retry_delay(deadline, now) {
1604 total_wait += delay;
1605 now += delay;
1606 }
1607 }
1608
1609 assert_eq!(total_wait, QUEUE_RETRY_TIMEOUT);
1610 assert_eq!(now, deadline);
1611 assert_eq!(queue_retry_delay(deadline, now), None);
1612 }
1613
1614 #[test]
1615 fn later_record_cannot_restart_an_expired_queue_retry_budget() {
1616 let start = Instant::now();
1617 let deadline = start + QUEUE_RETRY_TIMEOUT;
1618
1619 assert_eq!(
1620 queue_retry_delay(
1621 deadline,
1622 deadline.checked_sub(Duration::from_millis(25)).unwrap(),
1623 ),
1624 Some(Duration::from_millis(25))
1625 );
1626 assert_eq!(
1627 queue_retry_delay(deadline, deadline + Duration::from_secs(1)),
1628 None
1629 );
1630 }
1631
1632 #[test]
1633 fn partial_or_ambiguous_batch_requires_generation_retirement() {
1634 let error =
1635 unresolved_delivery_error("produce", 3, 1, 2, 0, Some("rejected".into()), false);
1636 assert!(error.is_outcome_unknown());
1637 assert!(!error.is_transient());
1638
1639 let error =
1640 unresolved_delivery_error("produce", 1, 0, 0, 1, Some("timed out".into()), true);
1641 assert!(error.is_outcome_unknown());
1642 assert!(error.is_transient());
1643
1644 let error =
1645 unresolved_delivery_error("produce", 1, 0, 1, 0, Some("too large".into()), false);
1646 assert!(matches!(error, ConnectorError::ConfigurationError(_)));
1647 }
1648
1649 #[test]
1650 fn upsert_contract_requires_singleton_writer() {
1651 let mut config = test_config();
1652 config.envelope = SinkEnvelope::Upsert;
1653 config.key_column = Some("id".into());
1654 let sink = KafkaSink::new(test_schema(), config, None);
1655
1656 let contract = sink.contract(&ConnectorConfig::new("kafka")).unwrap();
1657
1658 assert_eq!(contract.consistency, SinkConsistency::DurableAtLeastOnce);
1659 assert_eq!(contract.topology, SinkTopology::Singleton);
1660 assert_eq!(contract.input_mode, SinkInputMode::FullChangelog);
1661 }
1662
1663 #[test]
1664 fn test_serializer_selection_json() {
1665 let sink = KafkaSink::new(test_schema(), test_config(), None);
1666 assert_eq!(sink.serializer.format(), Format::Json);
1667 }
1668
1669 #[test]
1670 fn test_serializer_selection_avro() {
1671 let mut cfg = test_config();
1672 cfg.format = Format::Avro;
1673 let sink = KafkaSink::new(test_schema(), cfg, None);
1674 assert_eq!(sink.serializer.format(), Format::Avro);
1675 }
1676
1677 #[test]
1678 fn test_with_schema_registry() {
1679 let sr = SchemaRegistryClient::new("http://localhost:8081", None).unwrap();
1680 let mut cfg = test_config();
1681 cfg.format = Format::Avro;
1682 cfg.schema_registry_url = Some("http://localhost:8081".into());
1683
1684 let sink = KafkaSink::with_schema_registry(test_schema(), cfg, sr);
1685 assert!(sink.has_schema_registry());
1686 assert_eq!(sink.serializer.format(), Format::Avro);
1687 }
1688
1689 #[test]
1690 fn test_debug_output() {
1691 let sink = KafkaSink::new(test_schema(), test_config(), None);
1692 let debug = format!("{sink:?}");
1693 assert!(debug.contains("KafkaSink"));
1694 assert!(debug.contains("output-events"));
1695 }
1696
1697 #[test]
1698 fn test_extract_keys_no_key_column() {
1699 let sink = KafkaSink::new(test_schema(), test_config(), None);
1700 let batch = arrow_array::RecordBatch::try_new(
1701 test_schema(),
1702 vec![
1703 Arc::new(Int64Array::from(vec![1, 2])),
1704 Arc::new(StringArray::from(vec!["a", "b"])),
1705 ],
1706 )
1707 .unwrap();
1708 assert!(sink.extract_keys(&batch).unwrap().is_none());
1709 }
1710
1711 #[test]
1712 fn test_extract_keys_with_key_column() {
1713 let mut cfg = test_config();
1714 cfg.key_column = Some("value".into());
1715 let sink = KafkaSink::new(test_schema(), cfg, None);
1716 let batch = arrow_array::RecordBatch::try_new(
1717 test_schema(),
1718 vec![
1719 Arc::new(Int64Array::from(vec![1, 2])),
1720 Arc::new(StringArray::from(vec!["key-a", "key-b"])),
1721 ],
1722 )
1723 .unwrap();
1724 let keys = sink.extract_keys(&batch).unwrap().unwrap();
1725 assert_eq!(keys.len(), 2);
1726 assert_eq!(&keys[0], b"key-a");
1727 assert_eq!(&keys[1], b"key-b");
1728 }
1729}