From 1777be676d2324578902faad1f8355b8d6ed2fdf Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Thu, 8 Oct 2026 21:04:33 +0100 Subject: [PATCH 1/4] fix(rpc): hide endpoint API keys in cross-check failure logs The latest-block cross-check logged each endpoint's raw error, which names its URL, and hosted endpoints carry their API key in the URL. It now logs the same redacted description other calls use. --- .../src/chain_client/rpc_provider.rs | 36 ++++++++++++++++--- 1 file changed, 31 insertions(+), 5 deletions(-) diff --git a/bin/dipper-service/src/chain_client/rpc_provider.rs b/bin/dipper-service/src/chain_client/rpc_provider.rs index 836c6b98..d6e8bf56 100644 --- a/bin/dipper-service/src/chain_client/rpc_provider.rs +++ b/bin/dipper-service/src/chain_client/rpc_provider.rs @@ -53,6 +53,16 @@ fn describe_failure(url: &Url, error: &TransportError) -> String { .replace(url.as_str().trim_end_matches('/'), &name) } +/// 1 endpoint's latest block, asked once. A failure is described without the URL, since +/// hosted endpoints carry their API key in it. +async fn latest_block(http: reqwest::Client, url: &Url) -> Result { + let provider = ProviderBuilder::new().connect_reqwest(http, url.clone()); + provider + .get_block_number() + .await + .map_err(|err| describe_failure(url, &err)) +} + /// Error text that indicates a transient failure worth retrying, used only for faults /// that arrive as prose rather than as a status code or JSON-RPC error object. const RETRYABLE_ERROR_PATTERNS: &[&str] = &[ @@ -171,17 +181,16 @@ impl RpcProviderPool { pub async fn latest_blocks(&self) -> Vec { let mut asks = tokio::task::JoinSet::new(); for url in &self.providers { - let provider = ProviderBuilder::new().connect_reqwest(self.http.clone(), url.clone()); - let endpoint = endpoint_name(url); - asks.spawn(async move { (endpoint, provider.get_block_number().await) }); + let (http, url) = (self.http.clone(), url.clone()); + asks.spawn(async move { (endpoint_name(&url), latest_block(http, &url).await) }); } let mut heads = Vec::with_capacity(self.providers.len()); while let Some(answer) = asks.join_next().await { match answer { Ok((_, Ok(head))) => heads.push(head), - Ok((endpoint, Err(err))) => tracing::debug!( + Ok((endpoint, Err(reason))) => tracing::debug!( provider = %endpoint, - error = %err, + error = %reason, "RPC endpoint didn't give its latest block for a cross-check" ), Err(err) => tracing::warn!(error = %err, "Latest-block cross-check task failed"), @@ -592,6 +601,23 @@ mod tests { ); } + /// The latest-block cross-check logs each endpoint's failure, so it must hide the key too. + #[tokio::test] + async fn a_cross_check_failure_hides_the_api_key() { + let keyed: Url = "http://127.0.0.1:1/v2/super-secret-key" + .parse() + .expect("keyed endpoint URL"); + + let reason = latest_block(reqwest::Client::new(), &keyed) + .await + .expect_err("nothing is listening, so the ask fails"); + + assert!( + !reason.contains("super-secret-key"), + "the API key must not appear in the failure: {reason}" + ); + } + /// An endpoint that answers can only describe its own refusal, so it never repeats the /// URL. One that never answers is described by the HTTP client instead, which says which /// URL it was reaching for, and that is where the key sits. Nothing listens on port 1. From 82613d545fc03933c4a4c4a9c8e4b8a70e95939f Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Thu, 8 Oct 2026 21:04:33 +0100 Subject: [PATCH 2/4] fix(registry): keep the offer hash of an agreement being cancelled An agreement can be marked for cancelling while its offer is mining, and the offer can still land. Its hash is now stored then, without moving the time that paces its cancel retry, and an ended agreement reports that nothing was stored instead of skipping silently. --- bin/dipper-service/src/registry/agreement.rs | 7 ++- .../src/worker/handlers/submit_offer.rs | 13 +++-- dipper-pgregistry/src/postgres.rs | 25 ++++------ .../tests/it_registry_postgres.rs | 50 +++++++++++++++++++ 4 files changed, 73 insertions(+), 22 deletions(-) diff --git a/bin/dipper-service/src/registry/agreement.rs b/bin/dipper-service/src/registry/agreement.rs index 06e0ed96..95da2b73 100644 --- a/bin/dipper-service/src/registry/agreement.rs +++ b/bin/dipper-service/src/registry/agreement.rs @@ -247,10 +247,9 @@ pub trait AgreementRegistry { id: &IndexingAgreementId, ) -> RegistryResult<()>; - /// Record the on-chain tx hash of the most recent `offer()` submission - /// for this agreement. Observability-only; does not transition status. - /// Called once per submit (including resubmits after a dropped tx) so - /// the DB reflects the live hash rather than an evicted one. + /// Record the hash of the latest `offer()` transaction, unless the agreement has ended, so + /// a resubmit replaces a dropped one. [`NoRecordUpdated`](Error::NoRecordsUpdated) when no + /// row took it. async fn update_offer_tx_hash( &self, id: &IndexingAgreementId, diff --git a/bin/dipper-service/src/worker/handlers/submit_offer.rs b/bin/dipper-service/src/worker/handlers/submit_offer.rs index f77fa2ae..f5ee5088 100644 --- a/bin/dipper-service/src/worker/handlers/submit_offer.rs +++ b/bin/dipper-service/src/worker/handlers/submit_offer.rs @@ -121,19 +121,24 @@ where tx_hash = %tx_hash, "Offer submitted on-chain successfully" ); - // Observability only: record which tx hash actually mined. // Any failure here is non-fatal to the overall flow. - if let Err(err) = ctx + match ctx .registry .update_offer_tx_hash(agreement_id, tx_hash.as_ref()) .await { - tracing::warn!( + Ok(()) => {} + Err(crate::registry::Error::NoRecordsUpdated) => tracing::debug!( + agreement_id = %agreement_id, + tx_hash = %tx_hash, + "Agreement ended while its offer was mining, so its offer_tx_hash isn't stored" + ), + Err(err) => tracing::warn!( agreement_id = %agreement_id, tx_hash = %tx_hash, error = %err, "Failed to persist offer_tx_hash; continuing" - ); + ), } } Err(err @ ChainClientError::TxDropped { .. }) => { diff --git a/dipper-pgregistry/src/postgres.rs b/dipper-pgregistry/src/postgres.rs index 6af8d05a..d1bc8194 100644 --- a/dipper-pgregistry/src/postgres.rs +++ b/dipper-pgregistry/src/postgres.rs @@ -914,37 +914,34 @@ impl PgRegistry { Ok(()) } - /// Persist the on-chain tx hash of the most recent `offer()` submission - /// for this agreement. Overwrites any prior value, so a resubmit after - /// mempool eviction records the live hash rather than the dropped one. - /// Observability-only: no status transition is performed here. - /// - /// Guarded on `status IN (Created, AcceptedOnChain)` so a delayed - /// receipt-confirmation cannot stamp `offer_tx_hash` onto a row that - /// has since transitioned to `Expired`, `Unresponsive`, `Rejected`, - /// or one of the cancel states. The caller treats any failure here - /// as non-fatal and just logs; a no-match result is also non-fatal - /// and silently skipped. + /// Record the hash of the latest `offer()` transaction, unless the agreement has ended. A + /// `Cancelling` row keeps its `updated_at`, which says when it was marked and paces its + /// cancel retry. Returns [`Error::NoRecordsUpdated`] when no row took the hash. pub async fn update_offer_tx_hash( &self, agreement_id: &IndexingAgreementId, tx_hash: &[u8; 32], ) -> Result<(), Error> { - sqlx::query( + let updated = sqlx::query( r#" UPDATE dipper_reg_indexing_agreements SET offer_tx_hash = $1, - updated_at = timezone('UTC', now()) - WHERE id = $2 AND status IN ($3, $4) + updated_at = CASE WHEN status = $5 THEN updated_at + ELSE timezone('UTC', now()) END + WHERE id = $2 AND status IN ($3, $4, $5) "#, ) .bind(&tx_hash[..]) .bind(agreement_id) .bind(IndexingAgreementStatus::Created) .bind(IndexingAgreementStatus::AcceptedOnChain) + .bind(IndexingAgreementStatus::Cancelling) .execute(&self.pool) .await?; + if updated.rows_affected() == 0 { + return Err(Error::NoRecordsUpdated); + } Ok(()) } diff --git a/dipper-pgregistry/tests/it_registry_postgres.rs b/dipper-pgregistry/tests/it_registry_postgres.rs index 2f4c194f..e5bd0f3a 100644 --- a/dipper-pgregistry/tests/it_registry_postgres.rs +++ b/dipper-pgregistry/tests/it_registry_postgres.rs @@ -3834,3 +3834,53 @@ async fn a_cancelling_agreement_stays_live_and_unannounced_until_it_ends() { .expect("terminated query"); assert!(terminated.iter().any(|p| p.agreement_id == cancelling)); } + +/// Reassess can mark an agreement `Cancelling` while its offer is mining, and the offer can +/// still land, so its hash is worth keeping. Once the agreement has ended it isn't. +#[tokio::test] +async fn an_offer_mined_after_the_cancel_began_keeps_its_hash() { + let (db, _temp_db) = temp_registry_db().await; + run_fixture( + &db, + include_str!("fixtures/0003_multi_indexer_agreements.sql"), + ) + .await + .expect("Failed to run fixture"); + let registry = PgRegistry::new(db.clone()); + let id = fixture_agreement(0xaa); + registry + .mark_indexing_agreement_as_cancelling(&id) + .await + .expect("mark cancelling"); + let stored = async || { + sqlx::query_as::<_, (Option>, time::OffsetDateTime)>( + "SELECT offer_tx_hash, updated_at FROM dipper_reg_indexing_agreements WHERE id = $1", + ) + .bind(id) + .fetch_one(&db) + .await + .expect("read the agreement") + }; + let (_, marked_at) = stored().await; + + registry + .update_offer_tx_hash(&id, &[0x11; 32]) + .await + .expect("a cancelling agreement takes the hash"); + assert_eq!( + stored().await, + (Some(vec![0x11; 32]), marked_at), + "hash stored, and the time it was marked kept" + ); + + sqlx::query("UPDATE dipper_reg_indexing_agreements SET status = 5 WHERE id = $1") + .bind(id) + .execute(&db) + .await + .expect("expire the agreement"); + let err = registry + .update_offer_tx_hash(&id, &[0x22; 32]) + .await + .expect_err("an ended agreement doesn't take the hash"); + assert!(matches!(err, Error::NoRecordsUpdated)); +} From 8463c5db0aef1e211d7c2e26261f53cf4d2acdd1 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Thu, 8 Oct 2026 21:04:33 +0100 Subject: [PATCH 3/4] perf(alerts): read only the event tag of lines no alert is set up for Every warning passes through the Slack alert layer, which formatted all of a line's fields before checking its tag. It now reads the tag first and formats the rest only for listed events. --- bin/dipper-service/src/alerts.rs | 54 ++++++++++++++++++++++++++++---- 1 file changed, 48 insertions(+), 6 deletions(-) diff --git a/bin/dipper-service/src/alerts.rs b/bin/dipper-service/src/alerts.rs index d2460b08..8147ceef 100644 --- a/bin/dipper-service/src/alerts.rs +++ b/bin/dipper-service/src/alerts.rs @@ -68,11 +68,13 @@ pub fn layer(config: &AlertsConfig) -> Option { impl Layer for AlertLayer { fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) { - let mut line = LogLine::default(); - event.record(&mut line); - let Some(tag) = line.event.filter(|tag| self.events.contains(tag)) else { + let mut tag = EventTag::default(); + event.record(&mut tag); + let Some(tag) = tag.0.filter(|tag| self.events.contains(tag)) else { return; }; + let mut line = LogLine::default(); + event.record(&mut line); let alert = Alert { event: tag, level: *event.metadata().level(), @@ -92,10 +94,27 @@ impl Layer for AlertLayer { } } -/// A log line's `event` tag, message and other fields. +/// Only a log line's `event` tag, read first so a line no alert is set up for costs no more. +#[derive(Default)] +struct EventTag(Option); + +impl Visit for EventTag { + fn record_str(&mut self, field: &Field, value: &str) { + if field.name() == "event" { + self.0 = Some(value.to_owned()); + } + } + + fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) { + if field.name() == "event" { + self.0 = Some(format!("{value:?}")); + } + } +} + +/// A log line's message and its fields other than the `event` tag. #[derive(Default)] struct LogLine { - event: Option, message: String, fields: String, } @@ -103,7 +122,7 @@ struct LogLine { impl LogLine { fn record(&mut self, field: &Field, value: String) { match field.name() { - "event" => self.event = Some(value), + "event" => {} "message" => self.message = value, name => { if !self.fields.is_empty() { @@ -355,6 +374,29 @@ mod tests { assert!(alerts.try_recv().is_err(), "nothing else"); } + /// Every warning in dipper passes through this layer, so 1 it won't post shouldn't cost + /// formatting its fields. + #[test] + fn leaves_the_fields_of_lines_it_wont_post_unformatted() { + struct Counted<'a>(&'a std::sync::atomic::AtomicUsize); + impl std::fmt::Debug for Counted<'_> { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + self.0.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + f.write_str("counted") + } + } + let formatted = std::sync::atomic::AtomicUsize::new(0); + let (layer, _alerts) = test_layer(&["agreement_cancel_stuck"]); + let subscriber = tracing_subscriber::registry().with(layer); + + tracing::subscriber::with_default(subscriber, || { + tracing::warn!(event = "something_else", value = ?Counted(&formatted), "Not listed"); + tracing::warn!(value = ?Counted(&formatted), "Not tagged"); + }); + + assert_eq!(formatted.load(std::sync::atomic::Ordering::Relaxed), 0); + } + #[test] fn counts_alerts_dropped_when_the_queue_is_full() { let (queue, _alerts) = mpsc::channel(1); From 18752e32768396afa5007a73aa0692f1db9469a1 Mon Sep 17 00:00:00 2001 From: MoonBoi9001 Date: Thu, 8 Oct 2026 21:04:33 +0100 Subject: [PATCH 4/4] docs(cancel): drop comments saying ConfigError switches off cancels 2 comments said the cancel path reads ConfigError as the chain client being switched off, which no code does any more. --- bin/dipper-service/src/cancel_dispatch.rs | 5 ++--- bin/dipper-service/src/chain_client/client.rs | 5 ++--- 2 files changed, 4 insertions(+), 6 deletions(-) diff --git a/bin/dipper-service/src/cancel_dispatch.rs b/bin/dipper-service/src/cancel_dispatch.rs index 54e50521..bd3674db 100644 --- a/bin/dipper-service/src/cancel_dispatch.rs +++ b/bin/dipper-service/src/cancel_dispatch.rs @@ -530,9 +530,8 @@ pub(crate) mod tests { #[tokio::test] async fn manager_cancel_missing_hash_is_distinct_error_and_sends_nothing() { - // eh-1: a missing hash must be the distinct MissingTermsVersionHash, not - // a ConfigError the liveness checker reads as "chain client disabled" - // and would silently abandon while the agreement stays live on-chain. + // A missing hash must be the distinct MissingTermsVersionHash: no cancel can be sent + // without it, so the cancel retry spends every attempt at once and raises the alert. let client = RecordingChainClient::default(); let ag = agreement(IndexingAgreementStatus::AcceptedOnChain, None); diff --git a/bin/dipper-service/src/chain_client/client.rs b/bin/dipper-service/src/chain_client/client.rs index 7ca8a9c5..69122dea 100644 --- a/bin/dipper-service/src/chain_client/client.rs +++ b/bin/dipper-service/src/chain_client/client.rs @@ -734,9 +734,8 @@ impl AlloyChainClient { /// already does. Signing happens once, up front, so every endpoint is offered the same /// bytes under one hash and the hash is known before anyone is asked to accept them. async fn send_transaction(&self, tx: &TransactionRequest) -> Result { - // Nothing fills a field in on this path any more, and a request that names no chain is - // signed for chain 1 rather than refused, so check before the signature exists. Not a - // `ConfigError`: the cancel path reads that as the chain client being switched off. + // Nothing fills a field in on this path, and a request that names no chain is signed + // for chain 1 rather than refused, so check before the signature exists. if tx.chain_id() != Some(self.inner.chain_id) { return Err(ChainClientError::SubmitFailed(anyhow::anyhow!( "refusing to sign for chain {:?} while configured for chain {}",