1use std::collections::BTreeMap;
5use std::path::{Path, PathBuf};
6use std::sync::Arc;
7use std::time::Duration;
8
9use jiff::{SignedDuration, Timestamp};
10use serde::Deserialize;
11use sha2::{Digest, Sha256};
12use subtle::ConstantTimeEq;
13use tollgate_auth::{CredentialIssuer, CredentialVerifier, HmacRegistry};
14use zeroize::Zeroizing;
15
16use crate::google::GoogleVerifier;
17use crate::security::{
18 ControlIdentity, ProvisionerLimits, ProvisionerPolicyTemplate, Role, SecurityError,
19 SecurityPolicy, ServerSecurity,
20};
21use crate::transport::{TlsConfig, certificate_fingerprint};
22use tollgate_core::CostUnits;
23
24#[derive(Deserialize)]
25#[serde(deny_unknown_fields)]
26struct Manifest {
27 tls: Option<TlsFiles>,
28 #[serde(default)]
29 bearers: Vec<BearerFile>,
30 #[serde(default)]
31 certificates: Vec<CertificateFile>,
32 google: Option<GoogleConfig>,
33 issuer: Option<IssuerFile>,
34 #[serde(default)]
35 policy_templates: BTreeMap<String, ProvisionerPolicyTemplate>,
36}
37
38#[derive(Deserialize)]
42#[serde(deny_unknown_fields)]
43struct IssuerFile {
44 secret_file: PathBuf,
45}
46
47#[derive(Deserialize)]
48#[serde(deny_unknown_fields)]
49struct TlsFiles {
50 certificate: PathBuf,
51 private_key: PathBuf,
52 client_ca: Option<PathBuf>,
53}
54#[derive(Deserialize)]
55#[serde(deny_unknown_fields)]
56struct BearerFile {
57 identity: String,
58 role: Role,
59 token_file: PathBuf,
60 max_budget_allowance: Option<CostUnits>,
62 allowed_policy_templates: Option<Vec<String>>,
64}
65#[derive(Deserialize)]
66#[serde(deny_unknown_fields)]
67struct CertificateFile {
68 identity: String,
69 role: Role,
70 certificate: PathBuf,
71 max_budget_allowance: Option<CostUnits>,
73 allowed_policy_templates: Option<Vec<String>>,
75}
76#[derive(Deserialize)]
77#[serde(deny_unknown_fields)]
78struct GoogleConfig {
79 audience: String,
80 subjects: Vec<GoogleSubject>,
81}
82#[derive(Deserialize)]
83#[serde(deny_unknown_fields)]
84struct GoogleSubject {
85 subject: String,
86 identity: String,
87 role: Role,
88 max_budget_allowance: Option<CostUnits>,
90 allowed_policy_templates: Option<Vec<String>>,
92}
93
94pub struct SecurityLoader {
104 path: PathBuf,
105 keys: Option<(Vec<u8>, Timestamp)>,
106 next_key_attempt: tokio::time::Instant,
107 digest: Option<[u8; 32]>,
108 issuer: Option<Arc<HmacRegistry>>,
109 issuer_fingerprint: Option<[u8; 32]>,
110 pending_issuer: Option<Option<[u8; 32]>>,
111}
112
113pub struct LoadedSecurity {
119 policy: SecurityPolicy,
120 tls: Option<TlsConfig>,
121 issuer: Option<IssuerSecret>,
122 digest: [u8; 32],
123}
124
125struct IssuerSecret(Zeroizing<[u8; ISSUER_SECRET_LEN]>);
129
130const ISSUER_SECRET_LEN: usize = 64;
131
132impl IssuerSecret {
133 fn parse(bytes: &[u8]) -> Result<Self, SecurityError> {
134 let text = bytes
135 .strip_suffix(b"\n")
136 .map(|line| line.strip_suffix(b"\r").unwrap_or(line))
137 .unwrap_or(bytes);
138 let secret: [u8; ISSUER_SECRET_LEN] = text
139 .try_into()
140 .ok()
141 .filter(|text: &[u8; ISSUER_SECRET_LEN]| {
142 text.iter().all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f'))
143 })
144 .ok_or(SecurityError(
145 "issuer secret must be exactly 64 lowercase hexadecimal characters",
146 ))?;
147 Ok(Self(Zeroizing::new(secret)))
148 }
149
150 fn fingerprint(&self) -> [u8; 32] {
152 let mut digest = Sha256::new();
153 digest.update(b"tollgate-issuer-v1");
154 digest.update(self.0.as_slice());
155 digest.finalize().into()
156 }
157}
158
159impl SecurityLoader {
160 pub fn new(path: impl Into<PathBuf>) -> Self {
164 Self {
165 path: path.into(),
166 keys: None,
167 next_key_attempt: tokio::time::Instant::now(),
168 digest: None,
169 issuer: None,
170 issuer_fingerprint: None,
171 pending_issuer: None,
172 }
173 }
174
175 pub fn start(&mut self, loaded: LoadedSecurity) -> Result<Arc<ServerSecurity>, SecurityError> {
180 let security = ServerSecurity::new(loaded.policy, loaded.tls)?;
181 self.issuer_fingerprint = loaded.issuer.as_ref().map(IssuerSecret::fingerprint);
182 self.issuer = loaded
183 .issuer
184 .map(|secret| Arc::new(HmacRegistry::new(secret.0.as_slice())));
185 self.pending_issuer = None;
186 self.digest = Some(loaded.digest);
187 Ok(security)
188 }
189
190 pub fn install(
194 &mut self,
195 loaded: LoadedSecurity,
196 security: &ServerSecurity,
197 ) -> Result<(), SecurityError> {
198 security.replace(loaded.policy, loaded.tls)?;
199 let staged = loaded.issuer.as_ref().map(IssuerSecret::fingerprint);
200 if staged == self.issuer_fingerprint {
201 if self.pending_issuer.take().is_some() {
202 tracing::info!(
203 reason = "issuer-change-withdrawn",
204 "staged credential issuer matches the live one again"
205 );
206 }
207 } else if self.pending_issuer != Some(staged) {
208 self.pending_issuer = Some(staged);
209 tracing::warn!(
210 reason = "issuer-change-requires-restart",
211 "credential issuer unchanged; restart to apply the staged issuer"
212 );
213 }
214 self.digest = Some(loaded.digest);
215 Ok(())
216 }
217
218 pub fn issuer(&self) -> Option<Arc<dyn CredentialIssuer + Send + Sync>> {
222 self.issuer
223 .as_ref()
224 .map(|issuer| Arc::clone(issuer) as Arc<dyn CredentialIssuer + Send + Sync>)
225 }
226
227 pub fn issuer_change_pending(&self) -> bool {
230 self.pending_issuer.is_some()
231 }
232
233 pub async fn load(&mut self, now: Timestamp) -> Result<Option<LoadedSecurity>, SecurityError> {
261 self.load_with_keys(now, || {
262 crate::google::fetch_keys("https://www.googleapis.com/oauth2/v3/certs")
263 })
264 .await
265 }
266
267 async fn load_with_keys<F: Future<Output = Result<(Vec<u8>, Duration), SecurityError>>>(
268 &mut self,
269 now: Timestamp,
270 fetch: impl FnOnce() -> F,
271 ) -> Result<Option<LoadedSecurity>, SecurityError> {
272 let bytes = read(&self.path).await?;
273 let manifest: Manifest = serde_json::from_slice(&bytes)
274 .map_err(|_| SecurityError("invalid security manifest"))?;
275 for template in manifest.policy_templates.values() {
276 template.validate()?;
277 }
278 let directory = self.path.parent().unwrap_or(Path::new("."));
279 let mut digest = Sha256::new();
280 record_digest(&mut digest, bytes.as_slice());
281 let mut policy = SecurityPolicy::new();
282 let mut secret = Zeroizing::new([0u8; 32]);
283 getrandom::fill(secret.as_mut())
284 .map_err(|_| SecurityError("control-plane verifier entropy unavailable"))?;
285 let registry = Arc::new(HmacRegistry::new(secret.as_ref()));
286 let mut tokens = Vec::new();
287 let mut actors = Vec::new();
288 for bearer in manifest.bearers {
289 let bytes = read(&directory.join(bearer.token_file)).await?;
290 record_digest(&mut digest, bytes.as_slice());
291 let token = Zeroizing::new(
292 std::str::from_utf8(&bytes)
293 .map_err(|_| SecurityError("invalid bearer file encoding"))?
294 .trim_end_matches(['\r', '\n'])
295 .to_owned(),
296 );
297 if token.len() < 32
298 || token.len() > 16 * 1024 - 7
299 || !token.bytes().all(|b| b.is_ascii_graphic())
300 {
301 return Err(SecurityError(
302 "static bearer credentials must contain 32..=16377 visible ASCII bytes",
303 ));
304 }
305 actors.push(identity(
306 bearer.identity,
307 bearer.role,
308 bearer.max_budget_allowance,
309 bearer.allowed_policy_templates,
310 &manifest.policy_templates,
311 )?);
312 tokens.push(token);
313 }
314 registry.install_credentials(tokens.iter().map(|token| token.as_bytes()));
315 let mut mapped = Vec::new();
316 for (token, identity) in tokens.iter().zip(actors) {
317 let principal = registry
318 .verify(token.as_bytes())
319 .ok_or(SecurityError("installed credential failed verification"))?
320 .principal;
321 mapped.push((principal, identity));
322 }
323 policy = policy.with_bearer(registry, mapped)?;
324 let issuer = match manifest.issuer {
325 Some(issuer) => {
326 let bytes = read(&directory.join(issuer.secret_file)).await?;
327 record_digest(&mut digest, bytes.as_slice());
328 let secret = IssuerSecret::parse(&bytes)?;
329 let collides = tokens
330 .iter()
331 .any(|token| bool::from(token.as_bytes().ct_eq(secret.0.as_slice())));
332 if collides {
333 return Err(SecurityError(
334 "issuer secret must differ from every bearer credential",
335 ));
336 }
337 Some(secret)
338 }
339 None => None,
340 };
341 for certificate in manifest.certificates {
342 let pem = read(&directory.join(certificate.certificate)).await?;
343 record_digest(&mut digest, pem.as_slice());
344 policy = policy.with_certificate(
345 certificate_fingerprint(&pem)?,
346 identity(
347 certificate.identity,
348 certificate.role,
349 certificate.max_budget_allowance,
350 certificate.allowed_policy_templates,
351 &manifest.policy_templates,
352 )?,
353 )?;
354 }
355 let tls = match manifest.tls {
356 Some(tls) => {
357 let certificates = read(&directory.join(tls.certificate)).await?;
358 let key = read(&directory.join(tls.private_key)).await?;
359 let ca = match tls.client_ca {
360 Some(path) => Some(read(&directory.join(path)).await?),
361 None => None,
362 };
363 record_digest(&mut digest, certificates.as_slice());
364 record_digest(&mut digest, key.as_slice());
365 if let Some(ca) = &ca {
366 record_digest(&mut digest, ca.as_slice());
367 }
368 Some(TlsConfig::from_pem(
369 &certificates,
370 &key,
371 ca.as_ref().map(|ca| ca.as_slice()),
372 )?)
373 }
374 None => None,
375 };
376 if let Some(google) = manifest.google {
377 if tokio::time::Instant::now() >= self.next_key_attempt {
378 self.next_key_attempt = tokio::time::Instant::now() + Duration::from_secs(300);
379 match fetch().await {
380 Ok((keys, lifetime)) => {
381 self.next_key_attempt = tokio::time::Instant::now()
382 + (lifetime / 2)
383 .clamp(Duration::from_secs(5), Duration::from_secs(300));
384 let until = now
385 .checked_add(
386 SignedDuration::try_from(lifetime)
387 .map_err(|_| SecurityError("signing-key lifetime overflow"))?,
388 )
389 .map_err(|_| SecurityError("signing-key expiry overflow"))?;
390 GoogleVerifier::from_jwks(&google.audience, &keys, until)?;
392 self.keys = Some((keys, until));
393 }
394 Err(error) if self.keys.is_some() => {
395 tracing::warn!(%error, "Google key refresh failed; previous expiry remains authoritative")
396 }
397 Err(error) => return Err(error),
398 }
399 }
400 let (keys, until) = self
401 .keys
402 .as_ref()
403 .ok_or(SecurityError("Google signing keys unavailable"))?;
404 record_digest(&mut digest, keys);
405 record_digest(&mut digest, &until.as_second().to_be_bytes());
406 let verifier = Arc::new(GoogleVerifier::from_jwks(&google.audience, keys, *until)?);
407 let mut mapped = Vec::new();
408 for subject in google.subjects {
409 if subject.subject.is_empty() {
410 return Err(SecurityError("Google subject must not be empty"));
411 }
412 mapped.push((
413 GoogleVerifier::principal(&subject.subject),
414 identity(
415 subject.identity,
416 subject.role,
417 subject.max_budget_allowance,
418 subject.allowed_policy_templates,
419 &manifest.policy_templates,
420 )?,
421 ));
422 }
423 policy = policy.with_bearer(verifier, mapped)?;
424 }
425 let digest = digest.finalize().into();
426 if self.digest == Some(digest) {
427 return Ok(None);
428 }
429 Ok(Some(LoadedSecurity {
430 policy,
431 tls,
432 issuer,
433 digest,
434 }))
435 }
436}
437
438fn identity(
442 name: String,
443 role: Role,
444 max_budget_allowance: Option<CostUnits>,
445 allowed_policy_templates: Option<Vec<String>>,
446 templates: &BTreeMap<String, ProvisionerPolicyTemplate>,
447) -> Result<ControlIdentity, SecurityError> {
448 match (role, max_budget_allowance, allowed_policy_templates) {
449 (Role::Provisioner, Some(max), Some(names)) => {
450 let approved = names
451 .into_iter()
452 .map(|name| {
453 templates
454 .get(&name)
455 .cloned()
456 .ok_or(SecurityError("unknown provisioner policy template"))
457 })
458 .collect::<Result<Vec<_>, _>>()?;
459 ControlIdentity::provisioner(name, ProvisionerLimits::new(max, approved)?)
460 }
461 (Role::Provisioner, _, _) => Err(SecurityError(
462 "a provisioner identity requires max_budget_allowance and allowed_policy_templates",
463 )),
464 (Role::Instance | Role::Operator, None, None) => ControlIdentity::new(name, role),
465 (Role::Instance | Role::Operator, _, _) => Err(SecurityError(
466 "provisioner limits apply only to a provisioner identity",
467 )),
468 }
469}
470
471async fn read(path: &Path) -> Result<Zeroizing<Vec<u8>>, SecurityError> {
472 tokio::fs::read(path)
473 .await
474 .map(Zeroizing::new)
475 .map_err(|_| SecurityError("cannot read security manifest or referenced credential file"))
476}
477
478fn record_digest(digest: &mut Sha256, bytes: &[u8]) {
479 digest.update(bytes.len().to_be_bytes());
480 digest.update(bytes);
481}
482
483#[must_use]
485pub struct SecurityReloader {
486 task: tokio::task::JoinHandle<()>,
487 stopping: Arc<std::sync::atomic::AtomicBool>,
488}
489
490struct ReloadExit(Arc<std::sync::atomic::AtomicBool>);
494
495impl Drop for ReloadExit {
496 fn drop(&mut self) {
497 if !self.0.load(std::sync::atomic::Ordering::Acquire) {
498 tracing::error!(
499 operation = "security-reload",
500 reason = "unexpected-exit",
501 "security reload task stopped; configuration will not refresh"
502 );
503 }
504 }
505}
506
507impl SecurityReloader {
508 pub fn spawn(
522 mut loader: SecurityLoader,
523 security: Arc<ServerSecurity>,
524 clock: Arc<dyn tollgate_store::Clock>,
525 ) -> Self {
526 let stopping = Arc::new(std::sync::atomic::AtomicBool::new(false));
527 let exit = ReloadExit(Arc::clone(&stopping));
528 let task = tokio::spawn(async move {
529 let _exit = exit;
530 loop {
531 tokio::time::sleep(Duration::from_secs(5)).await;
532 match loader.load(clock.now()).await {
533 Ok(Some(loaded)) => match loader.install(loaded, &security) {
534 Ok(()) => tracing::info!("control-plane security configuration replaced"),
535 Err(error) => {
536 tracing::warn!(%error, "security replacement refused; previous configuration retained")
537 }
538 },
539 Ok(None) => {}
540 Err(error) => {
541 tracing::warn!(%error, "security reload failed; previous configuration retained")
542 }
543 }
544 }
545 });
546 Self { task, stopping }
547 }
548}
549
550impl Drop for SecurityReloader {
551 fn drop(&mut self) {
552 self.stopping
553 .store(true, std::sync::atomic::Ordering::Release);
554 self.task.abort();
555 }
556}
557
558#[cfg(test)]
559mod tests {
560 use super::*;
561
562 fn template() -> ProvisionerPolicyTemplate {
563 ProvisionerPolicyTemplate {
564 cost_table: Arc::new(
565 tollgate_core::CostTable::builder(CostUnits(1), CostUnits(1)).build(),
566 ),
567 limits: tollgate_core::ResolvedLimits::new(1),
568 permissions: tollgate_core::PermissionBits(0),
569 policy_revision: tollgate_core::PolicyRevision::UNSTATED,
570 }
571 }
572
573 #[tokio::test(start_paused = true)]
574 async fn dropping_the_reloader_releases_its_owned_task_and_clock() {
575 let directory = tempfile::tempdir().unwrap();
576 let path = directory.path().join("security.json");
577 std::fs::write(&path, "{}").unwrap();
578 let now = Timestamp::from_second(100).unwrap();
579 let mut loader = SecurityLoader::new(path);
580 let loaded = loader.load(now).await.unwrap().unwrap();
581 let security = loader.start(loaded).unwrap();
582 let clock = Arc::new(tollgate_store::ManualClock::new(now));
583 let owned = Arc::downgrade(&clock);
584 let reloader = SecurityReloader::spawn(loader, security, clock);
585 tokio::task::yield_now().await;
586 assert!(owned.upgrade().is_some());
587 drop(reloader);
588 tokio::task::yield_now().await;
589 assert!(
590 owned.upgrade().is_none(),
591 "a dropped reloader cannot retain its task's clock"
592 );
593 }
594
595 #[tokio::test(start_paused = true)]
596 async fn signing_key_refresh_honors_success_cadence_and_failure_backoff() {
597 for (lifetime, delay) in [(2, 5), (30, 15), (3600, 300)] {
598 let directory = tempfile::tempdir().unwrap();
599 let path = directory.path().join("security.json");
600 std::fs::write(
601 &path,
602 r#"{"google":{"audience":"https://control.example.test","subjects":[]}}"#,
603 )
604 .unwrap();
605 let fixture: serde_json::Value =
606 serde_json::from_str(include_str!("../tests/fixtures/google-tokens.json")).unwrap();
607 let keys = serde_json::to_vec(&fixture["jwks"]).unwrap();
608 let now = Timestamp::from_second(1_700_000_100).unwrap();
609 let mut loader = SecurityLoader::new(path);
610 let loaded = loader
611 .load_with_keys(now, || async {
612 Ok((keys.clone(), Duration::from_secs(lifetime)))
613 })
614 .await
615 .unwrap()
616 .unwrap();
617 let _security = loader.start(loaded).unwrap();
618 tokio::time::advance(Duration::from_secs(delay - 1)).await;
619 assert!(
620 loader
621 .load_with_keys(now, || async { panic!("successful fetch refreshed early") })
622 .await
623 .unwrap()
624 .is_none()
625 );
626 tokio::time::advance(Duration::from_secs(1)).await;
627 let mut attempted = false;
628 assert!(
629 loader
630 .load_with_keys(now, || async {
631 attempted = true;
632 Err(SecurityError("fixture outage"))
633 })
634 .await
635 .unwrap()
636 .is_none()
637 );
638 assert!(attempted);
639 tokio::time::advance(Duration::from_secs(299)).await;
640 assert!(
641 loader
642 .load_with_keys(now, || async {
643 panic!("failed fetch retried before backoff")
644 })
645 .await
646 .unwrap()
647 .is_none()
648 );
649 tokio::time::advance(Duration::from_secs(1)).await;
650 let mut retried = false;
651 loader
652 .load_with_keys(now, || async {
653 retried = true;
654 Ok((keys, Duration::from_secs(lifetime)))
655 })
656 .await
657 .unwrap();
658 assert!(retried);
659 }
660 }
661
662 #[tokio::test]
663 async fn a_google_provisioner_subject_requires_its_ceiling() {
664 let fixture: serde_json::Value =
665 serde_json::from_str(include_str!("../tests/fixtures/google-tokens.json")).unwrap();
666 let keys = serde_json::to_vec(&fixture["jwks"]).unwrap();
667 let now = Timestamp::from_second(1_700_000_100).unwrap();
668 for (subject, valid) in [
669 (
670 serde_json::json!({"subject": "1", "identity": "signup", "role": "provisioner", "max_budget_allowance": 1000, "allowed_policy_templates": ["standard"]}),
671 true,
672 ),
673 (
674 serde_json::json!({"subject": "1", "identity": "signup", "role": "provisioner"}),
675 false,
676 ),
677 (
678 serde_json::json!({"subject": "1", "identity": "ops", "role": "operator", "max_budget_allowance": 1000}),
679 false,
680 ),
681 (
682 serde_json::json!({"subject": "1", "identity": "signup", "role": "provisioner", "max_budget_allowance": 1000}),
683 false,
684 ),
685 (
686 serde_json::json!({"subject": "1", "identity": "signup", "role": "provisioner", "max_budget_allowance": 1000, "allowed_policy_templates": []}),
687 false,
688 ),
689 (
690 serde_json::json!({"subject": "1", "identity": "signup", "role": "provisioner", "max_budget_allowance": 1000, "allowed_policy_templates": ["missing"]}),
691 false,
692 ),
693 (
694 serde_json::json!({"subject": "1", "identity": "ops", "role": "operator", "allowed_policy_templates": ["standard"]}),
695 false,
696 ),
697 ] {
698 let directory = tempfile::tempdir().unwrap();
699 let path = directory.path().join("security.json");
700 std::fs::write(
701 &path,
702 serde_json::json!({"policy_templates": {"standard": template()}, "google": {"audience": "https://control.example.test", "subjects": [subject]}})
703 .to_string(),
704 )
705 .unwrap();
706 let loaded = SecurityLoader::new(path)
707 .load_with_keys(now, || async {
708 Ok((keys.clone(), Duration::from_secs(60)))
709 })
710 .await;
711 assert_eq!(loaded.is_ok(), valid, "{subject}");
712 }
713 }
714
715 #[test]
716 fn a_manifest_ceiling_reaches_only_the_provisioner_identity() {
717 let provisioner = identity(
718 "signup".into(),
719 Role::Provisioner,
720 Some(CostUnits(1000)),
721 Some(vec!["standard".into()]),
722 &BTreeMap::from([("standard".into(), template())]),
723 )
724 .unwrap();
725 assert_eq!(
726 provisioner.provisioner_limits(),
727 Some(ProvisionerLimits::new(CostUnits(1000), vec![template()]).unwrap())
728 );
729 for role in [Role::Instance, Role::Operator] {
730 let other = identity("ops".into(), role, None, None, &BTreeMap::new()).unwrap();
731 assert_eq!(other.provisioner_limits(), None);
732 }
733 }
734
735 #[tokio::test]
736 async fn initial_signing_key_failure_preserves_the_dependency_error() {
737 let directory = tempfile::tempdir().unwrap();
738 let path = directory.path().join("security.json");
739 std::fs::write(
740 &path,
741 r#"{"google":{"audience":"https://control.example.test","subjects":[]}}"#,
742 )
743 .unwrap();
744 let error = SecurityLoader::new(path)
745 .load_with_keys(Timestamp::from_second(100).unwrap(), || async {
746 Err(SecurityError("fixture initial outage"))
747 })
748 .await
749 .err()
750 .unwrap();
751 assert_eq!(error.0, "fixture initial outage");
752 }
753
754 #[tokio::test]
755 async fn failed_key_refresh_never_extends_verified_identity_validity() {
756 use tollgate_auth::CredentialVerifier;
757 let directory = tempfile::tempdir().unwrap();
758 let path = directory.path().join("security.json");
759 std::fs::write(
760 &path,
761 r#"{"google":{"audience":"https://control.example.test","subjects":[]}}"#,
762 )
763 .unwrap();
764 let fixture: serde_json::Value =
765 serde_json::from_str(include_str!("../tests/fixtures/google-tokens.json")).unwrap();
766 let keys = serde_json::to_vec(&fixture["jwks"]).unwrap();
767 let token = fixture["tokens"]["valid"].as_str().unwrap();
768 let now = Timestamp::from_second(1_700_000_100).unwrap();
769 let mut loader = SecurityLoader::new(path);
770 let loaded = loader
771 .load_with_keys(now, || async {
772 Ok((keys.clone(), Duration::from_secs(30)))
773 })
774 .await
775 .unwrap()
776 .unwrap();
777 let security = loader.start(loaded).unwrap();
778 let expiry = now.checked_add(SignedDuration::from_secs(30)).unwrap();
779 let verified = |loader: &SecurityLoader| {
780 let (keys, until) = loader.keys.as_ref().unwrap();
781 GoogleVerifier::from_jwks("https://control.example.test", keys, *until)
782 .unwrap()
783 .verify(token.as_bytes())
784 .unwrap()
785 };
786 assert!(verified(&loader).is_reusable_at(now));
787 assert!(!verified(&loader).is_reusable_at(expiry));
788 loader.next_key_attempt = tokio::time::Instant::now();
789 assert!(
790 loader
791 .load_with_keys(expiry, || async { Err(SecurityError("fixture outage")) })
792 .await
793 .unwrap()
794 .is_none()
795 );
796 assert!(!verified(&loader).is_reusable_at(expiry));
797 loader.next_key_attempt = tokio::time::Instant::now();
799 assert!(
800 loader
801 .load_with_keys(expiry, || async {
802 Ok((b"{}".to_vec(), Duration::from_secs(3600)))
803 })
804 .await
805 .is_err()
806 );
807 assert!(!verified(&loader).is_reusable_at(expiry));
808 loader.next_key_attempt = tokio::time::Instant::now();
809 let loaded = loader
810 .load_with_keys(expiry, || async { Ok((keys, Duration::from_secs(60))) })
811 .await
812 .unwrap()
813 .unwrap();
814 loader.install(loaded, &security).unwrap();
815 assert!(verified(&loader).is_reusable_at(expiry));
816 assert!(
817 !verified(&loader)
818 .is_reusable_at(expiry.checked_add(SignedDuration::from_secs(60)).unwrap())
819 );
820 }
821}