Skip to main content

risingwave_connector/source/kafka/
enumerator.rs

1// Copyright 2025 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{HashMap, HashSet};
16use std::sync::{Arc, LazyLock, Weak};
17use std::time::Duration;
18
19use anyhow::{Context, anyhow};
20use async_trait::async_trait;
21use moka::future::Cache as MokaCache;
22use moka::ops::compute::Op;
23use rdkafka::admin::{AdminClient, AdminOptions};
24use rdkafka::consumer::BaseConsumer;
25#[cfg(not(madsim))]
26use rdkafka::consumer::Consumer;
27use rdkafka::error::{KafkaError, KafkaResult};
28use rdkafka::types::RDKafkaErrorCode;
29use rdkafka::{ClientConfig, Offset, TopicPartitionList};
30use risingwave_common::bail;
31use risingwave_common::id::FragmentId;
32use risingwave_common::metrics::LabelGuardedIntGauge;
33use thiserror_ext::AsReport;
34
35use crate::connector_common::read_kafka_log_level;
36use crate::error::{ConnectorError, ConnectorResult};
37use crate::source::SourceEnumeratorContextRef;
38use crate::source::base::SplitEnumerator;
39use crate::source::kafka::split::KafkaSplit;
40use crate::source::kafka::{
41    KAFKA_ISOLATION_LEVEL, KafkaConnectionProps, KafkaContextCommon, KafkaProperties,
42    RwConsumerContext,
43};
44
45type KafkaConsumer = BaseConsumer<RwConsumerContext>;
46type KafkaAdmin = AdminClient<RwConsumerContext>;
47
48/// Consumer client is shared, and the cache doesn't manage the lifecycle, so we store `Weak` and no eviction.
49pub static SHARED_KAFKA_CONSUMER: LazyLock<MokaCache<KafkaConnectionProps, Weak<KafkaConsumer>>> =
50    LazyLock::new(|| moka::future::Cache::builder().build());
51/// Admin client is short-lived, so we store `Arc` and sets a time-to-idle eviction policy.
52pub static SHARED_KAFKA_ADMIN: LazyLock<MokaCache<KafkaConnectionProps, Arc<KafkaAdmin>>> =
53    LazyLock::new(|| {
54        moka::future::Cache::builder()
55            .time_to_idle(Duration::from_secs(5 * 60))
56            .build()
57    });
58
59#[derive(Debug, Copy, Clone, Eq, PartialEq)]
60pub enum KafkaEnumeratorOffset {
61    Earliest,
62    Latest,
63    Timestamp(i64),
64    None,
65}
66
67pub struct KafkaSplitEnumerator {
68    context: SourceEnumeratorContextRef,
69    broker_address: String,
70    topic: String,
71    client: Arc<KafkaConsumer>,
72    start_offset: KafkaEnumeratorOffset,
73    /// Start offsets already resolved from a `Timestamp` start offset, keyed by partition.
74    ///
75    /// The timestamp is fixed for the lifetime of the enumerator, and `offsets_for_times` is
76    /// expensive for the broker when the timestamp is older than the local log retention (it has
77    /// to consult tiered storage), so each partition is resolved only on the tick that first
78    /// observes it.
79    resolved_start_offsets: HashMap<i32, Option<i64>>,
80
81    // maybe used in the future for batch processing
82    stop_offset: KafkaEnumeratorOffset,
83
84    sync_call_timeout: Duration,
85    high_watermark_metrics: HashMap<i32, LabelGuardedIntGauge>,
86
87    properties: KafkaProperties,
88    config: rdkafka::ClientConfig,
89}
90
91impl KafkaSplitEnumerator {
92    fn report_consumer_group_delete_failure(&self, group_id: &str) {
93        let source_id = self.context.info.source_id.to_string();
94        self.context
95            .metrics
96            .kafka_consumer_group_delete_failure_count
97            .with_label_values(&[&source_id, group_id])
98            .inc();
99    }
100
101    async fn drop_consumer_groups(&self, fragment_ids: Vec<FragmentId>) -> ConnectorResult<()> {
102        let admin = Box::pin(SHARED_KAFKA_ADMIN.try_get_with_by_ref(
103            &self.properties.connection,
104            async {
105                tracing::info!("build new kafka admin for {}", self.broker_address);
106                Ok(Arc::new(
107                    build_kafka_admin(&self.config, &self.properties).await?,
108                ))
109            },
110        ))
111        .await?;
112
113        let group_ids = fragment_ids
114            .iter()
115            .map(|fragment_id| self.properties.group_id(*fragment_id))
116            .collect::<Vec<_>>();
117        let group_id_refs = group_ids
118            .iter()
119            .map(|group_id| group_id.as_str())
120            .collect::<Vec<_>>();
121
122        let res = match admin
123            .delete_groups(&group_id_refs, &AdminOptions::default())
124            .await
125        {
126            Ok(res) => res,
127            Err(err) => {
128                for group_id in &group_ids {
129                    self.report_consumer_group_delete_failure(group_id);
130                }
131                tracing::warn!(
132                    error = %err.as_report(),
133                    topic = self.topic,
134                    ?group_ids,
135                    "failed to delete Kafka consumer groups"
136                );
137                return Err(err.into());
138            }
139        };
140
141        let mut failure_count = 0;
142        for result in &res {
143            if let Err((group_id, error_code)) = result {
144                failure_count += 1;
145                self.report_consumer_group_delete_failure(group_id);
146                tracing::warn!(
147                    topic = self.topic,
148                    group_id,
149                    error = %error_code.as_report(),
150                    "failed to delete Kafka consumer group"
151                );
152            }
153        }
154        tracing::debug!(
155            topic = self.topic,
156            ?fragment_ids,
157            ?res,
158            failure_count,
159            "delete groups result"
160        );
161        Ok(())
162    }
163}
164
165#[async_trait]
166impl SplitEnumerator for KafkaSplitEnumerator {
167    type Properties = KafkaProperties;
168    type Split = KafkaSplit;
169
170    async fn new(
171        properties: KafkaProperties,
172        context: SourceEnumeratorContextRef,
173    ) -> ConnectorResult<KafkaSplitEnumerator> {
174        let mut config = rdkafka::ClientConfig::new();
175        let common_props = &properties.common;
176
177        let broker_address = properties.connection.brokers.clone();
178        let topic = common_props.topic.clone();
179        config.set("bootstrap.servers", &broker_address);
180        config.set("isolation.level", KAFKA_ISOLATION_LEVEL);
181        if let Some(log_level) = read_kafka_log_level() {
182            config.set_log_level(log_level);
183        }
184        properties.connection.set_security_properties(&mut config);
185        properties.set_client(&mut config);
186        // The meta-side split enumerator does not export librdkafka native stats, so disable
187        // periodic statistics callbacks here even if the source properties enable them for
188        // compute-side readers.
189        config.set("statistics.interval.ms", "0");
190        let mut scan_start_offset = match properties
191            .scan_startup_mode
192            .as_ref()
193            .map(|s| s.to_lowercase())
194            .as_deref()
195        {
196            Some("earliest") => KafkaEnumeratorOffset::Earliest,
197            Some("latest") => KafkaEnumeratorOffset::Latest,
198            None => KafkaEnumeratorOffset::Earliest,
199            _ => bail!(
200                "properties `scan_startup_mode` only supports earliest and latest or leaving it empty"
201            ),
202        };
203
204        if let Some(s) = &properties.time_offset {
205            let time_offset = s.parse::<i64>().map_err(|e| anyhow!(e))?;
206            scan_start_offset = KafkaEnumeratorOffset::Timestamp(time_offset)
207        }
208
209        let mut client: Option<Arc<KafkaConsumer>> = None;
210        SHARED_KAFKA_CONSUMER
211            .entry_by_ref(&properties.connection)
212            .and_try_compute_with::<_, _, ConnectorError>(|maybe_entry| async {
213                if let Some(entry) = maybe_entry {
214                    let entry_value = entry.into_value();
215                    if let Some(client_) = entry_value.upgrade() {
216                        // return if the client is already built
217                        tracing::info!("reuse existing kafka client for {}", broker_address);
218                        client = Some(client_);
219                        return Ok(Op::Nop);
220                    }
221                }
222                tracing::info!("build new kafka client for {}", broker_address);
223                client = Some(build_kafka_client(&config, &properties).await?);
224                Ok(Op::Put(Arc::downgrade(client.as_ref().unwrap())))
225            })
226            .await?;
227
228        Ok(Self {
229            context,
230            broker_address,
231            topic,
232            client: client.unwrap(),
233            start_offset: scan_start_offset,
234            resolved_start_offsets: HashMap::new(),
235            stop_offset: KafkaEnumeratorOffset::None,
236            sync_call_timeout: properties.common.sync_call_timeout,
237            high_watermark_metrics: HashMap::new(),
238            properties,
239            config,
240        })
241    }
242
243    async fn list_splits(&mut self) -> ConnectorResult<Vec<KafkaSplit>> {
244        // `KafkaSplitEnumerator` uses `BaseConsumer`, which does not have a background polling
245        // thread. Poll once per `list_splits` invocation so meta's periodic source-manager tick
246        // can serve queued callbacks like librdkafka statistics events.
247        //
248        // This meta-side enumerator does not have a fragment id, so it cannot derive the
249        // compute-side consumer group id. Polling this no-group client may therefore return
250        // `UnknownGroup`, which is expected and intentionally filtered below. Other poll errors
251        // are still logged as warnings.
252        if let Some(Err(poll_err)) = {
253            #[cfg(not(madsim))]
254            {
255                self.client.poll(Duration::ZERO)
256            }
257            #[cfg(madsim)]
258            {
259                self.client.poll(Duration::ZERO).await
260            }
261        } && !is_expected_no_group_poll_error(&poll_err)
262        {
263            tracing::warn!(
264                error = %poll_err.as_report(),
265                topic = self.topic,
266                broker_address = self.broker_address,
267                "failed to poll kafka client");
268        }
269
270        let topic_partitions = self.fetch_topic_partition().await.with_context(|| {
271            format!(
272                "failed to fetch metadata from kafka ({})",
273                self.broker_address
274            )
275        })?;
276
277        let watermarks = self.get_watermarks(topic_partitions.as_ref()).await?;
278        let mut start_offsets = self
279            .fetch_start_offset(topic_partitions.as_ref(), &watermarks)
280            .await?;
281
282        let mut stop_offsets = self
283            .fetch_stop_offset(topic_partitions.as_ref(), &watermarks)
284            .await?;
285
286        let ret: Vec<_> = topic_partitions
287            .into_iter()
288            .map(|partition| KafkaSplit {
289                topic: self.topic.clone(),
290                partition,
291                start_offset: start_offsets.remove(&partition).unwrap(),
292                stop_offset: stop_offsets.remove(&partition).unwrap(),
293            })
294            .collect();
295
296        Ok(ret)
297    }
298
299    async fn on_drop_fragments(&mut self, fragment_ids: Vec<FragmentId>) -> ConnectorResult<()> {
300        self.drop_consumer_groups(fragment_ids).await
301    }
302
303    async fn on_finish_backfill(&mut self, fragment_ids: Vec<FragmentId>) -> ConnectorResult<()> {
304        self.drop_consumer_groups(fragment_ids).await
305    }
306}
307
308fn is_expected_no_group_poll_error(error: &KafkaError) -> bool {
309    matches!(
310        error,
311        KafkaError::MessageConsumption(RDKafkaErrorCode::UnknownGroup)
312    )
313}
314
315async fn build_kafka_client(
316    config: &ClientConfig,
317    properties: &KafkaProperties,
318) -> ConnectorResult<Arc<KafkaConsumer>> {
319    let ctx_common = KafkaContextCommon::new(
320        properties.privatelink_common.broker_rewrite_map.clone(),
321        None,
322        None,
323        properties.aws_auth_props.clone(),
324        properties.connection.is_aws_msk_iam(),
325    )
326    .await?;
327    let client_ctx = RwConsumerContext::new(ctx_common);
328    let client: KafkaConsumer = config.create_with_context(client_ctx).await?;
329
330    // Note that before any SASL/OAUTHBEARER broker connection can succeed the application must call
331    // rd_kafka_oauthbearer_set_token() once – either directly or, more typically, by invoking either
332    // rd_kafka_poll(), rd_kafka_consumer_poll(), rd_kafka_queue_poll(), etc, in order to cause retrieval
333    // of an initial token to occur.
334    // https://docs.confluent.io/platform/current/clients/librdkafka/html/rdkafka_8h.html#a988395722598f63396d7a1bedb22adaf
335    if properties.connection.is_aws_msk_iam() {
336        #[cfg(not(madsim))]
337        client.poll(Duration::from_secs(10)); // note: this is a blocking call
338        #[cfg(madsim)]
339        client.poll(Duration::from_secs(10)).await;
340    }
341    Ok(Arc::new(client))
342}
343async fn build_kafka_admin(
344    config: &ClientConfig,
345    properties: &KafkaProperties,
346) -> ConnectorResult<KafkaAdmin> {
347    let ctx_common = KafkaContextCommon::new(
348        properties.privatelink_common.broker_rewrite_map.clone(),
349        None,
350        None,
351        properties.aws_auth_props.clone(),
352        properties.connection.is_aws_msk_iam(),
353    )
354    .await?;
355    let client_ctx = RwConsumerContext::new(ctx_common);
356    let client: KafkaAdmin = config.create_with_context(client_ctx).await?;
357    // AdminClient calls start_poll_thread on creation, so the additional poll seems not needed. (And currently no API for this.)
358    Ok(client)
359}
360
361impl KafkaSplitEnumerator {
362    async fn get_watermarks(
363        &mut self,
364        partitions: &[i32],
365    ) -> KafkaResult<HashMap<i32, (i64, i64)>> {
366        let mut map = HashMap::new();
367        for partition in partitions {
368            let (low, high) = self
369                .client
370                .fetch_watermarks(self.topic.as_str(), *partition, self.sync_call_timeout)
371                .await?;
372            self.report_high_watermark(*partition, high);
373            map.insert(*partition, (low, high));
374        }
375        tracing::debug!("fetch kafka watermarks: {map:?}");
376        Ok(map)
377    }
378
379    pub async fn list_splits_batch(
380        &mut self,
381        expect_start_timestamp_millis: Option<i64>,
382        expect_stop_timestamp_millis: Option<i64>,
383    ) -> ConnectorResult<Vec<KafkaSplit>> {
384        let topic_partitions = self.fetch_topic_partition().await.with_context(|| {
385            format!(
386                "failed to fetch metadata from kafka ({})",
387                self.broker_address
388            )
389        })?;
390
391        // Watermark here has nothing to do with watermark in streaming processing. Watermark
392        // here means smallest/largest offset available for reading.
393        let mut watermarks = self.get_watermarks(topic_partitions.as_ref()).await?;
394
395        // here we are getting the start offset and end offset for each partition with the given
396        // timestamp if the timestamp is None, we will use the low watermark and high
397        // watermark as the start and end offset if the timestamp is provided, we will use
398        // the watermark to narrow down the range
399        let mut expect_start_offset = if let Some(ts) = expect_start_timestamp_millis {
400            Some(
401                self.fetch_offset_for_time(topic_partitions.as_ref(), ts, &watermarks)
402                    .await?,
403            )
404        } else {
405            None
406        };
407
408        let mut expect_stop_offset = if let Some(ts) = expect_stop_timestamp_millis {
409            Some(
410                self.fetch_offset_for_time(topic_partitions.as_ref(), ts, &watermarks)
411                    .await?,
412            )
413        } else {
414            None
415        };
416
417        Ok(topic_partitions
418            .iter()
419            .map(|partition| {
420                let (low, high) = watermarks.remove(partition).unwrap();
421                let start_offset = {
422                    let earliest_offset = low - 1;
423                    let start = expect_start_offset
424                        .as_mut()
425                        .map(|m| m.remove(partition).flatten().unwrap_or(earliest_offset))
426                        .unwrap_or(earliest_offset);
427                    i64::max(start, earliest_offset)
428                };
429                let stop_offset = {
430                    let stop = expect_stop_offset
431                        .as_mut()
432                        .map(|m| m.remove(partition).unwrap_or(Some(high)))
433                        .unwrap_or(Some(high))
434                        .unwrap_or(high);
435                    i64::min(stop, high)
436                };
437
438                if start_offset > stop_offset {
439                    tracing::warn!(
440                        "Skipping topic {} partition {}: requested start offset {} is greater than stop offset {}",
441                        self.topic,
442                        partition,
443                        start_offset,
444                        stop_offset
445                    );
446                }
447                KafkaSplit {
448                    topic: self.topic.clone(),
449                    partition: *partition,
450                    start_offset: Some(start_offset),
451                    stop_offset: Some(stop_offset),
452                }
453            })
454            .collect::<Vec<KafkaSplit>>())
455    }
456
457    async fn fetch_stop_offset(
458        &self,
459        partitions: &[i32],
460        watermarks: &HashMap<i32, (i64, i64)>,
461    ) -> KafkaResult<HashMap<i32, Option<i64>>> {
462        match self.stop_offset {
463            KafkaEnumeratorOffset::Earliest => unreachable!(),
464            KafkaEnumeratorOffset::Latest => {
465                let mut map = HashMap::new();
466                for partition in partitions {
467                    let (_, high_watermark) = watermarks.get(partition).unwrap();
468                    map.insert(*partition, Some(*high_watermark));
469                }
470                Ok(map)
471            }
472            KafkaEnumeratorOffset::Timestamp(time) => {
473                self.fetch_offset_for_time(partitions, time, watermarks)
474                    .await
475            }
476            KafkaEnumeratorOffset::None => partitions
477                .iter()
478                .map(|partition| Ok((*partition, None)))
479                .collect(),
480        }
481    }
482
483    async fn fetch_start_offset(
484        &mut self,
485        partitions: &[i32],
486        watermarks: &HashMap<i32, (i64, i64)>,
487    ) -> KafkaResult<HashMap<i32, Option<i64>>> {
488        match self.start_offset {
489            KafkaEnumeratorOffset::Earliest | KafkaEnumeratorOffset::Latest => {
490                let mut map = HashMap::new();
491                for partition in partitions {
492                    let (low_watermark, high_watermark) = watermarks.get(partition).unwrap();
493                    let offset = match self.start_offset {
494                        KafkaEnumeratorOffset::Earliest => low_watermark - 1,
495                        KafkaEnumeratorOffset::Latest => high_watermark - 1,
496                        _ => unreachable!(),
497                    };
498                    map.insert(*partition, Some(offset));
499                }
500                Ok(map)
501            }
502            KafkaEnumeratorOffset::Timestamp(time) => {
503                let unresolved =
504                    sync_resolved_partitions(&mut self.resolved_start_offsets, partitions);
505                if !unresolved.is_empty() {
506                    let resolved = self
507                        .fetch_offset_for_time(&unresolved, time, watermarks)
508                        .await?;
509                    self.resolved_start_offsets.extend(resolved);
510                }
511                Ok(self.resolved_start_offsets.clone())
512            }
513            KafkaEnumeratorOffset::None => partitions
514                .iter()
515                .map(|partition| Ok((*partition, None)))
516                .collect(),
517        }
518    }
519
520    async fn fetch_offset_for_time(
521        &self,
522        partitions: &[i32],
523        time: i64,
524        watermarks: &HashMap<i32, (i64, i64)>,
525    ) -> KafkaResult<HashMap<i32, Option<i64>>> {
526        let mut tpl = TopicPartitionList::new();
527
528        for partition in partitions {
529            tpl.add_partition_offset(self.topic.as_str(), *partition, Offset::Offset(time))?;
530        }
531
532        let offsets = self
533            .client
534            .offsets_for_times(tpl, self.sync_call_timeout)
535            .await?;
536
537        let mut result = HashMap::with_capacity(partitions.len());
538
539        for elem in offsets.elements_for_topic(self.topic.as_str()) {
540            match elem.offset() {
541                Offset::Offset(offset) => {
542                    // XXX(rc): currently in RW source, `offset` means the last consumed offset, so we need to subtract 1
543                    result.insert(elem.partition(), Some(offset - 1));
544                }
545                Offset::End => {
546                    let (_, high_watermark) = watermarks.get(&elem.partition()).unwrap();
547                    tracing::info!(
548                        source_id = %self.context.info.source_id,
549                        "no message found before timestamp {} (ms) for partition {}, start from latest",
550                        time,
551                        elem.partition()
552                    );
553                    result.insert(elem.partition(), Some(high_watermark - 1)); // align to Latest
554                }
555                Offset::Invalid => {
556                    // special case for madsim test
557                    // For a read Kafka, it returns `Offset::Latest` when the timestamp is later than the latest message in the partition
558                    // But in madsim, it returns `Offset::Invalid`
559                    // So we align to Latest here
560                    tracing::info!(
561                        source_id = %self.context.info.source_id,
562                        "got invalid offset for partition  {} at timestamp {}, align to latest",
563                        elem.partition(),
564                        time
565                    );
566                    let (_, high_watermark) = watermarks.get(&elem.partition()).unwrap();
567                    result.insert(elem.partition(), Some(high_watermark - 1)); // align to Latest
568                }
569                Offset::Beginning => {
570                    let (low, _) = watermarks.get(&elem.partition()).unwrap();
571                    tracing::info!(
572                        source_id = %self.context.info.source_id,
573                        "all message in partition {} is after timestamp {} (ms), start from earliest",
574                        elem.partition(),
575                        time,
576                    );
577                    result.insert(elem.partition(), Some(low - 1)); // align to Earliest
578                }
579                err_offset @ Offset::Stored | err_offset @ Offset::OffsetTail(_) => {
580                    tracing::error!(
581                        source_id = %self.context.info.source_id,
582                        "got invalid offset for partition {}: {err_offset:?}",
583                        elem.partition(),
584                        err_offset = err_offset,
585                    );
586                    return Err(KafkaError::OffsetFetch(RDKafkaErrorCode::NoOffset));
587                }
588            }
589        }
590
591        Ok(result)
592    }
593
594    #[inline]
595    fn report_high_watermark(&mut self, partition: i32, offset: i64) {
596        let high_watermark_metrics =
597            self.high_watermark_metrics
598                .entry(partition)
599                .or_insert_with(|| {
600                    self.context
601                        .metrics
602                        .high_watermark
603                        .with_guarded_label_values(&[
604                            &self.context.info.source_id.to_string(),
605                            &partition.to_string(),
606                        ])
607                });
608        high_watermark_metrics.set(offset);
609    }
610
611    pub async fn check_reachability(&self) -> ConnectorResult<()> {
612        let _ = self
613            .client
614            .fetch_metadata(Some(self.topic.as_str()), self.sync_call_timeout)
615            .await?;
616        Ok(())
617    }
618
619    async fn fetch_topic_partition(&self) -> ConnectorResult<Vec<i32>> {
620        // for now, we only support one topic
621        let metadata = self
622            .client
623            .fetch_metadata(Some(self.topic.as_str()), self.sync_call_timeout)
624            .await?;
625
626        let topic_meta = match metadata.topics() {
627            [meta] => meta,
628            _ => bail!("topic {} not found", self.topic),
629        };
630
631        if topic_meta.partitions().is_empty() {
632            bail!("topic {} not found", self.topic);
633        }
634
635        Ok(topic_meta
636            .partitions()
637            .iter()
638            .map(|partition| partition.id())
639            .collect())
640    }
641}
642
643/// Drops cached start offsets of partitions that are no longer present and returns the
644/// partitions whose start offset still has to be resolved.
645fn sync_resolved_partitions(
646    resolved: &mut HashMap<i32, Option<i64>>,
647    partitions: &[i32],
648) -> Vec<i32> {
649    // `partitions` holds distinct ids, so equal length plus full containment means the set is
650    // unchanged, which is the case on almost every tick.
651    let unchanged = resolved.len() == partitions.len()
652        && partitions
653            .iter()
654            .all(|partition| resolved.contains_key(partition));
655    if unchanged {
656        return Vec::new();
657    }
658    let current: HashSet<i32> = partitions.iter().copied().collect();
659    resolved.retain(|partition, _| current.contains(partition));
660    partitions
661        .iter()
662        .copied()
663        .filter(|partition| !resolved.contains_key(partition))
664        .collect()
665}
666
667#[cfg(test)]
668mod tests {
669    use super::*;
670
671    #[test]
672    fn test_sync_resolved_partitions() {
673        let mut resolved = HashMap::new();
674
675        // First tick: every partition needs resolving.
676        assert_eq!(sync_resolved_partitions(&mut resolved, &[0, 1]), vec![0, 1]);
677        resolved.extend([(0, Some(10)), (1, Some(20))]);
678
679        // Stable partition set: nothing to resolve, cache untouched.
680        assert!(sync_resolved_partitions(&mut resolved, &[0, 1]).is_empty());
681        assert_eq!(resolved, HashMap::from([(0, Some(10)), (1, Some(20))]));
682
683        // New partition: only the new one needs resolving.
684        assert_eq!(sync_resolved_partitions(&mut resolved, &[0, 1, 2]), vec![2]);
685
686        // Removed partition: its cached entry is dropped; nothing else to resolve.
687        assert!(sync_resolved_partitions(&mut resolved, &[1]).is_empty());
688        assert_eq!(resolved, HashMap::from([(1, Some(20))]));
689    }
690}