Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions changelog.d/ibd-getdata-reserve-max.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
Changed

- IBD prices a new getdata hash against the assign-stop at the largest
wire length among the most recent 32 requested bodies, and never above
4 MiB. Until eight bodies have been recorded the charge stays 4 MiB.
A hash already in flight is still issued. The per-peer cap stays the
median of recent bodies (1 KiB until eight), then 64 blocks or 16 MiB.
149 changes: 146 additions & 3 deletions crates/rbitcoin-net/src/ibd/assign.rs
Original file line number Diff line number Diff line change
Expand Up @@ -697,16 +697,29 @@ pub(crate) fn issue_one(
issue_batch(st, pid, vec![h], room, issued)
}

/// Bytes reserved per outstanding getdata hash when the peer did not announce a size.
/// Ceiling on the per-hash assign-stop charge. A block cannot be larger.
pub(crate) const GETDATA_RESERVE_BYTES: u64 = 4 * 1024 * 1024;

/// `inflight_after` unique hashes, each counted at [`GETDATA_RESERVE_BYTES`],
/// Bytes reserved for one new getdata hash against the assign-stop.
///
/// Fewer than [`BLOCK_LEN_MIN_SAMPLES`] recorded bodies still cost the 4 MiB
/// ceiling. After that the charge is the largest wire length in the ring,
/// and never above that ceiling.
fn intake_reserve_per_hash(st: &IbdWorkState) -> u64 {
if st.block_lens.len() < BLOCK_LEN_MIN_SAMPLES {
return GETDATA_RESERVE_BYTES;
}
let max = st.block_lens.iter().copied().max().unwrap_or(u32::MAX);
u64::from(max).min(GETDATA_RESERVE_BYTES)
}

/// `inflight_after` unique hashes, each counted at [`intake_reserve_per_hash`],
/// plus the snapshotted queue, fit in the assign-stop budget.
fn intake_reserve_fits(st: &IbdWorkState, inflight_after: usize) -> bool {
if st.intake_stop == u64::MAX {
return true;
}
let reserved = (inflight_after as u64).saturating_mul(GETDATA_RESERVE_BYTES);
let reserved = (inflight_after as u64).saturating_mul(intake_reserve_per_hash(st));
st.intake_queued.saturating_add(reserved) <= st.intake_stop
}

Expand Down Expand Up @@ -2770,6 +2783,136 @@ pub(in crate::ibd) mod tests {
let _ = std::fs::remove_dir_all(dir);
}

/// Assign-stop charge is the max of the last 32 recorded bodies, capped
/// at 4 MiB. Under 8 samples it stays 4 MiB. A hash already in flight
/// does not need another charge.
#[test]
fn intake_reserve_charges_max_of_recent_bodies() {
let mut st = IbdWorkState::new(vec![dummy_slot(0), dummy_slot(1)], None, None);
let mut room = 64usize;
let mut issued = 0u64;

for _ in 0..7 {
note_block_len(&mut st, 50_000);
}
st.intake_queued = 0;
st.intake_stop = GETDATA_RESERVE_BYTES - 1;
assert!(
!issue_one(&mut st, 0, h(1), &mut room, &mut issued),
"seven samples still need a full 4 MiB"
);
st.intake_stop = GETDATA_RESERVE_BYTES;
assert!(
issue_one(&mut st, 0, h(1), &mut room, &mut issued),
"exactly 4 MiB fits one new hash before eight samples"
);
assert!(
!issue_one(&mut st, 0, h(2), &mut room, &mut issued),
"a second new hash still needs another 4 MiB"
);

st.inflight.clear();
for slot in &mut st.slots {
slot.in_flight.clear();
}
st.block_lens.clear();
let samples = [
100_000u32, 50_000, 200_000, 80_000, 10_000, 90_000, 120_000, 40_000,
];
for n in samples {
note_block_len(&mut st, n as usize);
}
let max_len = u64::from(*samples.iter().max().unwrap());
let queued = 1_000_000u64;
let headroom = max_len * 5 + max_len / 2;
st.intake_queued = queued;
st.intake_stop = queued + headroom;
room = 64;
let mut fit = 0u32;
for i in 0..8u32 {
let before = st.inflight.len();
let _ = issue_one(&mut st, 0, h(100 + i), &mut room, &mut issued);
if st.inflight.len() > before {
fit += 1;
}
}
let by_max = (headroom / max_len) as u32;
let by_flat = (headroom / GETDATA_RESERVE_BYTES) as u32;
assert_eq!(
fit, by_max,
"headroom {headroom} at max {max_len} fits {by_max}, flat 4 MiB would fit {by_flat}"
);

st.inflight.clear();
for slot in &mut st.slots {
slot.in_flight.clear();
}
st.block_lens.clear();
note_block_len(&mut st, 4_000_000);
for _ in 0..32 {
note_block_len(&mut st, 100_000);
}
let window_max = 100_000u64;
st.intake_queued = 0;
st.intake_stop = window_max * 3 + window_max / 2;
room = 16;
fit = 0;
for i in 0..6u32 {
let before = st.inflight.len();
let _ = issue_one(&mut st, 0, h(400 + i), &mut room, &mut issued);
if st.inflight.len() > before {
fit += 1;
}
}
assert_eq!(
fit,
(st.intake_stop / window_max) as u32,
"a max that aged out of the 32-ring must not keep the charge"
);

st.inflight.clear();
for slot in &mut st.slots {
slot.in_flight.clear();
}
st.block_lens.clear();
for _ in 0..8 {
note_block_len(&mut st, 5 * 1024 * 1024);
}
st.intake_queued = 0;
st.intake_stop = GETDATA_RESERVE_BYTES;
room = 8;
assert!(
issue_one(&mut st, 0, h(200), &mut room, &mut issued),
"a recorded length above 4 MiB still charges 4 MiB"
);
assert!(
!issue_one(&mut st, 0, h(201), &mut room, &mut issued),
"the ceiling is one 4 MiB charge, not the recorded length"
);

st.inflight.clear();
for slot in &mut st.slots {
slot.in_flight.clear();
}
st.block_lens.clear();
for _ in 0..8 {
note_block_len(&mut st, max_len as usize);
}
st.intake_queued = 0;
st.intake_stop = max_len;
room = 8;
assert!(issue_one(&mut st, 0, h(300), &mut room, &mut issued));
st.intake_stop = max_len - 1;
assert!(
issue_one(&mut st, 1, h(300), &mut room, &mut issued),
"a hash already in flight is issued with no new charge"
);
assert!(
!issue_one(&mut st, 1, h(301), &mut room, &mut issued),
"a different hash still needs a full charge"
);
}

#[test]
fn hostile_peer_session() {
let _g = BQ_ASSIGN_STOP_ENV_LOCK
Expand Down
4 changes: 2 additions & 2 deletions docs/ibd-memory.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ sidecars and SH BDZ3 occupancy prefixes are **read-only mmap** exceptions

| Structure | Cap / bound | Production clear / evict |
|-----------|-------------|---------------------------|
| **In-RAM body queue** | Soft densify assign (no hysteresis): under ~100 MiB free densify ahead; over ~100 MiB only heights confirm will consume in the next ~1 min at tip rate; at the 1 GiB assign-stop (`RBITCOIN_BLOCK_QUEUE_GB` / `_BYTES`, `0` = unlimited) no new getdata is issued. **Intake budget:** body-queue `bytes()` plus outstanding getdata hashes × `GETDATA_RESERVE_BYTES` (4 MiB each; no per-request announced size) must fit in that same assign-stop before a **new** getdata hash is issued. A hash already requested is never refused. A frame whose hash is not in inflight, or is already queued, is dropped before the payload copy. `bytes()` is **raw only** | Peer **BlockFramed** enqueues **raw** frame payload and stamps Σ `tx.input` via a CompactSize walk (no `Block` decode). Lookup packs/holds on that count; **dequeues** after load-batch send. Decoded `Arc<Block>` + `TxPrecompute` live on **loadq** (cap 14), then scriptq/writeq. **Have-body** (hole / densify / receive) is confirmed ∨ BQ hash ∨ `H ≤ lookup_taken_hi`. **Never both** raw and decoded. **RAM-only by design**. Restart empties BQ+loadq. Logs: `bq soft=n/win RAM=` (`bytes()`) and inflight length × 4 MiB. `loadq=n/14`. |
| **In-RAM body queue** | Soft densify assign (no hysteresis): under ~100 MiB free densify ahead; over ~100 MiB only heights confirm will consume in the next ~1 min at tip rate; at the 1 GiB assign-stop (`RBITCOIN_BLOCK_QUEUE_GB` / `_BYTES`, `0` = unlimited) no new getdata is issued. **Intake budget:** body-queue `bytes()` plus outstanding getdata hashes × the per-hash charge must fit in that same assign-stop before a **new** getdata hash is issued. The charge is 4 MiB until 8 requested bodies have been recorded, then the maximum wire length among the most recent 32, and never above 4 MiB. A hash already requested is never refused. A frame whose hash is not in inflight, or is already queued, is dropped before the payload copy. `bytes()` is **raw only** | Peer **BlockFramed** enqueues **raw** frame payload and stamps Σ `tx.input` via a CompactSize walk (no `Block` decode). Lookup packs/holds on that count; **dequeues** after load-batch send. Decoded `Arc<Block>` + `TxPrecompute` live on **loadq** (cap 14), then scriptq/writeq. **Have-body** (hole / densify / receive) is confirmed ∨ BQ hash ∨ `H ≤ lookup_taken_hi`. **Never both** raw and decoded. **RAM-only by design**. Restart empties BQ+loadq. Logs: `bq soft=n/win RAM=` (`bytes()`). `loadq=n/14`. |
| **Body densify height horizon** | `CONTIG_DENSIFY_AHEAD` (64 k past tip) | Safety max walk/receive; primary gate is soft assign (100 MiB free / 1 min confirm window). |
| **Confirm feed** | readiness (height/hash), no wire retain | Load pack / lookup-wave caps: [`concurrency.md`](./concurrency.md). Requeue / finish on outcome |

Expand Down Expand Up @@ -86,7 +86,7 @@ is over target.
|---------|-----------|
| Limit **densify getdata assign** when BQ payload is over ~100 MiB to heights confirm will consume in the next ~1 min at tip rate | Await a soft gate **before** the next TCP read on a peer |
| At the assign-stop, issue no new getdata. Between ~100 MiB and that stop, densify only the ~1 min confirm window | Drop a body we already **requested** solely for soft budget |
| Do not issue a new getdata hash when `bytes()` + (`inflight` + 1) × 4 MiB would pass assign-stop | Count outstanding getdata by scanning payloads |
| Do not issue a new getdata hash when `bytes()` + (`inflight` + 1) × the per-hash charge would pass assign-stop. The charge is 4 MiB until 8 requested bodies are recorded, then the max of the most recent 32 wire lengths, capped at 4 MiB | Count outstanding getdata by scanning payloads |
| Drop an unsolicited or already-queued frame before copying it into the queue | Refuse `block_queue_offer` for a hash that is still in inflight |
| Free densify ahead while BQ payload is under ~100 MiB | Make healthy peers look stalled by parking the reader on soft backpressure |
| Overshoot soft limits while in-flight requests complete; accept those bodies via `block_queue_offer` | Bound process RAM by refusing peer bytes already requested |
Expand Down
Loading