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