Sitelet https://github.com/init4tech/node-components/commit/ef8c2463495387f616c1616324b36d3ad73f1784
Skip to content

Commit ef8c246

Browse files
Evalirclaude
andauthored
fix(rpc): release shard lock before cold-storage awaits in get_filter_changes (#145)
Closes ENG-2139 get_filter_changes held a DashMap RefMut (a parking_lot RwLock write guard) across cold-storage .await points. On a current_thread tokio runtime, two concurrent polls that landed on the same shard could deadlock: the second task parked the OS thread waiting for the lock, leaving no way for the first task to resume. Refactored into snapshot -> cold I/O -> commit, with the RefMut scoped to two short critical sections (no .await inside either). Verified by temporarily forcing DashMap to 2 shards (~50% collision): - Pre-fix: test_rpc_filter_edge_cases hung 10 of 30 runs. - Post-fix: 0 of 30. Default shard count (~128 on dev, ~8 on 2-core CI) made the hang rare enough to slip through into main. Also adds .config/nextest.toml with a 5-minute per-test timeout so any future hang fails fast instead of burning the GitHub Actions 6-hour job ceiling (the failure mode of node-components#134). Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
1 parent ac7c174 commit ef8c246

3 files changed

Lines changed: 62 additions & 37 deletions

File tree

‎.config/nextest.toml‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
# Per-test timeout. Any test that runs longer than `period * terminate-after`
2+
# (= 5 minutes) is killed, instead of being allowed to consume the full
3+
# GitHub Actions job timeout (6 hours). See node-components#134 for the
4+
# incident that motivated this guard.
5+
[profile.default]
6+
slow-timeout = { period = "60s", terminate-after = 5 }

‎crates/rpc/src/eth/endpoints.rs‎

Lines changed: 50 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -1057,54 +1057,59 @@ where
10571057
{
10581058
let task = async move {
10591059
let fm = ctx.filter_manager();
1060-
let mut entry = fm
1061-
.get_mut(id)
1062-
.ok_or_else(|| EthError::InvalidParams(format!("filter not found: {id}")))?;
1063-
1064-
// Scan the global reorg ring buffer for notifications received
1065-
// since this filter's last poll, then lazily compute removed logs
1066-
// and rewind `next_start_block`.
1067-
let reorgs = fm.reorgs_since(entry.last_poll_time());
1068-
let removed = entry.compute_removed_logs(&reorgs);
1069-
if !removed.is_empty() {
1070-
trace!(count = removed.len(), "computed removed logs from reorg ring buffer");
1071-
}
1060+
1061+
// Snapshot filter state under the shard lock, then drop it before
1062+
// any cold-storage `.await`. Holding a `DashMap` `RefMut` across
1063+
// an await can deadlock the current_thread runtime when another
1064+
// task tries to acquire the same shard (parking_lot's RwLock
1065+
// parks the OS thread on contention).
1066+
let (is_block, stored_filter, start, removed, empty_out) = {
1067+
let mut entry = fm
1068+
.get_mut(id)
1069+
.ok_or_else(|| EthError::InvalidParams(format!("filter not found: {id}")))?;
1070+
1071+
// Scan the global reorg ring buffer for notifications received
1072+
// since this filter's last poll, then lazily compute removed
1073+
// logs and rewind `next_start_block`.
1074+
let reorgs = fm.reorgs_since(entry.last_poll_time());
1075+
let removed = entry.compute_removed_logs(&reorgs);
1076+
if !removed.is_empty() {
1077+
trace!(count = removed.len(), "computed removed logs from reorg ring buffer");
1078+
}
1079+
1080+
(
1081+
entry.is_block(),
1082+
entry.as_filter().cloned(),
1083+
entry.next_start_block(),
1084+
removed,
1085+
entry.empty_output(),
1086+
)
1087+
};
10721088

10731089
let latest = ctx.tags().latest();
1074-
let start = entry.next_start_block();
1075-
1076-
// Implicit reorg detection: if latest has moved backward past our
1077-
// window, a reorg occurred that we missed (e.g. broadcast lagged).
1078-
// Return any removed logs we do have, then reset.
1079-
if latest + 1 < start {
1080-
trace!(latest, start, "implicit reorg detected, resetting filter");
1081-
entry.touch_poll_time();
1082-
return Ok(if removed.is_empty() {
1083-
entry.empty_output()
1084-
} else {
1085-
FilterOutput::from(removed)
1086-
});
1087-
}
10881090

1091+
// Early returns: implicit reorg (latest moved backward past our
1092+
// window) or no new blocks since last poll. Either way, just
1093+
// update the poll timestamp and return.
10891094
if start > latest {
1090-
entry.touch_poll_time();
1091-
return Ok(if removed.is_empty() {
1092-
entry.empty_output()
1093-
} else {
1094-
FilterOutput::from(removed)
1095-
});
1095+
if latest + 1 < start {
1096+
trace!(latest, start, "implicit reorg detected, resetting filter");
1097+
}
1098+
if let Some(mut entry) = fm.get_mut(id) {
1099+
entry.touch_poll_time();
1100+
}
1101+
return Ok(if removed.is_empty() { empty_out } else { FilterOutput::from(removed) });
10961102
}
10971103

10981104
let cold = ctx.cold();
10991105

1100-
if entry.is_block() {
1106+
let result = if is_block {
11011107
let specs: Vec<_> = (start..=latest).map(HeaderSpecifier::Number).collect();
11021108
let headers = cold.get_headers(specs).await?;
11031109
let hashes: Vec<B256> = headers.into_iter().flatten().map(|h| h.hash()).collect();
1104-
entry.mark_polled(latest);
1105-
Ok(FilterOutput::from(hashes))
1110+
FilterOutput::from(hashes)
11061111
} else {
1107-
let stored = entry.as_filter().cloned().unwrap();
1112+
let stored = stored_filter.expect("log filter must carry a stored filter spec");
11081113
let resolved = Filter {
11091114
block_option: alloy::rpc::types::FilterBlockOption::Range {
11101115
from_block: Some(BlockNumberOrTag::Number(start)),
@@ -1128,9 +1133,17 @@ where
11281133
logs = combined;
11291134
}
11301135

1136+
FilterOutput::from(logs)
1137+
};
1138+
1139+
// Commit the poll cursor. Tolerate a concurrent uninstall: the
1140+
// caller already received `result`; future polls on this id will
1141+
// fail with "filter not found".
1142+
if let Some(mut entry) = fm.get_mut(id) {
11311143
entry.mark_polled(latest);
1132-
Ok(FilterOutput::from(logs))
11331144
}
1145+
1146+
Ok(result)
11341147
};
11351148

11361149
await_handler!(hctx.spawn(task), EthError::task_panic())

‎crates/rpc/src/interest/filters.rs‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -160,6 +160,12 @@ impl FilterManagerInner {
160160
}
161161

162162
/// Get a filter by ID.
163+
///
164+
/// The returned [`RefMut`] holds a `parking_lot` write lock on a
165+
/// [`DashMap`] shard. Do not hold it across `.await`: on a
166+
/// current_thread runtime a colliding `get_mut` from another task
167+
/// will park the OS thread and deadlock the runtime. (Same family
168+
/// as the [`DashMap::retain`] hazard documented on [`FilterManager`].)
163169
pub(crate) fn get_mut(&self, id: FilterId) -> Option<RefMut<'_, U64, ActiveFilter>> {
164170
self.filters.get_mut(&id)
165171
}

0 commit comments

Comments
 (0)