Skip to main content

hashtree_lmdb/pool/
maintenance.rs

1use super::maintenance_batch::{
2    MoveBatchResult, MovePlan, DEFAULT_MAINTENANCE_BATCH_ITEMS, MAX_MAINTENANCE_BATCH_BYTES,
3};
4use super::{LocationRecord, PoolMaintenanceReport, PoolMemberId, PoolMemberState, PoolStore};
5use hashtree_core::store::StoreError;
6use std::collections::{BTreeMap, HashMap};
7
8impl PoolStore {
9    pub fn maintain(&self, max_items: usize) -> Result<PoolMaintenanceReport, StoreError> {
10        self.maintain_with_batch_items(max_items, DEFAULT_MAINTENANCE_BATCH_ITEMS)
11    }
12
13    pub fn maintain_with_batch_items(
14        &self,
15        max_items: usize,
16        batch_items: usize,
17    ) -> Result<PoolMaintenanceReport, StoreError> {
18        let mut report = PoolMaintenanceReport::default();
19        if max_items == 0 {
20            return Ok(report);
21        }
22        if batch_items == 0 {
23            return Err(StoreError::Other(
24                "pool maintenance batch items must be non-zero".into(),
25            ));
26        }
27
28        let cleanups = self
29            .active_move_cleanups(max_items)?
30            .into_iter()
31            .filter_map(|(hash, location)| {
32                let LocationRecord::Moving {
33                    source,
34                    target,
35                    size,
36                } = location
37                else {
38                    return None;
39                };
40                Some(MovePlan {
41                    hash,
42                    source,
43                    target,
44                    size,
45                    expected: location,
46                })
47            })
48            .collect::<Vec<_>>();
49        report.examined = report.examined.saturating_add(cleanups.len());
50        self.execute_plans_bounded(cleanups, batch_items, true, &mut report)?;
51        if report.examined >= max_items || !report.failed.is_empty() {
52            return Ok(report);
53        }
54
55        let active = self
56            .active_moves(max_items.saturating_sub(report.examined))?
57            .into_iter()
58            .filter_map(|(hash, location)| {
59                let LocationRecord::Moving {
60                    source,
61                    target,
62                    size,
63                } = location
64                else {
65                    return None;
66                };
67                Some(MovePlan {
68                    hash,
69                    source,
70                    target,
71                    size,
72                    expected: location,
73                })
74            })
75            .collect::<Vec<_>>();
76        report.examined = report.examined.saturating_add(active.len());
77        self.execute_plans_bounded(active, batch_items, false, &mut report)?;
78        if report.examined >= max_items || !report.failed.is_empty() {
79            return Ok(report);
80        }
81
82        let draining = self
83            .read_manifest()?
84            .members
85            .into_iter()
86            .filter(|member| member.state == PoolMemberState::Draining)
87            .map(|member| member.id)
88            .collect::<Vec<_>>();
89
90        for source in draining {
91            let hashes = self.member_hashes(source, max_items.saturating_sub(report.examined))?;
92            let mut reserved_bytes = HashMap::new();
93            let mut plans = Vec::with_capacity(hashes.len());
94            for hash in hashes {
95                if report.examined >= max_items {
96                    break;
97                }
98                report.examined += 1;
99                let Some(location) = self.read_location(&hash)? else {
100                    continue;
101                };
102                let (target, size) = match location {
103                    LocationRecord::Pending { member, size }
104                    | LocationRecord::Stored { member, size }
105                        if member == source =>
106                    {
107                        match self.choose_write_member_with_reserved(
108                            size,
109                            Some(source),
110                            &reserved_bytes,
111                        ) {
112                            Ok(target) => {
113                                let reserved = reserved_bytes.entry(target).or_insert(0u64);
114                                *reserved = reserved.saturating_add(size);
115                                (target, size)
116                            }
117                            Err(error) => {
118                                report.failed.push(format!("{hash:?}: {error}"));
119                                continue;
120                            }
121                        }
122                    }
123                    LocationRecord::Moving {
124                        source: actual_source,
125                        target,
126                        size,
127                    } if actual_source == source => (target, size),
128                    _ => continue,
129                };
130                plans.push(MovePlan {
131                    hash,
132                    source,
133                    target,
134                    size,
135                    expected: location,
136                });
137            }
138            self.execute_plans_bounded(plans, batch_items, false, &mut report)?;
139            if report.examined >= max_items || !report.failed.is_empty() {
140                return Ok(report);
141            }
142        }
143
144        while report.examined < max_items {
145            let Some((source, target)) = self.rebalance_pair()? else {
146                break;
147            };
148            let hashes = self.member_hashes(source, max_items - report.examined)?;
149            if hashes.is_empty() {
150                break;
151            }
152            let mut progressed = false;
153            for hash in hashes {
154                if report.examined >= max_items {
155                    break;
156                }
157                report.examined += 1;
158                let Some(location) = self.read_location(&hash)? else {
159                    continue;
160                };
161                if !self.move_improves_balance(source, target, location.size())? {
162                    continue;
163                }
164                match self.move_blob(source, target, hash) {
165                    Ok(Some(bytes)) => {
166                        report.moved += 1;
167                        report.bytes_moved = report.bytes_moved.saturating_add(bytes);
168                        progressed = true;
169                    }
170                    Ok(None) => {}
171                    Err(error) => report.failed.push(format!("{hash:?}: {error}")),
172                }
173            }
174            if !progressed {
175                break;
176            }
177        }
178        Ok(report)
179    }
180
181    fn execute_plans_bounded(
182        &self,
183        plans: Vec<MovePlan>,
184        batch_items: usize,
185        cleanup_only: bool,
186        report: &mut PoolMaintenanceReport,
187    ) -> Result<(), StoreError> {
188        let mut groups = BTreeMap::<(PoolMemberId, PoolMemberId), Vec<MovePlan>>::new();
189        for plan in plans {
190            groups
191                .entry((plan.source, plan.target))
192                .or_default()
193                .push(plan);
194        }
195        for (_, plans) in groups {
196            let mut batch = Vec::new();
197            let mut batch_bytes = 0u64;
198            for plan in plans {
199                let exceeds_items = batch.len() >= batch_items;
200                let exceeds_bytes = !batch.is_empty()
201                    && batch_bytes.saturating_add(plan.size) > MAX_MAINTENANCE_BATCH_BYTES;
202                if exceeds_items || exceeds_bytes {
203                    self.execute_one_bounded_batch(&batch, cleanup_only, report)?;
204                    batch.clear();
205                    batch_bytes = 0;
206                }
207                if !cleanup_only && plan.size > MAX_MAINTENANCE_BATCH_BYTES {
208                    match self.move_blob(plan.source, plan.target, plan.hash) {
209                        Ok(Some(bytes)) => {
210                            report.moved += 1;
211                            report.bytes_moved = report.bytes_moved.saturating_add(bytes);
212                        }
213                        Ok(None) => {}
214                        Err(error) => report.failed.push(format!("{:?}: {error}", plan.hash)),
215                    }
216                    continue;
217                }
218                batch_bytes = batch_bytes.saturating_add(plan.size);
219                batch.push(plan);
220            }
221            self.execute_one_bounded_batch(&batch, cleanup_only, report)?;
222        }
223        Ok(())
224    }
225
226    fn execute_one_bounded_batch(
227        &self,
228        plans: &[MovePlan],
229        cleanup_only: bool,
230        report: &mut PoolMaintenanceReport,
231    ) -> Result<(), StoreError> {
232        if plans.is_empty() {
233            return Ok(());
234        }
235        let result = if cleanup_only {
236            self.execute_move_cleanup_batch(plans)?
237        } else {
238            self.execute_move_batch(plans)?
239        };
240        apply_batch_result(report, result);
241        Ok(())
242    }
243}
244
245fn apply_batch_result(report: &mut PoolMaintenanceReport, result: MoveBatchResult) {
246    for plan in result.moved {
247        report.moved = report.moved.saturating_add(1);
248        report.bytes_moved = report.bytes_moved.saturating_add(plan.size);
249    }
250    report.failed.extend(
251        result
252            .failed
253            .into_iter()
254            .map(|(hash, error)| format!("{hash:?}: {error}")),
255    );
256}