Skip to main content

laminar_connectors/kafka/
sink.rs

1//! Kafka sink connector.
2
3use 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
34/// One queue-full retry deadline is shared by every record sent through a
35/// producer phase. A non-record failure stops new enqueue work.
36const 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    /// `FutureProducer::send_result` returned the record, proving that this
81    /// attempt never entered librdkafka's queue.
82    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    /// rdkafka 0.39's `FutureProducer` discards native
93    /// `rd_kafka_message_status` when it creates an owned delivery result.
94    /// Error codes alone cannot prove non-persistence after driver retries.
95    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
121/// Error scope and retryability are independent of persistence certainty. This
122/// is deliberately a positive transient list: fatal, unknown, and future codes
123/// fail closed instead of creating an unbounded restart loop.
124fn 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
233/// Contiguous key buffer — stores all key bytes in a single allocation
234/// with per-row `(offset, length)` pairs. Avoids N separate heap
235/// allocations for N rows.
236struct 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
278/// Kafka sink connector that writes Arrow `RecordBatch` data to Kafka topics.
279///
280/// Operates in Ring 1 (background) receiving data from Ring 0 via the
281/// subscription API.
282///
283/// # Lifecycle
284///
285/// 1. Create with [`KafkaSink::new`]
286/// 2. Call `open()` to create the producer and connect to Kafka
287/// 3. `write_batch()` serializes and produces records; checkpoint flushes provide
288///    durable at-least-once delivery
289/// 4. Call `close()` for clean shutdown
290///
291/// This connector deliberately rejects exactly-once admission because it does
292/// not expose a coordinated external checkpoint namespace/cursor.
293pub struct KafkaSink {
294    /// rdkafka producer (set during `open()`).
295    producer: Option<FutureProducer>,
296    /// Parsed Kafka sink configuration.
297    config: KafkaSinkConfig,
298    /// Format-specific serializer.
299    serializer: Box<dyn RecordSerializer>,
300    /// Partitioner for determining target partitions.
301    partitioner: Box<dyn KafkaPartitioner>,
302    /// Connector lifecycle state.
303    state: ConnectorState,
304    /// Dead letter queue producer (separate, non-transactional).
305    dlq_producer: Option<FutureProducer>,
306    /// Production metrics.
307    metrics: KafkaSinkMetrics,
308    /// Arrow schema for input batches.
309    schema: SchemaRef,
310    /// Optional Schema Registry client.
311    schema_registry: Option<Arc<SchemaRegistryClient>>,
312    /// Shared Avro schema ID (updated after SR registration).
313    avro_schema_id: Arc<std::sync::atomic::AtomicU32>,
314    /// Cached topic partition count (queried from broker metadata after open).
315    topic_partition_count: Option<i32>,
316    /// Sole admission authority for detached producer destruction.
317    task_owner: ConnectorTaskOwner,
318    /// Cloneable terminal observer returned to the connector runtime.
319    task_tracker: ConnectorTaskTracker,
320}
321
322impl KafkaSink {
323    /// Creates a new Kafka sink connector with explicit schema.
324    ///
325    /// # Panics
326    ///
327    /// Panics if `config.format` is not a supported serialization format.
328    /// Call [`KafkaSinkConfig::validate`] first to catch this at config time.
329    #[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    /// Creates a new Kafka sink with Schema Registry integration.
360    ///
361    /// # Panics
362    ///
363    /// Panics if `config.format` is not a supported serialization format.
364    /// Call [`KafkaSinkConfig::validate`] first to catch this at config time.
365    #[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    /// Lifecycle state (Created → Running → Closed).
401    #[must_use]
402    pub fn state(&self) -> ConnectorState {
403        self.state
404    }
405
406    /// Whether Avro schema registration is available.
407    #[must_use]
408    pub fn has_schema_registry(&self) -> bool {
409        self.schema_registry.is_some()
410    }
411
412    /// Destroy the final producer references away from Tokio workers. rdkafka
413    /// purges, flushes for up to 500 ms, and joins its polling thread in Drop.
414    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            // The owner is a field on this live connector, so sealing before
426            // Drop is an invariant violation rather than a recoverable state.
427            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            // This branch cannot run on a Tokio worker. The closure has already
444            // been dropped by `spawn`, so only report the resource failure.
445            tracing::error!(role, %error, "failed to start Kafka producer teardown thread");
446        }
447    }
448
449    /// Ensures the sink schema and SR registration match the actual data.
450    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        // Register with SR *before* advancing schema/serializer so a failure
463        // doesn't leave avro_schema_id stale while the serializer already
464        // encodes with the new schema.
465        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    /// Contiguous key buffer: all key bytes in one allocation with per-row offsets.
504    ///
505    /// Returns `None` if no key column is configured.
506    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        // Try to get string values; fall back to display representation.
525        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    /// Enqueues a definitely rejected record to the dead letter queue. The
560    /// caller retains and drains the returned delivery future with all other
561    /// accepted main/DLQ records from the write.
562    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    /// Synchronously enqueue a record with a short retry on `QueueFull`.
618    /// Uses `send_result` rather than `send`, because the latter is
619    /// `async fn` in rdkafka 0.39+ and only enqueues when polled — which
620    /// would defeat the Vec-of-futures pipelining in `write_batch`.
621    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            // The deadline limits waiting, not the immediate enqueue attempt.
628            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    /// Upsert-envelope produce: collapse the Z-set changelog to one record per merge key, then
643    /// emit a keyed value for a live group (`_op = U`) or a null-value tombstone for a removed group
644    /// (`_op = D`). The topic must be log-compacted and keyed on the merge key for the tombstones to
645    /// GC and for the latest-per-key state to be recoverable from offset 0.
646    #[allow(clippy::cast_possible_truncation, clippy::too_many_lines)] // matches write_batch
647    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        // The value is the collapsed row without the `_op` tag — i.e. the plain MV row.
671        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        // Reject empty/NULL merge keys before producing ANY record: a compacted topic can't
686        // represent an unkeyed row, and a mid-loop bail would leave earlier rows already enqueued.
687        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            // No `.payload()` for a delete → null value = Kafka tombstone.
719            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        // Kafka acknowledges these idempotent writes with all in-sync replicas,
829        // but this connector does not expose an external checkpoint cursor.
830        let consistency = SinkConsistency::DurableAtLeastOnce;
831        // Append records from independent writers compose safely. Upsert records do not carry a
832        // fenced generation, so an old writer can otherwise overwrite a newer value for the same
833        // compacted key after ownership handoff.
834        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        // Schema registration is deferred to the first write_batch(), where the real pipeline
885        // output schema is known — the factory default is a placeholder that would pollute the
886        // registry and break compat checks.
887        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        // Keep DLQ production decoupled from the main producer. Once a producer
903        // exists it lives in `self`, so every later error and cancellation is
904        // routed through tracked off-runtime teardown.
905        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        // Clear stale metadata before lookup. Custom routing must never run
917        // against an assumed partition count after a failed reopen.
918        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)] // Record batch row/byte counts fit in narrower types
977    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        // Upsert envelope: collapse the Z-set changelog per merge key and produce keyed records
989        // (live groups) + null-value tombstones (removed groups). See `write_upsert_batch`.
990        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        // Phase 1: enqueue every record into librdkafka's bounded internal queue.
1018        // A record-local terminal rejection may continue only when it can be
1019        // routed to DLQ. Every other enqueue error stops new dispatch, but all
1020        // delivery futures already accepted by librdkafka are still drained.
1021        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        // Enqueue every DLQ candidate before awaiting any delivery report. Both
1062        // producers then make progress concurrently on their polling threads,
1063        // so draining the future vectors below does not multiply the driver
1064        // delivery deadline by the number of rejected rows.
1065        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        // The runtime deadline must dominate librdkafka's delivery deadline.
1231        // Main and DLQ deliveries are concurrent, so only bounded queue retry
1232        // and scheduling headroom are added rather than a per-row multiplier.
1233        kafka_write_timeout(self.config.delivery_timeout)
1234    }
1235
1236    async fn flush(&mut self) -> Result<(), ConnectorError> {
1237        // Every successful write_batch awaits every delivery report (including
1238        // DLQ delivery). There can be no acknowledged connector write with
1239        // producer work still in flight, so a blocking librdkafka flush adds no
1240        // checkpoint durability and is unsafe to cancel on generation retirement.
1241        Ok(())
1242    }
1243
1244    async fn close(&mut self) -> Result<(), ConnectorError> {
1245        info!("closing Kafka sink connector");
1246
1247        // Completed writes have already observed delivery reports. A cancelled
1248        // write retires the whole producer generation and is replayed from the
1249        // engine checkpoint, so close must not start another blocking flush.
1250        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
1273/// Selects the appropriate serializer for the given format.
1274///
1275/// For Avro, uses the shared `schema_id` handle so that Schema Registry
1276/// registration updates are visible to the serializer.
1277fn 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
1295/// Selects the appropriate partitioner for the given strategy.
1296fn 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}