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}