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