From 2d7c0a6087e1d70e739774c0b9a34e9573a223fe Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Arma=C3=ABl=20Gu=C3=A9neau?= Date: Wed, 13 May 2026 14:37:05 +0200 Subject: [PATCH] table: also apply repair-on-read the first time a value is set --- src/table/table.rs | 64 +++++++++++++++++++++++++++++++--------------- 1 file changed, 43 insertions(+), 21 deletions(-) diff --git a/src/table/table.rs b/src/table/table.rs index 8ddd8378..3fc0d5ee 100644 --- a/src/table/table.rs +++ b/src/table/table.rs @@ -326,25 +326,34 @@ impl Table { let mut ret = None; let mut not_all_same = false; - for resp in resps { - if let TableRpc::ReadEntryResponse(value) = resp { - if let Some(v_bytes) = value { - let v = self.data.decode_entry(v_bytes.as_slice())?; - ret = match ret { - None => Some(v), - Some(mut x) => { - if x != v { - not_all_same = true; - x.merge(&v); + { + let mut vals_nb = 0; + for resp in &resps { + if let TableRpc::ReadEntryResponse(value) = resp { + if let Some(v_bytes) = value { + vals_nb += 1; + let v = self.data.decode_entry(v_bytes.as_slice())?; + ret = match ret { + None => Some(v), + Some(mut x) => { + if x != v { + not_all_same = true; + x.merge(&v); + } + Some(x) } - Some(x) } } + } else { + return Err(Error::Message("Invalid return value to read".to_string())); } - } else { - return Err(Error::Message("Invalid return value to read".to_string())); + } + // Only some nodes store this value; we must propagate it during repair + if vals_nb < resps.len() { + not_all_same = true; } } + if let Some(ret_entry) = &ret { if not_all_same { let self2 = self.clone(); @@ -421,11 +430,26 @@ impl Table { let mut ret: BTreeMap, F::E> = BTreeMap::new(); let mut to_repair = BTreeSet::new(); - for resp in resps { - if let TableRpc::Update(entries) = resp { - for entry_bytes in entries.iter() { - let entry = self.data.decode_entry(entry_bytes.as_slice())?; - let entry_key = self.data.tree_key(entry.partition_key(), entry.sort_key()); + { + let mut all_entries: BTreeMap, Vec> = BTreeMap::new(); + for resp in &resps { + if let TableRpc::Update(entries) = resp { + for entry_bytes in entries.iter() { + let entry = self.data.decode_entry(entry_bytes.as_slice())?; + let entry_key = self.data.tree_key(entry.partition_key(), entry.sort_key()); + all_entries.entry(entry_key).or_default().push(entry); + } + } else { + return Err(Error::unexpected_rpc_message(resp)); + } + } + for (entry_key, entries) in all_entries { + // Only some nodes store this entry; we must propagate it during repair + if entries.len() < resps.len() { + to_repair.insert(entry_key.clone()); + } + // Merge all entries for this key together + for entry in entries { match ret.get_mut(&entry_key) { Some(e) => { if *e != entry { @@ -434,12 +458,10 @@ impl Table { } } None => { - ret.insert(entry_key, entry); + ret.insert(entry_key.clone(), entry); } } } - } else { - return Err(Error::unexpected_rpc_message(resp)); } }