Skip to main content

risingwave_meta/serving/
mod.rs

1// Copyright 2023 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;
17
18use parking_lot::RwLock;
19use risingwave_common::bitmap::Bitmap;
20use risingwave_common::hash::WorkerSlotMapping;
21use risingwave_common::vnode_mapping::vnode_placement::place_vnode;
22use risingwave_meta_model::{TableId, WorkerId};
23use risingwave_pb::common::{WorkerNode, WorkerType};
24use risingwave_pb::meta::serving_table_vnode_mappings::PbServingTableVnodeMapping;
25use risingwave_pb::meta::subscribe_response::{Info, Operation};
26use risingwave_pb::meta::table_fragments::fragment::FragmentDistributionType;
27use risingwave_pb::meta::{
28    FragmentWorkerSlotMapping, FragmentWorkerSlotMappings, PbServingTableVnodeMappings,
29    PbTableRefillRuntimeConfig,
30};
31use tokio::sync::oneshot::Sender;
32use tokio::task::JoinHandle;
33
34use crate::MetaResult;
35use crate::controller::fragment::FragmentServingInfo;
36use crate::controller::session_params::SessionParamsControllerRef;
37use crate::manager::{LocalNotification, MetadataManager, NotificationManagerRef, WorkerKey};
38use crate::model::FragmentId;
39
40pub type ServingVnodeMappingRef = Arc<ServingVnodeMapping>;
41
42#[derive(Default)]
43pub struct ServingVnodeMapping {
44    serving_vnode_mappings: RwLock<HashMap<FragmentId, WorkerSlotMapping>>,
45}
46
47impl ServingVnodeMapping {
48    pub fn all(&self) -> HashMap<FragmentId, WorkerSlotMapping> {
49        self.serving_vnode_mappings.read().clone()
50    }
51
52    pub(crate) fn table_vnode_mappings_by_worker(
53        &self,
54        worker_ids: impl IntoIterator<Item = WorkerId>,
55        fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
56    ) -> HashMap<WorkerId, HashMap<TableId, Bitmap>> {
57        let mut result = worker_ids
58            .into_iter()
59            .map(|worker_id| (worker_id, HashMap::new()))
60            .collect::<HashMap<_, _>>();
61        let serving_vnode_mappings = self.all();
62
63        for (fragment_id, info) in fragment_serving_infos {
64            let Some(table_id) = info.result_table_id else {
65                continue;
66            };
67            let Some(mapping) = serving_vnode_mappings.get(fragment_id) else {
68                continue;
69            };
70
71            for (worker_slot_id, bitmap) in mapping.to_bitmaps() {
72                let Some(table_mappings) = result.get_mut(&worker_slot_id.worker_id()) else {
73                    continue;
74                };
75                table_mappings
76                    .entry(table_id)
77                    .and_modify(|current| *current |= &bitmap)
78                    .or_insert(bitmap);
79            }
80        }
81        result
82    }
83
84    /// Upsert mapping for given fragments according to the latest `workers`.
85    /// Returns (successful updates, failed updates).
86    pub fn upsert(
87        &self,
88        fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
89        workers: &[WorkerNode],
90        max_serving_parallelism: Option<u64>,
91    ) -> (HashMap<FragmentId, WorkerSlotMapping>, HashSet<FragmentId>) {
92        let mut serving_vnode_mappings = self.serving_vnode_mappings.write();
93        Self::upsert_locked(
94            &mut serving_vnode_mappings,
95            fragment_serving_infos,
96            workers,
97            max_serving_parallelism,
98        )
99    }
100
101    /// Rebuild mappings from a complete catalog snapshot.
102    ///
103    /// Unlike [`Self::upsert`], this removes mappings for fragments absent from the snapshot.
104    pub(crate) fn reconcile(
105        &self,
106        fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
107        workers: &[WorkerNode],
108        max_serving_parallelism: Option<u64>,
109    ) -> (HashMap<FragmentId, WorkerSlotMapping>, HashSet<FragmentId>) {
110        let mut serving_vnode_mappings = self.serving_vnode_mappings.write();
111        serving_vnode_mappings
112            .retain(|fragment_id, _| fragment_serving_infos.contains_key(fragment_id));
113        Self::upsert_locked(
114            &mut serving_vnode_mappings,
115            fragment_serving_infos,
116            workers,
117            max_serving_parallelism,
118        )
119    }
120
121    fn upsert_locked(
122        serving_vnode_mappings: &mut HashMap<FragmentId, WorkerSlotMapping>,
123        fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
124        workers: &[WorkerNode],
125        max_serving_parallelism: Option<u64>,
126    ) -> (HashMap<FragmentId, WorkerSlotMapping>, HashSet<FragmentId>) {
127        let mut upserted: HashMap<FragmentId, WorkerSlotMapping> = HashMap::default();
128        let mut failed: HashSet<FragmentId> = HashSet::default();
129        for (fragment_id, info) in fragment_serving_infos {
130            let new_mapping = {
131                let old_mapping = serving_vnode_mappings.get(fragment_id);
132                let max_parallelism = match info.distribution_type {
133                    FragmentDistributionType::Unspecified => unreachable!(),
134                    FragmentDistributionType::Single => Some(1),
135                    FragmentDistributionType::Hash => None,
136                }
137                .or_else(|| max_serving_parallelism.map(|p| p as usize));
138                place_vnode(old_mapping, workers, max_parallelism, info.vnode_count)
139            };
140            match new_mapping {
141                None => {
142                    serving_vnode_mappings.remove(fragment_id);
143                    failed.insert(*fragment_id);
144                }
145                Some(mapping) => {
146                    serving_vnode_mappings.insert(*fragment_id, mapping.clone());
147                    upserted.insert(*fragment_id, mapping);
148                }
149            }
150        }
151        (upserted, failed)
152    }
153
154    fn remove(&self, fragment_ids: &[FragmentId]) {
155        let mut mappings = self.serving_vnode_mappings.write();
156        for fragment_id in fragment_ids {
157            mappings.remove(fragment_id);
158        }
159    }
160}
161
162pub(crate) fn to_fragment_worker_slot_mapping(
163    mappings: &HashMap<FragmentId, WorkerSlotMapping>,
164) -> Vec<FragmentWorkerSlotMapping> {
165    mappings
166        .iter()
167        .map(|(&fragment_id, mapping)| FragmentWorkerSlotMapping {
168            fragment_id,
169            mapping: Some(mapping.to_protobuf()),
170        })
171        .collect()
172}
173
174pub(crate) fn to_deleted_fragment_worker_slot_mapping(
175    fragment_ids: impl Iterator<Item = FragmentId>,
176) -> Vec<FragmentWorkerSlotMapping> {
177    fragment_ids
178        .map(|fragment_id| FragmentWorkerSlotMapping {
179            fragment_id,
180            mapping: None,
181        })
182        .collect()
183}
184
185pub async fn on_meta_start(
186    notification_manager: NotificationManagerRef,
187    metadata_manager: &MetadataManager,
188    serving_vnode_mapping: ServingVnodeMappingRef,
189    max_serving_parallelism: Option<u64>,
190) {
191    let (serving_compute_nodes, fragment_serving_infos) = fetch_serving_infos(metadata_manager)
192        .await
193        .expect("fail to fetch serving infos");
194    let (mappings, failed) = serving_vnode_mapping.reconcile(
195        &fragment_serving_infos,
196        &serving_compute_nodes,
197        max_serving_parallelism,
198    );
199    tracing::debug!(
200        "Initialize serving vnode mapping snapshot for fragments {:?}.",
201        mappings.keys()
202    );
203    if !failed.is_empty() {
204        tracing::warn!(
205            "Fail to update serving vnode mapping for fragments {:?}.",
206            failed
207        );
208    }
209    notification_manager.notify_frontend_without_version(
210        Operation::Snapshot,
211        Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings {
212            mappings: to_fragment_worker_slot_mapping(&mappings),
213        }),
214    );
215}
216
217pub(crate) async fn fetch_serving_infos(
218    metadata_manager: &MetadataManager,
219) -> MetaResult<(Vec<WorkerNode>, HashMap<FragmentId, FragmentServingInfo>)> {
220    let fragment_serving_infos = metadata_manager
221        .catalog_controller
222        .fragment_serving_infos()
223        .await?;
224    let serving_compute_nodes = metadata_manager
225        .cluster_controller
226        .list_active_serving_workers()
227        .await?;
228    Ok((
229        serving_compute_nodes,
230        fragment_serving_infos
231            .into_iter()
232            .map(|(fragment_id, info)| (fragment_id as FragmentId, info))
233            .collect(),
234    ))
235}
236
237pub(crate) async fn sync_serving_table_vnode_mappings_to_hummock(
238    notification_manager: &NotificationManagerRef,
239    serving_vnode_mapping: &ServingVnodeMappingRef,
240    workers: &[WorkerNode],
241    fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
242) {
243    let mut table_vnode_mappings = serving_vnode_mapping.table_vnode_mappings_by_worker(
244        workers.iter().map(|worker| worker.id),
245        fragment_serving_infos,
246    );
247    for worker in workers {
248        let Some(host) = worker.host.clone() else {
249            tracing::warn!(worker_id = %worker.id, "serving worker host not found");
250            continue;
251        };
252        let table_vnode_mapping = table_vnode_mappings
253            .remove(&worker.id)
254            .expect("requested worker must have a table vnode mapping");
255        let config = PbTableRefillRuntimeConfig {
256            serving_table_vnode_mappings: Some(to_pb_serving_table_vnode_mappings(
257                &table_vnode_mapping,
258            )),
259            ..Default::default()
260        };
261        notification_manager
262            .notify_hummock_targeted_update(WorkerKey(host), Info::TableRefillRuntimeConfig(config))
263            .await;
264    }
265}
266
267pub fn start_serving_vnode_mapping_worker(
268    notification_manager: NotificationManagerRef,
269    metadata_manager: MetadataManager,
270    serving_vnode_mapping: ServingVnodeMappingRef,
271    session_params: SessionParamsControllerRef,
272) -> (JoinHandle<()>, Sender<()>) {
273    let (local_notification_tx, mut local_notification_rx) = tokio::sync::mpsc::unbounded_channel();
274    let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
275    notification_manager.insert_local_sender(local_notification_tx);
276    let join_handle = tokio::spawn(async move {
277        let reset = || async {
278            let (workers, fragment_serving_infos) = fetch_serving_infos(&metadata_manager)
279                .await
280                .expect("fail to fetch serving infos");
281            let max_serving_parallelism = session_params
282                .get_params()
283                .await
284                .batch_parallelism()
285                .map(|p| p.get());
286            let (mappings, failed) = serving_vnode_mapping.reconcile(
287                &fragment_serving_infos,
288                &workers,
289                max_serving_parallelism,
290            );
291            tracing::debug!(
292                "Update serving vnode mapping snapshot for fragments {:?}.",
293                mappings.keys()
294            );
295            if !failed.is_empty() {
296                tracing::warn!(
297                    "Fail to update serving vnode mapping for fragments {:?}.",
298                    failed
299                );
300            }
301            notification_manager.notify_frontend_without_version(
302                Operation::Snapshot,
303                Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings {
304                    mappings: to_fragment_worker_slot_mapping(&mappings),
305                }),
306            );
307            sync_serving_table_vnode_mappings_to_hummock(
308                &notification_manager,
309                &serving_vnode_mapping,
310                &workers,
311                &fragment_serving_infos,
312            )
313            .await;
314        };
315        loop {
316            tokio::select! {
317                notification = local_notification_rx.recv() => {
318                    match notification {
319                        Some(notification) => {
320                            match notification {
321                                LocalNotification::WorkerNodeActivated(w) | LocalNotification::WorkerNodeDeleted(w) =>  {
322                                    if w.r#type() != WorkerType::ComputeNode || !w.property.as_ref().is_some_and(|p| p.is_serving) {
323                                        continue;
324                                    }
325                                    reset().await;
326                                }
327                                LocalNotification::BatchParallelismChange => {
328                                    reset().await;
329                                }
330                                LocalNotification::ServingFragmentMappingsUpsert(fragment_ids) => {
331                                    if fragment_ids.is_empty() {
332                                        continue;
333                                    }
334                                    let (workers, fragment_serving_infos) = fetch_serving_infos(&metadata_manager)
335                                        .await
336                                        .expect("fail to fetch serving infos");
337                                    let filtered_fragment_serving_infos = fragment_ids.iter().filter_map(|frag_id| {
338                                        match fragment_serving_infos.get(frag_id) {
339                                            Some(info) => Some((*frag_id, info.clone())),
340                                            None => {
341                                                tracing::warn!(fragment_id = %frag_id, "fragment serving info not found");
342                                                None
343                                            }
344                                        }
345                                    }).collect();
346                                    let max_serving_parallelism = session_params
347                                        .get_params()
348                                        .await
349                                        .batch_parallelism()
350                                        .map(|p|p.get());
351                                    let (upserted, failed) = serving_vnode_mapping.upsert(&filtered_fragment_serving_infos, &workers, max_serving_parallelism);
352                                    if !upserted.is_empty() {
353                                        tracing::debug!("Update serving vnode mapping for fragments {:?}.", upserted.keys());
354                                        notification_manager.notify_frontend_without_version(Operation::Update, Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings{ mappings: to_fragment_worker_slot_mapping(&upserted) }));
355                                    }
356                                    if !failed.is_empty() {
357                                        tracing::warn!("Fail to update serving vnode mapping for fragments {:?}.", failed);
358                                        notification_manager.notify_frontend_without_version(Operation::Delete, Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings{ mappings: to_deleted_fragment_worker_slot_mapping(failed.iter().cloned())}));
359                                    }
360                                    sync_serving_table_vnode_mappings_to_hummock(
361                                        &notification_manager,
362                                        &serving_vnode_mapping,
363                                        &workers,
364                                        &fragment_serving_infos,
365                                    )
366                                    .await;
367                                }
368                                LocalNotification::ServingFragmentMappingsDelete(fragment_ids) => {
369                                    if fragment_ids.is_empty() {
370                                        continue;
371                                    }
372
373                                    tracing::debug!("Delete serving vnode mapping for fragments {:?}.", fragment_ids);
374
375                                    serving_vnode_mapping.remove(&fragment_ids);
376                                    notification_manager.notify_frontend_without_version(Operation::Delete, Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings{ mappings: to_deleted_fragment_worker_slot_mapping(fragment_ids.iter().cloned()) }));
377                                    let (workers, fragment_serving_infos) = fetch_serving_infos(&metadata_manager)
378                                        .await
379                                        .expect("fail to fetch serving infos");
380                                    sync_serving_table_vnode_mappings_to_hummock(
381                                        &notification_manager,
382                                        &serving_vnode_mapping,
383                                        &workers,
384                                        &fragment_serving_infos,
385                                    )
386                                    .await;
387                                }
388                                _ => {}
389                            }
390                        }
391                        None => {
392                            return;
393                        }
394                    }
395                }
396                _ = &mut shutdown_rx => {
397                    return;
398                }
399            }
400        }
401    });
402    (join_handle, shutdown_tx)
403}
404
405pub(crate) fn to_pb_serving_table_vnode_mappings(
406    table_vnode_mapping: &HashMap<TableId, Bitmap>,
407) -> PbServingTableVnodeMappings {
408    PbServingTableVnodeMappings {
409        mappings: table_vnode_mapping
410            .iter()
411            .map(|(table_id, bitmap)| PbServingTableVnodeMapping {
412                table_id: table_id.as_raw_id(),
413                bitmap: Some(bitmap.to_protobuf()),
414            })
415            .collect(),
416    }
417}
418
419#[cfg(test)]
420mod tests {
421    use risingwave_common::hash::{VirtualNode, WorkerSlotId};
422    use risingwave_pb::common::{HostAddress, worker_node};
423    use risingwave_pb::meta::subscribe_response::Info;
424    use risingwave_pb::meta::{SubscribeResponse, SubscribeType};
425    use tokio::sync::mpsc;
426
427    use super::*;
428    use crate::controller::SqlMetaStore;
429    use crate::manager::NotificationManager;
430
431    fn serving_worker(id: u32) -> WorkerNode {
432        WorkerNode {
433            id: id.into(),
434            r#type: WorkerType::ComputeNode.into(),
435            host: Some(HostAddress {
436                host: format!("host{}", id),
437                port: id as i32,
438            }),
439            state: worker_node::State::Running as i32,
440            property: Some(worker_node::Property {
441                is_serving: true,
442                parallelism: 1,
443                ..Default::default()
444            }),
445            ..Default::default()
446        }
447    }
448
449    #[test]
450    fn test_table_vnode_mappings_merge_worker_slots() {
451        let worker1 = WorkerId::new(1);
452        let worker2 = WorkerId::new(2);
453        let worker3 = WorkerId::new(3);
454        let slot1 = WorkerSlotId::new(worker1, 0);
455        let slot2 = WorkerSlotId::new(worker1, 1);
456        let slot3 = WorkerSlotId::new(worker2, 0);
457        let fragment_id = FragmentId::new(233);
458        let table_id = TableId::new(234);
459        let mapping = WorkerSlotMapping::new_uniform(
460            [slot1, slot2, slot3].into_iter(),
461            VirtualNode::COUNT_FOR_TEST,
462        );
463        let slot_bitmaps = mapping.to_bitmaps();
464        let serving_vnode_mapping = ServingVnodeMapping {
465            serving_vnode_mappings: RwLock::new(HashMap::from([(fragment_id, mapping)])),
466        };
467        let fragment_serving_infos = HashMap::from([(
468            fragment_id,
469            FragmentServingInfo {
470                result_table_id: Some(table_id),
471                distribution_type: FragmentDistributionType::Hash,
472                vnode_count: VirtualNode::COUNT_FOR_TEST,
473            },
474        )]);
475
476        let mappings = serving_vnode_mapping
477            .table_vnode_mappings_by_worker([worker1, worker2, worker3], &fragment_serving_infos);
478        let mut worker1_bitmap = slot_bitmaps[&slot1].clone();
479        worker1_bitmap |= &slot_bitmaps[&slot2];
480        assert_eq!(mappings[&worker1][&table_id], worker1_bitmap);
481        assert_eq!(mappings[&worker2][&table_id], slot_bitmaps[&slot3]);
482        assert!(mappings[&worker3].is_empty());
483    }
484
485    #[tokio::test]
486    async fn test_sync_serving_table_vnode_mappings_to_hummock_targets_only_result_tables() {
487        let notification_manager =
488            Arc::new(NotificationManager::new(SqlMetaStore::for_test().await).await);
489        let worker1 = serving_worker(1);
490        let worker2 = serving_worker(2);
491        let workers = vec![worker1.clone(), worker2.clone()];
492        let worker_key1 = WorkerKey(worker1.host.clone().unwrap());
493        let worker_key2 = WorkerKey(worker2.host.clone().unwrap());
494        let (tx1, mut rx1) = mpsc::unbounded_channel();
495        let (tx2, mut rx2) = mpsc::unbounded_channel();
496        notification_manager.insert_sender(SubscribeType::Hummock, worker_key1, tx1);
497        notification_manager.insert_sender(SubscribeType::Hummock, worker_key2, tx2);
498
499        let table_id = TableId::new(233);
500        let fragment_id = FragmentId::new(234);
501        let private_fragment_id = FragmentId::new(235);
502        let serving_vnode_mapping = Arc::new(ServingVnodeMapping::default());
503        let fragment_serving_infos = HashMap::from([
504            (
505                fragment_id,
506                FragmentServingInfo {
507                    result_table_id: Some(table_id),
508                    distribution_type: FragmentDistributionType::Hash,
509                    vnode_count: VirtualNode::COUNT_FOR_TEST,
510                },
511            ),
512            (
513                private_fragment_id,
514                FragmentServingInfo {
515                    result_table_id: None,
516                    distribution_type: FragmentDistributionType::Hash,
517                    vnode_count: VirtualNode::COUNT_FOR_TEST,
518                },
519            ),
520        ]);
521        let (upserted, failed) =
522            serving_vnode_mapping.upsert(&fragment_serving_infos, &workers, None);
523        assert_eq!(upserted.len(), 2);
524        assert!(failed.is_empty());
525        let expected_bitmaps = upserted[&fragment_id]
526            .to_bitmaps()
527            .into_iter()
528            .map(|(worker_slot_id, bitmap)| (worker_slot_id.worker_id(), bitmap))
529            .collect::<HashMap<_, _>>();
530        // Give the non-result fragment a distinct mapping. If it were incorrectly
531        // associated with `table_id`, it would expand the table's ownership to all vnodes.
532        serving_vnode_mapping.serving_vnode_mappings.write().insert(
533            private_fragment_id,
534            WorkerSlotMapping::new_single(WorkerSlotId::new(worker1.id, 0)),
535        );
536
537        sync_serving_table_vnode_mappings_to_hummock(
538            &notification_manager,
539            &serving_vnode_mapping,
540            &workers,
541            &fragment_serving_infos,
542        )
543        .await;
544
545        let response1 = rx1.recv().await.unwrap().unwrap();
546        let response2 = rx2.recv().await.unwrap().unwrap();
547        assert!(rx1.try_recv().is_err());
548        assert!(rx2.try_recv().is_err());
549
550        let serving_bitmap = |response: SubscribeResponse| {
551            let Some(Info::TableRefillRuntimeConfig(config)) = response.info else {
552                panic!("expect table refill runtime config");
553            };
554            let mappings = config.serving_table_vnode_mappings.unwrap().mappings;
555            assert_eq!(mappings.len(), 1);
556            assert_eq!(mappings[0].table_id, table_id.as_raw_id());
557            Bitmap::from(mappings[0].bitmap.clone().unwrap())
558        };
559        let bitmap1 = serving_bitmap(response1);
560        let bitmap2 = serving_bitmap(response2);
561        assert_eq!(bitmap1, expected_bitmaps[&worker1.id]);
562        assert_eq!(bitmap2, expected_bitmaps[&worker2.id]);
563        assert_eq!(
564            bitmap1.count_ones() + bitmap2.count_ones(),
565            VirtualNode::COUNT_FOR_TEST
566        );
567    }
568
569    fn hash_serving_info() -> FragmentServingInfo {
570        FragmentServingInfo {
571            result_table_id: None,
572            distribution_type: FragmentDistributionType::Hash,
573            vnode_count: VirtualNode::COUNT_FOR_TEST,
574        }
575    }
576
577    #[test]
578    fn test_reconcile_exactly_matches_full_snapshot_and_removes_failed_placements() {
579        let mapping = ServingVnodeMapping::default();
580        let worker = serving_worker(1);
581        let stale_fragment = FragmentId::new(1);
582        let retained_fragment = FragmentId::new(2);
583        let added_fragment = FragmentId::new(3);
584
585        let initial_snapshot = HashMap::from([
586            (stale_fragment, hash_serving_info()),
587            (retained_fragment, hash_serving_info()),
588        ]);
589        mapping.upsert(&initial_snapshot, std::slice::from_ref(&worker), None);
590        assert!(mapping.all().contains_key(&stale_fragment));
591        let retained_mapping = mapping.all()[&retained_fragment].clone();
592
593        let current_snapshot = HashMap::from([
594            (retained_fragment, hash_serving_info()),
595            (added_fragment, hash_serving_info()),
596        ]);
597        let (reconciled, failed) = mapping.reconcile(&current_snapshot, &[worker], None);
598        assert!(failed.is_empty());
599        assert_eq!(reconciled.len(), 2);
600        assert!(reconciled.contains_key(&retained_fragment));
601        assert!(reconciled.contains_key(&added_fragment));
602
603        let mappings = mapping.all();
604        assert_eq!(mappings.len(), 2);
605        assert!(mappings.contains_key(&retained_fragment));
606        assert!(mappings.contains_key(&added_fragment));
607        assert_eq!(mappings[&retained_fragment], retained_mapping);
608
609        let final_snapshot = HashMap::from([(retained_fragment, hash_serving_info())]);
610        let (reconciled, failed) = mapping.reconcile(&final_snapshot, &[], None);
611        assert!(reconciled.is_empty());
612        assert_eq!(failed, HashSet::from([retained_fragment]));
613        assert!(mapping.all().is_empty());
614    }
615
616    #[test]
617    fn test_upsert_keeps_fragments_absent_from_incremental_update() {
618        let mapping = ServingVnodeMapping::default();
619        let worker = serving_worker(1);
620        let first_fragment = FragmentId::new(1);
621        let second_fragment = FragmentId::new(2);
622
623        let first_update = HashMap::from([(first_fragment, hash_serving_info())]);
624        mapping.upsert(&first_update, std::slice::from_ref(&worker), None);
625        let second_update = HashMap::from([(second_fragment, hash_serving_info())]);
626        mapping.upsert(&second_update, &[worker], None);
627
628        let mappings = mapping.all();
629        assert!(mappings.contains_key(&first_fragment));
630        assert!(mappings.contains_key(&second_fragment));
631    }
632}