Skip to main content

uqa_client/notifications/http/
options.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7use super::HttpNotificationError;
8use crate::notifications::{TimerLimits, MAX_NOTIFICATION_WIRE_BYTES};
9use std::{
10    num::{NonZeroU64, NonZeroUsize},
11    time::Duration,
12};
13use tokio::time::Instant;
14
15/// Native client timer range: one full period of Tokio's six-level, six-bit millisecond wheel. This is a supported range, not a deployment timeout default.
16pub(super) const MAX_TIMER_MS: u64 = (1 << 36) - 1;
17
18/// Required caller budgets. No capacity or deployment timeout is inferred from SQL settings.
19#[derive(Clone, Debug)]
20pub struct HttpNotificationOptions {
21    pub max_channels: usize,
22    pub max_queued_events: usize,
23    pub max_queued_bytes: usize,
24    /// Bounds the one HTTP data chunk retained by the worker, separately from the per-frame wire limit.
25    pub max_transport_chunk_bytes: usize,
26    pub connect_timeout: Duration,
27    /// Bounds the entire initial attempt through ready, including connection and registration.
28    pub ready_timeout: Duration,
29    /// Reject a server's incompatible idle budget before exposing ready.
30    pub max_idle_timeout: Duration,
31    /// None explicitly disables reconnection. A supplied policy manages each post-ready loss episode.
32    pub retry: Option<NotificationRetryOptions>,
33}
34
35#[derive(Clone, Debug)]
36pub struct NotificationRetryOptions {
37    /// Replacement attempts in one loss episode; the initial ready stream is not an attempt in this count.
38    pub max_attempts: u32,
39    pub episode_timeout: Duration,
40    pub initial_backoff: Duration,
41    pub max_backoff: Duration,
42    pub max_retry_after: Duration,
43}
44
45impl HttpNotificationOptions {
46    pub(super) fn validate(&self) -> Result<(), HttpNotificationError> {
47        if self.max_channels == 0
48            || self.max_queued_events == 0
49            || self.max_queued_bytes == 0
50            || self.max_transport_chunk_bytes < MAX_NOTIFICATION_WIRE_BYTES
51        {
52            return Err(HttpNotificationError::invalid_options());
53        }
54        for duration in [
55            self.connect_timeout,
56            self.ready_timeout,
57            self.max_idle_timeout,
58        ] {
59            validate_duration(duration)?;
60        }
61        if self.connect_timeout > self.ready_timeout {
62            return Err(HttpNotificationError::invalid_options());
63        }
64        if let Some(retry) = &self.retry {
65            if retry.max_attempts == 0 || retry.initial_backoff > retry.max_backoff {
66                return Err(HttpNotificationError::invalid_options());
67            }
68            for duration in [
69                retry.episode_timeout,
70                retry.initial_backoff,
71                retry.max_backoff,
72                retry.max_retry_after,
73            ] {
74                validate_duration(duration)?;
75            }
76        }
77        Ok(())
78    }
79
80    pub(super) fn channel_limit(&self) -> NonZeroUsize {
81        NonZeroUsize::new(self.max_channels).expect("validated channel limit")
82    }
83
84    pub(super) fn timer_limits(&self) -> TimerLimits {
85        TimerLimits::new(
86            NonZeroU64::new(MAX_TIMER_MS).unwrap(),
87            NonZeroU64::new(self.max_idle_timeout.as_millis() as u64),
88        )
89    }
90}
91
92fn validate_duration(duration: Duration) -> Result<(), HttpNotificationError> {
93    if duration.is_zero()
94        || duration.as_millis() == 0
95        || duration.as_millis() > u128::from(MAX_TIMER_MS)
96        || !duration.subsec_nanos().is_multiple_of(1_000_000)
97        || Instant::now().checked_add(duration).is_none()
98    {
99        return Err(HttpNotificationError::invalid_options());
100    }
101    Ok(())
102}
103
104pub(super) fn deadline(duration: Duration) -> Result<Instant, HttpNotificationError> {
105    Instant::now()
106        .checked_add(duration)
107        .ok_or_else(HttpNotificationError::invalid_options)
108}