1use std::path::{Path, PathBuf};
5use std::sync::Arc;
6use std::time::Duration;
7
8use jiff::{SignedDuration, Timestamp};
9use serde::Deserialize;
10use sha2::{Digest, Sha256};
11use subtle::ConstantTimeEq;
12use tollgate_auth::{CredentialIssuer, CredentialVerifier, HmacRegistry};
13use zeroize::Zeroizing;
14
15use crate::google::GoogleVerifier;
16use crate::security::{ControlIdentity, Role, SecurityError, SecurityPolicy, ServerSecurity};
17use crate::transport::{TlsConfig, certificate_fingerprint};
18
19#[derive(Deserialize)]
20#[serde(deny_unknown_fields)]
21struct Manifest {
22 tls: Option<TlsFiles>,
23 #[serde(default)]
24 bearers: Vec<BearerFile>,
25 #[serde(default)]
26 certificates: Vec<CertificateFile>,
27 google: Option<GoogleConfig>,
28 issuer: Option<IssuerFile>,
29}
30
31#[derive(Deserialize)]
35#[serde(deny_unknown_fields)]
36struct IssuerFile {
37 secret_file: PathBuf,
38}
39
40#[derive(Deserialize)]
41#[serde(deny_unknown_fields)]
42struct TlsFiles {
43 certificate: PathBuf,
44 private_key: PathBuf,
45 client_ca: Option<PathBuf>,
46}
47#[derive(Deserialize)]
48#[serde(deny_unknown_fields)]
49struct BearerFile {
50 identity: String,
51 role: Role,
52 token_file: PathBuf,
53}
54#[derive(Deserialize)]
55#[serde(deny_unknown_fields)]
56struct CertificateFile {
57 identity: String,
58 role: Role,
59 certificate: PathBuf,
60}
61#[derive(Deserialize)]
62#[serde(deny_unknown_fields)]
63struct GoogleConfig {
64 audience: String,
65 subjects: Vec<GoogleSubject>,
66}
67#[derive(Deserialize)]
68#[serde(deny_unknown_fields)]
69struct GoogleSubject {
70 subject: String,
71 identity: String,
72 role: Role,
73}
74
75pub struct SecurityLoader {
85 path: PathBuf,
86 keys: Option<(Vec<u8>, Timestamp)>,
87 next_key_attempt: tokio::time::Instant,
88 digest: Option<[u8; 32]>,
89 issuer: Option<Arc<HmacRegistry>>,
90 issuer_fingerprint: Option<[u8; 32]>,
91 pending_issuer: Option<Option<[u8; 32]>>,
92}
93
94pub struct LoadedSecurity {
100 policy: SecurityPolicy,
101 tls: Option<TlsConfig>,
102 issuer: Option<IssuerSecret>,
103 digest: [u8; 32],
104}
105
106struct IssuerSecret(Zeroizing<[u8; ISSUER_SECRET_LEN]>);
110
111const ISSUER_SECRET_LEN: usize = 64;
112
113impl IssuerSecret {
114 fn parse(bytes: &[u8]) -> Result<Self, SecurityError> {
115 let text = bytes
116 .strip_suffix(b"\n")
117 .map(|line| line.strip_suffix(b"\r").unwrap_or(line))
118 .unwrap_or(bytes);
119 let secret: [u8; ISSUER_SECRET_LEN] = text
120 .try_into()
121 .ok()
122 .filter(|text: &[u8; ISSUER_SECRET_LEN]| {
123 text.iter().all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f'))
124 })
125 .ok_or(SecurityError(
126 "issuer secret must be exactly 64 lowercase hexadecimal characters",
127 ))?;
128 Ok(Self(Zeroizing::new(secret)))
129 }
130
131 fn fingerprint(&self) -> [u8; 32] {
133 let mut digest = Sha256::new();
134 digest.update(b"tollgate-issuer-v1");
135 digest.update(self.0.as_slice());
136 digest.finalize().into()
137 }
138}
139
140impl SecurityLoader {
141 pub fn new(path: impl Into<PathBuf>) -> Self {
145 Self {
146 path: path.into(),
147 keys: None,
148 next_key_attempt: tokio::time::Instant::now(),
149 digest: None,
150 issuer: None,
151 issuer_fingerprint: None,
152 pending_issuer: None,
153 }
154 }
155
156 pub fn start(&mut self, loaded: LoadedSecurity) -> Result<Arc<ServerSecurity>, SecurityError> {
161 let security = ServerSecurity::new(loaded.policy, loaded.tls)?;
162 self.issuer_fingerprint = loaded.issuer.as_ref().map(IssuerSecret::fingerprint);
163 self.issuer = loaded
164 .issuer
165 .map(|secret| Arc::new(HmacRegistry::new(secret.0.as_slice())));
166 self.pending_issuer = None;
167 self.digest = Some(loaded.digest);
168 Ok(security)
169 }
170
171 pub fn install(
175 &mut self,
176 loaded: LoadedSecurity,
177 security: &ServerSecurity,
178 ) -> Result<(), SecurityError> {
179 security.replace(loaded.policy, loaded.tls)?;
180 let staged = loaded.issuer.as_ref().map(IssuerSecret::fingerprint);
181 if staged == self.issuer_fingerprint {
182 if self.pending_issuer.take().is_some() {
183 tracing::info!(
184 reason = "issuer-change-withdrawn",
185 "staged credential issuer matches the live one again"
186 );
187 }
188 } else if self.pending_issuer != Some(staged) {
189 self.pending_issuer = Some(staged);
190 tracing::warn!(
191 reason = "issuer-change-requires-restart",
192 "credential issuer unchanged; restart to apply the staged issuer"
193 );
194 }
195 self.digest = Some(loaded.digest);
196 Ok(())
197 }
198
199 pub fn issuer(&self) -> Option<Arc<dyn CredentialIssuer + Send + Sync>> {
203 self.issuer
204 .as_ref()
205 .map(|issuer| Arc::clone(issuer) as Arc<dyn CredentialIssuer + Send + Sync>)
206 }
207
208 pub fn issuer_change_pending(&self) -> bool {
211 self.pending_issuer.is_some()
212 }
213
214 pub async fn load(&mut self, now: Timestamp) -> Result<Option<LoadedSecurity>, SecurityError> {
242 self.load_with_keys(now, || {
243 crate::google::fetch_keys("https://www.googleapis.com/oauth2/v3/certs")
244 })
245 .await
246 }
247
248 async fn load_with_keys<F: Future<Output = Result<(Vec<u8>, Duration), SecurityError>>>(
249 &mut self,
250 now: Timestamp,
251 fetch: impl FnOnce() -> F,
252 ) -> Result<Option<LoadedSecurity>, SecurityError> {
253 let bytes = read(&self.path).await?;
254 let manifest: Manifest = serde_json::from_slice(&bytes)
255 .map_err(|_| SecurityError("invalid security manifest"))?;
256 let directory = self.path.parent().unwrap_or(Path::new("."));
257 let mut digest = Sha256::new();
258 record_digest(&mut digest, bytes.as_slice());
259 let mut policy = SecurityPolicy::new();
260 let mut secret = Zeroizing::new([0u8; 32]);
261 getrandom::fill(secret.as_mut())
262 .map_err(|_| SecurityError("control-plane verifier entropy unavailable"))?;
263 let registry = Arc::new(HmacRegistry::new(secret.as_ref()));
264 let mut tokens = Vec::new();
265 let mut actors = Vec::new();
266 for bearer in manifest.bearers {
267 let bytes = read(&directory.join(bearer.token_file)).await?;
268 record_digest(&mut digest, bytes.as_slice());
269 let token = Zeroizing::new(
270 std::str::from_utf8(&bytes)
271 .map_err(|_| SecurityError("invalid bearer file encoding"))?
272 .trim_end_matches(['\r', '\n'])
273 .to_owned(),
274 );
275 if token.len() < 32
276 || token.len() > 16 * 1024 - 7
277 || !token.bytes().all(|b| b.is_ascii_graphic())
278 {
279 return Err(SecurityError(
280 "static bearer credentials must contain 32..=16377 visible ASCII bytes",
281 ));
282 }
283 actors.push(ControlIdentity::new(bearer.identity, bearer.role)?);
284 tokens.push(token);
285 }
286 registry.install_credentials(tokens.iter().map(|token| token.as_bytes()));
287 let mut mapped = Vec::new();
288 for (token, identity) in tokens.iter().zip(actors) {
289 let principal = registry
290 .verify(token.as_bytes())
291 .ok_or(SecurityError("installed credential failed verification"))?
292 .principal;
293 mapped.push((principal, identity));
294 }
295 policy = policy.with_bearer(registry, mapped)?;
296 let issuer = match manifest.issuer {
297 Some(issuer) => {
298 let bytes = read(&directory.join(issuer.secret_file)).await?;
299 record_digest(&mut digest, bytes.as_slice());
300 let secret = IssuerSecret::parse(&bytes)?;
301 let collides = tokens
302 .iter()
303 .any(|token| bool::from(token.as_bytes().ct_eq(secret.0.as_slice())));
304 if collides {
305 return Err(SecurityError(
306 "issuer secret must differ from every bearer credential",
307 ));
308 }
309 Some(secret)
310 }
311 None => None,
312 };
313 for certificate in manifest.certificates {
314 let pem = read(&directory.join(certificate.certificate)).await?;
315 record_digest(&mut digest, pem.as_slice());
316 policy = policy.with_certificate(
317 certificate_fingerprint(&pem)?,
318 ControlIdentity::new(certificate.identity, certificate.role)?,
319 )?;
320 }
321 let tls = match manifest.tls {
322 Some(tls) => {
323 let certificates = read(&directory.join(tls.certificate)).await?;
324 let key = read(&directory.join(tls.private_key)).await?;
325 let ca = match tls.client_ca {
326 Some(path) => Some(read(&directory.join(path)).await?),
327 None => None,
328 };
329 record_digest(&mut digest, certificates.as_slice());
330 record_digest(&mut digest, key.as_slice());
331 if let Some(ca) = &ca {
332 record_digest(&mut digest, ca.as_slice());
333 }
334 Some(TlsConfig::from_pem(
335 &certificates,
336 &key,
337 ca.as_ref().map(|ca| ca.as_slice()),
338 )?)
339 }
340 None => None,
341 };
342 if let Some(google) = manifest.google {
343 if tokio::time::Instant::now() >= self.next_key_attempt {
344 self.next_key_attempt = tokio::time::Instant::now() + Duration::from_secs(300);
345 match fetch().await {
346 Ok((keys, lifetime)) => {
347 self.next_key_attempt = tokio::time::Instant::now()
348 + (lifetime / 2)
349 .clamp(Duration::from_secs(5), Duration::from_secs(300));
350 let until = now
351 .checked_add(
352 SignedDuration::try_from(lifetime)
353 .map_err(|_| SecurityError("signing-key lifetime overflow"))?,
354 )
355 .map_err(|_| SecurityError("signing-key expiry overflow"))?;
356 GoogleVerifier::from_jwks(&google.audience, &keys, until)?;
358 self.keys = Some((keys, until));
359 }
360 Err(error) if self.keys.is_some() => {
361 tracing::warn!(%error, "Google key refresh failed; previous expiry remains authoritative")
362 }
363 Err(error) => return Err(error),
364 }
365 }
366 let (keys, until) = self
367 .keys
368 .as_ref()
369 .ok_or(SecurityError("Google signing keys unavailable"))?;
370 record_digest(&mut digest, keys);
371 record_digest(&mut digest, &until.as_second().to_be_bytes());
372 let verifier = Arc::new(GoogleVerifier::from_jwks(&google.audience, keys, *until)?);
373 let mut mapped = Vec::new();
374 for subject in google.subjects {
375 if subject.subject.is_empty() {
376 return Err(SecurityError("Google subject must not be empty"));
377 }
378 mapped.push((
379 GoogleVerifier::principal(&subject.subject),
380 ControlIdentity::new(subject.identity, subject.role)?,
381 ));
382 }
383 policy = policy.with_bearer(verifier, mapped)?;
384 }
385 let digest = digest.finalize().into();
386 if self.digest == Some(digest) {
387 return Ok(None);
388 }
389 Ok(Some(LoadedSecurity {
390 policy,
391 tls,
392 issuer,
393 digest,
394 }))
395 }
396}
397
398async fn read(path: &Path) -> Result<Zeroizing<Vec<u8>>, SecurityError> {
399 tokio::fs::read(path)
400 .await
401 .map(Zeroizing::new)
402 .map_err(|_| SecurityError("cannot read security manifest or referenced credential file"))
403}
404
405fn record_digest(digest: &mut Sha256, bytes: &[u8]) {
406 digest.update(bytes.len().to_be_bytes());
407 digest.update(bytes);
408}
409
410#[must_use]
412pub struct SecurityReloader {
413 task: tokio::task::JoinHandle<()>,
414 stopping: Arc<std::sync::atomic::AtomicBool>,
415}
416
417struct ReloadExit(Arc<std::sync::atomic::AtomicBool>);
421
422impl Drop for ReloadExit {
423 fn drop(&mut self) {
424 if !self.0.load(std::sync::atomic::Ordering::Acquire) {
425 tracing::error!(
426 operation = "security-reload",
427 reason = "unexpected-exit",
428 "security reload task stopped; configuration will not refresh"
429 );
430 }
431 }
432}
433
434impl SecurityReloader {
435 pub fn spawn(
449 mut loader: SecurityLoader,
450 security: Arc<ServerSecurity>,
451 clock: Arc<dyn tollgate_store::Clock>,
452 ) -> Self {
453 let stopping = Arc::new(std::sync::atomic::AtomicBool::new(false));
454 let exit = ReloadExit(Arc::clone(&stopping));
455 let task = tokio::spawn(async move {
456 let _exit = exit;
457 loop {
458 tokio::time::sleep(Duration::from_secs(5)).await;
459 match loader.load(clock.now()).await {
460 Ok(Some(loaded)) => match loader.install(loaded, &security) {
461 Ok(()) => tracing::info!("control-plane security configuration replaced"),
462 Err(error) => {
463 tracing::warn!(%error, "security replacement refused; previous configuration retained")
464 }
465 },
466 Ok(None) => {}
467 Err(error) => {
468 tracing::warn!(%error, "security reload failed; previous configuration retained")
469 }
470 }
471 }
472 });
473 Self { task, stopping }
474 }
475}
476
477impl Drop for SecurityReloader {
478 fn drop(&mut self) {
479 self.stopping
480 .store(true, std::sync::atomic::Ordering::Release);
481 self.task.abort();
482 }
483}
484
485#[cfg(test)]
486mod tests {
487 use super::*;
488
489 #[tokio::test(start_paused = true)]
490 async fn dropping_the_reloader_releases_its_owned_task_and_clock() {
491 let directory = tempfile::tempdir().unwrap();
492 let path = directory.path().join("security.json");
493 std::fs::write(&path, "{}").unwrap();
494 let now = Timestamp::from_second(100).unwrap();
495 let mut loader = SecurityLoader::new(path);
496 let loaded = loader.load(now).await.unwrap().unwrap();
497 let security = loader.start(loaded).unwrap();
498 let clock = Arc::new(tollgate_store::ManualClock::new(now));
499 let owned = Arc::downgrade(&clock);
500 let reloader = SecurityReloader::spawn(loader, security, clock);
501 tokio::task::yield_now().await;
502 assert!(owned.upgrade().is_some());
503 drop(reloader);
504 tokio::task::yield_now().await;
505 assert!(
506 owned.upgrade().is_none(),
507 "a dropped reloader cannot retain its task's clock"
508 );
509 }
510
511 #[tokio::test(start_paused = true)]
512 async fn signing_key_refresh_honors_success_cadence_and_failure_backoff() {
513 for (lifetime, delay) in [(2, 5), (30, 15), (3600, 300)] {
514 let directory = tempfile::tempdir().unwrap();
515 let path = directory.path().join("security.json");
516 std::fs::write(
517 &path,
518 r#"{"google":{"audience":"https://control.example.test","subjects":[]}}"#,
519 )
520 .unwrap();
521 let fixture: serde_json::Value =
522 serde_json::from_str(include_str!("../tests/fixtures/google-tokens.json")).unwrap();
523 let keys = serde_json::to_vec(&fixture["jwks"]).unwrap();
524 let now = Timestamp::from_second(1_700_000_100).unwrap();
525 let mut loader = SecurityLoader::new(path);
526 let loaded = loader
527 .load_with_keys(now, || async {
528 Ok((keys.clone(), Duration::from_secs(lifetime)))
529 })
530 .await
531 .unwrap()
532 .unwrap();
533 let _security = loader.start(loaded).unwrap();
534 tokio::time::advance(Duration::from_secs(delay - 1)).await;
535 assert!(
536 loader
537 .load_with_keys(now, || async { panic!("successful fetch refreshed early") })
538 .await
539 .unwrap()
540 .is_none()
541 );
542 tokio::time::advance(Duration::from_secs(1)).await;
543 let mut attempted = false;
544 assert!(
545 loader
546 .load_with_keys(now, || async {
547 attempted = true;
548 Err(SecurityError("fixture outage"))
549 })
550 .await
551 .unwrap()
552 .is_none()
553 );
554 assert!(attempted);
555 tokio::time::advance(Duration::from_secs(299)).await;
556 assert!(
557 loader
558 .load_with_keys(now, || async {
559 panic!("failed fetch retried before backoff")
560 })
561 .await
562 .unwrap()
563 .is_none()
564 );
565 tokio::time::advance(Duration::from_secs(1)).await;
566 let mut retried = false;
567 loader
568 .load_with_keys(now, || async {
569 retried = true;
570 Ok((keys, Duration::from_secs(lifetime)))
571 })
572 .await
573 .unwrap();
574 assert!(retried);
575 }
576 }
577
578 #[tokio::test]
579 async fn initial_signing_key_failure_preserves_the_dependency_error() {
580 let directory = tempfile::tempdir().unwrap();
581 let path = directory.path().join("security.json");
582 std::fs::write(
583 &path,
584 r#"{"google":{"audience":"https://control.example.test","subjects":[]}}"#,
585 )
586 .unwrap();
587 let error = SecurityLoader::new(path)
588 .load_with_keys(Timestamp::from_second(100).unwrap(), || async {
589 Err(SecurityError("fixture initial outage"))
590 })
591 .await
592 .err()
593 .unwrap();
594 assert_eq!(error.0, "fixture initial outage");
595 }
596
597 #[tokio::test]
598 async fn failed_key_refresh_never_extends_verified_identity_validity() {
599 use tollgate_auth::CredentialVerifier;
600 let directory = tempfile::tempdir().unwrap();
601 let path = directory.path().join("security.json");
602 std::fs::write(
603 &path,
604 r#"{"google":{"audience":"https://control.example.test","subjects":[]}}"#,
605 )
606 .unwrap();
607 let fixture: serde_json::Value =
608 serde_json::from_str(include_str!("../tests/fixtures/google-tokens.json")).unwrap();
609 let keys = serde_json::to_vec(&fixture["jwks"]).unwrap();
610 let token = fixture["tokens"]["valid"].as_str().unwrap();
611 let now = Timestamp::from_second(1_700_000_100).unwrap();
612 let mut loader = SecurityLoader::new(path);
613 let loaded = loader
614 .load_with_keys(now, || async {
615 Ok((keys.clone(), Duration::from_secs(30)))
616 })
617 .await
618 .unwrap()
619 .unwrap();
620 let security = loader.start(loaded).unwrap();
621 let expiry = now.checked_add(SignedDuration::from_secs(30)).unwrap();
622 let verified = |loader: &SecurityLoader| {
623 let (keys, until) = loader.keys.as_ref().unwrap();
624 GoogleVerifier::from_jwks("https://control.example.test", keys, *until)
625 .unwrap()
626 .verify(token.as_bytes())
627 .unwrap()
628 };
629 assert!(verified(&loader).is_reusable_at(now));
630 assert!(!verified(&loader).is_reusable_at(expiry));
631 loader.next_key_attempt = tokio::time::Instant::now();
632 assert!(
633 loader
634 .load_with_keys(expiry, || async { Err(SecurityError("fixture outage")) })
635 .await
636 .unwrap()
637 .is_none()
638 );
639 assert!(!verified(&loader).is_reusable_at(expiry));
640 loader.next_key_attempt = tokio::time::Instant::now();
642 assert!(
643 loader
644 .load_with_keys(expiry, || async {
645 Ok((b"{}".to_vec(), Duration::from_secs(3600)))
646 })
647 .await
648 .is_err()
649 );
650 assert!(!verified(&loader).is_reusable_at(expiry));
651 loader.next_key_attempt = tokio::time::Instant::now();
652 let loaded = loader
653 .load_with_keys(expiry, || async { Ok((keys, Duration::from_secs(60))) })
654 .await
655 .unwrap()
656 .unwrap();
657 loader.install(loaded, &security).unwrap();
658 assert!(verified(&loader).is_reusable_at(expiry));
659 assert!(
660 !verified(&loader)
661 .is_reusable_at(expiry.checked_add(SignedDuration::from_secs(60)).unwrap())
662 );
663 }
664}