1use std::{
2 sync::{Arc, Mutex, MutexGuard},
3 time::{Duration, Instant},
4};
5
6use synd_protocol::{
7 CapabilitySet, capability,
8 daemon::{DaemonIdleShutdownStatus, DaemonSessionStatus},
9 session::{
10 CloseSessionErrorResponse, CloseSessionRequest, CloseSessionResponse,
11 OpenSessionErrorResponse, OpenSessionRequest, OpenSessionResponse,
12 RenewSessionErrorResponse, RenewSessionRequest, RenewSessionResponse, SessionLease,
13 },
14};
15use tracing::debug;
16
17mod decision;
18mod idle_shutdown;
19mod state;
20mod sweeper;
21
22pub use idle_shutdown::SessionIdleShutdown;
23pub(crate) use sweeper::DaemonSessionSweeper;
24
25use decision::{
26 SessionCloseContext, SessionCloseDecision, SessionOpenContext, SessionOpenDecision,
27 SessionRenewContext, SessionRenewDecision, SessionSweepDecision,
28};
29use state::{
30 DaemonSession, DaemonSessionsState, SessionIdIssuer, SessionLeaseDeadline, SessionSweepOutcome,
31};
32
33pub const DEFAULT_DAEMON_IDLE_SHUTDOWN_GRACE: Duration = Duration::from_secs(30);
34pub const DEFAULT_DAEMON_SESSION_LEASE_DURATION: Duration = Duration::from_secs(30);
35pub const DEFAULT_DAEMON_SESSION_SWEEP_INTERVAL: Duration = Duration::from_secs(5);
36
37#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub struct DaemonSessionConfig {
40 lease_policy: DaemonSessionLeasePolicy,
41 idle_shutdown_grace: Duration,
42}
43
44impl DaemonSessionConfig {
45 #[must_use]
46 pub fn new(lease_policy: DaemonSessionLeasePolicy, idle_shutdown_grace: Duration) -> Self {
47 Self {
48 lease_policy,
49 idle_shutdown_grace,
50 }
51 }
52
53 #[must_use]
54 pub fn with_lease_duration(self, lease_duration: Duration) -> Self {
55 Self {
56 lease_policy: DaemonSessionLeasePolicy::new(
57 lease_duration,
58 self.lease_policy.sweep_interval(),
59 ),
60 ..self
61 }
62 }
63
64 #[must_use]
65 pub fn with_idle_shutdown_grace(self, idle_shutdown_grace: Duration) -> Self {
66 Self {
67 idle_shutdown_grace,
68 ..self
69 }
70 }
71
72 #[must_use]
73 pub fn lease_policy(self) -> DaemonSessionLeasePolicy {
74 self.lease_policy
75 }
76
77 #[must_use]
78 pub fn idle_shutdown_grace(self) -> Duration {
79 self.idle_shutdown_grace
80 }
81}
82
83impl Default for DaemonSessionConfig {
84 fn default() -> Self {
85 Self::new(
86 DaemonSessionLeasePolicy::default(),
87 DEFAULT_DAEMON_IDLE_SHUTDOWN_GRACE,
88 )
89 }
90}
91
92#[derive(Debug, Clone)]
94pub struct DaemonSessions {
95 inner: Arc<DaemonSessionsInner>,
96}
97
98impl Default for DaemonSessions {
99 fn default() -> Self {
100 Self::new(capability::local_api_capabilities())
101 }
102}
103
104impl DaemonSessions {
105 pub fn new(supported_capabilities: CapabilitySet) -> Self {
106 Self::from_parts(
107 supported_capabilities,
108 DaemonSessionLeasePolicy::default(),
109 None,
110 )
111 }
112
113 #[must_use]
114 pub fn with_idle_shutdown(&self, idle_shutdown: SessionIdleShutdown) -> Self {
115 Self::from_parts(
116 self.supported_capabilities().clone(),
117 self.inner.lease_policy,
118 Some(idle_shutdown),
119 )
120 }
121
122 #[must_use]
123 pub fn with_lease_policy(&self, lease_policy: DaemonSessionLeasePolicy) -> Self {
124 Self::from_parts(
125 self.supported_capabilities().clone(),
126 lease_policy,
127 self.inner.idle_shutdown.clone(),
128 )
129 }
130
131 fn from_parts(
132 supported_capabilities: CapabilitySet,
133 lease_policy: DaemonSessionLeasePolicy,
134 idle_shutdown: Option<SessionIdleShutdown>,
135 ) -> Self {
136 Self {
137 inner: Arc::new(DaemonSessionsInner {
138 supported_capabilities,
139 lease_policy,
140 id_issuer: SessionIdIssuer::default(),
141 idle_shutdown,
142 state: Mutex::new(DaemonSessionsState::default()),
143 }),
144 }
145 }
146
147 pub fn open(
148 &self,
149 request: &OpenSessionRequest,
150 ) -> Result<OpenSessionResponse, OpenSessionErrorResponse> {
151 self.open_at(request, Instant::now())
152 }
153
154 fn open_at(
155 &self,
156 request: &OpenSessionRequest,
157 now: Instant,
158 ) -> Result<OpenSessionResponse, OpenSessionErrorResponse> {
159 let context = SessionOpenContext::new(
160 request.required_capabilities().clone(),
161 self.supported_capabilities().clone(),
162 );
163
164 match SessionOpenDecision::from(context) {
165 SessionOpenDecision::Accept {
166 available_capabilities,
167 } => {
168 let session_id = self.inner.id_issuer.issue();
169 let lease = self.inner.lease_policy.lease();
170 let lease_deadline = self.inner.lease_policy.deadline_from(now);
171 let (effect, active_sessions) = {
172 let mut state = self.lock_state();
173 let effect = state.insert(DaemonSession::new(
174 session_id.clone(),
175 request.required_capabilities().clone(),
176 lease_deadline,
177 ));
178
179 (effect, state.active_session_count())
180 };
181 effect.apply(self.inner.idle_shutdown.as_ref());
182 debug!(
183 %session_id,
184 active_sessions,
185 lease_duration_ms = lease.duration().as_millis(),
186 "Opened daemon session"
187 );
188
189 Ok(OpenSessionResponse::with_lease(
190 session_id,
191 available_capabilities,
192 lease,
193 ))
194 }
195 SessionOpenDecision::RejectMissingCapabilities {
196 missing_capabilities,
197 } => {
198 debug!(
199 missing_capabilities = ?missing_capabilities,
200 "Rejected daemon session open"
201 );
202 Err(OpenSessionErrorResponse::from_missing_capabilities(
203 missing_capabilities,
204 ))
205 }
206 }
207 }
208
209 pub fn renew(
210 &self,
211 request: &RenewSessionRequest,
212 ) -> Result<RenewSessionResponse, RenewSessionErrorResponse> {
213 self.renew_at(request, Instant::now())
214 }
215
216 fn renew_at(
217 &self,
218 request: &RenewSessionRequest,
219 now: Instant,
220 ) -> Result<RenewSessionResponse, RenewSessionErrorResponse> {
221 let session_id = request.session_id().clone();
222 let lease = self.inner.lease_policy.lease();
223 let lease_deadline = self.inner.lease_policy.deadline_from(now);
224 let context = {
225 let mut state = self.lock_state();
226 let renew = state.renew(
227 &session_id,
228 now,
229 lease_deadline,
230 self.inner.idle_shutdown.is_some(),
231 );
232
233 SessionRenewContext::new(session_id, lease, lease_deadline, renew)
234 };
235
236 match SessionRenewDecision::from(context) {
237 SessionRenewDecision::Accept {
238 session_id,
239 lease,
240 lease_deadline,
241 } => {
242 debug!(
243 %session_id,
244 lease_expires_in_ms = lease_deadline.remaining_from(now).as_millis(),
245 "Renewed daemon session"
246 );
247 Ok(RenewSessionResponse::new(session_id, lease))
248 }
249 SessionRenewDecision::RejectUnknownSession { session_id } => {
250 debug!(%session_id, "Rejected daemon session renew");
251 Err(RenewSessionErrorResponse::unknown_session(session_id))
252 }
253 SessionRenewDecision::RejectExpiredSession { session_id, effect } => {
254 effect.apply(self.inner.idle_shutdown.as_ref());
255 debug!(%session_id, "Expired daemon session during renew");
256 Err(RenewSessionErrorResponse::unknown_session(session_id))
257 }
258 }
259 }
260
261 pub fn close(
262 &self,
263 request: &CloseSessionRequest,
264 ) -> Result<CloseSessionResponse, CloseSessionErrorResponse> {
265 let session_id = request.session_id().clone();
266 let context = {
267 let mut state = self.lock_state();
268 let removed = state.remove(&session_id, self.inner.idle_shutdown.is_some());
269
270 SessionCloseContext::new(session_id.clone(), removed.known_session, removed.effect)
271 };
272
273 match SessionCloseDecision::from(context) {
274 SessionCloseDecision::Accept { effect } => {
275 effect.apply(self.inner.idle_shutdown.as_ref());
276 debug!(%session_id, "Closed daemon session");
277 Ok(CloseSessionResponse::new())
278 }
279 SessionCloseDecision::RejectUnknownSession { session_id } => {
280 debug!(%session_id, "Rejected daemon session close");
281 Err(CloseSessionErrorResponse::unknown_session(session_id))
282 }
283 }
284 }
285
286 fn sweep_interval(&self) -> Duration {
287 self.inner.lease_policy.sweep_interval()
288 }
289
290 fn sweep_expired(&self) -> SessionSweepOutcome {
291 self.sweep_expired_at(Instant::now())
292 }
293
294 fn sweep_expired_at(&self, now: Instant) -> SessionSweepOutcome {
295 let facts = {
296 let mut state = self.lock_state();
297 state.sweep_expired(now, self.inner.idle_shutdown.is_some())
298 };
299
300 match SessionSweepDecision::from(facts) {
301 SessionSweepDecision::NoExpiredSessions { active_sessions } => {
302 SessionSweepOutcome::new(0, active_sessions)
303 }
304 SessionSweepDecision::ExpiredSessions {
305 expired_sessions,
306 active_sessions,
307 effect,
308 } => {
309 let expired_session_count = expired_sessions.len();
310 effect.apply(self.inner.idle_shutdown.as_ref());
311 debug!(
312 expired_sessions = ?expired_sessions,
313 active_sessions,
314 "Expired daemon sessions during sweep"
315 );
316
317 SessionSweepOutcome::new(expired_session_count, active_sessions)
318 }
319 }
320 }
321
322 pub fn status(&self) -> DaemonSessionStatus {
323 let state = self.lock_state();
324 let idle_shutdown = match &self.inner.idle_shutdown {
325 Some(idle_shutdown) => DaemonIdleShutdownStatus::enabled(
326 idle_shutdown.grace(),
327 state.idle_shutdown_pending(),
328 ),
329 None => DaemonIdleShutdownStatus::disabled(),
330 };
331
332 DaemonSessionStatus::new(
333 state.active_session_count(),
334 self.inner.lease_policy.lease_duration(),
335 self.inner.lease_policy.sweep_interval(),
336 idle_shutdown,
337 )
338 }
339
340 fn supported_capabilities(&self) -> &CapabilitySet {
341 &self.inner.supported_capabilities
342 }
343
344 fn lock_state(&self) -> MutexGuard<'_, DaemonSessionsState> {
345 self.inner
346 .state
347 .lock()
348 .unwrap_or_else(std::sync::PoisonError::into_inner)
349 }
350
351 #[cfg(test)]
352 fn active_session_count(&self) -> usize {
353 self.lock_state().active_session_count()
354 }
355}
356
357#[derive(Debug)]
358struct DaemonSessionsInner {
359 supported_capabilities: CapabilitySet,
360 lease_policy: DaemonSessionLeasePolicy,
361 id_issuer: SessionIdIssuer,
362 idle_shutdown: Option<SessionIdleShutdown>,
363 state: Mutex<DaemonSessionsState>,
364}
365
366#[derive(Debug, Clone, Copy, PartialEq, Eq)]
368pub struct DaemonSessionLeasePolicy {
369 lease_duration: Duration,
370 sweep_interval: Duration,
371}
372
373impl DaemonSessionLeasePolicy {
374 #[must_use]
375 pub fn new(lease_duration: Duration, sweep_interval: Duration) -> Self {
376 Self {
377 lease_duration,
378 sweep_interval,
379 }
380 }
381
382 #[must_use]
383 fn lease(self) -> SessionLease {
384 SessionLease::new(self.lease_duration)
385 }
386
387 #[must_use]
388 fn deadline_from(self, now: Instant) -> SessionLeaseDeadline {
389 SessionLeaseDeadline::new(now + self.lease_duration)
390 }
391
392 #[must_use]
393 pub fn lease_duration(self) -> Duration {
394 self.lease_duration
395 }
396
397 #[must_use]
398 pub fn sweep_interval(self) -> Duration {
399 self.sweep_interval
400 }
401}
402
403impl Default for DaemonSessionLeasePolicy {
404 fn default() -> Self {
405 Self::new(
406 DEFAULT_DAEMON_SESSION_LEASE_DURATION,
407 DEFAULT_DAEMON_SESSION_SWEEP_INTERVAL,
408 )
409 }
410}
411
412#[cfg(test)]
413mod tests {
414 use std::{
415 sync::{
416 Arc,
417 atomic::{AtomicBool, Ordering},
418 },
419 time::{Duration, Instant},
420 };
421
422 use synd_protocol::{
423 CapabilitySet,
424 session::{CloseSessionRequest, OpenSessionRequest, RenewSessionRequest, SessionId},
425 };
426
427 use crate::shutdown::Shutdown;
428
429 use super::{
430 DaemonSessionLeasePolicy, DaemonSessions, SessionIdleShutdown,
431 decision::{
432 SessionCloseContext, SessionCloseDecision, SessionOpenContext, SessionOpenDecision,
433 SessionRenewContext, SessionRenewDecision, SessionSweepDecision,
434 },
435 idle_shutdown::DaemonSessionsEffect,
436 state::{SessionLeaseDeadline, SessionRenewChange, SessionSweepChange},
437 };
438
439 #[test]
440 fn selects_open_decision_from_capabilities() {
441 let cases = [
442 (
443 SessionOpenContext::new(
444 CapabilitySet::new(["timeline.read"]),
445 CapabilitySet::new(["timeline.read", "subscription.write"]),
446 ),
447 SessionOpenDecision::Accept {
448 available_capabilities: CapabilitySet::new([
449 "timeline.read",
450 "subscription.write",
451 ]),
452 },
453 ),
454 (
455 SessionOpenContext::new(
456 CapabilitySet::new(["timeline.read", "subscription.write"]),
457 CapabilitySet::new(["timeline.read"]),
458 ),
459 SessionOpenDecision::RejectMissingCapabilities {
460 missing_capabilities: CapabilitySet::new(["subscription.write"]),
461 },
462 ),
463 ];
464
465 for (context, expected) in cases {
466 assert_eq!(SessionOpenDecision::from(context), expected);
467 }
468 }
469
470 #[test]
471 fn opens_and_closes_session() {
472 let sessions = DaemonSessions::new(CapabilitySet::new(["timeline.read"]));
473
474 let opened = sessions
475 .open(&OpenSessionRequest::new(CapabilitySet::new([
476 "timeline.read",
477 ])))
478 .unwrap();
479
480 assert_eq!(sessions.active_session_count(), 1);
481
482 sessions
483 .close(&CloseSessionRequest::new(opened.session_id().clone()))
484 .unwrap();
485
486 assert_eq!(sessions.active_session_count(), 0);
487 }
488
489 #[test]
490 fn opens_session_with_lease_policy() {
491 let now = Instant::now();
492 let policy = DaemonSessionLeasePolicy::new(Duration::from_secs(12), Duration::from_secs(2));
493 let sessions = DaemonSessions::default().with_lease_policy(policy);
494
495 let opened = sessions
496 .open_at(&OpenSessionRequest::new(CapabilitySet::default()), now)
497 .unwrap();
498
499 assert_eq!(opened.lease().duration(), policy.lease_duration());
500 assert_eq!(sessions.active_session_count(), 1);
501 }
502
503 #[test]
504 fn rejects_session_when_required_capability_is_missing() {
505 let sessions = DaemonSessions::default();
506
507 let error = sessions
508 .open(&OpenSessionRequest::new(CapabilitySet::new([
509 "unsupported.capability",
510 ])))
511 .unwrap_err();
512
513 assert_eq!(
514 error.missing_capabilities().names(),
515 ["unsupported.capability"]
516 );
517 assert_eq!(sessions.active_session_count(), 0);
518 }
519
520 #[test]
521 fn renews_session_before_lease_deadline() {
522 let now = Instant::now();
523 let policy = DaemonSessionLeasePolicy::new(Duration::from_secs(30), Duration::from_secs(5));
524 let sessions = DaemonSessions::default().with_lease_policy(policy);
525 let opened = sessions
526 .open_at(&OpenSessionRequest::new(CapabilitySet::default()), now)
527 .unwrap();
528
529 let renewed = sessions
530 .renew_at(
531 &RenewSessionRequest::new(opened.session_id().clone()),
532 now + Duration::from_secs(10),
533 )
534 .unwrap();
535
536 assert_eq!(renewed.session_id(), opened.session_id());
537 assert_eq!(renewed.lease().duration(), policy.lease_duration());
538 assert_eq!(sessions.active_session_count(), 1);
539 }
540
541 #[test]
542 fn rejects_and_removes_expired_session_during_renew() {
543 let now = Instant::now();
544 let policy = DaemonSessionLeasePolicy::new(Duration::from_secs(30), Duration::from_secs(5));
545 let sessions = DaemonSessions::default().with_lease_policy(policy);
546 let opened = sessions
547 .open_at(&OpenSessionRequest::new(CapabilitySet::default()), now)
548 .unwrap();
549
550 let error = sessions
551 .renew_at(
552 &RenewSessionRequest::new(opened.session_id().clone()),
553 now + Duration::from_secs(31),
554 )
555 .unwrap_err();
556
557 assert_eq!(error.session_id(), opened.session_id());
558 assert_eq!(sessions.active_session_count(), 0);
559 }
560
561 #[test]
562 fn rejects_unknown_session_during_renew() {
563 let sessions = DaemonSessions::default();
564 let session_id = SessionId::new("session-unknown");
565
566 let error = sessions
567 .renew(&RenewSessionRequest::new(session_id.clone()))
568 .unwrap_err();
569
570 assert_eq!(error.session_id(), &session_id);
571 }
572
573 #[test]
574 fn selects_renew_decision_from_change() {
575 let now = Instant::now();
576 let lease = synd_protocol::session::SessionLease::new(Duration::from_secs(30));
577 let lease_deadline = SessionLeaseDeadline::new(now + Duration::from_secs(30));
578 let cases = [
579 (
580 SessionRenewContext::new(
581 SessionId::new("session-1"),
582 lease,
583 lease_deadline,
584 SessionRenewChange::renewed(),
585 ),
586 SessionRenewDecision::Accept {
587 session_id: SessionId::new("session-1"),
588 lease,
589 lease_deadline,
590 },
591 ),
592 (
593 SessionRenewContext::new(
594 SessionId::new("session-1"),
595 lease,
596 lease_deadline,
597 SessionRenewChange::expired(DaemonSessionsEffect::None),
598 ),
599 SessionRenewDecision::RejectExpiredSession {
600 session_id: SessionId::new("session-1"),
601 effect: DaemonSessionsEffect::None,
602 },
603 ),
604 (
605 SessionRenewContext::new(
606 SessionId::new("session-1"),
607 lease,
608 lease_deadline,
609 SessionRenewChange::unknown(),
610 ),
611 SessionRenewDecision::RejectUnknownSession {
612 session_id: SessionId::new("session-1"),
613 },
614 ),
615 ];
616
617 for (context, expected) in cases {
618 assert_eq!(SessionRenewDecision::from(context), expected);
619 }
620 }
621
622 #[test]
623 fn selects_close_decision_from_known_session() {
624 let cases = [
625 (
626 SessionCloseContext::new(
627 SessionId::new("session-1"),
628 true,
629 DaemonSessionsEffect::None,
630 ),
631 SessionCloseDecision::Accept {
632 effect: DaemonSessionsEffect::None,
633 },
634 ),
635 (
636 SessionCloseContext::new(
637 SessionId::new("session-1"),
638 false,
639 DaemonSessionsEffect::None,
640 ),
641 SessionCloseDecision::RejectUnknownSession {
642 session_id: SessionId::new("session-1"),
643 },
644 ),
645 ];
646
647 for (context, expected) in cases {
648 assert_eq!(SessionCloseDecision::from(context), expected);
649 }
650 }
651
652 #[test]
653 fn sweeps_expired_sessions() {
654 let now = Instant::now();
655 let policy = DaemonSessionLeasePolicy::new(Duration::from_secs(30), Duration::from_secs(5));
656 let sessions = DaemonSessions::default().with_lease_policy(policy);
657 sessions
658 .open_at(&OpenSessionRequest::new(CapabilitySet::default()), now)
659 .unwrap();
660 sessions
661 .open_at(
662 &OpenSessionRequest::new(CapabilitySet::default()),
663 now + Duration::from_secs(10),
664 )
665 .unwrap();
666
667 let outcome = sessions.sweep_expired_at(now + Duration::from_secs(31));
668
669 assert_eq!(outcome.expired_session_count(), 1);
670 assert_eq!(outcome.active_sessions(), 1);
671 assert_eq!(sessions.active_session_count(), 1);
672 }
673
674 #[test]
675 fn selects_sweep_decision_from_change() {
676 let cases = [
677 (
678 SessionSweepChange::new(vec![], 2, DaemonSessionsEffect::None),
679 SessionSweepDecision::NoExpiredSessions { active_sessions: 2 },
680 ),
681 (
682 SessionSweepChange::new(
683 vec![SessionId::new("session-1")],
684 0,
685 DaemonSessionsEffect::None,
686 ),
687 SessionSweepDecision::ExpiredSessions {
688 expired_sessions: vec![SessionId::new("session-1")],
689 active_sessions: 0,
690 effect: DaemonSessionsEffect::None,
691 },
692 ),
693 ];
694
695 for (facts, expected) in cases {
696 assert_eq!(SessionSweepDecision::from(facts), expected);
697 }
698 }
699
700 #[tokio::test]
701 async fn schedules_idle_shutdown_after_last_session_closes() {
702 let shutdown_called = Arc::new(AtomicBool::new(false));
703 let shutdown_called_for_hook = Arc::clone(&shutdown_called);
704 let shutdown = Shutdown::manual(move || {
705 shutdown_called_for_hook.store(true, Ordering::Relaxed);
706 });
707 let sessions = DaemonSessions::default().with_idle_shutdown(SessionIdleShutdown::new(
708 Duration::from_millis(10),
709 shutdown,
710 ));
711 let opened = sessions
712 .open(&OpenSessionRequest::new(CapabilitySet::default()))
713 .unwrap();
714
715 sessions
716 .close(&CloseSessionRequest::new(opened.session_id().clone()))
717 .unwrap();
718
719 tokio::time::timeout(Duration::from_secs(1), async {
720 loop {
721 if shutdown_called.load(Ordering::Relaxed) {
722 return;
723 }
724
725 tokio::time::sleep(Duration::from_millis(10)).await;
726 }
727 })
728 .await
729 .unwrap();
730 }
731
732 #[tokio::test]
733 async fn opening_session_cancels_pending_idle_shutdown() {
734 let shutdown_called = Arc::new(AtomicBool::new(false));
735 let shutdown_called_for_hook = Arc::clone(&shutdown_called);
736 let shutdown = Shutdown::manual(move || {
737 shutdown_called_for_hook.store(true, Ordering::Relaxed);
738 });
739 let sessions = DaemonSessions::default().with_idle_shutdown(SessionIdleShutdown::new(
740 Duration::from_millis(20),
741 shutdown,
742 ));
743 let first = sessions
744 .open(&OpenSessionRequest::new(CapabilitySet::default()))
745 .unwrap();
746
747 sessions
748 .close(&CloseSessionRequest::new(first.session_id().clone()))
749 .unwrap();
750 let _second = sessions
751 .open(&OpenSessionRequest::new(CapabilitySet::default()))
752 .unwrap();
753 tokio::time::sleep(Duration::from_millis(60)).await;
754
755 assert!(!shutdown_called.load(Ordering::Relaxed));
756 }
757
758 #[tokio::test]
759 async fn schedules_idle_shutdown_after_last_session_expires_during_sweep() {
760 let now = Instant::now();
761 let shutdown_called = Arc::new(AtomicBool::new(false));
762 let shutdown_called_for_hook = Arc::clone(&shutdown_called);
763 let shutdown = Shutdown::manual(move || {
764 shutdown_called_for_hook.store(true, Ordering::Relaxed);
765 });
766 let policy =
767 DaemonSessionLeasePolicy::new(Duration::from_millis(10), Duration::from_millis(5));
768 let sessions = DaemonSessions::default()
769 .with_lease_policy(policy)
770 .with_idle_shutdown(SessionIdleShutdown::new(
771 Duration::from_millis(10),
772 shutdown,
773 ));
774 sessions
775 .open_at(&OpenSessionRequest::new(CapabilitySet::default()), now)
776 .unwrap();
777
778 let outcome = sessions.sweep_expired_at(now + Duration::from_millis(11));
779
780 assert_eq!(outcome.expired_session_count(), 1);
781 assert_eq!(outcome.active_sessions(), 0);
782 tokio::time::timeout(Duration::from_secs(1), async {
783 loop {
784 if shutdown_called.load(Ordering::Relaxed) {
785 return;
786 }
787
788 tokio::time::sleep(Duration::from_millis(10)).await;
789 }
790 })
791 .await
792 .unwrap();
793 }
794}