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
2035 fn snapshot_phase(
2036 lifecycle: &RuntimeAuthLeaseHandle,
2037 target: &AuthBindingRef,
2038 ) -> Option<AuthLeasePhase> {
2039 lifecycle
2040 .snapshot(&LeaseKey::from_auth_binding(target))
2041 .phase
2042 }
2043
2044 #[test]
2045 fn browser_flow_only_machine_stays_reauth_required_until_credentials_commit() {
2046 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2047 let authority =
2048 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
2049 let target = target();
2050 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2051 let redirect_uri = "http://127.0.0.1/callback";
2052
2053 let state = authority
2054 .start(
2055 target.clone(),
2056 provider,
2057 redirect_uri.to_string(),
2058 "verifier".to_string(),
2059 )
2060 .expect("browser flow admitted");
2061 assert_eq!(
2062 snapshot_phase(&lifecycle, &target),
2063 Some(AuthLeasePhase::ReauthRequired)
2064 );
2065
2066 authority
2067 .verify(&state, &target, provider, redirect_uri)
2068 .expect("browser flow verifies");
2069 assert_eq!(
2070 snapshot_phase(&lifecycle, &target),
2071 Some(AuthLeasePhase::ReauthRequired)
2072 );
2073
2074 authority
2075 .consume(&state, &target, provider, redirect_uri)
2076 .expect("browser flow consumes");
2077 assert_eq!(
2078 snapshot_phase(&lifecycle, &target),
2079 Some(AuthLeasePhase::ReauthRequired)
2080 );
2081 }
2082
2083 #[test]
2084 fn missing_browser_projection_cannot_overwrite_authmachine_flow() {
2085 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2086 let authority =
2087 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
2088 let target = target();
2089 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2090 let redirect_uri = "http://127.0.0.1/callback";
2091 let state = authority
2092 .start(
2093 target.clone(),
2094 provider,
2095 redirect_uri.to_string(),
2096 "verifier".to_string(),
2097 )
2098 .expect("browser flow admitted");
2099
2100 authority
2101 .registry
2102 .consume(&state, &target, provider, redirect_uri)
2103 .expect("test removes only the local registry payload");
2104
2105 assert!(matches!(
2106 authority.verify(&state, &target, provider, redirect_uri),
2107 Err(OAuthFlowError::RegistryProjectionMissing {
2108 operation: "verify_oauth_browser_flow"
2109 })
2110 ));
2111 assert!(
2112 lifecycle.has_oauth_browser_flow_for_test(&target, &state),
2113 "missing process-local registry payload must not expire canonical AuthMachine membership"
2114 );
2115
2116 assert!(matches!(
2117 authority.consume(&state, &target, provider, redirect_uri),
2118 Err(OAuthFlowError::RegistryProjectionMissing {
2119 operation: "consume_oauth_browser_flow"
2120 })
2121 ));
2122 assert!(
2123 lifecycle.has_oauth_browser_flow_for_test(&target, &state),
2124 "terminal consume must fail closed instead of converting payload loss to not-found"
2125 );
2126 assert_eq!(
2127 snapshot_phase(&lifecycle, &target),
2128 Some(AuthLeasePhase::ReauthRequired)
2129 );
2130 }
2131
2132 #[test]
2133 fn missing_device_poll_projection_cannot_overwrite_authmachine_flow() {
2134 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2135 let authority =
2136 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
2137 let target = target();
2138 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2139 let device_code = "provider-device-code";
2140
2141 authority
2142 .admit_device_code(
2143 target.clone(),
2144 provider,
2145 device_code.to_string(),
2146 Duration::from_secs(60),
2147 )
2148 .expect("device flow admitted");
2149 let poll = authority
2150 .begin_device_code_poll(device_code, &target, provider)
2151 .expect("device poll begins");
2152 authority
2153 .registry
2154 .expire_device_code(device_code, &target, provider)
2155 .expect("test removes only the local registry payload");
2156
2157 assert!(matches!(
2158 poll.consume(),
2159 Err(OAuthFlowError::RegistryProjectionMissing {
2160 operation: "consume_oauth_device_flow"
2161 })
2162 ));
2163 assert!(
2164 lifecycle.has_oauth_device_flow_for_test(&target, device_code),
2165 "missing process-local poll payload must not expire canonical AuthMachine membership"
2166 );
2167 }
2168
2169 #[test]
2170 fn browser_admit_persistence_failure_rolls_back_unreturned_flow() {
2171 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2172 let store = Arc::new(FailingOAuthSnapshotStore::default());
2173 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2174 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2175 Duration::from_secs(60),
2176 lifecycle.clone(),
2177 &store_dyn,
2178 );
2179 let failed_target = target();
2180 let successful_target = alternate_target();
2181 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2182 let redirect_uri = "http://127.0.0.1/callback";
2183
2184 store.fail_oauth_persist();
2185 assert!(matches!(
2186 authority.start(
2187 failed_target.clone(),
2188 provider,
2189 redirect_uri.to_string(),
2190 "failed-verifier".to_string(),
2191 ),
2192 Err(OAuthFlowError::PersistenceFailed { .. })
2193 ));
2194 assert!(
2195 authority
2196 .registry
2197 .snapshot_for_persistence(current_time_millis())
2198 .browser
2199 .is_empty(),
2200 "failed browser admission must not leave an unreturned registry payload"
2201 );
2202
2203 store.allow_oauth_persist();
2204 let successful_state = authority
2205 .start(
2206 successful_target,
2207 provider,
2208 redirect_uri.to_string(),
2209 "successful-verifier".to_string(),
2210 )
2211 .expect("subsequent browser admit persists after store recovers");
2212 let snapshot_json = store
2213 .load_auth_oauth_flow_snapshot()
2214 .expect("durable OAuth snapshot loads")
2215 .expect("durable OAuth snapshot exists");
2216 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2217 .expect("durable OAuth snapshot decodes");
2218 assert_eq!(
2219 snapshot
2220 .browser
2221 .iter()
2222 .map(|flow| flow.state.as_str())
2223 .collect::<Vec<_>>(),
2224 vec![successful_state.as_str()],
2225 "a later successful admit must not persist a previously failed unreturned flow"
2226 );
2227 }
2228
2229 #[test]
2230 fn device_admit_persistence_failure_rolls_back_unreturned_flow() {
2231 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
2232 let store = Arc::new(FailingOAuthSnapshotStore::default());
2233 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2234 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2235 Duration::from_secs(60),
2236 lifecycle.clone(),
2237 &store_dyn,
2238 );
2239 let target = target();
2240 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2241 let failed_device_code = "failed-device-code";
2242 let successful_device_code = "successful-device-code";
2243
2244 store.fail_oauth_persist();
2245 assert!(matches!(
2246 authority.admit_device_code(
2247 target.clone(),
2248 provider,
2249 failed_device_code.to_string(),
2250 Duration::from_secs(60),
2251 ),
2252 Err(OAuthFlowError::PersistenceFailed { .. })
2253 ));
2254 assert!(matches!(
2255 authority
2256 .registry
2257 .verify_device_code(failed_device_code, &target, provider),
2258 Err(OAuthFlowError::Missing)
2259 ));
2260 assert!(!lifecycle.has_oauth_device_flow_for_test(&target, failed_device_code));
2261
2262 store.allow_oauth_persist();
2263 authority
2264 .admit_device_code(
2265 target,
2266 provider,
2267 successful_device_code.to_string(),
2268 Duration::from_secs(60),
2269 )
2270 .expect("subsequent device admit persists after store recovers");
2271 let snapshot_json = store
2272 .load_auth_oauth_flow_snapshot()
2273 .expect("durable OAuth snapshot loads")
2274 .expect("durable OAuth snapshot exists");
2275 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2276 .expect("durable OAuth snapshot decodes");
2277 assert_eq!(
2278 snapshot
2279 .device
2280 .iter()
2281 .map(|flow| flow.device_code.as_str())
2282 .collect::<Vec<_>>(),
2283 vec![successful_device_code],
2284 "a later successful admit must not persist a previously failed unreturned device flow"
2285 );
2286 }
2287
2288 #[cfg(feature = "sqlite-store")]
2289 #[test]
2290 fn persistent_oauth_snapshot_merges_independent_authority_writes() {
2291 let temp_dir = tempfile::tempdir().expect("tempdir");
2292 let store_path = temp_dir.path().join("runtime.sqlite");
2293 let store_one: Arc<dyn RuntimeStore> =
2294 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2295 let store_two: Arc<dyn RuntimeStore> =
2296 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2297 let first_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2298 Duration::from_secs(60),
2299 Arc::new(RuntimeAuthLeaseHandle::new()),
2300 &store_one,
2301 );
2302 let second_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2303 Duration::from_secs(60),
2304 Arc::new(RuntimeAuthLeaseHandle::new()),
2305 &store_two,
2306 );
2307 let first_target = target();
2308 let second_target = alternate_target();
2309 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2310
2311 let first_state = first_authority
2312 .start(
2313 first_target.clone(),
2314 provider,
2315 "http://127.0.0.1/callback".to_string(),
2316 "verifier-1".to_string(),
2317 )
2318 .expect("first process admits browser flow");
2319 let second_state = second_authority
2320 .start(
2321 second_target.clone(),
2322 provider,
2323 "http://127.0.0.1/other-callback".to_string(),
2324 "verifier-2".to_string(),
2325 )
2326 .expect("second process admits browser flow");
2327
2328 let store_three: Arc<dyn RuntimeStore> =
2329 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2330 let snapshot_json = store_three
2331 .load_auth_oauth_flow_snapshot()
2332 .expect("durable OAuth snapshot loads")
2333 .expect("durable OAuth snapshot exists");
2334 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2335 .expect("durable OAuth snapshot decodes");
2336 assert!(
2337 snapshot
2338 .browser
2339 .iter()
2340 .any(|flow| flow.state == first_state),
2341 "the first independent authority write must survive the second write"
2342 );
2343 assert!(
2344 snapshot
2345 .browser
2346 .iter()
2347 .any(|flow| flow.state == second_state),
2348 "the second independent authority write must be persisted"
2349 );
2350
2351 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2352 Duration::from_secs(60),
2353 Arc::new(RuntimeAuthLeaseHandle::new()),
2354 &store_three,
2355 );
2356 restarted
2357 .consume(
2358 &first_state,
2359 &first_target,
2360 provider,
2361 "http://127.0.0.1/callback",
2362 )
2363 .expect("first independent flow rehydrates");
2364 restarted
2365 .consume(
2366 &second_state,
2367 &second_target,
2368 provider,
2369 "http://127.0.0.1/other-callback",
2370 )
2371 .expect("second independent flow rehydrates after first consume");
2372 }
2373
2374 #[test]
2375 fn persistent_oauth_browser_admit_does_not_resurrect_consumed_between_sync_and_persist() {
2376 let store = Arc::new(FailingOAuthSnapshotStore::default());
2377 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2378 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2379 Duration::from_secs(60),
2380 Arc::new(RuntimeAuthLeaseHandle::new()),
2381 &store_dyn,
2382 );
2383 let target = target();
2384 let replacement_target = alternate_target();
2385 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2386 let redirect_uri = "http://127.0.0.1/callback";
2387 let replacement_redirect_uri = "http://127.0.0.1/replacement-callback";
2388 let consumed_state = creator
2389 .start(
2390 target.clone(),
2391 provider,
2392 redirect_uri.to_string(),
2393 "consumed-verifier".to_string(),
2394 )
2395 .expect("creator admits browser flow");
2396
2397 let stale_authority = Arc::new(
2398 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2399 Duration::from_secs(60),
2400 Arc::new(RuntimeAuthLeaseHandle::new()),
2401 &store_dyn,
2402 ),
2403 );
2404 let consumer = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2405 Duration::from_secs(60),
2406 Arc::new(RuntimeAuthLeaseHandle::new()),
2407 &store_dyn,
2408 );
2409
2410 store.block_next_oauth_persist();
2411 let stale_admit = std::thread::spawn({
2412 let stale_authority = Arc::clone(&stale_authority);
2413 let replacement_target = replacement_target.clone();
2414 move || {
2415 stale_authority.start(
2416 replacement_target,
2417 provider,
2418 replacement_redirect_uri.to_string(),
2419 "replacement-verifier".to_string(),
2420 )
2421 }
2422 });
2423 store.wait_for_blocked_oauth_persist();
2424 consumer
2425 .consume(&consumed_state, &target, provider, redirect_uri)
2426 .expect("independent authority consumes browser flow between sync and persist");
2427 store.release_blocked_oauth_persist();
2428 let replacement_state = stale_admit
2429 .join()
2430 .expect("stale browser admit thread should not panic")
2431 .expect("stale authority admits replacement browser flow");
2432
2433 let snapshot_json = store
2434 .load_auth_oauth_flow_snapshot()
2435 .expect("durable OAuth snapshot loads")
2436 .expect("durable OAuth snapshot exists");
2437 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2438 .expect("durable OAuth snapshot decodes");
2439 assert!(
2440 !snapshot
2441 .browser
2442 .iter()
2443 .any(|flow| flow.state == consumed_state),
2444 "a stale admission must not resurrect a browser flow consumed after pre-sync"
2445 );
2446 assert!(
2447 snapshot
2448 .browser
2449 .iter()
2450 .any(|flow| flow.state == replacement_state),
2451 "the stale authority's newly admitted browser flow should still persist"
2452 );
2453
2454 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2455 Duration::from_secs(60),
2456 Arc::new(RuntimeAuthLeaseHandle::new()),
2457 &store_dyn,
2458 );
2459 assert!(matches!(
2460 restarted.consume(&consumed_state, &target, provider, redirect_uri),
2461 Err(OAuthFlowError::LifecycleRejected {
2462 operation: "verify_oauth_browser_flow",
2463 ..
2464 })
2465 ));
2466 restarted
2467 .consume(
2468 &replacement_state,
2469 &replacement_target,
2470 provider,
2471 replacement_redirect_uri,
2472 )
2473 .expect("new stale-authority flow survives restart");
2474 }
2475
2476 #[test]
2477 fn persistent_oauth_device_admit_does_not_resurrect_consumed_between_sync_and_persist() {
2478 let store = Arc::new(FailingOAuthSnapshotStore::default());
2479 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2480 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2481 Duration::from_secs(60),
2482 Arc::new(RuntimeAuthLeaseHandle::new()),
2483 &store_dyn,
2484 );
2485 let target = target();
2486 let replacement_target = alternate_target();
2487 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2488 let consumed_device_code = "consumed-device-code";
2489 let replacement_device_code = "replacement-device-code";
2490 creator
2491 .admit_device_code(
2492 target.clone(),
2493 provider,
2494 consumed_device_code.to_string(),
2495 Duration::from_secs(60),
2496 )
2497 .expect("creator admits device flow");
2498
2499 let stale_authority = Arc::new(
2500 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2501 Duration::from_secs(60),
2502 Arc::new(RuntimeAuthLeaseHandle::new()),
2503 &store_dyn,
2504 ),
2505 );
2506 let consumer = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2507 Duration::from_secs(60),
2508 Arc::new(RuntimeAuthLeaseHandle::new()),
2509 &store_dyn,
2510 );
2511
2512 store.block_next_oauth_persist();
2513 let stale_admit = std::thread::spawn({
2514 let stale_authority = Arc::clone(&stale_authority);
2515 let replacement_target = replacement_target.clone();
2516 move || {
2517 stale_authority.admit_device_code(
2518 replacement_target,
2519 provider,
2520 replacement_device_code.to_string(),
2521 Duration::from_secs(60),
2522 )
2523 }
2524 });
2525 store.wait_for_blocked_oauth_persist();
2526 consumer
2527 .begin_device_code_poll(consumed_device_code, &target, provider)
2528 .expect("independent authority begins device poll")
2529 .consume()
2530 .expect("independent authority consumes device flow between sync and persist");
2531 store.release_blocked_oauth_persist();
2532 stale_admit
2533 .join()
2534 .expect("stale device admit thread should not panic")
2535 .expect("stale authority admits replacement device flow");
2536
2537 let snapshot_json = store
2538 .load_auth_oauth_flow_snapshot()
2539 .expect("durable OAuth snapshot loads")
2540 .expect("durable OAuth snapshot exists");
2541 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2542 .expect("durable OAuth snapshot decodes");
2543 assert!(
2544 !snapshot
2545 .device
2546 .iter()
2547 .any(|flow| flow.device_code == consumed_device_code),
2548 "a stale admission must not resurrect a device flow consumed after pre-sync"
2549 );
2550 assert!(
2551 snapshot
2552 .device
2553 .iter()
2554 .any(|flow| flow.device_code == replacement_device_code),
2555 "the stale authority's newly admitted device flow should still persist"
2556 );
2557
2558 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2559 Duration::from_secs(60),
2560 Arc::new(RuntimeAuthLeaseHandle::new()),
2561 &store_dyn,
2562 );
2563 assert!(matches!(
2564 restarted.verify_device_code(consumed_device_code, &target, provider),
2565 Err(OAuthFlowError::LifecycleRejected {
2566 operation: "verify_oauth_device_flow",
2567 ..
2568 })
2569 ));
2570 restarted
2571 .verify_device_code(replacement_device_code, &replacement_target, provider)
2572 .expect("new stale-authority device flow survives restart");
2573 }
2574
2575 #[cfg(feature = "sqlite-store")]
2576 #[test]
2577 fn persistent_oauth_device_poll_finish_does_not_resurrect_consumed_payload() {
2578 let temp_dir = tempfile::tempdir().expect("tempdir");
2579 let store_path = temp_dir.path().join("runtime.sqlite");
2580 let creator_store: Arc<dyn RuntimeStore> =
2581 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2582 let stale_store: Arc<dyn RuntimeStore> =
2583 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2584 let consumer_store: Arc<dyn RuntimeStore> =
2585 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2586 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2587 Duration::from_secs(60),
2588 Arc::new(RuntimeAuthLeaseHandle::new()),
2589 &creator_store,
2590 );
2591 let target = target();
2592 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2593 let device_code = "pending-finish-device-code";
2594 creator
2595 .admit_device_code(
2596 target.clone(),
2597 provider,
2598 device_code.to_string(),
2599 Duration::from_secs(60),
2600 )
2601 .expect("creator admits device flow");
2602
2603 let stale_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2604 Duration::from_secs(60),
2605 Arc::new(RuntimeAuthLeaseHandle::new()),
2606 &stale_store,
2607 );
2608 let stale_poll = stale_authority
2609 .begin_device_code_poll(device_code, &target, provider)
2610 .expect("stale authority begins pending poll");
2611 let consumer = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2612 Duration::from_secs(60),
2613 Arc::new(RuntimeAuthLeaseHandle::new()),
2614 &consumer_store,
2615 );
2616 consumer
2617 .begin_device_code_poll(device_code, &target, provider)
2618 .expect("consumer begins independent poll")
2619 .consume()
2620 .expect("consumer consumes durable payload");
2621
2622 stale_poll
2623 .finish()
2624 .expect("stale pending poll finish is local cleanup only");
2625
2626 let snapshot_json = stale_store
2627 .load_auth_oauth_flow_snapshot()
2628 .expect("durable OAuth snapshot loads")
2629 .expect("durable OAuth snapshot exists");
2630 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
2631 .expect("durable OAuth snapshot decodes");
2632 assert!(
2633 !snapshot
2634 .device
2635 .iter()
2636 .any(|flow| flow.device_code == device_code),
2637 "a stale pending poll finish must not resurrect a consumed device flow"
2638 );
2639
2640 let restarted_store: Arc<dyn RuntimeStore> =
2641 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2642 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2643 Duration::from_secs(60),
2644 Arc::new(RuntimeAuthLeaseHandle::new()),
2645 &restarted_store,
2646 );
2647 assert!(matches!(
2648 restarted.verify_device_code(device_code, &target, provider),
2649 Err(OAuthFlowError::LifecycleRejected {
2650 operation: "verify_oauth_device_flow",
2651 ..
2652 })
2653 ));
2654 }
2655
2656 #[cfg(feature = "sqlite-store")]
2657 #[test]
2658 fn persistent_oauth_browser_sync_prunes_stale_capacity_before_admit() {
2659 let temp_dir = tempfile::tempdir().expect("tempdir");
2660 let store_path = temp_dir.path().join("runtime.sqlite");
2661 let creator_store: Arc<dyn RuntimeStore> =
2662 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2663 let stale_store: Arc<dyn RuntimeStore> =
2664 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2665 let consumer_store: Arc<dyn RuntimeStore> =
2666 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2667 let creator = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2668 Duration::from_secs(60),
2669 1,
2670 Arc::new(RuntimeAuthLeaseHandle::new()),
2671 Some(Arc::downgrade(&creator_store)),
2672 );
2673 let target = target();
2674 let replacement_target = alternate_target();
2675 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2676 let redirect_uri = "http://127.0.0.1/callback";
2677 let replacement_redirect_uri = "http://127.0.0.1/replacement-callback";
2678 let consumed_state = creator
2679 .start(
2680 target.clone(),
2681 provider,
2682 redirect_uri.to_string(),
2683 "consumed-verifier".to_string(),
2684 )
2685 .expect("creator admits browser flow");
2686
2687 let stale_authority = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2688 Duration::from_secs(60),
2689 1,
2690 Arc::new(RuntimeAuthLeaseHandle::new()),
2691 Some(Arc::downgrade(&stale_store)),
2692 );
2693 let consumer = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2694 Duration::from_secs(60),
2695 1,
2696 Arc::new(RuntimeAuthLeaseHandle::new()),
2697 Some(Arc::downgrade(&consumer_store)),
2698 );
2699 consumer
2700 .consume(&consumed_state, &target, provider, redirect_uri)
2701 .expect("independent authority consumes browser flow");
2702
2703 let replacement_state = stale_authority
2704 .start(
2705 replacement_target.clone(),
2706 provider,
2707 replacement_redirect_uri.to_string(),
2708 "replacement-verifier".to_string(),
2709 )
2710 .expect("stale capacity is pruned before browser admit");
2711 stale_authority
2712 .consume(
2713 &replacement_state,
2714 &replacement_target,
2715 provider,
2716 replacement_redirect_uri,
2717 )
2718 .expect("replacement browser flow remains usable");
2719 }
2720
2721 #[cfg(feature = "sqlite-store")]
2722 #[test]
2723 fn persistent_oauth_device_sync_prunes_stale_capacity_before_admit() {
2724 let temp_dir = tempfile::tempdir().expect("tempdir");
2725 let store_path = temp_dir.path().join("runtime.sqlite");
2726 let creator_store: Arc<dyn RuntimeStore> =
2727 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2728 let stale_store: Arc<dyn RuntimeStore> =
2729 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2730 let consumer_store: Arc<dyn RuntimeStore> =
2731 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
2732 let creator = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2733 Duration::from_secs(60),
2734 1,
2735 Arc::new(RuntimeAuthLeaseHandle::new()),
2736 Some(Arc::downgrade(&creator_store)),
2737 );
2738 let target = target();
2739 let replacement_target = alternate_target();
2740 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2741 let consumed_device_code = "consumed-capacity-device-code";
2742 let replacement_device_code = "replacement-device-code";
2743 creator
2744 .admit_device_code(
2745 target.clone(),
2746 provider,
2747 consumed_device_code.to_string(),
2748 Duration::from_secs(60),
2749 )
2750 .expect("creator admits device flow");
2751
2752 let stale_authority = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2753 Duration::from_secs(60),
2754 1,
2755 Arc::new(RuntimeAuthLeaseHandle::new()),
2756 Some(Arc::downgrade(&stale_store)),
2757 );
2758 let consumer = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2759 Duration::from_secs(60),
2760 1,
2761 Arc::new(RuntimeAuthLeaseHandle::new()),
2762 Some(Arc::downgrade(&consumer_store)),
2763 );
2764 consumer
2765 .begin_device_code_poll(consumed_device_code, &target, provider)
2766 .expect("independent authority begins device poll")
2767 .consume()
2768 .expect("independent authority consumes device flow");
2769
2770 stale_authority
2771 .admit_device_code(
2772 replacement_target.clone(),
2773 provider,
2774 replacement_device_code.to_string(),
2775 Duration::from_secs(60),
2776 )
2777 .expect("stale capacity is pruned before device admit");
2778 stale_authority
2779 .verify_device_code(replacement_device_code, &replacement_target, provider)
2780 .expect("replacement device flow remains usable");
2781 }
2782
2783 #[test]
2784 fn concurrent_persistent_browser_consumes_require_fresh_durable_claim() {
2785 let store = Arc::new(FailingOAuthSnapshotStore::default());
2786 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2787 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2788 Duration::from_secs(60),
2789 Arc::new(RuntimeAuthLeaseHandle::new()),
2790 &store_dyn,
2791 );
2792 let target = target();
2793 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2794 let redirect_uri = "http://127.0.0.1/callback";
2795 let state = creator
2796 .start(
2797 target.clone(),
2798 provider,
2799 redirect_uri.to_string(),
2800 "verifier".to_string(),
2801 )
2802 .expect("browser flow admitted");
2803 let first = Arc::new(
2804 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2805 Duration::from_secs(60),
2806 Arc::new(RuntimeAuthLeaseHandle::new()),
2807 &store_dyn,
2808 ),
2809 );
2810 let second = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2811 Duration::from_secs(60),
2812 Arc::new(RuntimeAuthLeaseHandle::new()),
2813 &store_dyn,
2814 );
2815
2816 store.block_next_oauth_persist();
2817 let first_consume = std::thread::spawn({
2818 let first = Arc::clone(&first);
2819 let state = state.clone();
2820 let target = target.clone();
2821 move || first.consume(&state, &target, provider, redirect_uri)
2822 });
2823 store.wait_for_blocked_oauth_persist();
2824
2825 second
2826 .consume(&state, &target, provider, redirect_uri)
2827 .expect("second authority wins durable consume race");
2828 store.release_blocked_oauth_persist();
2829 assert!(matches!(
2830 first_consume
2831 .join()
2832 .expect("first consume thread should not panic"),
2833 Err(OAuthFlowError::RegistryProjectionMissing {
2834 operation: "consume_oauth_browser_flow"
2835 })
2836 ));
2837 }
2838
2839 #[test]
2840 fn concurrent_persistent_device_consumes_require_fresh_durable_claim() {
2841 let store = Arc::new(FailingOAuthSnapshotStore::default());
2842 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2843 let creator = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2844 Duration::from_secs(60),
2845 Arc::new(RuntimeAuthLeaseHandle::new()),
2846 &store_dyn,
2847 );
2848 let target = target();
2849 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2850 let device_code = "race-device-code";
2851 creator
2852 .admit_device_code(
2853 target.clone(),
2854 provider,
2855 device_code.to_string(),
2856 Duration::from_secs(60),
2857 )
2858 .expect("device flow admitted");
2859 let first = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2860 Duration::from_secs(60),
2861 Arc::new(RuntimeAuthLeaseHandle::new()),
2862 &store_dyn,
2863 );
2864 let second = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
2865 Duration::from_secs(60),
2866 Arc::new(RuntimeAuthLeaseHandle::new()),
2867 &store_dyn,
2868 );
2869 let first_poll = first
2870 .begin_device_code_poll(device_code, &target, provider)
2871 .expect("first authority begins poll");
2872 let second_poll = second
2873 .begin_device_code_poll(device_code, &target, provider)
2874 .expect("second authority begins poll");
2875
2876 store.block_next_oauth_persist();
2877 let first_consume = std::thread::spawn(move || first_poll.consume());
2878 store.wait_for_blocked_oauth_persist();
2879
2880 second_poll
2881 .consume()
2882 .expect("second authority wins durable device consume race");
2883 store.release_blocked_oauth_persist();
2884 assert!(matches!(
2885 first_consume
2886 .join()
2887 .expect("first device consume thread should not panic"),
2888 Err(OAuthFlowError::RegistryProjectionMissing {
2889 operation: "consume_oauth_device_flow"
2890 })
2891 ));
2892 }
2893
2894 #[test]
2895 fn concurrent_persistent_browser_admits_confirm_fresh_generated_capacity() {
2896 let store = Arc::new(FailingOAuthSnapshotStore::default());
2897 let first_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2898 let second_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2899 let first = Arc::new(
2900 RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2901 Duration::from_secs(60),
2902 1,
2903 Arc::new(RuntimeAuthLeaseHandle::new()),
2904 Some(Arc::downgrade(&first_store)),
2905 ),
2906 );
2907 let second = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2908 Duration::from_secs(60),
2909 1,
2910 Arc::new(RuntimeAuthLeaseHandle::new()),
2911 Some(Arc::downgrade(&second_store)),
2912 );
2913 let first_target = target();
2914 let second_target = alternate_target();
2915 let provider = OAuthProviderIdentity::OpenAiChatGpt;
2916 let first_redirect_uri = "http://127.0.0.1/first-callback";
2917 let second_redirect_uri = "http://127.0.0.1/second-callback";
2918
2919 store.block_next_oauth_persist();
2920 let first_admit = std::thread::spawn({
2921 let first = Arc::clone(&first);
2922 let first_target = first_target.clone();
2923 move || {
2924 first.start(
2925 first_target,
2926 provider,
2927 first_redirect_uri.to_string(),
2928 "first-verifier".to_string(),
2929 )
2930 }
2931 });
2932 store.wait_for_blocked_oauth_persist();
2933
2934 second
2935 .start(
2936 second_target,
2937 provider,
2938 second_redirect_uri.to_string(),
2939 "second-verifier".to_string(),
2940 )
2941 .expect("second authority wins durable browser admission race");
2942 store.release_blocked_oauth_persist();
2943 assert!(matches!(
2944 first_admit
2945 .join()
2946 .expect("first browser admit thread should not panic"),
2947 Err(OAuthFlowError::LifecycleRejected {
2948 operation: "admit_oauth_browser_flow",
2949 ..
2950 })
2951 ));
2952 }
2953
2954 #[test]
2955 fn concurrent_persistent_device_admits_confirm_fresh_generated_capacity() {
2956 let store = Arc::new(FailingOAuthSnapshotStore::default());
2957 let first_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2958 let second_store = Arc::clone(&store) as Arc<dyn RuntimeStore>;
2959 let first = Arc::new(
2960 RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2961 Duration::from_secs(60),
2962 1,
2963 Arc::new(RuntimeAuthLeaseHandle::new()),
2964 Some(Arc::downgrade(&first_store)),
2965 ),
2966 );
2967 let second = RuntimeOAuthFlowHandle::new_with_capacity_auth_lease_and_store(
2968 Duration::from_secs(60),
2969 1,
2970 Arc::new(RuntimeAuthLeaseHandle::new()),
2971 Some(Arc::downgrade(&second_store)),
2972 );
2973 let first_target = target();
2974 let second_target = alternate_target();
2975 let provider = OAuthProviderIdentity::GoogleCodeAssist;
2976
2977 store.block_next_oauth_persist();
2978 let first_admit = std::thread::spawn({
2979 let first = Arc::clone(&first);
2980 let first_target = first_target.clone();
2981 move || {
2982 first.admit_device_code(
2983 first_target,
2984 provider,
2985 "first-device-code".to_string(),
2986 Duration::from_secs(60),
2987 )
2988 }
2989 });
2990 store.wait_for_blocked_oauth_persist();
2991
2992 second
2993 .admit_device_code(
2994 second_target,
2995 provider,
2996 "second-device-code".to_string(),
2997 Duration::from_secs(60),
2998 )
2999 .expect("second authority wins durable device admission race");
3000 store.release_blocked_oauth_persist();
3001 assert!(matches!(
3002 first_admit
3003 .join()
3004 .expect("first device admit thread should not panic"),
3005 Err(OAuthFlowError::LifecycleRejected {
3006 operation: "admit_oauth_device_flow",
3007 ..
3008 })
3009 ));
3010 }
3011
3012 #[test]
3013 fn concurrent_browser_admits_preserve_newer_durable_snapshot() {
3014 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3015 let store = Arc::new(FailingOAuthSnapshotStore::default());
3016 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3017 let authority = Arc::new(
3018 RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3019 Duration::from_secs(60),
3020 lifecycle,
3021 &store_dyn,
3022 ),
3023 );
3024 let first_target = target();
3025 let second_target = alternate_target();
3026 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3027 let first_redirect_uri = "http://127.0.0.1/callback";
3028 let second_redirect_uri = "http://127.0.0.1/other-callback";
3029
3030 store.block_next_oauth_persist();
3031 let first_admit = std::thread::spawn({
3032 let authority = Arc::clone(&authority);
3033 let target = first_target.clone();
3034 move || {
3035 authority.start(
3036 target,
3037 provider,
3038 first_redirect_uri.to_string(),
3039 "verifier-1".to_string(),
3040 )
3041 }
3042 });
3043 store.wait_for_blocked_oauth_persist();
3044
3045 let (second_done_tx, second_done_rx) = std::sync::mpsc::channel();
3046 std::thread::spawn({
3047 let authority = Arc::clone(&authority);
3048 let target = second_target.clone();
3049 move || {
3050 let result = authority.start(
3051 target,
3052 provider,
3053 second_redirect_uri.to_string(),
3054 "verifier-2".to_string(),
3055 );
3056 let _ = second_done_tx.send(result);
3057 }
3058 });
3059 let second_before_release = second_done_rx.recv_timeout(Duration::from_millis(100)).ok();
3060 store.release_blocked_oauth_persist();
3061 first_admit
3062 .join()
3063 .expect("first admit thread should not panic")
3064 .expect("first browser flow admitted");
3065 let second_state = second_before_release
3066 .unwrap_or_else(|| {
3067 second_done_rx
3068 .recv_timeout(Duration::from_secs(1))
3069 .expect("second admit should finish after first durable write is released")
3070 })
3071 .expect("second browser flow admitted");
3072
3073 let snapshot_json = store
3074 .load_auth_oauth_flow_snapshot()
3075 .expect("durable OAuth snapshot loads")
3076 .expect("durable OAuth snapshot exists");
3077 let snapshot = serde_json::from_slice::<OAuthFlowRegistrySnapshot>(&snapshot_json)
3078 .expect("durable OAuth snapshot decodes");
3079 assert!(
3080 snapshot
3081 .browser
3082 .iter()
3083 .any(|flow| flow.state == second_state),
3084 "durable snapshot must retain flow admitted by a concurrent newer write"
3085 );
3086
3087 let restarted_lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3088 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3089 Duration::from_secs(60),
3090 restarted_lifecycle,
3091 &store_dyn,
3092 );
3093 let record = restarted
3094 .consume(&second_state, &second_target, provider, second_redirect_uri)
3095 .expect("newer durable flow should survive restart");
3096 assert_eq!(record.pkce_verifier, "verifier-2");
3097 }
3098
3099 #[test]
3100 fn browser_consume_persistence_failure_keeps_flow_retryable() {
3101 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3102 let store = Arc::new(FailingOAuthSnapshotStore::default());
3103 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3104 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3105 Duration::from_secs(60),
3106 lifecycle.clone(),
3107 &store_dyn,
3108 );
3109 let target = target();
3110 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3111 let redirect_uri = "http://127.0.0.1/callback";
3112 let state = authority
3113 .start(
3114 target.clone(),
3115 provider,
3116 redirect_uri.to_string(),
3117 "verifier".to_string(),
3118 )
3119 .expect("browser flow admitted");
3120
3121 store.fail_oauth_persist();
3122 assert!(matches!(
3123 authority.consume(&state, &target, provider, redirect_uri),
3124 Err(OAuthFlowError::PersistenceFailed { .. })
3125 ));
3126 assert!(lifecycle.has_oauth_browser_flow_for_test(&target, &state));
3127
3128 store.allow_oauth_persist();
3129 authority
3130 .consume(&state, &target, provider, redirect_uri)
3131 .expect("failed durable consume remains retryable");
3132 }
3133
3134 #[test]
3135 fn device_consume_persistence_failure_keeps_flow_retryable() {
3136 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3137 let store = Arc::new(FailingOAuthSnapshotStore::default());
3138 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3139 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3140 Duration::from_secs(60),
3141 lifecycle.clone(),
3142 &store_dyn,
3143 );
3144 let target = target();
3145 let provider = OAuthProviderIdentity::GoogleCodeAssist;
3146 let device_code = "provider-device-code";
3147 authority
3148 .admit_device_code(
3149 target.clone(),
3150 provider,
3151 device_code.to_string(),
3152 Duration::from_secs(60),
3153 )
3154 .expect("device flow admitted");
3155 let poll = authority
3156 .begin_device_code_poll(device_code, &target, provider)
3157 .expect("device poll begins");
3158
3159 store.fail_oauth_persist();
3160 assert!(matches!(
3161 poll.consume(),
3162 Err(OAuthFlowError::PersistenceFailed { .. })
3163 ));
3164 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3165
3166 store.allow_oauth_persist();
3167 let retry = authority
3168 .begin_device_code_poll(device_code, &target, provider)
3169 .expect("failed durable consume keeps device flow retryable");
3170 retry
3171 .consume()
3172 .expect("retry consumes after durable persistence recovers");
3173 }
3174
3175 #[test]
3176 fn release_persistence_failure_keeps_released_flows_retryable() {
3177 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3178 let store = Arc::new(FailingOAuthSnapshotStore::default());
3179 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3180 let authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3181 Duration::from_secs(60),
3182 lifecycle.clone(),
3183 &store_dyn,
3184 );
3185 let target = target();
3186 let lease_key = LeaseKey::from_auth_binding(&target);
3187 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3188 let redirect_uri = "http://127.0.0.1/callback";
3189 let state = authority
3190 .start(
3191 target.clone(),
3192 provider,
3193 redirect_uri.to_string(),
3194 "verifier".to_string(),
3195 )
3196 .expect("browser flow admitted");
3197
3198 store.fail_oauth_persist();
3199 assert!(
3200 lifecycle.release_lease(&lease_key).is_err(),
3201 "release should fail closed when durable OAuth cleanup cannot persist"
3202 );
3203 assert!(lifecycle.has_oauth_browser_flow_for_test(&target, &state));
3204
3205 store.allow_oauth_persist();
3206 authority
3207 .consume(&state, &target, provider, redirect_uri)
3208 .expect("failed durable release leaves browser flow retryable");
3209 }
3210
3211 #[test]
3212 fn stale_release_persistence_failure_does_not_install_released_authority() {
3213 let releasing_lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3214 let store = Arc::new(FailingOAuthSnapshotStore::default());
3215 let store_dyn = Arc::clone(&store) as Arc<dyn RuntimeStore>;
3216 let releasing_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3217 Duration::from_secs(60),
3218 releasing_lifecycle.clone(),
3219 &store_dyn,
3220 );
3221 let admitting_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3222 Duration::from_secs(60),
3223 Arc::new(RuntimeAuthLeaseHandle::new()),
3224 &store_dyn,
3225 );
3226 let target = target();
3227 let lease_key = LeaseKey::from_auth_binding(&target);
3228 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3229 let redirect_uri = "http://127.0.0.1/callback";
3230 let state = admitting_authority
3231 .start(
3232 target.clone(),
3233 provider,
3234 redirect_uri.to_string(),
3235 "verifier".to_string(),
3236 )
3237 .expect("other authority admits browser flow");
3238 assert!(
3239 !releasing_lifecycle.has_oauth_browser_flow_for_test(&target, &state),
3240 "releasing authority starts stale and has no local machine membership"
3241 );
3242
3243 store.fail_oauth_persist();
3244 assert!(
3245 releasing_lifecycle.release_lease(&lease_key).is_err(),
3246 "release should fail closed when stale durable OAuth cleanup cannot persist"
3247 );
3248 assert_eq!(
3249 releasing_lifecycle.snapshot(&lease_key).phase,
3250 None,
3251 "failed stale release must not synthesize a local AuthMachine authority"
3252 );
3253
3254 store.allow_oauth_persist();
3255 releasing_authority
3256 .consume(&state, &target, provider, redirect_uri)
3257 .expect("failed stale durable release must leave browser flow retryable");
3258 }
3259
3260 #[cfg(feature = "sqlite-store")]
3261 #[test]
3262 fn persistent_release_prunes_durable_flows_from_stale_authority() {
3263 let temp_dir = tempfile::tempdir().expect("tempdir");
3264 let store_path = temp_dir.path().join("runtime.sqlite");
3265 let releasing_store: Arc<dyn RuntimeStore> =
3266 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
3267 let admitting_store: Arc<dyn RuntimeStore> =
3268 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
3269 let releasing_lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3270 let releasing_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3271 Duration::from_secs(60),
3272 releasing_lifecycle.clone(),
3273 &releasing_store,
3274 );
3275 let admitting_authority = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3276 Duration::from_secs(60),
3277 Arc::new(RuntimeAuthLeaseHandle::new()),
3278 &admitting_store,
3279 );
3280 let target = target();
3281 let lease_key = LeaseKey::from_auth_binding(&target);
3282 let browser_provider = OAuthProviderIdentity::OpenAiChatGpt;
3283 let redirect_uri = "http://127.0.0.1/callback";
3284 let browser_state = admitting_authority
3285 .start(
3286 target.clone(),
3287 browser_provider,
3288 redirect_uri.to_string(),
3289 "browser-verifier".to_string(),
3290 )
3291 .expect("other authority admits browser flow");
3292 let device_provider = OAuthProviderIdentity::GoogleCodeAssist;
3293 let device_code = "released-device-code";
3294 admitting_authority
3295 .admit_device_code(
3296 target.clone(),
3297 device_provider,
3298 device_code.to_string(),
3299 Duration::from_secs(60),
3300 )
3301 .expect("other authority admits device flow");
3302 assert!(
3303 !releasing_lifecycle.has_oauth_browser_flow_for_test(&target, &browser_state),
3304 "releasing authority starts stale and does not know the browser flow locally"
3305 );
3306 assert!(
3307 !releasing_lifecycle.has_oauth_device_flow_for_test(&target, device_code),
3308 "releasing authority starts stale and does not know the device flow locally"
3309 );
3310
3311 releasing_lifecycle
3312 .release_lease(&lease_key)
3313 .expect("stale release succeeds");
3314
3315 let restarted_store: Arc<dyn RuntimeStore> =
3316 Arc::new(crate::store::sqlite::SqliteRuntimeStore::new(&store_path).unwrap());
3317 let restarted = RuntimeOAuthFlowHandle::new_with_persistent_store_and_auth_lease(
3318 Duration::from_secs(60),
3319 Arc::new(RuntimeAuthLeaseHandle::new()),
3320 &restarted_store,
3321 );
3322 let browser_after_release =
3323 restarted.consume(&browser_state, &target, browser_provider, redirect_uri);
3324 let device_after_release =
3325 restarted.verify_device_code(device_code, &target, device_provider);
3326 assert!(
3327 matches!(
3328 browser_after_release,
3329 Err(OAuthFlowError::LifecycleRejected {
3330 operation: "verify_oauth_browser_flow",
3331 ..
3332 })
3333 ),
3334 "release from a stale authority must prune durable browser flow, got {browser_after_release:?}"
3335 );
3336 assert!(
3337 matches!(
3338 device_after_release,
3339 Err(OAuthFlowError::LifecycleRejected {
3340 operation: "verify_oauth_device_flow",
3341 ..
3342 })
3343 ),
3344 "release from a stale authority must prune durable device flow, got {device_after_release:?}"
3345 );
3346 drop(releasing_authority);
3347 }
3348
3349 #[test]
3350 fn oauth_flow_membership_does_not_advance_credential_generation() {
3351 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3352 let authority =
3353 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3354 let target = target();
3355 let lease_key = LeaseKey::from_auth_binding(&target);
3356 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3357 let redirect_uri = "http://127.0.0.1/callback";
3358 let transition = lifecycle
3359 .acquire_lease(&lease_key, 4_200)
3360 .expect("credential lifecycle acquired");
3361
3362 let state = authority
3363 .start(
3364 target.clone(),
3365 provider,
3366 redirect_uri.to_string(),
3367 "verifier".to_string(),
3368 )
3369 .expect("browser flow admitted");
3370 authority
3371 .verify(&state, &target, provider, redirect_uri)
3372 .expect("browser flow verifies");
3373 authority
3374 .consume(&state, &target, provider, redirect_uri)
3375 .expect("browser flow consumes");
3376
3377 let snapshot = lifecycle.snapshot(&lease_key);
3378 assert_eq!(snapshot.generation, transition.generation());
3379 assert_eq!(
3380 snapshot.credential_published_at_millis,
3381 transition.credential_published_at_millis()
3382 );
3383 }
3384
3385 #[test]
3386 fn global_browser_expiry_preserves_reauth_required_phase() {
3387 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3388 let authority = RuntimeOAuthFlowHandle::new_with_auth_lease(
3389 Duration::from_millis(1),
3390 lifecycle.clone(),
3391 );
3392 let target = target();
3393 let other_target = alternate_target();
3394 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3395 let redirect_uri = "http://127.0.0.1/callback";
3396
3397 let expired_state = authority
3398 .start(
3399 target.clone(),
3400 provider,
3401 redirect_uri.to_string(),
3402 "verifier-old".to_string(),
3403 )
3404 .expect("browser flow admitted");
3405 assert_eq!(
3406 snapshot_phase(&lifecycle, &target),
3407 Some(AuthLeasePhase::ReauthRequired)
3408 );
3409 std::thread::sleep(Duration::from_millis(10));
3410
3411 authority
3412 .start(
3413 other_target,
3414 provider,
3415 redirect_uri.to_string(),
3416 "verifier-new".to_string(),
3417 )
3418 .expect("new browser flow admitted after pruning expired flow");
3419
3420 assert!(
3421 !lifecycle.has_oauth_browser_flow_for_test(&target, &expired_state),
3422 "passive registry expiry must remove stale AuthMachine browser membership"
3423 );
3424 assert_eq!(
3425 snapshot_phase(&lifecycle, &target),
3426 Some(AuthLeasePhase::ReauthRequired),
3427 "global OAuth expiry cleanup must not change credential lifecycle truth"
3428 );
3429 }
3430
3431 #[test]
3432 fn browser_passive_expiry_clears_lifecycle_membership_on_next_admit() {
3433 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3434 let authority = RuntimeOAuthFlowHandle::new_with_auth_lease(
3435 Duration::from_millis(1),
3436 lifecycle.clone(),
3437 );
3438 let target = target();
3439 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3440 let redirect_uri = "http://127.0.0.1/callback";
3441
3442 let expired_state = authority
3443 .start(
3444 target.clone(),
3445 provider,
3446 redirect_uri.to_string(),
3447 "verifier-old".to_string(),
3448 )
3449 .expect("browser flow admitted");
3450 assert!(lifecycle.has_oauth_browser_flow_for_test(&target, &expired_state));
3451 std::thread::sleep(Duration::from_millis(10));
3452
3453 authority
3454 .start(
3455 target.clone(),
3456 provider,
3457 redirect_uri.to_string(),
3458 "verifier-new".to_string(),
3459 )
3460 .expect("new browser flow admitted after pruning expired flow");
3461
3462 assert!(
3463 !lifecycle.has_oauth_browser_flow_for_test(&target, &expired_state),
3464 "passive registry expiry must remove stale AuthMachine browser membership"
3465 );
3466 }
3467
3468 #[test]
3469 fn device_admit_rejects_registry_pruned_canonical_membership() {
3470 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3471 let authority =
3472 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3473 let target = target();
3474 let provider = OAuthProviderIdentity::GoogleCodeAssist;
3475 let device_code = "provider-device-code";
3476
3477 authority
3478 .admit_device_code(
3479 target.clone(),
3480 provider,
3481 device_code.to_string(),
3482 Duration::from_secs(60),
3483 )
3484 .expect("device flow admitted");
3485 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3486 authority
3487 .registry
3488 .expire_device_code(device_code, &target, provider)
3489 .expect("test removes registry record without lifecycle cleanup");
3490 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3491
3492 assert!(matches!(
3493 authority.admit_device_code(
3494 target.clone(),
3495 provider,
3496 device_code.to_string(),
3497 Duration::from_secs(60),
3498 ),
3499 Err(OAuthFlowError::LifecycleRejected {
3500 operation: "admit_oauth_device_flow",
3501 ..
3502 })
3503 ));
3504
3505 assert!(
3506 lifecycle.has_oauth_device_flow_for_test(&target, device_code),
3507 "registry-only loss must not expire canonical AuthMachine device membership"
3508 );
3509 assert!(matches!(
3510 authority.verify_device_code(device_code, &target, provider),
3511 Err(OAuthFlowError::RegistryProjectionMissing {
3512 operation: "verify_oauth_device_flow"
3513 })
3514 ));
3515 assert!(matches!(
3516 authority.begin_device_code_poll(device_code, &target, provider),
3517 Err(OAuthFlowError::RegistryProjectionMissing {
3518 operation: "begin_oauth_device_poll"
3519 })
3520 ));
3521 assert!(
3522 lifecycle.has_oauth_device_flow_for_test(&target, device_code),
3523 "missing process-local device payload must fail closed without removing the flow"
3524 );
3525 }
3526
3527 #[test]
3528 fn duplicate_device_admit_preserves_active_lifecycle_membership() {
3529 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3530 let authority =
3531 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3532 let target = target();
3533 let provider = OAuthProviderIdentity::GoogleCodeAssist;
3534 let device_code = "provider-device-code";
3535
3536 authority
3537 .admit_device_code(
3538 target.clone(),
3539 provider,
3540 device_code.to_string(),
3541 Duration::from_secs(60),
3542 )
3543 .expect("device flow admitted");
3544 let duplicate = authority.admit_device_code(
3545 target.clone(),
3546 provider,
3547 device_code.to_string(),
3548 Duration::from_secs(60),
3549 );
3550
3551 assert!(matches!(
3552 duplicate,
3553 Err(OAuthFlowError::LifecycleRejected {
3554 operation: "admit_oauth_device_flow",
3555 ..
3556 })
3557 ));
3558 assert!(lifecycle.has_oauth_device_flow_for_test(&target, device_code));
3559 authority
3560 .begin_device_code_poll(device_code, &target, provider)
3561 .expect("duplicate admit must not orphan active lifecycle membership");
3562 }
3563
3564 #[test]
3565 fn stale_registry_payloads_do_not_block_authmachine_admission_after_lifecycle_release() {
3566 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3567 let authority = RuntimeOAuthFlowHandle::new_with_capacity_and_auth_lease(
3568 Duration::from_secs(60),
3569 1,
3570 lifecycle.clone(),
3571 );
3572 let target = target();
3573 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3574
3575 authority
3576 .start(
3577 target.clone(),
3578 provider,
3579 "http://127.0.0.1/callback".to_string(),
3580 "verifier-1".to_string(),
3581 )
3582 .expect("first browser flow admitted");
3583 lifecycle
3584 .release_lease(&LeaseKey::from_auth_binding(&target))
3585 .expect("credential lifecycle release succeeds");
3586
3587 authority
3588 .start(
3589 alternate_target(),
3590 provider,
3591 "http://127.0.0.1/other-callback".to_string(),
3592 "verifier-2".to_string(),
3593 )
3594 .expect("AuthMachine release must clear stale registry payload capacity");
3595 }
3596
3597 #[test]
3598 fn release_observer_does_not_prune_flow_admitted_after_release_acceptance() {
3599 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3600 let authority = Arc::new(RuntimeOAuthFlowHandle::new_with_auth_lease(
3601 Duration::from_secs(60),
3602 lifecycle.clone(),
3603 ));
3604 let target = target_with_binding("release_after_accept_openai");
3605 let lease_key = LeaseKey::from_auth_binding(&target);
3606 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3607 let redirect_uri = "http://127.0.0.1/callback";
3608 let old_state = authority
3609 .start(
3610 target.clone(),
3611 provider,
3612 redirect_uri.to_string(),
3613 "old-verifier".to_string(),
3614 )
3615 .expect("old browser flow admitted");
3616 let new_state = Arc::new(std::sync::Mutex::new(None));
3617 let new_state_for_hook = Arc::clone(&new_state);
3618 let authority_for_hook = Arc::clone(&authority);
3619 let target_for_hook = target.clone();
3620 let lease_key_for_hook = lease_key.clone();
3621 let _hook_guard = crate::handles::auth_lease::install_release_after_accept_hook_for_test(
3622 Arc::new(move |released_key| {
3623 if released_key != &lease_key_for_hook {
3624 return;
3625 }
3626 let admitted = authority_for_hook
3627 .start(
3628 target_for_hook.clone(),
3629 provider,
3630 redirect_uri.to_string(),
3631 "new-verifier".to_string(),
3632 )
3633 .expect("new browser flow admitted after release acceptance");
3634 *new_state_for_hook
3635 .lock()
3636 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(admitted);
3637 }),
3638 );
3639
3640 lifecycle
3641 .release_lease(&lease_key)
3642 .expect("credential lifecycle release succeeds");
3643
3644 let new_state = new_state
3645 .lock()
3646 .unwrap_or_else(std::sync::PoisonError::into_inner)
3647 .clone()
3648 .expect("hook admitted replacement flow");
3649 let flow = authority
3650 .consume(&new_state, &target, provider, redirect_uri)
3651 .expect("release observer must not prune newly admitted flow");
3652 assert_eq!(flow.pkce_verifier, "new-verifier");
3653 assert!(matches!(
3654 authority.consume(&old_state, &target, provider, redirect_uri),
3655 Err(OAuthFlowError::LifecycleRejected {
3656 operation: "verify_oauth_browser_flow",
3657 ..
3658 })
3659 ));
3660 }
3661
3662 #[test]
3663 fn release_precommit_admission_waits_for_release_commit() {
3664 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3665 let authority = Arc::new(RuntimeOAuthFlowHandle::new_with_auth_lease(
3666 Duration::from_secs(60),
3667 lifecycle.clone(),
3668 ));
3669 let target = target_with_binding("release_before_commit_openai");
3670 let lease_key = LeaseKey::from_auth_binding(&target);
3671 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3672 let redirect_uri = "http://127.0.0.1/callback";
3673 let old_state = authority
3674 .start(
3675 target.clone(),
3676 provider,
3677 redirect_uri.to_string(),
3678 "old-verifier".to_string(),
3679 )
3680 .expect("old browser flow admitted");
3681 let (done_tx, done_rx) = mpsc::channel();
3682 let admission_finished = Arc::new(AtomicBool::new(false));
3683 let admission_finished_for_hook = Arc::clone(&admission_finished);
3684 let authority_for_hook = Arc::clone(&authority);
3685 let target_for_hook = target.clone();
3686 let lease_key_for_hook = lease_key.clone();
3687 let _hook_guard = crate::handles::auth_lease::install_release_before_commit_hook_for_test(
3688 Arc::new(move |released_key| {
3689 if released_key != &lease_key_for_hook {
3690 return;
3691 }
3692 let authority_for_thread = Arc::clone(&authority_for_hook);
3693 let target_for_thread = target_for_hook.clone();
3694 let done_tx = done_tx.clone();
3695 let admission_finished_for_thread = Arc::clone(&admission_finished_for_hook);
3696 std::thread::spawn(move || {
3697 let admitted = authority_for_thread.start(
3698 target_for_thread,
3699 provider,
3700 redirect_uri.to_string(),
3701 "new-verifier".to_string(),
3702 );
3703 admission_finished_for_thread.store(true, Ordering::Release);
3704 let _ = done_tx.send(admitted);
3705 });
3706 std::thread::sleep(Duration::from_millis(50));
3707 assert!(
3708 !admission_finished_for_hook.load(Ordering::Acquire),
3709 "OAuth admission must wait until release commits"
3710 );
3711 }),
3712 );
3713
3714 lifecycle
3715 .release_lease(&lease_key)
3716 .expect("credential lifecycle release succeeds");
3717
3718 let new_state = done_rx
3719 .recv_timeout(Duration::from_secs(1))
3720 .expect("pre-commit admission should finish after release commits")
3721 .expect("new browser flow admitted after release commit");
3722 let flow = authority
3723 .consume(&new_state, &target, provider, redirect_uri)
3724 .expect("release observer must not prune newly admitted flow");
3725 assert_eq!(flow.pkce_verifier, "new-verifier");
3726 assert!(matches!(
3727 authority.consume(&old_state, &target, provider, redirect_uri),
3728 Err(OAuthFlowError::LifecycleRejected {
3729 operation: "verify_oauth_browser_flow",
3730 ..
3731 })
3732 ));
3733 }
3734
3735 #[test]
3736 fn browser_capacity_rejection_comes_from_authmachine_lifecycle() {
3737 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3738 let authority = RuntimeOAuthFlowHandle::new_with_capacity_and_auth_lease(
3739 Duration::from_secs(60),
3740 1,
3741 lifecycle,
3742 );
3743 let target = target();
3744 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3745 let redirect_uri = "http://127.0.0.1/callback";
3746
3747 authority
3748 .start(
3749 target.clone(),
3750 provider,
3751 redirect_uri.to_string(),
3752 "verifier-1".to_string(),
3753 )
3754 .expect("first browser flow admitted");
3755
3756 assert!(matches!(
3757 authority.start(
3758 alternate_target(),
3759 provider,
3760 "http://127.0.0.1/other-callback".to_string(),
3761 "verifier-2".to_string(),
3762 ),
3763 Err(OAuthFlowError::LifecycleRejected {
3764 operation: "admit_oauth_browser_flow",
3765 ..
3766 })
3767 ));
3768 }
3769
3770 #[test]
3771 fn browser_provider_mismatch_rejection_comes_from_authmachine_lifecycle() {
3772 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3773 let authority =
3774 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle);
3775 let target = target();
3776 let redirect_uri = "http://127.0.0.1/callback";
3777 let state = authority
3778 .start(
3779 target.clone(),
3780 OAuthProviderIdentity::OpenAiChatGpt,
3781 redirect_uri.to_string(),
3782 "verifier".to_string(),
3783 )
3784 .expect("browser flow admitted");
3785
3786 assert!(matches!(
3787 authority.verify(
3788 &state,
3789 &target,
3790 OAuthProviderIdentity::GoogleCodeAssist,
3791 redirect_uri,
3792 ),
3793 Err(OAuthFlowError::LifecycleRejected {
3794 operation: "verify_oauth_browser_flow",
3795 ..
3796 })
3797 ));
3798 }
3799
3800 #[test]
3813 fn release_lease_terminally_cancels_pending_browser_flow() {
3814 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3815 let authority =
3816 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3817 let target = target();
3818 let lease_key = LeaseKey::from_auth_binding(&target);
3819 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3820 let redirect_uri = "http://127.0.0.1/callback";
3821
3822 let state = authority
3823 .start(
3824 target.clone(),
3825 provider,
3826 redirect_uri.to_string(),
3827 "verifier".to_string(),
3828 )
3829 .expect("browser flow admitted");
3830
3831 lifecycle
3832 .release_lease(&lease_key)
3833 .expect("release with a pending flow must drain it, not fail");
3834 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3835 assert!(
3836 !lifecycle.has_oauth_browser_flow_for_test(&target, &state),
3837 "pending flow must be terminally cancelled by the release drain"
3838 );
3839
3840 assert!(matches!(
3842 authority.consume(&state, &target, provider, redirect_uri),
3843 Err(OAuthFlowError::LifecycleRejected { .. })
3844 ));
3845
3846 lifecycle
3850 .apply_oauth_input(
3851 &target,
3852 auth_dsl::AuthMachineInput::ExpireOAuthBrowserFlow {
3853 flow_id: state.clone(),
3854 },
3855 "late_prune_expire_browser",
3856 false,
3857 )
3858 .expect("post-release expire must be a benign no-op");
3859 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3860 }
3861
3862 #[test]
3865 fn release_lease_terminally_cancels_pending_device_flow() {
3866 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3867 let authority =
3868 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3869 let target = target();
3870 let lease_key = LeaseKey::from_auth_binding(&target);
3871 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3872
3873 authority
3874 .admit_device_code(
3875 target.clone(),
3876 provider,
3877 "device-code-1".to_string(),
3878 Duration::from_secs(60),
3879 )
3880 .expect("device flow admitted");
3881
3882 lifecycle
3883 .release_lease(&lease_key)
3884 .expect("release with a pending device flow must drain it, not fail");
3885 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3886 assert!(
3887 !lifecycle.has_oauth_device_flow_for_test(&target, "device-code-1"),
3888 "pending device flow must be terminally cancelled by the release drain"
3889 );
3890
3891 lifecycle
3893 .expire_device_flow(&target, "device-code-1")
3894 .expect("post-release device expire must be a benign no-op");
3895 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3896 }
3897
3898 #[test]
3903 fn post_release_oauth_observations_through_shell_dispatch_are_benign() {
3904 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3905 let target = target();
3906 let lease_key = LeaseKey::from_auth_binding(&target);
3907
3908 lifecycle
3909 .release_lease(&lease_key)
3910 .expect("releasing an empty lease succeeds");
3911 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3914
3915 lifecycle
3916 .apply_oauth_input(
3917 &target,
3918 auth_dsl::AuthMachineInput::ExpireOAuthBrowserFlow {
3919 flow_id: "ghost-browser".to_string(),
3920 },
3921 "late_prune_expire_browser",
3922 false,
3923 )
3924 .expect("post-release browser expire must be a benign no-op");
3925 lifecycle
3926 .expire_device_flow(&target, "ghost-device")
3927 .expect("post-release device expire must be a benign no-op");
3928 lifecycle
3929 .finish_device_poll(&target, "ghost-poll")
3930 .expect("post-release poll finish must be a benign no-op");
3931 lifecycle
3932 .confirm_oauth_durable_admission(&target, 0, 16, "late_confirm_oauth_durable_admission")
3933 .expect("post-release durable-admission confirmation must be a benign no-op");
3934
3935 assert_eq!(snapshot_phase(&lifecycle, &target), None);
3937 }
3938
3939 #[test]
3944 fn late_prune_expire_at_release_acceptance_is_benign() {
3945 let lifecycle = Arc::new(RuntimeAuthLeaseHandle::new());
3946 let authority =
3947 RuntimeOAuthFlowHandle::new_with_auth_lease(Duration::from_secs(60), lifecycle.clone());
3948 let target = target();
3949 let lease_key = LeaseKey::from_auth_binding(&target);
3950 let provider = OAuthProviderIdentity::OpenAiChatGpt;
3951
3952 let state = authority
3953 .start(
3954 target.clone(),
3955 provider,
3956 "http://127.0.0.1/callback".to_string(),
3957 "verifier".to_string(),
3958 )
3959 .expect("browser flow admitted");
3960
3961 let hook_result: Arc<StdMutex<Option<Result<(), String>>>> = Arc::new(StdMutex::new(None));
3962 let hook_result_for_hook = Arc::clone(&hook_result);
3963 let lifecycle_for_hook = Arc::clone(&lifecycle);
3964 let target_for_hook = target.clone();
3965 let lease_key_for_hook = lease_key.clone();
3966 let flow_for_hook = state.clone();
3967 let _hook_guard = crate::handles::auth_lease::install_release_after_accept_hook_for_test(
3968 Arc::new(move |released_key| {
3969 if released_key != &lease_key_for_hook {
3970 return;
3971 }
3972 let outcome = lifecycle_for_hook
3973 .apply_oauth_input(
3974 &target_for_hook,
3975 auth_dsl::AuthMachineInput::ExpireOAuthBrowserFlow {
3976 flow_id: flow_for_hook.clone(),
3977 },
3978 "late_prune_expire_at_release_acceptance",
3979 false,
3980 )
3981 .map_err(|err| err.to_string());
3982 *hook_result_for_hook
3983 .lock()
3984 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(outcome);
3985 }),
3986 );
3987
3988 lifecycle
3989 .release_lease(&lease_key)
3990 .expect("release with a pending flow must drain it, not fail");
3991
3992 let outcome = hook_result
3993 .lock()
3994 .unwrap_or_else(std::sync::PoisonError::into_inner)
3995 .clone()
3996 .expect("release acceptance hook must have fired");
3997 assert_eq!(
3998 outcome,
3999 Ok(()),
4000 "expire fired at release acceptance must be a benign no-op"
4001 );
4002 assert_eq!(snapshot_phase(&lifecycle, &target), None);
4003 }
4004}