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
19type AsyncRollback =
20    Box<dyn FnOnce() -> Pin<Box<dyn Future<Output = ()> + Send + 'static>> + Send + 'static>;
21
22/// A reservation whose runner slot already exists but whose task has not yet
23/// been launched. The router publishes the logical owner before calling
24/// `launch`, so even an immediately-completing task participates in the
25/// finalization handshake.
26pub 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    /// Build a launch backed by an already-reserved external runner slot.
68    ///
69    /// If the launch is dropped before `launch` commits, `rollback` must release
70    /// that exact reservation. This makes cancellation between the spawner
71    /// returning and the router publishing ownership recoverable.
72    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    /// Build a launch whose exact external runner reservation is released
83    /// asynchronously if publication is cancelled. The router does not release
84    /// its coalescing token until this future completes, preventing a newer
85    /// activation from adopting the still-present unlaunched slot.
86    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    /// Existing runner reservation succeeded. The returned launch must be
147    /// invoked exactly once.
148    Reserved(SessionActivationLaunch),
149    /// Some non-router path already owns the existing runner reservation.
150    AlreadyRunning { run_id: String },
151    /// The target disappeared before reservation.
152    NotFound,
153    /// The inbox was drained by another valid owner before reservation.
154    NoWork,
155}
156
157/// A logical session already belongs to a different exact activation run.
158///
159/// Registration is the last line of defence between independently scheduled
160/// entry points. A collision is returned to the later caller without mutating
161/// the existing owner.
162#[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    /// The exact run that already owns the logical session.
177    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/// Adapter to the owning runtime's existing runner reservation and spawn path.
187///
188/// Implementations must reserve through the same per-session runner registry as
189/// user/resume execution. They must not start a task before returning
190/// [`SessionActivationReserveOutcome::Reserved`].
191#[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    /// Present only while an external actor driver is alive. The payload is a
206    /// wake generation; that driver remains the single consumer of the
207    /// canonical FileSessionInbox and forwards the claimed typed envelope.
208    delivery_sink: Option<mpsc::UnboundedSender<u64>>,
209}
210
211#[derive(Debug)]
212struct TargetActivationState {
213    latest_generation: u64,
214    /// Highest generation for which this process already launched an execution.
215    /// Adopting an independently running owner does not consume this retry. If a
216    /// launched execution leaves the same poison claim pending,
217    /// terminal finalization must not hot-loop provider runs forever. A newer
218    /// generation may still trigger one fresh activation; process restart also
219    /// resets this bounded retry guard.
220    last_dispatched_generation: u64,
221    owner: Option<ActiveOwner>,
222    activation_reserved: bool,
223    /// Identity of the current reservation attempt. A cancellation cleanup or
224    /// late spawner result may only release the token it acquired.
225    activation_token: u64,
226    /// Incremented whenever an in-flight reservation attempt releases the
227    /// reservation. A delivery that coalesces behind that attempt waits for
228    /// this edge and re-evaluates ownership instead of returning a success that
229    /// could be stranded when the earlier attempt reports NoWork or fails.
230    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
266/// Cancellation lease for one router reservation attempt. Every await after
267/// publishing `activation_reserved` is covered; dropping the caller releases
268/// only its exact token and wakes coalesced deliveries. Finalization-owned
269/// attempts additionally schedule one bounded retry because their racing
270/// producer already returned a coalesced success and is no longer waiting.
271struct 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/// One router per owning runtime/AppState. All keys are logical Session ids.
330#[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
340/// Exact-run ownership lease returned by
341/// [`SessionActivationRouter::register_run`].
342///
343/// The lease deliberately owns the safe-point receiver. Dropping it before the
344/// normal finalization handshake schedules exact-run cleanup, reconciles the
345/// durable activation watermark, and gives the router one chance to launch a
346/// successor. A stale lease can never clear a newer owner.
347pub 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    /// Run host-specific exact-reservation cleanup before an abandoned owner
364    /// asks the router to reserve a successor.
365    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    /// Synchronously abandon this exact ownership lease.
380    ///
381    /// Most cancellation paths can rely on [`Drop`], which schedules the same
382    /// cleanup. Startup paths that must return a truthful retryable response
383    /// use this method so the runner slot and router owner are released before
384    /// the response becomes observable. Cleanup runs in a detached owned task:
385    /// cancelling the caller's wait cannot strand an already-taken host
386    /// cleanup closure or this router owner.
387    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    /// Complete the normal exact-owner handshake. The Drop fallback remains
414    /// armed across the await and is disabled only after finalization returns.
415    pub async fn finish(
416        mut self,
417        admitted_generation: u64,
418    ) -> Result<Option<SessionActivationDisposition>, SessionActivationError> {
419        // The router compacts caught-up routing state only after the last
420        // safe-point receiver is gone.
421        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    /// Late-bind the real server/SDK reservation adapter after its dependency
460    /// graph has been assembled.
461    pub async fn set_spawner(&self, spawner: Arc<dyn SessionActivationSpawner>) {
462        *self.spawner.write().await = Some(spawner);
463    }
464
465    /// Bind the durable inbox used by abandoned-run reconciliation.
466    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    /// Register an execution created by any entry point. The returned receiver
503    /// is a safe-point wake signal; the loop also drains the durable inbox at
504    /// every round boundary, so missed/coalesced notifications do not lose data.
505    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                        // `dispatch_reserved` publishes the exact owner before
521                        // making the task runnable. The task converts that
522                        // placeholder into the one live registration here.
523                        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                        // The router's two-phase activation already owns the
535                        // right to publish the next exact owner. A direct/manual
536                        // entry must wait for that reservation to publish or
537                        // roll back; superseding its token can strand a raw
538                        // external runner behind a phantom router owner.
539                        (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    /// Bind an external actor driver to the current logical-session owner.
569    ///
570    /// Returns the active run id used to correlate every forwarded claim and
571    /// confirmation. A pending generation is pushed immediately, closing the
572    /// delivery-before-bind race without making this in-memory signal durable.
573    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    /// Remove a driver only when it still belongs to the same activation run.
594    /// A stale driver can therefore never unbind its successor.
595    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    /// True only while this exact activation run remains the logical owner.
609    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    /// Mark the current run as finalizing before its runner slot becomes
618    /// available. Deliveries racing this interval are retained in
619    /// `latest_generation` and handed to one successor by
620    /// [`finish_finalization`](Self::finish_finalization).
621    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        // Delivery commits and its activation watermark are durable before the
659        // producer asks the in-memory router to wake an owner. If that producer
660        // or this execution is cancelled between those steps, recover the
661        // authoritative eligible prefix rather than trusting only RAM.
662        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    /// Release the terminal owner and, if the durable transcript cursor is
730    /// behind a concurrently delivered inbox generation, reserve and launch
731    /// exactly one successor using the injected real runtime adapter.
732    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                // Both execution paths drop their safe-point receiver before
758                // calling finish_finalization. Once the durable cursor has
759                // caught up and no reservation exists, this target carries no
760                // live routing state and must not remain in AppState forever.
761                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    /// Release an activation attempt whose caller was cancelled or panicked.
792    ///
793    /// Finalization-racing producers already received a coalesced success while
794    /// the old owner was visible, so no external waiter remains to retry after
795    /// the reservation token is released. The first abandoned attempt therefore
796    /// gets one detached, router-level retry for the newest still-undispatched
797    /// generation. That retry is deliberately bounded: its own Drop releases
798    /// the token but does not recursively hot-loop a broken spawner.
799    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                        // Cancellation recovery or a newer exact reservation
887                        // may have superseded this attempt while its spawner
888                        // was preparing. Direct/manual registration is
889                        // serialized behind `activation_reserved`, so it
890                        // cannot create an unrelated owner in this window.
891                        // Dropping the unlaunched value rolls back its external
892                        // slot.
893                        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                // The existing runner slot is already reserved. Publish owner
913                // state first, then make the task runnable exactly once.
914                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                // A coalesced caller does not claim success until the preceding
1006                // reservation publishes an owner or releases its slot. If that
1007                // attempt returns NoWork/error, this caller wakes and reserves
1008                // its own (newer) generation, so no durable delivery is left
1009                // behind a false-positive coalesced response.
1010                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        // Model the gap between the server runner reservation and its delayed
1630        // router registration: an independent direct SDK entry point wins the
1631        // logical owner first.
1632        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            // Reserved actor-child runs use the raw router handshake rather
1733            // than SessionRunRegistration, so recovery cannot depend on that
1734            // guard's Drop implementation.
1735            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        // The spawner can now return its reserved slot, but publishing the
1867        // logical owner is blocked. Aborting in this exact window must drop the
1868        // launch and invoke its external rollback.
1869        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        // Force the first request to yield after receiving the rollback-capable
1918        // real registry reservation but before publishing router ownership.
1919        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        // Simulate the successor hitting the same permanently failing
2046        // checkpoint/poison claim. The same generation remains inspectable but
2047        // cannot recursively launch provider loops.
2048        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        // A genuinely newer delivery gets one new bounded attempt.
2058        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}