diff --git a/changelog.d/post-070-review.md b/changelog.d/post-070-review.md new file mode 100644 index 000000000..637af182b --- /dev/null +++ b/changelog.d/post-070-review.md @@ -0,0 +1,34 @@ +Fixed + +- **Esplora `after_txid` is 422 when that tx is not in the script's history.** + A cursor that exists somewhere else on the chain used to restart page 1. + `/txs`, `/txs/chain`, `/txs/summary`, the address routes, and a multi + POST now return `after_txid not found` and no rows. +- **`estimatesmartfee`, `estimaterawfee`, and Electrum `blockchain.estimatefee` + use the 2-block rate for target 2.** Target 0 is still the next-block + horizon. Target 1 stays the 1-block rate. +- **A shared orphan survives the other announcer's reserve.** Evicting + peer B drops only B. Peer A's copy is still delivered when the parent + arrives. +- **A tip shrink clamps the spend-durable marker.** Open revalidation and + spend replay lower a marker that sits above the surviving tip, so a + later confirm still annotates spends. A checkpoint cannot publish the + pre-disconnect height over that clamp. +- **A newest-first scripthash page stops at the page edge.** An unspent + tail no longer re-reads every older create. +- **A tip or compact block with a repeated transaction pair is not a + block.** The merkle root can still match (CVE-2012-2459). Tip follow + disconnects that peer. Compact reconstruction returns the hash to + `getdata`. +- **`submitblock` of a sibling that spends a coin the tip also spent is + inconclusive.** That header is not cached as `duplicate-invalid`. +- **A refused local I2P SAM port does not rotate the session.** The dial + error is no longer classified as a dead `STREAM CONNECT`. A SAM reply + of `INVALID_ID` still is. +- **Outbound dial keeps one onion or I2P seat when clearnet fills the + batch.** A dead overlay peer is recorded on its real address. An + unspecified version socket is not inserted into addrman. +- **An Esplora singleflight join does not put an older scripthash back + over a newer last-1** for the same client. That includes a leader that + is still inside its handler when the newer script finishes, and a waiter + that resumes after it. diff --git a/crates/rbitcoin-consensus/src/confirm_run/write.rs b/crates/rbitcoin-consensus/src/confirm_run/write.rs index 0d5784027..28e030c46 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/write.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/write.rs @@ -399,6 +399,20 @@ pub fn replay_spend_annotations(query: &Query) -> Result { let Some(tip) = query.tip_height().map(|h| h.0) else { return Ok(0); }; + // A marker above the tip is not a cursor for this chain. Shrink and a + // crash between the tip flush and the clamp both leave that file behind. + // Lower it before the `a == tip` skip, or the next confirm's spends are + // never rewritten. + let annotated = query + .store() + .spend_annotated_through() + .map_err(ConsensusError::from)?; + if annotated.is_some_and(|h| h > tip) { + query + .store() + .clamp_spend_durable() + .map_err(ConsensusError::from)?; + } let annotated = query .store() .spend_annotated_through() diff --git a/crates/rbitcoin-electrum/src/server.rs b/crates/rbitcoin-electrum/src/server.rs index 775bb219d..9657536a3 100644 --- a/crates/rbitcoin-electrum/src/server.rs +++ b/crates/rbitcoin-electrum/src/server.rs @@ -1785,14 +1785,17 @@ fn dispatch_pinned( sh_join, |q, view| { q.scripthash_history_filtered_in(&sh, &filter, view) + .map(|page| page.rows) .map_err(|e| e.to_string()) }, |q, slot, view| { q.scripthash_history_filtered_slot_in(&sh, &filter, slot, view) + .map(|page| page.rows) .map_err(|e| e.to_string()) }, |q, slot| { q.scripthash_history_filtered_slot(&sh, &filter, slot) + .map(|page| page.rows) .map_err(|e| e.to_string()) }, )?; diff --git a/crates/rbitcoin-esplora/src/esplora_http_journey.rs b/crates/rbitcoin-esplora/src/esplora_http_journey.rs index 5654e7add..fcb644589 100644 --- a/crates/rbitcoin-esplora/src/esplora_http_journey.rs +++ b/crates/rbitcoin-esplora/src/esplora_http_journey.rs @@ -112,6 +112,27 @@ async fn wallet_pages_after_txid(addr: SocketAddr, sh1: &str) { } let (st, body) = http_get(addr, &format!("/scripthash/{sh1}/txs?after_txid=zz")).await; assert_eq!(st, 422, "{body}"); + let (st, body) = http_get(addr, &format!("/scripthash/{sh1}/txs/chain/{t3}")).await; + assert_eq!(st, 200, "{body}"); + let page: Vec = serde_json::from_str(&body).unwrap(); + let ids: Vec<&str> = page.iter().filter_map(|v| v["txid"].as_str()).collect(); + assert_eq!(ids, vec![t1.as_str()]); + + foreign_chain_cursor_is_not_this_page(addr, sh1, &t1, &t3).await; +} + +async fn foreign_chain_cursor_is_not_this_page(addr: SocketAddr, sh1: &str, t1: &str, t3: &str) { + let t2 = block_hash_hex(&pay_txid(0x22)); + for path in [ + format!("/scripthash/{sh1}/txs?after_txid={t2}"), + format!("/scripthash/{sh1}/txs/summary?after_txid={t2}"), + format!("/scripthash/{sh1}/txs/chain/{t2}"), + ] { + let (st, body) = http_get(addr, &path).await; + assert_eq!(st, 422, "{path}: {body}"); + assert!(body.contains("after_txid not found"), "{body}"); + assert!(!body.contains(t1) && !body.contains(t3), "{path}: {body}"); + } } async fn wallet_posts_scripthashes(addr: SocketAddr, a: [&str; 2], sh: [&str; 2]) { @@ -150,6 +171,7 @@ async fn wallet_posts_scripthashes(addr: SocketAddr, a: [&str; 2], sh: [&str; 2] .await; assert_eq!(st, 422, "{resp}"); assert!(resp.contains("after_txid not found"), "{resp}"); + cursor_outside_one_script_is_422(addr, a[0], sh[0], &t1, &t2, &t3).await; let too: Vec = (0..301).map(|_| "aa".repeat(32)).collect(); let (st, resp) = http_post( addr, @@ -161,6 +183,24 @@ async fn wallet_posts_scripthashes(addr: SocketAddr, a: [&str; 2], sh: [&str; 2] assert!(resp.contains("body too long"), "{resp}"); } +async fn cursor_outside_one_script_is_422( + addr: SocketAddr, + address: &str, + sh: &str, + t1: &str, + t2: &str, + t3: &str, +) { + let only = serde_json::to_vec(&json!([sh])).unwrap(); + let (st, resp) = http_post(addr, &format!("/scripthashes/txs?after_txid={t2}"), &only).await; + assert_eq!(st, 422, "{resp}"); + assert!(resp.contains("after_txid not found"), "{resp}"); + assert!(!resp.contains(t1) && !resp.contains(t3), "{resp}"); + let (st, resp) = http_get(addr, &format!("/address/{address}/txs?after_txid={t2}")).await; + assert_eq!(st, 422, "{resp}"); + assert!(resp.contains("after_txid not found"), "{resp}"); +} + /// Packed SH `/txs` runs on `spawn_blocking`, so tip height still answers on /// the test's single worker. async fn tip_height_overlaps_scripthash_txs(pad: &HttpPad) { diff --git a/crates/rbitcoin-esplora/src/handlers.rs b/crates/rbitcoin-esplora/src/handlers.rs index 487439da7..23b6a18d0 100644 --- a/crates/rbitcoin-esplora/src/handlers.rs +++ b/crates/rbitcoin-esplora/src/handlers.rs @@ -1067,19 +1067,22 @@ fn summary_page_sh( } } let filter = HistoryFilter::esplora_chain_page(after); - let (items, view) = match sh_at_view( + let (page, view) = match sh_at_view( st, sh, asof, |q, view| q.scripthash_history_summary_filtered_in(sh, &filter, view), |q, slot, view| q.scripthash_history_summary_filtered_slot_in(sh, &filter, slot, view), - Vec::new(), + cursor_page(after), client, None, ) { Ok(x) => x, Err(r) => return r, }; + let Some(items) = rows_or_cursor_missing(page) else { + return after_txid_not_found(); + }; maybe_attach_view( match summaries_json(&st.query, &items) { Ok(v) => Json(v).into_response(), @@ -1137,19 +1140,22 @@ fn chain_page_sh( client: Option<&str>, ) -> Response { let filter = HistoryFilter::esplora_chain_page(after); - let (items, view) = match sh_at_view( + let (page, view) = match sh_at_view( st, sh, asof, |q, view| q.scripthash_history_filtered_in(sh, &filter, view), |q, slot, view| q.scripthash_history_filtered_slot_in(sh, &filter, slot, view), - Vec::new(), + cursor_page(after), client, None, ) { Ok(x) => x, Err(r) => return r, }; + let Some(items) = rows_or_cursor_missing(page) else { + return after_txid_not_found(); + }; maybe_attach_view( match history_items_to_tx_json(&st.query, &items, st.network) { Ok(v) => Json(v).into_response(), @@ -1171,33 +1177,42 @@ fn combined_txs( return after_txid_not_found(); } } + // Membership is this script's mempool rows. A tx sitting in the pool for + // some other script must not restart this script's first page. + let mempool_rows = if asof.is_none() { + mempool_txs_json(st, sh) + } else { + Vec::new() + }; let after_in_mempool = after.is_some_and(|id| { - let tid = Txid::from_byte_array(id); - st.mempool.as_ref().is_some_and(|m| m.contains(&tid)) + let hex = block_hash_hex(&id); + mempool_rows.iter().any(|v| v["txid"] == hex) }); let mut out = Vec::new(); if asof.is_none() && (after.is_none() || after_in_mempool) { - let rows = mempool_txs_json(st, sh); out.extend(match after { - Some(id) if after_in_mempool => skip_mempool_after(rows, &id), - _ => rows, + Some(id) if after_in_mempool => skip_mempool_after(mempool_rows, &id), + _ => mempool_rows, }); } let chain_after = if after_in_mempool { None } else { after }; let filter = HistoryFilter::esplora_chain_page(chain_after); - let (items, view) = match sh_at_view( + let (page, view) = match sh_at_view( st, sh, asof, |q, view| q.scripthash_history_filtered_in(sh, &filter, view), |q, slot, view| q.scripthash_history_filtered_slot_in(sh, &filter, slot, view), - Vec::new(), + cursor_page(chain_after), client, None, ) { Ok(x) => x, Err(r) => return r, }; + let Some(items) = rows_or_cursor_missing(page) else { + return after_txid_not_found(); + }; maybe_attach_view( match history_items_to_tx_json(&st.query, &items, st.network) { Ok(chain) => { @@ -1274,14 +1289,29 @@ fn sort_txs_newest_first(rows: &mut [Value]) { }); } -fn skip_rows_after(rows: Vec, after: Option<[u8; 32]>) -> Vec { +/// `None`: `after` was set and is not in `rows`. The first page is not a +/// stand-in for a cursor this response does not contain. +fn skip_rows_after(rows: Vec, after: Option<[u8; 32]>) -> Option> { let Some(id) = after else { - return rows; + return Some(rows); }; let hex = block_hash_hex(&id); - match rows.iter().position(|v| v["txid"] == hex) { - Some(i) => rows[i.saturating_add(1)..].to_vec(), - None => rows, + let i = rows.iter().position(|v| v["txid"] == hex)?; + Some(rows[i.saturating_add(1)..].to_vec()) +} + +fn cursor_page(after: Option<[u8; 32]>) -> rbitcoin_query::FilteredHistory { + rbitcoin_query::FilteredHistory { + rows: Vec::new(), + cursor_missing: after.is_some(), + } +} + +fn rows_or_cursor_missing(page: rbitcoin_query::FilteredHistory) -> Option> { + if page.cursor_missing { + None + } else { + Some(page.rows) } } @@ -1311,16 +1341,17 @@ fn combined_tx_vec( out.extend(mempool_txs_json(st, sh)); } let filter = HistoryFilter::esplora_chain_page(None); - let (items, _) = sh_at_view( + let (page, _) = sh_at_view( st, sh, asof, |q, view| q.scripthash_history_filtered_in(sh, &filter, view), |q, slot, view| q.scripthash_history_filtered_slot_in(sh, &filter, slot, view), - Vec::new(), + cursor_page(None), None, bag, )?; + let items = page.rows; let chain = history_items_to_tx_json(&st.query, &items, st.network).map_err(store_err)?; out.extend(chain); Ok(out) @@ -1334,16 +1365,17 @@ fn summary_vec( bag: Option<&mut JoinBag>, ) -> Result, Response> { let filter = HistoryFilter::esplora_chain_page(None); - let (items, _) = sh_at_view( + let (page, _) = sh_at_view( st, sh, asof, |q, view| q.scripthash_history_summary_filtered_in(sh, &filter, view), |q, slot, view| q.scripthash_history_summary_filtered_slot_in(sh, &filter, slot, view), - Vec::new(), + cursor_page(None), None, bag, )?; + let items = page.rows; match summaries_json(&st.query, &items) { Ok(Value::Array(v)) => Ok(v), Ok(_) => Ok(Vec::new()), @@ -1375,7 +1407,9 @@ fn multi_txs( st.promote_bulk(client, bag); let mut rows = dedup_txid(rows); sort_txs_newest_first(&mut rows); - let rows = skip_rows_after(rows, after); + let Some(rows) = skip_rows_after(rows, after) else { + return after_txid_not_found(); + }; Json(rows).into_response() } @@ -1412,7 +1446,9 @@ fn multi_summary( .cmp(a["txid"].as_str().unwrap_or("")) }) }); - let rows = skip_rows_after(rows, after); + let Some(rows) = skip_rows_after(rows, after) else { + return after_txid_not_found(); + }; Json(rows).into_response() } diff --git a/crates/rbitcoin-esplora/src/server.rs b/crates/rbitcoin-esplora/src/server.rs index 242bc9634..fbf310ef9 100644 --- a/crates/rbitcoin-esplora/src/server.rs +++ b/crates/rbitcoin-esplora/src/server.rs @@ -402,6 +402,8 @@ struct InflightJoin { /// `None` = still running. `Some(slot)` = finished (`slot` may be empty). done: Mutex>>>, cv: std::sync::Condvar, + /// Start order of the leader. Waiters publish this stamp. + seq: u64, } impl InflightJoin { @@ -437,8 +439,30 @@ impl Drop for InflightGuard { } } +/// Publish last-1 only when this join started at or after the one stored. +/// +/// The leader stamps the counter before `f` runs. A waiter publishes that +/// same stamp, so waking after a newer script is not a newer start. +/// Matching on the scripthash alone would refuse a later sequential script. +fn publish_last_sh(c: &mut ClientJoins, sh: &[u8; 32], slot: Arc, seq: u64) { + let stored = c + .last_sh + .as_ref() + .map(|(_, _, stored_seq)| *stored_seq) + .unwrap_or(0); + if seq >= stored { + c.last_sh = Some((*sh, slot, seq)); + } +} + +fn take_join_seq(c: &mut ClientJoins) -> u64 { + c.join_seq = c.join_seq.saturating_add(1); + c.join_seq +} + struct ClientJoins { - last_sh: Option<([u8; 32], Arc)>, + last_sh: Option<([u8; 32], Arc, u64)>, + join_seq: u64, last_bulk: HashMap<[u8; 32], Arc>, last_req: Instant, inflight: HashMap<[u8; 32], Arc>, @@ -448,6 +472,7 @@ impl Default for ClientJoins { fn default() -> Self { Self { last_sh: None, + join_seq: 0, last_bulk: HashMap::new(), last_req: Instant::now(), inflight: HashMap::new(), @@ -465,7 +490,7 @@ impl JoinCache { fn last_sh_key(&self, id: &str) -> Option<[u8; 32]> { self.clients .get(id) - .and_then(|c| c.last_sh.as_ref().map(|(k, _)| *k)) + .and_then(|c| c.last_sh.as_ref().map(|(k, _, _)| *k)) } #[cfg(test)] @@ -494,7 +519,7 @@ fn cap_bulk(c: &mut ClientJoins) { let last_sh_bytes = c .last_sh .as_ref() - .map(|(_, s)| s.packed_bytes()) + .map(|(_, s, _)| s.packed_bytes()) .unwrap_or(0); let mut bytes: usize = last_sh_bytes.saturating_add(c.last_bulk.values().map(|s| s.packed_bytes()).sum()); @@ -515,7 +540,7 @@ fn cap_bulk(c: &mut ClientJoins) { fn retain_join_budget(c: &mut ClientJoins) { if c.last_sh .as_ref() - .is_some_and(|(_, s)| s.packed_bytes() > JOIN_BULK_CAP) + .is_some_and(|(_, s, _)| s.packed_bytes() > JOIN_BULK_CAP) { c.last_sh = None; } @@ -574,14 +599,15 @@ impl AppState { sweep_clients(&mut g.clients, Instant::now()); let c = g.clients.entry(id.to_string()).or_default(); c.last_req = Instant::now(); - if c.last_sh.as_ref().is_some_and(|(k, _)| k == sh) { - let mut slot = c.last_sh.as_ref().map(|(_, s)| s.clone()); + if c.last_sh.as_ref().is_some_and(|(k, _, _)| k == sh) { + let seq = take_join_seq(c); + let mut slot = c.last_sh.as_ref().map(|(_, s, _)| s.clone()); drop(g); let r = f(&mut slot); if let Some(s) = slot { let mut g = self.sh_join.lock().unwrap_or_else(|p| p.into_inner()); if let Some(c) = g.clients.get_mut(id) { - c.last_sh = Some((*sh, s)); + publish_last_sh(c, sh, s, seq); c.last_req = Instant::now(); retain_join_budget(c); } @@ -589,18 +615,21 @@ impl AppState { return r; } if let Some(inf) = c.inflight.get(sh).cloned() { - Err(inf) + let seq = inf.seq; + Err((inf, seq)) } else { + let seq = take_join_seq(c); let inf = Arc::new(InflightJoin { done: Mutex::new(None), cv: std::sync::Condvar::new(), + seq, }); c.inflight.insert(*sh, Arc::clone(&inf)); - Ok(inf) + Ok((inf, seq)) } }; - let inf = match inflight { - Err(inf) => { + let (inf, seq) = match inflight { + Err((inf, seq)) => { let mut d = inf.done.lock().unwrap_or_else(|p| p.into_inner()); while d.is_none() { d = inf.cv.wait(d).unwrap_or_else(|p| p.into_inner()); @@ -611,14 +640,14 @@ impl AppState { if let Some(s) = slot { let mut g = self.sh_join.lock().unwrap_or_else(|p| p.into_inner()); if let Some(c) = g.clients.get_mut(id) { - c.last_sh = Some((*sh, s)); + publish_last_sh(c, sh, s, seq); c.last_req = Instant::now(); retain_join_budget(c); } } return r; } - Ok(inf) => inf, + Ok((inf, seq)) => (inf, seq), }; let mut slot = { let g = self.sh_join.lock().unwrap_or_else(|p| p.into_inner()); @@ -638,7 +667,7 @@ impl AppState { c.last_req = Instant::now(); c.inflight.remove(sh); if let Some(s) = slot { - c.last_sh = Some((*sh, s)); + publish_last_sh(c, sh, s, seq); } retain_join_budget(c); } @@ -1465,6 +1494,70 @@ mod tests { ); } + #[test] + fn late_waiter_keeps_a_newer_last_sh() { + let (_dir, q) = temp_query("join-late-waiter-last-sh"); + let cache = Arc::new(Mutex::new(JoinCache::default())); + let st = Arc::new(join_only_state(Arc::new(q), Arc::clone(&cache))); + let sh_old = [0x11u8; 32]; + let sh_new = [0x22u8; 32]; + let slot = rbitcoin_query::testutil::sh_join_slot_small(); + let (leader_in, leader_in_rx) = std::sync::mpsc::channel::<()>(); + let (release, release_rx) = std::sync::mpsc::channel::<()>(); + let st_l = Arc::clone(&st); + let slot_l = Arc::clone(&slot); + let leader = std::thread::spawn(move || { + st_l.with_sh_join(Some("c1"), &sh_old, |s| { + *s = Some(slot_l); + let _ = leader_in.send(()); + let _ = release_rx.recv(); + }); + }); + leader_in_rx.recv().expect("leader entered f"); + let slot_b = Arc::clone(&slot); + st.with_sh_join(Some("c1"), &sh_new, |n| { + *n = Some(slot_b); + }); + assert_eq!( + cache.lock().unwrap().last_sh_key("c1"), + Some(sh_new), + "a newer script that finishes while the older leader is inside f is last-1" + ); + let (waiter_in, waiter_in_rx) = std::sync::mpsc::channel::<()>(); + let (waiter_go, waiter_go_rx) = std::sync::mpsc::channel::<()>(); + let st_w = Arc::clone(&st); + let slot_w = Arc::clone(&slot); + let waiter = std::thread::spawn(move || { + st_w.with_sh_join(Some("c1"), &sh_old, |s| { + *s = Some(slot_w); + let _ = waiter_in.send(()); + let _ = waiter_go_rx.recv(); + }); + }); + std::thread::sleep(Duration::from_millis(50)); + assert!( + waiter_in_rx.try_recv().is_err(), + "waiter must block on the leader, not run the join itself" + ); + let _ = release.send(()); + leader.join().expect("leader"); + assert_eq!( + cache.lock().unwrap().last_sh_key("c1"), + Some(sh_new), + "an older leader that finishes after a newer script must not replace last-1" + ); + waiter_in_rx + .recv_timeout(Duration::from_secs(2)) + .expect("waiter entered f after the leader finished"); + let _ = waiter_go.send(()); + waiter.join().expect("waiter"); + assert_eq!( + cache.lock().unwrap().last_sh_key("c1"), + Some(sh_new), + "a late waiter must not put the old scripthash back over a newer last-1" + ); + } + #[cfg(unix)] #[cfg(unix)] #[tokio::test] diff --git a/crates/rbitcoin-mempool/src/orphanage.rs b/crates/rbitcoin-mempool/src/orphanage.rs index 9641d7515..c685b5073 100644 --- a/crates/rbitcoin-mempool/src/orphanage.rs +++ b/crates/rbitcoin-mempool/src/orphanage.rs @@ -250,6 +250,8 @@ impl Orphanage { .any(|p| self.peer_orphan_weight(*p) <= ORPHAN_RESERVED_WEIGHT_PER_PEER) } + /// Drop `peer` from the oldest orphan it announced. The tx stays when + /// another peer still announces it; only that peer's weight is released. fn evict_one_from_peer(&mut self, peer: u64) -> bool { let victim = self .fifo @@ -263,7 +265,23 @@ impl Orphanage { let Some(txid) = victim else { return false; }; - self.remove_txid(&txid); + let weight = { + let Some(e) = self.by_txid.get_mut(&txid) else { + return false; + }; + if !e.announcers.remove(&peer) { + return false; + } + e.weight + }; + self.sub_peer_weight(peer, weight); + let last = self + .by_txid + .get(&txid) + .is_some_and(|e| e.announcers.is_empty()); + if last { + self.remove_txid(&txid); + } self.fifo.retain(|t| self.by_txid.contains_key(t)); true } @@ -506,6 +524,37 @@ mod tests { assert!(peer2_kept, "peer 2 must keep its orphan"); } + /// Peer B re-announces peer A's orphan, then overflows its own reserve. + /// The eviction drops B only. The parent still delivers A's orphan. + #[test] + fn shared_announcer_eviction_keeps_the_other_peer() { + let mut o = Orphanage::new(); + let parent = txid_n(8); + let mut miss = BTreeSet::new(); + miss.insert(parent); + let tx_a = make_orphan(parent, 1); + let tid = tx_a.compute_txid(); + assert!(o.insert_from(tx_a, miss.clone(), Some(1))); + assert!(o.add_announcer(&tid, 2)); + for i in 0..80u8 { + let mut tx = make_orphan(parent, i.wrapping_add(3)); + tx.input[0].witness = Witness::from_slice(&[vec![i.wrapping_add(3); 20_000]]); + tx.lock_time = LockTime::from_height(i as u32).unwrap(); + o.insert_from(tx, miss.clone(), Some(2)); + } + assert!(o.peer_orphan_weight(2) <= ORPHAN_RESERVED_WEIGHT_PER_PEER); + assert!( + o.contains(&tid), + "peer 1's orphan must survive peer 2's reserve" + ); + assert!(o.announcers_of(&tid).contains(&1)); + let kids = o.take_children_of(&parent); + assert!( + kids.iter().any(|tx| tx.compute_txid() == tid), + "the parent still delivers peer 1's orphan" + ); + } + #[test] fn insert_and_take_by_parent() { let mut o = Orphanage::new(); diff --git a/crates/rbitcoin-net/src/compact.rs b/crates/rbitcoin-net/src/compact.rs index 9ed32a358..f06383ab9 100644 --- a/crates/rbitcoin-net/src/compact.rs +++ b/crates/rbitcoin-net/src/compact.rs @@ -263,9 +263,22 @@ pub(crate) fn apply_block_transactions( /// BIP152: filled slots must merkle to the compact header and not be a /// Core `IsBlockMutated` body (empty → getdata). +/// Equal adjacent txids keep `Block::check_merkle_root` (the odd-node duplicate +/// matches CVE-2012-2459) and still must not be treated as this header's block. +pub(crate) fn merkle_body_mutated(txdata: &[Transaction]) -> bool { + let leaves: Vec<[u8; 32]> = txdata + .iter() + .map(|tx| tx.compute_txid().to_byte_array()) + .collect(); + rbitcoin_store::merkle_root_mutated(&leaves).1 +} + fn finish_reconstructed(header: Header, txdata: Vec) -> Result> { let block = Block { header, txdata }; - if !block.check_merkle_root() || rbitcoin_consensus::block_mutated_without_coinbase(&block) { + if !block.check_merkle_root() + || merkle_body_mutated(&block.txdata) + || rbitcoin_consensus::block_mutated_without_coinbase(&block) + { return Err(Vec::new()); } Ok(block) @@ -668,6 +681,16 @@ mod tests { ); } + #[test] + fn repeated_pair_with_matching_root_is_not_a_block() { + let cb = coinbase(); + let block = sealed_block(vec![cb.clone(), cb]); + assert!(block.check_merkle_root()); + assert!(merkle_body_mutated(&block.txdata)); + let err = finish_reconstructed(block.header, block.txdata).expect_err("mutated pair"); + assert!(err.is_empty(), "getdata the hash, got {err:?}"); + } + #[test] fn prefilled_64_byte_body_without_coinbase_is_not_a_block() { let inner = crate::chain::sixty_four_byte_body(BlockHash::all_zeros(), 1); diff --git a/crates/rbitcoin-net/src/i2p_sam.rs b/crates/rbitcoin-net/src/i2p_sam.rs index f472a33c0..301cd41ba 100644 --- a/crates/rbitcoin-net/src/i2p_sam.rs +++ b/crates/rbitcoin-net/src/i2p_sam.rs @@ -72,6 +72,13 @@ fn installed() -> Result { .ok_or_else(|| NetError::Encode("i2p dial requires SAM (--i2p-sam)".into())) } +/// Error from the SAM TCP connect that opens a stream. A refused local +/// port uses this text. It must not contain `STREAM CONNECT:`, which +/// `session_dead` treats as a dead installed session. +fn sam_dial_error(err: &std::io::Error) -> NetError { + NetError::Encode(format!("i2p sam dial: {err}")) +} + fn session_dead(err: &NetError) -> bool { if stream_socket_dead(err) { return true; @@ -316,7 +323,7 @@ impl I2pDialer { async fn stream_connect_once(&self, dest_b32: &str) -> Result { let mut s = TcpStream::connect(self.sam_addr) .await - .map_err(|e| NetError::Encode(format!("i2p sam stream connect: {e}")))?; + .map_err(|e| sam_dial_error(&e))?; hello(&mut s).await?; write_line( &mut s, @@ -513,6 +520,36 @@ async fn read_line(s: &mut TcpStream) -> Result { #[cfg(test)] mod tests { use super::*; + + #[tokio::test(flavor = "current_thread")] + async fn refused_local_sam_dial_is_not_a_dead_session() { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let sam_addr = listener.local_addr().expect("addr"); + drop(listener); + let dialer = I2pDialer { + sam_addr, + session_id: "refused".into(), + destination: "dest".into(), + forward_port: None, + }; + let refused = dialer + .stream_connect_once("abcdef.b32.i2p") + .await + .expect_err("closed SAM port"); + assert!( + !session_dead(&refused), + "a refused local SAM port must not rotate the installed session: {refused}" + ); + assert!( + matches!(refused, NetError::Encode(ref s) if s.starts_with("i2p sam dial:")), + "the dial path must use sam_dial_error: {refused}" + ); + let reply = NetError::Encode("i2p sam stream: STREAM STATUS RESULT=CANT_REACH_PEER".into()); + assert!(!session_dead(&reply)); + let invalid = NetError::Encode("i2p sam stream: SESSION STATUS RESULT=INVALID_ID".into()); + assert!(session_dead(&invalid)); + } + use std::collections::HashSet; use std::sync::{Arc, Mutex}; use std::time::{SystemTime, UNIX_EPOCH}; diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index d98410404..5b22e3408 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -577,10 +577,12 @@ pub(crate) fn disconnect_peer( let Some(idx) = slots.iter().position(|s| s.id == peer && s.alive) else { return; }; - let addr = slots[idx].addr; + let net = slots[idx].net; slots[idx].alive = false; let _ = slots[idx].cmd_tx.send(PeerCmd::Shutdown); - record_stall_kick(addr_cooldown, addr_strikes, addr, Instant::now()); + if let Some(sock) = net.socket_addr().filter(|s| !s.ip().is_unspecified()) { + record_stall_kick(addr_cooldown, addr_strikes, sock, Instant::now()); + } } /// Bump the process-local strike count and set `addr_cooldown`. @@ -604,26 +606,30 @@ pub(crate) fn record_stall_kick( pub(crate) fn note_dead_without_block_bytes( book: &mut AddrMan, addr_cooldown: &mut HashMap, - addr: SocketAddr, + addr: crate::NetAddr, first_data_ms: u64, now: Instant, ) { if first_data_ms != 0 { return; } - book.note_connect_failed(addr, false); - addr_cooldown.insert(addr, now + STALL_ADDR_COOLDOWN); + book.note_connect_failed_addr(addr, false); + if let Some(sock) = addr.socket_addr().filter(|s| !s.ip().is_unspecified()) { + addr_cooldown.insert(sock, now + STALL_ADDR_COOLDOWN); + } } /// Misbehavior-threshold death. Cools the dial even after a block body was counted. pub(crate) fn note_misbehavior_dead( book: &mut AddrMan, addr_cooldown: &mut HashMap, - addr: SocketAddr, + addr: crate::NetAddr, now: Instant, ) { - book.note_connect_failed(addr, false); - addr_cooldown.insert(addr, now + STALL_ADDR_COOLDOWN); + book.note_connect_failed_addr(addr, false); + if let Some(sock) = addr.socket_addr().filter(|s| !s.ip().is_unspecified()) { + addr_cooldown.insert(sock, now + STALL_ADDR_COOLDOWN); + } } /// One stall rule: if a peer has outstanding block getdata and no **block** @@ -985,7 +991,7 @@ mod tests { book.note_connected(lemon); let mut cooldown = HashMap::new(); let now = Instant::now(); - note_dead_without_block_bytes(&mut book, &mut cooldown, lemon, 0, now); + note_dead_without_block_bytes(&mut book, &mut cooldown, crate::NetAddr::Ip(lemon), 0, now); assert!( book.flags(&lemon).failed_last_connect(), "no block bytes → last-resort" @@ -997,7 +1003,7 @@ mod tests { let good = addr(5); book.note_connected(good); - note_dead_without_block_bytes(&mut book, &mut cooldown, good, 42, now); + note_dead_without_block_bytes(&mut book, &mut cooldown, crate::NetAddr::Ip(good), 42, now); assert!( !book.flags(&good).failed_last_connect(), "peer that sent block bytes keeps its connected rank" @@ -1012,7 +1018,7 @@ mod tests { book.note_connected(lemon); let mut cooldown = HashMap::new(); let now = Instant::now(); - note_misbehavior_dead(&mut book, &mut cooldown, lemon, now); + note_misbehavior_dead(&mut book, &mut cooldown, crate::NetAddr::Ip(lemon), now); assert!( book.flags(&lemon).failed_last_connect(), "a misbehavior death is a failed connect" diff --git a/crates/rbitcoin-net/src/ibd/events/mod.rs b/crates/rbitcoin-net/src/ibd/events/mod.rs index dfe4fd90b..1204a5a3a 100644 --- a/crates/rbitcoin-net/src/ibd/events/mod.rs +++ b/crates/rbitcoin-net/src/ibd/events/mod.rs @@ -690,23 +690,22 @@ fn apply_peer_dead(st: &mut IbdWorkState, peer_book: &mut AddrMan, peer: usize, super::header_walk::forget_walk_peer(st, peer); if let Some(s) = st.slots.iter().find(|s| s.id == peer) { if reason == "peer misbehavior threshold" { - note_misbehavior_dead(peer_book, &mut st.addr_cooldown, s.addr, Instant::now()); + note_misbehavior_dead(peer_book, &mut st.addr_cooldown, s.net, Instant::now()); } else { note_dead_without_block_bytes( peer_book, &mut st.addr_cooldown, - s.addr, + s.net, s.first_data_ms, Instant::now(), ); } let lat = s.first_data_ms.saturating_sub(s.connected_ms); - peer_book.apply_ibd_dead_speed( - s.addr, - lat, - s.rate.bps(), - st.addr_cooldown.contains_key(&s.addr), - ); + let cooled = s + .net + .socket_addr() + .is_some_and(|sock| st.addr_cooldown.contains_key(&sock)); + peer_book.apply_ibd_dead_speed_addr(s.net, lat, s.rate.bps(), cooled); } let freed = release_peer_block_work(&mut st.slots, &mut st.inflight, &mut st.body, peer); st.reopen_for_densify(&freed); diff --git a/crates/rbitcoin-net/src/overlay_addrman_journey.rs b/crates/rbitcoin-net/src/overlay_addrman_journey.rs index a2cb36c84..5dcf17e3d 100644 --- a/crates/rbitcoin-net/src/overlay_addrman_journey.rs +++ b/crates/rbitcoin-net/src/overlay_addrman_journey.rs @@ -155,6 +155,39 @@ fn only_net_dials_and_peers_file(am: &Mutex, overlays: [c only.set_only_net(vec![OnlyNet::Cjdns]); assert_eq!(only.take_dial_candidates(8, &HashSet::new(), &[]), vec![cjdns_sock]); + let mut mixed = crate::seeds::AddrMan::new(); + for i in 0..8u8 { + mixed.add(SocketAddr::from((Ipv4Addr::new(10, 0, 0, i), 8333))); + } + mixed.add_addr(onion); + let got = mixed.take_dial_candidates_net(8, &HashSet::new(), &[]); + assert_eq!(got.len(), 8, "{got:?}"); + assert!( + got.contains(&onion), + "a full clearnet batch still dials one onion: {got:?}" + ); + let mut failed = mixed.clone(); + failed.note_connect_failed_addr(onion, false); + let fresh = crate::NetAddr::Onion { + pk: [0x44; 32], + port: 8333, + }; + failed.add_addr(fresh); + let got = failed.take_dial_candidates_net(8, &HashSet::new(), &[]); + assert!( + got.contains(&fresh) && !got.contains(&onion), + "a failed onion loses to a fresh one: {got:?}" + ); + let placeholder = SocketAddr::from((Ipv4Addr::UNSPECIFIED, 8333)); + failed.note_connect_failed(placeholder, false); + assert!( + failed + .entries() + .iter() + .all(|e| e.addr != crate::NetAddr::Ip(placeholder)), + "a version-message placeholder is not a peer" + ); + let dir = rbitcoin_query::testutil::TempDir::labeled("overlay-peers").unwrap(); let path = dir.join("peers"); book.save(&path).unwrap(); diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index e15eb30d3..6e8379836 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -3523,6 +3523,12 @@ async fn on_block( punish_disconnect(&mut follow.ban_score, session); return Ok(()); } + if crate::compact::merkle_body_mutated(&block.txdata) { + rbitcoin_log::info!("Block mutated: bad-txns-duplicate"); + take_requested_block(hub, &mut follow.requested_blocks, &hash); + punish_disconnect(&mut follow.ban_score, session); + return Ok(()); + } if rbitcoin_consensus::block_mutated_without_coinbase(block) { rbitcoin_log::info!("Block mutated: 64-byte transaction without a coinbase"); take_requested_block(hub, &mut follow.requested_blocks, &hash); diff --git a/crates/rbitcoin-net/src/peer_header_dos_journey.rs b/crates/rbitcoin-net/src/peer_header_dos_journey.rs index 5044198e4..1923abd10 100644 --- a/crates/rbitcoin-net/src/peer_header_dos_journey.rs +++ b/crates/rbitcoin-net/src/peer_header_dos_journey.rs @@ -941,6 +941,79 @@ async fn merkle_mismatch_forgets_the_ask( assert_eq!(hub.tip_hash(), Some(tip)); } +/// CVE-2012-2459: two identical transactions keep `check_merkle_root` and are +/// still not a block. Tip follow disconnects and does not cache the hash. +async fn duplicate_pair_disconnects( + hub: &crate::chain::ChainHub, + peers: &std::sync::Arc, +) { + use bitcoin::absolute::LockTime; + use bitcoin::transaction::{OutPoint, Sequence, TxIn, TxOut, Version}; + use bitcoin::{Amount, ScriptBuf, Transaction, Witness}; + + let cb = Transaction { + version: Version::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint::null(), + script_sig: ScriptBuf::from_bytes(vec![0x01, 0x01]), + sequence: Sequence::MAX, + witness: Witness::new(), + }], + output: vec![TxOut { + value: Amount::from_sat(50_0000_0000), + script_pubkey: ScriptBuf::from_bytes(vec![0x51]), + }], + }; + let mut block = bitcoin::Block { + header: bitcoin::block::Header { + version: bitcoin::block::Version::TWO, + prev_blockhash: hub.tip_hash().unwrap(), + merkle_root: bitcoin::TxMerkleNode::all_zeros(), + time: hub.tip_header().unwrap().time.saturating_add(3), + bits: bitcoin::CompactTarget::from_consensus(0x207fffff), + nonce: 0, + }, + txdata: vec![cb.clone(), cb], + }; + block.header.merkle_root = block.compute_merkle_root().expect("pair"); + assert!( + block.check_merkle_root(), + "the duplicated pair still matches the header root" + ); + let hash = block.block_hash(); + let tip = hub.tip_hash().unwrap(); + let (out_tx, _rx) = mpsc::unbounded_channel(); + let sender = live_peer(peers, 18476, 25, true); + let mut follow = PeerFollowState::new(); + follow.requested_blocks.insert(hash); + hub.note_asked_block(hash); + rbitcoin_log::capture_logs(true); + on_block( + hub, + &out_tx, + &mut follow, + Some(sender.as_ref()), + &block, + ) + .await + .unwrap(); + assert!( + logs_have("Block mutated: bad-txns-duplicate"), + "the drop names the duplicate pair" + ); + assert!( + sender.stop.load(Ordering::SeqCst), + "tip follow disconnects a duplicate-pair body" + ); + assert!( + !hub.is_block_invalid(&hash), + "the real body for this hash must stay acceptable" + ); + assert!(!follow.requested_blocks.contains(&hash)); + assert_eq!(hub.tip_hash(), Some(tip)); +} + /// A block whose header fails contextual checks still logs Core's reject reason. async fn rejected_header_logs_core_reason(hub: &crate::chain::ChainHub) { use bitcoin::block::Version; @@ -1195,6 +1268,7 @@ async fn peer_header_dos_and_self_announce() { noban_bad_block_is_not_punished(&hub, &peers).await; sixty_four_byte_body_is_mutated(&hub, &peers).await; merkle_mismatch_forgets_the_ask(&hub, &peers).await; + duplicate_pair_disconnects(&hub, &peers).await; rejected_header_logs_core_reason(&hub).await; let tip = hub.tip_height().unwrap(); diff --git a/crates/rbitcoin-net/src/seeds.rs b/crates/rbitcoin-net/src/seeds.rs index 5e4f79e22..12caf417b 100644 --- a/crates/rbitcoin-net/src/seeds.rs +++ b/crates/rbitcoin-net/src/seeds.rs @@ -20,6 +20,10 @@ use crate::asmap::AsMap; use crate::netaddr::{addr_allowed, NetAddr, OnlyNet}; use crate::netgroup::{netgroup, select_diverse}; +fn net_addr_unspecified(addr: NetAddr) -> bool { + matches!(addr, NetAddr::Ip(s) if s.ip().is_unspecified()) +} + /// Skip a recently dialed addr while any other candidate remains. pub(crate) const DIAL_ATTEMPT_RECENT: Duration = Duration::from_secs(10 * 60); @@ -360,24 +364,51 @@ impl AddrMan { .copied() .filter_map(NetAddr::socket_addr) .collect(); - let mut out: Vec = self + let ip: Vec = self .take_dial_candidates(max, &ip_ex, occupied) .into_iter() .map(NetAddr::from_socket) .collect(); - if out.len() >= max { - return out; + let now = Instant::now(); + let mut overlay: Vec<(u8, NetAddr)> = self + .order + .iter() + .copied() + .filter(|a| matches!(a, NetAddr::Onion { .. } | NetAddr::I2p { .. })) + .filter(|a| self.dialable(*a) && !exclude.contains(a)) + .map(|a| (self.flags_of(&a).dial_tier(), a)) + .collect(); + overlay.sort_by_key(|(tier, _)| *tier); + let ip_has_fresh = ip.iter().any(|a| self.flags_of(a).dial_tier() < 2); + if ip_has_fresh { + overlay.retain(|(tier, a)| { + *tier < 2 + && !self.flags_of(a).is_incompatible() + && !self.recently_attempted_addr(*a, now) + }); + } else if overlay + .iter() + .any(|(_, a)| !self.flags_of(a).is_incompatible()) + { + overlay.retain(|(_, a)| !self.flags_of(a).is_incompatible()); + if overlay + .iter() + .any(|(_, a)| !self.recently_attempted_addr(*a, now)) + { + overlay.retain(|(_, a)| !self.recently_attempted_addr(*a, now)); + } + } + let mut out = ip; + // A full clearnet batch used to return before any onion or I2P row. + // Keep one slot for the best eligible overlay when the batch has room + // for more than a single peer. + if !overlay.is_empty() && out.len() >= max && max >= 2 { + out.pop(); } - for &a in &self.order { + for (_, a) in overlay { if out.len() >= max { break; } - if !matches!(a, NetAddr::Onion { .. } | NetAddr::I2p { .. }) - || !self.dialable(a) - || exclude.contains(&a) - { - continue; - } if !out.contains(&a) { out.push(a); } @@ -385,6 +416,12 @@ impl AddrMan { out } + fn recently_attempted_addr(&self, addr: NetAddr, now: Instant) -> bool { + self.last_attempt + .get(&addr) + .is_some_and(|&t| now.saturating_duration_since(t) < DIAL_ATTEMPT_RECENT) + } + pub fn add(&mut self, addr: SocketAddr) { self.add_addr(NetAddr::from_socket(addr)); } @@ -618,6 +655,9 @@ impl AddrMan { } pub fn note_connect_failed_addr(&mut self, addr: NetAddr, incompatible: bool) { + if net_addr_unspecified(addr) { + return; + } self.add_addr(addr); if let Some(f) = self.by_addr.get_mut(&addr) { if incompatible { @@ -631,8 +671,15 @@ impl AddrMan { /// Throughput / latency sample from an active session. pub fn note_speed(&mut self, addr: SocketAddr, latency_ms: u64, bytes_per_sec: u64) { - self.add(addr); - if let Some(f) = self.by_addr.get_mut(&NetAddr::from_socket(addr)) { + self.note_speed_addr(NetAddr::from_socket(addr), latency_ms, bytes_per_sec); + } + + pub fn note_speed_addr(&mut self, addr: NetAddr, latency_ms: u64, bytes_per_sec: u64) { + if net_addr_unspecified(addr) { + return; + } + self.add_addr(addr); + if let Some(f) = self.by_addr.get_mut(&addr) { f.insert(PeerFlags::HAS_CONNECTED); f.apply_speed_sample(latency_ms, bytes_per_sec); } @@ -641,8 +688,15 @@ impl AddrMan { /// IBD stall / relative-slow: force `SLOW`, clear `FAST` (mid-range samples /// would otherwise leave a prior FAST bit and keep `dial_tier` 0). pub fn note_ibd_slow(&mut self, addr: SocketAddr) { - self.add(addr); - if let Some(f) = self.by_addr.get_mut(&NetAddr::from_socket(addr)) { + self.note_ibd_slow_addr(NetAddr::from_socket(addr)); + } + + pub fn note_ibd_slow_addr(&mut self, addr: NetAddr) { + if net_addr_unspecified(addr) { + return; + } + self.add_addr(addr); + if let Some(f) = self.by_addr.get_mut(&addr) { f.insert(PeerFlags::HAS_CONNECTED); f.insert(PeerFlags::SLOW); f.remove(PeerFlags::FAST); @@ -657,11 +711,24 @@ impl AddrMan { bps: Option, ibd_outlier: bool, ) { + self.apply_ibd_dead_speed_addr(NetAddr::from_socket(addr), latency_ms, bps, ibd_outlier); + } + + pub fn apply_ibd_dead_speed_addr( + &mut self, + addr: NetAddr, + latency_ms: u64, + bps: Option, + ibd_outlier: bool, + ) { + if net_addr_unspecified(addr) { + return; + } if let Some(bps) = bps { - self.note_speed(addr, latency_ms, bps); + self.note_speed_addr(addr, latency_ms, bps); } if ibd_outlier { - self.note_ibd_slow(addr); + self.note_ibd_slow_addr(addr); } } diff --git a/crates/rbitcoin-net/src/tx_relay.rs b/crates/rbitcoin-net/src/tx_relay.rs index ca9fca141..e995d1108 100644 --- a/crates/rbitcoin-net/src/tx_relay.rs +++ b/crates/rbitcoin-net/src/tx_relay.rs @@ -2069,9 +2069,10 @@ impl MempoolHub { self.tx_snap_dirty.store(true, Ordering::Release); } - /// Map API target blocks → engine depth (0–2 → default horizon of 1). + /// Map API target blocks → engine depth. `0` is the 1-block horizon. + /// Targets 1 and 2 stay on the curve (`fee-target-curve`). fn fee_depth(target_blocks: u32) -> u32 { - if target_blocks == 0 || target_blocks <= 2 { + if target_blocks == 0 { Self::DEFAULT_HORIZON_BLOCKS } else { target_blocks @@ -3746,6 +3747,36 @@ mod tests { std::env::temp_dir().join(format!("rbitcoin-txrelay-{n}-{seq}")) } + /// `estimatesmartfee` / Electrum `blockchain.estimatefee` forward this rate. + /// Target 2 is the depth-2 curve point, not the 1-block horizon. + #[test] + fn confirm_target_two_is_the_two_block_rate() { + let store_dir = tmp(); + let mp_dir = tmp(); + let q = Query::open_or_create_tiny(&store_dir).unwrap(); + let hub = MempoolHub::open(&mp_dir, Arc::new(q)).unwrap(); + let mut snap = FeeSnapshot::empty(Instant::now()); + let rates: Vec> = FEE_SNAPSHOT_DEPTHS + .iter() + .enumerate() + .map(|(i, _)| Some(10_000 - i as u64 * 100)) + .collect(); + snap.depth_rates_sat_kvb = rates.clone(); + hub.fee_dirty.store(false, Ordering::Release); + hub.fee_snapshot.store(Arc::new(snap)); + let one = hub.estimate_fee_btc_per_kb(1); + let two = hub.estimate_fee_btc_per_kb(2); + let expect = + fee_at_target_sat_kvb(FEE_SNAPSHOT_DEPTHS, &rates, 2).unwrap() as f64 / 100_000_000.0; + assert!((two - expect).abs() < 1e-12, "two={two} expect={expect}"); + assert!( + (one - two).abs() > 1e-12, + "target 2 must not collapse onto target 1 ({one})" + ); + let _ = std::fs::remove_dir_all(&mp_dir); + let _ = std::fs::remove_dir_all(&store_dir); + } + #[test] fn recent_reject_at_the_cap_does_not_clear() { let dir = tmp(); diff --git a/crates/rbitcoin-query/src/lib.rs b/crates/rbitcoin-query/src/lib.rs index a7ede537f..4b6ff80fc 100644 --- a/crates/rbitcoin-query/src/lib.rs +++ b/crates/rbitcoin-query/src/lib.rs @@ -178,8 +178,8 @@ pub use connect::{spawn_sh_writebehind, ConfirmPrepared}; pub use id_map::{IdMap, OutPointHasher, OutPointSet, TxidHasher, TxidMap, TxidSet}; pub use in_flight::InFlight; pub use scripthash::{ - BlockTouch, HistoryFilter, HistoryOrder, ScanUtxo, ScriptHashBalance, ScriptHashChainStats, - ScriptHashHistoryItem, ScriptHashTxSummary, ScriptHashUtxo, ShJoinSlot, + BlockTouch, FilteredHistory, HistoryFilter, HistoryOrder, ScanUtxo, ScriptHashBalance, + ScriptHashChainStats, ScriptHashHistoryItem, ScriptHashTxSummary, ScriptHashUtxo, ShJoinSlot, }; pub use spend_sync::SpendSync; pub use stamp::{ diff --git a/crates/rbitcoin-query/src/query_sh_caps_journey.rs b/crates/rbitcoin-query/src/query_sh_caps_journey.rs index b336522fb..5b5ca410a 100644 --- a/crates/rbitcoin-query/src/query_sh_caps_journey.rs +++ b/crates/rbitcoin-query/src/query_sh_caps_journey.rs @@ -1,3 +1,24 @@ +fn assert_unspent_newest_page_stops( + q: &Query, + quiet_sh: [u8; 32], + tip_cb_txid: [u8; 32], + quiet_old: rbitcoin_primitives::Fk, +) { + use crate::scripthash::HistoryOrder::NewestFirst; + + q.store().reset_txid_get_many(); + assert_eq!( + sh_page(q, &quiet_sh, NewestFirst, None).unwrap(), + [tip_cb_txid], + "newest-first page is the tip create" + ); + let scanned = q.store().txid_get_many_fks(); + assert!( + !scanned.contains(&quiet_old.0), + "an unspent newest-first page stops before older creates: {scanned:?}" + ); +} + fn sh_page( q: &Query, sh: &[u8; 32], @@ -11,6 +32,7 @@ fn sh_page( ..crate::scripthash::HistoryFilter::open() }; Ok(q.scripthash_history_filtered(sh, &filter)? + .rows .iter() .map(|i| i.txid) .collect()) @@ -29,8 +51,11 @@ fn sh_history_caps() { let probe_sh = script_hash(&[0x52]); let add_probe_output = |ta: &mut TxApply| { ta.tx.output_count += 1; - ta.outputs - .push(OutputRecord::unspent(1, vec![0x52])); + ta.outputs.push(OutputRecord::unspent(1, vec![0x52])); + }; + let add_quiet_output = |ta: &mut TxApply| { + ta.tx.output_count += 1; + ta.outputs.push(OutputRecord::unspent(1, vec![0x53])); }; let mut prev = Fk::NULL; let mut parent = None; @@ -38,6 +63,9 @@ fn sh_history_caps() { for h in 0..5u32 { let (header, mut ta) = coinbase_block(h, prev, parent); add_probe_output(&mut ta); + if h >= 1 { + add_quiet_output(&mut ta); + } parent = Some(header.hash); cb_txids.push(ta.tx.txid); prev = q.connect_block(Height(h), &header, &[ta]).unwrap(); @@ -46,6 +74,10 @@ fn sh_history_caps() { let (h5, mut cb5) = coinbase_block(5, prev, parent); add_probe_output(&mut cb5); + add_quiet_output(&mut cb5); + let tip_cb_txid = cb5.tx.txid; + let quiet_sh = script_hash(&[0x53]); + let quiet_old = q.block_tx_fks(Height(1)).unwrap()[0]; let mut spend = cb5.clone(); spend.tx.output_count = 1; spend.outputs.truncate(1); @@ -108,13 +140,15 @@ fn sh_history_caps() { [spend_txid], "the newest row spends the oldest create" ); + assert_unspent_newest_page_stops(&q, quiet_sh, tip_cb_txid, quiet_old); let view = q.pin_sh_chain_view().unwrap().expect("sh view"); let page = crate::scripthash::HistoryFilter::esplora_chain_page(None); let mut slot = None; let rows = q .scripthash_history_filtered_slot_in(&sh, &page, &mut slot, &view) - .expect("paged history skips the unpaged cap"); + .expect("paged history skips the unpaged cap") + .rows; assert!(!rows.is_empty()); assert!( slot.is_none(), @@ -122,7 +156,8 @@ fn sh_history_caps() { ); let sums = q .scripthash_history_summary_filtered_slot_in(&sh, &page, &mut slot, &view) - .expect("paged summary skips the unpaged cap"); + .expect("paged summary skips the unpaged cap") + .rows; assert!(!sums.is_empty()); assert!(slot.is_none()); let err = q diff --git a/crates/rbitcoin-query/src/scripthash.rs b/crates/rbitcoin-query/src/scripthash.rs index 2cc4a710a..8f2ac1699 100644 --- a/crates/rbitcoin-query/src/scripthash.rs +++ b/crates/rbitcoin-query/src/scripthash.rs @@ -119,6 +119,18 @@ impl HistoryFilter { } } +/// Rows after [`apply_history_filter`]. +/// +/// `cursor_missing` is set when `after_txid` was given and that tx is not in +/// the history. `rows` is then empty, so a caller that serves them does not +/// restart the first page. A cursor that is the last row is not missing: +/// `rows` is empty and `cursor_missing` is false (the next page). +#[derive(Debug)] +pub struct FilteredHistory { + pub rows: Vec, + pub cursor_missing: bool, +} + /// Apply [`HistoryFilter`] to an already-built history list (no store I/O). /// /// Does not re-sort input beyond the filter's [`HistoryOrder`]. Window is applied @@ -126,7 +138,7 @@ impl HistoryFilter { pub fn apply_history_filter( items: &[ScriptHashHistoryItem], filter: &HistoryFilter, -) -> Vec { +) -> FilteredHistory { let from = i64::from(filter.from_height); let mut out: Vec = items .iter() @@ -153,11 +165,16 @@ pub fn apply_history_filter( } } + let mut cursor_missing = false; if let Some(after) = filter.after_txid { if let Some(pos) = out.iter().position(|i| i.txid == after) { out = out.split_off(pos.saturating_add(1)); + } else { + // Not this history. Empty, not the first page: a wallet that + // retries the same cursor must not loop. + cursor_missing = true; + out.clear(); } - // If after_txid not found, Esplora-like behavior: return from start (no skip). } if let Some(lim) = filter.limit { @@ -165,13 +182,16 @@ pub fn apply_history_filter( out.truncate(lim); } } - out + FilteredHistory { + rows: out, + cursor_missing, + } } fn history_items_from_joined( joined: &[ShJoinedOut], filter: &HistoryFilter, -) -> Vec { +) -> FilteredHistory { let mut by_txid: BTreeMap<[u8; 32], (i64, Fk)> = BTreeMap::new(); let to_excl = filter.to_height; for rec in joined { @@ -223,8 +243,11 @@ fn history_items_from_joined( fn summaries_from_joined( joined: &[ShJoinedOut], filter: &HistoryFilter, -) -> Vec { - let items = history_items_from_joined(joined, filter); +) -> FilteredHistory { + let FilteredHistory { + rows: items, + cursor_missing, + } = history_items_from_joined(joined, filter); let mut net: HashMap = HashMap::new(); for rec in joined { net.entry(rec.out.create_tx_fk) @@ -236,15 +259,18 @@ fn summaries_from_joined( .or_insert(0i64.saturating_sub(rec.out.value)); } } - items - .into_iter() - .map(|it| ScriptHashTxSummary { - txid: it.txid, - value: net.get(&it.tx_fk).copied().unwrap_or(0), - height: it.height, - tx_fk: it.tx_fk, - }) - .collect() + FilteredHistory { + rows: items + .into_iter() + .map(|it| ScriptHashTxSummary { + txid: it.txid, + value: net.get(&it.tx_fk).copied().unwrap_or(0), + height: it.height, + tx_fk: it.tx_fk, + }) + .collect(), + cursor_missing, + } } #[derive(Clone, Debug, PartialEq, Eq, Default)] @@ -493,19 +519,12 @@ impl Query { Ok(()) } - /// True when more creates cannot change the already-full page. - pub(crate) fn history_page_closed( - &self, - joined: &[ShJoinedOut], - filter: &HistoryFilter, - rest: &[Fk], - ) -> Result { + /// The limited page already has its rows, and `after_txid` is in the join + /// when the caller asked for one. + fn history_page_full(joined: &[ShJoinedOut], filter: &HistoryFilter) -> bool { let Some(limit) = filter.limit else { - return Ok(false); + return false; }; - if rest.is_empty() { - return Ok(false); - } if let Some(after) = filter.after_txid { let open = HistoryFilter { limit: None, @@ -513,14 +532,27 @@ impl Query { ..filter.clone() }; let seen = history_items_from_joined(joined, &open); - if !seen.iter().any(|i| i.txid == after) { - return Ok(false); + if !seen.rows.iter().any(|i| i.txid == after) { + return false; } } - let page = history_items_from_joined(joined, filter); - if page.len() < limit { - return Ok(false); + history_items_from_joined(joined, filter).rows.len() >= limit + } + + /// How many leading `rest` creates can still change a full page. + /// + /// One height read and one spent-range read for the whole tail. The join + /// loop calls this once, so a full page does not re-read that tail per wave. + fn history_tail_keep( + &self, + joined: &[ShJoinedOut], + filter: &HistoryFilter, + rest: &[Fk], + ) -> Result { + if rest.is_empty() { + return Ok(0); } + let page = history_items_from_joined(joined, filter).rows; let edge = page.last().map(|i| i.height).unwrap_or(0); let heights = self.store.tx_height_get_batch(rest)?; if heights.len() != rest.len() { @@ -529,13 +561,30 @@ impl Query { )); } match filter.order { - HistoryOrder::HeightAsc => Ok(heights.iter().all(|h| i64::from(h.unwrap_or(0)) > edge)), + HistoryOrder::HeightAsc => Ok(heights + .iter() + .take_while(|h| i64::from(h.unwrap_or(0)) <= edge) + .count()), HistoryOrder::NewestFirst => { let ranges = self.store.tx_spent_range_batch(rest)?; - Ok(heights - .iter() - .zip(ranges) - .all(|(h, range)| i64::from(h.unwrap_or(0)) < edge && range.is_none())) + if ranges.len() != rest.len() { + return Err(StoreError::Corrupt( + "invariant: SH spent-range batch length", + )); + } + let reach = self.store.spent_ranges_reach(&ranges, edge)?; + if reach.len() != rest.len() { + return Err(StoreError::Corrupt( + "invariant: SH spent-range batch length", + )); + } + let mut keep = 0usize; + for (i, (h, hot)) in heights.iter().zip(reach).enumerate() { + if i64::from(h.unwrap_or(0)) >= edge || hot { + keep = i.saturating_add(1); + } + } + Ok(keep) } } } @@ -612,6 +661,7 @@ impl Query { let mut class_a_us = 0u128; let mut spends_us = 0u128; let mut offset = 0usize; + let mut stop_at: Option = None; for wave in sh_join_waves(&fks, wave_n) { let t_a = std::time::Instant::now(); let creates = self.expand_create_fks_wave(scripthash, wave, need)?; @@ -620,7 +670,15 @@ impl Query { out.extend(self.join_spends_wave(&creates, need, view)?); spends_us = spends_us.saturating_add(t_s.elapsed().as_micros()); offset = offset.saturating_add(wave.len()); - if paging && self.history_page_closed(&out, page.expect("paging"), &fks[offset..])? { + if !paging { + continue; + } + let filter = page.expect("paging"); + if stop_at.is_none() && Self::history_page_full(&out, filter) { + let keep = self.history_tail_keep(&out, filter, &fks[offset..])?; + stop_at = Some(offset.saturating_add(keep)); + } + if stop_at.is_some_and(|end| offset >= end) { break; } } @@ -962,7 +1020,9 @@ impl Query { &self, scripthash: &[u8; 32], ) -> Result, QueryError> { - self.scripthash_history_filtered(scripthash, &HistoryFilter::open()) + Ok(self + .scripthash_history_filtered(scripthash, &HistoryFilter::open())? + .rows) } /// Confirmed history for `scripthash` as of `view` (open filter). @@ -971,7 +1031,9 @@ impl Query { scripthash: &[u8; 32], view: &ChainView, ) -> Result, QueryError> { - self.scripthash_history_filtered_in(scripthash, &HistoryFilter::open(), view) + Ok(self + .scripthash_history_filtered_in(scripthash, &HistoryFilter::open(), view)? + .rows) } /// Confirmed history for a scripthash, filtered by height window / limit / cursor. @@ -985,9 +1047,12 @@ impl Query { &self, scripthash: &[u8; 32], filter: &HistoryFilter, - ) -> Result, QueryError> { + ) -> Result, QueryError> { let Some(view) = self.pin_sh_chain_view()? else { - return Ok(Vec::new()); + return Ok(FilteredHistory { + rows: Vec::new(), + cursor_missing: filter.after_txid.is_some(), + }); }; self.scripthash_history_filtered_in(scripthash, filter, &view) } @@ -997,7 +1062,7 @@ impl Query { scripthash: &[u8; 32], filter: &HistoryFilter, view: &ChainView, - ) -> Result, QueryError> { + ) -> Result, QueryError> { let joined = self.sh_join_limited( scripthash, ShJoinNeed::HISTORY, @@ -1015,7 +1080,9 @@ impl Query { scripthash: &[u8; 32], slot: &mut Option>, ) -> Result, QueryError> { - self.scripthash_history_filtered_slot(scripthash, &HistoryFilter::open(), slot) + Ok(self + .scripthash_history_filtered_slot(scripthash, &HistoryFilter::open(), slot)? + .rows) } /// Slot-aware [`Self::scripthash_history_filtered`]. @@ -1024,10 +1091,13 @@ impl Query { scripthash: &[u8; 32], filter: &HistoryFilter, slot: &mut Option>, - ) -> Result, QueryError> { + ) -> Result, QueryError> { let Some(view) = self.pin_sh_chain_view()? else { *slot = None; - return Ok(Vec::new()); + return Ok(FilteredHistory { + rows: Vec::new(), + cursor_missing: filter.after_txid.is_some(), + }); }; self.scripthash_history_filtered_slot_in(scripthash, filter, slot, &view) } @@ -1039,7 +1109,7 @@ impl Query { filter: &HistoryFilter, slot: &mut Option>, view: &ChainView, - ) -> Result, QueryError> { + ) -> Result, QueryError> { if self.page_without_full_slot(scripthash, filter, slot, view)? { return self.scripthash_history_filtered_in(scripthash, filter, view); } @@ -1058,7 +1128,7 @@ impl Query { scripthash: &[u8; 32], filter: &HistoryFilter, view: &ChainView, - ) -> Result, QueryError> { + ) -> Result, QueryError> { let joined = self.sh_join_limited( scripthash, ShJoinNeed::HISTORY, @@ -1075,7 +1145,7 @@ impl Query { filter: &HistoryFilter, slot: &mut Option>, view: &ChainView, - ) -> Result, QueryError> { + ) -> Result, QueryError> { if self.page_without_full_slot(scripthash, filter, slot, view)? { return self.scripthash_history_summary_filtered_in(scripthash, filter, view); } @@ -1488,10 +1558,11 @@ mod history_filter_tests { fn open_filter_keeps_all_height_asc() { let items = vec![item(10, 1), item(5, 2), item(20, 3)]; let got = apply_history_filter(&items, &HistoryFilter::open()); - assert_eq!(got.len(), 3); - assert_eq!(got[0].height, 5); - assert_eq!(got[1].height, 10); - assert_eq!(got[2].height, 20); + assert!(!got.cursor_missing); + assert_eq!(got.rows.len(), 3); + assert_eq!(got.rows[0].height, 5); + assert_eq!(got.rows[1].height, 10); + assert_eq!(got.rows[2].height, 20); } #[test] @@ -1500,7 +1571,7 @@ mod history_filter_tests { let f = HistoryFilter::height_window(5, Some(15)); let got = apply_history_filter(&items, &f); assert_eq!( - got.iter().map(|i| i.height).collect::>(), + got.rows.iter().map(|i| i.height).collect::>(), vec![5, 10] ); } @@ -1510,8 +1581,8 @@ mod history_filter_tests { let items = vec![item(1, 1), item(100, 2)]; let f = HistoryFilter::height_window(50, None); let got = apply_history_filter(&items, &f); - assert_eq!(got.len(), 1); - assert_eq!(got[0].height, 100); + assert_eq!(got.rows.len(), 1); + assert_eq!(got.rows[0].height, 100); } #[test] @@ -1521,7 +1592,7 @@ mod history_filter_tests { f.order = HistoryOrder::NewestFirst; let got = apply_history_filter(&items, &f); assert_eq!( - got.iter().map(|i| i.height).collect::>(), + got.rows.iter().map(|i| i.height).collect::>(), vec![3, 2, 1] ); } @@ -1540,13 +1611,17 @@ mod history_filter_tests { order: HistoryOrder::NewestFirst, }; let got = apply_history_filter(&items, &f); - assert_eq!(got.len(), 1); - assert_eq!(got[0].height, 20); - assert_eq!(got[0].txid[0], 2); + assert!(!got.cursor_missing); + assert_eq!(got.rows.len(), 1); + assert_eq!(got.rows[0].height, 20); + assert_eq!(got.rows[0].txid[0], 2); } + /// A cursor that is not in this history is not the first page. Serving + /// those rows made `?after_txid=` of some other confirmed tx restart + /// the wallet at the newest row. #[test] - fn after_txid_unknown_does_not_skip() { + fn after_txid_missing_is_not_the_first_page() { let items = vec![item(10, 1), item(20, 2)]; let f = HistoryFilter { from_height: 0, @@ -1556,8 +1631,8 @@ mod history_filter_tests { order: HistoryOrder::NewestFirst, }; let got = apply_history_filter(&items, &f); - assert_eq!(got.len(), 1); - assert_eq!(got[0].height, 20); + assert!(got.cursor_missing); + assert!(got.rows.is_empty()); } #[test] @@ -1579,8 +1654,8 @@ mod history_filter_tests { order: HistoryOrder::HeightAsc, }; let got = apply_history_filter(&items, &f); - assert_eq!(got.len(), 3); - assert_eq!(got[0].height, 1); - assert_eq!(got[2].height, 3); + assert_eq!(got.rows.len(), 3); + assert_eq!(got.rows[0].height, 1); + assert_eq!(got.rows[2].height, 3); } } diff --git a/crates/rbitcoin-query/src/testutil.rs b/crates/rbitcoin-query/src/testutil.rs index 0cad632ef..17dd8b58d 100644 --- a/crates/rbitcoin-query/src/testutil.rs +++ b/crates/rbitcoin-query/src/testutil.rs @@ -112,6 +112,11 @@ impl FixtureChain for Query { } } +/// One-output join that fits Esplora's last-1 cap. +pub fn sh_join_slot_small() -> Arc { + ShJoinSlot::with_spender_fk_count(0) +} + /// SH join whose packed size exceeds Esplora's 16 MiB last-1 + last-bulk cap. pub fn sh_join_slot_over_16mib() -> Arc { const CAP: usize = 16 * 1024 * 1024; diff --git a/crates/rbitcoin-rpc/src/methods/mine.rs b/crates/rbitcoin-rpc/src/methods/mine.rs index 5b396b492..378e4aa60 100644 --- a/crates/rbitcoin-rpc/src/methods/mine.rs +++ b/crates/rbitcoin-rpc/src/methods/mine.rs @@ -1076,6 +1076,22 @@ pub fn submit_received_block(hub: &rbitcoin_net::ChainHub, block: Block) -> Subm } } +/// True when `block`'s parent header is the active tip. A known ancestor that +/// is not the tip is a competing block; tip spentness does not apply. +fn parent_extends_tip( + query: &rbitcoin_query::Query, + block: &Block, +) -> Result { + use rbitcoin_store::StoreError; + let prev = block.header.prev_blockhash.to_byte_array(); + let parent_fk = match query.get_header_by_hash(&prev) { + Ok(Some((fk, _))) => fk, + Ok(None) | Err(StoreError::NotFound) => return Ok(false), + Err(e) => return Err(e), + }; + Ok(query.tip_header_fk()? == Some(parent_fk)) +} + /// Spend height of `block` when its parent is confirmed. The tip parent is a /// header-fk compare. Any other parent uses the best-chain hash index. fn cheap_spend_height( @@ -1125,6 +1141,7 @@ fn cheap_immature_coinbase( fn read_confirmed_prevout( query: &rbitcoin_query::Query, op: bitcoin::OutPoint, + at_tip: bool, ) -> Result, rbitcoin_store::StoreError> { use rbitcoin_store::StoreError; let tid = op.txid.to_byte_array(); @@ -1138,7 +1155,10 @@ fn read_confirmed_prevout( if op.vout >= rec.output_count { return Ok(None); } - if query.is_outpoint_spent(&tid, op.vout)? { + // Tip spentness is the parent chain only when this block extends the tip. + // A sibling can spend a coin the current tip also spent; caching that as + // invalid drops the competing block for the rest of the process. + if at_tip && query.is_outpoint_spent(&tid, op.vout)? { return Ok(None); } let out = match query @@ -1170,6 +1190,7 @@ fn cheap_submit_tx_reject( ) -> Result, rbitcoin_store::StoreError> { use bitcoin::{OutPoint, TxOut}; let spend_height = cheap_spend_height(query, block)?; + let parent_is_tip = parent_extends_tip(query, block)?; // Core CheckMerkleRoot runs before every body rule. A repeat that keeps // the root (CVE-2012-2459) is the only bad-txns-duplicate. Any other // repeat spends a coin its first copy spent. @@ -1228,7 +1249,7 @@ fn cheap_submit_tx_reject( let txout = if let Some(o) = created.get(&op) { o.clone() } else { - let Some((fk, out)) = read_confirmed_prevout(query, op)? else { + let Some((fk, out)) = read_confirmed_prevout(query, op, parent_is_tip)? else { return Ok(Some("bad-txns-inputs-missingorspent".into())); }; if let Some(h) = spend_height { diff --git a/crates/rbitcoin-rpc/src/methods_tests.rs b/crates/rbitcoin-rpc/src/methods_tests.rs index 57cf41355..85fe712aa 100644 --- a/crates/rbitcoin-rpc/src/methods_tests.rs +++ b/crates/rbitcoin-rpc/src/methods_tests.rs @@ -4354,6 +4354,58 @@ fn submitblock_store_fault_does_not_cache_block_invalid() { let _ = std::fs::remove_dir_all(&dir); } +/// A sibling that spends the coin the current tip also spent is a competing +/// block, not a permanent `duplicate-invalid`. +#[test] +fn submitblock_sibling_spend_is_not_cached_invalid() { + use rbitcoin_primitives::Height; + + let (ctx, dir, hub) = ctx_regtest_hub(); + let op_true = ScriptBuf::from_bytes(vec![0x51]); + dispatch(&ctx, "generate", vec![json!(101)]).unwrap(); + let mature = hub.query.reconstruct_block_at_height(Height(1)).unwrap(); + let spend = Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: bitcoin::absolute::LockTime::ZERO, + input: vec![bitcoin::TxIn { + previous_output: bitcoin::OutPoint { + txid: mature.txdata[0].compute_txid(), + vout: 0, + }, + script_sig: ScriptBuf::new(), + sequence: bitcoin::Sequence::MAX, + witness: bitcoin::Witness::new(), + }], + output: vec![bitcoin::TxOut { + value: Amount::from_sat(1), + script_pubkey: op_true.clone(), + }], + }; + let parent = hub.tip_hash().unwrap(); + let t0 = hub.tip_header().unwrap().time; + let height = hub.tip_height().unwrap() + 1; + let b2 = rbitcoin_consensus::mine_regtest_paying( + parent, + t0 + 1, + height, + op_true.clone(), + vec![spend.clone()], + ); + let r = dispatch(&ctx, "submitblock", vec![json!(block_hex(&b2))]).unwrap(); + assert!(r.is_null(), "{r}"); + let sibling = + rbitcoin_consensus::mine_regtest_paying(parent, t0 + 2, height, op_true, vec![spend]); + let r = dispatch(&ctx, "submitblock", vec![json!(block_hex(&sibling))]).unwrap(); + assert_eq!(r, "inconclusive", "{r}"); + assert!( + !hub.is_block_invalid(&sibling.block_hash()), + "a competing spend must stay submittable" + ); + let again = dispatch(&ctx, "submitblock", vec![json!(block_hex(&sibling))]).unwrap(); + assert_ne!(again, "duplicate-invalid", "{again}"); + let _ = std::fs::remove_dir_all(&dir); +} + /// Core `CheckTxInputs`: a confirmed txid with no output at `vout` is a /// missing input, not a store fault. #[test] diff --git a/crates/rbitcoin-store/src/integrity.rs b/crates/rbitcoin-store/src/integrity.rs index a89ea8914..f8083b091 100644 --- a/crates/rbitcoin-store/src/integrity.rs +++ b/crates/rbitcoin-store/src/integrity.rs @@ -491,6 +491,9 @@ impl Store { self.flush_confirmed_only()?; self.rebuild_height_fence()?; self.persist_class_c_repair()?; + // Disconnect clamps in `Query`. Open revalidation shrinks here and + // then replay treats a marker above the new tip as already annotated. + self.clamp_spend_durable()?; report.tip_shrunk = true; report.tip_after = self.confirmed.tip_height().map(|h| h.0); Ok(()) @@ -674,12 +677,24 @@ mod tests { // Steal tip height 2 to point at G (false conf edge). s.confirmed.set(Height(2), g_fk).unwrap(); s.flush_class_c_tip().unwrap(); + crate::spend_durable::SpendDurable::new(5, 5) + .store(s.path()) + .unwrap(); let r = s.revalidate_tip_window_n(6).unwrap(); assert!(r.tip_shrunk, "must shrink: {r:?}"); assert_eq!(r.first_bad_reason, Some("prev_fk != confirmed parent")); assert_eq!(s.confirmed.tip_height(), Some(Height(1))); assert_eq!(r.tip_after, Some(1)); + let m = crate::spend_durable::SpendDurable::load(s.path()) + .unwrap() + .unwrap(); + assert_eq!(m.annotated_through(), 1); + assert_eq!( + m.durable_through(), + 1, + "a marker above the shrunk tip is clamped" + ); let _ = std::fs::remove_dir_all(&dir); } diff --git a/crates/rbitcoin-store/src/store.rs b/crates/rbitcoin-store/src/store.rs index bc2b0fb1c..9bdf0ecf6 100644 --- a/crates/rbitcoin-store/src/store.rs +++ b/crates/rbitcoin-store/src/store.rs @@ -281,6 +281,8 @@ pub struct Store { mtp_ring: std::sync::RwLock, /// Latest confirm height plus one. Zero means no snapshot yet. spend_snapshot: std::sync::atomic::AtomicU64, + /// Serializes spend-marker publishes so a checkpoint cannot overwrite a clamp. + spend_marker: std::sync::Mutex<()>, /// First height of a confirm write whose spend annotate has not finished, /// plus one. Zero means none. spend_annotate_from: std::sync::atomic::AtomicU64, @@ -373,6 +375,7 @@ impl Store { mtp_ring: std::sync::RwLock::new(MtpRing::empty()), spend_snapshot: std::sync::atomic::AtomicU64::new(0), spend_annotate_from: std::sync::atomic::AtomicU64::new(0), + spend_marker: std::sync::Mutex::new(()), path, cold_path, head_scale: layout.head_scale, @@ -435,6 +438,7 @@ impl Store { mtp_ring: std::sync::RwLock::new(MtpRing::empty()), spend_snapshot: std::sync::atomic::AtomicU64::new(0), spend_annotate_from: std::sync::atomic::AtomicU64::new(0), + spend_marker: std::sync::Mutex::new(()), path, cold_path, head_scale: layout.head_scale, @@ -1135,6 +1139,78 @@ impl Store { .collect()) } + /// True when every spent slot in `ranges` is unspent or its spender's + /// height is below `edge`. A multi-spender overflow is not below: the + /// caller must keep scanning. + pub fn spent_ranges_below(&self, ranges: &[(u64, u64)], edge: i64) -> Result { + let opts: Vec> = ranges.iter().copied().map(Some).collect(); + Ok(self.spent_ranges_reach(&opts, edge)?.iter().all(|hot| !hot)) + } + + /// One flag per range. True when a slot is a multi-spender or a single + /// spender's height is at or above `edge`. `None` and an empty span are cold. + pub fn spent_ranges_reach( + &self, + ranges: &[Option<(u64, u64)>], + edge: i64, + ) -> Result, StoreError> { + const SLOT: u64 = 8; + let mut reach = vec![false; ranges.len()]; + let mut offs = Vec::new(); + let mut owner = Vec::new(); + for (ri, range) in ranges.iter().enumerate() { + let Some((start, len)) = *range else { + continue; + }; + if len == 0 { + continue; + } + if !len.is_multiple_of(SLOT) { + return Err(StoreError::Corrupt("invariant: spent range length")); + } + let n = len / SLOT; + for i in 0..n { + offs.push(start.saturating_add(i.saturating_mul(SLOT))); + owner.push(ri); + } + } + let mut spenders: Vec<(usize, Fk)> = Vec::new(); + for (chunk, own) in offs.chunks(4096).zip(owner.chunks(4096)) { + let metas = self.get_spender_meta_at_abs_batch(chunk)?; + if metas.len() != chunk.len() { + return Err(StoreError::Corrupt("invariant: spent meta batch length")); + } + for (meta, ri) in metas.into_iter().zip(own.iter().copied()) { + let Some((fk, flags, _)) = meta else { + continue; + }; + if flags & crate::compact::output_flags::MULTI_SPENDER != 0 { + reach[ri] = true; + continue; + } + if !fk.is_null() { + spenders.push((ri, fk)); + } + } + } + if spenders.is_empty() { + return Ok(reach); + } + let fks: Vec = spenders.iter().map(|(_, fk)| *fk).collect(); + let heights = self.tx_height_get_batch(&fks)?; + if heights.len() != spenders.len() { + return Err(StoreError::Corrupt( + "invariant: spender height batch length", + )); + } + for ((ri, _), h) in spenders.iter().zip(heights) { + if i64::from(h.unwrap_or(0)) >= edge { + reach[*ri] = true; + } + } + Ok(reach) + } + /// Completion-driven loc→body io_uring pipeline (confirm load / prep). /// /// Jobs with pre-known `range` skip loc fill when `n_out` is already set. @@ -1580,10 +1656,18 @@ impl Store { }; self.txs.sync_replay_bodies()?; self.spenders.flush()?; - crate::spend_durable::SpendDurable::new(tip, tip).store(self.path())?; + self.store_spend_marker(tip, tip)?; Ok(t.elapsed().as_nanos() as u64) } + /// Publish the marker at `min(requested, confirmed tip)` under [`Self::spend_marker`]. + fn store_spend_marker(&self, annotated: u32, durable: u32) -> Result<(), StoreError> { + let _g = self.spend_marker.lock().unwrap_or_else(|e| e.into_inner()); + let tip = self.confirmed.tip_height().map(|h| h.0).unwrap_or(0); + crate::spend_durable::SpendDurable::new(annotated.min(tip), durable.min(tip)) + .store(self.path()) + } + /// Record the confirmed height whose annotations have been written. /// /// This is not a durability claim. The checkpoint thread reads it before @@ -1714,7 +1798,7 @@ impl Store { }, _ => height, }; - crate::spend_durable::SpendDurable::new(height, height).store(self.path()) + self.store_spend_marker(height, height) } /// A disconnect below the marker lowers both heights to the new tip. @@ -1733,6 +1817,7 @@ impl Store { .map_or(0, |h| cur.min(u64::from(h.0) + 1)); (low != cur).then_some(low) }); + let _g = self.spend_marker.lock().unwrap_or_else(|e| e.into_inner()); let Some(marker) = crate::spend_durable::SpendDurable::load(self.path())? else { return Ok(()); }; diff --git a/crates/rbitcoin-test/tests/scenarios.rs b/crates/rbitcoin-test/tests/scenarios.rs index 114d8696f..644c437d2 100644 --- a/crates/rbitcoin-test/tests/scenarios.rs +++ b/crates/rbitcoin-test/tests/scenarios.rs @@ -2078,12 +2078,14 @@ fn pin_scripthash_views_on_pad( ); assert_eq!( q.scripthash_history_filtered(&sh, &HistoryFilter::open()) - .unwrap(), + .unwrap() + .rows, full ); let heights = |f: &HistoryFilter| -> Vec { q.scripthash_history_filtered(&sh, f) .unwrap() + .rows .iter() .map(|i| i.height) .collect() @@ -2110,7 +2112,8 @@ fn pin_scripthash_views_on_pad( &HistoryFilter::esplora_chain_page(None), &view, ) - .unwrap(); + .unwrap() + .rows; assert_eq!(rows.len(), 25); let value_of = |txid: [u8; 32]| rows.iter().find(|r| r.txid == txid).unwrap().value; assert_eq!(value_of(tip_cb), 50_0000_0000);