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,6 +32,8 @@ 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
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ use crate::task::streaming::CheckpointBroadcaster;
use crate::task::streaming::ProcessedCheckpoint;
use crate::task::streaming::StreamedCaches;
use crate::task::streaming::SubscriptionBroadcast;
use crate::task::streaming::SubscriptionLifecycleGuard;
use crate::task::streaming::SubscriptionTerminationReason;
use crate::task::streaming::broadcast_error;
use crate::task::streaming::reconnect_error;
use crate::task::streaming::wait_for_pipelines_catching_up_at;
Expand All @@ -55,6 +57,9 @@ 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 @@ -137,6 +142,7 @@ 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 All @@ -159,6 +165,7 @@ pub(super) fn subscribe<S: Subscribable>(
let scanned = match scanned {
Ok(scanned) => scanned,
Err(e) => {
guard.terminate(SubscriptionTerminationReason::BackfillError);
yield Err(e);
return;
}
Expand All @@ -179,6 +186,7 @@ pub(super) fn subscribe<S: Subscribable>(
break;
}
yield Ok(edge);
guard.record_backfill_delivered();
}
// A coverage marker (its `checkpoint` is the fully-scanned frontier): stop once
// it has covered the handoff.
Expand All @@ -194,7 +202,9 @@ pub(super) fn subscribe<S: Subscribable>(

// Phase 2: follow live from `handoff + 1` (a fresh receiver if there was no backfill).
let receiver = pending_receiver.unwrap_or_else(|| broadcast.broadcaster().resubscribe());
for await edge in live::<S>(receiver, last_checkpoint, caches, resolver_limits, filter) {
for await edge in
live::<S>(receiver, last_checkpoint, caches, resolver_limits, filter, guard)
{
yield edge;
}
}
Expand Down Expand Up @@ -271,6 +281,7 @@ fn live<S: Subscribable>(
caches: Arc<StreamedCaches>,
resolver_limits: sui_package_resolver::Limits,
filter: S::Filter,
guard: SubscriptionLifecycleGuard,
) -> impl Stream<Item = Result<Edge<String, S::Item, EmptyFields>, RpcError>> {
stream! {
let mut delivered_live = false;
Expand All @@ -289,26 +300,37 @@ fn live<S: Subscribable>(
received = seq,
"Unexpected gap between scan and live; disconnecting"
);
guard.terminate(SubscriptionTerminationReason::UnexpectedGap);
yield Err(reconnect_error());
return;
}
}
// Deliver each matching item as its own payload, ordered within the checkpoint.
// Empty checkpoints yield nothing.
let edges = S::matching_edges(&checkpoint, &caches, &resolver_limits, &filter)?;
let edges = match S::matching_edges(&checkpoint, &caches, &resolver_limits, &filter) {
Ok(edges) => edges,
Err(e) => {
guard.terminate(SubscriptionTerminationReason::Error);
yield Err(e);
return;
}
};
for edge in edges {
yield Ok(edge);
guard.record_delivered(checkpoint.summary.timestamp_ms);
}
last_checkpoint = Some(seq);
delivered_live = true;
}
// A lag before the first live checkpoint is catch-up overflow (likely kv-rpc lag).
Err(broadcast::error::RecvError::Lagged(missed)) if !delivered_live => {
warn!(missed, "Subscriber fell behind during catch-up; disconnecting");
guard.terminate(SubscriptionTerminationReason::Lagged);
yield Err(reconnect_error());
return;
}
Err(e) => {
guard.terminate(SubscriptionTerminationReason::from_recv_error(&e));
yield Err(broadcast_error(e));
return;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ 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
29 changes: 19 additions & 10 deletions crates/sui-indexer-alt-graphql/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -303,10 +303,14 @@ struct IdeEnabled(bool);
#[derive(Clone, Copy)]
struct SubscriptionsEnabled(bool);

/// Per-subscriber delivery rate in output nodes per second, surfaced to the subscription handler so
/// it can pace each payload by its cost. `0` disables pacing.
#[derive(Clone, Copy)]
struct SubscriptionThrottleRate(u32);
/// Per-subscriber delivery throttle settings surfaced to the subscription handler: the rate in
/// output nodes per second (`0` disables pacing) and the metric each payload's pacing delay is
/// observed into.
#[derive(Clone)]
struct SubscriptionThrottle {
nodes_per_second: u32,
delay_metric: prometheus::Histogram,
}

/// Set-up and run the RPC service, using the provided arguments (expected to be extracted from the
/// command-line).
Expand Down Expand Up @@ -394,6 +398,8 @@ pub async fn start_rpc(
metrics.clone(),
);

let subscription_metrics = Arc::new(SubscriptionMetrics::new(registry));

let streaming_setup = match subscription_args.checkpoint_stream_url {
Some(uri) => {
let ledger_grpc = ledger_grpc_reader
Expand All @@ -417,7 +423,7 @@ pub async fn start_rpc(
readiness.clone(),
ledger_grpc.clone(),
watermark_task.watermarks_rx(),
Arc::new(SubscriptionMetrics::new(registry)),
subscription_metrics.clone(),
);
let caches = Arc::new(StreamedCaches::new(
streaming_packages,
Expand All @@ -435,6 +441,7 @@ pub async fn start_rpc(
None => None,
};

let throttle_delay_metric = subscription_metrics.subscriber_throttle_delay.clone();
let mut rpc = rpc
.route(GRAPHQL_PATH, post(graphql).get(graphiql))
.route(GRAPHQL_SUBSCRIPTIONS_PATH, post(graphql_subscriptions))
Expand Down Expand Up @@ -473,11 +480,12 @@ pub async fn start_rpc(

let subscriptions_enabled = streaming_setup.is_some();
rpc = rpc.layer(SubscriptionsEnabled(subscriptions_enabled));
rpc = rpc.layer(SubscriptionThrottleRate(
config
rpc = rpc.layer(SubscriptionThrottle {
nodes_per_second: config
.subscription
.per_subscriber_max_output_nodes_per_second,
));
delay_metric: throttle_delay_metric,
});

// The transaction subscription backfill waits on pipeline watermarks to gate delivery, so it
// needs a live view of them. Captured before the watermark task is consumed by `run()`.
Expand Down Expand Up @@ -513,6 +521,7 @@ pub async fn start_rpc(
let subscription_broadcast = Arc::new(SubscriptionBroadcast::new(
_broadcaster,
first_live_checkpoint,
subscription_metrics.clone(),
));
rpc = rpc
.data(subscription_broadcast)
Expand Down Expand Up @@ -584,7 +593,7 @@ async fn graphql_subscriptions(
ConnectInfo(addr): ConnectInfo<SocketAddr>,
Extension(schema): Extension<Schema<Query, Mutation, Subscription>>,
Extension(SubscriptionsEnabled(subscriptions_enabled)): Extension<SubscriptionsEnabled>,
Extension(SubscriptionThrottleRate(nodes_per_second)): Extension<SubscriptionThrottleRate>,
Extension(throttle_cfg): Extension<SubscriptionThrottle>,
Extension(watermark): Extension<WatermarksLock>,
request: GraphQLRequest,
) -> axum::response::Response {
Expand All @@ -600,7 +609,7 @@ async fn graphql_subscriptions(
// Query depth is computed once by the query-limits extension during validation and stashed here,
// so the throttle can add its depth surcharge to each payload's cost.
let query_depth = QueryDepth::default();
let throttle = Throttle::new(nodes_per_second);
let throttle = Throttle::new(throttle_cfg.nodes_per_second, throttle_cfg.delay_metric);
let req = request
.into_inner()
.data(Session::new(addr))
Expand Down
83 changes: 83 additions & 0 deletions crates/sui-indexer-alt-graphql/src/metrics/subscription.rs
Original file line number Diff line number Diff line change
@@ -1,12 +1,14 @@
// Copyright (c) Mysten Labs, Inc.
// SPDX-License-Identifier: Apache-2.0

use prometheus::Histogram;
use prometheus::HistogramVec;
use prometheus::IntCounter;
use prometheus::IntCounterVec;
use prometheus::IntGaugeVec;
use prometheus::Registry;
use prometheus::register_histogram_vec_with_registry;
use prometheus::register_histogram_with_registry;
use prometheus::register_int_counter_vec_with_registry;
use prometheus::register_int_counter_with_registry;
use prometheus::register_int_gauge_vec_with_registry;
Expand All @@ -17,6 +19,22 @@ const LAG_SEC_BUCKETS: &[f64] = &[
0.95, 1.0, 2.0, 3.0, 4.0, 5.0, 10.0, 20.0, 50.0, 100.0, 1000.0,
];

/// Histogram buckets for subscription lifetime, in seconds: sub-second churn through multi-hour
/// holds, with day/week edges to size the long-lived tail (subscriptions may run indefinitely).
const DURATION_SEC_BUCKETS: &[f64] = &[
1.0, 5.0, 15.0, 30.0, // seconds
60.0, 300.0, 900.0, 1800.0, // 1m, 5m, 15m, 30m
3600.0, 7200.0, 21600.0, // 1h, 2h, 6h
43200.0, 86400.0, // 12h, 24h
172800.0, 604800.0, // 2d, 7d
];

/// Histogram buckets for the per-payload throttle pause, in seconds: sub-millisecond when the rate
/// budget is generous, up to a few seconds for a heavy payload against a tight budget.
const THROTTLE_DELAY_SEC_BUCKETS: &[f64] = &[
0.0005, 0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0,
];

/// Metrics specific to the streaming subscription feature. Subscription payloads also record into
/// the shared field-resolution (`fields_*`) and query-complexity (`input_nodes`, `output_nodes`)
/// metrics on `RpcMetrics`.
Expand All @@ -31,6 +49,15 @@ pub struct SubscriptionMetrics {
pub upstream_latest_processed_checkpoint_timestamp_ms: IntGaugeVec,
pub upstream_gap_recoveries: IntCounter,
pub upstream_malformed_checkpoints: IntCounter,

// Metrics aggregated across all subscribers.
pub active_subscriptions: IntGaugeVec,
pub subscriptions_opened: IntCounterVec,
pub subscription_terminations: IntCounterVec,
pub subscription_duration: HistogramVec,
pub payloads_delivered: IntCounterVec,
pub live_payload_delivery_checkpoint_timestamp_lag: HistogramVec,
pub subscriber_throttle_delay: Histogram,
}

impl SubscriptionMetrics {
Expand Down Expand Up @@ -97,6 +124,57 @@ impl SubscriptionMetrics {
registry,
)
.unwrap(),
active_subscriptions: register_int_gauge_vec_with_registry!(
"graphql_subscription_active_subscriptions",
"Number of subscriptions currently active, by type",
&["type"],
registry,
)
.unwrap(),
subscriptions_opened: register_int_counter_vec_with_registry!(
"graphql_subscription_opened",
"Total subscriptions opened, by type",
&["type"],
registry,
)
.unwrap(),
subscription_terminations: register_int_counter_vec_with_registry!(
"graphql_subscription_terminations",
"Total subscriptions that have ended, by type and reason",
&["type", "reason"],
registry,
)
.unwrap(),
subscription_duration: register_histogram_vec_with_registry!(
"graphql_subscription_duration_seconds",
"Lifetime of each subscription from open to close, in seconds, by type",
&["type"],
DURATION_SEC_BUCKETS.to_vec(),
registry,
)
.unwrap(),
payloads_delivered: register_int_counter_vec_with_registry!(
"graphql_subscription_payloads_delivered",
"Total payloads delivered to subscribers, by type and phase (live or backfill)",
&["type", "phase"],
registry,
)
.unwrap(),
live_payload_delivery_checkpoint_timestamp_lag: register_histogram_vec_with_registry!(
"graphql_subscription_live_payload_delivery_checkpoint_timestamp_lag",
"Seconds from a live payload's checkpoint timestamp to when it was delivered to a subscriber, by type (server-side; backfill deliveries excluded, since their lag is catch-up distance not staleness)",
&["type"],
LAG_SEC_BUCKETS.to_vec(),
registry,
)
.unwrap(),
subscriber_throttle_delay: register_histogram_with_registry!(
"graphql_subscription_throttle_delay_seconds",
"Seconds the delivery throttle paused before the next payload, observed once per payload across all subscribers and phases (zero while the rate budget is not binding)",
THROTTLE_DELAY_SEC_BUCKETS.to_vec(),
registry,
)
.unwrap(),
}
}

Expand Down Expand Up @@ -142,6 +220,11 @@ impl SubscriptionMetrics {
.with_label_values(&[phase])
.observe(lag_ms as f64 / 1000.0);
}

#[cfg(test)]
pub(crate) fn new_for_test() -> std::sync::Arc<Self> {
std::sync::Arc::new(Self::new(&Registry::new()))
}
}

#[cfg(test)]
Expand Down
Loading
Loading