From 8ced249a2900742267cbe3dd1ed977b0dac7b9e3 Mon Sep 17 00:00:00 2001 From: Gospelsam019 Date: Sat, 26 Sep 2026 18:01:31 +0000 Subject: [PATCH] fix: address webhook, config, index, and ingest issues --- crates/api/src/routes/sponsorship.rs | 10 ++++++++- crates/api/src/routes/whitelist.rs | 7 ++++++- crates/ingest/src/lib.rs | 24 +++++++++++++++------ crates/store/src/lib.rs | 14 ++++++++++++- crates/webhooks/src/lib.rs | 31 +++++++++++++++++++++++++--- docs/api.md | 2 ++ 6 files changed, 76 insertions(+), 12 deletions(-) diff --git a/crates/api/src/routes/sponsorship.rs b/crates/api/src/routes/sponsorship.rs index af38095..1b3fc9c 100644 --- a/crates/api/src/routes/sponsorship.rs +++ b/crates/api/src/routes/sponsorship.rs @@ -62,7 +62,7 @@ pub async fn get_config( Ok(Envelope::ok(view)) } -/// `PUT /v1/wallets/:id/sponsorship` +/// `PUT /v1/wallets/:id/sponsorship`; requires at least one field (see `docs/api.md`). pub async fn put_config( State(state): State, Path(wallet_id): Path, @@ -71,6 +71,14 @@ pub async fn put_config( ) -> ApiResult>> { authorize_wallet(&headers, &state, wallet_id).await?; let req: SponsorshipConfigRequest = parse_optional(&body)?; + if req.enabled.is_none() + && req.per_tx_fee_cap_stroops.is_none() + && req.daily_budget_stroops.is_none() + { + return Err(ApiError::BadRequest( + "at least one field must be provided".into(), + )); + } let enabled = req.enabled.unwrap_or(false); if let Some(cap) = req.per_tx_fee_cap_stroops { diff --git a/crates/api/src/routes/whitelist.rs b/crates/api/src/routes/whitelist.rs index 38ba18b..c174d9e 100644 --- a/crates/api/src/routes/whitelist.rs +++ b/crates/api/src/routes/whitelist.rs @@ -57,7 +57,7 @@ pub async fn get_config( Ok(Envelope::ok(AllowlistConfigView { enabled })) } -/// `PUT /v1/wallets/:id/whitelist/config` +/// `PUT /v1/wallets/:id/whitelist/config`; requires at least one field (see `docs/api.md`). pub async fn put_config( State(state): State, Path(wallet_id): Path, @@ -66,6 +66,11 @@ pub async fn put_config( ) -> ApiResult>> { owned_wallet(&state, &headers, wallet_id).await?; let req: AllowlistConfigRequest = parse_optional(&body)?; + if req.enabled.is_none() { + return Err(ApiError::BadRequest( + "at least one field must be provided".into(), + )); + } let enabled = req.enabled.unwrap_or(false); // Enabling with an empty list would lock the wallet out of every destination — including diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 1ed809e..cc57cce 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -589,8 +589,10 @@ impl Supervisor { .await?; let semaphore = Arc::new(tokio::sync::Semaphore::new(Self::MAX_CONCURRENT_POLLS)); let mut tasks = tokio::task::JoinSet::new(); + let mut task_wallets = HashMap::new(); for w in wallets { + let wallet_id = w.id; let store = self.store.clone(); let store_for_mark = self.store.clone(); let horizon_url = self.horizon_url.clone(); @@ -599,7 +601,7 @@ impl Supervisor { let retry = self.retry.clone(); let circuit = self.circuit.clone(); let semaphore = semaphore.clone(); - tasks.spawn(async move { + let task_id = tasks.spawn(async move { // Held for the duration of this wallet's poll; bounds how many Horizon requests // are in flight at once without limiting how many wallets we *queue*. let _permit = semaphore.acquire_owned().await; @@ -617,18 +619,28 @@ impl Supervisor { // Record the attempt regardless of outcome, so a wallet whose polls keep failing // still backs off instead of being retried at full rate forever. let _ = store_for_mark.mark_polled(w.id).await; - (w.id, result) + (wallet_id, result) }); + task_wallets.insert(task_id.id(), wallet_id); } let mut total = 0; - while let Some(joined) = tasks.join_next().await { + while let Some(joined) = tasks.join_next_with_id().await { match joined { - Ok((_wallet_id, Ok(n))) => total += n, - Ok((wallet_id, Err(e))) => { + Ok((task_id, (_wallet_id, Ok(n)))) => { + task_wallets.remove(&task_id); + total += n; + } + Ok((task_id, (wallet_id, Err(e)))) => { + task_wallets.remove(&task_id); tracing::warn!(wallet = %wallet_id, error = ?e, "wallet poll failed") } - Err(e) => tracing::warn!(error = ?e, "wallet poll task panicked"), + Err(e) => match task_wallets.remove(&e.id()) { + Some(wallet_id) => { + tracing::error!(wallet = %wallet_id, error = ?e, "wallet poll task panicked; it will be retried on a later tick") + } + None => tracing::error!(error = ?e, "wallet poll task panicked; wallet id unavailable"), + }, } } Ok(total) diff --git a/crates/store/src/lib.rs b/crates/store/src/lib.rs index 97f981a..cf3173f 100644 --- a/crates/store/src/lib.rs +++ b/crates/store/src/lib.rs @@ -1283,7 +1283,8 @@ impl Store { .ok_or(StoreError::NotFound) } - /// Public lookup by slug — no wallet scoping, this is the pay-page entry point. + /// Public lookup by slug — the UNIQUE constraint supplies the index for this equality lookup. + /// No wallet scoping; this is the pay-page entry point. pub async fn get_payment_link_by_slug(&self, slug: &str) -> Result { sqlx::query_as::<_, PaymentLink>("SELECT * FROM payment_links WHERE slug = $1") .bind(slug) @@ -1855,6 +1856,17 @@ impl Store { .ok_or(StoreError::NotFound) } + /// Check whether a webhook endpoint is still active using its indexed id. + pub async fn is_webhook_endpoint_active(&self, id: Uuid) -> Result { + sqlx::query_scalar::<_, bool>( + "SELECT EXISTS (SELECT 1 FROM webhook_endpoints WHERE id = $1 AND active = true)", + ) + .bind(id) + .fetch_one(&self.pool) + .await + .map_err(StoreError::from) + } + /// An endpoint's delivery history, newest first, capped at `limit` rows. pub async fn list_webhook_deliveries( &self, diff --git a/crates/webhooks/src/lib.rs b/crates/webhooks/src/lib.rs index 40fc3c2..68c1f08 100644 --- a/crates/webhooks/src/lib.rs +++ b/crates/webhooks/src/lib.rs @@ -105,6 +105,17 @@ impl WebhookSender { let res = tokio::time::timeout(self.delivery_timeout, async { for attempt in 1..=self.max_attempts { + match self.store.is_webhook_endpoint_active(ep.id).await { + Ok(true) => {} + Ok(false) => { + tracing::info!(endpoint_id = %ep.id, "stopping webhook retries for inactive endpoint"); + return Ok(None); + } + Err(error) => { + tracing::warn!(endpoint_id = %ep.id, error = ?error, "stopping webhook retries because endpoint status could not be checked"); + return Err(()); + } + } attempts_made = attempt; let resp = self .http @@ -120,7 +131,7 @@ impl WebhookSender { let code = r.status().as_u16() as i32; last_code = Some(code); if r.status().is_success() { - return Ok(attempt); + return Ok(Some(attempt)); } } Err(_) => last_code = None, @@ -136,7 +147,7 @@ impl WebhookSender { .await; match res { - Ok(Ok(successful_attempt)) => { + Ok(Ok(Some(successful_attempt))) => { let _ = self .store .log_webhook_delivery( @@ -150,6 +161,20 @@ impl WebhookSender { .await; true } + Ok(Ok(None)) => { + let _ = self + .store + .log_webhook_delivery( + ep.id, + event_type, + body, + "failed", + attempts_made as i32, + last_code, + ) + .await; + false + } Ok(Err(())) => { let _ = self .store @@ -158,7 +183,7 @@ impl WebhookSender { event_type, body, "failed", - self.max_attempts as i32, + attempts_made as i32, last_code, ) .await; diff --git a/docs/api.md b/docs/api.md index b18cace..885415e 100644 --- a/docs/api.md +++ b/docs/api.md @@ -77,6 +77,8 @@ carries fee float only — the one server-held key in the system, bounded by you gets `401`). Idempotent: a second call returns the existing tank. - `GET /v1/wallets/{id}/sponsorship` / `PUT` — read/update `enabled`, the per-transaction fee cap, and the daily budget. +- Sponsorship and whitelist config `PUT` requests require at least one field; an empty body or + `{}` returns `400 Bad Request`. - `POST /v1/wallets/{id}/sponsor` — fee-bump a user's **already-signed** inner transaction. The gas tank signs only the outer fee-bump envelope; the inner transaction is passed through untouched. Over budget → `429`; duplicate inner tx → `409`.