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
37 changes: 37 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions crates/discovery-core/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ starknet-providers = "0.16"
tracing = "0.1"
serde = { version = "1.0", features = ["derive"] }
url = "2"
futures = "0.3"
zeroize = "1"

[dev-dependencies]
Expand Down
116 changes: 116 additions & 0 deletions crates/discovery-core/src/discovery/cursor.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
//! Cursor types for paginated discovery.
//!
//! These cursors track progress across paginated discovery calls, allowing
//! callers to resume discovery from where they left off.

use std::collections::HashMap;

use serde::{Deserialize, Serialize};
use starknet_types_core::felt::Felt;

/// Top-level cursor for channel discovery (shared by incoming and outgoing).
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct DiscoveryCursor {
/// All channels have been enumerated. Set by the discovery service once
/// the sentinel channel is reached. When `true`, no further channel
/// discovery is attempted — only channels already in the cursor are
/// processed.
#[serde(default)]
pub channel_discovery_complete: bool,

/// Total number of channels (cached from `get_num_of_channels` for incoming).
/// Used as optimization to avoid redundant RPC calls.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub total_n_channels: Option<u64>,

/// Last fully processed channel index. `None` = start from index 0.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_channel_index: Option<u64>,

/// Channels with pending subchannel/note discovery.
/// - Incoming: keyed by sender address.
/// - Outgoing: keyed by recipient address.
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub channels: HashMap<Felt, ChannelCursor>,
}

impl DiscoveryCursor {
/// Returns `true` when all discovery levels are complete: channels,
/// subchannels within each channel, and notes within each subchannel.
pub fn is_complete(&self) -> bool {
self.channel_discovery_complete && self.all_channels_processed()
}

/// Returns `true` when every channel currently in the cursor has
/// completed subchannel and note discovery. Also returns `true` when
/// the cursor has no channels (vacuously).
///
/// Used by sync orchestrators to decide whether to discover new
/// channels vs. process pending subchannel/note work.
pub fn all_channels_processed(&self) -> bool {
self.channels.values().all(ChannelCursor::is_complete)
}
}

/// Cursor state for a single channel (shared by incoming and outgoing).
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ChannelCursor {
// TODO: Consider encrypting/masking channel_key in the serialized cursor
// to avoid exposing it in plaintext (sensitive value).
/// The channel key for this channel.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub channel_key: Option<Felt>,

/// All subchannels have been enumerated. Set by the discovery service
/// once the sentinel subchannel is reached. When `true`, no further
/// subchannel discovery is attempted — only subchannels already in the
/// cursor are processed.
#[serde(default)]
pub subchannel_discovery_complete: bool,

/// Last fully processed subchannel index. `None` = start from index 0.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_subchannel_index: Option<u64>,

/// Subchannels with pending note discovery, keyed by token address.
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub subchannels: HashMap<Felt, SubchannelCursor>,
}

impl ChannelCursor {
/// Returns `true` when subchannel discovery is complete and all
/// subchannels have finished note discovery.
pub fn is_complete(&self) -> bool {
self.subchannel_discovery_complete
&& self
.subchannels
.values()
.all(|sc| sc.note_discovery_complete)
}
}

/// Cursor state for a single subchannel (shared by incoming and outgoing).
///
/// For incoming (linear scan): only `last_note_index` is used.
/// For outgoing (exponential search): `last_note_index` = lo, `max_note_index` = hi.
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct SubchannelCursor {
/// All notes in this subchannel have been discovered. Set when the
/// note discovery scan completes without budget exhaustion.
#[serde(default)]
pub note_discovery_complete: bool,

/// Last note index where a note exists.
/// - Incoming: last scanned index.
/// - Outgoing: lower bound (lo) for exponential search.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_note_index: Option<u64>,

/// - Incoming: last index confirmed to exist by exponential probe. Linear
/// scan reads notes up to this index. Kept after scan — used to bound the
/// next exponential probe range (`max_note_index * 2`). Re-probe triggers
/// when `last_note_index == max_note_index`.
/// - Outgoing: first index confirmed empty (hi); `Some` = bisection phase.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_note_index: Option<u64>,
}
27 changes: 22 additions & 5 deletions crates/discovery-core/src/discovery/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,21 +5,32 @@ use thiserror::Error;
use crate::privacy_pool::decryption::DecryptionError;
use crate::storage_backend::StorageError;

pub mod cursor;
pub use cursor::{ChannelCursor, DiscoveryCursor, SubchannelCursor};
pub mod incoming_channels;
pub mod notes;
pub mod subchannels;

/// Cost for `get_num_of_channels` (1 storage slot read).
const COST_NUM_CHANNELS: usize = 1;
pub const COST_NUM_CHANNELS: usize = 1;

/// Cost for `get_channel_info` (3 storage slot reads).
const COST_CHANNEL_INFO: usize = 3;
pub const COST_CHANNEL_INFO: usize = 3;

/// Cost for `get_subchannel_info` (2 storage slot reads).
const COST_SUBCHANNEL_INFO: usize = 2;
pub const COST_SUBCHANNEL_INFO: usize = 2;

/// Cost for `get_note` (1 storage slot read).
const COST_NOTE: usize = 1;
/// Cost for `get_note` + `nullifier_exists` (2 storage slot reads).
pub const COST_NOTE: usize = 2;

/// Cost for outgoing channel info (3 storage slot reads: salt + enc_recipient_addr + public_key).
pub const COST_OUTGOING_CHANNEL_INFO: usize = 3;

/// Cost for a single note existence probe (1 `get_note` read, no nullifier check).
pub const COST_NOTE_PROBING: usize = 1;

/// Cost for a single `get_public_key` (1 storage slot read).
pub const COST_PUBLIC_KEY: usize = 1;

/// Errors that can occur during channel discovery.
#[derive(Debug, Error)]
Expand All @@ -36,4 +47,10 @@ pub enum DiscoveryError {
#[source]
source: DecryptionError,
},
/// A spawned task panicked.
#[error("spawned task panicked: {0}")]
TaskPanicked(String),
/// Invalid cursor data provided by client.
#[error("invalid cursor: {0}")]
InvalidCursor(String),
}
2 changes: 1 addition & 1 deletion crates/discovery-core/src/discovery/notes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -267,7 +267,7 @@ mod tests {
.await
.expect("Alice's channel should have at least one subchannel");

// Budget exhausted before starting (COST_NOTE = 1)
// Budget exhausted before starting (COST_NOTE = 2)
let budget = IoBudget::new(0);
let result = discover_notes(&backend, channel_key, token, 0, &budget)
.await
Expand Down
98 changes: 98 additions & 0 deletions crates/discovery-core/src/io_budget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,41 @@ impl IoBudget {
})
.is_ok()
}

/// Atomically consumes as many whole items as the budget allows.
///
/// Returns `(num_consumed_items, budget_exhausted)`:
/// - `num_consumed_items`: number of items consumed (0..=`max_items`).
/// - `budget_exhausted`: `true` when the budget limited how many items could
/// be consumed (i.e. `consumed < max_items`).
/// Callers use this to derive `has_more` for pagination cursors.
///
/// When `max_items == 0`, returns `(0, false)` — the request is trivially
/// satisfied, not budget-limited.
///
/// # Panics
///
/// Panics if `cost_per_item == 0`. All callers use compile-time cost constants;
/// a zero cost is always a programmer bug.
pub fn consume_up_to(&self, max_items: usize, cost_per_item: usize) -> (usize, bool) {
assert!(
cost_per_item > 0,
"cost_per_item must be positive; zero-cost items are not supported"
);
if max_items == 0 {
return (0, false);
}
let num_consumed_items = self
.remaining
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |current| {
let num_items = (current / cost_per_item).min(max_items);
current.checked_sub(num_items * cost_per_item)
})
.map(|old| (old / cost_per_item).min(max_items))
.unwrap_or(0);
let budget_exhausted = num_consumed_items < max_items;
(num_consumed_items, budget_exhausted)
}
}

#[cfg(test)]
Expand Down Expand Up @@ -90,6 +125,69 @@ mod tests {
assert_eq!(budget.remaining(), 0);
}

#[test]
fn test_consume_up_to_exact_budget() {
let budget = IoBudget::new(9);
// 9 / 3 = 3 items, but max is 2 → not budget-limited (got what we asked for)
assert_eq!(budget.consume_up_to(2, 3), (2, false));
assert_eq!(budget.remaining(), 3);
// 3 / 3 = 1 item, asked for 5 → budget-limited
assert_eq!(budget.consume_up_to(5, 3), (1, true));
assert_eq!(budget.remaining(), 0);
}

#[test]
fn test_consume_up_to_partial_budget() {
let budget = IoBudget::new(7);
// 7 / 3 = 2 items, asked for 10 → budget-limited
assert_eq!(budget.consume_up_to(10, 3), (2, true));
assert_eq!(budget.remaining(), 1);
// 1 / 3 = 0 items → budget-limited
assert_eq!(budget.consume_up_to(10, 3), (0, true));
assert_eq!(budget.remaining(), 1); // unchanged
}

#[test]
fn test_consume_up_to_zero_budget() {
let budget = IoBudget::new(0);
// Asked for 5 but budget is 0 → budget-limited
assert_eq!(budget.consume_up_to(5, 3), (0, true));
assert_eq!(budget.remaining(), 0);
}

#[test]
#[should_panic(expected = "cost_per_item must be positive")]
fn test_consume_up_to_zero_cost_panics() {
let budget = IoBudget::new(10);
budget.consume_up_to(5, 0);
}

#[test]
fn test_consume_up_to_zero_max_items() {
let budget = IoBudget::new(10);
// Asked for 0 → trivially satisfied, not budget-limited
assert_eq!(budget.consume_up_to(0, 3), (0, false));
assert_eq!(budget.remaining(), 10); // unchanged
}

#[test]
fn test_consume_up_to_concurrent() {
let budget = IoBudget::new(100);
let mut handles = vec![];

// 10 threads each trying to consume up to 5 items at cost 2 = 10 per thread
for _ in 0..10 {
let budget = budget.clone();
handles.push(thread::spawn(move || budget.consume_up_to(5, 2).0));
}

let total: usize = handles.into_iter().map(|h| h.join().unwrap()).sum();

// 100 / 2 = 50 items total possible
assert_eq!(total, 50);
assert_eq!(budget.remaining(), 0);
}

#[test]
fn test_concurrent_consume() {
let budget = IoBudget::new(1000);
Expand Down
Loading