mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-06 07:28:23 +00:00
chore(sync): classify transient watchdog recovery
Production at 57bc76ea forced exactly three relay.ngit.dev transient subscriptions through the 120-second watchdog every two minutes: 30 CLOSEs in the first 20 minutes. Subscription accounting remained within bounds, but the held-permit state and warning retained only random subscription IDs, so the stalled lifecycle could not be distinguished without guessing.
Require every auto-close REQ caller to provide one of six bounded request classes, retain that class with its per-session permit, and expose it in watchdog logs and a class/outcome Prometheus counter. Pagination distinguishes ordinary continuation from NIP-11 hint verification using its existing state; exact-ID purgatory polling remains correctly outside these labels because it uses the separately timed fetch_events path.
This is observability only: timeout, retry, pagination, request admission, filter contents, and slot-release behaviour are unchanged. Relay URLs and subscription IDs are deliberately excluded from metric labels, and no configuration surface is added.
Validated by the 646-test library suite, including fixed-label uniqueness and hint-verification classification tests, plus rustfmt and git diff --check.
This commit is contained in:
@@ -11,6 +11,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
- Expanded `grasp-audit` with stable machine-readable results, full JSON
|
- Expanded `grasp-audit` with stable machine-readable results, full JSON
|
||||||
reports, explicit audit identities, and hardened discovered-server probes.
|
reports, explicit audit identities, and hardened discovered-server probes.
|
||||||
|
- Added bounded request-class labels to transient-sync watchdog logs and
|
||||||
|
metrics so operators can identify which historic recovery path is stalling.
|
||||||
|
|
||||||
### Changed
|
### Changed
|
||||||
|
|
||||||
|
|||||||
@@ -372,7 +372,11 @@ So concurrency is not a free scaling axis; it is the residual of the ledger:
|
|||||||
five-request class cap inside the shared ledger: a slot is acquired when
|
five-request class cap inside the shared ledger: a slot is acquired when
|
||||||
the auto-close REQ is sent and released when its EOSE or CLOSED arrives
|
the auto-close REQ is sent and released when its EOSE or CLOSED arrives
|
||||||
(with a 120 s watchdog that sends CLOSE for only the unresponsive
|
(with a 120 s watchdog that sends CLOSE for only the unresponsive
|
||||||
subscription before releasing its slot). Live subscriptions are ledgered first, so NEG,
|
subscription before releasing its slot). Each held permit retains one of a
|
||||||
|
fixed set of request classes (historic page, pagination page, hint
|
||||||
|
verification, negentropy hydration, retry, or semantic fallback); watchdog
|
||||||
|
logs and metrics expose that class without using relay URLs or subscription
|
||||||
|
IDs as metric labels. Live subscriptions are ledgered first, so NEG,
|
||||||
transient REQ,
|
transient REQ,
|
||||||
and purgatory exact-ID polling share only the remaining capacity.
|
and purgatory exact-ID polling share only the remaining capacity.
|
||||||
- Permit acquisition checks relay health first: while a rate-limit or
|
- Permit acquisition checks relay health first: while a rate-limit or
|
||||||
|
|||||||
@@ -202,6 +202,26 @@ lazy_static! {
|
|||||||
.expect("register purgatory git fetch oids metric");
|
.expect("register purgatory git fetch oids metric");
|
||||||
metric
|
metric
|
||||||
};
|
};
|
||||||
|
static ref TRANSIENT_REQ_WATCHDOG_TOTAL: CounterVec = {
|
||||||
|
let metric = CounterVec::new(
|
||||||
|
Opts::new(
|
||||||
|
"ngit_sync_transient_req_watchdog_total",
|
||||||
|
"Transient REQ watchdog outcomes by bounded request class",
|
||||||
|
),
|
||||||
|
&["class", "outcome"],
|
||||||
|
)
|
||||||
|
.expect("build transient REQ watchdog metric");
|
||||||
|
REGISTRY
|
||||||
|
.register(Box::new(metric.clone()))
|
||||||
|
.expect("register transient REQ watchdog metric");
|
||||||
|
metric
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn record_transient_req_watchdog(class: &str, outcome: &str) {
|
||||||
|
TRANSIENT_REQ_WATCHDOG_TOTAL
|
||||||
|
.with_label_values(&[class, outcome])
|
||||||
|
.inc();
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Record one completed outbound purgatory git fetch pass.
|
/// Record one completed outbound purgatory git fetch pass.
|
||||||
|
|||||||
+49
-7
@@ -34,7 +34,9 @@ pub use rejected_index::{EventType, RejectionReason};
|
|||||||
// Current code still uses the simple HashSet type alias below
|
// Current code still uses the simple HashSet type alias below
|
||||||
|
|
||||||
// Re-export relay connection types
|
// Re-export relay connection types
|
||||||
pub use relay_connection::{NegentropySyncResult, RelayConnection, RelayEvent};
|
pub use relay_connection::{
|
||||||
|
NegentropySyncResult, RelayConnection, RelayEvent, TransientRequestClass,
|
||||||
|
};
|
||||||
|
|
||||||
// Re-export self-subscriber types
|
// Re-export self-subscriber types
|
||||||
pub use self_subscriber::SelfSubscriber;
|
pub use self_subscriber::SelfSubscriber;
|
||||||
@@ -466,6 +468,18 @@ impl PaginationState {
|
|||||||
.collect()
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn request_class(&self) -> TransientRequestClass {
|
||||||
|
if self
|
||||||
|
.filters
|
||||||
|
.iter()
|
||||||
|
.any(|state| state.verification_baseline.is_some())
|
||||||
|
{
|
||||||
|
TransientRequestClass::PaginationVerification
|
||||||
|
} else {
|
||||||
|
TransientRequestClass::PaginationPage
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn next_page(mut self, session: &mut RelayPaginationSession) -> Option<Self> {
|
fn next_page(mut self, session: &mut RelayPaginationSession) -> Option<Self> {
|
||||||
let mut next = Vec::new();
|
let mut next = Vec::new();
|
||||||
for mut state in self.filters.drain(..) {
|
for mut state in self.filters.drain(..) {
|
||||||
@@ -1491,7 +1505,10 @@ impl SyncManager {
|
|||||||
|
|
||||||
let mut next_page_started = false;
|
let mut next_page_started = false;
|
||||||
if let Some(conn) = self.connections.get(&relay_url_for_pagination) {
|
if let Some(conn) = self.connections.get(&relay_url_for_pagination) {
|
||||||
match conn.subscribe_filters(next_filters.clone(), true).await {
|
match conn
|
||||||
|
.subscribe_filters(next_filters.clone(), next_page.request_class())
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(new_sub_id) => {
|
Ok(new_sub_id) => {
|
||||||
let mut pending = self.pending_sync_index.write().await;
|
let mut pending = self.pending_sync_index.write().await;
|
||||||
if let Some(batches) = pending.get_mut(&relay_url_for_pagination) {
|
if let Some(batches) = pending.get_mut(&relay_url_for_pagination) {
|
||||||
@@ -1637,7 +1654,13 @@ impl SyncManager {
|
|||||||
let mut new_sub_ids = HashSet::new();
|
let mut new_sub_ids = HashSet::new();
|
||||||
if let Some(conn) = self.connections.get(&relay_url_for_fallback) {
|
if let Some(conn) = self.connections.get(&relay_url_for_fallback) {
|
||||||
for filter_group in group_filters_for_req(&fallback_filters) {
|
for filter_group in group_filters_for_req(&fallback_filters) {
|
||||||
match conn.subscribe_filters(filter_group, true).await {
|
match conn
|
||||||
|
.subscribe_filters(
|
||||||
|
filter_group,
|
||||||
|
TransientRequestClass::NegentropyFallback,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(sub_id) => {
|
Ok(sub_id) => {
|
||||||
new_sub_ids.insert(sub_id);
|
new_sub_ids.insert(sub_id);
|
||||||
}
|
}
|
||||||
@@ -1808,7 +1831,10 @@ impl SyncManager {
|
|||||||
let mut new_sub_ids = HashSet::new();
|
let mut new_sub_ids = HashSet::new();
|
||||||
if let Some(conn) = self.connections.get(&relay_url_for_retry) {
|
if let Some(conn) = self.connections.get(&relay_url_for_retry) {
|
||||||
for filter in retry_filters {
|
for filter in retry_filters {
|
||||||
match conn.subscribe_filter(filter, true).await {
|
match conn
|
||||||
|
.subscribe_filter(filter, TransientRequestClass::NegentropyRetry)
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(sub_id) => {
|
Ok(sub_id) => {
|
||||||
new_sub_ids.insert(sub_id);
|
new_sub_ids.insert(sub_id);
|
||||||
}
|
}
|
||||||
@@ -1958,7 +1984,7 @@ impl SyncManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
match connection
|
match connection
|
||||||
.subscribe_filters(next_filters.clone(), true)
|
.subscribe_filters(next_filters.clone(), next_page.request_class())
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(new_sub_id) => {
|
Ok(new_sub_id) => {
|
||||||
@@ -5573,7 +5599,13 @@ impl SyncManager {
|
|||||||
let mut subscription_ids = HashSet::new();
|
let mut subscription_ids = HashSet::new();
|
||||||
for (idx, filter) in ids_filters.iter().enumerate() {
|
for (idx, filter) in ids_filters.iter().enumerate() {
|
||||||
if let Some(conn) = self.connections.get(relay_url) {
|
if let Some(conn) = self.connections.get(relay_url) {
|
||||||
match conn.subscribe_filter(filter.clone(), true).await {
|
match conn
|
||||||
|
.subscribe_filter(
|
||||||
|
filter.clone(),
|
||||||
|
TransientRequestClass::NegentropyHydration,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(sub_id) => {
|
Ok(sub_id) => {
|
||||||
subscription_ids.insert(sub_id);
|
subscription_ids.insert(sub_id);
|
||||||
}
|
}
|
||||||
@@ -5646,7 +5678,13 @@ impl SyncManager {
|
|||||||
|
|
||||||
if let Some(conn) = self.connections.get(relay_url) {
|
if let Some(conn) = self.connections.get(relay_url) {
|
||||||
let grouped_filters = filter_group;
|
let grouped_filters = filter_group;
|
||||||
match conn.subscribe_filters(grouped_filters.clone(), true).await {
|
match conn
|
||||||
|
.subscribe_filters(
|
||||||
|
grouped_filters.clone(),
|
||||||
|
TransientRequestClass::HistoricPage,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(sub_id) => {
|
Ok(sub_id) => {
|
||||||
subscription_ids.insert(sub_id.clone());
|
subscription_ids.insert(sub_id.clone());
|
||||||
pagination_state.insert(sub_id, PaginationState::new(grouped_filters));
|
pagination_state.insert(sub_id, PaginationState::new(grouped_filters));
|
||||||
@@ -6195,6 +6233,10 @@ mod tests {
|
|||||||
.next_page(&mut session)
|
.next_page(&mut session)
|
||||||
.expect("a suspiciously short hinted page needs verification");
|
.expect("a suspiciously short hinted page needs verification");
|
||||||
assert_eq!(session.hint, PaginationHint::Verifying(1000));
|
assert_eq!(session.hint, PaginationHint::Verifying(1000));
|
||||||
|
assert_eq!(
|
||||||
|
verification.request_class(),
|
||||||
|
TransientRequestClass::PaginationVerification
|
||||||
|
);
|
||||||
|
|
||||||
let unseen = EventBuilder::new(Kind::TextNote, "older unseen event")
|
let unseen = EventBuilder::new(Kind::TextNote, "older unseen event")
|
||||||
.custom_created_at(Timestamp::from_secs(99))
|
.custom_created_at(Timestamp::from_secs(99))
|
||||||
|
|||||||
+104
-27
@@ -119,8 +119,43 @@ struct HeldTransientPermits {
|
|||||||
_class_cap: tokio::sync::OwnedSemaphorePermit,
|
_class_cap: tokio::sync::OwnedSemaphorePermit,
|
||||||
_ledger_slot: tokio::sync::OwnedSemaphorePermit,
|
_ledger_slot: tokio::sync::OwnedSemaphorePermit,
|
||||||
generation: u64,
|
generation: u64,
|
||||||
|
request_class: TransientRequestClass,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Bounded origin of an auto-close REQ, retained for watchdog diagnosis.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum TransientRequestClass {
|
||||||
|
HistoricPage,
|
||||||
|
PaginationPage,
|
||||||
|
PaginationVerification,
|
||||||
|
NegentropyHydration,
|
||||||
|
NegentropyRetry,
|
||||||
|
NegentropyFallback,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TransientRequestClass {
|
||||||
|
pub const fn as_str(self) -> &'static str {
|
||||||
|
match self {
|
||||||
|
Self::HistoricPage => "historic_page",
|
||||||
|
Self::PaginationPage => "pagination_page",
|
||||||
|
Self::PaginationVerification => "pagination_verification",
|
||||||
|
Self::NegentropyHydration => "negentropy_hydration",
|
||||||
|
Self::NegentropyRetry => "negentropy_retry",
|
||||||
|
Self::NegentropyFallback => "negentropy_fallback",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
const TRANSIENT_REQUEST_CLASSES: [TransientRequestClass; 6] = [
|
||||||
|
TransientRequestClass::HistoricPage,
|
||||||
|
TransientRequestClass::PaginationPage,
|
||||||
|
TransientRequestClass::PaginationVerification,
|
||||||
|
TransientRequestClass::NegentropyHydration,
|
||||||
|
TransientRequestClass::NegentropyRetry,
|
||||||
|
TransientRequestClass::NegentropyFallback,
|
||||||
|
];
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
struct SubscriptionLedgerSession {
|
struct SubscriptionLedgerSession {
|
||||||
generation: u64,
|
generation: u64,
|
||||||
@@ -655,7 +690,10 @@ impl RelayConnection {
|
|||||||
/// Acquire both transient constraints without waiting while holding only
|
/// Acquire both transient constraints without waiting while holding only
|
||||||
/// one of them. This prevents class-cap waiters from pinning ledger slots
|
/// one of them. This prevents class-cap waiters from pinning ledger slots
|
||||||
/// and ledger waiters from pinning the class cap.
|
/// and ledger waiters from pinning the class cap.
|
||||||
async fn acquire_transient_permits(&self) -> Result<HeldTransientPermits, String> {
|
async fn acquire_transient_permits(
|
||||||
|
&self,
|
||||||
|
request_class: TransientRequestClass,
|
||||||
|
) -> Result<HeldTransientPermits, String> {
|
||||||
loop {
|
loop {
|
||||||
let ledger_slot = self.acquire_subscription_slots(1).await?;
|
let ledger_slot = self.acquire_subscription_slots(1).await?;
|
||||||
if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() {
|
if let Ok(class_cap) = self.transient_req_permits.clone().try_acquire_owned() {
|
||||||
@@ -663,6 +701,7 @@ impl RelayConnection {
|
|||||||
_class_cap: class_cap,
|
_class_cap: class_cap,
|
||||||
generation: ledger_slot.generation,
|
generation: ledger_slot.generation,
|
||||||
_ledger_slot: ledger_slot.permit,
|
_ledger_slot: ledger_slot.permit,
|
||||||
|
request_class,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
drop(ledger_slot);
|
drop(ledger_slot);
|
||||||
@@ -678,6 +717,7 @@ impl RelayConnection {
|
|||||||
_class_cap: class_cap,
|
_class_cap: class_cap,
|
||||||
generation: ledger_slot.generation,
|
generation: ledger_slot.generation,
|
||||||
_ledger_slot: ledger_slot.permit,
|
_ledger_slot: ledger_slot.permit,
|
||||||
|
request_class,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
drop(class_cap);
|
drop(class_cap);
|
||||||
@@ -741,7 +781,7 @@ impl RelayConnection {
|
|||||||
filter_groups: Vec<Vec<Filter>>,
|
filter_groups: Vec<Vec<Filter>>,
|
||||||
) -> Result<Vec<SubscriptionId>, String> {
|
) -> Result<Vec<SubscriptionId>, String> {
|
||||||
self.replace_live_filter_groups_with(filter_groups, |filters, permit| {
|
self.replace_live_filter_groups_with(filter_groups, |filters, permit| {
|
||||||
self.subscribe_filters_with_live_permit(filters, false, Some(permit))
|
self.subscribe_filters_with_live_permit(filters, None, Some(permit))
|
||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
@@ -1009,7 +1049,7 @@ impl RelayConnection {
|
|||||||
///
|
///
|
||||||
/// # Arguments
|
/// # Arguments
|
||||||
/// * `filter` - The filter to subscribe to
|
/// * `filter` - The filter to subscribe to
|
||||||
/// * `auto_close` - If true, subscription automatically closes after EOSE (for historic sync). If false, stays open for new events (for live sync).
|
/// * `request_class` - Bounded historic-sync origin retained for watchdog diagnostics.
|
||||||
///
|
///
|
||||||
/// # Returns
|
/// # Returns
|
||||||
/// * `Ok(SubscriptionId)` - The subscription ID on success
|
/// * `Ok(SubscriptionId)` - The subscription ID on success
|
||||||
@@ -1017,9 +1057,9 @@ impl RelayConnection {
|
|||||||
pub async fn subscribe_filter(
|
pub async fn subscribe_filter(
|
||||||
&self,
|
&self,
|
||||||
filter: Filter,
|
filter: Filter,
|
||||||
auto_close: bool,
|
request_class: TransientRequestClass,
|
||||||
) -> Result<SubscriptionId, String> {
|
) -> Result<SubscriptionId, String> {
|
||||||
self.subscribe_filters(vec![filter], auto_close).await
|
self.subscribe_filters(vec![filter], request_class).await
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Subscribe to several OR filters under one NIP-01 subscription ID.
|
/// Subscribe to several OR filters under one NIP-01 subscription ID.
|
||||||
@@ -1030,9 +1070,9 @@ impl RelayConnection {
|
|||||||
pub async fn subscribe_filters(
|
pub async fn subscribe_filters(
|
||||||
&self,
|
&self,
|
||||||
filters: Vec<Filter>,
|
filters: Vec<Filter>,
|
||||||
auto_close: bool,
|
request_class: TransientRequestClass,
|
||||||
) -> Result<SubscriptionId, String> {
|
) -> Result<SubscriptionId, String> {
|
||||||
self.subscribe_filters_with_live_permit(filters, auto_close, None)
|
self.subscribe_filters_with_live_permit(filters, Some(request_class), None)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1044,7 +1084,7 @@ impl RelayConnection {
|
|||||||
filter_groups: Vec<Vec<Filter>>,
|
filter_groups: Vec<Vec<Filter>>,
|
||||||
) -> Result<Vec<SubscriptionId>, String> {
|
) -> Result<Vec<SubscriptionId>, String> {
|
||||||
self.subscribe_live_filter_groups_with(filter_groups, |filters, permit| {
|
self.subscribe_live_filter_groups_with(filter_groups, |filters, permit| {
|
||||||
self.subscribe_filters_with_live_permit(filters, false, Some(permit))
|
self.subscribe_filters_with_live_permit(filters, None, Some(permit))
|
||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
@@ -1101,7 +1141,7 @@ impl RelayConnection {
|
|||||||
async fn subscribe_filters_with_live_permit(
|
async fn subscribe_filters_with_live_permit(
|
||||||
&self,
|
&self,
|
||||||
filters: Vec<Filter>,
|
filters: Vec<Filter>,
|
||||||
auto_close: bool,
|
transient_class: Option<TransientRequestClass>,
|
||||||
reserved_live_permit: Option<SessionPermit>,
|
reserved_live_permit: Option<SessionPermit>,
|
||||||
) -> Result<SubscriptionId, String> {
|
) -> Result<SubscriptionId, String> {
|
||||||
if filters.is_empty() {
|
if filters.is_empty() {
|
||||||
@@ -1112,7 +1152,7 @@ impl RelayConnection {
|
|||||||
relay = %self.url,
|
relay = %self.url,
|
||||||
filter_count = filters.len(),
|
filter_count = filters.len(),
|
||||||
filters = ?filters,
|
filters = ?filters,
|
||||||
auto_close = auto_close,
|
auto_close = transient_class.is_some(),
|
||||||
"subscribe_filters called"
|
"subscribe_filters called"
|
||||||
);
|
);
|
||||||
|
|
||||||
@@ -1121,7 +1161,7 @@ impl RelayConnection {
|
|||||||
// exceeding relay subscription budgets. The permit is registered
|
// exceeding relay subscription budgets. The permit is registered
|
||||||
// against the subscription id on success and released when the
|
// against the subscription id on success and released when the
|
||||||
// subscription's EOSE or CLOSED arrives (see run_event_loop).
|
// subscription's EOSE or CLOSED arrives (see run_event_loop).
|
||||||
let (transient_permit, live_permit) = if auto_close {
|
let (transient_permit, live_permit) = if let Some(request_class) = transient_class {
|
||||||
if self.historic_capacity_consumed_by_live() {
|
if self.historic_capacity_consumed_by_live() {
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
relay = %self.url,
|
relay = %self.url,
|
||||||
@@ -1132,7 +1172,10 @@ impl RelayConnection {
|
|||||||
self.url
|
self.url
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
(Some(self.acquire_transient_permits().await?), None)
|
(
|
||||||
|
Some(self.acquire_transient_permits(request_class).await?),
|
||||||
|
None,
|
||||||
|
)
|
||||||
} else {
|
} else {
|
||||||
let acquired = match reserved_live_permit {
|
let acquired = match reserved_live_permit {
|
||||||
Some(permit) => Ok(permit),
|
Some(permit) => Ok(permit),
|
||||||
@@ -1162,7 +1205,7 @@ impl RelayConnection {
|
|||||||
// Transient permits are acquired through the same session helper; a
|
// Transient permits are acquired through the same session helper; a
|
||||||
// reset closes queued acquisitions and the helper checks generation.
|
// reset closes queued acquisitions and the helper checks generation.
|
||||||
let retained_filters = filters.clone();
|
let retained_filters = filters.clone();
|
||||||
let output = if auto_close {
|
let output = if transient_class.is_some() {
|
||||||
self.client
|
self.client
|
||||||
.subscribe(filters)
|
.subscribe(filters)
|
||||||
.close_on(
|
.close_on(
|
||||||
@@ -1184,7 +1227,7 @@ impl RelayConnection {
|
|||||||
return Err(format!("Failed to subscribe on {}: {}", self.url, failures));
|
return Err(format!("Failed to subscribe on {}: {}", self.url, failures));
|
||||||
}
|
}
|
||||||
|
|
||||||
if auto_close {
|
if transient_class.is_some() {
|
||||||
if let Some(permit) = transient_permit {
|
if let Some(permit) = transient_permit {
|
||||||
self.hold_transient_req_permit(output.value.clone(), permit);
|
self.hold_transient_req_permit(output.value.clone(), permit);
|
||||||
}
|
}
|
||||||
@@ -1323,33 +1366,38 @@ impl RelayConnection {
|
|||||||
let relay_url = self.url.clone();
|
let relay_url = self.url.clone();
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
tokio::time::sleep(TRANSIENT_REQ_PERMIT_TIMEOUT).await;
|
tokio::time::sleep(TRANSIENT_REQ_PERMIT_TIMEOUT).await;
|
||||||
let generation = held
|
let held_request = held
|
||||||
.lock()
|
.lock()
|
||||||
.expect("transient permit map poisoned")
|
.expect("transient permit map poisoned")
|
||||||
.get(&sub_id)
|
.get(&sub_id)
|
||||||
.map(|permits| permits.generation);
|
.map(|permits| (permits.generation, permits.request_class));
|
||||||
if let (Some(generation), Ok(Some(relay))) = (generation, client.relay(&url).await) {
|
if let (Some((generation, request_class)), Ok(Some(relay))) =
|
||||||
|
(held_request, client.relay(&url).await)
|
||||||
|
{
|
||||||
if relay
|
if relay
|
||||||
.send_msg(ClientMessage::close(sub_id.clone()))
|
.send_msg(ClientMessage::close(sub_id.clone()))
|
||||||
.await
|
.await
|
||||||
.is_ok()
|
.is_ok()
|
||||||
{
|
{
|
||||||
Self::release_transient_req_permit_for_generation(
|
Self::release_transient_req_permit_for_generation(&held, &sub_id, generation);
|
||||||
&held,
|
|
||||||
&sub_id,
|
|
||||||
generation,
|
|
||||||
);
|
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
relay = %relay_url,
|
relay = %relay_url,
|
||||||
sub_id = %sub_id,
|
sub_id = %sub_id,
|
||||||
|
request_class = request_class.as_str(),
|
||||||
"Transient REQ watchdog sent CLOSE after missing EOSE/CLOSED"
|
"Transient REQ watchdog sent CLOSE after missing EOSE/CLOSED"
|
||||||
);
|
);
|
||||||
|
crate::metrics::record_transient_req_watchdog(request_class.as_str(), "closed");
|
||||||
} else {
|
} else {
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
relay = %relay_url,
|
relay = %relay_url,
|
||||||
sub_id = %sub_id,
|
sub_id = %sub_id,
|
||||||
|
request_class = request_class.as_str(),
|
||||||
"Transient REQ watchdog could not send CLOSE; retaining subscription slot"
|
"Transient REQ watchdog could not send CLOSE; retaining subscription slot"
|
||||||
);
|
);
|
||||||
|
crate::metrics::record_transient_req_watchdog(
|
||||||
|
request_class.as_str(),
|
||||||
|
"close_failed",
|
||||||
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
@@ -2313,7 +2361,10 @@ mod tests {
|
|||||||
assert_eq!(ids.len(), 8);
|
assert_eq!(ids.len(), 8);
|
||||||
assert_eq!(connection.subscription_budget().available_permits(), 0);
|
assert_eq!(connection.subscription_budget().available_permits(), 0);
|
||||||
let error = connection
|
let error = connection
|
||||||
.subscribe_filter(Filter::new().kind(Kind::TextNote), true)
|
.subscribe_filter(
|
||||||
|
Filter::new().kind(Kind::TextNote),
|
||||||
|
TransientRequestClass::HistoricPage,
|
||||||
|
)
|
||||||
.await
|
.await
|
||||||
.expect_err("history must defer when complete live coverage fills usable capacity");
|
.expect_err("history must defer when complete live coverage fills usable capacity");
|
||||||
assert!(
|
assert!(
|
||||||
@@ -2368,6 +2419,7 @@ mod tests {
|
|||||||
let generation = connection.current_subscription_generation();
|
let generation = connection.current_subscription_generation();
|
||||||
let held = HeldTransientPermits {
|
let held = HeldTransientPermits {
|
||||||
generation,
|
generation,
|
||||||
|
request_class: TransientRequestClass::HistoricPage,
|
||||||
_class_cap: connection
|
_class_cap: connection
|
||||||
.transient_req_permits
|
.transient_req_permits
|
||||||
.clone()
|
.clone()
|
||||||
@@ -2449,8 +2501,11 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
.expect("occupy transient class cap");
|
.expect("occupy transient class cap");
|
||||||
let class_waiter_connection = connection.clone();
|
let class_waiter_connection = connection.clone();
|
||||||
let mut class_waiter =
|
let mut class_waiter = tokio::spawn(async move {
|
||||||
tokio::spawn(async move { class_waiter_connection.acquire_transient_permits().await });
|
class_waiter_connection
|
||||||
|
.acquire_transient_permits(TransientRequestClass::HistoricPage)
|
||||||
|
.await
|
||||||
|
});
|
||||||
assert!(
|
assert!(
|
||||||
tokio::time::timeout(Duration::from_millis(50), &mut class_waiter)
|
tokio::time::timeout(Duration::from_millis(50), &mut class_waiter)
|
||||||
.await
|
.await
|
||||||
@@ -2471,8 +2526,11 @@ mod tests {
|
|||||||
.await
|
.await
|
||||||
.expect("occupy subscription ledger");
|
.expect("occupy subscription ledger");
|
||||||
let ledger_waiter_connection = connection.clone();
|
let ledger_waiter_connection = connection.clone();
|
||||||
let mut ledger_waiter =
|
let mut ledger_waiter = tokio::spawn(async move {
|
||||||
tokio::spawn(async move { ledger_waiter_connection.acquire_transient_permits().await });
|
ledger_waiter_connection
|
||||||
|
.acquire_transient_permits(TransientRequestClass::HistoricPage)
|
||||||
|
.await
|
||||||
|
});
|
||||||
assert!(
|
assert!(
|
||||||
tokio::time::timeout(Duration::from_millis(50), &mut ledger_waiter)
|
tokio::time::timeout(Duration::from_millis(50), &mut ledger_waiter)
|
||||||
.await
|
.await
|
||||||
@@ -2600,6 +2658,25 @@ mod tests {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn transient_request_metric_labels_are_bounded_and_unique() {
|
||||||
|
let labels = TRANSIENT_REQUEST_CLASSES.map(TransientRequestClass::as_str);
|
||||||
|
let unique = labels.into_iter().collect::<std::collections::HashSet<_>>();
|
||||||
|
|
||||||
|
assert_eq!(unique.len(), TRANSIENT_REQUEST_CLASSES.len());
|
||||||
|
assert_eq!(
|
||||||
|
unique,
|
||||||
|
std::collections::HashSet::from([
|
||||||
|
"historic_page",
|
||||||
|
"pagination_page",
|
||||||
|
"pagination_verification",
|
||||||
|
"negentropy_hydration",
|
||||||
|
"negentropy_retry",
|
||||||
|
"negentropy_fallback",
|
||||||
|
])
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_normalize_url_real_world_example() {
|
fn test_normalize_url_real_world_example() {
|
||||||
// Test the exact case from the bug report
|
// Test the exact case from the bug report
|
||||||
|
|||||||
Reference in New Issue
Block a user