1use std::collections::{HashMap, HashSet, VecDeque};
51use std::fmt;
52use std::fs::{File, OpenOptions};
53use std::io::{Read, Write};
54#[cfg(unix)]
55use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
56use std::path::{Path, PathBuf};
57use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering};
58use std::sync::{Arc, Weak};
59use std::time::{Duration, SystemTime, UNIX_EPOCH};
60
61use parking_lot::{Mutex, RwLock};
62use rand::RngExt;
63use ring::aead::{AES_256_GCM, Aad, LessSafeKey, Nonce, UnboundKey};
64use ring::hkdf::{HKDF_SHA256, Salt};
65use serde::{Deserialize, Serialize};
66use subtle::ConstantTimeEq;
67use tokio::sync::{mpsc, oneshot};
68use tokio_util::sync::CancellationToken;
69
70use pb_mapper_core::checksum::{
71 AesKeyType, Credential, ENV_MSG_HEADER_KEY, MACHINE_MSG_HEADER_KEY_PATH,
72 encode_temporary_credential, env_safe_admin_key_error, get_process_credential,
73 is_env_safe_admin_key, parse_credential, set_process_msg_header_key,
74};
75
76pub const ADMIN_NAMESPACE: u64 = ADMIN_KEY_ID.as_u64();
79pub const DEFAULT_AUTH_STATE_DIR: &str = "/var/lib/pb-mapper/auth";
80pub const DEFAULT_TEMP_KEY_CAPACITY: usize = 65_536;
81pub const MAX_TEMP_KEY_CAPACITY: usize = 1_048_576;
82pub const DEFAULT_MAX_TEMP_KEY_TTL: Duration = Duration::from_secs(30 * 24 * 60 * 60);
83pub const MIN_TEMP_KEY_TTL: Duration = Duration::from_secs(10);
84pub const MAX_TEMP_KEY_TTL: Duration = Duration::from_secs(365 * 24 * 60 * 60);
85const TOMBSTONE_RETENTION: Duration = Duration::from_secs(60);
86const MAX_SCHEDULABLE_DELAY: Duration =
89 Duration::from_secs(MAX_TEMP_KEY_TTL.as_secs() + TOMBSTONE_RETENTION.as_secs());
90const SNAPSHOT_COMPACTION_INTERVAL: Duration = Duration::from_secs(5 * 60);
91const SNAPSHOT_SCHEMA_VERSION: u16 = 1;
92const STATE_BLOB_MAGIC: &[u8; 5] = b"PBAS1";
93const STATE_AAD: &[u8] = b"pb-mapper-auth-state-v1";
94const INSTANCE_ID_LEN: usize = 16;
95const ADMIN_REPLAY_RETENTION: Duration = Duration::from_secs(10 * 60);
96const ADMIN_REPLAY_CAPACITY: usize = 65_536;
97const AUDIT_RECORD_CAPACITY: usize = 4096;
98
99#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
100#[serde(rename_all = "snake_case")]
101pub enum LegacyProtocolPolicy {
102 Allow,
103 Deny,
104}
105
106impl LegacyProtocolPolicy {
107 pub fn is_allowed(self) -> bool {
108 matches!(self, Self::Allow)
109 }
110}
111
112#[derive(Clone, Debug)]
113pub struct AuthConfig {
114 pub state_dir: PathBuf,
115 pub max_temporary_keys: usize,
116 pub max_temporary_key_ttl: Duration,
117 pub legacy_protocol: LegacyProtocolPolicy,
118}
119
120#[derive(Clone, Debug, Serialize, Deserialize)]
121pub struct AuthFailure {
122 pub code: String,
123 pub message: String,
124 pub retryable: bool,
125}
126
127impl AuthFailure {
128 pub fn new(code: impl Into<String>, message: impl Into<String>, retryable: bool) -> Self {
129 Self {
130 code: code.into(),
131 message: message.into(),
132 retryable,
133 }
134 }
135
136 pub fn internal(message: impl Into<String>) -> Self {
137 Self::new("auth_internal_error", message, false)
138 }
139}
140
141impl fmt::Display for AuthFailure {
142 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
143 write!(formatter, "{}: {}", self.code, self.message)
144 }
145}
146
147impl std::error::Error for AuthFailure {}
148
149const LEASE_CANCEL_NONE: u8 = 0;
150const LEASE_CANCEL_EXPIRED: u8 = 1;
151const LEASE_CANCEL_REVOKED: u8 = 2;
152const LEASE_CANCEL_ROTATED: u8 = 3;
153
154#[derive(Debug)]
155pub struct AuthLease {
156 key_id: KeyId,
157 expires_at: AtomicU64,
158 cancellation: CancellationToken,
159 cancel_reason: AtomicU8,
160}
161
162impl AuthLease {
163 fn new(key_id: KeyId, expires_at: u64) -> Self {
164 Self {
165 key_id,
166 expires_at: AtomicU64::new(expires_at),
167 cancellation: CancellationToken::new(),
168 cancel_reason: AtomicU8::new(LEASE_CANCEL_NONE),
169 }
170 }
171
172 pub fn key_id(&self) -> KeyId {
173 self.key_id
174 }
175
176 pub fn expires_at(&self) -> u64 {
177 self.expires_at.load(Ordering::Acquire)
178 }
179
180 pub fn cancellation_token(&self) -> CancellationToken {
181 self.cancellation.clone()
182 }
183
184 fn record_cancel(&self, reason: u8) {
185 let _ = self.cancel_reason.compare_exchange(
186 LEASE_CANCEL_NONE,
187 reason,
188 Ordering::AcqRel,
189 Ordering::Acquire,
190 );
191 self.cancellation.cancel();
192 }
193
194 pub(crate) fn cancel_expired(&self) {
195 self.record_cancel(LEASE_CANCEL_EXPIRED);
196 }
197
198 pub(crate) fn cancel_revoked(&self) {
199 self.record_cancel(LEASE_CANCEL_REVOKED);
200 }
201
202 pub(crate) fn cancel_rotated(&self) {
203 self.record_cancel(LEASE_CANCEL_ROTATED);
204 }
205
206 #[cfg(test)]
207 pub(crate) fn expire_now(&self) {
208 self.expires_at.store(0, Ordering::Release);
209 }
210}
211
212#[derive(Clone, Debug)]
213pub struct AuthContext {
214 pub key_id: KeyId,
215 pub namespace: u64,
216 pub is_admin: bool,
217 lease: Weak<AuthLease>,
218}
219
220impl AuthContext {
221 fn from_lease(key_id: KeyId, is_admin: bool, lease: &Arc<AuthLease>) -> Self {
222 Self {
223 key_id,
224 namespace: if is_admin {
225 ADMIN_NAMESPACE
226 } else {
227 key_id.as_u64()
228 },
229 is_admin,
230 lease: Arc::downgrade(lease),
231 }
232 }
233
234 pub fn ensure_active(&self) -> Result<Arc<AuthLease>, AuthFailure> {
235 let lease = self.lease.upgrade().ok_or_else(|| {
236 AuthFailure::new(
237 if self.is_admin {
238 "administrator_key_rotated"
239 } else {
240 "temporary_key_inactive"
241 },
242 "credential lease is no longer active",
243 false,
244 )
245 })?;
246 if lease.cancellation.is_cancelled() {
247 return Err(cancelled_lease_failure(self.is_admin, &lease));
248 }
249 if !self.is_admin && lease.expires_at() <= unix_seconds() {
250 lease.cancel_expired();
251 return Err(AuthFailure::new(
252 "temporary_key_expired",
253 "temporary key has expired",
254 false,
255 ));
256 }
257 Ok(lease)
258 }
259
260 pub fn cancellation_token(&self) -> Result<CancellationToken, AuthFailure> {
261 Ok(self.ensure_active()?.cancellation_token())
262 }
263
264 pub fn admin_cancellation_token(&self) -> Result<CancellationToken, AuthFailure> {
265 self.require_admin()?;
266 self.cancellation_token()
267 }
268
269 fn admin_authority(&self) -> Result<Weak<AuthLease>, AuthFailure> {
270 self.require_admin()?;
271 self.ensure_active()?;
272 Ok(self.lease.clone())
273 }
274
275 fn require_admin(&self) -> Result<(), AuthFailure> {
276 if self.is_admin {
277 Ok(())
278 } else {
279 Err(AuthFailure::new(
280 "admin_permission_required",
281 "administrator credential is required for this operation",
282 false,
283 ))
284 }
285 }
286}
287
288fn cancelled_lease_failure(is_admin: bool, lease: &AuthLease) -> AuthFailure {
289 if is_admin {
290 return AuthFailure::new(
291 "administrator_key_rotated",
292 "credential lease has been cancelled",
293 false,
294 );
295 }
296 match lease.cancel_reason.load(Ordering::Acquire) {
297 LEASE_CANCEL_EXPIRED => {
298 AuthFailure::new("temporary_key_expired", "temporary key has expired", false)
299 }
300 LEASE_CANCEL_ROTATED => AuthFailure::new(
301 "temporary_key_rotated",
302 "temporary credential was invalidated by administrator root rotation or auth-state reset",
303 false,
304 ),
305 LEASE_CANCEL_REVOKED => {
306 AuthFailure::new("temporary_key_revoked", "temporary key was revoked", false)
307 }
308 _ => AuthFailure::new(
309 "temporary_key_inactive",
310 "credential lease has been cancelled",
311 false,
312 ),
313 }
314}
315
316#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
386#[serde(rename_all = "snake_case")]
387enum SlotState {
388 Free,
390 Active,
392 Expired,
394 Revoked,
395}
396
397#[derive(Debug)]
399struct SlotHot {
400 generation: Generation,
401 state: SlotState,
402 expires_at: u64,
403 lease: Weak<AuthLease>,
406}
407
408impl SlotHot {
409 fn holds(&self, key_id: KeyId) -> bool {
413 self.generation == key_id.generation() && self.state != SlotState::Free
414 }
415
416 fn is_collectable(&self, now: u64) -> bool {
419 match self.state {
420 SlotState::Expired | SlotState::Revoked => true,
421 SlotState::Active => self.expires_at <= now,
422 SlotState::Free => false,
423 }
424 }
425
426 fn retire(&mut self) {
429 *self = Self {
430 generation: self.generation,
431 ..Self::default()
432 };
433 }
434}
435
436impl Default for SlotHot {
437 fn default() -> Self {
438 Self {
439 generation: Generation::FIRST,
440 state: SlotState::Free,
441 expires_at: 0,
442 lease: Weak::new(),
443 }
444 }
445}
446
447#[derive(Debug)]
448struct AdminState {
449 key: AesKeyType,
450 lease: Weak<AuthLease>,
451}
452
453#[derive(Clone, Debug)]
454struct PreviousRoot {
455 admin_key: AesKeyType,
456 instance_id: [u8; INSTANCE_ID_LEN],
457}
458
459#[derive(Debug)]
460struct AuthStateInner {
461 admin: RwLock<AdminState>,
462 sync_process_credential: bool,
463 instance_id: RwLock<[u8; INSTANCE_ID_LEN]>,
464 slots: RwLock<Box<[SlotHot]>>,
467 high_slot_generations: RwLock<Vec<Generation>>,
480 high_slot_entries: RwLock<Vec<PersistedEntry>>,
481 cold: RwLock<HashMap<KeyId, ColdMetadata>>,
485 safe_mode: AtomicBool,
486 legacy_protocol_allowed: AtomicBool,
487 active_legacy_connections: AtomicU64,
488 last_legacy_connection_at: AtomicU64,
489 auth_successes: AtomicU64,
490 auth_failures: AtomicU64,
491 root_epoch: AtomicU64,
492 previous_root: RwLock<Option<PreviousRoot>>,
493 audit_records: RwLock<VecDeque<AuditRecord>>,
494}
495
496impl AuthStateInner {
497 fn high(&self) -> parking_lot::RwLockReadGuard<'_, Vec<PersistedEntry>> {
500 self.high_slot_entries.read()
501 }
502
503 fn high_mut(&self) -> parking_lot::RwLockWriteGuard<'_, Vec<PersistedEntry>> {
504 self.high_slot_entries.write()
505 }
506
507 fn slots(&self) -> parking_lot::RwLockReadGuard<'_, Box<[SlotHot]>> {
508 self.slots.read()
509 }
510
511 fn slots_mut(&self) -> parking_lot::RwLockWriteGuard<'_, Box<[SlotHot]>> {
512 self.slots.write()
513 }
514
515 fn cold(&self) -> parking_lot::RwLockReadGuard<'_, HashMap<KeyId, ColdMetadata>> {
516 self.cold.read()
517 }
518
519 fn cold_mut(&self) -> parking_lot::RwLockWriteGuard<'_, HashMap<KeyId, ColdMetadata>> {
520 self.cold.write()
521 }
522
523 fn admin_key(&self) -> AesKeyType {
524 self.admin.read().key
525 }
526
527 fn instance_id(&self) -> [u8; INSTANCE_ID_LEN] {
528 *self.instance_id.read()
529 }
530}
531
532#[derive(Clone)]
533pub struct AuthRuntime {
534 inner: Weak<AuthStateInner>,
535 command_tx: mpsc::Sender<AuthCommand>,
536 config: AuthConfig,
537 _state_lock: Arc<File>,
538 actor: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>,
539 actor_abort: tokio::task::AbortHandle,
540}
541
542#[derive(Clone, Debug, Serialize, Deserialize)]
543pub struct TemporaryKeyMetadata {
544 pub key_id: KeyId,
545 pub state: String,
546 pub issued_at: u64,
547 pub expires_at: u64,
548 pub label: Option<String>,
549}
550
551#[derive(Clone, Debug, Serialize, Deserialize)]
552pub struct IssuedTemporaryKey {
553 #[serde(flatten)]
554 pub metadata: TemporaryKeyMetadata,
555 pub credential: String,
556}
557
558#[derive(Clone, Debug, Serialize, Deserialize)]
559pub struct KeyPage {
560 pub schema_version: u16,
561 pub items: Vec<TemporaryKeyMetadata>,
562 pub next_page: Option<u32>,
563}
564
565#[derive(Clone, Debug, Serialize, Deserialize)]
566pub struct AuthStatus {
567 pub schema_version: u16,
568 pub safe_mode: bool,
569 pub capacity: usize,
570 pub active_keys: usize,
571 pub expired_keys: usize,
572 pub revoked_keys: usize,
573 pub legacy_protocol: LegacyProtocolPolicy,
574 pub active_legacy_connections: u64,
575 pub last_legacy_connection_at: Option<u64>,
576 pub auth_successes: u64,
577 pub auth_failures: u64,
578 pub server_instance_id: String,
579}
580
581#[derive(Clone, Debug)]
582struct ColdMetadata {
583 issued_at: u64,
584 label: Option<String>,
585 tombstoned_at: u64,
586}
587
588enum AuthCommand {
589 ClaimAdminMutation {
590 authority: Weak<AuthLease>,
591 fingerprint: [u8; 32],
592 client_timestamp: u64,
593 response: oneshot::Sender<Result<(), AuthFailure>>,
594 },
595 Issue {
596 authority: Weak<AuthLease>,
597 ttl: Duration,
598 label: Option<String>,
599 response: oneshot::Sender<Result<IssuedTemporaryKey, AuthFailure>>,
600 },
601 List {
602 authority: Weak<AuthLease>,
603 page: u32,
604 page_size: u16,
605 response: oneshot::Sender<Result<KeyPage, AuthFailure>>,
606 },
607 Show {
608 authority: Weak<AuthLease>,
609 key_id: KeyId,
610 reveal: bool,
611 response: oneshot::Sender<Result<IssuedTemporaryKey, AuthFailure>>,
612 },
613 Renew {
614 authority: Weak<AuthLease>,
615 key_id: KeyId,
616 ttl: Duration,
617 response: oneshot::Sender<Result<IssuedTemporaryKey, AuthFailure>>,
618 },
619 Revoke {
620 authority: Weak<AuthLease>,
621 key_id: KeyId,
622 response: oneshot::Sender<Result<TemporaryKeyMetadata, AuthFailure>>,
623 },
624 Gc {
625 authority: Weak<AuthLease>,
626 response: oneshot::Sender<Result<u64, AuthFailure>>,
627 },
628 Reset {
629 authority: Weak<AuthLease>,
630 response: oneshot::Sender<Result<(), AuthFailure>>,
631 },
632 RotateRoot {
633 authority: Weak<AuthLease>,
634 new_key: AesKeyType,
635 response: oneshot::Sender<Result<(), AuthFailure>>,
636 },
637 SetLegacyProtocol {
638 authority: Weak<AuthLease>,
639 policy: LegacyProtocolPolicy,
640 response: oneshot::Sender<Result<(), AuthFailure>>,
641 },
642 Status {
643 authority: Weak<AuthLease>,
644 response: oneshot::Sender<Result<AuthStatus, AuthFailure>>,
645 },
646 Audit {
647 authority: Weak<AuthLease>,
648 action: String,
649 key_id: Option<KeyId>,
650 detail: Option<String>,
651 response: oneshot::Sender<Result<(), AuthFailure>>,
652 },
653 Shutdown {
654 response: oneshot::Sender<()>,
655 },
656}
657
658mod config;
659pub use config::default_auth_state_dir;
660#[cfg(all(test, not(any(windows, target_os = "macos"))))]
661pub(crate) use config::linux_default_auth_state_dir;
662#[cfg(test)]
663pub(crate) use config::parse_legacy_protocol_policy;
664#[cfg(test)]
665pub(crate) use config::platform_default_auth_state_dir;
666#[cfg(all(test, not(any(windows, target_os = "macos"))))]
667pub(crate) use config::{linux_system_auth_dir_usable, unix_effective_uid};
668mod keys;
669pub use keys::derive_temporary_key;
670#[cfg(test)]
671pub(crate) use keys::recover_admin_key_after_rotation;
672pub(crate) use keys::{load_isolated_server_admin_credential, load_server_admin_credential};
673mod runtime;
674
675pub struct LegacyConnectionGuard {
676 inner: Weak<AuthStateInner>,
677}
678
679impl Drop for LegacyConnectionGuard {
680 fn drop(&mut self) {
681 if let Some(inner) = self.inner.upgrade() {
682 inner
683 .active_legacy_connections
684 .fetch_sub(1, Ordering::AcqRel);
685 }
686 }
687}
688
689#[derive(Clone, Debug, Serialize, Deserialize)]
690struct PersistedEntry {
691 key_id: KeyId,
692 state: SlotState,
693 issued_at: u64,
694 expires_at: u64,
695 label: Option<String>,
696 #[serde(default)]
697 tombstoned_at: Option<u64>,
698}
699
700#[derive(Clone, Debug, Serialize, Deserialize)]
701struct PersistedSnapshot {
702 schema_version: u16,
703 instance_id: [u8; INSTANCE_ID_LEN],
704 generations: Vec<Generation>,
705 entries: Vec<PersistedEntry>,
706 legacy_protocol: LegacyProtocolPolicy,
707 #[serde(default)]
708 admin_replays: Vec<AdminReplayRecord>,
709 #[serde(default)]
710 audit_records: VecDeque<AuditRecord>,
711 #[serde(default)]
712 root_epoch: u64,
713}
714
715#[derive(Clone, Debug, Serialize, Deserialize)]
716struct AdminReplayRecord {
717 fingerprint: [u8; 32],
718 client_timestamp: u64,
719 #[serde(default)]
722 accepted_at: u64,
723}
724
725impl AdminReplayRecord {
726 fn within_retention(&self, now: u64) -> bool {
727 let anchor = if self.accepted_at == 0 {
728 self.client_timestamp
729 } else {
730 self.accepted_at
731 };
732 now.saturating_sub(anchor) <= ADMIN_REPLAY_RETENTION.as_secs()
733 }
734}
735
736#[derive(Clone, Debug, Serialize, Deserialize)]
737enum StateMutation {
738 Issue(PersistedEntry),
739 Renew { key_id: KeyId, expires_at: u64 },
740 Revoke { key_id: KeyId, at: u64 },
741 LegacyProtocol(LegacyProtocolPolicy),
742}
743
744#[derive(Clone, Debug, Serialize, Deserialize)]
745struct AuditRecord {
746 at: u64,
747 action: String,
748 key_id: Option<KeyId>,
749 label: Option<String>,
750}
751
752#[derive(Clone, Debug, Serialize, Deserialize)]
753enum WalRecord {
754 Mutation {
755 mutation: StateMutation,
756 audit: AuditRecord,
757 },
758 Audit(AuditRecord),
759 AdminReplay(AdminReplayRecord),
760}
761
762mod actor;
763use actor::{AuthActorState, run_auth_actor};
764mod persistence;
765pub use persistence::*;
766pub(crate) use persistence::{
767 append_audit, append_mutation, append_wal, atomic_write, auth_snapshot_path, build_snapshot,
768 cancel_all_temporary_leases, compaction_is_allowed, empty_snapshot,
769 fail_closed_on_uncertain_wal, hex, key_matches_existing_state, load_or_create_instance_id,
770 load_persisted_state, normalize_tombstone_times, open_blob, prepare_state_dir_and_lock,
771 push_audit_record, push_persisted_audit, random_instance_id, recover_instance_id_after_reset,
772 reset_already_installed, rotation_already_installed, split_high_slot_state, truncate_auth_wal,
773 unix_seconds, write_admin_key, write_snapshot_and_truncate_wal,
774};
775#[cfg(test)]
776pub(crate) use persistence::{prepare_state_dir, read_instance_id_file, try_load_persisted_state};
777mod ids;
778pub use ids::{ADMIN_KEY_ID, Generation, KeyId, SlotIndex};
779mod leases;
780use leases::Leases;
781mod timing_wheel;
782use timing_wheel::{Timer, TimingWheel};
783#[cfg(test)]
784mod tests;