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
10 changes: 9 additions & 1 deletion crates/api/src/routes/sponsorship.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<AppState>,
Path(wallet_id): Path<Uuid>,
Expand All @@ -71,6 +71,14 @@ pub async fn put_config(
) -> ApiResult<Json<Envelope<SponsorshipConfigView>>> {
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 {
Expand Down
7 changes: 6 additions & 1 deletion crates/api/src/routes/whitelist.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<AppState>,
Path(wallet_id): Path<Uuid>,
Expand All @@ -66,6 +66,11 @@ pub async fn put_config(
) -> ApiResult<Json<Envelope<AllowlistConfigView>>> {
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
Expand Down
22 changes: 21 additions & 1 deletion crates/ingest/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -693,8 +693,10 @@ impl Supervisor {
// of the shared concurrency budget. Paged by id so memory doesn't scale with wallet count.
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();
Expand All @@ -703,7 +705,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;
Expand Down Expand Up @@ -741,6 +743,7 @@ impl Supervisor {
}
(w.id, result)
});
task_wallets.insert(task_id.id(), wallet_id);
}

let mut total = 0;
Expand Down Expand Up @@ -794,6 +797,23 @@ impl Supervisor {
}
}

/// Queue one wallet's poll on `tasks`, gated by the shared concurrency `semaphore`.
fn spawn_poll(
&self,
tasks: &mut tokio::task::JoinSet<(Uuid, Result<usize, IngestError>)>,
semaphore: &Arc<tokio::sync::Semaphore>,
}
}

while let Some(joined) = tasks.join_next().await {
total += Self::tally(joined);
}
match fetch_error {
Some(e) => Err(e.into()),
None => Ok(total),
}
}

/// Queue one wallet's poll on `tasks`, gated by the shared concurrency `semaphore`.
fn spawn_poll(
&self,
Expand Down
Loading
Loading