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
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,6 @@ use crate::task::streaming::StreamedCaches;
use super::scan_then_live::Subscribable;

impl Subscribable for Event {
const SUBSCRIPTION_TYPE: &'static str = "events";

type Item = Self;
type Cursor = CEvent;
type Filter = EventFilter;
Expand Down
23 changes: 22 additions & 1 deletion crates/sui-indexer-alt-graphql/src/api/subscription/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,12 @@ use crate::config::Limits;
use crate::config::SubscriptionConfig;
use crate::error::RpcError;
use crate::error::bad_user_input;
use crate::error::upcast;
use crate::scope::Scope;
use crate::task::streaming::StreamedCaches;
use crate::task::streaming::SubscriberLimit;
use crate::task::streaming::SubscriptionBroadcast;
use crate::task::streaming::SubscriptionLifecycleGuard;
use crate::task::watermark::Watermarks;

mod events;
Expand Down Expand Up @@ -79,13 +82,14 @@ impl Subscription {
(a, b) => a.or(b),
};
reject_if_start_too_far_ahead(start_from, broadcast, config)?;
let guard = admit_subscription(ctx, broadcast, "checkpoints")?;

let caches = caches.clone();
let resolver_limits = limits.package_resolver();

let stream = broadcast
.clone()
.subscribe(start_from, fetcher.clone(), config);
.subscribe(start_from, fetcher.clone(), config, guard);

Ok(stream.map(move |item| {
item.map(|processed| {
Expand Down Expand Up @@ -148,6 +152,7 @@ impl Subscription {
(a, b) => a.or(b),
};
reject_if_start_too_far_ahead(start_from, broadcast, config)?;
let guard = admit_subscription(ctx, broadcast, "transactions")?;

Ok(subscribe::<Transaction>(
reader.clone(),
Expand All @@ -159,6 +164,7 @@ impl Subscription {
after,
after_checkpoint,
config.clone(),
guard,
))
}

Expand Down Expand Up @@ -202,6 +208,7 @@ impl Subscription {
(a, b) => a.or(b),
};
reject_if_start_too_far_ahead(start_from, broadcast, config)?;
let guard = admit_subscription(ctx, broadcast, "events")?;

Ok(subscribe::<Event>(
reader.clone(),
Expand All @@ -213,10 +220,24 @@ impl Subscription {
after,
after_checkpoint,
config.clone(),
guard,
))
}
}

/// Admit a new subscription of `subscription_type` by claiming a concurrency slot, or return an
/// at-capacity error refusing it. The returned guard holds the slot and the per-subscriber metric
/// handles for the subscription's lifetime; it is moved into the subscription's stream driver.
fn admit_subscription(
ctx: &Context<'_>,
broadcast: &SubscriptionBroadcast,
subscription_type: &'static str,
) -> Result<SubscriptionLifecycleGuard, RpcError<Error>> {
let subscriber_limit: &SubscriberLimit = ctx.data()?;
SubscriptionLifecycleGuard::new(subscription_type, broadcast.metrics(), subscriber_limit)
.map_err(upcast)
}

/// Reject a start point sitting more than `max_ahead` checkpoints past the chain tip. There is
/// nothing to backfill ahead of the tip, so such a request would only wait for the chain to reach
/// it, and a far-future one would hold the connection open indefinitely.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,9 +57,6 @@ const BACKFILL_POLL_INTERVAL: Duration = Duration::from_millis(100);
/// GraphQL type (transactions, events) supplies its cursor, filter, and matching, and the free
/// functions here drive the shared machinery over it.
pub(super) trait Subscribable {
/// The subscription kind, used to label this feed's aggregate subscriber metrics.
const SUBSCRIPTION_TYPE: &'static str;

/// The GraphQL node delivered in each edge.
type Item: OutputType + Send + 'static;
/// The opaque cursor minted for resumption.
Expand Down Expand Up @@ -130,6 +127,7 @@ pub(super) fn subscribe<S: Subscribable>(
after: Option<S::Cursor>,
after_checkpoint: Option<u64>,
config: SubscriptionConfig,
guard: SubscriptionLifecycleGuard,
) -> impl Stream<Item = Result<Edge<String, S::Item, EmptyFields>, RpcError>> {
// Size the backfill scan page to the resolve concurrency. Scans are sequential (each needs the
// previous page's cursor), so feeding one window of `n` concurrent resolutions takes ceil(n /
Expand All @@ -142,7 +140,6 @@ pub(super) fn subscribe<S: Subscribable>(
let handoff_threshold = config.broadcast_buffer as u64 / 2;

stream! {
let guard = SubscriptionLifecycleGuard::new(S::SUBSCRIPTION_TYPE, broadcast.metrics());
let mut pending_receiver = None;
let mut handoff: Option<u64> = None;
let mut last_checkpoint: Option<u64> = None;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,6 @@ use crate::task::streaming::StreamedCaches;
use super::scan_then_live::Subscribable;

impl Subscribable for Transaction {
const SUBSCRIPTION_TYPE: &'static str = "transactions";

type Item = Self;
type Cursor = CTransaction;
type Filter = TransactionFilter;
Expand Down
8 changes: 8 additions & 0 deletions crates/sui-indexer-alt-graphql/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,10 @@ pub struct SubscriptionConfig {
/// is rejected instead. Raising this admits starts further past the tip, at the cost of
/// connections parked waiting longer.
pub max_start_checkpoints_ahead_of_tip: u64,

/// Maximum number of concurrent subscriptions the server admits at once. A subscription opened
/// while this many are already active is rejected so it can retry later.
pub max_subscribers: usize,
}

impl Default for SubscriptionConfig {
Expand All @@ -313,6 +317,7 @@ impl Default for SubscriptionConfig {
per_subscriber_max_output_nodes_per_second: 1_000_000,
// About a minute at the average checkpoint rate.
max_start_checkpoints_ahead_of_tip: 300,
max_subscribers: 1024,
}
}
}
Expand All @@ -328,6 +333,7 @@ pub struct SubscriptionLayer {
pub max_concurrent_resolutions: Option<usize>,
pub per_subscriber_max_output_nodes_per_second: Option<u32>,
pub max_start_checkpoints_ahead_of_tip: Option<u64>,
pub max_subscribers: Option<usize>,
}

impl SubscriptionLayer {
Expand Down Expand Up @@ -355,6 +361,7 @@ impl SubscriptionLayer {
max_start_checkpoints_ahead_of_tip: self
.max_start_checkpoints_ahead_of_tip
.unwrap_or(base.max_start_checkpoints_ahead_of_tip),
max_subscribers: self.max_subscribers.unwrap_or(base.max_subscribers),
}
}
}
Expand Down Expand Up @@ -680,6 +687,7 @@ impl From<SubscriptionConfig> for SubscriptionLayer {
value.per_subscriber_max_output_nodes_per_second,
),
max_start_checkpoints_ahead_of_tip: Some(value.max_start_checkpoints_ahead_of_tip),
max_subscribers: Some(value.max_subscribers),
}
}
}
Expand Down
7 changes: 6 additions & 1 deletion crates/sui-indexer-alt-graphql/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,8 @@ use task::streaming::StreamedObjectStore;
use task::streaming::StreamedTransactionStore;
use task::streaming::StreamingPackageStore;
#[cfg(feature = "staging")]
use task::streaming::SubscriberLimit;
#[cfg(feature = "staging")]
use task::streaming::SubscriptionBroadcast;
use task::streaming::SubscriptionReadiness;
use task::watermark::WatermarkTask;
Expand Down Expand Up @@ -506,6 +508,8 @@ pub async fn start_rpc(
readiness,
)) = streaming_setup
{
#[cfg(feature = "staging")]
let max_subscribers = config.subscription.max_subscribers;
rpc = rpc.data(caches).data(config.subscription);
let s_stream = stream_task.run();
let s_eviction = eviction_task.run();
Expand All @@ -525,7 +529,8 @@ pub async fn start_rpc(
));
rpc = rpc
.data(subscription_broadcast)
.data(subscription_watermarks_rx);
.data(subscription_watermarks_rx)
.data(SubscriberLimit::new(max_subscribers));
}
Some((s_stream, s_eviction))
} else {
Expand Down
8 changes: 8 additions & 0 deletions crates/sui-indexer-alt-graphql/src/metrics/subscription.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ pub struct SubscriptionMetrics {
// Metrics aggregated across all subscribers.
pub active_subscriptions: IntGaugeVec,
pub subscriptions_opened: IntCounterVec,
pub subscriptions_rejected: IntCounterVec,
pub subscription_terminations: IntCounterVec,
pub subscription_duration: HistogramVec,
pub payloads_delivered: IntCounterVec,
Expand Down Expand Up @@ -138,6 +139,13 @@ impl SubscriptionMetrics {
registry,
)
.unwrap(),
subscriptions_rejected: register_int_counter_vec_with_registry!(
"graphql_subscription_rejected",
"Total subscriptions refused before opening, by type and reason (e.g. the server was at its concurrent-subscription capacity)",
&["type", "reason"],
registry,
)
.unwrap(),
subscription_terminations: register_int_counter_vec_with_registry!(
"graphql_subscription_terminations",
"Total subscriptions that have ended, by type and reason",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,8 @@ use super::checkpoint_resume::scan_checkpoints;
#[cfg(any(feature = "staging", test))]
use super::gap_recovery::CheckpointFetcher;
use super::gap_recovery::recover_gap;
#[cfg(test)]
use super::lifecycle::SubscriberLimit;
#[cfg(any(feature = "staging", test))]
use super::lifecycle::SubscriptionLifecycleGuard;
#[cfg(any(feature = "staging", test))]
Expand Down Expand Up @@ -242,6 +244,7 @@ impl SubscriptionBroadcast {
resume_from: Option<u64>,
fetcher: F,
config: &SubscriptionConfig,
guard: SubscriptionLifecycleGuard,
) -> impl Stream<Item = Result<Arc<ProcessedCheckpoint>, RpcError>> + 'static {
// Resubscribe and pin the handoff once the scan is within this many checkpoints of the tip.
// Half the buffer leaves room for checkpoints arriving during the handoff, so it won't lag.
Expand All @@ -250,7 +253,6 @@ impl SubscriptionBroadcast {
let config = config.clone();

stream! {
let guard = SubscriptionLifecycleGuard::new("checkpoints", &self.metrics);
let mut last_yielded: Option<u64> = resume_from;
let mut receiver = self.broadcaster.resubscribe();
// `resubscribe` is future-only (delivers `handoff + 1`), so Phase 1 stops exactly at
Expand Down Expand Up @@ -854,6 +856,15 @@ mod tests {
tx.send(Arc::new(processed)).ok();
}

fn test_guard() -> SubscriptionLifecycleGuard {
SubscriptionLifecycleGuard::new(
"checkpoints",
&SubscriptionMetrics::new_for_test(),
&SubscriberLimit::new(10),
)
.unwrap()
}

#[tokio::test]
async fn subscribe_no_resume_yields_live_only() {
use futures::FutureExt;
Expand All @@ -862,7 +873,8 @@ mod tests {
// Fetcher is unused since resume_from is None.
let fetcher = MockFetcher::success_for_range(0..=0);

let stream = broadcast.subscribe(None, fetcher, &SubscriptionConfig::default());
let stream =
broadcast.subscribe(None, fetcher, &SubscriptionConfig::default(), test_guard());
tokio::pin!(stream);

// Poll once so the receiver gets pinned at tail=0 before any sends.
Expand All @@ -886,7 +898,12 @@ mod tests {

// resume_from = 2 → Phase 1 yields 3, 4, 5; then Phase 2 picks up live items.
let fetcher = MockFetcher::success_for_range(3..=5);
let stream = broadcast.subscribe(Some(2), fetcher, &SubscriptionConfig::default());
let stream = broadcast.subscribe(
Some(2),
fetcher,
&SubscriptionConfig::default(),
test_guard(),
);
tokio::pin!(stream);

// Phase 1 catches up via scan.
Expand All @@ -902,7 +919,8 @@ mod tests {
async fn subscribe_yields_error_when_channel_closes() {
let (tx, broadcast) = test_broadcast(/* first_live_checkpoint */ 1);
let fetcher = MockFetcher::success_for_range(0..=0);
let stream = broadcast.subscribe(None, fetcher, &SubscriptionConfig::default());
let stream =
broadcast.subscribe(None, fetcher, &SubscriptionConfig::default(), test_guard());
tokio::pin!(stream);

// Dropping the sender closes the channel; subscriber should yield an error and end.
Expand Down
Loading
Loading