Skip to content
Merged
10 changes: 9 additions & 1 deletion crates/sui-rpc/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,11 @@ sui-crypto = { version = "0.3.0", path = "../sui-crypto", default-features = fal

# dependencies for the protobuf and gRPC definitions
bytes = "1.10"
# Not used directly; this floors hyper at 1.10, whose h2 requirement
# excludes the h2 versions with stream-cancellation flow-control
# accounting bugs (hyperium/h2#893, #896, #897, #898, and #913) that can
# permanently wedge a multiplexed gRPC connection.
hyper = { version = "1.10", default-features = false }
tonic = { version = "0.14.2", default-features = false, features = ["channel", "codegen", "tls-ring", "tls-webpki-roots", "zstd"] }
tonic-prost = "0.14.2"
prost = "0.14.1"
Expand All @@ -57,7 +62,10 @@ proptest = { version = "1.8.0", default-features = false, features = ["std"] }
test-strategy = { version = "0.4" }
sui-sdk-types = { version = "0.3.0", path = "../sui-sdk-types", default-features = false, features = ["proptest", "serde", "hash"] }
serde_json = { version = "1.0.145" }
tokio = { version = "1.40", features = ["rt"] }
tokio = { version = "1.40", features = ["rt", "rt-multi-thread", "macros", "net", "test-util"] }
# The starvation integration tests host a mock gRPC server built from this
# crate's own generated server stubs.
tonic = { version = "0.14.2", default-features = false, features = ["server", "router"] }

[lints.rust]
unexpected_cfgs = { level = "warn", check-cfg = ['cfg(doc_cfg)'] }
96 changes: 64 additions & 32 deletions crates/sui-rpc/src/client/lists.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,12 +37,17 @@ impl Client {
request, // request (page_token will be updated as we paginate)
client, // client for making requests
),
move |(mut iter, has_next_page, mut request, mut client)| async move {
move |(mut iter, mut has_next_page, mut request, mut client)| async move {
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}

if has_next_page {
// A page may be empty while still carrying a next_page_token
// (for example, a filtered scan that hit a server-side
// budget), so keep following the token until a page yields an
// item or pagination ends. Each response's token reflects
// server-side scan progress, so this terminates.
while has_next_page {
let new_request = tonic::Request::from_parts(
request.metadata().clone(),
request.extensions().clone(),
Expand All @@ -54,21 +59,24 @@ impl Client {
let response = response.into_inner();
let mut iter = response.objects.into_iter();

let has_next_page = response.next_page_token.is_some();
has_next_page = response.next_page_token.is_some();
request.get_mut().page_token = response.next_page_token;

iter.next()
.map(|item| (Ok(item), (iter, has_next_page, request, client)))
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}
}
Err(e) => {
// Return error and terminate stream
request.get_mut().page_token = None;
Some((Err(e), (Vec::new().into_iter(), false, request, client)))
return Some((
Err(e),
(Vec::new().into_iter(), false, request, client),
));
}
}
} else {
None
}
None
},
)
}
Expand Down Expand Up @@ -98,12 +106,17 @@ impl Client {
request, // request (page_token will be updated as we paginate)
client, // client for making requests
),
move |(mut iter, has_next_page, mut request, mut client)| async move {
move |(mut iter, mut has_next_page, mut request, mut client)| async move {
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}

if has_next_page {
// A page may be empty while still carrying a next_page_token
// (for example, a filtered scan that hit a server-side
// budget), so keep following the token until a page yields an
// item or pagination ends. Each response's token reflects
// server-side scan progress, so this terminates.
while has_next_page {
let new_request = tonic::Request::from_parts(
request.metadata().clone(),
request.extensions().clone(),
Expand All @@ -115,21 +128,24 @@ impl Client {
let response = response.into_inner();
let mut iter = response.dynamic_fields.into_iter();

let has_next_page = response.next_page_token.is_some();
has_next_page = response.next_page_token.is_some();
request.get_mut().page_token = response.next_page_token;

iter.next()
.map(|item| (Ok(item), (iter, has_next_page, request, client)))
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}
}
Err(e) => {
// Return error and terminate stream
request.get_mut().page_token = None;
Some((Err(e), (Vec::new().into_iter(), false, request, client)))
return Some((
Err(e),
(Vec::new().into_iter(), false, request, client),
));
}
}
} else {
None
}
None
},
)
}
Expand Down Expand Up @@ -159,12 +175,17 @@ impl Client {
request, // request (page_token will be updated as we paginate)
client, // client for making requests
),
move |(mut iter, has_next_page, mut request, mut client)| async move {
move |(mut iter, mut has_next_page, mut request, mut client)| async move {
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}

if has_next_page {
// A page may be empty while still carrying a next_page_token
// (for example, a filtered scan that hit a server-side
// budget), so keep following the token until a page yields an
// item or pagination ends. Each response's token reflects
// server-side scan progress, so this terminates.
while has_next_page {
let new_request = tonic::Request::from_parts(
request.metadata().clone(),
request.extensions().clone(),
Expand All @@ -176,21 +197,24 @@ impl Client {
let response = response.into_inner();
let mut iter = response.balances.into_iter();

let has_next_page = response.next_page_token.is_some();
has_next_page = response.next_page_token.is_some();
request.get_mut().page_token = response.next_page_token;

iter.next()
.map(|item| (Ok(item), (iter, has_next_page, request, client)))
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}
}
Err(e) => {
// Return error and terminate stream
request.get_mut().page_token = None;
Some((Err(e), (Vec::new().into_iter(), false, request, client)))
return Some((
Err(e),
(Vec::new().into_iter(), false, request, client),
));
}
}
} else {
None
}
None
},
)
}
Expand Down Expand Up @@ -220,12 +244,17 @@ impl Client {
request, // request (page_token will be updated as we paginate)
client, // client for making requests
),
move |(mut iter, has_next_page, mut request, mut client)| async move {
move |(mut iter, mut has_next_page, mut request, mut client)| async move {
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}

if has_next_page {
// A page may be empty while still carrying a next_page_token
// (for example, a filtered scan that hit a server-side
// budget), so keep following the token until a page yields an
// item or pagination ends. Each response's token reflects
// server-side scan progress, so this terminates.
while has_next_page {
let new_request = tonic::Request::from_parts(
request.metadata().clone(),
request.extensions().clone(),
Expand All @@ -241,21 +270,24 @@ impl Client {
let response = response.into_inner();
let mut iter = response.versions.into_iter();

let has_next_page = response.next_page_token.is_some();
has_next_page = response.next_page_token.is_some();
request.get_mut().page_token = response.next_page_token;

iter.next()
.map(|item| (Ok(item), (iter, has_next_page, request, client)))
if let Some(item) = iter.next() {
return Some((Ok(item), (iter, has_next_page, request, client)));
}
}
Err(e) => {
// Return error and terminate stream
request.get_mut().page_token = None;
Some((Err(e), (Vec::new().into_iter(), false, request, client)))
return Some((
Err(e),
(Vec::new().into_iter(), false, request, client),
));
}
}
} else {
None
}
None
},
)
}
Expand Down
Loading