1use std::collections::BTreeSet;
10use std::sync::{Arc, Mutex, Weak};
11use std::time::{Duration, Instant};
12
13use meerkat_auth_core::oauth_flow::{
14 OAuthDeviceFlowRecord, OAuthDevicePollLease, OAuthDevicePollLifecycle, OAuthFlowAuthority,
15 OAuthFlowError, OAuthFlowRecord, OAuthFlowRegistry, OAuthFlowRegistrySnapshot,
16 OAuthProviderIdentity, OAuthPrunedFlows, PersistedOAuthBrowserFlow, PersistedOAuthDeviceFlow,
17};
18use meerkat_core::AuthBindingRef;
19use meerkat_core::handles::{DslTransitionError, LeaseKey};
20use meerkat_core::time_compat::{SystemTime, UNIX_EPOCH};
21
22use crate::auth_machine::dsl as auth_dsl;
23use crate::store::RuntimeStore;
24
25use super::RuntimeAuthLeaseHandle;
26use super::auth_lease::{AuthLeaseReleaseObserver, AuthLeaseReleasePermit, ReleasedOAuthFlows};
27
28type StoreSlot = Arc<Mutex<Option<Weak<dyn RuntimeStore>>>>;
29type PayloadLock = Arc<Mutex<()>>;
30type RemovedBrowserSnapshotKeys = Vec<BrowserSnapshotKey>;
31type RemovedDeviceSnapshotKeys = Vec<DeviceSnapshotKey>;
32
33fn current_time_millis() -> u64 {
34 SystemTime::now()
35 .duration_since(UNIX_EPOCH)
36 .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
37 .unwrap_or(0)
38}
39
40fn expires_at_millis(duration: Duration) -> Result<u64, OAuthFlowError> {
41 let duration_millis =
42 u64::try_from(duration.as_millis()).map_err(|_| OAuthFlowError::DeviceExpiryOutOfRange)?;
43 current_time_millis()
44 .checked_add(duration_millis)
45 .ok_or(OAuthFlowError::DeviceExpiryOutOfRange)
46}
47
48fn load_oauth_snapshot_for_release(
49 store: &StoreSlot,
50 operation: &'static str,
51) -> Result<Option<OAuthFlowRegistrySnapshot>, DslTransitionError> {
52 let store = store
53 .lock()
54 .unwrap_or_else(std::sync::PoisonError::into_inner)
55 .clone();
56 let Some(store) = store else {
57 return Ok(None);
58 };
59 let store = store.upgrade().ok_or_else(|| {
60 DslTransitionError::no_matching(operation, "runtime store is no longer available")
61 })?;
62 let Some(bytes) = store
63 .load_auth_oauth_flow_snapshot()
64 .map_err(|err| DslTransitionError::no_matching(operation, err.to_string()))?
65 else {
66 return Ok(None);
67 };
68 serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&bytes)
69 .map(Some)
70 .map_err(|err| DslTransitionError::no_matching(operation, err.to_string()))
71}
72
73#[derive(Debug)]
74pub struct RuntimeOAuthFlowHandle {
75 registry: Arc<OAuthFlowRegistry>,
76 lifecycle: Arc<RuntimeAuthLeaseHandle>,
77 store: StoreSlot,
78 payload_lock: PayloadLock,
79 _release_observer: Option<Arc<OAuthPayloadReleaseObserver>>,
80}
81
82#[derive(Debug)]
83struct OAuthPayloadReleaseObserver {
84 registry: Arc<OAuthFlowRegistry>,
85 store: StoreSlot,
86 payload_lock: PayloadLock,
87}
88
89struct OAuthPayloadReleasePermit<'a> {
90 _payload_guard: std::sync::MutexGuard<'a, ()>,
91}
92
93impl AuthLeaseReleasePermit for OAuthPayloadReleasePermit<'_> {}
94
95impl AuthLeaseReleaseObserver for OAuthPayloadReleaseObserver {
96 fn begin_auth_lease_release<'a>(
97 &'a self,
98 _lease_key: &LeaseKey,
99 ) -> Result<Option<Box<dyn AuthLeaseReleasePermit + 'a>>, DslTransitionError> {
100 let payload_guard = self
101 .payload_lock
102 .lock()
103 .unwrap_or_else(std::sync::PoisonError::into_inner);
104 Ok(Some(Box::new(OAuthPayloadReleasePermit {
105 _payload_guard: payload_guard,
106 })))
107 }
108
109 fn oauth_flows_for_release(
110 &self,
111 lease_key: &LeaseKey,
112 ) -> Result<ReleasedOAuthFlows, DslTransitionError> {
113 let target = AuthBindingRef {
114 realm: lease_key.realm.clone(),
115 binding: lease_key.binding.clone(),
116 profile: lease_key.profile.clone(),
117 origin: meerkat_core::connection::BindingOrigin::Configured,
118 };
119 let Some(snapshot) =
120 load_oauth_snapshot_for_release(&self.store, "collect_oauth_flow_payloads")?
121 else {
122 return Ok(ReleasedOAuthFlows {
123 lease_key: lease_key.clone(),
124 browser_flow_ids: Vec::new(),
125 device_flow_ids: Vec::new(),
126 });
127 };
128 let now_millis = current_time_millis();
129 Ok(ReleasedOAuthFlows {
130 lease_key: lease_key.clone(),
131 browser_flow_ids: snapshot
132 .browser
133 .iter()
134 .filter(|flow| flow.target == target && flow.expires_at_millis > now_millis)
135 .map(|flow| flow.state.clone())
136 .collect(),
137 device_flow_ids: snapshot
138 .device
139 .iter()
140 .filter(|flow| flow.target == target && flow.expires_at_millis > now_millis)
141 .map(|flow| flow.device_code.clone())
142 .collect(),
143 })
144 }
145
146 fn auth_lease_released(&self, released: &ReleasedOAuthFlows) -> Result<(), DslTransitionError> {
147 let target = AuthBindingRef {
148 realm: released.lease_key.realm.clone(),
149 binding: released.lease_key.binding.clone(),
150 profile: released.lease_key.profile.clone(),
151 origin: meerkat_core::connection::BindingOrigin::Configured,
152 };
153 let browser_flow_ids = released
154 .browser_flow_ids
155 .iter()
156 .map(String::as_str)
157 .collect::<BTreeSet<_>>();
158 let device_flow_ids = released
159 .device_flow_ids
160 .iter()
161 .map(String::as_str)
162 .collect::<BTreeSet<_>>();
163 let now_millis = current_time_millis();
164 let mut snapshot = self.registry.snapshot_for_persistence(now_millis);
165 let mut removed_browser = snapshot
166 .browser
167 .iter()
168 .filter(|flow| flow.target == target && browser_flow_ids.contains(flow.state.as_str()))
169 .map(persisted_browser_snapshot_key)
170 .collect::<BTreeSet<_>>();
171 let mut removed_device = snapshot
172 .device
173 .iter()
174 .filter(|flow| {
175 flow.target == target && device_flow_ids.contains(flow.device_code.as_str())
176 })
177 .map(persisted_device_snapshot_key)
178 .collect::<BTreeSet<_>>();
179 if let Some(durable) =
180 load_oauth_snapshot_for_release(&self.store, "release_oauth_flow_payloads")?
181 {
182 removed_browser.extend(
183 durable
184 .browser
185 .iter()
186 .filter(|flow| {
187 flow.target == target
188 && browser_flow_ids.contains(flow.state.as_str())
189 && flow.expires_at_millis > now_millis
190 })
191 .map(persisted_browser_snapshot_key),
192 );
193 removed_device.extend(
194 durable
195 .device
196 .iter()
197 .filter(|flow| {
198 flow.target == target
199 && device_flow_ids.contains(flow.device_code.as_str())
200 && flow.expires_at_millis > now_millis
201 })
202 .map(persisted_device_snapshot_key),
203 );
204 }
205 let removed_browser = removed_browser.into_iter().collect::<Vec<_>>();
206 let removed_device = removed_device.into_iter().collect::<Vec<_>>();
207 snapshot.browser.retain(|flow| {
208 !(flow.target == target && browser_flow_ids.contains(flow.state.as_str()))
209 });
210 snapshot.device.retain(|flow| {
211 !(flow.target == target && device_flow_ids.contains(flow.device_code.as_str()))
212 });
213 persist_registry_snapshot(
214 &snapshot,
215 &self.store,
216 "release_oauth_flow_payloads",
217 &removed_browser,
218 &removed_device,
219 now_millis,
220 SnapshotPersistPolicy::merge(),
221 )
222 .map_err(|err| {
223 DslTransitionError::no_matching(
224 "AuthLeaseReleaseObserver::release_oauth_flow_payloads",
225 err.to_string(),
226 )
227 })?;
228 let _ = self.registry.retain_flows_with_lifecycle(
229 |record_target, flow_id| {
230 !(record_target == &target && browser_flow_ids.contains(flow_id))
231 },
232 |record_target, device_code| {
233 !(record_target == &target && device_flow_ids.contains(device_code))
234 },
235 );
236 Ok(())
237 }
238}
239
240impl RuntimeOAuthFlowHandle {
241 pub fn new(ttl: Duration) -> Self {
242 Self::new_with_auth_lease(ttl, Arc::new(RuntimeAuthLeaseHandle::new()))
243 }
244
245 pub fn new_with_auth_lease(ttl: Duration, lifecycle: Arc<RuntimeAuthLeaseHandle>) -> Self {
246 Self::new_with_capacity_auth_lease_and_store(ttl, 1024, lifecycle, None)
247 }
248
249 pub fn new_with_capacity(ttl: Duration, max_outstanding: usize) -> Self {
250 Self::new_with_capacity_and_auth_lease(
251 ttl,
252 max_outstanding,
253 Arc::new(RuntimeAuthLeaseHandle::new()),
254 )
255 }
256
257 pub fn new_with_capacity_and_auth_lease(
258 ttl: Duration,
259 max_outstanding: usize,
260 lifecycle: Arc<RuntimeAuthLeaseHandle>,
261 ) -> Self {
262 Self::new_with_capacity_auth_lease_and_store(ttl, max_outstanding, lifecycle, None)
263 }
264
265 pub fn new_with_persistent_store_and_auth_lease(
266 ttl: Duration,
267 lifecycle: Arc<RuntimeAuthLeaseHandle>,
268 store: &Arc<dyn RuntimeStore>,
269 ) -> Self {
270 Self::new_with_capacity_auth_lease_and_store(
271 ttl,
272 1024,
273 lifecycle,
274 Some(Arc::downgrade(store)),
275 )
276 }
277
278 fn new_with_capacity_auth_lease_and_store(
279 ttl: Duration,
280 max_outstanding: usize,
281 lifecycle: Arc<RuntimeAuthLeaseHandle>,
282 store: Option<Weak<dyn RuntimeStore>>,
283 ) -> Self {
284 let registry = Arc::new(OAuthFlowRegistry::new_with_capacity(ttl, max_outstanding));
285 let store = Arc::new(Mutex::new(store));
286 let payload_lock = Arc::new(Mutex::new(()));
287 let release_observer = Arc::new(OAuthPayloadReleaseObserver {
288 registry: Arc::clone(®istry),
289 store: Arc::clone(&store),
290 payload_lock: Arc::clone(&payload_lock),
291 });
292 let release_observer_dyn: Arc<dyn AuthLeaseReleaseObserver> = release_observer.clone();
293 lifecycle.add_release_observer(Arc::downgrade(&release_observer_dyn));
294 let handle = Self {
295 registry,
296 lifecycle,
297 store,
298 payload_lock,
299 _release_observer: Some(release_observer),
300 };
301 handle.rehydrate_persisted_payloads();
302 handle
303 }
304
305 pub(crate) fn bind_persistent_store(&self, store: &Arc<dyn RuntimeStore>) {
306 *self
307 .store
308 .lock()
309 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Arc::downgrade(store));
310 }
311
312 fn apply(
313 &self,
314 target: &AuthBindingRef,
315 input: auth_dsl::AuthMachineInput,
316 operation: &'static str,
317 create_if_missing: bool,
318 ) -> Result<(), OAuthFlowError> {
319 self.lifecycle
320 .apply_oauth_input(target, input, operation, create_if_missing)
321 .map_err(|err| OAuthFlowError::LifecycleRejected {
322 operation,
323 detail: err.to_string(),
324 })
325 }
326
327 fn admit_browser(
328 &self,
329 target: &AuthBindingRef,
330 flow_id: &str,
331 provider: OAuthProviderIdentity,
332 redirect_uri: &str,
333 expires_at_millis: u64,
334 ) -> Result<(), OAuthFlowError> {
335 self.apply(
336 target,
337 auth_dsl::AuthMachineInput::AdmitOAuthBrowserFlow {
338 flow_id: flow_id.to_string(),
339 provider: provider.canonical_alias().to_string(),
340 redirect_uri: redirect_uri.to_string(),
341 expires_at_millis,
342 max_outstanding_flows: self.registry.max_outstanding() as u64,
343 observed_global_outstanding_flows: 0,
344 },
345 "admit_oauth_browser_flow",
346 true,
347 )
348 }
349
350 fn verify_browser(
351 &self,
352 target: &AuthBindingRef,
353 flow_id: &str,
354 provider: OAuthProviderIdentity,
355 redirect_uri: &str,
356 ) -> Result<(), OAuthFlowError> {
357 self.apply(
358 target,
359 auth_dsl::AuthMachineInput::VerifyOAuthBrowserFlow {
360 flow_id: flow_id.to_string(),
361 provider: provider.canonical_alias().to_string(),
362 redirect_uri: redirect_uri.to_string(),
363 now_millis: current_time_millis(),
364 },
365 "verify_oauth_browser_flow",
366 false,
367 )
368 }
369
370 fn consume_browser(
371 &self,
372 target: &AuthBindingRef,
373 flow_id: &str,
374 provider: OAuthProviderIdentity,
375 redirect_uri: &str,
376 ) -> Result<(), OAuthFlowError> {
377 self.apply(
378 target,
379 auth_dsl::AuthMachineInput::ConsumeOAuthBrowserFlow {
380 flow_id: flow_id.to_string(),
381 provider: provider.canonical_alias().to_string(),
382 redirect_uri: redirect_uri.to_string(),
383 now_millis: current_time_millis(),
384 },
385 "consume_oauth_browser_flow",
386 false,
387 )
388 }
389
390 fn expire_browser(&self, target: &AuthBindingRef, flow_id: &str) -> Result<(), OAuthFlowError> {
391 self.apply(
392 target,
393 auth_dsl::AuthMachineInput::ExpireOAuthBrowserFlow {
394 flow_id: flow_id.to_string(),
395 },
396 "expire_oauth_browser_flow",
397 false,
398 )
399 }
400
401 fn admit_device(
402 &self,
403 target: &AuthBindingRef,
404 flow_id: &str,
405 provider: OAuthProviderIdentity,
406 expires_at_millis: u64,
407 ) -> Result<(), OAuthFlowError> {
408 self.apply(
409 target,
410 auth_dsl::AuthMachineInput::AdmitOAuthDeviceFlow {
411 flow_id: flow_id.to_string(),
412 provider: provider.canonical_alias().to_string(),
413 expires_at_millis,
414 max_outstanding_flows: self.registry.max_outstanding() as u64,
415 observed_global_outstanding_flows: 0,
416 },
417 "admit_oauth_device_flow",
418 true,
419 )
420 }
421
422 fn verify_device(
423 &self,
424 target: &AuthBindingRef,
425 flow_id: &str,
426 provider: OAuthProviderIdentity,
427 ) -> Result<(), OAuthFlowError> {
428 self.apply(
429 target,
430 auth_dsl::AuthMachineInput::VerifyOAuthDeviceFlow {
431 flow_id: flow_id.to_string(),
432 provider: provider.canonical_alias().to_string(),
433 now_millis: current_time_millis(),
434 },
435 "verify_oauth_device_flow",
436 false,
437 )
438 }
439
440 fn begin_device_poll(
441 &self,
442 target: &AuthBindingRef,
443 flow_id: &str,
444 provider: OAuthProviderIdentity,
445 ) -> Result<(), OAuthFlowError> {
446 self.apply(
447 target,
448 auth_dsl::AuthMachineInput::BeginOAuthDevicePoll {
449 flow_id: flow_id.to_string(),
450 provider: provider.canonical_alias().to_string(),
451 now_millis: current_time_millis(),
452 },
453 "begin_oauth_device_poll",
454 false,
455 )
456 }
457
458 fn expire_pruned_flows(&self) {
459 self.expire_collected_flows(OAuthPrunedFlows {
460 browser: self.registry.prune_expired_browser_flows(),
461 device: self.registry.prune_expired_device_flows(),
462 });
463 }
464
465 fn retain_registry_payloads_with_lifecycle(
466 &self,
467 ) -> (OAuthPrunedFlows, OAuthFlowRegistrySnapshot) {
468 let before = self
469 .registry
470 .snapshot_for_persistence(current_time_millis());
471 let pruned = self.registry.retain_flows_with_lifecycle(
472 |target, flow_id| self.lifecycle.has_oauth_browser_flow(target, flow_id),
473 |target, flow_id| self.lifecycle.has_oauth_device_flow(target, flow_id),
474 );
475 (pruned, before)
476 }
477
478 fn expire_collected_flows(&self, pruned: OAuthPrunedFlows) {
479 for (flow_id, target) in pruned.browser {
480 if let Err(err) = self.expire_browser(&target, &flow_id) {
481 tracing::debug!(
482 target: "meerkat::auth::oauth",
483 binding_target = ?target, %flow_id,
484 "expire_collected_flows: browser flow expiry no-op (legitimate interleaving): {err}"
485 );
486 }
487 }
488 for (device_code, target) in pruned.device {
489 if let Err(err) = self.lifecycle.expire_device_flow(&target, &device_code) {
490 tracing::debug!(
491 target: "meerkat::auth::oauth",
492 binding_target = ?target, %device_code,
493 "expire_collected_flows: device flow expiry no-op (legitimate interleaving): {err}"
494 );
495 }
496 }
497 }
498
499 fn removed_snapshot_keys_from_pruned(
500 snapshot: &OAuthFlowRegistrySnapshot,
501 pruned: &OAuthPrunedFlows,
502 ) -> (RemovedBrowserSnapshotKeys, RemovedDeviceSnapshotKeys) {
503 let pruned_browser = pruned
504 .browser
505 .iter()
506 .map(|(flow_id, target)| browser_snapshot_key(target, flow_id))
507 .collect::<BTreeSet<_>>();
508 let pruned_device = pruned
509 .device
510 .iter()
511 .map(|(device_code, target)| device_snapshot_key(target, device_code))
512 .collect::<BTreeSet<_>>();
513 let browser = snapshot
514 .browser
515 .iter()
516 .filter(|flow| pruned_browser.contains(&persisted_browser_snapshot_key(flow)))
517 .map(persisted_browser_snapshot_key)
518 .collect();
519 let device = snapshot
520 .device
521 .iter()
522 .filter(|flow| pruned_device.contains(&persisted_device_snapshot_key(flow)))
523 .map(persisted_device_snapshot_key)
524 .collect();
525 (browser, device)
526 }
527
528 fn store(&self) -> Option<Arc<dyn RuntimeStore>> {
529 self.store
530 .lock()
531 .unwrap_or_else(std::sync::PoisonError::into_inner)
532 .as_ref()
533 .and_then(Weak::upgrade)
534 }
535
536 fn persist_registry_payloads_removing(
537 &self,
538 operation: &'static str,
539 removed_browser: &[BrowserSnapshotKey],
540 removed_device: &[DeviceSnapshotKey],
541 ) -> Result<(), OAuthFlowError> {
542 persist_registry_payloads(
543 &self.registry,
544 &self.store,
545 operation,
546 removed_browser,
547 removed_device,
548 )
549 }
550
551 fn persist_registry_payloads_claiming_removal(
552 &self,
553 operation: &'static str,
554 removed_browser: &[BrowserSnapshotKey],
555 removed_device: &[DeviceSnapshotKey],
556 ) -> Result<(), OAuthFlowError> {
557 persist_registry_payloads_claiming_removal(
558 &self.registry,
559 &self.store,
560 operation,
561 removed_browser,
562 removed_device,
563 )
564 }
565
566 fn persist_registry_payloads_claiming_admission(
567 &self,
568 operation: &'static str,
569 target: &AuthBindingRef,
570 removed_browser: &[BrowserSnapshotKey],
571 removed_device: &[DeviceSnapshotKey],
572 admitted_browser: &[BrowserSnapshotKey],
573 admitted_device: &[DeviceSnapshotKey],
574 ) -> Result<(), OAuthFlowError> {
575 persist_registry_payloads_claiming_admission(
576 &self.registry,
577 &self.lifecycle,
578 &self.store,
579 operation,
580 removed_browser,
581 removed_device,
582 SnapshotAdmissionClaim {
583 target,
584 admitted_browser,
585 admitted_device,
586 },
587 )
588 }
589
590 fn browser_record_expires_at_millis(
591 &self,
592 record: &OAuthFlowRecord,
593 ) -> Result<u64, OAuthFlowError> {
594 let remaining = self
595 .registry
596 .ttl()
597 .checked_sub(record.created_at.elapsed())
598 .ok_or(OAuthFlowError::Missing)?;
599 expires_at_millis(remaining)
600 }
601
602 fn restore_browser_flow(
603 &self,
604 state: &str,
605 record: &OAuthFlowRecord,
606 ) -> Result<(), OAuthFlowError> {
607 let expires_at_millis = self.browser_record_expires_at_millis(record)?;
608 self.admit_browser(
609 &record.target,
610 state,
611 record.provider,
612 &record.redirect_uri,
613 expires_at_millis,
614 )?;
615 self.registry.insert_restored_browser_flow(
616 state.to_string(),
617 record.target.clone(),
618 record.provider,
619 record.redirect_uri.clone(),
620 record.pkce_verifier.clone(),
621 record.created_at,
622 )
623 }
624
625 fn rehydrate_persisted_payloads(&self) {
626 let Some(store) = self.store() else {
627 return;
628 };
629 let Ok(Some(bytes)) = store.load_auth_oauth_flow_snapshot() else {
630 return;
631 };
632 let Ok(snapshot) = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&bytes) else {
633 return;
634 };
635 let now_millis = current_time_millis();
636 let now_instant = Instant::now();
637
638 for persisted in snapshot.browser.iter().cloned() {
639 self.restore_browser_payload(persisted, now_millis, now_instant);
640 }
641 for persisted in snapshot.device.iter().cloned() {
642 self.restore_device_payload(persisted, now_millis, now_instant);
643 }
644 let current = self.registry.snapshot_for_persistence(now_millis);
645 let current_browser = current
646 .browser
647 .iter()
648 .map(persisted_browser_snapshot_key)
649 .collect::<BTreeSet<_>>();
650 let current_device = current
651 .device
652 .iter()
653 .map(persisted_device_snapshot_key)
654 .collect::<BTreeSet<_>>();
655 let removed_browser = snapshot
656 .browser
657 .iter()
658 .filter(|flow| !current_browser.contains(&persisted_browser_snapshot_key(flow)))
659 .map(persisted_browser_snapshot_key)
660 .collect::<Vec<_>>();
661 let removed_device = snapshot
662 .device
663 .iter()
664 .filter(|flow| !current_device.contains(&persisted_device_snapshot_key(flow)))
665 .map(persisted_device_snapshot_key)
666 .collect::<Vec<_>>();
667 let _ = self.persist_registry_payloads_removing(
668 "rehydrate_oauth_flows",
669 &removed_browser,
670 &removed_device,
671 );
672 }
673
674 fn sync_persisted_payloads(&self, operation: &'static str) -> Result<(), OAuthFlowError> {
675 let Some(store) = self.store() else {
676 return Ok(());
677 };
678 let snapshot =
679 match store.load_auth_oauth_flow_snapshot().map_err(|err| {
680 OAuthFlowError::PersistenceFailed {
681 operation,
682 detail: err.to_string(),
683 }
684 })? {
685 Some(bytes) => serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&bytes)
686 .map_err(|err| OAuthFlowError::PersistenceFailed {
687 operation,
688 detail: err.to_string(),
689 })?,
690 None => OAuthFlowRegistrySnapshot::default(),
691 };
692 let now_millis = current_time_millis();
693 let now_instant = Instant::now();
694 let durable_browser = snapshot
695 .browser
696 .iter()
697 .filter(|flow| flow.expires_at_millis > now_millis)
698 .map(persisted_browser_snapshot_key)
699 .collect::<BTreeSet<_>>();
700 let durable_device = snapshot
701 .device
702 .iter()
703 .filter(|flow| flow.expires_at_millis > now_millis)
704 .map(persisted_device_snapshot_key)
705 .collect::<BTreeSet<_>>();
706 let current = self.registry.snapshot_for_persistence(now_millis);
707 for flow in current
708 .browser
709 .iter()
710 .filter(|flow| !durable_browser.contains(&persisted_browser_snapshot_key(flow)))
711 {
712 let _ =
713 self.registry
714 .consume(&flow.state, &flow.target, flow.provider, &flow.redirect_uri);
715 if let Err(err) = self.expire_browser(&flow.target, &flow.state) {
716 tracing::debug!(
717 target: "meerkat::auth::oauth",
718 binding_target = ?flow.target, flow_id = %flow.state,
719 "sync_persisted_payloads: stale browser expiry no-op (legitimate interleaving): {err}"
720 );
721 }
722 }
723 for flow in current
724 .device
725 .iter()
726 .filter(|flow| !durable_device.contains(&persisted_device_snapshot_key(flow)))
727 {
728 let _ =
729 self.registry
730 .expire_device_code(&flow.device_code, &flow.target, flow.provider);
731 if let Err(err) = self
732 .lifecycle
733 .expire_device_flow(&flow.target, &flow.device_code)
734 {
735 tracing::debug!(
736 target: "meerkat::auth::oauth",
737 binding_target = ?flow.target, device_code = %flow.device_code,
738 "sync_persisted_payloads: stale device expiry no-op (legitimate interleaving): {err}"
739 );
740 }
741 }
742 let current = self.registry.snapshot_for_persistence(now_millis);
743 let current_browser = current
744 .browser
745 .iter()
746 .map(persisted_browser_snapshot_key)
747 .collect::<BTreeSet<_>>();
748 let current_device = current
749 .device
750 .iter()
751 .map(persisted_device_snapshot_key)
752 .collect::<BTreeSet<_>>();
753 for persisted in snapshot
754 .browser
755 .iter()
756 .filter(|flow| {
757 flow.expires_at_millis > now_millis
758 && !current_browser.contains(&persisted_browser_snapshot_key(flow))
759 })
760 .cloned()
761 {
762 self.restore_browser_payload(persisted, now_millis, now_instant);
763 }
764 for persisted in snapshot
765 .device
766 .iter()
767 .filter(|flow| {
768 flow.expires_at_millis > now_millis
769 && !current_device.contains(&persisted_device_snapshot_key(flow))
770 })
771 .cloned()
772 {
773 self.restore_device_payload(persisted, now_millis, now_instant);
774 }
775 Ok(())
776 }
777
778 fn restore_browser_payload(
779 &self,
780 persisted: PersistedOAuthBrowserFlow,
781 now_millis: u64,
782 now_instant: Instant,
783 ) {
784 if persisted.expires_at_millis <= now_millis {
785 return;
786 }
787 let provider = persisted.provider;
788 let remaining = Duration::from_millis(persisted.expires_at_millis - now_millis);
789 let elapsed = self.registry.ttl().saturating_sub(remaining);
790 let created_at = now_instant.checked_sub(elapsed).unwrap_or(now_instant);
791 if self
792 .admit_browser(
793 &persisted.target,
794 &persisted.state,
795 provider,
796 &persisted.redirect_uri,
797 persisted.expires_at_millis,
798 )
799 .is_err()
800 {
801 return;
802 }
803 if self
804 .registry
805 .insert_restored_browser_flow(
806 persisted.state.clone(),
807 persisted.target.clone(),
808 provider,
809 persisted.redirect_uri.clone(),
810 persisted.pkce_verifier.clone(),
811 created_at,
812 )
813 .is_err()
814 && let Err(err) = self.expire_browser(&persisted.target, &persisted.state)
815 {
816 tracing::debug!(
817 target: "meerkat::auth::oauth",
818 binding_target = ?persisted.target, flow_id = %persisted.state,
819 "restore_browser_payload: browser expiry compensation no-op (legitimate interleaving): {err}"
820 );
821 }
822 }
823
824 fn restore_device_payload(
825 &self,
826 persisted: PersistedOAuthDeviceFlow,
827 now_millis: u64,
828 now_instant: Instant,
829 ) {
830 if persisted.expires_at_millis <= now_millis {
831 return;
832 }
833 let provider = persisted.provider;
834 let remaining = Duration::from_millis(persisted.expires_at_millis - now_millis);
835 let expires_at = now_instant.checked_add(remaining).unwrap_or(now_instant);
836 let elapsed = Duration::from_millis(now_millis.saturating_sub(persisted.created_at_millis));
837 let created_at = now_instant.checked_sub(elapsed).unwrap_or(now_instant);
838 if self
839 .admit_device(
840 &persisted.target,
841 &persisted.device_code,
842 provider,
843 persisted.expires_at_millis,
844 )
845 .is_err()
846 {
847 return;
848 }
849 if self
850 .registry
851 .insert_restored_device_flow(
852 persisted.target.clone(),
853 provider,
854 persisted.device_code.clone(),
855 created_at,
856 expires_at,
857 )
858 .is_err()
859 && let Err(err) = self
860 .lifecycle
861 .expire_device_flow(&persisted.target, &persisted.device_code)
862 {
863 tracing::debug!(
864 target: "meerkat::auth::oauth",
865 binding_target = ?persisted.target, device_code = %persisted.device_code,
866 "restore_device_payload: device expiry compensation no-op (legitimate interleaving): {err}"
867 );
868 }
869 }
870}
871
872fn persist_registry_payloads(
873 registry: &OAuthFlowRegistry,
874 store: &StoreSlot,
875 operation: &'static str,
876 removed_browser: &[BrowserSnapshotKey],
877 removed_device: &[DeviceSnapshotKey],
878) -> Result<(), OAuthFlowError> {
879 let now_millis = current_time_millis();
880 let snapshot = registry.snapshot_for_persistence(now_millis);
881 persist_registry_snapshot(
882 &snapshot,
883 store,
884 operation,
885 removed_browser,
886 removed_device,
887 now_millis,
888 SnapshotPersistPolicy::merge(),
889 )
890}
891
892fn persist_existing_registry_payloads(
893 registry: &OAuthFlowRegistry,
894 store: &StoreSlot,
895 operation: &'static str,
896) -> Result<(), OAuthFlowError> {
897 let now_millis = current_time_millis();
898 let snapshot = registry.snapshot_for_persistence(now_millis);
899 persist_registry_snapshot(
900 &snapshot,
901 store,
902 operation,
903 &[],
904 &[],
905 now_millis,
906 SnapshotPersistPolicy::merge_existing(),
907 )
908}
909
910fn persist_registry_payloads_claiming_removal(
911 registry: &OAuthFlowRegistry,
912 store: &StoreSlot,
913 operation: &'static str,
914 removed_browser: &[BrowserSnapshotKey],
915 removed_device: &[DeviceSnapshotKey],
916) -> Result<(), OAuthFlowError> {
917 let now_millis = current_time_millis();
918 let snapshot = registry.snapshot_for_persistence(now_millis);
919 persist_registry_snapshot(
920 &snapshot,
921 store,
922 operation,
923 removed_browser,
924 removed_device,
925 now_millis,
926 SnapshotPersistPolicy::claim_removal(),
927 )
928}
929
930#[derive(Clone, Copy)]
931struct SnapshotAdmissionClaim<'a> {
932 target: &'a AuthBindingRef,
933 admitted_browser: &'a [BrowserSnapshotKey],
934 admitted_device: &'a [DeviceSnapshotKey],
935}
936
937fn persist_registry_payloads_claiming_admission(
938 registry: &OAuthFlowRegistry,
939 lifecycle: &RuntimeAuthLeaseHandle,
940 store: &StoreSlot,
941 operation: &'static str,
942 removed_browser: &[BrowserSnapshotKey],
943 removed_device: &[DeviceSnapshotKey],
944 claim: SnapshotAdmissionClaim<'_>,
945) -> Result<(), OAuthFlowError> {
946 let now_millis = current_time_millis();
947 let snapshot = registry.snapshot_for_persistence(now_millis);
948 persist_registry_snapshot(
949 &snapshot,
950 store,
951 operation,
952 removed_browser,
953 removed_device,
954 now_millis,
955 SnapshotPersistPolicy::claim_admission(
956 registry.max_outstanding(),
957 lifecycle,
958 claim.target,
959 claim.admitted_browser,
960 claim.admitted_device,
961 ),
962 )
963}
964
965#[derive(Clone, Copy, Debug, Eq, PartialEq)]
966enum SnapshotRemovalMode {
967 Merge,
968 Claim,
969}
970
971#[derive(Clone, Copy)]
972struct SnapshotAdmissionConfirmation<'a> {
973 max_outstanding: usize,
974 lifecycle: &'a RuntimeAuthLeaseHandle,
975 target: &'a AuthBindingRef,
976}
977
978#[derive(Clone, Copy)]
979struct SnapshotPersistPolicy<'a> {
980 removal_mode: SnapshotRemovalMode,
981 admission_confirmation: Option<SnapshotAdmissionConfirmation<'a>>,
982 admitted_browser: &'a [BrowserSnapshotKey],
983 admitted_device: &'a [DeviceSnapshotKey],
984}
985
986impl<'a> SnapshotPersistPolicy<'a> {
987 fn merge() -> Self {
988 Self {
989 removal_mode: SnapshotRemovalMode::Merge,
990 admission_confirmation: None,
991 admitted_browser: &[],
992 admitted_device: &[],
993 }
994 }
995
996 fn merge_existing() -> Self {
997 Self::merge()
998 }
999
1000 fn claim_removal() -> Self {
1001 Self {
1002 removal_mode: SnapshotRemovalMode::Claim,
1003 admission_confirmation: None,
1004 admitted_browser: &[],
1005 admitted_device: &[],
1006 }
1007 }
1008
1009 fn claim_admission(
1010 max_outstanding: usize,
1011 lifecycle: &'a RuntimeAuthLeaseHandle,
1012 target: &'a AuthBindingRef,
1013 admitted_browser: &'a [BrowserSnapshotKey],
1014 admitted_device: &'a [DeviceSnapshotKey],
1015 ) -> Self {
1016 Self {
1017 removal_mode: SnapshotRemovalMode::Merge,
1018 admission_confirmation: Some(SnapshotAdmissionConfirmation {
1019 max_outstanding,
1020 lifecycle,
1021 target,
1022 }),
1023 admitted_browser,
1024 admitted_device,
1025 }
1026 }
1027}
1028
1029fn persist_registry_snapshot(
1030 snapshot: &OAuthFlowRegistrySnapshot,
1031 store: &StoreSlot,
1032 operation: &'static str,
1033 removed_browser: &[BrowserSnapshotKey],
1034 removed_device: &[DeviceSnapshotKey],
1035 now_millis: u64,
1036 policy: SnapshotPersistPolicy<'_>,
1037) -> Result<(), OAuthFlowError> {
1038 let store = store
1039 .lock()
1040 .unwrap_or_else(std::sync::PoisonError::into_inner)
1041 .clone();
1042 let Some(store) = store else {
1043 return Ok(());
1044 };
1045 let Some(store) = store.upgrade() else {
1046 return Err(OAuthFlowError::PersistenceFailed {
1047 operation,
1048 detail: "runtime store is no longer available".to_string(),
1049 });
1050 };
1051 let mut durable_admission_rejection: Option<DslTransitionError> = None;
1052 let mut update = |current: Option<&[u8]>| -> Result<Vec<u8>, crate::store::RuntimeStoreError> {
1053 let merged = match merge_oauth_registry_snapshot(
1054 current,
1055 snapshot,
1056 removed_browser,
1057 removed_device,
1058 now_millis,
1059 policy,
1060 ) {
1061 Ok(merged) => merged,
1062 Err(OAuthSnapshotMergeError::Store(err)) => return Err(err),
1063 Err(OAuthSnapshotMergeError::AdmissionRejected(err)) => {
1064 durable_admission_rejection = Some(err);
1065 return Err(crate::store::RuntimeStoreError::Internal(
1066 DURABLE_OAUTH_ADMISSION_REJECTED.to_string(),
1067 ));
1068 }
1069 };
1070 serde_json::to_vec(&merged)
1071 .map_err(|err| crate::store::RuntimeStoreError::WriteFailed(err.to_string()))
1072 };
1073 match store.update_auth_oauth_flow_snapshot(&mut update) {
1074 Ok(()) => Ok(()),
1075 Err(crate::store::RuntimeStoreError::NotFound(_))
1076 if policy.removal_mode == SnapshotRemovalMode::Claim =>
1077 {
1078 Err(OAuthFlowError::RegistryProjectionMissing { operation })
1079 }
1080 Err(crate::store::RuntimeStoreError::Internal(detail))
1081 if detail == DURABLE_OAUTH_ADMISSION_REJECTED =>
1082 {
1083 let detail = durable_admission_rejection
1084 .take()
1085 .map_or_else(|| detail.clone(), |err| err.to_string());
1086 Err(OAuthFlowError::LifecycleRejected { operation, detail })
1087 }
1088 Err(err) => Err(OAuthFlowError::PersistenceFailed {
1089 operation,
1090 detail: err.to_string(),
1091 }),
1092 }
1093}
1094
1095type BrowserSnapshotKey = (String, String, Option<String>, String);
1096type DeviceSnapshotKey = (String, String, Option<String>, String);
1097const DURABLE_OAUTH_ADMISSION_REJECTED: &str =
1098 "oauth durable admission rejected by generated authority";
1099
1100fn target_snapshot_key(target: &AuthBindingRef) -> (String, String, Option<String>) {
1101 (
1102 target.realm.to_string(),
1103 target.binding.to_string(),
1104 target.profile.as_ref().map(ToString::to_string),
1105 )
1106}
1107
1108fn browser_snapshot_key(target: &AuthBindingRef, state: &str) -> BrowserSnapshotKey {
1109 let (realm, binding, profile) = target_snapshot_key(target);
1110 (realm, binding, profile, state.to_string())
1111}
1112
1113fn device_snapshot_key(target: &AuthBindingRef, device_code: &str) -> DeviceSnapshotKey {
1114 let (realm, binding, profile) = target_snapshot_key(target);
1115 (realm, binding, profile, device_code.to_string())
1116}
1117
1118fn persisted_browser_snapshot_key(flow: &PersistedOAuthBrowserFlow) -> BrowserSnapshotKey {
1119 browser_snapshot_key(&flow.target, &flow.state)
1120}
1121
1122fn persisted_device_snapshot_key(flow: &PersistedOAuthDeviceFlow) -> DeviceSnapshotKey {
1123 device_snapshot_key(&flow.target, &flow.device_code)
1124}
1125
1126fn ensure_removed_flows_are_active(
1127 current: &OAuthFlowRegistrySnapshot,
1128 removed_browser: &BTreeSet<BrowserSnapshotKey>,
1129 removed_device: &BTreeSet<DeviceSnapshotKey>,
1130 now_millis: u64,
1131) -> Result<(), crate::store::RuntimeStoreError> {
1132 let active_browser = current
1133 .browser
1134 .iter()
1135 .filter(|flow| flow.expires_at_millis > now_millis)
1136 .map(persisted_browser_snapshot_key)
1137 .collect::<BTreeSet<_>>();
1138 for key in removed_browser {
1139 if !active_browser.contains(key) {
1140 return Err(crate::store::RuntimeStoreError::NotFound(
1141 "oauth browser flow was already consumed".to_string(),
1142 ));
1143 }
1144 }
1145
1146 let active_device = current
1147 .device
1148 .iter()
1149 .filter(|flow| flow.expires_at_millis > now_millis)
1150 .map(persisted_device_snapshot_key)
1151 .collect::<BTreeSet<_>>();
1152 for key in removed_device {
1153 if !active_device.contains(key) {
1154 return Err(crate::store::RuntimeStoreError::NotFound(
1155 "oauth device flow was already consumed".to_string(),
1156 ));
1157 }
1158 }
1159 Ok(())
1160}
1161
1162enum OAuthSnapshotMergeError {
1163 Store(crate::store::RuntimeStoreError),
1164 AdmissionRejected(DslTransitionError),
1165}
1166
1167impl From<crate::store::RuntimeStoreError> for OAuthSnapshotMergeError {
1168 fn from(err: crate::store::RuntimeStoreError) -> Self {
1169 Self::Store(err)
1170 }
1171}
1172
1173fn merge_oauth_registry_snapshot(
1174 current: Option<&[u8]>,
1175 local: &OAuthFlowRegistrySnapshot,
1176 removed_browser: &[BrowserSnapshotKey],
1177 removed_device: &[DeviceSnapshotKey],
1178 now_millis: u64,
1179 policy: SnapshotPersistPolicy<'_>,
1180) -> Result<OAuthFlowRegistrySnapshot, OAuthSnapshotMergeError> {
1181 let mut merged = match current {
1182 Some(bytes) => serde_json::from_slice::<OAuthFlowRegistrySnapshot>(bytes)
1183 .map_err(|err| crate::store::RuntimeStoreError::WriteFailed(err.to_string()))?,
1184 None => OAuthFlowRegistrySnapshot::default(),
1185 };
1186 let removed_browser = removed_browser.iter().cloned().collect::<BTreeSet<_>>();
1187 let removed_device = removed_device.iter().cloned().collect::<BTreeSet<_>>();
1188 let admitted_browser = policy
1189 .admitted_browser
1190 .iter()
1191 .cloned()
1192 .collect::<BTreeSet<_>>();
1193 let admitted_device = policy
1194 .admitted_device
1195 .iter()
1196 .cloned()
1197 .collect::<BTreeSet<_>>();
1198 if policy.removal_mode == SnapshotRemovalMode::Claim {
1199 ensure_removed_flows_are_active(&merged, &removed_browser, &removed_device, now_millis)?;
1200 }
1201 let current_browser = merged
1202 .browser
1203 .iter()
1204 .filter(|flow| flow.expires_at_millis > now_millis)
1205 .map(persisted_browser_snapshot_key)
1206 .collect::<BTreeSet<_>>();
1207 let current_device = merged
1208 .device
1209 .iter()
1210 .filter(|flow| flow.expires_at_millis > now_millis)
1211 .map(persisted_device_snapshot_key)
1212 .collect::<BTreeSet<_>>();
1213 if let Some(confirmation) = policy.admission_confirmation {
1214 let observed_browser = current_browser
1215 .iter()
1216 .filter(|key| !admitted_browser.contains(*key))
1217 .count();
1218 let observed_device = current_device
1219 .iter()
1220 .filter(|key| !admitted_device.contains(*key))
1221 .count();
1222 let observed_global_outstanding_flows =
1223 u64::try_from(observed_browser.saturating_add(observed_device)).unwrap_or(u64::MAX);
1224 let max_outstanding_flows = u64::try_from(confirmation.max_outstanding).unwrap_or(u64::MAX);
1225 confirmation
1226 .lifecycle
1227 .confirm_oauth_durable_admission(
1228 confirmation.target,
1229 observed_global_outstanding_flows,
1230 max_outstanding_flows,
1231 "confirm_oauth_durable_admission",
1232 )
1233 .map_err(OAuthSnapshotMergeError::AdmissionRejected)?;
1234 }
1235 let local_browser = local
1236 .browser
1237 .iter()
1238 .map(persisted_browser_snapshot_key)
1239 .collect::<BTreeSet<_>>();
1240 let local_device = local
1241 .device
1242 .iter()
1243 .map(persisted_device_snapshot_key)
1244 .collect::<BTreeSet<_>>();
1245
1246 merged.browser.retain(|flow| {
1247 flow.expires_at_millis > now_millis
1248 && !removed_browser.contains(&persisted_browser_snapshot_key(flow))
1249 && !local_browser.contains(&persisted_browser_snapshot_key(flow))
1250 });
1251 merged.device.retain(|flow| {
1252 flow.expires_at_millis > now_millis
1253 && !removed_device.contains(&persisted_device_snapshot_key(flow))
1254 && !local_device.contains(&persisted_device_snapshot_key(flow))
1255 });
1256 merged.browser.extend(
1257 local
1258 .browser
1259 .iter()
1260 .filter(|flow| {
1261 let key = persisted_browser_snapshot_key(flow);
1262 flow.expires_at_millis > now_millis
1263 && !removed_browser.contains(&key)
1264 && (current_browser.contains(&key) || admitted_browser.contains(&key))
1265 })
1266 .cloned(),
1267 );
1268 merged.device.extend(
1269 local
1270 .device
1271 .iter()
1272 .filter(|flow| {
1273 let key = persisted_device_snapshot_key(flow);
1274 flow.expires_at_millis > now_millis
1275 && !removed_device.contains(&key)
1276 && (current_device.contains(&key) || admitted_device.contains(&key))
1277 })
1278 .cloned(),
1279 );
1280 merged.browser.sort_by_key(persisted_browser_snapshot_key);
1281 merged.device.sort_by_key(persisted_device_snapshot_key);
1282 Ok(merged)
1283}
1284
1285struct RuntimeOAuthDevicePollLifecycle {
1286 lifecycle: Arc<RuntimeAuthLeaseHandle>,
1287 registry: Arc<OAuthFlowRegistry>,
1288 store: StoreSlot,
1289}
1290
1291impl Default for RuntimeOAuthFlowHandle {
1292 fn default() -> Self {
1293 Self::new(Duration::from_secs(10 * 60))
1294 }
1295}
1296
1297impl OAuthDevicePollLifecycle for RuntimeAuthLeaseHandle {
1298 fn device_flow_state_is_authmachine_owned(&self) -> bool {
1299 true
1300 }
1301
1302 fn finish_device_poll(
1303 &self,
1304 target: &AuthBindingRef,
1305 device_code: &str,
1306 ) -> Result<(), OAuthFlowError> {
1307 self.apply_oauth_input(
1308 target,
1309 auth_dsl::AuthMachineInput::FinishOAuthDevicePoll {
1310 flow_id: device_code.to_string(),
1311 },
1312 "finish_oauth_device_poll",
1313 false,
1314 )
1315 .map_err(|err| OAuthFlowError::LifecycleRejected {
1316 operation: "finish_oauth_device_poll",
1317 detail: err.to_string(),
1318 })
1319 }
1320
1321 fn consume_device_flow(
1322 &self,
1323 target: &AuthBindingRef,
1324 device_code: &str,
1325 provider: OAuthProviderIdentity,
1326 ) -> Result<(), OAuthFlowError> {
1327 self.apply_oauth_input(
1328 target,
1329 auth_dsl::AuthMachineInput::ConsumeOAuthDeviceFlow {
1330 flow_id: device_code.to_string(),
1331 provider: provider.canonical_alias().to_string(),
1332 now_millis: current_time_millis(),
1333 },
1334 "consume_oauth_device_flow",
1335 false,
1336 )
1337 .map_err(|err| OAuthFlowError::LifecycleRejected {
1338 operation: "consume_oauth_device_flow",
1339 detail: err.to_string(),
1340 })
1341 }
1342
1343 fn expire_device_flow(
1344 &self,
1345 target: &AuthBindingRef,
1346 device_code: &str,
1347 ) -> Result<(), OAuthFlowError> {
1348 self.apply_oauth_input(
1349 target,
1350 auth_dsl::AuthMachineInput::ExpireOAuthDeviceFlow {
1351 flow_id: device_code.to_string(),
1352 },
1353 "expire_oauth_device_flow",
1354 false,
1355 )
1356 .map_err(|err| OAuthFlowError::LifecycleRejected {
1357 operation: "expire_oauth_device_flow",
1358 detail: err.to_string(),
1359 })
1360 }
1361
1362 fn restore_device_flow(&self, record: &OAuthDeviceFlowRecord) -> Result<(), OAuthFlowError> {
1363 let remaining = record
1364 .expires_at
1365 .checked_duration_since(Instant::now())
1366 .ok_or(OAuthFlowError::Missing)?;
1367 let expires_at_millis = expires_at_millis(remaining)?;
1368 self.apply_oauth_input(
1369 &record.target,
1370 auth_dsl::AuthMachineInput::AdmitOAuthDeviceFlow {
1371 flow_id: record.device_code.clone(),
1372 provider: record.provider.canonical_alias().to_string(),
1373 expires_at_millis,
1374 max_outstanding_flows: u64::MAX,
1375 observed_global_outstanding_flows: 0,
1376 },
1377 "restore_oauth_device_flow",
1378 true,
1379 )
1380 .map_err(|err| OAuthFlowError::LifecycleRejected {
1381 operation: "restore_oauth_device_flow",
1382 detail: err.to_string(),
1383 })
1384 }
1385}
1386
1387impl OAuthDevicePollLifecycle for RuntimeOAuthDevicePollLifecycle {
1388 fn device_flow_state_is_authmachine_owned(&self) -> bool {
1389 true
1390 }
1391
1392 fn finish_device_poll(
1393 &self,
1394 target: &AuthBindingRef,
1395 device_code: &str,
1396 ) -> Result<(), OAuthFlowError> {
1397 self.lifecycle.finish_device_poll(target, device_code)
1398 }
1399
1400 fn consume_device_flow(
1401 &self,
1402 target: &AuthBindingRef,
1403 device_code: &str,
1404 provider: OAuthProviderIdentity,
1405 ) -> Result<(), OAuthFlowError> {
1406 self.lifecycle
1407 .consume_device_flow(target, device_code, provider)
1408 }
1409
1410 fn expire_device_flow(
1411 &self,
1412 target: &AuthBindingRef,
1413 device_code: &str,
1414 ) -> Result<(), OAuthFlowError> {
1415 self.lifecycle.expire_device_flow(target, device_code)
1416 }
1417
1418 fn restore_device_flow(&self, record: &OAuthDeviceFlowRecord) -> Result<(), OAuthFlowError> {
1419 let remaining = record
1420 .expires_at
1421 .checked_duration_since(Instant::now())
1422 .ok_or(OAuthFlowError::Missing)?;
1423 let expires_at_millis = expires_at_millis(remaining)?;
1424 self.lifecycle
1425 .apply_oauth_input(
1426 &record.target,
1427 auth_dsl::AuthMachineInput::AdmitOAuthDeviceFlow {
1428 flow_id: record.device_code.clone(),
1429 provider: record.provider.canonical_alias().to_string(),
1430 expires_at_millis,
1431 max_outstanding_flows: self.registry.max_outstanding() as u64,
1432 observed_global_outstanding_flows: 0,
1433 },
1434 "restore_oauth_device_flow",
1435 true,
1436 )
1437 .map_err(|err| OAuthFlowError::LifecycleRejected {
1438 operation: "restore_oauth_device_flow",
1439 detail: err.to_string(),
1440 })
1441 }
1442
1443 fn device_flow_payloads_changed(&self) -> Result<(), OAuthFlowError> {
1444 persist_existing_registry_payloads(
1445 &self.registry,
1446 &self.store,
1447 "persist_oauth_device_flow_payloads",
1448 )
1449 }
1450
1451 fn device_flow_payload_removed(
1452 &self,
1453 record: &OAuthDeviceFlowRecord,
1454 ) -> Result<(), OAuthFlowError> {
1455 let removed = [device_snapshot_key(&record.target, &record.device_code)];
1456 persist_registry_payloads_claiming_removal(
1457 &self.registry,
1458 &self.store,
1459 "consume_oauth_device_flow",
1460 &[],
1461 &removed,
1462 )
1463 }
1464}
1465
1466impl OAuthFlowAuthority for RuntimeOAuthFlowHandle {
1467 fn terminal_flow_state_is_authmachine_owned(&self) -> bool {
1468 true
1469 }
1470
1471 fn start(
1472 &self,
1473 target: AuthBindingRef,
1474 provider: OAuthProviderIdentity,
1475 redirect_uri: String,
1476 pkce_verifier: String,
1477 ) -> Result<String, OAuthFlowError> {
1478 let _payload_guard = self
1479 .payload_lock
1480 .lock()
1481 .unwrap_or_else(std::sync::PoisonError::into_inner);
1482 self.sync_persisted_payloads("admit_oauth_browser_flow")?;
1483 self.expire_pruned_flows();
1484 let state = OAuthFlowRegistry::new_state()?;
1485 let expires_at = expires_at_millis(self.registry.ttl())?;
1486 self.admit_browser(&target, &state, provider, &redirect_uri, expires_at)?;
1487 let (lifecycle_pruned, lifecycle_pruned_snapshot) =
1488 self.retain_registry_payloads_with_lifecycle();
1489 let inserted = self.registry.insert_browser_flow_with_pruned(
1490 state.clone(),
1491 target.clone(),
1492 provider,
1493 redirect_uri.clone(),
1494 pkce_verifier,
1495 );
1496 let pruned = match inserted {
1497 Ok(pruned) => pruned,
1498 Err(err) => {
1499 if let Err(expire_err) = self.expire_browser(&target, &state) {
1500 tracing::debug!(
1501 target: "meerkat::auth::oauth",
1502 binding_target = ?target, flow_id = %state,
1503 "start: browser expiry compensation no-op after insert failure (legitimate interleaving): {expire_err}"
1504 );
1505 }
1506 return Err(err);
1507 }
1508 };
1509 let (removed_browser, removed_device) =
1510 Self::removed_snapshot_keys_from_pruned(&lifecycle_pruned_snapshot, &lifecycle_pruned);
1511 self.expire_collected_flows(pruned);
1512 let admitted_browser = [browser_snapshot_key(&target, &state)];
1513 if let Err(err) = self.persist_registry_payloads_claiming_admission(
1514 "admit_oauth_browser_flow",
1515 &target,
1516 &removed_browser,
1517 &removed_device,
1518 &admitted_browser,
1519 &[],
1520 ) {
1521 let _ = self
1522 .registry
1523 .consume(&state, &target, provider, &redirect_uri);
1524 if let Err(expire_err) = self.expire_browser(&target, &state) {
1525 tracing::debug!(
1526 target: "meerkat::auth::oauth",
1527 binding_target = ?target, flow_id = %state,
1528 "start: browser expiry compensation no-op after persist failure (legitimate interleaving): {expire_err}"
1529 );
1530 }
1531 return Err(err);
1532 }
1533 Ok(state)
1534 }
1535
1536 fn verify(
1537 &self,
1538 state: &str,
1539 target: &AuthBindingRef,
1540 provider: OAuthProviderIdentity,
1541 redirect_uri: &str,
1542 ) -> Result<OAuthFlowRecord, OAuthFlowError> {
1543 let _payload_guard = self
1544 .payload_lock
1545 .lock()
1546 .unwrap_or_else(std::sync::PoisonError::into_inner);
1547 self.sync_persisted_payloads("verify_oauth_browser_flow")?;
1548 self.expire_pruned_flows();
1549 self.verify_browser(target, state, provider, redirect_uri)?;
1550 match self.registry.verify(state, target, provider, redirect_uri) {
1551 Ok(record) => Ok(record),
1552 Err(OAuthFlowError::Missing) => Err(OAuthFlowError::RegistryProjectionMissing {
1553 operation: "verify_oauth_browser_flow",
1554 }),
1555 Err(err) => Err(err),
1556 }
1557 }
1558
1559 fn consume(
1560 &self,
1561 state: &str,
1562 target: &AuthBindingRef,
1563 provider: OAuthProviderIdentity,
1564 redirect_uri: &str,
1565 ) -> Result<OAuthFlowRecord, OAuthFlowError> {
1566 let _payload_guard = self
1567 .payload_lock
1568 .lock()
1569 .unwrap_or_else(std::sync::PoisonError::into_inner);
1570 self.sync_persisted_payloads("consume_oauth_browser_flow")?;
1571 self.expire_pruned_flows();
1572 self.verify_browser(target, state, provider, redirect_uri)?;
1573 let record = match self.registry.verify(state, target, provider, redirect_uri) {
1574 Ok(record) => record,
1575 Err(OAuthFlowError::Missing) => {
1576 return Err(OAuthFlowError::RegistryProjectionMissing {
1577 operation: "consume_oauth_browser_flow",
1578 });
1579 }
1580 Err(err) => return Err(err),
1581 };
1582 self.consume_browser(target, state, provider, redirect_uri)?;
1583 if let Err(err) = self.registry.consume(state, target, provider, redirect_uri) {
1584 let _ = self.restore_browser_flow(state, &record);
1585 return Err(match err {
1586 OAuthFlowError::Missing => OAuthFlowError::RegistryProjectionMissing {
1587 operation: "consume_oauth_browser_flow",
1588 },
1589 other => other,
1590 });
1591 }
1592 let removed_browser = [browser_snapshot_key(target, state)];
1593 if let Err(err) = self.persist_registry_payloads_claiming_removal(
1594 "consume_oauth_browser_flow",
1595 &removed_browser,
1596 &[],
1597 ) {
1598 if matches!(
1599 err,
1600 OAuthFlowError::Missing | OAuthFlowError::RegistryProjectionMissing { .. }
1601 ) {
1602 return Err(err);
1603 }
1604 let _ = self.restore_browser_flow(state, &record);
1605 return Err(err);
1606 }
1607 Ok(record)
1608 }
1609
1610 fn admit_device_code(
1611 &self,
1612 target: AuthBindingRef,
1613 provider: OAuthProviderIdentity,
1614 device_code: String,
1615 expires_in: Duration,
1616 ) -> Result<(), OAuthFlowError> {
1617 let _payload_guard = self
1618 .payload_lock
1619 .lock()
1620 .unwrap_or_else(std::sync::PoisonError::into_inner);
1621 self.sync_persisted_payloads("admit_oauth_device_flow")?;
1622 self.expire_pruned_flows();
1623 let machine_expires_at = expires_at_millis(expires_in)?;
1624 self.admit_device(&target, &device_code, provider, machine_expires_at)?;
1625 let (lifecycle_pruned, lifecycle_pruned_snapshot) =
1626 self.retain_registry_payloads_with_lifecycle();
1627 let inserted = self.registry.admit_device_code_with_pruned(
1628 target.clone(),
1629 provider,
1630 device_code.clone(),
1631 expires_in,
1632 );
1633 let pruned = match inserted {
1634 Ok(pruned) => pruned,
1635 Err(err) => {
1636 if let Err(expire_err) = self.lifecycle.expire_device_flow(&target, &device_code) {
1637 tracing::debug!(
1638 target: "meerkat::auth::oauth",
1639 binding_target = ?target, %device_code,
1640 "admit_device_code: device expiry compensation no-op after insert failure (legitimate interleaving): {expire_err}"
1641 );
1642 }
1643 return Err(err);
1644 }
1645 };
1646 let (removed_browser, removed_device) =
1647 Self::removed_snapshot_keys_from_pruned(&lifecycle_pruned_snapshot, &lifecycle_pruned);
1648 self.expire_collected_flows(pruned);
1649 let admitted_device = [device_snapshot_key(&target, &device_code)];
1650 if let Err(err) = self.persist_registry_payloads_claiming_admission(
1651 "admit_oauth_device_flow",
1652 &target,
1653 &removed_browser,
1654 &removed_device,
1655 &[],
1656 &admitted_device,
1657 ) {
1658 let _ = self
1659 .registry
1660 .expire_device_code(&device_code, &target, provider);
1661 if let Err(expire_err) = self.lifecycle.expire_device_flow(&target, &device_code) {
1662 tracing::debug!(
1663 target: "meerkat::auth::oauth",
1664 binding_target = ?target, %device_code,
1665 "admit_device_code: device expiry compensation no-op after persist failure (legitimate interleaving): {expire_err}"
1666 );
1667 }
1668 return Err(err);
1669 }
1670 Ok(())
1671 }
1672
1673 fn verify_device_code(
1674 &self,
1675 device_code: &str,
1676 target: &AuthBindingRef,
1677 provider: OAuthProviderIdentity,
1678 ) -> Result<OAuthDeviceFlowRecord, OAuthFlowError> {
1679 let _payload_guard = self
1680 .payload_lock
1681 .lock()
1682 .unwrap_or_else(std::sync::PoisonError::into_inner);
1683 self.sync_persisted_payloads("verify_oauth_device_flow")?;
1684 self.expire_pruned_flows();
1685 self.verify_device(target, device_code, provider)?;
1686 match self
1687 .registry
1688 .verify_device_code(device_code, target, provider)
1689 {
1690 Ok(record) => Ok(record),
1691 Err(OAuthFlowError::Missing) => Err(OAuthFlowError::RegistryProjectionMissing {
1692 operation: "verify_oauth_device_flow",
1693 }),
1694 Err(err) => Err(err),
1695 }
1696 }
1697
1698 fn begin_device_code_poll(
1699 &self,
1700 device_code: &str,
1701 target: &AuthBindingRef,
1702 provider: OAuthProviderIdentity,
1703 ) -> Result<OAuthDevicePollLease, OAuthFlowError> {
1704 let _payload_guard = self
1705 .payload_lock
1706 .lock()
1707 .unwrap_or_else(std::sync::PoisonError::into_inner);
1708 self.sync_persisted_payloads("begin_oauth_device_poll")?;
1709 self.expire_pruned_flows();
1710 self.begin_device_poll(target, device_code, provider)?;
1711 let poll = match self
1712 .registry
1713 .begin_device_code_poll(device_code, target, provider)
1714 {
1715 Ok(poll) => poll,
1716 Err(OAuthFlowError::Missing) => {
1717 if let Err(finish_err) = self.lifecycle.finish_device_poll(target, device_code) {
1718 tracing::debug!(
1719 target: "meerkat::auth::oauth",
1720 binding_target = ?target, %device_code,
1721 "begin_device_code_poll: poll finish compensation no-op after missing registry entry (legitimate interleaving): {finish_err}"
1722 );
1723 }
1724 return Err(OAuthFlowError::RegistryProjectionMissing {
1725 operation: "begin_oauth_device_poll",
1726 });
1727 }
1728 Err(err) => {
1729 if let Err(finish_err) = self.lifecycle.finish_device_poll(target, device_code) {
1730 tracing::debug!(
1731 target: "meerkat::auth::oauth",
1732 binding_target = ?target, %device_code,
1733 "begin_device_code_poll: poll finish compensation no-op after poll error (legitimate interleaving): {finish_err}"
1734 );
1735 }
1736 return Err(err);
1737 }
1738 };
1739 let lifecycle: Arc<dyn OAuthDevicePollLifecycle> =
1740 Arc::new(RuntimeOAuthDevicePollLifecycle {
1741 lifecycle: Arc::clone(&self.lifecycle),
1742 registry: Arc::clone(&self.registry),
1743 store: Arc::clone(&self.store),
1744 });
1745 Ok(poll
1746 .with_lifecycle(lifecycle)
1747 .with_operation_lock(Arc::clone(&self.payload_lock)))
1748 }
1749}
1750
1751#[cfg(test)]
1752mod tests {
1753 use std::sync::{
1754 Arc, Condvar, Mutex as StdMutex,
1755 atomic::{AtomicBool, Ordering},
1756 mpsc,
1757 };
1758
1759 use meerkat_core::handles::{AuthLeaseHandle, AuthLeasePhase, LeaseKey};
1760 use meerkat_core::lifecycle::run_primitive::RunApplyBoundary;
1761 use meerkat_core::lifecycle::{InputId, RunBoundaryReceipt, RunId};
1762 use meerkat_core::types::SessionId;
1763
1764 use super::*;
1765 use crate::identifiers::LogicalRuntimeId;
1766 use crate::input_state::{InputStatePersistenceRecord, StoredInputState};
1767 use crate::runtime_state::RuntimeState;
1768 use crate::store::{RuntimeStore, RuntimeStoreError, SessionDelta};
1769
1770 fn target_with_binding(binding: &str) -> AuthBindingRef {
1771 AuthBindingRef {
1772 realm: meerkat_core::RealmId::parse("dev").expect("valid realm"),
1773 binding: meerkat_core::BindingId::parse(binding).expect("valid binding"),
1774 profile: None,
1775 origin: meerkat_core::connection::BindingOrigin::Configured,
1776 }
1777 }
1778
1779 fn target() -> AuthBindingRef {
1780 target_with_binding("default_openai")
1781 }
1782
1783 fn alternate_target() -> AuthBindingRef {
1784 target_with_binding("secondary_openai")
1785 }
1786
1787 #[derive(Debug, Default)]
1788 struct BlockingOAuthPersistState {
1789 armed: bool,
1790 blocked: bool,
1791 released: bool,
1792 }
1793
1794 #[derive(Debug, Default)]
1795 struct FailingOAuthSnapshotStore {
1796 snapshot: StdMutex<Option<Vec<u8>>>,
1797 fail_oauth_persist: AtomicBool,
1798 blocking_oauth_persist: StdMutex<BlockingOAuthPersistState>,
1799 blocking_oauth_persist_cv: Condvar,
1800 }
1801
1802 impl FailingOAuthSnapshotStore {
1803 fn block_next_oauth_persist(&self) {
1804 let mut state = self
1805 .blocking_oauth_persist
1806 .lock()
1807 .expect("blocking persist state lock");
1808 state.armed = true;
1809 state.blocked = false;
1810 state.released = false;
1811 }
1812
1813 fn wait_for_blocked_oauth_persist(&self) {
1814 let mut state = self
1815 .blocking_oauth_persist
1816 .lock()
1817 .expect("blocking persist state lock");
1818 while !state.blocked {
1819 let (next, timeout) = self
1820 .blocking_oauth_persist_cv
1821 .wait_timeout(state, Duration::from_secs(1))
1822 .expect("blocking persist state wait");
1823 assert!(
1824 !timeout.timed_out(),
1825 "expected OAuth snapshot persist to block"
1826 );
1827 state = next;
1828 }
1829 }
1830
1831 fn release_blocked_oauth_persist(&self) {
1832 let mut state = self
1833 .blocking_oauth_persist
1834 .lock()
1835 .expect("blocking persist state lock");
1836 state.released = true;
1837 self.blocking_oauth_persist_cv.notify_all();
1838 }
1839
1840 fn wait_if_oauth_persist_blocked(&self) {
1841 let mut state = self
1842 .blocking_oauth_persist
1843 .lock()
1844 .expect("blocking persist state lock");
1845 if !state.armed {
1846 return;
1847 }
1848 state.armed = false;
1849 state.blocked = true;
1850 self.blocking_oauth_persist_cv.notify_all();
1851 while !state.released {
1852 state = self
1853 .blocking_oauth_persist_cv
1854 .wait(state)
1855 .expect("blocking persist state wait");
1856 }
1857 }
1858
1859 fn fail_oauth_persist(&self) {
1860 self.fail_oauth_persist.store(true, Ordering::SeqCst);
1861 }
1862
1863 fn allow_oauth_persist(&self) {
1864 self.fail_oauth_persist.store(false, Ordering::SeqCst);
1865 }
1866 }
1867
1868 #[async_trait::async_trait]
1869 impl RuntimeStore for FailingOAuthSnapshotStore {
1870 fn persist_auth_oauth_flow_snapshot(
1871 &self,
1872 snapshot_json: &[u8],
1873 ) -> Result<(), RuntimeStoreError> {
1874 self.wait_if_oauth_persist_blocked();
1875 if self.fail_oauth_persist.load(Ordering::SeqCst) {
1876 return Err(RuntimeStoreError::WriteFailed(
1877 "injected oauth snapshot failure".to_string(),
1878 ));
1879 }
1880 *self
1881 .snapshot
1882 .lock()
1883 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))? =
1884 Some(snapshot_json.to_vec());
1885 Ok(())
1886 }
1887
1888 fn load_auth_oauth_flow_snapshot(&self) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
1889 self.snapshot
1890 .lock()
1891 .map(|snapshot| snapshot.clone())
1892 .map_err(|err| RuntimeStoreError::ReadFailed(err.to_string()))
1893 }
1894
1895 fn update_auth_oauth_flow_snapshot(
1896 &self,
1897 update: &mut crate::store::AuthOAuthFlowSnapshotUpdate<'_>,
1898 ) -> Result<(), RuntimeStoreError> {
1899 self.wait_if_oauth_persist_blocked();
1900 if self.fail_oauth_persist.load(Ordering::SeqCst) {
1901 return Err(RuntimeStoreError::WriteFailed(
1902 "injected oauth snapshot failure".to_string(),
1903 ));
1904 }
1905 let mut snapshot = self
1906 .snapshot
1907 .lock()
1908 .map_err(|err| RuntimeStoreError::WriteFailed(err.to_string()))?;
1909 let next = update(snapshot.as_deref())?;
1910 *snapshot = Some(next);
1911 Ok(())
1912 }
1913
1914 async fn commit_session_snapshot(
1915 &self,
1916 _runtime_id: &LogicalRuntimeId,
1917 _session_delta: SessionDelta,
1918 ) -> Result<(), RuntimeStoreError> {
1919 Err(RuntimeStoreError::Unsupported(
1920 "commit_session_snapshot".to_string(),
1921 ))
1922 }
1923
1924 async fn atomic_apply(
1925 &self,
1926 _runtime_id: &LogicalRuntimeId,
1927 _session_delta: Option<SessionDelta>,
1928 _receipt: RunBoundaryReceipt,
1929 _input_updates: Vec<InputStatePersistenceRecord>,
1930 _session_store_key: Option<SessionId>,
1931 ) -> Result<(), RuntimeStoreError> {
1932 Err(RuntimeStoreError::Unsupported("atomic_apply".to_string()))
1933 }
1934
1935 async fn load_input_states(
1936 &self,
1937 _runtime_id: &LogicalRuntimeId,
1938 ) -> Result<Vec<StoredInputState>, RuntimeStoreError> {
1939 Err(RuntimeStoreError::Unsupported(
1940 "load_input_states".to_string(),
1941 ))
1942 }
1943
1944 async fn load_boundary_receipt(
1945 &self,
1946 _runtime_id: &LogicalRuntimeId,
1947 _run_id: &RunId,
1948 _sequence: u64,
1949 ) -> Result<Option<RunBoundaryReceipt>, RuntimeStoreError> {
1950 Err(RuntimeStoreError::Unsupported(
1951 "load_boundary_receipt".to_string(),
1952 ))
1953 }
1954
1955 async fn load_session_snapshot(
1956 &self,
1957 _runtime_id: &LogicalRuntimeId,
1958 ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
1959 Err(RuntimeStoreError::Unsupported(
1960 "load_session_snapshot".to_string(),
1961 ))
1962 }
1963
1964 async fn clear_session_snapshot(
1965 &self,
1966 _runtime_id: &LogicalRuntimeId,
1967 ) -> Result<(), RuntimeStoreError> {
1968 Err(RuntimeStoreError::Unsupported(
1969 "clear_session_snapshot".to_string(),
1970 ))
1971 }
1972
1973 async fn replace_session_snapshot_if_current(
1974 &self,
1975 _runtime_id: &LogicalRuntimeId,
1976 _expected_current: &[u8],
1977 _replacement: Vec<u8>,
1978 ) -> Result<bool, RuntimeStoreError> {
1979 Err(RuntimeStoreError::Unsupported(
1980 "replace_session_snapshot_if_current".to_string(),
1981 ))
1982 }
1983
1984 async fn clear_session_snapshot_if_current(
1985 &self,
1986 _runtime_id: &LogicalRuntimeId,
1987 _expected_current: &[u8],
1988 ) -> Result<bool, RuntimeStoreError> {
1989 Err(RuntimeStoreError::Unsupported(
1990 "clear_session_snapshot_if_current".to_string(),
1991 ))
1992 }
1993
1994 async fn persist_input_state(
1995 &self,
1996 _runtime_id: &LogicalRuntimeId,
1997 _state: &InputStatePersistenceRecord,
1998 ) -> Result<(), RuntimeStoreError> {
1999 Err(RuntimeStoreError::Unsupported(
2000 "persist_input_state".to_string(),
2001 ))
2002 }
2003
2004 async fn load_input_state(
2005 &self,
2006 _runtime_id: &LogicalRuntimeId,
2007 _input_id: &InputId,
2008 ) -> Result<Option<StoredInputState>, RuntimeStoreError> {
2009 Err(RuntimeStoreError::Unsupported(
2010 "load_input_state".to_string(),
2011 ))
2012 }
2013
2014 async fn load_machine_lifecycle_record(
2015 &self,
2016 _runtime_id: &LogicalRuntimeId,
2017 ) -> Result<Option<Vec<u8>>, RuntimeStoreError> {
2018 Err(RuntimeStoreError::Unsupported(
2019 "load_machine_lifecycle_record".to_string(),
2020 ))
2021 }
2022
2023 async fn commit_machine_lifecycle(
2024 &self,
2025 _runtime_id: &LogicalRuntimeId,
2026 _commit: crate::store::MachineLifecycleCommit,
2027 _input_states: &[InputStatePersistenceRecord],
2028 ) -> Result<(), RuntimeStoreError> {
2029 Err(RuntimeStoreError::Unsupported(
2030 "commit_machine_lifecycle".to_string(),
2031 ))
2032 }
2033
2034 async fn commit_unregister_finalization(
2035 &self,
2036 _runtime_id: &LogicalRuntimeId,
2037 _finalization: crate::store::UnregisterFinalizationCommit,
2038 ) -> Result<(), RuntimeStoreError> {
2039 Err(RuntimeStoreError::Unsupported(
2040 "commit_unregister_finalization".to_string(),
2041 ))
2042 }
2043 }
2044
2045 fn snapshot_phase(
2046 lifecycle: &RuntimeAuthLeaseHandle,
2047 target: &AuthBindingRef,
2048 ) -> Option<AuthLeasePhase> {
2049 lifecycle
2050 .snapshot(&LeaseKey::from_auth_binding(target))
2051 .phase
2052 }
2053
2054 #[test]
2055 fn browser_flow_only_machine_stays_reauth_required_until_credentials_commit() {
2056 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2057 let authority =
2058 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
2059 let target = target();
2060 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2061 let redirect_uri = "http://127.0.0.1/callback";
2062
2063 let state = authority
2064 .start(
2065 target.clone(),
2066 provider,
2067 redirect_uri.to_string(),
2068 "verifier".to_string(),
2069 )
2070 .expect("browser flow admitted");
2071 assert_eq!(
2072 snapshot_phase(&lifecycle, &target),
2073 Some(AuthLeasePhase::ReauthRequired)
2074 );
2075
2076 authority
2077 .verify(&state, &target, provider, redirect_uri)
2078 .expect("browser flow verifies");
2079 assert_eq!(
2080 snapshot_phase(&lifecycle, &target),
2081 Some(AuthLeasePhase::ReauthRequired)
2082 );
2083
2084 authority
2085 .consume(&state, &target, provider, redirect_uri)
2086 .expect("browser flow consumes");
2087 assert_eq!(
2088 snapshot_phase(&lifecycle, &target),
2089 Some(AuthLeasePhase::ReauthRequired)
2090 );
2091 }
2092
2093 #[test]
2094 fn missing_browser_projection_cannot_overwrite_authmachine_flow() {
2095 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2096 let authority =
2097 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
2098 let target = target();
2099 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2100 let redirect_uri = "http://127.0.0.1/callback";
2101 let state = authority
2102 .start(
2103 target.clone(),
2104 provider,
2105 redirect_uri.to_string(),
2106 "verifier".to_string(),
2107 )
2108 .expect("browser flow admitted");
2109
2110 authority
2111 .registry
2112 .consume(&state, &target, provider, redirect_uri)
2113 .expect("test removes only the local registry payload");
2114
2115 assert!(matches!(
2116 authority.verify(&state, &target, provider, redirect_uri),
2117 Err(OAuthFlowError::RegistryProjectionMissing {
2118 operation: "verify_oauth_browser_flow"
2119 })
2120 ));
2121 assert!(
2122 lifecycle.has_oauth_browser_flow_for_test(&target, &state),
2123 "missing process-local registry payload must not expire canonical AuthMachine membership"
2124 );
2125
2126 assert!(matches!(
2127 authority.consume(&state, &target, provider, redirect_uri),
2128 Err(OAuthFlowError::RegistryProjectionMissing {
2129 operation: "consume_oauth_browser_flow"
2130 })
2131 ));
2132 assert!(
2133 lifecycle.has_oauth_browser_flow_for_test(&target, &state),
2134 "terminal consume must fail closed instead of converting payload loss to not-found"
2135 );
2136 assert_eq!(
2137 snapshot_phase(&lifecycle, &target),
2138 Some(AuthLeasePhase::ReauthRequired)
2139 );
2140 }
2141
2142 #[test]
2143 fn missing_device_poll_projection_cannot_overwrite_authmachine_flow() {
2144 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2145 let authority =
2146 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
2147 let target = target();
2148 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2149 let device_code = "provider-device-code";
2150
2151 authority
2152 .admit_device_code(
2153 target.clone(),
2154 provider,
2155 device_code.to_string(),
2156 Duration::from_secs(60),
2157 )
2158 .expect("device flow admitted");
2159 let poll = authority
2160 .begin_device_code_poll(device_code, &target, provider)
2161 .expect("device poll begins");
2162 authority
2163 .registry
2164 .expire_device_code(device_code, &target, provider)
2165 .expect("test removes only the local registry payload");
2166
2167 assert!(matches!(
2168 poll.consume(),
2169 Err(OAuthFlowError::RegistryProjectionMissing {
2170 operation: "consume_oauth_device_flow"
2171 })
2172 ));
2173 assert!(
2174 lifecycle.has_oauth_device_flow_for_test(&target, device_code),
2175 "missing process-local poll payload must not expire canonical AuthMachine membership"
2176 );
2177 }
2178
2179 #[test]
2180 fn browser_admit_persistence_failure_rolls_back_unreturned_flow() {
2181 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2182 let store = Arc::new(FailingOAuthSnapshotStore::default());
2183 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2184 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2185 Duration::from_secs(60),
2186 lifecycle.clone(),
2187 &store_dyn,
2188 );
2189 let failed_target = target();
2190 let successful_target = alternate_target();
2191 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2192 let redirect_uri = "http://127.0.0.1/callback";
2193
2194 store.fail_oauth_persist();
2195 assert!(matches!(
2196 authority.start(
2197 failed_target.clone(),
2198 provider,
2199 redirect_uri.to_string(),
2200 "failed-verifier".to_string(),
2201 ),
2202 Err(OAuthFlowError::PersistenceFailed { .. })
2203 ));
2204 assert!(
2205 authority
2206 .registry
2207 .snapshot_for_persistence(current_time_millis())
2208 .browser
2209 .is_empty(),
2210 "failed browser admission must not leave an unreturned registry payload"
2211 );
2212
2213 store.allow_oauth_persist();
2214 let successful_state = authority
2215 .start(
2216 successful_target,
2217 provider,
2218 redirect_uri.to_string(),
2219 "successful-verifier".to_string(),
2220 )
2221 .expect("subsequent browser admit persists after store recovers");
2222 let snapshot_json = store
2223 .load_auth_oauth_flow_snapshot()
2224 .expect("durable OAuth snapshot loads")
2225 .expect("durable OAuth snapshot exists");
2226 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2227 .expect("durable OAuth snapshot decodes");
2228 assert_eq!(
2229 snapshot
2230 .browser
2231 .iter()
2232 .map(|flow| flow.state.as_str())
2233 .collect::<Vec<_>>(),
2234 vec![successful_state.as_str()],
2235 "a later successful admit must not persist a previously failed unreturned flow"
2236 );
2237 }
2238
2239 #[test]
2240 fn device_admit_persistence_failure_rolls_back_unreturned_flow() {
2241 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2242 let store = Arc::new(FailingOAuthSnapshotStore::default());
2243 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2244 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2245 Duration::from_secs(60),
2246 lifecycle.clone(),
2247 &store_dyn,
2248 );
2249 let target = target();
2250 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2251 let failed_device_code = "failed-device-code";
2252 let successful_device_code = "successful-device-code";
2253
2254 store.fail_oauth_persist();
2255 assert!(matches!(
2256 authority.admit_device_code(
2257 target.clone(),
2258 provider,
2259 failed_device_code.to_string(),
2260 Duration::from_secs(60),
2261 ),
2262 Err(OAuthFlowError::PersistenceFailed { .. })
2263 ));
2264 assert!(matches!(
2265 authority
2266 .registry
2267 .verify_device_code(failed_device_code, &target, provider),
2268 Err(OAuthFlowError::Missing)
2269 ));
2270 assert!(!lifecycle.has_oauth_device_flow_for_test(&target, failed_device_code));
2271
2272 store.allow_oauth_persist();
2273 authority
2274 .admit_device_code(
2275 target,
2276 provider,
2277 successful_device_code.to_string(),
2278 Duration::from_secs(60),
2279 )
2280 .expect("subsequent device admit persists after store recovers");
2281 let snapshot_json = store
2282 .load_auth_oauth_flow_snapshot()
2283 .expect("durable OAuth snapshot loads")
2284 .expect("durable OAuth snapshot exists");
2285 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2286 .expect("durable OAuth snapshot decodes");
2287 assert_eq!(
2288 snapshot
2289 .device
2290 .iter()
2291 .map(|flow| flow.device_code.as_str())
2292 .collect::<Vec<_>>(),
2293 vec![successful_device_code],
2294 "a later successful admit must not persist a previously failed unreturned device flow"
2295 );
2296 }
2297
2298 #[cfg(feature = "sqlite-store")]
2299 #[test]
2300 fn persistent_oauth_snapshot_merges_independent_authority_writes() {
2301 let temp_dir = tempfile::tempdir().expect("tempdir");
2302 let store_path = temp_dir.path().join("runtime.sqlite");
2303 let store_one: Arc<dyn RuntimeStore> =
2304 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2305 let store_two: Arc<dyn RuntimeStore> =
2306 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2307 let first_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2308 Duration::from_secs(60),
2309 Arc::new(RuntimeAuthLeaseHandle::new()),
2310 &store_one,
2311 );
2312 let second_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2313 Duration::from_secs(60),
2314 Arc::new(RuntimeAuthLeaseHandle::new()),
2315 &store_two,
2316 );
2317 let first_target = target();
2318 let second_target = alternate_target();
2319 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2320
2321 let first_state = first_authority
2322 .start(
2323 first_target.clone(),
2324 provider,
2325 "http://127.0.0.1/callback".to_string(),
2326 "verifier-1".to_string(),
2327 )
2328 .expect("first process admits browser flow");
2329 let second_state = second_authority
2330 .start(
2331 second_target.clone(),
2332 provider,
2333 "http://127.0.0.1/other-callback".to_string(),
2334 "verifier-2".to_string(),
2335 )
2336 .expect("second process admits browser flow");
2337
2338 let store_three: Arc<dyn RuntimeStore> =
2339 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2340 let snapshot_json = store_three
2341 .load_auth_oauth_flow_snapshot()
2342 .expect("durable OAuth snapshot loads")
2343 .expect("durable OAuth snapshot exists");
2344 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2345 .expect("durable OAuth snapshot decodes");
2346 assert!(
2347 snapshot
2348 .browser
2349 .iter()
2350 .any(|flow| flow.state == first_state),
2351 "the first independent authority write must survive the second write"
2352 );
2353 assert!(
2354 snapshot
2355 .browser
2356 .iter()
2357 .any(|flow| flow.state == second_state),
2358 "the second independent authority write must be persisted"
2359 );
2360
2361 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2362 Duration::from_secs(60),
2363 Arc::new(RuntimeAuthLeaseHandle::new()),
2364 &store_three,
2365 );
2366 restarted
2367 .consume(
2368 &first_state,
2369 &first_target,
2370 provider,
2371 "http://127.0.0.1/callback",
2372 )
2373 .expect("first independent flow rehydrates");
2374 restarted
2375 .consume(
2376 &second_state,
2377 &second_target,
2378 provider,
2379 "http://127.0.0.1/other-callback",
2380 )
2381 .expect("second independent flow rehydrates after first consume");
2382 }
2383
2384 #[test]
2385 fn persistent_oauth_browser_admit_does_not_resurrect_consumed_between_sync_and_persist() {
2386 let store = Arc::new(FailingOAuthSnapshotStore::default());
2387 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2388 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2389 Duration::from_secs(60),
2390 Arc::new(RuntimeAuthLeaseHandle::new()),
2391 &store_dyn,
2392 );
2393 let target = target();
2394 let replacement_target = alternate_target();
2395 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2396 let redirect_uri = "http://127.0.0.1/callback";
2397 let replacement_redirect_uri = "http://127.0.0.1/replacement-callback";
2398 let consumed_state = creator
2399 .start(
2400 target.clone(),
2401 provider,
2402 redirect_uri.to_string(),
2403 "consumed-verifier".to_string(),
2404 )
2405 .expect("creator admits browser flow");
2406
2407 let stale_authority = Arc::new(
2408 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2409 Duration::from_secs(60),
2410 Arc::new(RuntimeAuthLeaseHandle::new()),
2411 &store_dyn,
2412 ),
2413 );
2414 let consumer = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2415 Duration::from_secs(60),
2416 Arc::new(RuntimeAuthLeaseHandle::new()),
2417 &store_dyn,
2418 );
2419
2420 store.block_next_oauth_persist();
2421 let stale_admit = std::thread::spawn({
2422 let stale_authority = Arc::clone(&stale_authority);
2423 let replacement_target = replacement_target.clone();
2424 move || {
2425 stale_authority.start(
2426 replacement_target,
2427 provider,
2428 replacement_redirect_uri.to_string(),
2429 "replacement-verifier".to_string(),
2430 )
2431 }
2432 });
2433 store.wait_for_blocked_oauth_persist();
2434 consumer
2435 .consume(&consumed_state, &target, provider, redirect_uri)
2436 .expect("independent authority consumes browser flow between sync and persist");
2437 store.release_blocked_oauth_persist();
2438 let replacement_state = stale_admit
2439 .join()
2440 .expect("stale browser admit thread should not panic")
2441 .expect("stale authority admits replacement browser flow");
2442
2443 let snapshot_json = store
2444 .load_auth_oauth_flow_snapshot()
2445 .expect("durable OAuth snapshot loads")
2446 .expect("durable OAuth snapshot exists");
2447 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2448 .expect("durable OAuth snapshot decodes");
2449 assert!(
2450 !snapshot
2451 .browser
2452 .iter()
2453 .any(|flow| flow.state == consumed_state),
2454 "a stale admission must not resurrect a browser flow consumed after pre-sync"
2455 );
2456 assert!(
2457 snapshot
2458 .browser
2459 .iter()
2460 .any(|flow| flow.state == replacement_state),
2461 "the stale authority's newly admitted browser flow should still persist"
2462 );
2463
2464 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2465 Duration::from_secs(60),
2466 Arc::new(RuntimeAuthLeaseHandle::new()),
2467 &store_dyn,
2468 );
2469 assert!(matches!(
2470 restarted.consume(&consumed_state, &target, provider, redirect_uri),
2471 Err(OAuthFlowError::LifecycleRejected {
2472 operation: "verify_oauth_browser_flow",
2473 ..
2474 })
2475 ));
2476 restarted
2477 .consume(
2478 &replacement_state,
2479 &replacement_target,
2480 provider,
2481 replacement_redirect_uri,
2482 )
2483 .expect("new stale-authority flow survives restart");
2484 }
2485
2486 #[test]
2487 fn persistent_oauth_device_admit_does_not_resurrect_consumed_between_sync_and_persist() {
2488 let store = Arc::new(FailingOAuthSnapshotStore::default());
2489 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2490 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2491 Duration::from_secs(60),
2492 Arc::new(RuntimeAuthLeaseHandle::new()),
2493 &store_dyn,
2494 );
2495 let target = target();
2496 let replacement_target = alternate_target();
2497 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2498 let consumed_device_code = "consumed-device-code";
2499 let replacement_device_code = "replacement-device-code";
2500 creator
2501 .admit_device_code(
2502 target.clone(),
2503 provider,
2504 consumed_device_code.to_string(),
2505 Duration::from_secs(60),
2506 )
2507 .expect("creator admits device flow");
2508
2509 let stale_authority = Arc::new(
2510 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2511 Duration::from_secs(60),
2512 Arc::new(RuntimeAuthLeaseHandle::new()),
2513 &store_dyn,
2514 ),
2515 );
2516 let consumer = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2517 Duration::from_secs(60),
2518 Arc::new(RuntimeAuthLeaseHandle::new()),
2519 &store_dyn,
2520 );
2521
2522 store.block_next_oauth_persist();
2523 let stale_admit = std::thread::spawn({
2524 let stale_authority = Arc::clone(&stale_authority);
2525 let replacement_target = replacement_target.clone();
2526 move || {
2527 stale_authority.admit_device_code(
2528 replacement_target,
2529 provider,
2530 replacement_device_code.to_string(),
2531 Duration::from_secs(60),
2532 )
2533 }
2534 });
2535 store.wait_for_blocked_oauth_persist();
2536 consumer
2537 .begin_device_code_poll(consumed_device_code, &target, provider)
2538 .expect("independent authority begins device poll")
2539 .consume()
2540 .expect("independent authority consumes device flow between sync and persist");
2541 store.release_blocked_oauth_persist();
2542 stale_admit
2543 .join()
2544 .expect("stale device admit thread should not panic")
2545 .expect("stale authority admits replacement device flow");
2546
2547 let snapshot_json = store
2548 .load_auth_oauth_flow_snapshot()
2549 .expect("durable OAuth snapshot loads")
2550 .expect("durable OAuth snapshot exists");
2551 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2552 .expect("durable OAuth snapshot decodes");
2553 assert!(
2554 !snapshot
2555 .device
2556 .iter()
2557 .any(|flow| flow.device_code == consumed_device_code),
2558 "a stale admission must not resurrect a device flow consumed after pre-sync"
2559 );
2560 assert!(
2561 snapshot
2562 .device
2563 .iter()
2564 .any(|flow| flow.device_code == replacement_device_code),
2565 "the stale authority's newly admitted device flow should still persist"
2566 );
2567
2568 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2569 Duration::from_secs(60),
2570 Arc::new(RuntimeAuthLeaseHandle::new()),
2571 &store_dyn,
2572 );
2573 assert!(matches!(
2574 restarted.verify_device_code(consumed_device_code, &target, provider),
2575 Err(OAuthFlowError::LifecycleRejected {
2576 operation: "verify_oauth_device_flow",
2577 ..
2578 })
2579 ));
2580 restarted
2581 .verify_device_code(replacement_device_code, &replacement_target, provider)
2582 .expect("new stale-authority device flow survives restart");
2583 }
2584
2585 #[cfg(feature = "sqlite-store")]
2586 #[test]
2587 fn persistent_oauth_device_poll_finish_does_not_resurrect_consumed_payload() {
2588 let temp_dir = tempfile::tempdir().expect("tempdir");
2589 let store_path = temp_dir.path().join("runtime.sqlite");
2590 let creator_store: Arc<dyn RuntimeStore> =
2591 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2592 let stale_store: Arc<dyn RuntimeStore> =
2593 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2594 let consumer_store: Arc<dyn RuntimeStore> =
2595 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2596 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2597 Duration::from_secs(60),
2598 Arc::new(RuntimeAuthLeaseHandle::new()),
2599 &creator_store,
2600 );
2601 let target = target();
2602 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2603 let device_code = "pending-finish-device-code";
2604 creator
2605 .admit_device_code(
2606 target.clone(),
2607 provider,
2608 device_code.to_string(),
2609 Duration::from_secs(60),
2610 )
2611 .expect("creator admits device flow");
2612
2613 let stale_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2614 Duration::from_secs(60),
2615 Arc::new(RuntimeAuthLeaseHandle::new()),
2616 &stale_store,
2617 );
2618 let stale_poll = stale_authority
2619 .begin_device_code_poll(device_code, &target, provider)
2620 .expect("stale authority begins pending poll");
2621 let consumer = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2622 Duration::from_secs(60),
2623 Arc::new(RuntimeAuthLeaseHandle::new()),
2624 &consumer_store,
2625 );
2626 consumer
2627 .begin_device_code_poll(device_code, &target, provider)
2628 .expect("consumer begins independent poll")
2629 .consume()
2630 .expect("consumer consumes durable payload");
2631
2632 stale_poll
2633 .finish()
2634 .expect("stale pending poll finish is local cleanup only");
2635
2636 let snapshot_json = stale_store
2637 .load_auth_oauth_flow_snapshot()
2638 .expect("durable OAuth snapshot loads")
2639 .expect("durable OAuth snapshot exists");
2640 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2641 .expect("durable OAuth snapshot decodes");
2642 assert!(
2643 !snapshot
2644 .device
2645 .iter()
2646 .any(|flow| flow.device_code == device_code),
2647 "a stale pending poll finish must not resurrect a consumed device flow"
2648 );
2649
2650 let restarted_store: Arc<dyn RuntimeStore> =
2651 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2652 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2653 Duration::from_secs(60),
2654 Arc::new(RuntimeAuthLeaseHandle::new()),
2655 &restarted_store,
2656 );
2657 assert!(matches!(
2658 restarted.verify_device_code(device_code, &target, provider),
2659 Err(OAuthFlowError::LifecycleRejected {
2660 operation: "verify_oauth_device_flow",
2661 ..
2662 })
2663 ));
2664 }
2665
2666 #[cfg(feature = "sqlite-store")]
2667 #[test]
2668 fn persistent_oauth_browser_sync_prunes_stale_capacity_before_admit() {
2669 let temp_dir = tempfile::tempdir().expect("tempdir");
2670 let store_path = temp_dir.path().join("runtime.sqlite");
2671 let creator_store: Arc<dyn RuntimeStore> =
2672 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2673 let stale_store: Arc<dyn RuntimeStore> =
2674 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2675 let consumer_store: Arc<dyn RuntimeStore> =
2676 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2677 let creator = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2678 Duration::from_secs(60),
2679 1,
2680 Arc::new(RuntimeAuthLeaseHandle::new()),
2681 Some(Arc::downgrade(&creator_store)),
2682 );
2683 let target = target();
2684 let replacement_target = alternate_target();
2685 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2686 let redirect_uri = "http://127.0.0.1/callback";
2687 let replacement_redirect_uri = "http://127.0.0.1/replacement-callback";
2688 let consumed_state = creator
2689 .start(
2690 target.clone(),
2691 provider,
2692 redirect_uri.to_string(),
2693 "consumed-verifier".to_string(),
2694 )
2695 .expect("creator admits browser flow");
2696
2697 let stale_authority = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2698 Duration::from_secs(60),
2699 1,
2700 Arc::new(RuntimeAuthLeaseHandle::new()),
2701 Some(Arc::downgrade(&stale_store)),
2702 );
2703 let consumer = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2704 Duration::from_secs(60),
2705 1,
2706 Arc::new(RuntimeAuthLeaseHandle::new()),
2707 Some(Arc::downgrade(&consumer_store)),
2708 );
2709 consumer
2710 .consume(&consumed_state, &target, provider, redirect_uri)
2711 .expect("independent authority consumes browser flow");
2712
2713 let replacement_state = stale_authority
2714 .start(
2715 replacement_target.clone(),
2716 provider,
2717 replacement_redirect_uri.to_string(),
2718 "replacement-verifier".to_string(),
2719 )
2720 .expect("stale capacity is pruned before browser admit");
2721 stale_authority
2722 .consume(
2723 &replacement_state,
2724 &replacement_target,
2725 provider,
2726 replacement_redirect_uri,
2727 )
2728 .expect("replacement browser flow remains usable");
2729 }
2730
2731 #[cfg(feature = "sqlite-store")]
2732 #[test]
2733 fn persistent_oauth_device_sync_prunes_stale_capacity_before_admit() {
2734 let temp_dir = tempfile::tempdir().expect("tempdir");
2735 let store_path = temp_dir.path().join("runtime.sqlite");
2736 let creator_store: Arc<dyn RuntimeStore> =
2737 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2738 let stale_store: Arc<dyn RuntimeStore> =
2739 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2740 let consumer_store: Arc<dyn RuntimeStore> =
2741 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2742 let creator = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2743 Duration::from_secs(60),
2744 1,
2745 Arc::new(RuntimeAuthLeaseHandle::new()),
2746 Some(Arc::downgrade(&creator_store)),
2747 );
2748 let target = target();
2749 let replacement_target = alternate_target();
2750 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2751 let consumed_device_code = "consumed-capacity-device-code";
2752 let replacement_device_code = "replacement-device-code";
2753 creator
2754 .admit_device_code(
2755 target.clone(),
2756 provider,
2757 consumed_device_code.to_string(),
2758 Duration::from_secs(60),
2759 )
2760 .expect("creator admits device flow");
2761
2762 let stale_authority = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2763 Duration::from_secs(60),
2764 1,
2765 Arc::new(RuntimeAuthLeaseHandle::new()),
2766 Some(Arc::downgrade(&stale_store)),
2767 );
2768 let consumer = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2769 Duration::from_secs(60),
2770 1,
2771 Arc::new(RuntimeAuthLeaseHandle::new()),
2772 Some(Arc::downgrade(&consumer_store)),
2773 );
2774 consumer
2775 .begin_device_code_poll(consumed_device_code, &target, provider)
2776 .expect("independent authority begins device poll")
2777 .consume()
2778 .expect("independent authority consumes device flow");
2779
2780 stale_authority
2781 .admit_device_code(
2782 replacement_target.clone(),
2783 provider,
2784 replacement_device_code.to_string(),
2785 Duration::from_secs(60),
2786 )
2787 .expect("stale capacity is pruned before device admit");
2788 stale_authority
2789 .verify_device_code(replacement_device_code, &replacement_target, provider)
2790 .expect("replacement device flow remains usable");
2791 }
2792
2793 #[test]
2794 fn concurrent_persistent_browser_consumes_require_fresh_durable_claim() {
2795 let store = Arc::new(FailingOAuthSnapshotStore::default());
2796 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2797 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2798 Duration::from_secs(60),
2799 Arc::new(RuntimeAuthLeaseHandle::new()),
2800 &store_dyn,
2801 );
2802 let target = target();
2803 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2804 let redirect_uri = "http://127.0.0.1/callback";
2805 let state = creator
2806 .start(
2807 target.clone(),
2808 provider,
2809 redirect_uri.to_string(),
2810 "verifier".to_string(),
2811 )
2812 .expect("browser flow admitted");
2813 let first = Arc::new(
2814 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2815 Duration::from_secs(60),
2816 Arc::new(RuntimeAuthLeaseHandle::new()),
2817 &store_dyn,
2818 ),
2819 );
2820 let second = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2821 Duration::from_secs(60),
2822 Arc::new(RuntimeAuthLeaseHandle::new()),
2823 &store_dyn,
2824 );
2825
2826 store.block_next_oauth_persist();
2827 let first_consume = std::thread::spawn({
2828 let first = Arc::clone(&first);
2829 let state = state.clone();
2830 let target = target.clone();
2831 move || first.consume(&state, &target, provider, redirect_uri)
2832 });
2833 store.wait_for_blocked_oauth_persist();
2834
2835 second
2836 .consume(&state, &target, provider, redirect_uri)
2837 .expect("second authority wins durable consume race");
2838 store.release_blocked_oauth_persist();
2839 assert!(matches!(
2840 first_consume
2841 .join()
2842 .expect("first consume thread should not panic"),
2843 Err(OAuthFlowError::RegistryProjectionMissing {
2844 operation: "consume_oauth_browser_flow"
2845 })
2846 ));
2847 }
2848
2849 #[test]
2850 fn concurrent_persistent_device_consumes_require_fresh_durable_claim() {
2851 let store = Arc::new(FailingOAuthSnapshotStore::default());
2852 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2853 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2854 Duration::from_secs(60),
2855 Arc::new(RuntimeAuthLeaseHandle::new()),
2856 &store_dyn,
2857 );
2858 let target = target();
2859 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2860 let device_code = "race-device-code";
2861 creator
2862 .admit_device_code(
2863 target.clone(),
2864 provider,
2865 device_code.to_string(),
2866 Duration::from_secs(60),
2867 )
2868 .expect("device flow admitted");
2869 let first = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2870 Duration::from_secs(60),
2871 Arc::new(RuntimeAuthLeaseHandle::new()),
2872 &store_dyn,
2873 );
2874 let second = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2875 Duration::from_secs(60),
2876 Arc::new(RuntimeAuthLeaseHandle::new()),
2877 &store_dyn,
2878 );
2879 let first_poll = first
2880 .begin_device_code_poll(device_code, &target, provider)
2881 .expect("first authority begins poll");
2882 let second_poll = second
2883 .begin_device_code_poll(device_code, &target, provider)
2884 .expect("second authority begins poll");
2885
2886 store.block_next_oauth_persist();
2887 let first_consume = std::thread::spawn(move || first_poll.consume());
2888 store.wait_for_blocked_oauth_persist();
2889
2890 second_poll
2891 .consume()
2892 .expect("second authority wins durable device consume race");
2893 store.release_blocked_oauth_persist();
2894 assert!(matches!(
2895 first_consume
2896 .join()
2897 .expect("first device consume thread should not panic"),
2898 Err(OAuthFlowError::RegistryProjectionMissing {
2899 operation: "consume_oauth_device_flow"
2900 })
2901 ));
2902 }
2903
2904 #[test]
2905 fn concurrent_persistent_browser_admits_confirm_fresh_generated_capacity() {
2906 let store = Arc::new(FailingOAuthSnapshotStore::default());
2907 let first_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2908 let second_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2909 let first = Arc::new(
2910 RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2911 Duration::from_secs(60),
2912 1,
2913 Arc::new(RuntimeAuthLeaseHandle::new()),
2914 Some(Arc::downgrade(&first_store)),
2915 ),
2916 );
2917 let second = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2918 Duration::from_secs(60),
2919 1,
2920 Arc::new(RuntimeAuthLeaseHandle::new()),
2921 Some(Arc::downgrade(&second_store)),
2922 );
2923 let first_target = target();
2924 let second_target = alternate_target();
2925 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2926 let first_redirect_uri = "http://127.0.0.1/first-callback";
2927 let second_redirect_uri = "http://127.0.0.1/second-callback";
2928
2929 store.block_next_oauth_persist();
2930 let first_admit = std::thread::spawn({
2931 let first = Arc::clone(&first);
2932 let first_target = first_target.clone();
2933 move || {
2934 first.start(
2935 first_target,
2936 provider,
2937 first_redirect_uri.to_string(),
2938 "first-verifier".to_string(),
2939 )
2940 }
2941 });
2942 store.wait_for_blocked_oauth_persist();
2943
2944 second
2945 .start(
2946 second_target,
2947 provider,
2948 second_redirect_uri.to_string(),
2949 "second-verifier".to_string(),
2950 )
2951 .expect("second authority wins durable browser admission race");
2952 store.release_blocked_oauth_persist();
2953 assert!(matches!(
2954 first_admit
2955 .join()
2956 .expect("first browser admit thread should not panic"),
2957 Err(OAuthFlowError::LifecycleRejected {
2958 operation: "admit_oauth_browser_flow",
2959 ..
2960 })
2961 ));
2962 }
2963
2964 #[test]
2965 fn concurrent_persistent_device_admits_confirm_fresh_generated_capacity() {
2966 let store = Arc::new(FailingOAuthSnapshotStore::default());
2967 let first_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2968 let second_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2969 let first = Arc::new(
2970 RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2971 Duration::from_secs(60),
2972 1,
2973 Arc::new(RuntimeAuthLeaseHandle::new()),
2974 Some(Arc::downgrade(&first_store)),
2975 ),
2976 );
2977 let second = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2978 Duration::from_secs(60),
2979 1,
2980 Arc::new(RuntimeAuthLeaseHandle::new()),
2981 Some(Arc::downgrade(&second_store)),
2982 );
2983 let first_target = target();
2984 let second_target = alternate_target();
2985 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2986
2987 store.block_next_oauth_persist();
2988 let first_admit = std::thread::spawn({
2989 let first = Arc::clone(&first);
2990 let first_target = first_target.clone();
2991 move || {
2992 first.admit_device_code(
2993 first_target,
2994 provider,
2995 "first-device-code".to_string(),
2996 Duration::from_secs(60),
2997 )
2998 }
2999 });
3000 store.wait_for_blocked_oauth_persist();
3001
3002 second
3003 .admit_device_code(
3004 second_target,
3005 provider,
3006 "second-device-code".to_string(),
3007 Duration::from_secs(60),
3008 )
3009 .expect("second authority wins durable device admission race");
3010 store.release_blocked_oauth_persist();
3011 assert!(matches!(
3012 first_admit
3013 .join()
3014 .expect("first device admit thread should not panic"),
3015 Err(OAuthFlowError::LifecycleRejected {
3016 operation: "admit_oauth_device_flow",
3017 ..
3018 })
3019 ));
3020 }
3021
3022 #[test]
3023 fn concurrent_browser_admits_preserve_newer_durable_snapshot() {
3024 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3025 let store = Arc::new(FailingOAuthSnapshotStore::default());
3026 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3027 let authority = Arc::new(
3028 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3029 Duration::from_secs(60),
3030 lifecycle,
3031 &store_dyn,
3032 ),
3033 );
3034 let first_target = target();
3035 let second_target = alternate_target();
3036 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3037 let first_redirect_uri = "http://127.0.0.1/callback";
3038 let second_redirect_uri = "http://127.0.0.1/other-callback";
3039
3040 store.block_next_oauth_persist();
3041 let first_admit = std::thread::spawn({
3042 let authority = Arc::clone(&authority);
3043 let target = first_target.clone();
3044 move || {
3045 authority.start(
3046 target,
3047 provider,
3048 first_redirect_uri.to_string(),
3049 "verifier-1".to_string(),
3050 )
3051 }
3052 });
3053 store.wait_for_blocked_oauth_persist();
3054
3055 let (second_done_tx, second_done_rx) = std::sync::mpsc::channel();
3056 std::thread::spawn({
3057 let authority = Arc::clone(&authority);
3058 let target = second_target.clone();
3059 move || {
3060 let result = authority.start(
3061 target,
3062 provider,
3063 second_redirect_uri.to_string(),
3064 "verifier-2".to_string(),
3065 );
3066 let _ = second_done_tx.send(result);
3067 }
3068 });
3069 let second_before_release = second_done_rx.recv_timeout(Duration::from_millis(100)).ok();
3070 store.release_blocked_oauth_persist();
3071 first_admit
3072 .join()
3073 .expect("first admit thread should not panic")
3074 .expect("first browser flow admitted");
3075 let second_state = second_before_release
3076 .unwrap_or_else(|| {
3077 second_done_rx
3078 .recv_timeout(Duration::from_secs(1))
3079 .expect("second admit should finish after first durable write is released")
3080 })
3081 .expect("second browser flow admitted");
3082
3083 let snapshot_json = store
3084 .load_auth_oauth_flow_snapshot()
3085 .expect("durable OAuth snapshot loads")
3086 .expect("durable OAuth snapshot exists");
3087 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
3088 .expect("durable OAuth snapshot decodes");
3089 assert!(
3090 snapshot
3091 .browser
3092 .iter()
3093 .any(|flow| flow.state == second_state),
3094 "durable snapshot must retain flow admitted by a concurrent newer write"
3095 );
3096
3097 let restarted_lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3098 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3099 Duration::from_secs(60),
3100 restarted_lifecycle,
3101 &store_dyn,
3102 );
3103 let record = restarted
3104 .consume(&second_state, &second_target, provider, second_redirect_uri)
3105 .expect("newer durable flow should survive restart");
3106 assert_eq!(record.pkce_verifier, "verifier-2");
3107 }
3108
3109 #[test]
3110 fn browser_consume_persistence_failure_keeps_flow_retryable() {
3111 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3112 let store = Arc::new(FailingOAuthSnapshotStore::default());
3113 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3114 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3115 Duration::from_secs(60),
3116 lifecycle.clone(),
3117 &store_dyn,
3118 );
3119 let target = target();
3120 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3121 let redirect_uri = "http://127.0.0.1/callback";
3122 let state = authority
3123 .start(
3124 target.clone(),
3125 provider,
3126 redirect_uri.to_string(),
3127 "verifier".to_string(),
3128 )
3129 .expect("browser flow admitted");
3130
3131 store.fail_oauth_persist();
3132 assert!(matches!(
3133 authority.consume(&state, &target, provider, redirect_uri),
3134 Err(OAuthFlowError::PersistenceFailed { .. })
3135 ));
3136 assert!(lifecycle.has_oauth_browser_flow_for_test(&target, &state));
3137
3138 store.allow_oauth_persist();
3139 authority
3140 .consume(&state, &target, provider, redirect_uri)
3141 .expect("failed durable consume remains retryable");
3142 }
3143
3144 #[test]
3145 fn device_consume_persistence_failure_keeps_flow_retryable() {
3146 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3147 let store = Arc::new(FailingOAuthSnapshotStore::default());
3148 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3149 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3150 Duration::from_secs(60),
3151 lifecycle.clone(),
3152 &store_dyn,
3153 );
3154 let target = target();
3155 let provider = OAuthProviderIdentity::GoogleCodeAssist;
3156 let device_code = "provider-device-code";
3157 authority
3158 .admit_device_code(
3159 target.clone(),
3160 provider,
3161 device_code.to_string(),
3162 Duration::from_secs(60),
3163 )
3164 .expect("device flow admitted");
3165 let poll = authority
3166 .begin_device_code_poll(device_code, &target, provider)
3167 .expect("device poll begins");
3168
3169 store.fail_oauth_persist();
3170 assert!(matches!(
3171 poll.consume(),
3172 Err(OAuthFlowError::PersistenceFailed { .. })
3173 ));
3174 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3175
3176 store.allow_oauth_persist();
3177 let retry = authority
3178 .begin_device_code_poll(device_code, &target, provider)
3179 .expect("failed durable consume keeps device flow retryable");
3180 retry
3181 .consume()
3182 .expect("retry consumes after durable persistence recovers");
3183 }
3184
3185 #[test]
3186 fn release_persistence_failure_keeps_released_flows_retryable() {
3187 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3188 let store = Arc::new(FailingOAuthSnapshotStore::default());
3189 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3190 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3191 Duration::from_secs(60),
3192 lifecycle.clone(),
3193 &store_dyn,
3194 );
3195 let target = target();
3196 let lease_key = LeaseKey::from_auth_binding(&target);
3197 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3198 let redirect_uri = "http://127.0.0.1/callback";
3199 let state = authority
3200 .start(
3201 target.clone(),
3202 provider,
3203 redirect_uri.to_string(),
3204 "verifier".to_string(),
3205 )
3206 .expect("browser flow admitted");
3207
3208 store.fail_oauth_persist();
3209 assert!(
3210 lifecycle.release_lease(&lease_key).is_err(),
3211 "release should fail closed when durable OAuth cleanup cannot persist"
3212 );
3213 assert!(lifecycle.has_oauth_browser_flow_for_test(&target, &state));
3214
3215 store.allow_oauth_persist();
3216 authority
3217 .consume(&state, &target, provider, redirect_uri)
3218 .expect("failed durable release leaves browser flow retryable");
3219 }
3220
3221 #[test]
3222 fn stale_release_persistence_failure_does_not_install_released_authority() {
3223 let releasing_lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3224 let store = Arc::new(FailingOAuthSnapshotStore::default());
3225 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3226 let releasing_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3227 Duration::from_secs(60),
3228 releasing_lifecycle.clone(),
3229 &store_dyn,
3230 );
3231 let admitting_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3232 Duration::from_secs(60),
3233 Arc::new(RuntimeAuthLeaseHandle::new()),
3234 &store_dyn,
3235 );
3236 let target = target();
3237 let lease_key = LeaseKey::from_auth_binding(&target);
3238 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3239 let redirect_uri = "http://127.0.0.1/callback";
3240 let state = admitting_authority
3241 .start(
3242 target.clone(),
3243 provider,
3244 redirect_uri.to_string(),
3245 "verifier".to_string(),
3246 )
3247 .expect("other authority admits browser flow");
3248 assert!(
3249 !releasing_lifecycle.has_oauth_browser_flow_for_test(&target, &state),
3250 "releasing authority starts stale and has no local machine membership"
3251 );
3252
3253 store.fail_oauth_persist();
3254 assert!(
3255 releasing_lifecycle.release_lease(&lease_key).is_err(),
3256 "release should fail closed when stale durable OAuth cleanup cannot persist"
3257 );
3258 assert_eq!(
3259 releasing_lifecycle.snapshot(&lease_key).phase,
3260 None,
3261 "failed stale release must not synthesize a local AuthMachine authority"
3262 );
3263
3264 store.allow_oauth_persist();
3265 releasing_authority
3266 .consume(&state, &target, provider, redirect_uri)
3267 .expect("failed stale durable release must leave browser flow retryable");
3268 }
3269
3270 #[cfg(feature = "sqlite-store")]
3271 #[test]
3272 fn persistent_release_prunes_durable_flows_from_stale_authority() {
3273 let temp_dir = tempfile::tempdir().expect("tempdir");
3274 let store_path = temp_dir.path().join("runtime.sqlite");
3275 let releasing_store: Arc<dyn RuntimeStore> =
3276 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
3277 let admitting_store: Arc<dyn RuntimeStore> =
3278 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
3279 let releasing_lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3280 let releasing_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3281 Duration::from_secs(60),
3282 releasing_lifecycle.clone(),
3283 &releasing_store,
3284 );
3285 let admitting_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3286 Duration::from_secs(60),
3287 Arc::new(RuntimeAuthLeaseHandle::new()),
3288 &admitting_store,
3289 );
3290 let target = target();
3291 let lease_key = LeaseKey::from_auth_binding(&target);
3292 let browser_provider = OAuthProviderIdentity::OpenAiChatGpt;
3293 let redirect_uri = "http://127.0.0.1/callback";
3294 let browser_state = admitting_authority
3295 .start(
3296 target.clone(),
3297 browser_provider,
3298 redirect_uri.to_string(),
3299 "browser-verifier".to_string(),
3300 )
3301 .expect("other authority admits browser flow");
3302 let device_provider = OAuthProviderIdentity::GoogleCodeAssist;
3303 let device_code = "released-device-code";
3304 admitting_authority
3305 .admit_device_code(
3306 target.clone(),
3307 device_provider,
3308 device_code.to_string(),
3309 Duration::from_secs(60),
3310 )
3311 .expect("other authority admits device flow");
3312 assert!(
3313 !releasing_lifecycle.has_oauth_browser_flow_for_test(&target, &browser_state),
3314 "releasing authority starts stale and does not know the browser flow locally"
3315 );
3316 assert!(
3317 !releasing_lifecycle.has_oauth_device_flow_for_test(&target, device_code),
3318 "releasing authority starts stale and does not know the device flow locally"
3319 );
3320
3321 releasing_lifecycle
3322 .release_lease(&lease_key)
3323 .expect("stale release succeeds");
3324
3325 let restarted_store: Arc<dyn RuntimeStore> =
3326 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
3327 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3328 Duration::from_secs(60),
3329 Arc::new(RuntimeAuthLeaseHandle::new()),
3330 &restarted_store,
3331 );
3332 let browser_after_release =
3333 restarted.consume(&browser_state, &target, browser_provider, redirect_uri);
3334 let device_after_release =
3335 restarted.verify_device_code(device_code, &target, device_provider);
3336 assert!(
3337 matches!(
3338 browser_after_release,
3339 Err(OAuthFlowError::LifecycleRejected {
3340 operation: "verify_oauth_browser_flow",
3341 ..
3342 })
3343 ),
3344 "release from a stale authority must prune durable browser flow, got {browser_after_release:?}"
3345 );
3346 assert!(
3347 matches!(
3348 device_after_release,
3349 Err(OAuthFlowError::LifecycleRejected {
3350 operation: "verify_oauth_device_flow",
3351 ..
3352 })
3353 ),
3354 "release from a stale authority must prune durable device flow, got {device_after_release:?}"
3355 );
3356 drop(releasing_authority);
3357 }
3358
3359 #[test]
3360 fn oauth_flow_membership_does_not_advance_credential_generation() {
3361 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3362 let authority =
3363 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3364 let target = target();
3365 let lease_key = LeaseKey::from_auth_binding(&target);
3366 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3367 let redirect_uri = "http://127.0.0.1/callback";
3368 let transition = lifecycle
3369 .acquire_lease(&lease_key, 4_200)
3370 .expect("credential lifecycle acquired");
3371
3372 let state = authority
3373 .start(
3374 target.clone(),
3375 provider,
3376 redirect_uri.to_string(),
3377 "verifier".to_string(),
3378 )
3379 .expect("browser flow admitted");
3380 authority
3381 .verify(&state, &target, provider, redirect_uri)
3382 .expect("browser flow verifies");
3383 authority
3384 .consume(&state, &target, provider, redirect_uri)
3385 .expect("browser flow consumes");
3386
3387 let snapshot = lifecycle.snapshot(&lease_key);
3388 assert_eq!(snapshot.generation, transition.generation());
3389 assert_eq!(
3390 snapshot.credential_published_at_millis,
3391 transition.credential_published_at_millis()
3392 );
3393 }
3394
3395 #[test]
3396 fn global_browser_expiry_preserves_reauth_required_phase() {
3397 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3398 let authority = RuntimeOAuthFlowHandle::new_with_auth_lease(
3399 Duration::from_millis(1),
3400 lifecycle.clone(),
3401 );
3402 let target = target();
3403 let other_target = alternate_target();
3404 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3405 let redirect_uri = "http://127.0.0.1/callback";
3406
3407 let expired_state = authority
3408 .start(
3409 target.clone(),
3410 provider,
3411 redirect_uri.to_string(),
3412 "verifier-old".to_string(),
3413 )
3414 .expect("browser flow admitted");
3415 assert_eq!(
3416 snapshot_phase(&lifecycle, &target),
3417 Some(AuthLeasePhase::ReauthRequired)
3418 );
3419 std::thread::sleep(Duration::from_millis(10));
3420
3421 authority
3422 .start(
3423 other_target,
3424 provider,
3425 redirect_uri.to_string(),
3426 "verifier-new".to_string(),
3427 )
3428 .expect("new browser flow admitted after pruning expired flow");
3429
3430 assert!(
3431 !lifecycle.has_oauth_browser_flow_for_test(&target, &expired_state),
3432 "passive registry expiry must remove stale AuthMachine browser membership"
3433 );
3434 assert_eq!(
3435 snapshot_phase(&lifecycle, &target),
3436 Some(AuthLeasePhase::ReauthRequired),
3437 "global OAuth expiry cleanup must not change credential lifecycle truth"
3438 );
3439 }
3440
3441 #[test]
3442 fn browser_passive_expiry_clears_lifecycle_membership_on_next_admit() {
3443 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3444 let authority = RuntimeOAuthFlowHandle::new_with_auth_lease(
3445 Duration::from_millis(1),
3446 lifecycle.clone(),
3447 );
3448 let target = target();
3449 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3450 let redirect_uri = "http://127.0.0.1/callback";
3451
3452 let expired_state = authority
3453 .start(
3454 target.clone(),
3455 provider,
3456 redirect_uri.to_string(),
3457 "verifier-old".to_string(),
3458 )
3459 .expect("browser flow admitted");
3460 assert!(lifecycle.has_oauth_browser_flow_for_test(&target, &expired_state));
3461 std::thread::sleep(Duration::from_millis(10));
3462
3463 authority
3464 .start(
3465 target.clone(),
3466 provider,
3467 redirect_uri.to_string(),
3468 "verifier-new".to_string(),
3469 )
3470 .expect("new browser flow admitted after pruning expired flow");
3471
3472 assert!(
3473 !lifecycle.has_oauth_browser_flow_for_test(&target, &expired_state),
3474 "passive registry expiry must remove stale AuthMachine browser membership"
3475 );
3476 }
3477
3478 #[test]
3479 fn device_admit_rejects_registry_pruned_canonical_membership() {
3480 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3481 let authority =
3482 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3483 let target = target();
3484 let provider = OAuthProviderIdentity::GoogleCodeAssist;
3485 let device_code = "provider-device-code";
3486
3487 authority
3488 .admit_device_code(
3489 target.clone(),
3490 provider,
3491 device_code.to_string(),
3492 Duration::from_secs(60),
3493 )
3494 .expect("device flow admitted");
3495 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3496 authority
3497 .registry
3498 .expire_device_code(device_code, &target, provider)
3499 .expect("test removes registry record without lifecycle cleanup");
3500 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3501
3502 assert!(matches!(
3503 authority.admit_device_code(
3504 target.clone(),
3505 provider,
3506 device_code.to_string(),
3507 Duration::from_secs(60),
3508 ),
3509 Err(OAuthFlowError::LifecycleRejected {
3510 operation: "admit_oauth_device_flow",
3511 ..
3512 })
3513 ));
3514
3515 assert!(
3516 lifecycle.has_oauth_device_flow_for_test(&target, device_code),
3517 "registry-only loss must not expire canonical AuthMachine device membership"
3518 );
3519 assert!(matches!(
3520 authority.verify_device_code(device_code, &target, provider),
3521 Err(OAuthFlowError::RegistryProjectionMissing {
3522 operation: "verify_oauth_device_flow"
3523 })
3524 ));
3525 assert!(matches!(
3526 authority.begin_device_code_poll(device_code, &target, provider),
3527 Err(OAuthFlowError::RegistryProjectionMissing {
3528 operation: "begin_oauth_device_poll"
3529 })
3530 ));
3531 assert!(
3532 lifecycle.has_oauth_device_flow_for_test(&target, device_code),
3533 "missing process-local device payload must fail closed without removing the flow"
3534 );
3535 }
3536
3537 #[test]
3538 fn duplicate_device_admit_preserves_active_lifecycle_membership() {
3539 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3540 let authority =
3541 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3542 let target = target();
3543 let provider = OAuthProviderIdentity::GoogleCodeAssist;
3544 let device_code = "provider-device-code";
3545
3546 authority
3547 .admit_device_code(
3548 target.clone(),
3549 provider,
3550 device_code.to_string(),
3551 Duration::from_secs(60),
3552 )
3553 .expect("device flow admitted");
3554 let duplicate = authority.admit_device_code(
3555 target.clone(),
3556 provider,
3557 device_code.to_string(),
3558 Duration::from_secs(60),
3559 );
3560
3561 assert!(matches!(
3562 duplicate,
3563 Err(OAuthFlowError::LifecycleRejected {
3564 operation: "admit_oauth_device_flow",
3565 ..
3566 })
3567 ));
3568 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3569 authority
3570 .begin_device_code_poll(device_code, &target, provider)
3571 .expect("duplicate admit must not orphan active lifecycle membership");
3572 }
3573
3574 #[test]
3575 fn stale_registry_payloads_do_not_block_authmachine_admission_after_lifecycle_release() {
3576 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3577 let authority = RuntimeOAuthFlowHandle::new_with_capacity_and_auth_lease(
3578 Duration::from_secs(60),
3579 1,
3580 lifecycle.clone(),
3581 );
3582 let target = target();
3583 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3584
3585 authority
3586 .start(
3587 target.clone(),
3588 provider,
3589 "http://127.0.0.1/callback".to_string(),
3590 "verifier-1".to_string(),
3591 )
3592 .expect("first browser flow admitted");
3593 lifecycle
3594 .release_lease(&LeaseKey::from_auth_binding(&target))
3595 .expect("credential lifecycle release succeeds");
3596
3597 authority
3598 .start(
3599 alternate_target(),
3600 provider,
3601 "http://127.0.0.1/other-callback".to_string(),
3602 "verifier-2".to_string(),
3603 )
3604 .expect("AuthMachine release must clear stale registry payload capacity");
3605 }
3606
3607 #[test]
3608 fn release_observer_does_not_prune_flow_admitted_after_release_acceptance() {
3609 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3610 let authority = Arc::new(RuntimeOAuthFlowHandle::new_with_auth_lease(
3611 Duration::from_secs(60),
3612 lifecycle.clone(),
3613 ));
3614 let target = target_with_binding("release_after_accept_openai");
3615 let lease_key = LeaseKey::from_auth_binding(&target);
3616 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3617 let redirect_uri = "http://127.0.0.1/callback";
3618 let old_state = authority
3619 .start(
3620 target.clone(),
3621 provider,
3622 redirect_uri.to_string(),
3623 "old-verifier".to_string(),
3624 )
3625 .expect("old browser flow admitted");
3626 let new_state = Arc::new(std::sync::Mutex::new(None));
3627 let new_state_for_hook = Arc::clone(&new_state);
3628 let authority_for_hook = Arc::clone(&authority);
3629 let target_for_hook = target.clone();
3630 let lease_key_for_hook = lease_key.clone();
3631 let _hook_guard = crate::handles::auth_lease::install_release_after_accept_hook_for_test(
3632 Arc::new(move |released_key| {
3633 if released_key != &lease_key_for_hook {
3634 return;
3635 }
3636 let admitted = authority_for_hook
3637 .start(
3638 target_for_hook.clone(),
3639 provider,
3640 redirect_uri.to_string(),
3641 "new-verifier".to_string(),
3642 )
3643 .expect("new browser flow admitted after release acceptance");
3644 *new_state_for_hook
3645 .lock()
3646 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(admitted);
3647 }),
3648 );
3649
3650 lifecycle
3651 .release_lease(&lease_key)
3652 .expect("credential lifecycle release succeeds");
3653
3654 let new_state = new_state
3655 .lock()
3656 .unwrap_or_else(std::sync::PoisonError::into_inner)
3657 .clone()
3658 .expect("hook admitted replacement flow");
3659 let flow = authority
3660 .consume(&new_state, &target, provider, redirect_uri)
3661 .expect("release observer must not prune newly admitted flow");
3662 assert_eq!(flow.pkce_verifier, "new-verifier");
3663 assert!(matches!(
3664 authority.consume(&old_state, &target, provider, redirect_uri),
3665 Err(OAuthFlowError::LifecycleRejected {
3666 operation: "verify_oauth_browser_flow",
3667 ..
3668 })
3669 ));
3670 }
3671
3672 #[test]
3673 fn release_precommit_admission_waits_for_release_commit() {
3674 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3675 let authority = Arc::new(RuntimeOAuthFlowHandle::new_with_auth_lease(
3676 Duration::from_secs(60),
3677 lifecycle.clone(),
3678 ));
3679 let target = target_with_binding("release_before_commit_openai");
3680 let lease_key = LeaseKey::from_auth_binding(&target);
3681 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3682 let redirect_uri = "http://127.0.0.1/callback";
3683 let old_state = authority
3684 .start(
3685 target.clone(),
3686 provider,
3687 redirect_uri.to_string(),
3688 "old-verifier".to_string(),
3689 )
3690 .expect("old browser flow admitted");
3691 let (done_tx, done_rx) = mpsc::channel();
3692 let admission_finished = Arc::new(AtomicBool::new(false));
3693 let admission_finished_for_hook = Arc::clone(&admission_finished);
3694 let authority_for_hook = Arc::clone(&authority);
3695 let target_for_hook = target.clone();
3696 let lease_key_for_hook = lease_key.clone();
3697 let _hook_guard = crate::handles::auth_lease::install_release_before_commit_hook_for_test(
3698 Arc::new(move |released_key| {
3699 if released_key != &lease_key_for_hook {
3700 return;
3701 }
3702 let authority_for_thread = Arc::clone(&authority_for_hook);
3703 let target_for_thread = target_for_hook.clone();
3704 let done_tx = done_tx.clone();
3705 let admission_finished_for_thread = Arc::clone(&admission_finished_for_hook);
3706 std::thread::spawn(move || {
3707 let admitted = authority_for_thread.start(
3708 target_for_thread,
3709 provider,
3710 redirect_uri.to_string(),
3711 "new-verifier".to_string(),
3712 );
3713 admission_finished_for_thread.store(true, Ordering::Release);
3714 let _ = done_tx.send(admitted);
3715 });
3716 std::thread::sleep(Duration::from_millis(50));
3717 assert!(
3718 !admission_finished_for_hook.load(Ordering::Acquire),
3719 "OAuth admission must wait until release commits"
3720 );
3721 }),
3722 );
3723
3724 lifecycle
3725 .release_lease(&lease_key)
3726 .expect("credential lifecycle release succeeds");
3727
3728 let new_state = done_rx
3729 .recv_timeout(Duration::from_secs(1))
3730 .expect("pre-commit admission should finish after release commits")
3731 .expect("new browser flow admitted after release commit");
3732 let flow = authority
3733 .consume(&new_state, &target, provider, redirect_uri)
3734 .expect("release observer must not prune newly admitted flow");
3735 assert_eq!(flow.pkce_verifier, "new-verifier");
3736 assert!(matches!(
3737 authority.consume(&old_state, &target, provider, redirect_uri),
3738 Err(OAuthFlowError::LifecycleRejected {
3739 operation: "verify_oauth_browser_flow",
3740 ..
3741 })
3742 ));
3743 }
3744
3745 #[test]
3746 fn browser_capacity_rejection_comes_from_authmachine_lifecycle() {
3747 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3748 let authority = RuntimeOAuthFlowHandle::new_with_capacity_and_auth_lease(
3749 Duration::from_secs(60),
3750 1,
3751 lifecycle,
3752 );
3753 let target = target();
3754 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3755 let redirect_uri = "http://127.0.0.1/callback";
3756
3757 authority
3758 .start(
3759 target.clone(),
3760 provider,
3761 redirect_uri.to_string(),
3762 "verifier-1".to_string(),
3763 )
3764 .expect("first browser flow admitted");
3765
3766 assert!(matches!(
3767 authority.start(
3768 alternate_target(),
3769 provider,
3770 "http://127.0.0.1/other-callback".to_string(),
3771 "verifier-2".to_string(),
3772 ),
3773 Err(OAuthFlowError::LifecycleRejected {
3774 operation: "admit_oauth_browser_flow",
3775 ..
3776 })
3777 ));
3778 }
3779
3780 #[test]
3781 fn browser_provider_mismatch_rejection_comes_from_authmachine_lifecycle() {
3782 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3783 let authority =
3784 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle);
3785 let target = target();
3786 let redirect_uri = "http://127.0.0.1/callback";
3787 let state = authority
3788 .start(
3789 target.clone(),
3790 OAuthProviderIdentity::OpenAiChatGpt,
3791 redirect_uri.to_string(),
3792 "verifier".to_string(),
3793 )
3794 .expect("browser flow admitted");
3795
3796 assert!(matches!(
3797 authority.verify(
3798 &state,
3799 &target,
3800 OAuthProviderIdentity::GoogleCodeAssist,
3801 redirect_uri,
3802 ),
3803 Err(OAuthFlowError::LifecycleRejected {
3804 operation: "verify_oauth_browser_flow",
3805 ..
3806 })
3807 ));
3808 }
3809
3810 #[test]
3823 fn release_lease_terminally_cancels_pending_browser_flow() {
3824 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3825 let authority =
3826 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3827 let target = target();
3828 let lease_key = LeaseKey::from_auth_binding(&target);
3829 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3830 let redirect_uri = "http://127.0.0.1/callback";
3831
3832 let state = authority
3833 .start(
3834 target.clone(),
3835 provider,
3836 redirect_uri.to_string(),
3837 "verifier".to_string(),
3838 )
3839 .expect("browser flow admitted");
3840
3841 lifecycle
3842 .release_lease(&lease_key)
3843 .expect("release with a pending flow must drain it, not fail");
3844 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3845 assert!(
3846 !lifecycle.has_oauth_browser_flow_for_test(&target, &state),
3847 "pending flow must be terminally cancelled by the release drain"
3848 );
3849
3850 assert!(matches!(
3852 authority.consume(&state, &target, provider, redirect_uri),
3853 Err(OAuthFlowError::LifecycleRejected { .. })
3854 ));
3855
3856 lifecycle
3860 .apply_oauth_input(
3861 &target,
3862 auth_dsl::AuthMachineInput::ExpireOAuthBrowserFlow {
3863 flow_id: state.clone(),
3864 },
3865 "late_prune_expire_browser",
3866 false,
3867 )
3868 .expect("post-release expire must be a benign no-op");
3869 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3870 }
3871
3872 #[test]
3875 fn release_lease_terminally_cancels_pending_device_flow() {
3876 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3877 let authority =
3878 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3879 let target = target();
3880 let lease_key = LeaseKey::from_auth_binding(&target);
3881 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3882
3883 authority
3884 .admit_device_code(
3885 target.clone(),
3886 provider,
3887 "device-code-1".to_string(),
3888 Duration::from_secs(60),
3889 )
3890 .expect("device flow admitted");
3891
3892 lifecycle
3893 .release_lease(&lease_key)
3894 .expect("release with a pending device flow must drain it, not fail");
3895 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3896 assert!(
3897 !lifecycle.has_oauth_device_flow_for_test(&target, "device-code-1"),
3898 "pending device flow must be terminally cancelled by the release drain"
3899 );
3900
3901 lifecycle
3903 .expire_device_flow(&target, "device-code-1")
3904 .expect("post-release device expire must be a benign no-op");
3905 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3906 }
3907
3908 #[test]
3913 fn post_release_oauth_observations_through_shell_dispatch_are_benign() {
3914 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3915 let target = target();
3916 let lease_key = LeaseKey::from_auth_binding(&target);
3917
3918 lifecycle
3919 .release_lease(&lease_key)
3920 .expect("releasing an empty lease succeeds");
3921 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3924
3925 lifecycle
3926 .apply_oauth_input(
3927 &target,
3928 auth_dsl::AuthMachineInput::ExpireOAuthBrowserFlow {
3929 flow_id: "ghost-browser".to_string(),
3930 },
3931 "late_prune_expire_browser",
3932 false,
3933 )
3934 .expect("post-release browser expire must be a benign no-op");
3935 lifecycle
3936 .expire_device_flow(&target, "ghost-device")
3937 .expect("post-release device expire must be a benign no-op");
3938 lifecycle
3939 .finish_device_poll(&target, "ghost-poll")
3940 .expect("post-release poll finish must be a benign no-op");
3941 lifecycle
3942 .confirm_oauth_durable_admission(&target, 0, 16, "late_confirm_oauth_durable_admission")
3943 .expect("post-release durable-admission confirmation must be a benign no-op");
3944
3945 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3947 }
3948
3949 #[test]
3954 fn late_prune_expire_at_release_acceptance_is_benign() {
3955 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3956 let authority =
3957 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3958 let target = target();
3959 let lease_key = LeaseKey::from_auth_binding(&target);
3960 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3961
3962 let state = authority
3963 .start(
3964 target.clone(),
3965 provider,
3966 "http://127.0.0.1/callback".to_string(),
3967 "verifier".to_string(),
3968 )
3969 .expect("browser flow admitted");
3970
3971 let hook_result: Arc<StdMutex<Option<Result<(), String>>>> = Arc::new(StdMutex::new(None));
3972 let hook_result_for_hook = Arc::clone(&hook_result);
3973 let lifecycle_for_hook = Arc::clone(&lifecycle);
3974 let target_for_hook = target.clone();
3975 let lease_key_for_hook = lease_key.clone();
3976 let flow_for_hook = state.clone();
3977 let _hook_guard = crate::handles::auth_lease::install_release_after_accept_hook_for_test(
3978 Arc::new(move |released_key| {
3979 if released_key != &lease_key_for_hook {
3980 return;
3981 }
3982 let outcome = lifecycle_for_hook
3983 .apply_oauth_input(
3984 &target_for_hook,
3985 auth_dsl::AuthMachineInput::ExpireOAuthBrowserFlow {
3986 flow_id: flow_for_hook.clone(),
3987 },
3988 "late_prune_expire_at_release_acceptance",
3989 false,
3990 )
3991 .map_err(|err| err.to_string());
3992 *hook_result_for_hook
3993 .lock()
3994 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(outcome);
3995 }),
3996 );
3997
3998 lifecycle
3999 .release_lease(&lease_key)
4000 .expect("release with a pending flow must drain it, not fail");
4001
4002 let outcome = hook_result
4003 .lock()
4004 .unwrap_or_else(std::sync::PoisonError::into_inner)
4005 .clone()
4006 .expect("release acceptance hook must have fired");
4007 assert_eq!(
4008 outcome,
4009 Ok(()),
4010 "expire fired at release acceptance must be a benign no-op"
4011 );
4012 assert_eq!(snapshot_phase(&lifecycle, &target), None);
4013 }
4014}