Skip to main content

bamboo_engine/
session_activation.rs

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