Skip to main content

pb_mapper_client/sdk/admin/
mod.rs

1//! Administrator RPCs: the credential lifecycle, the relay's auth state, and
2//! the service and connection inventories.
3//!
4//! Split three ways — this file is the RPC surface, [`types`] holds the owned
5//! response types and their wire conversions, and [`transport`] owns the
6//! single-shot session each request runs over.
7
8mod transport;
9mod types;
10
11use std::path::Path;
12use std::sync::Arc;
13use std::time::Duration;
14
15use pb_mapper_auth::{
16    discard_staged_admin_key, generate_admin_key, stage_admin_key_candidate, write_admin_key_file,
17};
18use pb_mapper_core::checksum::parse_credential;
19use pb_mapper_core::config::control_io_timeout;
20// The largest page the relay will serve, taken from the relay's own definition
21// rather than restated here: an oversized request is rejected locally so it fails
22// with a clear message instead of a protocol error, and that check has to agree
23// with what the relay clamps to.
24use pb_mapper_core::paging::MAX_PAGE_SIZE;
25use pb_mapper_protocol::command::{AdminRequest, AdminResponse};
26
27use self::transport::send_admin_request;
28use self::types::Paged;
29use super::Error;
30use super::client::ClientInner;
31use super::error::Result;
32use super::types::LegacyProtocol;
33
34pub use self::types::{
35    AuthStatusInfo, ConnectionInfo, ConnectionPage, IssuedKey, KeyListPage, KeyMetadata,
36    ServiceInfo, ServicePage,
37};
38
39/// Page size used by the `*_all` helpers, which page on the caller's behalf.
40///
41/// The largest page the relay serves, not a smaller round number, for two
42/// reasons. A relay configured up to `MAX_TEMP_KEY_CAPACITY` (1,048,576) holds
43/// more credentials than [`MAX_PAGES`] pages of 100 could carry, so `*_all`
44/// would fail on a full inventory instead of returning it. And the relay
45/// re-collects and re-sorts its whole table for every page it serves, so the
46/// page count is what that work is multiplied by.
47const COLLECT_PAGE_SIZE: u16 = MAX_PAGE_SIZE;
48
49/// A cap on `*_all` pagination, so a relay that keeps handing back a `next_page`
50/// cursor cannot spin the caller forever.
51///
52/// Ample at [`COLLECT_PAGE_SIZE`]: the relay's hard capacity needs 1,049 pages.
53const MAX_PAGES: u32 = 10_000;
54
55/// Administrator RPCs. Constructed via [`super::Client::admin`].
56#[derive(Clone)]
57pub struct Admin {
58    pub(crate) inner: Arc<ClientInner>,
59}
60
61/// Declares an RPC that sends one request and expects one response variant.
62///
63/// Most administrator calls are exactly that: build the request, match the one
64/// response that answers it, and convert. Written out, each is five lines of
65/// which only two say anything, and the `unexpected` arm is easy to get subtly
66/// wrong — naming the wrong expected variant in the error text.
67macro_rules! admin_rpc {
68    (
69        $(#[$doc:meta])*
70        $name:ident($($arg:ident: $arg_ty:ty),* $(,)?)
71            -> $output:ty,
72        request: $request:expr,
73        response: $variant:ident($binding:pat) => $value:expr
74    ) => {
75        $(#[$doc])*
76        pub async fn $name(&self, $($arg: $arg_ty),*) -> Result<$output> {
77            match self.request($request).await? {
78                AdminResponse::$variant($binding) => Ok($value),
79                other => unexpected(stringify!($variant), &other),
80            }
81        }
82    };
83}
84
85/// The struct-variant form of [`admin_rpc`], for the two responses that carry
86/// named fields rather than a payload.
87macro_rules! admin_rpc_struct {
88    (
89        $(#[$doc:meta])*
90        $name:ident($($arg:ident: $arg_ty:ty),* $(,)?)
91            -> $output:ty,
92        request: $request:expr,
93        response: $variant:ident { $($binding:tt)* } => $value:expr
94    ) => {
95        $(#[$doc])*
96        pub async fn $name(&self, $($arg: $arg_ty),*) -> Result<$output> {
97            match self.request($request).await? {
98                AdminResponse::$variant { $($binding)* } => Ok($value),
99                other => unexpected(stringify!($variant), &other),
100            }
101        }
102    };
103}
104
105impl Admin {
106    /// One-shot administrator RPC. CLI output rendering can call this directly.
107    pub async fn request(&self, request: AdminRequest) -> Result<AdminResponse> {
108        self.request_with_timeout(request, control_io_timeout())
109            .await
110    }
111
112    /// [`Self::request`] with an explicit I/O bound, covering the connect as well
113    /// as the exchange.
114    pub async fn request_with_timeout(
115        &self,
116        request: AdminRequest,
117        io_timeout: Duration,
118    ) -> Result<AdminResponse> {
119        send_admin_request(&self.inner.server, self.credential(), request, io_timeout).await
120    }
121
122    admin_rpc!(
123        /// Mint a temporary credential valid for `ttl`.
124        issue_key(ttl: Duration, label: Option<String>) -> IssuedKey,
125        request: AdminRequest::KeyIssue { ttl_seconds: ttl.as_secs(), label },
126        response: KeyIssued(issued) => IssuedKey::from(issued)
127    );
128
129    admin_rpc!(
130        /// Metadata for one credential, without its secret.
131        show_key(key_id: u64) -> IssuedKey,
132        request: AdminRequest::KeyShow { key_id },
133        response: KeyShown(issued) => IssuedKey::from(issued)
134    );
135
136    admin_rpc!(
137        /// Metadata for one credential, *with* its secret.
138        reveal_key(key_id: u64) -> IssuedKey,
139        request: AdminRequest::KeyReveal { key_id },
140        response: KeyShown(issued) => IssuedKey::from(issued)
141    );
142
143    admin_rpc!(
144        /// Extend a credential's lifetime to `ttl` from now.
145        renew_key(key_id: u64, ttl: Duration) -> IssuedKey,
146        request: AdminRequest::KeyRenew { key_id, ttl_seconds: ttl.as_secs() },
147        response: KeyRenewed(issued) => IssuedKey::from(issued)
148    );
149
150    admin_rpc!(
151        /// Revoke a credential, taking effect on the relay immediately.
152        revoke_key(key_id: u64) -> KeyMetadata,
153        request: AdminRequest::KeyRevoke { key_id },
154        response: KeyRevoked(meta) => KeyMetadata::from(meta)
155    );
156
157    admin_rpc!(
158        /// The relay's credential-subsystem snapshot.
159        auth_status() -> AuthStatusInfo,
160        request: AdminRequest::AuthStatus,
161        response: AuthStatus(status) => AuthStatusInfo::from(status)
162    );
163
164    admin_rpc_struct!(
165        /// Drop expired and revoked credentials, returning how many were removed.
166        gc_keys() -> u64,
167        request: AdminRequest::KeyGc,
168        response: KeyGc { removed } => removed
169    );
170
171    admin_rpc_struct!(
172        /// Erase the relay's credential state. Every temporary credential stops
173        /// authenticating; the administrator key is unaffected.
174        reset_auth_state() -> (),
175        request: AdminRequest::AuthStateReset { confirm: true },
176        response: Ok { .. } => ()
177    );
178
179    admin_rpc_struct!(
180        /// Allow or deny pre-v2 framing on new connections.
181        set_legacy_protocol(policy: LegacyProtocol) -> (),
182        request: AdminRequest::LegacyProtocolSet { policy: policy.into() },
183        response: Ok { .. } => ()
184    );
185
186    admin_rpc!(
187        /// One page of temporary credentials.
188        list_keys(page: u32, page_size: u16) -> KeyListPage,
189        request: AdminRequest::KeyList { page, page_size: validate_page_size(page_size)? },
190        response: KeyList(page) => KeyListPage::from(page)
191    );
192
193    admin_rpc!(
194        /// One page of registered services, optionally scoped to one credential.
195        list_services(key_id: Option<u64>, page: u32, page_size: u16) -> ServicePage,
196        request: AdminRequest::ServiceList {
197            key_id,
198            page,
199            page_size: validate_page_size(page_size)?,
200        },
201        response: Services(page) => ServicePage::from(page)
202    );
203
204    admin_rpc!(
205        /// One page of live connections, optionally scoped to one credential.
206        list_connections(key_id: Option<u64>, page: u32, page_size: u16) -> ConnectionPage,
207        request: AdminRequest::ConnectionList {
208            key_id,
209            page,
210            page_size: validate_page_size(page_size)?,
211        },
212        response: Connections(page) => ConnectionPage::from(page)
213    );
214
215    admin_rpc_struct!(
216        /// Drop registered control connections the relay is still holding for a
217        /// service, returning how many it dropped.
218        ///
219        /// `conn_id` of `None` retires every connection the service has, which is
220        /// what frees a connection quota filled by connections that should have
221        /// gone away. A registration whose client is still healthy simply
222        /// reconnects, so this is a nudge rather than a shutdown.
223        retire_connections(key_id: Option<u64>, service_name: String, conn_id: Option<u32>) -> u32,
224        request: AdminRequest::ConnectionRetire { key_id, service_name, conn_id },
225        response: ConnectionsRetired { retired } => retired
226    );
227
228    /// Every temporary credential, paging until the relay stops handing back a
229    /// cursor.
230    pub async fn list_keys_all(&self) -> Result<Vec<KeyMetadata>> {
231        collect_pages(|page| self.list_keys(page, COLLECT_PAGE_SIZE)).await
232    }
233
234    /// Every registered service, optionally scoped to one credential.
235    pub async fn list_services_all(&self, key_id: Option<u64>) -> Result<Vec<ServiceInfo>> {
236        collect_pages(|page| self.list_services(key_id, page, COLLECT_PAGE_SIZE)).await
237    }
238
239    /// Every live connection, optionally scoped to one credential.
240    pub async fn list_connections_all(&self, key_id: Option<u64>) -> Result<Vec<ConnectionInfo>> {
241        collect_pages(|page| self.list_connections(key_id, page, COLLECT_PAGE_SIZE)).await
242    }
243
244    /// Rotate the relay's administrator key, returning the key now in force.
245    ///
246    /// Generates a key when `new_key` is `None`. Rotation is not idempotent, so a
247    /// caller who never learns the outcome cannot retry safely: if the relay
248    /// committed the change and the response was lost, only the new key still
249    /// authenticates. When the SDK generated that key, losing it locks the
250    /// operator out — so any inconclusive result is reported as
251    /// [`Error::RootRotationUncertain`], which carries the candidate. A caller
252    /// who supplied `new_key` already holds it and gets the underlying error
253    /// unchanged.
254    ///
255    /// Prefer [`Admin::rotate_root_key_to_file`], which persists the candidate
256    /// before the request is ever sent.
257    pub async fn rotate_root_key(&self, new_key: Option<String>) -> Result<String> {
258        let caller_supplied = new_key.is_some();
259        let new_key = new_key.unwrap_or_else(generate_admin_key);
260        let parsed = parse_credential(new_key.trim()).map_err(Error::invalid_config)?;
261        if !parsed.is_admin() {
262            return Err(Error::invalid_config(
263                "root rotation requires a 32-byte administrator key",
264            ));
265        }
266        // A generated key exists nowhere but this frame, so every exit that
267        // leaves the relay's state unknown has to hand it back.
268        let preserve = |error: Error| {
269            if caller_supplied {
270                return error;
271            }
272            Error::RootRotationUncertain {
273                candidate: new_key.clone(),
274                message: error.to_string(),
275            }
276        };
277        let response = self
278            .request(AdminRequest::RootKeyRotate {
279                new_admin_key: new_key.clone(),
280            })
281            .await
282            .map_err(preserve)?;
283        match response {
284            AdminResponse::Ok { .. } => {
285                *self
286                    .inner
287                    .credential
288                    .write()
289                    .unwrap_or_else(|poisoned| poisoned.into_inner()) = parsed;
290                Ok(new_key)
291            }
292            // A response that is not `Ok` still came from the relay, but nothing
293            // here proves the rotation did not take effect.
294            other => unexpected::<String>("Ok", &other).map_err(preserve),
295        }
296    }
297
298    /// Rotate the root key and persist it to `path`.
299    ///
300    /// The candidate is staged at a sibling file before the request is sent, and
301    /// `path` is only replaced once the relay has accepted the rotation and the new
302    /// key has passed a post-rotation status check. A rotation that fails after the
303    /// relay may have committed it therefore leaves both keys on disk rather than a
304    /// `path` holding a key the relay never installed.
305    ///
306    /// An unresolved candidate from an earlier attempt is never overwritten — see
307    /// [`stage_admin_key_candidate`] — so this fails until the operator establishes
308    /// which of the two keys the relay accepts. Retrying with the same `new_key`
309    /// is allowed.
310    pub async fn rotate_root_key_to_file(
311        &self,
312        path: &Path,
313        new_key: Option<String>,
314    ) -> Result<String> {
315        let new_key = new_key.unwrap_or_else(generate_admin_key);
316        let staged_path = stage_admin_key_candidate(path, &new_key).map_err(auth_file_error)?;
317        // Every failure past this point leaves the candidate on disk, and has to
318        // say so: it may be the key the relay is now running on.
319        let staged_note = || format!("the candidate key remains at `{}`", staged_path.display());
320
321        let rotated = self.rotate_root_key(Some(new_key)).await.map_err(|error| {
322            auth_file_message(format!(
323                "root rotation request failed; {}: {error}",
324                staged_note()
325            ))
326        })?;
327        // `rotate_root_key` already swapped the in-memory credential, so this
328        // proves the relay authenticates the key we are about to persist.
329        self.auth_status().await.map_err(|error| {
330            auth_file_message(format!(
331                "new administrator key did not pass the post-rotation status check; {}: {error}",
332                staged_note()
333            ))
334        })?;
335        write_admin_key_file(path, &rotated, true).map_err(|error| {
336            auth_file_message(format!(
337                "administrator key rotated and verified, but `{}` could not be updated; recover \
338                 the key from `{}`: {error}",
339                path.display(),
340                staged_path.display()
341            ))
342        })?;
343        discard_staged_admin_key(&staged_path);
344        Ok(rotated)
345    }
346
347    fn credential(&self) -> pb_mapper_core::checksum::Credential {
348        *self
349            .inner
350            .credential
351            .read()
352            .unwrap_or_else(|poisoned| poisoned.into_inner())
353    }
354}
355
356/// Lift an administrator key-file failure into the SDK's error type.
357fn auth_file_error(error: pb_mapper_auth::AuthFailure) -> Error {
358    auth_file_message(error.to_string())
359}
360
361fn auth_file_message(message: impl Into<String>) -> Error {
362    Error::AuthFile {
363        message: message.into(),
364    }
365}
366
367fn validate_page_size(page_size: u16) -> Result<u16> {
368    if !(1..=MAX_PAGE_SIZE).contains(&page_size) {
369        return Err(Error::invalid_config(format!(
370            "page_size must be between 1 and {MAX_PAGE_SIZE}"
371        )));
372    }
373    Ok(page_size)
374}
375
376fn unexpected<T>(expected: &str, actual: &AdminResponse) -> Result<T> {
377    Err(Error::protocol(format!(
378        "expected {expected}, got {actual:?}"
379    )))
380}
381
382/// Drain a paginated listing into one vector.
383///
384/// Bounded by [`MAX_PAGES`]: the cursor comes from the relay, and a listing that
385/// never terminates would otherwise hang the caller with no way to tell why.
386async fn collect_pages<P, F, Fut>(mut fetch: F) -> Result<Vec<P::Item>>
387where
388    P: Paged,
389    F: FnMut(u32) -> Fut,
390    Fut: std::future::Future<Output = Result<P>>,
391{
392    let mut page = 0_u32;
393    let mut items = Vec::new();
394    for _ in 0..MAX_PAGES {
395        let (chunk, next) = fetch(page).await?.into_parts();
396        items.extend(chunk);
397        match next {
398            Some(next_page) => page = next_page,
399            None => return Ok(items),
400        }
401    }
402    Err(Error::protocol(format!(
403        "pagination exceeded {MAX_PAGES} pages"
404    )))
405}
406
407#[cfg(test)]
408mod tests {
409    use super::*;
410    use crate::sdk::admin::types::KeyListPage;
411
412    fn page(next_page: Option<u32>, ids: &[u64]) -> KeyListPage {
413        KeyListPage {
414            schema_version: 1,
415            items: ids
416                .iter()
417                .map(|&key_id| KeyMetadata {
418                    key_id,
419                    state: "active".into(),
420                    issued_at: 0,
421                    expires_at: 0,
422                    label: None,
423                })
424                .collect(),
425            next_page,
426        }
427    }
428
429    /// The `*_all` helpers have to be able to drain a relay filled to its hard
430    /// capacity. At a smaller page size they cannot: 1,048,576 credentials over
431    /// pages of 100 is 10,486 pages, past the [`MAX_PAGES`] guard, so a full
432    /// inventory would come back as a pagination error instead of a listing.
433    #[test]
434    fn the_collect_page_size_can_drain_a_relay_at_capacity() {
435        assert_eq!(
436            COLLECT_PAGE_SIZE, MAX_PAGE_SIZE,
437            "paging below the relay's maximum multiplies its per-page re-sort \
438             and shrinks what `*_all` can drain"
439        );
440        let pages_at_capacity =
441            pb_mapper_auth::MAX_TEMP_KEY_CAPACITY.div_ceil(usize::from(COLLECT_PAGE_SIZE));
442        assert!(
443            pages_at_capacity <= MAX_PAGES as usize,
444            "a full inventory needs {pages_at_capacity} pages, over the {MAX_PAGES} cap"
445        );
446    }
447
448    /// `validate_page_size` rejects locally, so an oversized request never
449    /// reaches the relay to come back as an opaque protocol error.
450    #[test]
451    fn a_page_size_outside_the_relays_range_is_rejected_locally() {
452        assert!(validate_page_size(0).is_err());
453        assert!(validate_page_size(MAX_PAGE_SIZE + 1).is_err());
454        assert_eq!(validate_page_size(1).unwrap(), 1);
455        assert_eq!(validate_page_size(MAX_PAGE_SIZE).unwrap(), MAX_PAGE_SIZE);
456    }
457
458    #[tokio::test]
459    async fn collect_pages_follows_the_cursor_and_concatenates() {
460        let requested = std::sync::Mutex::new(Vec::new());
461        let items = collect_pages::<KeyListPage, _, _>(|page_number| {
462            requested.lock().expect("not poisoned").push(page_number);
463            async move {
464                Ok(match page_number {
465                    0 => page(Some(7), &[1, 2]),
466                    7 => page(Some(9), &[3]),
467                    _ => page(None, &[4]),
468                })
469            }
470        })
471        .await
472        .expect("three pages is well inside the cap");
473
474        assert_eq!(
475            items.iter().map(|item| item.key_id).collect::<Vec<_>>(),
476            vec![1, 2, 3, 4]
477        );
478        assert_eq!(
479            *requested.lock().expect("not poisoned"),
480            vec![0, 7, 9],
481            "the cursor the relay hands back is what gets asked for next"
482        );
483    }
484
485    /// A relay that always hands back a cursor must not spin the caller: the
486    /// cursor is remote input, so the loop is bounded here rather than trusted.
487    #[tokio::test]
488    async fn collect_pages_gives_up_on_a_cursor_that_never_ends() {
489        let calls = std::sync::atomic::AtomicU32::new(0);
490        let error = collect_pages::<KeyListPage, _, _>(|page_number| {
491            calls.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
492            async move { Ok(page(Some(page_number + 1), &[u64::from(page_number)])) }
493        })
494        .await
495        .expect_err("an endless cursor must fail rather than hang");
496
497        assert!(error.to_string().contains("pagination exceeded"));
498        assert_eq!(
499            calls.load(std::sync::atomic::Ordering::Relaxed),
500            MAX_PAGES,
501            "the cap is what stops it, and it stops exactly there"
502        );
503    }
504}