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
54 changes: 48 additions & 6 deletions bin/dipper-service/src/alerts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,11 +68,13 @@ pub fn layer(config: &AlertsConfig) -> Option<AlertLayer> {

impl<S: Subscriber> Layer<S> 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(),
Expand All @@ -92,18 +94,35 @@ impl<S: Subscriber> Layer<S> 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<String>);

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<String>,
message: String,
fields: String,
}

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() {
Expand Down Expand Up @@ -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);
Expand Down
5 changes: 2 additions & 3 deletions bin/dipper-service/src/cancel_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down
5 changes: 2 additions & 3 deletions bin/dipper-service/src/chain_client/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<B256, ChainClientError> {
// 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 {}",
Expand Down
36 changes: 31 additions & 5 deletions bin/dipper-service/src/chain_client/rpc_provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64, String> {
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] = &[
Expand Down Expand Up @@ -171,17 +181,16 @@ impl RpcProviderPool {
pub async fn latest_blocks(&self) -> Vec<u64> {
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"),
Expand Down Expand Up @@ -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.
Expand Down
7 changes: 3 additions & 4 deletions bin/dipper-service/src/registry/agreement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
13 changes: 9 additions & 4 deletions bin/dipper-service/src/worker/handlers/submit_offer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 { .. }) => {
Expand Down
25 changes: 11 additions & 14 deletions dipper-pgregistry/src/postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}

Expand Down
50 changes: 50 additions & 0 deletions dipper-pgregistry/tests/it_registry_postgres.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<u8>>, 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));
}
Loading