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
1 change: 1 addition & 0 deletions apps/indexer/src/provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ use decode::parse_hex_i64;
mod block_transaction;
mod code;
mod decode;
mod error;
mod http_client;
mod logs_receipts;
mod ops;
Expand Down
3 changes: 2 additions & 1 deletion apps/indexer/src/provider/block_transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use super::{
ProviderHeadHashSnapshot, ProviderHeadSnapshot, ProviderResolvedBlock,
RAW_PAYLOAD_KIND_FULL_BLOCK,
decode::{block_hash_from_value, normalize_hash},
error::format_provider_error,
logs_receipts::ProviderBlockLogFetch,
provider_batch_item_limit,
request::{JsonRpcBatchCall, is_retryable_provider_error},
Expand Down Expand Up @@ -138,7 +139,7 @@ impl JsonRpcProvider {
service = "indexer",
component = "provider",
tag,
error = %format!("{error:#}"),
error = %format_provider_error(&error),
"provider checkpoint tag is unavailable; degrading to absent optional head"
);
Ok(None)
Expand Down
146 changes: 146 additions & 0 deletions apps/indexer/src/provider/error.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
pub(super) fn format_provider_error(error: &anyhow::Error) -> String {
let mut rendered = format!("{error:#}");
for cause in error.chain() {
if let Some(error) = cause.downcast_ref::<reqwest::Error>() {
redact_reqwest_url(&mut rendered, error);
}
}
rendered
}

pub(super) fn format_provider_transport_error(error: &reqwest::Error) -> String {
let mut rendered = error.to_string();
redact_reqwest_url(&mut rendered, error);
rendered
}

pub(super) fn redact_provider_transport_error_url(error: &mut reqwest::Error) {
let Some(url) = error.url_mut() else {
return;
};
let _ = url.set_username("");
let _ = url.set_password(None);
let _ = url.set_port(None);
url.set_path("");
url.set_query(None);
url.set_fragment(None);
}

fn redact_reqwest_url(rendered: &mut String, error: &reqwest::Error) {
let Some(url) = error.url() else {
return;
};
let redacted_url = match url.host_str() {
Some(host) => format!("{}://{host}", url.scheme()),
None => format!("{}://<redacted-host>", url.scheme()),
};
*rendered = rendered.replace(url.as_str(), &redacted_url);
}

#[cfg(test)]
mod tests {
use std::{future::pending, time::Duration};

use anyhow::{Context, Result};
use tokio::net::TcpListener;

use super::{format_provider_error, format_provider_transport_error};
use crate::provider::{JsonRpcProvider, request::is_retryable_provider_error};

const SECRET_PATH: &str = "/provider-key-secret/v1";
const SECRET_QUERY: &str = "api_key=query-secret";

#[tokio::test]
async fn transport_log_error_keeps_host_and_redacts_path_and_query() -> Result<()> {
let (error, host) = url_bearing_timeout_error().await?;

let rendered = format_provider_transport_error(&error);

assert!(rendered.contains(&format!("for url (http://{host})")));
assert!(!rendered.contains(SECRET_PATH));
assert!(!rendered.contains(SECRET_QUERY));
Ok(())
}

#[tokio::test]
async fn retry_warning_error_is_redacted_and_remains_retryable() -> Result<()> {
let (error, host) = url_bearing_timeout_error().await?;
let error = anyhow::Error::new(error).context("failed to send JSON-RPC request for test");

let rendered = format_provider_error(&error);

assert!(rendered.contains(&format!("for url (http://{host})")));
assert!(!rendered.contains(SECRET_PATH));
assert!(!rendered.contains(SECRET_QUERY));
assert!(rendered.to_ascii_lowercase().contains("timed out"));
assert!(is_retryable_provider_error(&error));
assert!(is_retryable_provider_error(&anyhow::anyhow!(rendered)));
Ok(())
}

#[tokio::test]
async fn retry_exhaustion_error_is_safe_for_debug_log_sinks() -> Result<()> {
let listener = TcpListener::bind("127.0.0.1:0")
.await
.context("failed to bind retry exhaustion test server")?;
let address = listener
.local_addr()
.context("failed to read retry exhaustion test server address")?;
let server = tokio::spawn(async move {
while let Ok((connection, _)) = listener.accept().await {
tokio::spawn(async move {
let _connection = connection;
pending::<()>().await;
});
}
});
let endpoint = format!("http://{address}{SECRET_PATH}?{SECRET_QUERY}");
let provider =
JsonRpcProvider::new_with_request_timeout(&endpoint, Duration::from_millis(25))?;

let error = provider
.fetch_json_rpc_result("eth_chainId", Vec::new())
.await
.expect_err("the provider request must exhaust its timeout retries");
server.abort();
let rendered = format!("{error:?}");

assert!(rendered.contains(&format!("http://{}", address.ip())));
assert!(!rendered.contains(SECRET_PATH));
assert!(!rendered.contains(SECRET_QUERY));
assert!(is_retryable_provider_error(&error));
Ok(())
}

async fn url_bearing_timeout_error() -> Result<(reqwest::Error, String)> {
let listener = TcpListener::bind("127.0.0.1:0")
.await
.context("failed to bind timeout test server")?;
let address = listener
.local_addr()
.context("failed to read timeout test server address")?;
let server = tokio::spawn(async move {
let _connection = listener.accept().await;
pending::<()>().await;
});
let endpoint = format!("http://{address}{SECRET_PATH}?{SECRET_QUERY}");
let client = reqwest::Client::builder()
.timeout(Duration::from_millis(100))
.build()
.context("failed to build timeout test client")?;

let error = client
.get(&endpoint)
.send()
.await
.expect_err("the server must hold the response until the request times out");
server.abort();
assert!(error.is_timeout(), "expected timeout error: {error}");
assert_eq!(
error.url().map(reqwest::Url::as_str),
Some(endpoint.as_str())
);

Ok((error, address.ip().to_string()))
}
}
19 changes: 12 additions & 7 deletions apps/indexer/src/provider/request.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,9 @@ use tracing::warn;

use super::{
JsonRpcProvider,
error::{
format_provider_error, format_provider_transport_error, redact_provider_transport_error_url,
},
http_client::JSON_RPC_POOL_RESET_TIMEOUT_THRESHOLD,
payload_cache::{JsonRpcPayloadFingerprint, JsonRpcResultPayload},
};
Expand Down Expand Up @@ -87,7 +90,7 @@ impl JsonRpcProvider {
method,
attempt = attempt + 1,
max_attempts = MAX_JSON_RPC_ATTEMPTS,
error = %format!("{error:#}"),
error = %format_provider_error(&error),
"retrying transient JSON-RPC provider request"
);
sleep_json_rpc_backoff(attempt).await;
Expand Down Expand Up @@ -189,7 +192,7 @@ impl JsonRpcProvider {
request_context = "batch",
attempt = attempt + 1,
max_attempts = MAX_JSON_RPC_ATTEMPTS,
error = %format!("{error:#}"),
error = %format_provider_error(&error),
"retrying transient JSON-RPC provider batch"
);
sleep_json_rpc_backoff(attempt).await;
Expand Down Expand Up @@ -283,7 +286,8 @@ impl JsonRpcProvider {
.await
{
Ok(response) => response,
Err(error) => {
Err(mut error) => {
redact_provider_transport_error_url(&mut error);
self.record_json_rpc_transport_error(client_generation, request_context, &error);
return Err(error).with_context(|| {
format!("failed to send JSON-RPC request for {request_context}")
Expand All @@ -293,7 +297,8 @@ impl JsonRpcProvider {
let status = response.status();
let body = match response.bytes().await {
Ok(body) => body,
Err(error) => {
Err(mut error) => {
redact_provider_transport_error_url(&mut error);
self.record_json_rpc_transport_error(client_generation, request_context, &error);
return Err(error).context("failed to read JSON-RPC response body");
}
Expand Down Expand Up @@ -336,7 +341,7 @@ impl JsonRpcProvider {
timeout_threshold = JSON_RPC_POOL_RESET_TIMEOUT_THRESHOLD,
previous_client_generation = reset.previous_generation,
new_client_generation = reset.new_generation,
error = %error,
error = %format_provider_transport_error(error),
"rebuilt JSON-RPC HTTP client after a transport timeout"
),
Ok(None) => {}
Expand All @@ -345,8 +350,8 @@ impl JsonRpcProvider {
component = "provider",
request_context,
client_generation,
error = %error,
reset_error = %format!("{reset_error:#}"),
error = %format_provider_transport_error(error),
reset_error = %format_provider_error(&reset_error),
"failed to rebuild JSON-RPC HTTP client after a transport timeout"
),
}
Expand Down
7 changes: 4 additions & 3 deletions apps/indexer/src/provider/transaction_receipts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,8 @@ use tracing::warn;
use super::{
JsonRpcProvider, ProviderBlockBundle, ProviderReceipt, ProviderResolvedBlock,
ProviderTransaction, ProviderTransactionReceiptBundle, ProviderTransactionReceiptRequest,
provider_batch_item_limit, provider_batch_request_concurrency, request::JsonRpcBatchCall,
error::format_provider_error, provider_batch_item_limit, provider_batch_request_concurrency,
request::JsonRpcBatchCall,
};
use validation::{fallback_receipts_by_key, fallback_transactions_by_key};

Expand Down Expand Up @@ -139,7 +140,7 @@ impl JsonRpcProvider {
Ok(block_payloads) => block_payloads,
Err(error) => {
warn!(
error = %format!("{error:#}"),
error = %format_provider_error(&error),
selected_transaction_receipt_fallback_count = requests.len(),
"block-scoped selected transaction/receipt fallback failed; retrying direct lookup"
);
Expand Down Expand Up @@ -374,7 +375,7 @@ impl JsonRpcProvider {
Ok(fallback_bundles) => fallback_bundles,
Err(error) => {
warn!(
error = %format!("{error:#}"),
error = %format_provider_error(&error),
"fallback JSON-RPC provider failed selected transaction/receipt lookup"
);
return bundles;
Expand Down
Loading