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