1use std::collections::HashMap;
8use std::fmt;
9use std::future::Future;
10use std::pin::Pin;
11use std::sync::{Arc, RwLock as StdRwLock};
12
13use async_trait::async_trait;
14use bamboo_domain::{
15 SessionActivationDisposition, SessionActivationError, SessionActivationPort, SessionInboxPort,
16};
17use tokio::sync::{mpsc, watch, Mutex, RwLock};
18
19pub struct SessionActivationLaunch {
24 pub run_id: String,
25 launch: Option<Box<dyn FnOnce() + Send + 'static>>,
26 rollback: Option<
27 Box<dyn FnOnce() -> Pin<Box<dyn Future<Output = ()> + Send + 'static>> + Send + 'static>,
28 >,
29 rollback_completion: Option<Arc<RollbackCompletion>>,
30}
31
32struct RollbackCompletion {
33 completed: watch::Sender<bool>,
34}
35
36impl RollbackCompletion {
37 fn new() -> Arc<Self> {
38 let (completed, _receiver) = watch::channel(false);
39 Arc::new(Self { completed })
40 }
41
42 fn complete(&self) {
43 self.completed.send_replace(true);
44 }
45
46 async fn wait(&self) {
47 let mut completed = self.completed.subscribe();
48 while !*completed.borrow() {
49 if completed.changed().await.is_err() {
50 break;
51 }
52 }
53 }
54}
55
56impl SessionActivationLaunch {
57 pub fn new(run_id: impl Into<String>, launch: impl FnOnce() + Send + 'static) -> Self {
58 Self {
59 run_id: run_id.into(),
60 launch: Some(Box::new(launch)),
61 rollback: None,
62 rollback_completion: None,
63 }
64 }
65
66 pub fn new_with_rollback(
72 run_id: impl Into<String>,
73 launch: impl FnOnce() + Send + 'static,
74 rollback: impl FnOnce() + Send + 'static,
75 ) -> Self {
76 Self::new_with_async_rollback(run_id, launch, move || async move {
77 rollback();
78 })
79 }
80
81 pub fn new_with_async_rollback<F, Fut>(
86 run_id: impl Into<String>,
87 launch: impl FnOnce() + Send + 'static,
88 rollback: F,
89 ) -> Self
90 where
91 F: FnOnce() -> Fut + Send + 'static,
92 Fut: Future<Output = ()> + Send + 'static,
93 {
94 let rollback_completion = RollbackCompletion::new();
95 Self {
96 run_id: run_id.into(),
97 launch: Some(Box::new(launch)),
98 rollback: Some(Box::new(move || Box::pin(rollback()))),
99 rollback_completion: Some(rollback_completion),
100 }
101 }
102
103 fn rollback_completion(&self) -> Option<Arc<RollbackCompletion>> {
104 self.rollback_completion.clone()
105 }
106
107 fn launch(mut self) {
108 if let Some(launch) = self.launch.take() {
109 self.rollback = None;
110 self.rollback_completion = None;
111 launch();
112 }
113 }
114}
115
116impl Drop for SessionActivationLaunch {
117 fn drop(&mut self) {
118 if self.launch.is_some() {
119 if let Some(rollback) = self.rollback.take() {
120 let completion = self.rollback_completion.take();
121 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
122 runtime.spawn(async move {
123 rollback().await;
124 if let Some(completion) = completion {
125 completion.complete();
126 }
127 });
128 }
129 }
130 }
131 }
132}
133
134impl fmt::Debug for SessionActivationLaunch {
135 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
136 formatter
137 .debug_struct("SessionActivationLaunch")
138 .field("run_id", &self.run_id)
139 .finish_non_exhaustive()
140 }
141}
142
143#[derive(Debug)]
144pub enum SessionActivationReserveOutcome {
145 Reserved(SessionActivationLaunch),
148 AlreadyRunning { run_id: String },
150 NotFound,
152 NoWork,
154}
155
156#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
162pub enum SessionRunRegistrationError {
163 #[error(
164 "session activation owner collision for {target_session_id}: \
165 existing run {existing_run_id}, attempted run {attempted_run_id}"
166 )]
167 OwnerCollision {
168 target_session_id: String,
169 existing_run_id: String,
170 attempted_run_id: String,
171 },
172}
173
174impl SessionRunRegistrationError {
175 pub fn existing_run_id(&self) -> &str {
177 match self {
178 Self::OwnerCollision {
179 existing_run_id, ..
180 } => existing_run_id,
181 }
182 }
183}
184
185#[async_trait]
191pub trait SessionActivationSpawner: Send + Sync {
192 async fn reserve_activation(
193 &self,
194 target_session_id: &str,
195 inbox_generation: u64,
196 ) -> Result<SessionActivationReserveOutcome, SessionActivationError>;
197}
198
199#[derive(Debug, Clone)]
200struct ActiveOwner {
201 run_id: String,
202 finalizing: bool,
203 registrations: usize,
204 delivery_sink: Option<mpsc::UnboundedSender<u64>>,
208}
209
210#[derive(Debug)]
211struct TargetActivationState {
212 latest_generation: u64,
213 last_dispatched_generation: u64,
220 owner: Option<ActiveOwner>,
221 activation_reserved: bool,
222 activation_token: u64,
225 activation_epoch: watch::Sender<u64>,
230 notify: watch::Sender<u64>,
231}
232
233impl Default for TargetActivationState {
234 fn default() -> Self {
235 let (activation_epoch, _activation_receiver) = watch::channel(0);
236 let (notify, _receiver) = watch::channel(0);
237 Self {
238 latest_generation: 0,
239 last_dispatched_generation: 0,
240 owner: None,
241 activation_reserved: false,
242 activation_token: 0,
243 activation_epoch,
244 notify,
245 }
246 }
247}
248
249fn reserve_activation_token(state: &mut TargetActivationState) -> u64 {
250 state.activation_reserved = true;
251 state.activation_token = state.activation_token.wrapping_add(1).max(1);
252 state.activation_token
253}
254
255fn release_activation_token(state: &mut TargetActivationState, token: u64) -> bool {
256 if !state.activation_reserved || state.activation_token != token {
257 return false;
258 }
259 state.activation_reserved = false;
260 let next_epoch = (*state.activation_epoch.borrow()).wrapping_add(1);
261 state.activation_epoch.send_replace(next_epoch);
262 true
263}
264
265struct ActivationReservationLease {
271 router: SessionActivationRouter,
272 target_session_id: String,
273 token: u64,
274 recover_on_drop: bool,
275 armed: bool,
276 rollback_completion: Option<Arc<RollbackCompletion>>,
277}
278
279impl ActivationReservationLease {
280 fn new(
281 router: SessionActivationRouter,
282 target_session_id: &str,
283 token: u64,
284 recover_on_drop: bool,
285 ) -> Self {
286 Self {
287 router,
288 target_session_id: target_session_id.to_string(),
289 token,
290 recover_on_drop,
291 armed: true,
292 rollback_completion: None,
293 }
294 }
295
296 fn wait_for_rollback(&mut self, completion: Option<Arc<RollbackCompletion>>) {
297 self.rollback_completion = completion;
298 }
299
300 fn disarm(&mut self) {
301 self.armed = false;
302 }
303}
304
305impl Drop for ActivationReservationLease {
306 fn drop(&mut self) {
307 if !self.armed {
308 return;
309 }
310 let router = self.router.clone();
311 let target_session_id = self.target_session_id.clone();
312 let token = self.token;
313 let recover_on_drop = self.recover_on_drop;
314 let rollback_completion = self.rollback_completion.take();
315 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
316 runtime.spawn(async move {
317 if let Some(completion) = rollback_completion {
318 completion.wait().await;
319 }
320 router
321 .recover_cancelled_reservation(&target_session_id, token, recover_on_drop)
322 .await;
323 });
324 }
325 }
326}
327
328#[derive(Clone, Default)]
330pub struct SessionActivationRouter {
331 states: Arc<Mutex<HashMap<String, TargetActivationState>>>,
332 spawner: Arc<RwLock<Option<Arc<dyn SessionActivationSpawner>>>>,
333 inbox: Arc<StdRwLock<Option<Arc<dyn SessionInboxPort>>>>,
334}
335
336type AbortCleanup =
337 Box<dyn FnOnce() -> Pin<Box<dyn Future<Output = ()> + Send + 'static>> + Send + 'static>;
338
339pub struct SessionRunRegistration {
347 router: Arc<SessionActivationRouter>,
348 target_session_id: String,
349 run_id: String,
350 notifications: Option<watch::Receiver<u64>>,
351 abort_cleanup: Option<AbortCleanup>,
352 armed: bool,
353}
354
355impl SessionRunRegistration {
356 pub fn notifications_mut(&mut self) -> &mut watch::Receiver<u64> {
357 self.notifications
358 .as_mut()
359 .expect("run registration notifications are live before finalization")
360 }
361
362 pub fn set_abort_cleanup<F, Fut>(&mut self, cleanup: F)
365 where
366 F: FnOnce() -> Fut + Send + 'static,
367 Fut: Future<Output = ()> + Send + 'static,
368 {
369 self.abort_cleanup = Some(Box::new(move || Box::pin(cleanup())));
370 }
371
372 pub async fn begin_finalization(&mut self) {
373 self.router
374 .begin_finalization(&self.target_session_id, &self.run_id)
375 .await;
376 }
377
378 pub async fn abandon(mut self) {
387 self.armed = false;
388 let abort_cleanup = self.abort_cleanup.take();
389 let router = self.router.clone();
390 let target_session_id = self.target_session_id.clone();
391 let run_id = self.run_id.clone();
392 let cleanup_target_session_id = target_session_id.clone();
393 let cleanup_run_id = run_id.clone();
394 let cleanup = tokio::spawn(async move {
395 if let Some(cleanup) = abort_cleanup {
396 cleanup().await;
397 }
398 router
399 .cleanup_abandoned_registration(&cleanup_target_session_id, &cleanup_run_id)
400 .await;
401 });
402 if let Err(error) = cleanup.await {
403 tracing::error!(
404 %target_session_id,
405 %run_id,
406 %error,
407 "detached SessionInbox registration cleanup failed"
408 );
409 }
410 }
411
412 pub async fn finish(
415 mut self,
416 admitted_generation: u64,
417 ) -> Result<Option<SessionActivationDisposition>, SessionActivationError> {
418 drop(self.notifications.take());
421 let result = self
422 .router
423 .finish_finalization(&self.target_session_id, &self.run_id, admitted_generation)
424 .await;
425 self.armed = false;
426 self.abort_cleanup = None;
427 result
428 }
429}
430
431impl Drop for SessionRunRegistration {
432 fn drop(&mut self) {
433 if !self.armed {
434 return;
435 }
436 let router = self.router.clone();
437 let target_session_id = self.target_session_id.clone();
438 let run_id = self.run_id.clone();
439 let abort_cleanup = self.abort_cleanup.take();
440 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
441 runtime.spawn(async move {
442 if let Some(cleanup) = abort_cleanup {
443 cleanup().await;
444 }
445 router
446 .cleanup_abandoned_registration(&target_session_id, &run_id)
447 .await;
448 });
449 }
450 }
451}
452
453impl SessionActivationRouter {
454 pub fn new() -> Arc<Self> {
455 Arc::new(Self::default())
456 }
457
458 pub async fn set_spawner(&self, spawner: Arc<dyn SessionActivationSpawner>) {
461 *self.spawner.write().await = Some(spawner);
462 }
463
464 pub fn set_inbox(&self, inbox: Arc<dyn SessionInboxPort>) {
466 *self
467 .inbox
468 .write()
469 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(inbox);
470 }
471
472 #[cfg(test)]
473 pub(crate) async fn hold_state_lock_for_test(
474 &self,
475 entered: Arc<tokio::sync::Notify>,
476 release: Arc<tokio::sync::Notify>,
477 ) {
478 let _states = self.states.lock().await;
479 entered.notify_one();
480 release.notified().await;
481 }
482
483 #[cfg(test)]
484 pub(crate) async fn install_owner_placeholder_for_test(
485 &self,
486 target_session_id: &str,
487 run_id: &str,
488 ) {
489 let mut states = self.states.lock().await;
490 states
491 .entry(target_session_id.to_string())
492 .or_default()
493 .owner = Some(ActiveOwner {
494 run_id: run_id.to_string(),
495 finalizing: false,
496 registrations: 0,
497 delivery_sink: None,
498 });
499 }
500
501 pub async fn register_run(
505 self: &Arc<Self>,
506 target_session_id: &str,
507 run_id: &str,
508 ) -> Result<SessionRunRegistration, SessionRunRegistrationError> {
509 loop {
510 let (notifications, reservation_wait) = {
511 let mut states = self.states.lock().await;
512 let state = states.entry(target_session_id.to_string()).or_default();
513 match state.owner.as_mut() {
514 Some(owner)
515 if owner.run_id == run_id
516 && owner.registrations == 0
517 && !owner.finalizing =>
518 {
519 owner.registrations = 1;
523 (Some(state.notify.subscribe()), None)
524 }
525 Some(owner) => {
526 return Err(SessionRunRegistrationError::OwnerCollision {
527 target_session_id: target_session_id.to_string(),
528 existing_run_id: owner.run_id.clone(),
529 attempted_run_id: run_id.to_string(),
530 });
531 }
532 None if state.activation_reserved => {
533 (None, Some(state.activation_epoch.subscribe()))
539 }
540 None => {
541 state.owner = Some(ActiveOwner {
542 run_id: run_id.to_string(),
543 finalizing: false,
544 registrations: 1,
545 delivery_sink: None,
546 });
547 (Some(state.notify.subscribe()), None)
548 }
549 }
550 };
551
552 if let Some(mut reservation_wait) = reservation_wait {
553 let _ = reservation_wait.changed().await;
554 continue;
555 }
556 return Ok(SessionRunRegistration {
557 router: self.clone(),
558 target_session_id: target_session_id.to_string(),
559 run_id: run_id.to_string(),
560 notifications,
561 abort_cleanup: None,
562 armed: true,
563 });
564 }
565 }
566
567 pub async fn attach_delivery_sink(
573 &self,
574 target_session_id: &str,
575 sink: mpsc::UnboundedSender<u64>,
576 ) -> Option<String> {
577 let mut states = self.states.lock().await;
578 let state = states.get_mut(target_session_id)?;
579 let owner = state.owner.as_mut()?;
580 if owner.finalizing {
581 return None;
582 }
583 owner.delivery_sink = Some(sink.clone());
584 let run_id = owner.run_id.clone();
585 if state.latest_generation > 0 && sink.send(state.latest_generation).is_err() {
586 owner.delivery_sink = None;
587 return None;
588 }
589 Some(run_id)
590 }
591
592 pub async fn detach_delivery_sink(&self, target_session_id: &str, run_id: &str) {
595 let mut states = self.states.lock().await;
596 let Some(state) = states.get_mut(target_session_id) else {
597 return;
598 };
599 let Some(owner) = state.owner.as_mut() else {
600 return;
601 };
602 if owner.run_id == run_id {
603 owner.delivery_sink = None;
604 }
605 }
606
607 pub async fn owns_run(&self, target_session_id: &str, run_id: &str) -> bool {
609 let states = self.states.lock().await;
610 states
611 .get(target_session_id)
612 .and_then(|state| state.owner.as_ref())
613 .is_some_and(|owner| owner.run_id == run_id && !owner.finalizing)
614 }
615
616 pub async fn begin_finalization(&self, target_session_id: &str, run_id: &str) {
621 let mut states = self.states.lock().await;
622 let state = states.entry(target_session_id.to_string()).or_default();
623 match state.owner.as_mut() {
624 Some(owner) if owner.run_id == run_id => owner.finalizing = true,
625 Some(_) => {}
626 None => {
627 state.owner = Some(ActiveOwner {
628 run_id: run_id.to_string(),
629 finalizing: true,
630 registrations: 0,
631 delivery_sink: None,
632 });
633 }
634 }
635 }
636
637 async fn cleanup_abandoned_registration(
638 self: &Arc<Self>,
639 target_session_id: &str,
640 run_id: &str,
641 ) {
642 {
643 let mut states = self.states.lock().await;
644 let Some(state) = states.get_mut(target_session_id) else {
645 return;
646 };
647 let Some(owner) = state.owner.as_mut() else {
648 return;
649 };
650 if owner.run_id != run_id {
651 return;
652 }
653 owner.registrations = 0;
654 owner.finalizing = true;
655 }
656
657 let inbox = self
662 .inbox
663 .read()
664 .unwrap_or_else(std::sync::PoisonError::into_inner)
665 .clone();
666 let durable_generation = match inbox {
667 Some(inbox) => match inbox.inspect(target_session_id).await {
668 Ok(backlog) if backlog.activation_pending() => backlog.activation_generation,
669 Ok(_) => 0,
670 Err(error) => {
671 tracing::warn!(
672 %target_session_id,
673 %run_id,
674 %error,
675 "failed to reconcile durable SessionInbox generation for abandoned run"
676 );
677 0
678 }
679 },
680 None => 0,
681 };
682
683 let reservation_to_dispatch = {
684 let mut states = self.states.lock().await;
685 let Some(state) = states.get_mut(target_session_id) else {
686 return;
687 };
688 if !state.owner.as_ref().is_some_and(|owner| {
689 owner.run_id == run_id && owner.finalizing && owner.registrations == 0
690 }) {
691 return;
692 }
693 state.latest_generation = state.latest_generation.max(durable_generation);
694 state.owner = None;
695 if state.latest_generation > state.last_dispatched_generation
696 && !state.activation_reserved
697 {
698 let generation = state.latest_generation;
699 let token = reserve_activation_token(state);
700 Some((generation, token))
701 } else {
702 if state.latest_generation == 0
703 && !state.activation_reserved
704 && state.notify.receiver_count() == 0
705 && state.activation_epoch.receiver_count() == 0
706 {
707 states.remove(target_session_id);
708 }
709 None
710 }
711 };
712
713 if let Some((generation, token)) = reservation_to_dispatch {
714 if let Err(error) = self
715 .dispatch_reserved(target_session_id, generation, token)
716 .await
717 {
718 tracing::error!(
719 %target_session_id,
720 %run_id,
721 %error,
722 "failed to reserve successor for abandoned SessionInbox owner"
723 );
724 }
725 }
726 }
727
728 pub async fn finish_finalization(
732 &self,
733 target_session_id: &str,
734 run_id: &str,
735 admitted_generation: u64,
736 ) -> Result<Option<SessionActivationDisposition>, SessionActivationError> {
737 let reservation_to_dispatch = {
738 let mut states = self.states.lock().await;
739 let state = states.entry(target_session_id.to_string()).or_default();
740 if state
741 .owner
742 .as_ref()
743 .is_some_and(|owner| owner.run_id != run_id)
744 {
745 return Ok(None);
746 }
747 state.owner = None;
748 if state.latest_generation > admitted_generation
749 && state.latest_generation > state.last_dispatched_generation
750 && !state.activation_reserved
751 {
752 let generation = state.latest_generation;
753 let token = reserve_activation_token(state);
754 Some((generation, token))
755 } else {
756 if state.latest_generation <= admitted_generation
761 && state.owner.is_none()
762 && !state.activation_reserved
763 && state.notify.receiver_count() == 0
764 && state.activation_epoch.receiver_count() == 0
765 {
766 states.remove(target_session_id);
767 }
768 None
769 }
770 };
771
772 match reservation_to_dispatch {
773 Some((generation, token)) => self
774 .dispatch_reserved(target_session_id, generation, token)
775 .await
776 .map(Some),
777 None => Ok(None),
778 }
779 }
780
781 pub async fn subscribe(&self, target_session_id: &str) -> watch::Receiver<u64> {
782 let mut states = self.states.lock().await;
783 states
784 .entry(target_session_id.to_string())
785 .or_default()
786 .notify
787 .subscribe()
788 }
789
790 async fn recover_cancelled_reservation(
799 &self,
800 target_session_id: &str,
801 reservation_token: u64,
802 recover_on_drop: bool,
803 ) {
804 let reservation_to_dispatch = {
805 let mut states = self.states.lock().await;
806 let Some(state) = states.get_mut(target_session_id) else {
807 return;
808 };
809 if !release_activation_token(state, reservation_token) || !recover_on_drop {
810 return;
811 }
812 if state.owner.is_none()
813 && state.latest_generation > state.last_dispatched_generation
814 && !state.activation_reserved
815 {
816 let generation = state.latest_generation;
817 let token = reserve_activation_token(state);
818 Some((generation, token))
819 } else {
820 None
821 }
822 };
823
824 if let Some((generation, token)) = reservation_to_dispatch {
825 if let Err(error) = self
826 .dispatch_reserved_with_recovery(target_session_id, generation, token, false)
827 .await
828 {
829 tracing::error!(
830 %target_session_id,
831 %error,
832 "failed bounded retry after cancelled SessionInbox activation reservation"
833 );
834 }
835 }
836 }
837
838 async fn dispatch_reserved(
839 &self,
840 target_session_id: &str,
841 generation: u64,
842 reservation_token: u64,
843 ) -> Result<SessionActivationDisposition, SessionActivationError> {
844 self.dispatch_reserved_with_recovery(target_session_id, generation, reservation_token, true)
845 .await
846 }
847
848 async fn dispatch_reserved_with_recovery(
849 &self,
850 target_session_id: &str,
851 generation: u64,
852 reservation_token: u64,
853 recover_on_drop: bool,
854 ) -> Result<SessionActivationDisposition, SessionActivationError> {
855 let mut lease = ActivationReservationLease::new(
856 self.clone(),
857 target_session_id,
858 reservation_token,
859 recover_on_drop,
860 );
861 let Some(spawner) = self.spawner.read().await.clone() else {
862 let mut states = self.states.lock().await;
863 if let Some(state) = states.get_mut(target_session_id) {
864 release_activation_token(state, reservation_token);
865 }
866 lease.disarm();
867 return Err(SessionActivationError::Internal(
868 "session activation spawner is not configured".to_string(),
869 ));
870 };
871
872 let outcome = spawner
873 .reserve_activation(target_session_id, generation)
874 .await;
875 if let Ok(SessionActivationReserveOutcome::Reserved(launch)) = &outcome {
876 lease.wait_for_rollback(launch.rollback_completion());
877 }
878 match outcome {
879 Ok(SessionActivationReserveOutcome::Reserved(launch)) => {
880 let run_id = launch.run_id.clone();
881 {
882 let mut states = self.states.lock().await;
883 let state = states.entry(target_session_id.to_string()).or_default();
884 if !release_activation_token(state, reservation_token) {
885 lease.disarm();
893 return if state.owner.is_some() {
894 Ok(SessionActivationDisposition::ActiveNotified)
895 } else {
896 Err(SessionActivationError::Internal(
897 "session activation reservation was superseded".to_string(),
898 ))
899 };
900 }
901 state.owner = Some(ActiveOwner {
902 run_id,
903 finalizing: false,
904 registrations: 0,
905 delivery_sink: None,
906 });
907 state.last_dispatched_generation =
908 state.last_dispatched_generation.max(generation);
909 }
910 lease.disarm();
911 launch.launch();
914 Ok(SessionActivationDisposition::ActivationReserved)
915 }
916 Ok(SessionActivationReserveOutcome::AlreadyRunning { run_id }) => {
917 let mut states = self.states.lock().await;
918 let state = states.entry(target_session_id.to_string()).or_default();
919 if !release_activation_token(state, reservation_token) {
920 lease.disarm();
921 return Ok(SessionActivationDisposition::ActiveNotified);
922 }
923 state.owner = Some(ActiveOwner {
924 run_id,
925 finalizing: false,
926 registrations: 0,
927 delivery_sink: None,
928 });
929 state.notify.send_replace(generation);
930 lease.disarm();
931 Ok(SessionActivationDisposition::ActiveNotified)
932 }
933 Ok(SessionActivationReserveOutcome::NoWork) => {
934 let mut states = self.states.lock().await;
935 if let Some(state) = states.get_mut(target_session_id) {
936 release_activation_token(state, reservation_token);
937 }
938 lease.disarm();
939 Ok(SessionActivationDisposition::ActivationCoalesced)
940 }
941 Ok(SessionActivationReserveOutcome::NotFound) => {
942 let mut states = self.states.lock().await;
943 if states
944 .get(target_session_id)
945 .is_some_and(|state| state.activation_token == reservation_token)
946 {
947 states.remove(target_session_id);
948 }
949 lease.disarm();
950 Err(SessionActivationError::TargetNotFound(
951 target_session_id.to_string(),
952 ))
953 }
954 Err(error) => {
955 let mut states = self.states.lock().await;
956 if let Some(state) = states.get_mut(target_session_id) {
957 release_activation_token(state, reservation_token);
958 }
959 lease.disarm();
960 Err(error)
961 }
962 }
963 }
964}
965
966#[async_trait]
967impl SessionActivationPort for SessionActivationRouter {
968 async fn request_activation(
969 &self,
970 target_session_id: &str,
971 inbox_generation: u64,
972 ) -> Result<SessionActivationDisposition, SessionActivationError> {
973 loop {
974 let (reservation_wait, reservation_token) = {
975 let mut states = self.states.lock().await;
976 let state = states.entry(target_session_id.to_string()).or_default();
977 state.latest_generation = state.latest_generation.max(inbox_generation);
978 if let Some(owner) = state.owner.as_mut() {
979 if !owner.finalizing {
980 state.notify.send_replace(inbox_generation);
981 if owner
982 .delivery_sink
983 .as_ref()
984 .is_some_and(|sink| sink.send(inbox_generation).is_err())
985 {
986 owner.delivery_sink = None;
987 }
988 return Ok(SessionActivationDisposition::ActiveNotified);
989 }
990 return Ok(SessionActivationDisposition::ActivationCoalesced);
991 }
992 if inbox_generation <= state.last_dispatched_generation {
993 return Ok(SessionActivationDisposition::ActivationCoalesced);
994 }
995 if state.activation_reserved {
996 (Some(state.activation_epoch.subscribe()), None)
997 } else {
998 let token = reserve_activation_token(state);
999 (None, Some(token))
1000 }
1001 };
1002
1003 if let Some(mut reservation_wait) = reservation_wait {
1004 if reservation_wait.changed().await.is_err() {
1010 return Err(SessionActivationError::TargetNotFound(
1011 target_session_id.to_string(),
1012 ));
1013 }
1014 continue;
1015 }
1016
1017 let reservation_token =
1018 reservation_token.expect("non-waiting activation owns a reservation token");
1019 return self
1020 .dispatch_reserved_with_recovery(
1021 target_session_id,
1022 inbox_generation,
1023 reservation_token,
1024 false,
1025 )
1026 .await;
1027 }
1028 }
1029}
1030
1031#[cfg(test)]
1032mod tests {
1033 use super::*;
1034 use std::sync::atomic::{AtomicUsize, Ordering};
1035 use tokio::sync::{Barrier, Notify};
1036
1037 struct RecordingSpawner {
1038 reservations: AtomicUsize,
1039 launches: Arc<AtomicUsize>,
1040 entered: Option<Arc<Barrier>>,
1041 release: Option<Arc<Notify>>,
1042 }
1043
1044 #[async_trait]
1045 impl SessionActivationSpawner for RecordingSpawner {
1046 async fn reserve_activation(
1047 &self,
1048 _target_session_id: &str,
1049 _inbox_generation: u64,
1050 ) -> Result<SessionActivationReserveOutcome, SessionActivationError> {
1051 let ordinal = self.reservations.fetch_add(1, Ordering::SeqCst) + 1;
1052 if let Some(entered) = &self.entered {
1053 entered.wait().await;
1054 }
1055 if let Some(release) = &self.release {
1056 release.notified().await;
1057 }
1058 let launches = self.launches.clone();
1059 Ok(SessionActivationReserveOutcome::Reserved(
1060 SessionActivationLaunch::new(format!("run-{ordinal}"), move || {
1061 launches.fetch_add(1, Ordering::SeqCst);
1062 }),
1063 ))
1064 }
1065 }
1066
1067 fn spawner() -> Arc<RecordingSpawner> {
1068 Arc::new(RecordingSpawner {
1069 reservations: AtomicUsize::new(0),
1070 launches: Arc::new(AtomicUsize::new(0)),
1071 entered: None,
1072 release: None,
1073 })
1074 }
1075
1076 struct BlockingInspectInbox {
1077 entered: Arc<Notify>,
1078 release: Arc<Notify>,
1079 }
1080
1081 #[async_trait]
1082 impl SessionInboxPort for BlockingInspectInbox {
1083 async fn deliver(
1084 &self,
1085 _envelope: &bamboo_domain::SessionMessageEnvelope,
1086 ) -> Result<bamboo_domain::SessionInboxReceipt, bamboo_domain::SessionInboxError> {
1087 unreachable!("router cleanup only inspects the durable inbox")
1088 }
1089
1090 async fn mark_activation_eligible(
1091 &self,
1092 _target_session_id: &str,
1093 _generation: u64,
1094 _policy: bamboo_domain::SessionActivationPolicy,
1095 ) -> Result<(), bamboo_domain::SessionInboxError> {
1096 unreachable!("router cleanup only inspects the durable inbox")
1097 }
1098
1099 async fn claim(
1100 &self,
1101 _target_session_id: &str,
1102 _limit: usize,
1103 ) -> Result<Vec<bamboo_domain::SessionInboxClaim>, bamboo_domain::SessionInboxError>
1104 {
1105 unreachable!("router cleanup only inspects the durable inbox")
1106 }
1107
1108 async fn was_admitted(
1109 &self,
1110 _target_session_id: &str,
1111 _id: &bamboo_domain::SessionMessageId,
1112 ) -> Result<bool, bamboo_domain::SessionInboxError> {
1113 unreachable!("router cleanup only inspects the durable inbox")
1114 }
1115
1116 async fn ack(
1117 &self,
1118 _target_session_id: &str,
1119 _claim: &bamboo_domain::SessionInboxClaim,
1120 ) -> Result<(), bamboo_domain::SessionInboxError> {
1121 unreachable!("router cleanup only inspects the durable inbox")
1122 }
1123
1124 async fn inspect(
1125 &self,
1126 _target_session_id: &str,
1127 ) -> Result<bamboo_domain::SessionInboxBacklog, bamboo_domain::SessionInboxError> {
1128 self.entered.notify_one();
1129 self.release.notified().await;
1130 Ok(bamboo_domain::SessionInboxBacklog {
1131 pending: 1,
1132 claimed: 0,
1133 generation: 1,
1134 activation_generation: 1,
1135 interrupt_generation: 1,
1136 oldest_generation: Some(1),
1137 })
1138 }
1139 }
1140
1141 #[tokio::test]
1142 async fn rollback_completion_is_retained_before_waiter_subscribes() {
1143 let completion = RollbackCompletion::new();
1144 completion.complete();
1145 tokio::time::timeout(std::time::Duration::from_millis(100), completion.wait())
1146 .await
1147 .expect("completion-before-wait must not lose its wake");
1148 }
1149
1150 #[derive(Clone, Copy)]
1151 enum InjectedFirstOutcome {
1152 NoWork,
1153 Error,
1154 }
1155
1156 struct RetrySpawner {
1157 reservations: AtomicUsize,
1158 launches: Arc<AtomicUsize>,
1159 first_outcome: InjectedFirstOutcome,
1160 first_entered: Arc<Barrier>,
1161 release_first: Arc<Notify>,
1162 }
1163
1164 struct CancelFirstSpawner {
1165 reservations: AtomicUsize,
1166 launches: Arc<AtomicUsize>,
1167 first_entered: Arc<Notify>,
1168 }
1169
1170 struct PanicFirstSpawner {
1171 reservations: AtomicUsize,
1172 launches: Arc<AtomicUsize>,
1173 }
1174
1175 struct AlreadyRunningThenReserveSpawner {
1176 reservations: AtomicUsize,
1177 launches: Arc<AtomicUsize>,
1178 }
1179
1180 #[async_trait]
1181 impl SessionActivationSpawner for AlreadyRunningThenReserveSpawner {
1182 async fn reserve_activation(
1183 &self,
1184 _target_session_id: &str,
1185 _inbox_generation: u64,
1186 ) -> Result<SessionActivationReserveOutcome, SessionActivationError> {
1187 let ordinal = self.reservations.fetch_add(1, Ordering::SeqCst) + 1;
1188 if ordinal == 1 {
1189 return Ok(SessionActivationReserveOutcome::AlreadyRunning {
1190 run_id: "adopted-run".to_string(),
1191 });
1192 }
1193 let launches = self.launches.clone();
1194 Ok(SessionActivationReserveOutcome::Reserved(
1195 SessionActivationLaunch::new("successor-run", move || {
1196 launches.fetch_add(1, Ordering::SeqCst);
1197 }),
1198 ))
1199 }
1200 }
1201
1202 #[async_trait]
1203 impl SessionActivationSpawner for CancelFirstSpawner {
1204 async fn reserve_activation(
1205 &self,
1206 _target_session_id: &str,
1207 _inbox_generation: u64,
1208 ) -> Result<SessionActivationReserveOutcome, SessionActivationError> {
1209 let ordinal = self.reservations.fetch_add(1, Ordering::SeqCst) + 1;
1210 if ordinal == 1 {
1211 self.first_entered.notify_one();
1212 std::future::pending::<()>().await;
1213 }
1214 let launches = self.launches.clone();
1215 Ok(SessionActivationReserveOutcome::Reserved(
1216 SessionActivationLaunch::new(format!("run-{ordinal}"), move || {
1217 launches.fetch_add(1, Ordering::SeqCst);
1218 }),
1219 ))
1220 }
1221 }
1222
1223 #[async_trait]
1224 impl SessionActivationSpawner for PanicFirstSpawner {
1225 async fn reserve_activation(
1226 &self,
1227 _target_session_id: &str,
1228 _inbox_generation: u64,
1229 ) -> Result<SessionActivationReserveOutcome, SessionActivationError> {
1230 let ordinal = self.reservations.fetch_add(1, Ordering::SeqCst) + 1;
1231 if ordinal == 1 {
1232 panic!("injected first activation-spawner panic");
1233 }
1234 let launches = self.launches.clone();
1235 Ok(SessionActivationReserveOutcome::Reserved(
1236 SessionActivationLaunch::new(format!("run-{ordinal}"), move || {
1237 launches.fetch_add(1, Ordering::SeqCst);
1238 }),
1239 ))
1240 }
1241 }
1242
1243 struct RollbackSpawner {
1244 entered: Arc<Barrier>,
1245 release: Arc<Notify>,
1246 launch_ready: Arc<Notify>,
1247 rollbacks: Arc<AtomicUsize>,
1248 }
1249
1250 struct RealRegistryRollbackSpawner {
1251 runners: Arc<
1252 tokio::sync::RwLock<std::collections::HashMap<String, crate::execution::AgentRunner>>,
1253 >,
1254 senders: Arc<
1255 tokio::sync::RwLock<
1256 std::collections::HashMap<
1257 String,
1258 tokio::sync::broadcast::Sender<bamboo_agent_core::AgentEvent>,
1259 >,
1260 >,
1261 >,
1262 sender: tokio::sync::broadcast::Sender<bamboo_agent_core::AgentEvent>,
1263 reservations: AtomicUsize,
1264 first_reserved: Arc<tokio::sync::Notify>,
1265 allow_first_return: Arc<tokio::sync::Notify>,
1266 first_returning: Arc<tokio::sync::Notify>,
1267 rollback_started: Arc<tokio::sync::Notify>,
1268 allow_rollback: Arc<tokio::sync::Notify>,
1269 launches: Arc<AtomicUsize>,
1270 }
1271
1272 #[async_trait]
1273 impl SessionActivationSpawner for RealRegistryRollbackSpawner {
1274 async fn reserve_activation(
1275 &self,
1276 target_session_id: &str,
1277 _inbox_generation: u64,
1278 ) -> Result<SessionActivationReserveOutcome, SessionActivationError> {
1279 let ordinal = self.reservations.fetch_add(1, Ordering::SeqCst) + 1;
1280 match crate::execution::reserve_runner_core(
1281 &self.runners,
1282 &self.senders,
1283 target_session_id,
1284 &self.sender,
1285 )
1286 .await
1287 {
1288 crate::execution::ReserveOutcome::AlreadyRunning(run_id) => {
1289 Ok(SessionActivationReserveOutcome::AlreadyRunning { run_id })
1290 }
1291 crate::execution::ReserveOutcome::Reserved(reservation) => {
1292 let run_id = reservation.run_id.clone();
1293 if ordinal == 1 {
1294 self.first_reserved.notify_one();
1295 self.allow_first_return.notified().await;
1296 self.first_returning.notify_one();
1297 }
1298 let rollback_runners = self.runners.clone();
1299 let rollback_session_id = target_session_id.to_string();
1300 let rollback_run_id = run_id.clone();
1301 let rollback_started = self.rollback_started.clone();
1302 let allow_rollback = self.allow_rollback.clone();
1303 let launches = self.launches.clone();
1304 Ok(SessionActivationReserveOutcome::Reserved(
1305 SessionActivationLaunch::new_with_async_rollback(
1306 run_id,
1307 move || {
1308 launches.fetch_add(1, Ordering::SeqCst);
1309 },
1310 move || async move {
1311 rollback_started.notify_one();
1312 allow_rollback.notified().await;
1313 let mut runners = rollback_runners.write().await;
1314 if runners
1315 .get(&rollback_session_id)
1316 .is_some_and(|runner| runner.run_id == rollback_run_id)
1317 {
1318 runners.remove(&rollback_session_id);
1319 }
1320 },
1321 ),
1322 ))
1323 }
1324 }
1325 }
1326 }
1327
1328 #[async_trait]
1329 impl SessionActivationSpawner for RollbackSpawner {
1330 async fn reserve_activation(
1331 &self,
1332 _target_session_id: &str,
1333 _inbox_generation: u64,
1334 ) -> Result<SessionActivationReserveOutcome, SessionActivationError> {
1335 self.entered.wait().await;
1336 self.release.notified().await;
1337 let rollbacks = self.rollbacks.clone();
1338 let outcome = SessionActivationReserveOutcome::Reserved(
1339 SessionActivationLaunch::new_with_rollback(
1340 "reserved-run",
1341 || {},
1342 move || {
1343 rollbacks.fetch_add(1, Ordering::SeqCst);
1344 },
1345 ),
1346 );
1347 self.launch_ready.notify_one();
1348 Ok(outcome)
1349 }
1350 }
1351
1352 #[async_trait]
1353 impl SessionActivationSpawner for RetrySpawner {
1354 async fn reserve_activation(
1355 &self,
1356 _target_session_id: &str,
1357 _inbox_generation: u64,
1358 ) -> Result<SessionActivationReserveOutcome, SessionActivationError> {
1359 let ordinal = self.reservations.fetch_add(1, Ordering::SeqCst) + 1;
1360 if ordinal == 1 {
1361 self.first_entered.wait().await;
1362 self.release_first.notified().await;
1363 return match self.first_outcome {
1364 InjectedFirstOutcome::NoWork => Ok(SessionActivationReserveOutcome::NoWork),
1365 InjectedFirstOutcome::Error => Err(SessionActivationError::Internal(
1366 "injected first reservation failure".to_string(),
1367 )),
1368 };
1369 }
1370 let launches = self.launches.clone();
1371 Ok(SessionActivationReserveOutcome::Reserved(
1372 SessionActivationLaunch::new(format!("run-{ordinal}"), move || {
1373 launches.fetch_add(1, Ordering::SeqCst);
1374 }),
1375 ))
1376 }
1377 }
1378
1379 async fn assert_newer_generation_retries_after(
1380 first_outcome: InjectedFirstOutcome,
1381 ) -> Result<SessionActivationDisposition, SessionActivationError> {
1382 let router = SessionActivationRouter::new();
1383 let first_entered = Arc::new(Barrier::new(2));
1384 let release_first = Arc::new(Notify::new());
1385 let spawner = Arc::new(RetrySpawner {
1386 reservations: AtomicUsize::new(0),
1387 launches: Arc::new(AtomicUsize::new(0)),
1388 first_outcome,
1389 first_entered: first_entered.clone(),
1390 release_first: release_first.clone(),
1391 });
1392 router.set_spawner(spawner.clone()).await;
1393
1394 let first_router = router.clone();
1395 let first =
1396 tokio::spawn(async move { first_router.request_activation("session", 1).await });
1397 first_entered.wait().await;
1398
1399 let second_router = router.clone();
1400 let second =
1401 tokio::spawn(async move { second_router.request_activation("session", 2).await });
1402 tokio::task::yield_now().await;
1403 assert_eq!(
1404 spawner.reservations.load(Ordering::SeqCst),
1405 1,
1406 "generation 2 must coalesce while generation 1 owns the reservation"
1407 );
1408
1409 release_first.notify_one();
1410 let first_result = first.await.unwrap();
1411 assert_eq!(
1412 second.await.unwrap().unwrap(),
1413 SessionActivationDisposition::ActivationReserved,
1414 "the coalesced newer delivery must take the released reservation"
1415 );
1416 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 2);
1417 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
1418 first_result
1419 }
1420
1421 #[tokio::test]
1422 async fn active_owner_is_notified_without_second_reservation() {
1423 let router = SessionActivationRouter::new();
1424 let spawner = spawner();
1425 router.set_spawner(spawner.clone()).await;
1426 let mut registration = router.register_run("session", "run-live").await.unwrap();
1427
1428 assert_eq!(
1429 router.request_activation("session", 7).await.unwrap(),
1430 SessionActivationDisposition::ActiveNotified
1431 );
1432 registration.notifications_mut().changed().await.unwrap();
1433 assert_eq!(*registration.notifications_mut().borrow(), 7);
1434 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 0);
1435 registration.begin_finalization().await;
1436 assert_eq!(registration.finish(7).await.unwrap(), None);
1437 }
1438
1439 #[tokio::test]
1440 async fn different_run_registration_never_overwrites_live_owner() {
1441 let router = SessionActivationRouter::new();
1442 let mut first = router
1443 .register_run("session", "run-1")
1444 .await
1445 .expect("first run owns the logical session");
1446
1447 let error = router
1448 .register_run("session", "run-2")
1449 .await
1450 .err()
1451 .expect("a second exact run must collide");
1452 assert_eq!(
1453 error,
1454 SessionRunRegistrationError::OwnerCollision {
1455 target_session_id: "session".to_string(),
1456 existing_run_id: "run-1".to_string(),
1457 attempted_run_id: "run-2".to_string(),
1458 }
1459 );
1460 assert!(router.owns_run("session", "run-1").await);
1461 assert!(!router.owns_run("session", "run-2").await);
1462
1463 first.begin_finalization().await;
1464 assert_eq!(first.finish(0).await.unwrap(), None);
1465 }
1466
1467 #[tokio::test]
1468 async fn direct_registration_waits_for_activation_to_publish_its_exact_owner() {
1469 let router = SessionActivationRouter::new();
1470 let reservation_entered = Arc::new(Barrier::new(2));
1471 let allow_reservation = Arc::new(Notify::new());
1472 let launches = Arc::new(AtomicUsize::new(0));
1473 let spawner = Arc::new(RecordingSpawner {
1474 reservations: AtomicUsize::new(0),
1475 launches: launches.clone(),
1476 entered: Some(reservation_entered.clone()),
1477 release: Some(allow_reservation.clone()),
1478 });
1479 router.set_spawner(spawner.clone()).await;
1480
1481 let activation_router = router.clone();
1482 let activation =
1483 tokio::spawn(async move { activation_router.request_activation("session", 1).await });
1484 reservation_entered.wait().await;
1485
1486 let direct_router = router.clone();
1487 let direct =
1488 tokio::spawn(async move { direct_router.register_run("session", "manual-run").await });
1489 tokio::task::yield_now().await;
1490 assert!(
1491 !direct.is_finished(),
1492 "manual registration must not supersede an in-flight activation token"
1493 );
1494 assert!(!router.owns_run("session", "manual-run").await);
1495
1496 allow_reservation.notify_one();
1497 assert_eq!(
1498 activation.await.unwrap().unwrap(),
1499 SessionActivationDisposition::ActivationReserved
1500 );
1501 let collision = match direct.await.unwrap() {
1502 Ok(_) => panic!("manual registration must collide with the published activation"),
1503 Err(error) => error,
1504 };
1505 assert_eq!(
1506 collision,
1507 SessionRunRegistrationError::OwnerCollision {
1508 target_session_id: "session".to_string(),
1509 existing_run_id: "run-1".to_string(),
1510 attempted_run_id: "manual-run".to_string(),
1511 }
1512 );
1513 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 1);
1514 assert_eq!(launches.load(Ordering::SeqCst), 1);
1515 assert!(router.owns_run("session", "run-1").await);
1516 assert!(!router.owns_run("session", "manual-run").await);
1517
1518 let mut registration = router.register_run("session", "run-1").await.unwrap();
1519 registration.begin_finalization().await;
1520 assert_eq!(registration.finish(1).await.unwrap(), None);
1521 }
1522
1523 #[tokio::test]
1524 async fn cancelled_explicit_abandon_still_finishes_owned_registration_cleanup() {
1525 let router = SessionActivationRouter::new();
1526 let cleanup_entered = Arc::new(Notify::new());
1527 let cleanup_release = Arc::new(Notify::new());
1528 let cleanup_completed = Arc::new(AtomicUsize::new(0));
1529 let mut registration = router.register_run("session", "run-abandon").await.unwrap();
1530 let entered = cleanup_entered.clone();
1531 let release = cleanup_release.clone();
1532 let completed = cleanup_completed.clone();
1533 registration.set_abort_cleanup(move || async move {
1534 entered.notify_one();
1535 release.notified().await;
1536 completed.fetch_add(1, Ordering::SeqCst);
1537 });
1538
1539 let abandon = tokio::spawn(registration.abandon());
1540 tokio::time::timeout(
1541 std::time::Duration::from_secs(1),
1542 cleanup_entered.notified(),
1543 )
1544 .await
1545 .expect("explicit abandon must start its owned cleanup");
1546 abandon.abort();
1547 assert!(abandon.await.unwrap_err().is_cancelled());
1548 cleanup_release.notify_one();
1549
1550 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1551 loop {
1552 if cleanup_completed.load(Ordering::SeqCst) == 1
1553 && !router.owns_run("session", "run-abandon").await
1554 {
1555 break;
1556 }
1557 tokio::task::yield_now().await;
1558 }
1559 })
1560 .await
1561 .expect("caller cancellation must not cancel detached registration cleanup");
1562 }
1563
1564 #[tokio::test]
1565 async fn delayed_same_run_cannot_adopt_an_owner_during_abandoned_cleanup() {
1566 let router = SessionActivationRouter::new();
1567 let spawner = spawner();
1568 router.set_spawner(spawner.clone()).await;
1569 let inspect_entered = Arc::new(Notify::new());
1570 let inspect_release = Arc::new(Notify::new());
1571 router.set_inbox(Arc::new(BlockingInspectInbox {
1572 entered: inspect_entered.clone(),
1573 release: inspect_release.clone(),
1574 }));
1575
1576 let registration = router
1577 .register_run("session", "run-abandoned")
1578 .await
1579 .unwrap();
1580 drop(registration);
1581 tokio::time::timeout(
1582 std::time::Duration::from_secs(1),
1583 inspect_entered.notified(),
1584 )
1585 .await
1586 .expect("abandoned cleanup must reach durable inbox reconciliation");
1587
1588 let error = router
1589 .register_run("session", "run-abandoned")
1590 .await
1591 .err()
1592 .expect("a finalizing abandoned owner is not an adoptable launch placeholder");
1593 assert_eq!(error.existing_run_id(), "run-abandoned");
1594
1595 inspect_release.notify_one();
1596 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1597 while spawner.launches.load(Ordering::SeqCst) == 0 {
1598 tokio::task::yield_now().await;
1599 }
1600 })
1601 .await
1602 .expect("durable pending work must launch one successor");
1603 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 1);
1604 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
1605 assert!(!router.owns_run("session", "run-abandoned").await);
1606
1607 let mut successor = router.register_run("session", "run-1").await.unwrap();
1608 successor.begin_finalization().await;
1609 assert_eq!(successor.finish(1).await.unwrap(), None);
1610 }
1611
1612 #[tokio::test]
1613 async fn delayed_registry_owner_collision_rolls_back_only_its_exact_runner() {
1614 let router = SessionActivationRouter::new();
1615 let runners = Arc::new(tokio::sync::RwLock::new(std::collections::HashMap::new()));
1616 let senders = Arc::new(tokio::sync::RwLock::new(std::collections::HashMap::new()));
1617 let (sender, _receiver) = tokio::sync::broadcast::channel(8);
1618 let reserved =
1619 match crate::execution::reserve_runner_core(&runners, &senders, "session", &sender)
1620 .await
1621 {
1622 crate::execution::ReserveOutcome::Reserved(reservation) => reservation,
1623 crate::execution::ReserveOutcome::AlreadyRunning(_) => {
1624 panic!("fixture must reserve a fresh server runner")
1625 }
1626 };
1627
1628 let mut direct = router.register_run("session", "sdk-direct").await.unwrap();
1632 let collision = router
1633 .register_run("session", &reserved.run_id)
1634 .await
1635 .err()
1636 .expect("delayed server registration must collide");
1637 let collision_result = Err(bamboo_agent_core::AgentError::LLM(collision.to_string()));
1638 assert!(
1639 crate::execution::finalize_runner_exact(
1640 &runners,
1641 "session",
1642 &reserved.run_id,
1643 &collision_result,
1644 )
1645 .await
1646 );
1647 assert!(router.owns_run("session", "sdk-direct").await);
1648 assert!(matches!(
1649 runners
1650 .read()
1651 .await
1652 .get("session")
1653 .map(|runner| &runner.status),
1654 Some(crate::execution::AgentStatus::Error(_))
1655 ));
1656
1657 direct.begin_finalization().await;
1658 assert_eq!(direct.finish(0).await.unwrap(), None);
1659 }
1660
1661 #[tokio::test]
1662 async fn concurrent_idle_deliveries_coalesce_into_one_launch() {
1663 let router = SessionActivationRouter::new();
1664 let entered = Arc::new(Barrier::new(2));
1665 let release = Arc::new(Notify::new());
1666 let spawner = Arc::new(RecordingSpawner {
1667 reservations: AtomicUsize::new(0),
1668 launches: Arc::new(AtomicUsize::new(0)),
1669 entered: Some(entered.clone()),
1670 release: Some(release.clone()),
1671 });
1672 router.set_spawner(spawner.clone()).await;
1673
1674 let first_router = router.clone();
1675 let first =
1676 tokio::spawn(async move { first_router.request_activation("session", 1).await });
1677 entered.wait().await;
1678 let second_router = router.clone();
1679 let second =
1680 tokio::spawn(async move { second_router.request_activation("session", 2).await });
1681 release.notify_one();
1682 assert_eq!(
1683 first.await.unwrap().unwrap(),
1684 SessionActivationDisposition::ActivationReserved
1685 );
1686 assert_eq!(
1687 second.await.unwrap().unwrap(),
1688 SessionActivationDisposition::ActiveNotified
1689 );
1690 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 1);
1691 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
1692 }
1693
1694 #[tokio::test]
1695 async fn newer_generation_re_reserves_after_coalesced_no_work() {
1696 assert_eq!(
1697 assert_newer_generation_retries_after(InjectedFirstOutcome::NoWork)
1698 .await
1699 .unwrap(),
1700 SessionActivationDisposition::ActivationCoalesced
1701 );
1702 }
1703
1704 #[tokio::test]
1705 async fn newer_generation_re_reserves_after_coalesced_error() {
1706 let error = assert_newer_generation_retries_after(InjectedFirstOutcome::Error)
1707 .await
1708 .unwrap_err();
1709 assert!(matches!(error, SessionActivationError::Internal(_)));
1710 }
1711
1712 async fn assert_cancelled_finalization_recovers_without_redelivery(registered_owner: bool) {
1713 let router = SessionActivationRouter::new();
1714 let first_entered = Arc::new(Notify::new());
1715 let spawner = Arc::new(CancelFirstSpawner {
1716 reservations: AtomicUsize::new(0),
1717 launches: Arc::new(AtomicUsize::new(0)),
1718 first_entered: first_entered.clone(),
1719 });
1720 router.set_spawner(spawner.clone()).await;
1721
1722 let finish = if registered_owner {
1723 let mut registration = router.register_run("session", "run-old").await.unwrap();
1724 registration.begin_finalization().await;
1725 assert_eq!(
1726 router.request_activation("session", 1).await.unwrap(),
1727 SessionActivationDisposition::ActivationCoalesced
1728 );
1729 tokio::spawn(async move { registration.finish(0).await })
1730 } else {
1731 router.begin_finalization("session", "run-old").await;
1735 assert_eq!(
1736 router.request_activation("session", 1).await.unwrap(),
1737 SessionActivationDisposition::ActivationCoalesced
1738 );
1739 let finish_router = router.clone();
1740 tokio::spawn(async move {
1741 finish_router
1742 .finish_finalization("session", "run-old", 0)
1743 .await
1744 })
1745 };
1746 first_entered.notified().await;
1747 finish.abort();
1748 assert!(finish.await.unwrap_err().is_cancelled());
1749
1750 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1751 while spawner.launches.load(Ordering::SeqCst) == 0 {
1752 tokio::task::yield_now().await;
1753 }
1754 })
1755 .await
1756 .expect("cancelled finalization must retry pending work without redelivery");
1757 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 2);
1758 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
1759 assert!(router.owns_run("session", "run-2").await);
1760
1761 let mut successor = router.register_run("session", "run-2").await.unwrap();
1762 successor.begin_finalization().await;
1763 assert_eq!(successor.finish(1).await.unwrap(), None);
1764 }
1765
1766 #[tokio::test]
1767 async fn cancelled_registered_finalization_retries_without_another_delivery() {
1768 assert_cancelled_finalization_recovers_without_redelivery(true).await;
1769 }
1770
1771 #[tokio::test]
1772 async fn cancelled_reserved_child_finalization_retries_without_another_delivery() {
1773 assert_cancelled_finalization_recovers_without_redelivery(false).await;
1774 }
1775
1776 #[tokio::test]
1777 async fn panicking_finalization_spawner_gets_one_bounded_retry() {
1778 let router = SessionActivationRouter::new();
1779 let spawner = Arc::new(PanicFirstSpawner {
1780 reservations: AtomicUsize::new(0),
1781 launches: Arc::new(AtomicUsize::new(0)),
1782 });
1783 router.set_spawner(spawner.clone()).await;
1784 router.begin_finalization("session", "run-old").await;
1785 assert_eq!(
1786 router.request_activation("session", 1).await.unwrap(),
1787 SessionActivationDisposition::ActivationCoalesced
1788 );
1789
1790 let finish_router = router.clone();
1791 let finish = tokio::spawn(async move {
1792 finish_router
1793 .finish_finalization("session", "run-old", 0)
1794 .await
1795 });
1796 assert!(finish.await.unwrap_err().is_panic());
1797
1798 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1799 while spawner.launches.load(Ordering::SeqCst) == 0 {
1800 tokio::task::yield_now().await;
1801 }
1802 })
1803 .await
1804 .expect("router Drop recovery must survive one spawner panic");
1805 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 2);
1806 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
1807
1808 let mut successor = router.register_run("session", "run-2").await.unwrap();
1809 successor.begin_finalization().await;
1810 assert_eq!(successor.finish(1).await.unwrap(), None);
1811 }
1812
1813 #[tokio::test]
1814 async fn cancelled_reservation_releases_token_and_newer_generation_launches() {
1815 let router = SessionActivationRouter::new();
1816 let first_entered = Arc::new(Notify::new());
1817 let spawner = Arc::new(CancelFirstSpawner {
1818 reservations: AtomicUsize::new(0),
1819 launches: Arc::new(AtomicUsize::new(0)),
1820 first_entered: first_entered.clone(),
1821 });
1822 router.set_spawner(spawner.clone()).await;
1823
1824 let first_router = router.clone();
1825 let first =
1826 tokio::spawn(async move { first_router.request_activation("session", 1).await });
1827 first_entered.notified().await;
1828 let second_router = router.clone();
1829 let second =
1830 tokio::spawn(async move { second_router.request_activation("session", 2).await });
1831 tokio::task::yield_now().await;
1832 first.abort();
1833 assert!(first.await.unwrap_err().is_cancelled());
1834
1835 assert_eq!(
1836 tokio::time::timeout(std::time::Duration::from_secs(1), second)
1837 .await
1838 .expect("coalesced caller must not hang behind a cancelled reservation")
1839 .unwrap()
1840 .unwrap(),
1841 SessionActivationDisposition::ActivationReserved
1842 );
1843 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 2);
1844 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
1845 }
1846
1847 #[tokio::test]
1848 async fn cancellation_after_external_reservation_rolls_back_unlaunched_slot() {
1849 let router = SessionActivationRouter::new();
1850 let entered = Arc::new(Barrier::new(2));
1851 let release = Arc::new(Notify::new());
1852 let launch_ready = Arc::new(Notify::new());
1853 let rollbacks = Arc::new(AtomicUsize::new(0));
1854 let spawner = Arc::new(RollbackSpawner {
1855 entered: entered.clone(),
1856 release: release.clone(),
1857 launch_ready: launch_ready.clone(),
1858 rollbacks: rollbacks.clone(),
1859 });
1860 router.set_spawner(spawner).await;
1861
1862 let task_router = router.clone();
1863 let task = tokio::spawn(async move { task_router.request_activation("session", 1).await });
1864 entered.wait().await;
1865 let states = router.states.lock().await;
1869 release.notify_one();
1870 launch_ready.notified().await;
1871 task.abort();
1872 assert!(task.await.unwrap_err().is_cancelled());
1873 drop(states);
1874
1875 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1876 while rollbacks.load(Ordering::SeqCst) == 0 {
1877 tokio::task::yield_now().await;
1878 }
1879 })
1880 .await
1881 .expect("unlaunched reservation rollback");
1882 assert_eq!(rollbacks.load(Ordering::SeqCst), 1);
1883 }
1884
1885 #[tokio::test]
1886 async fn redelivery_waits_for_real_runner_rollback_before_reserving() {
1887 let router = SessionActivationRouter::new();
1888 let runners = Arc::new(tokio::sync::RwLock::new(std::collections::HashMap::new()));
1889 let senders = Arc::new(tokio::sync::RwLock::new(std::collections::HashMap::new()));
1890 let (sender, _receiver) = tokio::sync::broadcast::channel(8);
1891 let first_reserved = Arc::new(tokio::sync::Notify::new());
1892 let allow_first_return = Arc::new(tokio::sync::Notify::new());
1893 let rollback_started = Arc::new(tokio::sync::Notify::new());
1894 let first_returning = Arc::new(tokio::sync::Notify::new());
1895 let allow_rollback = Arc::new(tokio::sync::Notify::new());
1896 let launches = Arc::new(AtomicUsize::new(0));
1897 let spawner = Arc::new(RealRegistryRollbackSpawner {
1898 runners: runners.clone(),
1899 senders,
1900 sender,
1901 reservations: AtomicUsize::new(0),
1902 first_reserved: first_reserved.clone(),
1903 allow_first_return: allow_first_return.clone(),
1904 first_returning: first_returning.clone(),
1905 rollback_started: rollback_started.clone(),
1906 allow_rollback: allow_rollback.clone(),
1907 launches: launches.clone(),
1908 });
1909 router.set_spawner(spawner.clone()).await;
1910
1911 let first_router = router.clone();
1912 let first =
1913 tokio::spawn(async move { first_router.request_activation("session", 1).await });
1914 first_reserved.notified().await;
1915
1916 let state_guard = router.states.lock().await;
1919 allow_first_return.notify_one();
1920 first_returning.notified().await;
1921 first.abort();
1922 assert!(first.await.unwrap_err().is_cancelled());
1923 drop(state_guard);
1924 rollback_started.notified().await;
1925
1926 let second_router = router.clone();
1927 let second =
1928 tokio::spawn(async move { second_router.request_activation("session", 2).await });
1929 tokio::task::yield_now().await;
1930 assert_eq!(
1931 spawner.reservations.load(Ordering::SeqCst),
1932 1,
1933 "redelivery must remain coalesced until exact external rollback completes"
1934 );
1935
1936 allow_rollback.notify_one();
1937 assert_eq!(
1938 tokio::time::timeout(std::time::Duration::from_secs(1), second)
1939 .await
1940 .expect("redelivery must reserve after rollback")
1941 .unwrap()
1942 .unwrap(),
1943 SessionActivationDisposition::ActivationReserved
1944 );
1945 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 2);
1946 assert_eq!(launches.load(Ordering::SeqCst), 1);
1947 let current = runners
1948 .read()
1949 .await
1950 .get("session")
1951 .map(|runner| runner.run_id.clone())
1952 .expect("second real registry reservation remains");
1953 assert_ne!(
1954 current, "",
1955 "the redelivery must never adopt an empty/stale run id"
1956 );
1957 }
1958
1959 #[tokio::test]
1960 async fn delivery_during_finalization_launches_one_successor() {
1961 let router = SessionActivationRouter::new();
1962 let spawner = spawner();
1963 router.set_spawner(spawner.clone()).await;
1964 let mut registration = router.register_run("session", "run-old").await.unwrap();
1965 registration.begin_finalization().await;
1966
1967 assert_eq!(
1968 router.request_activation("session", 11).await.unwrap(),
1969 SessionActivationDisposition::ActivationCoalesced
1970 );
1971 let disposition = registration.finish(10).await.unwrap();
1972 assert_eq!(
1973 disposition,
1974 Some(SessionActivationDisposition::ActivationReserved)
1975 );
1976 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 1);
1977 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
1978 }
1979
1980 #[tokio::test]
1981 async fn finalization_does_not_restart_when_cursor_caught_up() {
1982 let router = SessionActivationRouter::new();
1983 let spawner = spawner();
1984 router.set_spawner(spawner.clone()).await;
1985 let mut registration = router.register_run("session", "run-old").await.unwrap();
1986 assert_eq!(
1987 router.request_activation("session", 4).await.unwrap(),
1988 SessionActivationDisposition::ActiveNotified
1989 );
1990 registration.begin_finalization().await;
1991 assert_eq!(registration.finish(4).await.unwrap(), None);
1992 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 0);
1993 assert!(
1994 router.states.lock().await.get("session").is_none(),
1995 "a caught-up completed session must not leak routing state"
1996 );
1997 }
1998
1999 #[tokio::test]
2000 async fn successor_finalization_removes_state_after_raced_generation_catches_up() {
2001 let router = SessionActivationRouter::new();
2002 let spawner = spawner();
2003 router.set_spawner(spawner).await;
2004 let mut registration = router.register_run("session", "run-old").await.unwrap();
2005 registration.begin_finalization().await;
2006
2007 assert_eq!(
2008 router.request_activation("session", 11).await.unwrap(),
2009 SessionActivationDisposition::ActivationCoalesced
2010 );
2011 assert_eq!(
2012 registration.finish(10).await.unwrap(),
2013 Some(SessionActivationDisposition::ActivationReserved)
2014 );
2015 let mut successor_registration = router.register_run("session", "run-1").await.unwrap();
2016 successor_registration.begin_finalization().await;
2017 assert_eq!(successor_registration.finish(11).await.unwrap(), None);
2018 assert!(
2019 router.states.lock().await.get("session").is_none(),
2020 "the successor's caught-up terminal state must be compacted"
2021 );
2022 }
2023
2024 #[tokio::test]
2025 async fn poison_generation_launches_only_one_successor_until_new_work_arrives() {
2026 let router = SessionActivationRouter::new();
2027 let spawner = spawner();
2028 router.set_spawner(spawner.clone()).await;
2029 let mut registration = router
2030 .register_run("session", "run-original")
2031 .await
2032 .unwrap();
2033 assert_eq!(
2034 router.request_activation("session", 7).await.unwrap(),
2035 SessionActivationDisposition::ActiveNotified
2036 );
2037 registration.begin_finalization().await;
2038 assert_eq!(
2039 registration.finish(0).await.unwrap(),
2040 Some(SessionActivationDisposition::ActivationReserved)
2041 );
2042 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 1);
2043
2044 let mut successor_registration = router.register_run("session", "run-1").await.unwrap();
2048 successor_registration.begin_finalization().await;
2049 assert_eq!(successor_registration.finish(0).await.unwrap(), None);
2050 assert_eq!(
2051 router.request_activation("session", 7).await.unwrap(),
2052 SessionActivationDisposition::ActivationCoalesced
2053 );
2054 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 1);
2055
2056 assert_eq!(
2058 router.request_activation("session", 8).await.unwrap(),
2059 SessionActivationDisposition::ActivationReserved
2060 );
2061 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 2);
2062 }
2063
2064 #[tokio::test]
2065 async fn adopted_terminal_run_preserves_generation_for_one_real_successor() {
2066 let router = SessionActivationRouter::new();
2067 let spawner = Arc::new(AlreadyRunningThenReserveSpawner {
2068 reservations: AtomicUsize::new(0),
2069 launches: Arc::new(AtomicUsize::new(0)),
2070 });
2071 router.set_spawner(spawner.clone()).await;
2072
2073 assert_eq!(
2074 router.request_activation("session", 12).await.unwrap(),
2075 SessionActivationDisposition::ActiveNotified
2076 );
2077 router.begin_finalization("session", "adopted-run").await;
2078 assert_eq!(
2079 router
2080 .finish_finalization("session", "adopted-run", 0)
2081 .await
2082 .unwrap(),
2083 Some(SessionActivationDisposition::ActivationReserved)
2084 );
2085 assert_eq!(spawner.reservations.load(Ordering::SeqCst), 2);
2086 assert_eq!(spawner.launches.load(Ordering::SeqCst), 1);
2087 }
2088}