Skip to content
Draft
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
162 changes: 140 additions & 22 deletions crates/discovery-core/src/history/transactions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,12 @@ use crate::io_budget::IoBudget;
use crate::privacy_pool::events::{IEvents, PrivacyPoolEventContent};
use crate::privacy_pool::views::IViews;

/// Blocks per `get_withdrawal_events` call within a gap window. Keeps each call
/// short enough that the budget's deadline is checked often: an RPC node pages
/// a key-filtered event query over a wide range, so one call over the whole
/// window can run far past the deadline.
const GAP_SCAN_SUB_WINDOW_BLOCKS: u64 = 64 * EVENTS_COST_CHUNK_SIZE as u64;

/// Outcome of one [`process_next_block`] iteration.
enum ScanStep {
/// Committed a gap window and its anchoring note block; keep scanning.
Expand All @@ -42,8 +48,9 @@ enum ScanStep {
/// via [`fetch_aggregated_block_events`].
///
/// Budget is consumed incrementally; the scan stops at the first iteration that
/// cannot make further progress within budget, with the cursor left at the
/// frontier so the next page continues. A page that makes no progress at all
/// cannot make further progress within budget, or at the first resumable point
/// after the budget's deadline, with the cursor left at the frontier so the
/// next page continues. A page that makes no progress at all
/// (no transaction and no cursor movement) returns `InsufficientBudget`
/// instead, since a same-budget retry would repeat it. Returns transactions
/// sorted by `block_number` descending.
Expand Down Expand Up @@ -86,6 +93,9 @@ pub async fn fetch_transactions<B: IViews + IEvents>(
)
.await
{
// The deadline is checked only after a committed step, so a page
// that reaches it has always moved the cursor.
Ok(ScanStep::Advanced) if budget.deadline_reached() => break,
Ok(ScanStep::Advanced) => continue,
Ok(ScanStep::Halted) => break,
Ok(ScanStep::Exhausted) => {
Expand Down Expand Up @@ -433,7 +443,8 @@ async fn fetch_aggregated_block_events<E: IEvents>(
/// Returns the lowest block scanned (`window_bottom`); the caller advances the
/// cursor to `window_bottom - 1`. When the full gap was covered in this call,
/// `window_bottom == from_block`; otherwise it is above `from_block` and the
/// caller resumes the remainder on the next page.
/// caller resumes the remainder on the next page. The scan also stops early, at
/// a sub-window boundary, once the budget's deadline is reached.
///
/// `to_block_id` is the `BlockId` used for the window top, so a fresh scan can
/// pass the snapshot's pinned tag (e.g. `PreConfirmed`) to include
Expand Down Expand Up @@ -472,29 +483,52 @@ async fn fetch_gap_withdrawals_chunked<E: IEvents>(
let granted_span = (chunks_granted as u64).saturating_mul(EVENTS_COST_CHUNK_SIZE as u64);
let window_bottom = to_block.saturating_sub(granted_span - 1).max(from_block);

let events = backend
.get_withdrawal_events(user_address, BlockId::Number(window_bottom), to_block_id)
.await?;
// Walk the window top-down in sub-windows so the budget's deadline can end
// the scan at a sub-window boundary. The first sub-window always runs, so a
// call reaching the deadline has still moved the frontier down.
let mut sub_window_top = to_block;
let mut sub_window_top_id = to_block_id;
let scanned_bottom = loop {
let sub_window_bottom = sub_window_top
.saturating_sub(GAP_SCAN_SUB_WINDOW_BLOCKS - 1)
.max(window_bottom);
let events = backend
.get_withdrawal_events(
user_address,
BlockId::Number(sub_window_bottom),
sub_window_top_id,
)
.await?;

trace!(
from_block,
to_block,
window_bottom,
chunks_granted,
num_withdrawal_events = events.len(),
"withdrawal_events: chunked window fetched"
);
trace!(
from_block,
to_block,
sub_window_bottom,
sub_window_top,
chunks_granted,
num_withdrawal_events = events.len(),
"withdrawal_events: gap sub-window fetched"
);

for event in events {
let tx = transactions
.entry(event.transaction_hash)
.or_insert_with(|| {
HistoryTransaction::new(event.block_number, event.transaction_hash)
});
if let PrivacyPoolEventContent::Withdrawal(withdrawal) = event.content {
tx.withdrawals.push(withdrawal);
}
}

for event in events {
let tx = transactions
.entry(event.transaction_hash)
.or_insert_with(|| HistoryTransaction::new(event.block_number, event.transaction_hash));
if let PrivacyPoolEventContent::Withdrawal(withdrawal) = event.content {
tx.withdrawals.push(withdrawal);
if sub_window_bottom == window_bottom || budget.deadline_reached() {
break sub_window_bottom;
}
}
sub_window_top = sub_window_bottom.saturating_sub(1);
sub_window_top_id = BlockId::Number(sub_window_top);
};

Ok(window_bottom)
Ok(scanned_bottom)
}

/// Fetches the synthetic registration transaction for the user.
Expand Down Expand Up @@ -552,6 +586,7 @@ mod tests {
use crate::privacy_pool::types::SecretFelt;
use crate::storage_backend::{MockBackend, RawStorageAccess, StorageError};
use starknet_core::types::StorageResult;
use std::time::Instant;

const ADDRESS: Felt = Felt::from_hex_unchecked("0xABCD");
const OTHER_ADDRESS: Felt = Felt::from_hex_unchecked("0x9999");
Expand Down Expand Up @@ -1336,6 +1371,89 @@ mod tests {
assert!(cursor.history_complete);
}

#[tokio::test]
async fn expired_deadline_ends_gap_scan_after_one_sub_window() {
const TX_HASH_FAR_WITHDRAWAL: Felt = Felt::from_hex_unchecked("0x1009");
const GAP_TOP: u64 = 1_100_000;

let (backend, mut cursor) = FixtureBuilder::new()
.note(0, 100, TX_HASH_1)
.withdrawal(50, 150, TX_HASH_FAR_WITHDRAWAL)
.build(Some(0));
cursor.begin_block_number = Some(GAP_TOP);
let expired_budget = || IoBudget::new(10_000).with_deadline(Instant::now());

let page1 = fetch_transactions(&backend, ADDRESS, &mut cursor, 10, &expired_budget())
.await
.unwrap();
assert!(page1.is_empty());
assert!(!cursor.history_complete);
assert_eq!(
cursor.begin_block_number,
Some(GAP_TOP - GAP_SCAN_SUB_WINDOW_BLOCKS),
"an expired deadline scans exactly one sub-window"
);

// Every page still moves the cursor, so the walk reaches the note.
let max_pages = GAP_TOP / GAP_SCAN_SUB_WINDOW_BLOCKS + 3;
let mut all_transactions = page1;
for _ in 0..max_pages {
if cursor.history_complete {
break;
}
let previous_bound = cursor.begin_block_number;
let page = fetch_transactions(&backend, ADDRESS, &mut cursor, 10, &expired_budget())
.await
.unwrap();
assert!(
!page.is_empty()
|| cursor.history_complete
|| cursor.begin_block_number < previous_bound,
"page made no progress"
);
all_transactions.extend(page);
}
assert!(cursor.history_complete);
let blocks: Vec<u64> = all_transactions.iter().map(|tx| tx.block_number).collect();
assert_eq!(blocks, vec![150, 100]);
}

#[tokio::test]
async fn expired_deadline_ends_page_after_note_block() {
let fixture = || {
FixtureBuilder::new()
.note(0, 10, TX_HASH_1)
.note(1, 20, TX_HASH_2)
.deposit(100, 10, TX_HASH_1)
.withdrawal(50, 20, TX_HASH_2)
.build(Some(1))
};

// Control: without a deadline both note blocks fit in one page.
let (backend, mut cursor) = fixture();
let unbounded_page =
fetch_transactions(&backend, ADDRESS, &mut cursor, 10, &IoBudget::new(1000))
.await
.unwrap();
assert_eq!(unbounded_page.len(), 2);

let (backend, mut cursor) = fixture();
let expired_budget = || IoBudget::new(1000).with_deadline(Instant::now());
let page1 = fetch_transactions(&backend, ADDRESS, &mut cursor, 10, &expired_budget())
.await
.unwrap();
assert_eq!(page1.len(), 1);
assert_eq!(page1[0].block_number, 20);
assert_eq!(cursor.begin_block_number, Some(19));
assert!(!cursor.history_complete);

let page2 = fetch_transactions(&backend, ADDRESS, &mut cursor, 10, &expired_budget())
.await
.unwrap();
assert_eq!(page2.len(), 1);
assert_eq!(page2[0].block_number, 10);
}

#[tokio::test]
async fn note_block_withdrawal_not_double_counted_across_pages() {
// A partial withdrawal emits a note and a Withdrawal in the same block.
Expand Down
37 changes: 36 additions & 1 deletion crates/discovery-core/src/io_budget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Instant;

/// Error returned when a budget consumption fails.
#[derive(Debug)]
Expand All @@ -21,19 +22,40 @@ pub struct InsufficientBudgetError {
///
/// `IoBudget` is cheap to clone - clones share the same underlying counter,
/// making it easy to pass across async tasks.
///
/// A budget may also carry a wall-clock `deadline`. The counter bounds how much
/// I/O a request may issue, but not how long that I/O takes, and the RPC node's
/// latency is outside the server's control. Scans that can stop at a resumable
/// point consult [`IoBudget::deadline_reached`] to end the page early.
#[derive(Debug, Clone)]
pub struct IoBudget {
remaining: Arc<AtomicUsize>,
deadline: Option<Instant>,
}

impl IoBudget {
/// Creates a new budget with the given limit.
/// Creates a new budget with the given limit and no deadline.
pub fn new(limit: usize) -> Self {
Self {
remaining: Arc::new(AtomicUsize::new(limit)),
deadline: None,
}
}

/// Returns this budget with a wall-clock `deadline`.
pub fn with_deadline(self, deadline: Instant) -> Self {
Self {
deadline: Some(deadline),
..self
}
}

/// Returns `true` once the deadline has passed. Always `false` without a deadline.
pub fn deadline_reached(&self) -> bool {
self.deadline
.is_some_and(|deadline| Instant::now() >= deadline)
}

/// Returns the current remaining budget.
pub fn remaining(&self) -> usize {
self.remaining.load(Ordering::Relaxed)
Expand Down Expand Up @@ -113,6 +135,19 @@ impl IoBudget {
mod tests {
use super::*;
use std::thread;
use std::time::Duration;

#[test]
fn test_deadline_reached() {
assert!(!IoBudget::new(10).deadline_reached());
assert!(IoBudget::new(10)
.with_deadline(Instant::now())
.deadline_reached());
let future_deadline = Instant::now() + Duration::from_secs(3600);
assert!(!IoBudget::new(10)
.with_deadline(future_deadline)
.deadline_reached());
}

#[test]
fn test_consume_success() {
Expand Down
2 changes: 2 additions & 0 deletions crates/discovery-service/specs/06-api-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -354,12 +354,14 @@ Retrieves paginated transaction history by scanning backward through note subcha
- A page may therefore return **few or zero transactions** while still advancing `begin_block_number` through a long stretch with no activity. Clients must keep paginating until `history_complete` is `true` rather than treating an empty page as the end.
- The gap window is always strictly **above** the note block it anchors, so a withdrawal in a note's own block is attributed once (via that block's events) and never re-scanned by a later page's gap.
- One gap window can attach every withdrawal it contains at once, so a page covering a withdrawal-dense range may return **more than `max_transactions`** transactions. `max_transactions` is a per-page target, not a hard cap on a single page's size.
- A page also ends once `history_time_limit` has elapsed, since the budget bounds how many RPC calls a page makes but not how long the node takes to answer them. The gap window is fetched in fixed-size sub-windows and the limit is checked after each sub-window and after each note block, so a page stops only after moving the cursor and may overrun the limit by one call. This keeps each response well inside a proxy or load balancer timeout, at the cost of more pages over a long idle gap.
- A page that can make **no** forward progress (the budget covers neither a gap chunk nor the next note step) returns `500 INTERNAL_ERROR` rather than an endless stream of empty `200`s. This only occurs when `server_budget` is set too low for the account's per-request cost — many subchannels, or many notes sharing one block, since draining a block refills once per note — and repeats until `server_budget` is raised.

**Validation limits:**

- `max_history_subchannels` (default: 256): Maximum number of subchannels in a history cursor.
- `max_history_transactions` (default: 100): Maximum allowed `max_transactions` value.
- `history_time_limit` (default: 10 seconds): Wall-clock time after which a page stops at the next resumable point (see Chunked gap scan).
- `max_transactions` must be at least `1`. A zero-size page can neither advance the scan nor mark it complete.

**Error responses:**
Expand Down
2 changes: 2 additions & 0 deletions crates/discovery-service/specs/18-configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ max_cursor_channels = 256
max_cursor_subchannels_per_channel = 64
max_outgoing_recipients = 64
server_budget = 10000
history_time_limit = 10
max_request_body_bytes = 102400
```

Expand All @@ -79,6 +80,7 @@ These env vars override the corresponding config file values at runtime:
| `RUST_LOG` | `logging.level` | `info` |
| `LOG_FORMAT` | `logging.format` | `text` (accepts `text` or `json`, case-insensitive) |
| `SERVER_BUDGET` | `limits.server_budget` | `10000` |
| `HISTORY_TIME_LIMIT_SECS` | `limits.history_time_limit` | `10` |
| `TLS_CERT_PATH` | `api.tls.cert_path` | — (TLS disabled) |
| `TLS_KEY_PATH` | `api.tls.key_path` | — (TLS disabled) |

Expand Down
5 changes: 3 additions & 2 deletions crates/discovery-service/src/api/handlers.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
//! API route handlers.

use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use std::time::{Instant, SystemTime, UNIX_EPOCH};

use axum::extract::State;
use axum::http::StatusCode;
Expand Down Expand Up @@ -354,7 +354,8 @@ where
"history request"
);

let budget = IoBudget::new(state.validation_limits.server_budget);
let budget = IoBudget::new(state.validation_limits.server_budget)
.with_deadline(Instant::now() + state.validation_limits.history_time_limit);
let transactions = discovery_core::history::transactions::fetch_transactions(
&snapshot,
request.user_address,
Expand Down
12 changes: 12 additions & 0 deletions crates/discovery-service/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,12 @@ pub struct ValidationLimits {
pub max_history_transactions: usize,
/// Server-controlled I/O budget per request.
pub server_budget: usize,
/// Wall-clock time after which a history page stops at the next resumable
/// point and returns its cursor. Keep it well below any proxy or load
/// balancer timeout in front of the service; a page may overrun it by one
/// RPC call.
#[serde(deserialize_with = "deserialize_secs")]
pub history_time_limit: Duration,
/// Maximum request body size in bytes.
pub max_request_body_bytes: usize,
/// Maximum number of entries in the public key cache.
Expand All @@ -173,6 +179,7 @@ impl Default for ValidationLimits {
max_history_subchannels: 256,
max_history_transactions: 100,
server_budget: 10_000,
history_time_limit: Duration::from_secs(10),
max_request_body_bytes: 102_400,
public_key_cache_capacity: 10_000,
}
Expand Down Expand Up @@ -268,6 +275,11 @@ impl ServiceConfig {
self.limits.server_budget = n;
}
}
if let Ok(v) = std::env::var("HISTORY_TIME_LIMIT_SECS") {
if let Ok(secs) = v.parse() {
self.limits.history_time_limit = Duration::from_secs(secs);
}
}

if let Ok(v) = std::env::var("OHTTP_ENABLED") {
self.ohttp.enabled = v == "true" || v == "1";
Expand Down
1 change: 1 addition & 0 deletions deploy/discovery-service/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ Precedence: env var > config file > code default.
| `API_HOST` | `api.host` | `0.0.0.0:8080` (Docker) / `127.0.0.1:8080` (native) |
| `RUST_LOG` | `logging.level` | `info` |
| `SERVER_BUDGET` | `limits.server_budget` | `10000` |
| `HISTORY_TIME_LIMIT_SECS` | `limits.history_time_limit` | `10` |
| `TLS_CERT_PATH` | `api.tls.cert_path` | — (TLS disabled) |
| `TLS_KEY_PATH` | `api.tls.key_path` | — (TLS disabled) |
| `OHTTP_ENABLED` | `ohttp.enabled` | `false` |
Expand Down
4 changes: 3 additions & 1 deletion deploy/discovery-service/config.example.toml
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
# Discovery Service configuration
# All sections and fields are optional — sensible defaults are used when omitted.
# Env vars override matching fields: RPC_URL, WS_URL, API_HOST, RUST_LOG, SERVER_BUDGET,
# TLS_CERT_PATH, TLS_KEY_PATH.
# HISTORY_TIME_LIMIT_SECS, TLS_CERT_PATH, TLS_KEY_PATH.
# Supports env var expansion: ${VAR} (required) or ${VAR:-default} (with fallback).
#
# Full spec: crates/discovery-service/specs/18-configuration.md
Expand Down Expand Up @@ -44,5 +44,7 @@ max_channels = 256
max_subchannels_per_channel = 64
max_outgoing_recipients = 64
server_budget = 10000
# Seconds after which a history page stops and returns its cursor.
history_time_limit = 10
max_request_body_bytes = 102400
public_key_cache_capacity = 10000
Loading