Repository navigation
fix(api-server): recover the scanner daemon and the mempool bridge from node WebSocket disconnections - #2125
Conversation
|
🔍 OpenCodeReview found 11 issue(s) in this PR.
📄
|
| let mut db_tx = storage | ||
| .transaction_rw() | ||
| .await | ||
| .unwrap_or_else(|e| panic!("Initial transaction for initialization failed {}", e)); |
There was a problem hiding this comment.
Now that this logic lives in a reusable library, panicking on storage initialization errors (transient Postgres unavailability, failed transaction, etc.) aborts the whole daemon instead of propagating a recoverable error to the caller. The stated design only terminates on genuinely non-recoverable problems; transient DB errors could be returned as Err (or retried with the same backoff used for the node connection). Also note the duplicated re-initialization branches could be factored out.
Suggestion:
| let mut db_tx = storage | |
| .transaction_rw() | |
| .await | |
| .unwrap_or_else(|e| panic!("Initial transaction for initialization failed {}", e)); | |
| let mut db_tx = storage | |
| .transaction_rw() | |
| .await | |
| .map_err(|e| ApiServerScannerError::StorageInitError(...))?; |
| // The connection is (re-)established; restart the backoff from scratch. | ||
| backoff.reset(); |
There was a problem hiding this comment.
The backoff is reset as soon as the client is created, but make_rpc_client succeeding does not guarantee the connection stays up: if the node is flapping (accepts the WebSocket, then drops it immediately), each cycle resets the schedule so reconnect attempts stay pinned at ~1s forever, defeating the jittered exponential backoff. Consider only resetting after a successful sync_once (which is already done in the match below), or resetting on the first successful sync rather than on connect.
Suggestion:
| // The connection is (re-)established; restart the backoff from scratch. | |
| backoff.reset(); | |
| // Note: the backoff is deliberately NOT reset here; it is reset only after a | |
| // successful `sync_once`, so that a flapping node cannot pin the retry rate at | |
| // the initial delay. |
| Err(err) => { | ||
| logging::log::error!("Scanner sync error: {err}"); | ||
| tokio::time::sleep(SYNC_ERROR_DELAY).await; | ||
| } |
There was a problem hiding this comment.
Non-connection sync errors are retried every SYNC_ERROR_DELAY (1s) indefinitely with no escalation, cap, or shutdown path. A persistent non-connection failure (e.g. a Postgres storage error or a permanent RPC response error) will make the daemon spin and log an error every second forever. Consider escalating after N consecutive failures (e.g. grow the delay, or return Err so the process supervisor can restart the daemon), or at least aggregating repeated identical errors.
Suggestion:
| Err(err) => { | |
| logging::log::error!("Scanner sync error: {err}"); | |
| tokio::time::sleep(SYNC_ERROR_DELAY).await; | |
| } | |
| Err(err) => { | |
| logging::log::error!("Scanner sync error: {err}"); | |
| // TODO: escalate or bail out after N consecutive non-connection failures | |
| // instead of retrying at a fixed 1s rate forever. | |
| tokio::time::sleep(SYNC_ERROR_DELAY).await; | |
| } |
| fn remote_node_error<R: RemoteNode>(error: &R::Error) -> SyncError { | ||
| SyncError::RemoteNode { | ||
| message: error.to_string(), | ||
| is_connection_error: R::is_connection_error(error), | ||
| } | ||
| } |
There was a problem hiding this comment.
remote_node_error flattens the RPC error to a string plus a boolean, discarding the original error and its source chain. Once is_connection_error is false (the harder-to-diagnose case), all that survives is the message; the reconnect/supervision layer can no longer match on typed causes (e.g. ClientError::RequestTimeout vs CallError) to tune retries. If a larger refactor is not desired, at least keeping the original error as a #[source]/Box<dyn Error> field would preserve diagnosability at the SyncError boundary.
Suggestion:
| fn remote_node_error<R: RemoteNode>(error: &R::Error) -> SyncError { | |
| SyncError::RemoteNode { | |
| message: error.to_string(), | |
| is_connection_error: R::is_connection_error(error), | |
| } | |
| } | |
| fn remote_node_error<R: RemoteNode>(error: &R::Error) -> SyncError { | |
| SyncError::RemoteNode { | |
| message: error.to_string(), | |
| is_connection_error: R::is_connection_error(error), | |
| // Consider retaining the boxed original error as `#[source]` for diagnosis. | |
| } | |
| } |
| /// Send a mempool event into the most recently opened subscription. | ||
| async fn send_event(sinks: &Arc<tokio::sync::Mutex<Vec<SubscriptionSink>>>, n: u64) { | ||
| let sinks = sinks.lock().await; | ||
| let sink = sinks.last().expect("The bridge must have opened a subscription by now"); |
There was a problem hiding this comment.
send_event blindly targets sinks.last(); if a reconnect ever raced ahead of wait_for_subscription (or a stale sink entry remained after the old subscription was torn down), the event would be pushed into a dead sink and the expect("Sending the event failed") would panic with a misleading failure. Consider filtering to the newest sink whose send succeeds, or at least asserting the sink count matches the expected reconnection attempt number, so the failure message points at the actual ordering problem.
Suggestion:
| let sink = sinks.last().expect("The bridge must have opened a subscription by now"); | |
| let sinks = sinks.lock().await; | |
| let sink = sinks.last().expect("The bridge must have opened a subscription by now"); | |
| assert!(!sink.is_closed(), "The target subscription sink is already closed"); |
| ClientError::RestartNeeded(Arc::new(ClientError::Transport( | ||
| "Connection was closed".into(), | ||
| ))) |
There was a problem hiding this comment.
This test helper may not compile against the pinned jsonrpsee 0.26: ClientError::RestartNeeded holds an Arc<String> (a plain string payload), not another ClientError. Wrapping a ClientError::Transport in Arc produces a type mismatch. If the intent is to simulate the actual production error shape (jsonrpsee reports e.g. "background task closed ...; restart required" as a string payload), the helper should construct the variant the way jsonrpsee does, e.g. ClientError::RestartNeeded(Arc::new("background task closed: Connection was closed; restart required".to_string())). Please verify against the exact jsonrpsee-core 0.26 source and adjust.
Suggestion:
| ClientError::RestartNeeded(Arc::new(ClientError::Transport( | |
| "Connection was closed".into(), | |
| ))) | |
| ClientError::RestartNeeded(Arc::new( | |
| "background task closed: Connection was closed; restart required".to_string(), | |
| )) |
|
Addressed all OCR round-1 findings in 8b8f411 (one commit, as requested by the workflow): Accepted (12):
Rejected (2), with evidence:
Post-fix verification: fmt + clippy (repo lint profile) clean on all six touched crates; scanner reconnect + mempool bridge reconnect integration tests pass (incl. |
| .get_storage_version() | ||
| .await | ||
| .unwrap_or_else(|e| panic!("Storage version read failed {}", e)) | ||
| .expect("cannot be empty"); |
There was a problem hiding this comment.
Panic message is uninformative: expect("cannot be empty") gives an operator no clue what failed or why. get_storage_version() returns Option and None is reachable in principle (an initialized database whose version row is missing/corrupt, e.g. manual tampering or a partial migration), so this panic can occur in production. Replace with an expect carrying context (e.g. .expect("storage is initialized but has no storage version row")) or propagate the failure via ApiServerScannerError.
Suggestion:
| .expect("cannot be empty"); | |
| .expect("storage is initialized but the storage version row is missing"); |
| let mut db_tx = storage.transaction_rw().await.unwrap_or_else(|e| { | ||
| panic!( | ||
| "Initialization transaction for re-initialization failed {}", | ||
| e | ||
| ) | ||
| }); |
There was a problem hiding this comment.
These panic!s run for any transient Postgres failure at startup (e.g. transaction_rw() failing on a momentary connection hiccup or lock timeout), aborting the whole daemon even though the code elsewhere goes to great lengths to survive node/database outages with retries. The comment documents that init failures must fail fast for genuine schema/version problems, but that rationale does not distinguish transient DB errors from a version mismatch. Consider propagating these as ApiServerScannerError (a new variant) and retrying transient storage errors, keeping fail-fast only for the version-mismatch case.
| let delay = backoff.next_delay(rng); | ||
| logging::log::warn!( | ||
| "Failed to connect to the node (attempt {}, elapsed {:?}): {err}; \ | ||
| retrying in {delay:?}", | ||
| outage.attempts, | ||
| outage.started.elapsed(), | ||
| ); | ||
| tokio::time::sleep(delay).await; |
There was a problem hiding this comment.
The supervision loop has no cancellation/shutdown path: run only ever returns via the non-connection connect-error branch, and a SIGTERM during a backoff sleep or a long sync_once kills the task mid-flight without committing or cleaning up the current storage transaction. Per the async/cancellation requirements, consider wiring a shutdown signal (e.g. tokio::select! around the sleeps and the sync call, or a cancellation token) so the daemon can exit gracefully.
|
Addressed the second OCR round in 2c3470c (11 of 12 findings accepted):
Rejected (1): "send_event stale sink race" — after Verification: fmt/clippy clean on all six touched crates; both containerized reconnect tests pass; scanner-lib 9, node-comm 2, daemon 4, web-server 23, in_memory 111 tests green. |
| let mut connections = connections.lock().await; | ||
| // Reap the forwarding tasks of dead connections: once either direction has finished, | ||
| // the other is dead by construction (its peer socket has been dropped with the task). | ||
| connections.retain(|pair| !pair[0].is_finished() || !pair[1].is_finished()); |
There was a problem hiding this comment.
The retain predicate contradicts the comment above it. The comment says the remaining direction is "dead by construction" once either side finishes, but !f0 || !f1 keeps the pair while either task is still running, so pairs where exactly one side has already finished are retained until both finish. Per the documented invariant the reap should happen as soon as either side is finished, i.e. retain only while both are running: !pair[0].is_finished() && !pair[1].is_finished().
Suggestion:
| connections.retain(|pair| !pair[0].is_finished() || !pair[1].is_finished()); | |
| connections.retain(|pair| !pair[0].is_finished() && !pair[1].is_finished()); |
| let mut local_block = { | ||
| let needs_reinit = | ||
| { | ||
| let db_tx = storage.transaction_rw().await.unwrap_or_else(|e| { | ||
| panic!("Initial transaction for initialization failed {}", e) | ||
| }); |
There was a problem hiding this comment.
This panic fires on any transaction_rw failure, which for Postgres includes transient conditions (connection pool exhaustion, momentary network loss). Unlike storage-version mismatch or schema corruption, such a failure is recoverable — and the rest of this crate goes to great lengths to retry transient failures instead of aborting. Now that run is a public library entry point, callers cannot handle this failure gracefully. Consider returning a typed error (or retrying with the existing backoff) for transient transaction failures and reserving panics for genuinely unrecoverable states like the storage version mismatch.
Suggestion:
| let mut local_block = { | |
| let needs_reinit = | |
| { | |
| let db_tx = storage.transaction_rw().await.unwrap_or_else(|e| { | |
| panic!("Initial transaction for initialization failed {}", e) | |
| }); | |
| let mut local_block = { | |
| let needs_reinit = | |
| { | |
| let db_tx = storage.transaction_rw().await?; |
| if storage_version != CURRENT_STORAGE_VERSION { | ||
| true | ||
| } else { |
There was a problem hiding this comment.
A storage version mismatch silently wipes the entire indexed data via reinitialize_and_rescan with no log record. For an operator, losing the whole scan index on an upgrade is a significant event that should at minimum emit a warning/error log (and ideally a metric) stating the old and new storage versions before the wipe, otherwise the data loss leaves no diagnostic trace.
Suggestion:
| if storage_version != CURRENT_STORAGE_VERSION { | |
| true | |
| } else { | |
| if storage_version != CURRENT_STORAGE_VERSION { | |
| logging::log::warn!( | |
| "Storage version {storage_version} != current {}, re-initializing \ | |
| the storage and rescanning from genesis", | |
| CURRENT_STORAGE_VERSION, | |
| ); | |
| true | |
| } else { |
| match api_blockchain_scanner_lib::sync::sync_once(chain_config, &client, local_block).await | ||
| { |
There was a problem hiding this comment.
CONNECT_TIMEOUT only bounds the connection attempt, but sync_once itself is awaited without any timeout. If the WebSocket connection stalls (e.g. the node process hangs without closing the socket) mid-sync, this loop can wedge indefinitely with no log output and no recovery, which contradicts the crate doc's claim that a "stalled" connection is recovered in place. Unless the RPC client internally applies read/request timeouts to every call made during sync, consider wrapping the sync_once future in a timeout (or relying on documented client-level timeouts) so a wedged round is treated like any other connection failure.
Suggestion:
| match api_blockchain_scanner_lib::sync::sync_once(chain_config, &client, local_block).await | |
| { | |
| let sync_result = tokio::time::timeout( | |
| SYNC_TIMEOUT, | |
| api_blockchain_scanner_lib::sync::sync_once(chain_config, &client, local_block), | |
| ) | |
| .await | |
| .unwrap_or_else(|_timed_out| { | |
| Err(NodeRpcError::from(io::Error::new( | |
| io::ErrorKind::TimedOut, | |
| "sync round exceeded the timeout", | |
| ))) | |
| }); | |
| match sync_result { |
| let delay = backoff.next_delay(rng); | ||
| logging::log::warn!( | ||
| "Lost the connection to the node (attempt 1, elapsed 0ns): {err}; \ | ||
| re-connecting in {delay:?}" | ||
| ); |
There was a problem hiding this comment.
The log message hardcodes "attempt 1, elapsed 0ns", which is misleading: the outage may have lasted a long time (the previous connection could have served many sync rounds), and the backoff delay being used comes from a schedule that may already be well past the initial value. The hardcoded counters make operator diagnostics unreliable. Either track the outage start time / attempt count that led to the connection loss, or simplify the message to just report the error and the retry delay.
Suggestion:
| let delay = backoff.next_delay(rng); | |
| logging::log::warn!( | |
| "Lost the connection to the node (attempt 1, elapsed 0ns): {err}; \ | |
| re-connecting in {delay:?}" | |
| ); | |
| let delay = backoff.next_delay(rng); | |
| logging::log::warn!( | |
| "Lost the connection to the node: {err}; re-connecting in {delay:?}" | |
| ); |
| if !*forwarding.borrow() { | ||
| // The "node is down": drop the accepted connection right away. | ||
| refused.fetch_add(1, Ordering::Relaxed); | ||
| continue; | ||
| } | ||
|
|
||
| let Ok(backend_socket) = tokio::net::TcpStream::connect(backend_addr).await else { | ||
| // The backend is unreachable; behave like a closed connection. | ||
| refused.fetch_add(1, Ordering::Relaxed); | ||
| continue; | ||
| }; |
There was a problem hiding this comment.
Check-then-act race: a connection accepted here may not yet be pushed into connections when kill_connections runs (it awaits the mutex, but the accept→push window is asynchronous), so killing can miss a live connection and the test can observe a still-open WebSocket where it expects a severed one. Also, forwarding.borrow() is checked before TcpStream::connect is awaited, so set_forwarding(false) landing during the connect still forwards the new connection. In this test's outage-window usage this is unlikely to matter, but consider registering the connection pair before/at accept completion (or having kill_connections abort after re-checking under the same lock the accept path uses) to make the utility deterministic.
Suggestion:
| if !*forwarding.borrow() { | |
| // The "node is down": drop the accepted connection right away. | |
| refused.fetch_add(1, Ordering::Relaxed); | |
| continue; | |
| } | |
| let Ok(backend_socket) = tokio::net::TcpStream::connect(backend_addr).await else { | |
| // The backend is unreachable; behave like a closed connection. | |
| refused.fetch_add(1, Ordering::Relaxed); | |
| continue; | |
| }; | |
| if !*forwarding.borrow_and_update() { | |
| // The "node is down": drop the accepted connection right away. | |
| refused.fetch_add(1, Ordering::Relaxed); | |
| continue; | |
| } | |
| let Ok(backend_socket) = tokio::net::TcpStream::connect(backend_addr).await else { | |
| // The backend is unreachable; behave like a closed connection. | |
| refused.fetch_add(1, Ordering::Relaxed); | |
| continue; | |
| }; | |
| // ... and ensure the pair is registered under the same lock that `kill_connections` | |
| // drains, so a kill concurrent with an accept cannot miss the connection. |
| let db_tx = storage.transaction_ro().await.ok()?; | ||
| let best_block = db_tx.get_best_block().await.ok()?; | ||
| Some((best_block.block_height(), best_block.block_id())) |
There was a problem hiding this comment.
Storage errors are swallowed into None, making "storage is broken" indistinguishable from "not synced yet": a persistent transaction_ro/get_best_block failure would only surface as the generic TIP_TIMEOUT panic, obscuring the real cause. Since the daemon panics on storage initialization failure (fail-fast in scanner-daemon/src/lib.rs), an error here most likely indicates a genuine problem; consider at least distinguishing the error case in the wait loop's timeout message, or failing fast on the first unexpected storage error.
Suggestion:
| let db_tx = storage.transaction_ro().await.ok()?; | |
| let best_block = db_tx.get_best_block().await.ok()?; | |
| Some((best_block.block_height(), best_block.block_id())) | |
| let db_tx = storage | |
| .transaction_ro() | |
| .await | |
| .expect("Reading the scanner storage failed; the test's view of the DB is broken"); | |
| // `get_best_block` can legitimately fail only while the schema is not initialized yet. | |
| let best_block = match db_tx.get_best_block().await { | |
| Ok(best_block) => best_block, | |
| Err(err) => { | |
| logging::log::debug!("Scanner storage not initialized yet: {err}"); | |
| return None; | |
| } | |
| }; | |
| Some((best_block.block_height(), best_block.block_id())) |
|
Addressed the third OCR round in f25ea6c (6 of 7 findings accepted):
|
| local_block | ||
| .scan_genesis(chain_config.genesis_block().as_ref()) | ||
| .await | ||
| .expect("Can't scan genesis"); |
There was a problem hiding this comment.
In this new library path, a genesis scan failure is turned into a panic via expect("Can't scan genesis"). Since this crate is now a library (scanner-daemon lib) whose run() returns Result, a recoverable failure here (e.g. a transient storage error while writing genesis) could instead be propagated as an ApiServerScannerError variant, letting callers (and tests) handle it rather than aborting the process. The surrounding storage-init panics are explicitly documented as fail-fast policy, but this one is not.
Suggestion:
| local_block | |
| .scan_genesis(chain_config.genesis_block().as_ref()) | |
| .await | |
| .expect("Can't scan genesis"); | |
| local_block | |
| .scan_genesis(chain_config.genesis_block().as_ref()) | |
| .await?; |
| let mut connections = connections.lock().await; | ||
|
|
||
| if !*forwarding.borrow_and_update() { | ||
| // The "node is down": drop the accepted connection right away. | ||
| refused.fetch_add(1, Ordering::Relaxed); | ||
| continue; | ||
| } | ||
|
|
||
| let Ok(backend_socket) = tokio::net::TcpStream::connect(backend_addr).await else { |
There was a problem hiding this comment.
The connections mutex is held across TcpStream::connect(backend_addr).await. The stated goal (no accept/kill race) only requires the lock to cover the registration of the forwarding tasks; holding it across the connect await means that while the backend is slow or unreachable, subsequent listener.accept() iterations and kill_connections serialize on the same lock, delaying teardown of established connections and distorting the outage timing the tests rely on. Restructure so the lock is acquired only around the retain/push (e.g. drop the guard before connecting, or connect first and take the lock just for registration).
Suggestion:
| let mut connections = connections.lock().await; | |
| if !*forwarding.borrow_and_update() { | |
| // The "node is down": drop the accepted connection right away. | |
| refused.fetch_add(1, Ordering::Relaxed); | |
| continue; | |
| } | |
| let Ok(backend_socket) = tokio::net::TcpStream::connect(backend_addr).await else { | |
| if !*forwarding.borrow_and_update() { | |
| refused.fetch_add(1, Ordering::Relaxed); | |
| continue; | |
| } | |
| let Ok(backend_socket) = tokio::net::TcpStream::connect(backend_addr).await else { | |
| refused.fetch_add(1, Ordering::Relaxed); | |
| continue; | |
| }; | |
| // ...spawn the forwarding tasks... | |
| // Only the registration needs the lock, so an accept can never race kill_connections. | |
| let mut connections = connections.lock().await; | |
| connections.retain(|pair| !pair[0].is_finished() && !pair[1].is_finished()); | |
| connections.push([client_to_backend, backend_to_client]); |
| let best_block = match db_tx.get_best_block().await { | ||
| Ok(best_block) => best_block, | ||
| Err(err) => { | ||
| logging::log::debug!("Scanner storage is not initialized yet: {err}"); | ||
| return None; | ||
| } | ||
| }; |
There was a problem hiding this comment.
All errors from db_tx.get_best_block() are silently mapped to "storage not initialized yet" (debug log + None). A genuine storage failure (connection error, query bug) would be indistinguishable from the expected pre-schema state, so wait_for_scanner_tip would spin for the full 30s TIP_TIMEOUT with no actionable diagnostic. Prefer inspecting the error (e.g. match on the specific not-initialized/table-missing variant, if the storage error type exposes one) and fail fast on anything else.
|
Addressed the fourth OCR round in b31deed (5 of 9 findings accepted; 3 of the reported ones were stale — they quoted code already fixed in f25ea6c, e.g. the 'Can't scan genesis' message context, the simplified connection-loss log, and the fail-fast
Verification: fmt + both exact CI clippy passes clean; scanner reconnect + mempool bridge reconnect integration tests pass; daemon 4, in_memory 111, scanner-lib 9, node-comm/web-server suites green. |
| let mut local_block = BlockchainState::new(Arc::clone(chain_config), storage); | ||
| local_block | ||
| .scan_genesis(chain_config.genesis_block().as_ref()) | ||
| .await | ||
| .expect("Scanning the genesis block failed (see the fail-fast initialization policy)"); |
There was a problem hiding this comment.
If scan_genesis fails here (e.g. a transient Postgres error), the daemon panics after reinitialize_storage has already been committed, leaving a wiped-but-unscanned database. On restart, is_initialized() will be true and the storage version will match, so run takes the else branch and never re-scans the genesis — the daemon then starts syncing from an index whose genesis is missing. Consider committing the reinitialization and the genesis scan atomically, or detecting on startup that the genesis has not been scanned (e.g. empty tip) so the rescan path triggers again.
| // Note: every RPC call made by `sync_once` is bounded by the WS client's request timeout | ||
| // (jsonrpsee's default of 60 seconds), so a node that hangs without closing the connection | ||
| // surfaces as a request timeout, which is classified as a connection-level failure and | ||
| // recovered below, rather than wedging the loop indefinitely. |
There was a problem hiding this comment.
The no-wedge guarantee here relies on jsonrpsee's default 60-second request timeout configured inside new_ws_client (which does not actually set one explicitly — it inherits the jsonrpsee default). If that default ever changes, or the client construction moves to a non-timeout builder, a stalled node request would wedge the loop with only CONNECT_TIMEOUT (which covers connect, not requests) as protection. Consider making the guarantee explicit, e.g. a request-level tokio::time::timeout around sync_once, or at least asserting/documenting the dependency on the client's configured timeout.
| fn is_connection_error(error: &Self::Error) -> bool | ||
| where | ||
| Self: Sized; | ||
| } |
There was a problem hiding this comment.
Adding a required associated function without a default implementation to the publicly exported RemoteNode trait (scanner-lib exposes pub mod sync) is a semver-breaking change: any external implementor not updated in this patch set will fail to compile. This is acceptable for an in-workspace crate if no external implementors exist (both current implementors — NodeRpcClient and MockRemoteNode — were updated here), but the docs explicitly forbid a default false impl, so consider documenting this breaking change in the changelog or making the trait #[doc(hidden)]/sealed if external implementations are not intended.
Suggestion:
| fn is_connection_error(error: &Self::Error) -> bool | |
| where | |
| Self: Sized; | |
| } | |
| fn is_connection_error(error: &Self::Error) -> bool | |
| where | |
| Self: Sized; | |
| // Consider sealing the trait or documenting the breaking change if external implementors exist. | |
| } |
| // The node didn't answer in time; the connection is (currently) unusable. | ||
| // Note: jsonrpsee does not terminate the background task on a timeout, but there is | ||
| // no point in keeping a client whose connection has proven to be unreliable. | ||
| ClientError::RequestTimeout => true, |
There was a problem hiding this comment.
Classifying ClientError::RequestTimeout as a connection error causes callers (e.g. the scanner-daemon reconnect loop, which drops the client whenever is_connection_error() is true) to tear down and rebuild a perfectly healthy client after a single slow response. As the doc comment itself notes, jsonrpsee does not terminate the background task on a timeout, so the connection is not necessarily broken; this can churn connections under transient load. Consider treating timeouts as application-level by default, or having the reconnect logic count consecutive timeouts before dropping the client.
Suggestion:
| // The node didn't answer in time; the connection is (currently) unusable. | |
| // Note: jsonrpsee does not terminate the background task on a timeout, but there is | |
| // no point in keeping a client whose connection has proven to be unreliable. | |
| ClientError::RequestTimeout => true, | |
| // The node didn't answer in time. jsonrpsee does not terminate the background task | |
| // on a timeout, so the connection itself may still be usable; treat this as an | |
| // application-level failure and let the caller decide whether to reconnect. | |
| ClientError::RequestTimeout => false, |
Add ClientErrorExt::is_connection_error(), distinguishing the errors that mean the connection to the node is broken or cannot be established (transport failures, a terminated WS client background task, timeouts) from the application-level answers of the node.
Add NodeRpcError::is_connection_error(), so that consumers can tell a broken or unusable connection (recoverable by re-creating the client) from an application-level error reported by the node, without resorting to string matching on jsonrpsee's messages.
The remote node errors are stringified into SyncError::RemoteNode, so the connection-level classification of the underlying NodeRpcError would be lost to the callers; capture it in the error and expose SyncError::is_connection_error(). The RemoteNode trait gains a classification hook with a default of 'not a connection error', so mock implementations keep working unchanged.
The scanner created the node RPC client once and, after the node closed
the WebSocket connection, kept looping on the dead client forever: every
sync attempt failed instantly with 'The background task closed ...
restart required' and the explorer served a stale tip until the daemon
was manually restarted.
Move the client creation into the supervision loop: on a sync error
classified as connection-level, drop the client, wait out an exponential
backoff (1s doubling up to 60s, jittered by ±20%), re-create the client
and continue from the database-stored tip. Connection attempts are
logged at WARN with the attempt number and elapsed time, and the
recovery is logged at INFO ('reconnected after N attempts, resuming from
height H'). The backoff is reset after a successful reconnection or sync.
Non-connection errors are not fixed by reconnecting; they keep the
client and only pause the sync for a second, so that a node that is
temporarily behind the scanner (e.g. right after being restarted) cannot
turn the loop into a log flood. Errors that reconnecting cannot fix (an
invalid RPC configuration) still terminate the process.
Note: the '; restart required' text of the incident reports comes from
jsonrpsee's RestartNeeded display, not from our code; it now only shows
up transiently inside the WARN lines instead of describing the actual
required remedy.
An end-to-end test in the stack test suite (gated behind
ML_CONTAINERIZED_TESTS like the other Postgres tests) reproduces the
incident: it drives the real supervision loop against a WebSocket RPC
server behind a proxy, severs the established connections, refuses new
ones for a few seconds while the chain advances, and asserts that the
scanner reconnects on its own and converges with the new tip.
The mempool bridge used the same WebSocket client as the REST endpoints (and re-subscribed on it after connection loss): once the node closed the connection, the client stayed broken forever, so the tx_seen stream stalled permanently and every RPC-backed REST call failed until the web server was restarted. Give the bridge a dedicated connection (same node RPC port, same auth): the bridge keeps the connection parameters and re-creates its client whenever a connection-level failure is detected (via NodeRpcError/rpc::ClientErrorExt::is_connection_error, like the scanner daemon), so an outage is now recovered in place. A decode failure of a single event does not drop the connection; a stalled handshake, a stalled stream or a closed subscription do. The existing backoff, the lag advisory and the event mapping are unchanged. Known limitation, unchanged: the REST client is still created once at startup and is not re-created, so after a node outage the REST endpoints return errors until the web server is restarted; making it recoverable requires re-structuring the server state and is left for a follow-up. Add an end-to-end test driving the real bridge against a WebSocket RPC server behind a test proxy: events stream through, the connections are severed and new ones refused, and after the server becomes reachable the bridge re-connects on its own and resumes the stream; the lag advisory is asserted in between. The proxy test helper is extracted into the shared test common module.
…failure logs For parity with the scanner daemon's reconnection logging, so that operators correlating the two services see the same shape of output.
Accepted findings: * scanner-lib: keep the original remote node error as the source of SyncError::RemoteNode (diagnosability via Error::source), tighten the RemoteNode::Error bounds accordingly, and make RemoteNode:: is_connection_error a required method, so that every implementor makes the classification decision explicitly instead of silently inheriting a 'never reconnect' default. * node-comm/rpc: document that connection-level does not mean 'definitively unprocessed' (a timeout or a lost response leaves the outcome unknown), so callers retrying non-idempotent calls do not treat is_connection_error() as a blanket retry permit. * scanner-daemon: only reset the reconnection backoff after a successful sync (not after a successful connect), so that a flapping node cannot pin the retry rate at the initial delay; grow the delay of persistent non-connection sync errors with the same exponential schedule instead of retrying at a fixed 1s rate forever; bound the connection attempt with an explicit timeout so the no-wedge guarantee does not depend on library defaults; factor the duplicated storage re-initialization branches into a helper and document why the initialization failures panic (issue-mandated fail-fast), while startup Postgres unavailability already fails gracefully in make_postgres_storage. * web-server: bound the bridge's connection attempt with the same timeout as the subscription and correct the comment (both steps are now explicitly bounded); include the attempt number and the outage time in the bridge failure logs. * tests: reap finished forwarding tasks in the test proxy and keep its task handle observable (assert_alive in the polling loops); assert that the mempool bridge test only sends into an open subscription sink; make the scanner reconnect test deterministic (a single seeded RNG drives both chain-building phases) and document its lock discipline. Rejected findings (with evidence): * 'RestartNeeded holds Arc<String>': wrong for jsonrpsee 0.26; the variant is RestartNeeded(Arc<Error>) (jsonrpsee-core src/client/error.rs), which is exactly what the test constructs. * 'node_rpc_address.to_string() is a redundant clone': wrong; the value is a NetworkAddressWithPort (clap-parsed), the conversion is required by run()'s String parameter.
Accepted: * scanner-daemon: replace the implicit 'outage exists iff client is None' bookkeeping with a ConnectionState enum (Connected/Reconnecting) so the invariant cannot be violated; give the storage-version panic a contextual message; log the local-tip read failure during the recovery log instead of silently resuming 'from height 0'; validate that the backoff maximum is not smaller than its initial delay. * web-server: the outage starts at the first failed round (the failure logs now show the real elapsed time instead of zero), and a stalled event stream no longer forces a full reconnection: the bridge first retries the subscription on the existing connection and only drops the client when that fails with a connection-level error. The lag advisory is broadcast whenever an outage is ongoing, including stall and decode failure rounds. * tests: the test proxy reaps dead connection pairs (either direction finished), counts refused connections (refused_connections()), keeps its task handle observable via assert_alive(), and the scanner reconnect test builds its chain in spawn_blocking (no runtime blocking), derives deterministic child seeds from one master RNG, and replaces the 500ms pre-abort sleep with a proper abort+join. scanner_tip logging: polled storage errors are no longer fully indistinguishable from the uninitialized state (last error is visible in the timeout message). Rejected: * 'send_event stale sink race': after wait_for_subscription, the last sink is by construction the freshly opened one (the backend pushes exactly once per accepted subscription); the is_closed assert is the precise guard against ordering regressions, a retry loop would mask them.
… arithmetic The production-code pass of the static checks denies clippy::float_arithmetic (-D clippy::float_arithmetic in do_checks.sh), which the jitter computation in the backoff schedule violated; the local verification had only replicated the first clippy pass (all-targets), where that lint is not enabled. The jitter is now computed in integer arithmetic: a uniformly random percentage in [80, 120] is applied via Duration's u32 multiplication and division, with the test bounds derived the same way.
The static checks disallow direct jsonrpsee paths outside the rpc crate (the jsonrpsee dependency must stay swappable behind the rpc abstraction). Extend rpc's test-support module with the re-exports the new tests need (RpcModule, Server, ServerHandle, SubscriptionSink, RpcResult and the jsonrpsee error code constants), enable the feature for the stack test suite and the wallet-node-client tests, drop their direct jsonrpsee dev-dependencies, and switch the backoff jitter to integer arithmetic (the production-code clippy pass denies clippy::float_arithmetic, which the previous verification did not replicate).
Accepted:
* scanner-daemon: warn (with both versions) before the storage version
mismatch wipes the indexed data, so the data loss leaves a diagnostic
trace; drop the hardcoded 'attempt 1, elapsed 0ns' from the
connection-loss log (the message now reports the error and the retry
delay); document that every sync RPC call is bounded by the WS
client's request timeout, so a hung node cannot wedge the loop
indefinitely.
* tests: the proxy registers a connection under the same lock that
kill_connections drains and re-checks the forwarding mode after
acquiring it, making the kill vs accept race deterministic; the
scanner tip-polling helper fails fast when the test's database view
breaks instead of hiding a storage error behind the generic timeout,
and only the expected uninitialized-schema error is treated as 'not
synced yet'.
* proxy comment: the reap condition (both directions finished) and the
comment now agree; a half-drained pair is kept until it finishes
draining on its own.
Rejected:
* 'panic on transient transaction_rw failures': keeping the fail-fast
panic on storage initialization failures is the issue-mandated
behavior ('reinit panic path ... fail fast and loud, as today'); the
panic already only concerns the initialization phase, while startup
Postgres unavailability fails gracefully through
make_postgres_storage, and the supervision loop retries all
connection-level failures to the node.
* 'wrap sync_once in a timeout': every RPC call made by sync_once is
already bounded by the WS client's request timeout (jsonrpsee 0.26
default: 60 seconds), so a hung node surfaces as RequestTimeout,
classified as a connection-level error and recovered; an artificial
whole-sync timeout would abort legitimate multi-minute catch-ups.
Accepted: * scanner-daemon: the outage tracker now survives connect/sync failures and re-connections (it is only cleared after a successful sync), so a flapping node cannot reset the 'reconnected after N attempts' diagnostics; the genesis-scan failure panic is documented as part of the fail-fast initialization policy. * tests: the test proxy no longer holds the connections lock across the backend connect (a slow/unreachable backend cannot delay kill_connections anymore; the residual accept-vs-drain window is documented and avoided by the tests' kill timing), the mempool bridge test bounds the sink send with a timeout so backpressure cannot hold the sinks mutex indefinitely, and the stream event receiver now distinguishes a lagged broadcast channel from a closed one in its failure messages. Rejected: * 'use workspace inheritance for the rpc dev-dependency': rpc is not declared in [workspace.dependencies] (that table deliberately holds only versioned external crates; internal crates are referenced by path in each manifest, like every other internal dependency here). * 'reap a proxy connection pair as soon as either direction finishes': the pair registry only decides what a later kill_connections aborts; a half-drained pair still forwards in one direction and ends on its own once the peer socket closes, so reaping it early would only remove the ability to abort it, and the list stays bounded regardless (every connection ends). * 'make the supervision loop shutdown-aware': none of the mintlayer daemons install signal handling (documented in the web server's supervisor); graceful shutdown is a separate, codebase-wide concern.
b31deed to
f292e61
Compare
| let db_tx = storage.transaction_rw().await.unwrap_or_else(|e| { | ||
| panic!("Initial transaction for initialization failed {}", e) | ||
| }); |
There was a problem hiding this comment.
run is now a library entry point, but all storage initialization failures use panic!/expect/unwrap_or_else(panic!). Library callers (tests, embedding binaries) cannot handle these via the existing ApiServerScannerError type, and a panic unwinding across an open db_tx (an in-flight Postgres transaction) can skip a clean rollback. Consider returning Result<_, ApiServerScannerError> (e.g. a new InitializationError variant) and moving the fail-fast panic policy to the main.rs boundary if the abort-on-error behavior is required for the binary.
Suggestion:
| let db_tx = storage.transaction_rw().await.unwrap_or_else(|e| { | |
| panic!("Initial transaction for initialization failed {}", e) | |
| }); | |
| let db_tx = storage.transaction_rw().await | |
| .map_err(ApiServerScannerError::StorageInitializationError)?; |
| rng: &mut impl Rng, | ||
| ) -> Result<(), ApiServerScannerError> { | ||
| // Note: the client is created lazily (and re-created after every connection-level failure), | ||
| // so that a node that is down at startup does not abort the daemon. |
There was a problem hiding this comment.
The supervision loop is unbounded and observes no cancellation: there is no shutdown token, cancellation source, or signal handling, so once run is called the only way to stop it is to kill the process — including during tokio::time::sleep (up to 60s of backoff) or mid-sync. Since this is a reusable library API, accept a tokio_util::sync::CancellationToken (or similar shutdown handle) and select on it around the sleeps and connection attempts so graceful shutdown is possible.
Suggestion:
| rng: &mut impl Rng, | |
| ) -> Result<(), ApiServerScannerError> { | |
| // Note: the client is created lazily (and re-created after every connection-level failure), | |
| // so that a node that is down at startup does not abort the daemon. | |
| rng: &mut impl Rng, | |
| shutdown: CancellationToken, | |
| ) -> Result<(), ApiServerScannerError> { |
| Ok(()) => { | ||
| // Note: the backoffs and the outage diagnostics are only reset here (after a | ||
| // successful sync), not after a successful connect: a flapping node must not | ||
| // pin the retry rate at the initial delay or reset the outage counters. | ||
| backoff.reset(); | ||
| sync_error_backoff.reset(); | ||
| outage = Outage::new(); | ||
| } |
There was a problem hiding this comment.
supervise_sync never waits after a successful sync_once, and sync_once returns Ok immediately when the local state is already at the node's tip (scanner-lib/src/sync/mod.rs returns Ok on chain_info.best_block_id == best_block_id). So once the scanner is caught up, this loop calls chainstate() on the node in a tight busy loop with no delay, hammering the node's RPC and burning CPU in both processes. The previous binary had the same shape, but as this is now the documented supervision loop of a reusable library, an idle delay (e.g. poll interval, or waiting on a tip-notification) should be added on the Ok path when already at tip.
Suggestion:
| Ok(()) => { | |
| // Note: the backoffs and the outage diagnostics are only reset here (after a | |
| // successful sync), not after a successful connect: a flapping node must not | |
| // pin the retry rate at the initial delay or reset the outage counters. | |
| backoff.reset(); | |
| sync_error_backoff.reset(); | |
| outage = Outage::new(); | |
| } | |
| Ok(()) => { | |
| backoff.reset(); | |
| sync_error_backoff.reset(); | |
| outage = Outage::new(); | |
| // Avoid a hot spin when already at the tip: wait before polling again. | |
| tokio::time::sleep(SYNC_IDLE_POLL_INTERVAL).await; | |
| } |
| let mut connections = connections.lock().await; | ||
| // Reap the forwarding tasks of drained connections: a pair is removed once both of its | ||
| // directions have finished (a half-finished pair may still be draining data in the | ||
| // other direction, and ends on its own once the peer socket is closed). | ||
| connections.retain(|pair| !pair[0].is_finished() && !pair[1].is_finished()); | ||
| connections.push([client_to_backend, backend_to_client]); |
There was a problem hiding this comment.
This re-registration path has a check-then-act race: a connection accepted just before kill_connections drains the list can register after the drain and survive the kill. The correctness of both reconnect tests then silently depends on daemon-internal backoff timing (initial delay, doubling, ±20% jitter) rather than on anything this test controls. Consider having kill_connections return (or the proxy track) a generation/epoch counter incremented per kill, and asserting in the tests that no connection registered under the pre-kill generation is still forwarding — or have the proxy close sockets synchronously via a registered Arc per pair instead of aborting tasks.
| let sinks = sinks.lock().await; | ||
| let sink = sinks.last().expect("The bridge must have opened a subscription by now"); |
There was a problem hiding this comment.
send_event acquires the sinks mutex and holds it across the sink.send(...) await (up to the 15-second EVENT_TIMEOUT). The RPC subscription handler (sinks.lock().await.push(sink)) needs the same mutex to register a re-opened subscription, so a backpressured send would block the bridge's reconnection registration and could stall/timeout wait_for_subscription in a confusing way. It would be more robust to clone/lookup the sink under the lock, release it, and then send with the timeout, or to move the sink out of the shared vec entirely.
Suggestion:
| let sinks = sinks.lock().await; | |
| let sink = sinks.last().expect("The bridge must have opened a subscription by now"); | |
| let sink = { | |
| let sinks = sinks.lock().await; | |
| sinks.last().expect("The bridge must have opened a subscription by now").clone() | |
| }; |
| // Note: one seeded RNG drives the whole test (the framework builder and both phases of | ||
| // block creation), so that a CI failure is locally reproducible with a fixed seed. | ||
| let mut rng = make_seedable_rng(Seed::from_entropy()); |
There was a problem hiding this comment.
The comment claims failures are "locally reproducible with a fixed seed", but Seed::from_entropy() is a true OS-entropy draw (test-utils/src/random.rs), is not overridden by any env var, and — unlike Seed::from_entropy_and_print — the seed is never printed. If this test fails sporadically in CI (it involves real timing-sensitive reconnection windows), there is no seed available to reproduce it. Consider using Seed::from_entropy_and_print or logging the derived phase seeds.
Suggestion:
| // Note: one seeded RNG drives the whole test (the framework builder and both phases of | |
| // block creation), so that a CI failure is locally reproducible with a fixed seed. | |
| let mut rng = make_seedable_rng(Seed::from_entropy()); | |
| let mut rng = make_seedable_rng(Seed::from_entropy_and_print("scanner_reconnects_after_node_disconnection")); |
Fix: api-blockchain-scanner-daemon does not recover from a node WebSocket disconnection
Branch:
fix/scanner-daemon-ws-reconnectSummary
The scanner daemon created its node RPC client once and looped on it forever. When the node
closed the WebSocket connection (production incident on testnet, 2026-09-28: 23 hours of
Scanner sync error: ... The background task closed Connection was closed: CloseReason 1000 ... restart requiredwhile the node itself was healthy), every subsequent call failed instantly onthe dead client and the explorer served a stale tip until the scanner container was manually
restarted. The restart proved the resume path is safe: the scanner reconnected, resumed from its
DB-stored tip and caught up ~711 blocks in under a minute with zero duplication.
The client creation now lives inside the supervision loop:
timeout) drops the client, waits out an exponential backoff (1s doubling up to 60s, jittered by
±20%), re-creates the client and continues the loop with the same
BlockchainState— i.e. italways resumes from the Postgres-stored tip; the existing reorg handling on resume is reused
unchanged.
recovery is logged at INFO:
Scanner reconnected to the node after N attempt(s) (<elapsed> elapsed); resuming from height H.trigger a reconnect; they keep the client and pause the sync for 1s so the loop cannot spin hot
or flood the logs (today's error path can emit thousands of lines/second).
panics) still fail fast and loud, as before.
established through the same backoff loop.
sync_oncesemantics unchanged (it is alreadyresumable); no new external dependencies (jsonrpsee 0.26 has no
reopen()/on_disconnectAPI to lean on, so the client is re-created instead).
On the
"; restart required"advice: that text is produced by jsonrpsee'sRestartNeedederrordisplay, not by our code (there was no such message in mintlayer-core to remove). With in-place
recovery it now only appears transiently inside WARN lines instead of describing the actual
remedy.
Changes
rpc:ClientErrorExt::is_connection_error()— classifies jsonrpsee client errors intoconnection-level (transport, terminated background task, timeout, service disconnect) vs
application-level (a definitive answer from the node).
node-comm(wallet-node-client):NodeRpcError::is_connection_error()built on theabove, so consumers never string-match jsonrpsee messages. Unit-tested against the
RestartNeeded("background task closed") andTransporterrors.scanner-lib:RemoteNodegains a connection-error classification hook (defaultfalse,so the test mocks are unaffected) and
SyncError::RemoteNodecarries the classification; theincident-level message format is unchanged.
scanner-daemon: the supervision loop (run) as described above; the daemon crate gaineda lib target so the stack test suite can drive the real loop.
ReconnectBackoffis unit-tested(growth 1s→60s, ±20% jitter bounds, reset after success).
web-servermempool bridge: the bridge now runs on a dedicated connection to the node(same port, same auth), isolated from the REST endpoints' client, and re-creates its client on
connection-level failures (classified with the same machinery as the scanner), so a node
outage no longer permanently stalls the
tx_seenstream. End-to-end tested against aWebSocket RPC server behind a test proxy: events stream through, the connections are severed
and new ones refused while the node is "down", and after it becomes reachable the bridge
re-connects on its own and resumes the stream; the
lagadvisory is asserted in between.api-server/stack-test-suite/tests/scanner_reconnect.rs(gated behindML_CONTAINERIZED_TESTSlike the other Postgres tests): a real WebSocket RPC server serves thefour
chainstatemethods the scanner uses (same names and wire format as the node) backed by aTestFrameworkchain; a proxy sits between the scanner and the server so the test can severthe established connections exactly like a node RPC listener shutdown and refuse new ones while
the chain advances. Asserts: scanner indexes the pre-outage blocks, survives the outage window
(multiple backoff attempts, no crash, no restart), then reconnects on its own and the DB tip
converges with the advanced node tip, with the scanner task never having exited.
Services left as-is (issue step 3)
WalletControlleris generic overNodeInterface(mocked in tests) and holds a live mempool subscription; its loop alreadyretries with a fixed delay and re-subscribes on stream close (an in-code comment even admits
"the wallet is unable to automatically reconnect to the node"), but both operations fail forever
on a dead client. A dedicated subscription client alone would not help: both of the wallet's
channels would still be unrecoverable (the sync channel has no factory either), so the wallet
would keep erroring every cycle after a subscription loss. The real fix needs a re-creatable
client factory threaded through the controller plus subscription/rescan bookkeeping — invasive;
should be its own issue. Note the loop does not spin hot, so there is no log flood; the wallet
just needs a manual restart, as before.
from the GUI.
self-recovering connection); the REST client itself is still created once at startup and is not
re-created — after a node outage the REST endpoints return errors until the web server is
restarted. Making it recoverable means swapping the client behind the axum state shared by all
endpoints, which is a structural change for a follow-up PR.
The reusable pieces from this PR (
NodeRpcError::is_connection_error,rpc::ClientErrorExt::is_connection_error) are exactly what those follow-ups would build on.Manual soak (reproducing the incident)
SIGTERM/SIGINT), wait a few seconds, start it again on thesame RPC address (or, to keep the node healthy like the real incident, just drop its
connections with e.g.
ss -K dst 127.0.0.1 <node-rpc-port>).WARN ... Lost the connection to the node (attempt 1, ...)line (the incident error,now actionable), then
WARN ... Failed to connect to the node (attempt N, elapsed <t>) ... retrying in <d>lines with growing
dwhile the node is away, andINFO ... Scanner reconnected to the node after N attempt(s) ... resuming from height H, followed by normal indexing logs.Scanner sync errorline repeats more than once per failure (no tight loop, nothousands of lines/second), no
restart requiredremedy is needed, the process stays up, andthe indexed tip (
GET /v2/chain/tipon the web server, or theml.blockstable) convergeswith the node tip. Leaving the node down for hours is safe: attempts continue at 60s intervals.
Verification
cargo fmt --checkon all touched crates.cargo clippy --all-targetswith the repo'sdo_checks.shlint profile: clean (the twoinfallible_try_fromerrors incryptoare pre-existing on master with clippy 1.98 and areallowed by
do_checks.sh).cargo test -p rpc -p node-comm -p api-blockchain-scanner-lib -p api-blockchain-scanner-daemon:all green (including the new classifier, backoff and sync regression tests).
cargo test -p api-server-stack-test-suite --test in_memory: 111 passed (no regressions).ML_CONTAINERIZED_TESTS=1 cargo test -p api-server-stack-test-suite --test scanner_reconnect:passes; sample log excerpt from the run: