pb_mapper_client/sdk/admin/
mod.rs1mod 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;
20use 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
39const COLLECT_PAGE_SIZE: u16 = MAX_PAGE_SIZE;
48
49const MAX_PAGES: u32 = 10_000;
54
55#[derive(Clone)]
57pub struct Admin {
58 pub(crate) inner: Arc<ClientInner>,
59}
60
61macro_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
85macro_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 pub async fn request(&self, request: AdminRequest) -> Result<AdminResponse> {
108 self.request_with_timeout(request, control_io_timeout())
109 .await
110 }
111
112 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 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 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 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 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_key(key_id: u64) -> KeyMetadata,
153 request: AdminRequest::KeyRevoke { key_id },
154 response: KeyRevoked(meta) => KeyMetadata::from(meta)
155 );
156
157 admin_rpc!(
158 auth_status() -> AuthStatusInfo,
160 request: AdminRequest::AuthStatus,
161 response: AuthStatus(status) => AuthStatusInfo::from(status)
162 );
163
164 admin_rpc_struct!(
165 gc_keys() -> u64,
167 request: AdminRequest::KeyGc,
168 response: KeyGc { removed } => removed
169 );
170
171 admin_rpc_struct!(
172 reset_auth_state() -> (),
175 request: AdminRequest::AuthStateReset { confirm: true },
176 response: Ok { .. } => ()
177 );
178
179 admin_rpc_struct!(
180 set_legacy_protocol(policy: LegacyProtocol) -> (),
182 request: AdminRequest::LegacyProtocolSet { policy: policy.into() },
183 response: Ok { .. } => ()
184 );
185
186 admin_rpc!(
187 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 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 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 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 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 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 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 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 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 other => unexpected::<String>("Ok", &other).map_err(preserve),
295 }
296 }
297
298 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 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 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
356fn 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
382async 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 #[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 #[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 #[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}