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