Skip to main content

mj_controller/session_manager/
standalone.rs

1use super::*;
2
3pub struct StandaloneSession {
4    pub(super) client: RelayClient,
5    pub(super) materialized: MaterializedSession,
6    pub(super) window: mj_core::state::ProjectionWindow,
7    pub(super) operational: RelayOperationalState,
8    pub(super) latest_credential_sync_signal: Option<CredentialSyncSignal>,
9    pub(super) project_memory: Option<ProjectMemorySyncTarget>,
10    pub(super) subagent_requests: Vec<mj_core::subagent::SubagentToolRequest>,
11    pub(super) subagent_results: Vec<mj_core::subagent::SubagentToolResult>,
12    history_jobs: tokio::task::JoinSet<mj_core::history::HistoryResult>,
13    history_active: std::collections::BTreeSet<String>,
14}
15
16impl StandaloneSession {
17    pub fn set_project_memory_target(&mut self, target: Option<ProjectMemorySyncTarget>) {
18        self.project_memory = target;
19    }
20
21    pub async fn connect(target: &RelaySessionTarget) -> Result<Self> {
22        // Reach the worker before reading the projection. A stored session can
23        // be tens of megabytes, and the reconnect loop would otherwise pay that
24        // whole synchronous read on every attempt against a worker that is down.
25        let mut client = RelayClient::connect(&target.spec, &target.session_id).await?;
26        let operational = client.status().await?;
27        let (materialized, window) = load_projection(&target.session_id).await?;
28        let mut connection = Self {
29            client,
30            materialized,
31            window,
32            operational,
33            latest_credential_sync_signal: None,
34            project_memory: target.project_memory.clone(),
35            subagent_requests: Vec::new(),
36            subagent_results: Vec::new(),
37            history_jobs: tokio::task::JoinSet::new(),
38            history_active: Default::default(),
39        };
40        connection.sync_in_place().await?;
41        Ok(connection)
42    }
43
44    pub async fn connect_command(spec: &CommandSpec, session_id: &str) -> Result<Self> {
45        Self::connect(&RelaySessionTarget {
46            session_id: session_id.to_owned(),
47            spec: spec.clone(),
48            worker_recovery: None,
49            project_memory: None,
50        })
51        .await
52    }
53
54    /// Protocol negotiated with the worker behind this connection. Lifecycle
55    /// operations use it to avoid sending a newly introduced command to an
56    /// older worker that cannot decode it.
57    pub fn protocol_version(&self) -> u32 {
58        self.client.protocol_version()
59    }
60
61    pub async fn reserve_idle(&mut self, command_id: String) -> Result<bool> {
62        self.client.reserve_idle(command_id).await
63    }
64
65    pub(super) async fn detach(self) -> Result<()> {
66        self.client.detach().await
67    }
68
69    pub async fn sync(&mut self) -> Result<ManagedSessionSnapshot> {
70        self.sync_in_place().await?;
71        Ok(self.snapshot())
72    }
73
74    /// Whether the actor should follow this worker at the fast cadence.
75    ///
76    /// True while the worker owns work whose events stream in, while it
77    /// classifies a turn that just ended (`mj prompt --wait` returns on that
78    /// verdict and automatic continuation acts on it, and it lands a moment
79    /// after the turn goes quiet), or while a history search this
80    /// connection runs still has a result to deliver on its next sync. A
81    /// classification that is still pending well after its turn ended is left
82    /// to the slow sweep, so no worker state can hold the fast cadence forever.
83    pub(super) fn follows_closely(&self) -> bool {
84        const FRESH_CLASSIFICATION_MS: i64 = 60_000;
85        let now = mj_core::clock::epoch_millis();
86        !self.history_active.is_empty()
87            || self.operational.has_work_in_flight()
88            || self
89                .operational
90                .assessment
91                .as_ref()
92                .is_some_and(|assessment| {
93                    assessment.current()
94                        && assessment.needs_classification(now)
95                        && now.saturating_sub(assessment.completed_at_ms) < FRESH_CLASSIFICATION_MS
96                })
97    }
98
99    pub(super) async fn sync_in_place(&mut self) -> Result<bool> {
100        self.sync_history().await?;
101        let original_ordinal = self.materialized.applied_event_ordinal;
102        let original_digest = self.materialized.applied_event_digest.clone();
103        let original_operational = self.operational.clone();
104        let mut repaired = false;
105        let mut repaired_frontiers = std::collections::HashSet::new();
106        loop {
107            let after_ordinal = self.materialized.applied_event_ordinal;
108            match self.catch_up_fixed_frontier().await {
109                Ok(()) => break,
110                Err(error) if error.downcast_ref::<ProjectionAdvancedError>().is_some() => {
111                    let (durable, window) = load_projection(&self.materialized.session_id).await?;
112                    if durable.applied_event_ordinal <= after_ordinal {
113                        return Err(error);
114                    }
115                    self.materialized = durable;
116                    self.window = window;
117                    continue;
118                }
119                Err(error) if relay_desynchronized(&error) => {
120                    self.repair_projection()
121                        .await
122                        .with_context(|| {
123                            format!(
124                                "controller projection for {} cannot catch up from ordinal {after_ordinal}: {error:#}",
125                                self.materialized.session_id
126                            )
127                        })?;
128                    repaired = true;
129                    // Repair rebuilds from the same durable checkpoint every
130                    // time. If catching up from that frontier still desyncs — as
131                    // it does when relay history is unreadable past the
132                    // checkpoint — repairing again lands on the same frontier and
133                    // would loop forever. Fail loudly on the second visit instead
134                    // of hanging; recovery got everything the checkpoint covers.
135                    let frontier = self.materialized.applied_event_ordinal;
136                    if !repaired_frontiers.insert(frontier) {
137                        bail!(
138                            "controller projection for {} cannot catch up: relay history is \
139                             unreadable and rebuilding from checkpoint frontier {frontier} does \
140                             not get past it",
141                            self.materialized.session_id
142                        );
143                    }
144                    continue;
145                }
146                Err(error) => return Err(error),
147            }
148        }
149        let previous_requests = self.subagent_requests.clone();
150        let previous_results = self.subagent_results.clone();
151        (self.subagent_requests, self.subagent_results) = self.client.subagent_requests().await?;
152        let changed = repaired
153            || self.materialized.applied_event_ordinal != original_ordinal
154            || self.materialized.applied_event_digest != original_digest
155            || self.operational != original_operational
156            || self.subagent_requests != previous_requests
157            || self.subagent_results != previous_results;
158        Ok(changed)
159    }
160
161    /// Poll only bounded messages here. Disk searches run independently of this
162    /// actor, and the owned JoinSet cancels them when the connection is retired.
163    async fn sync_history(&mut self) -> Result<()> {
164        while let Some(completed) = self.history_jobs.try_join_next() {
165            match completed {
166                Ok(result) => {
167                    self.history_active.remove(&result.request_id);
168                    self.client.complete_history_request(result).await?;
169                }
170                Err(error) => {
171                    tracing::error!(%error, "history task failed; pending requests will retry");
172                    self.history_jobs.abort_all();
173                    while let Some(result) = self.history_jobs.join_next().await {
174                        if let Err(error) = result
175                            && !error.is_cancelled()
176                        {
177                            tracing::error!(%error, "history task failed during cleanup");
178                        }
179                    }
180                    self.history_active.clear();
181                }
182            }
183        }
184        for request in self.client.history_requests().await? {
185            if self.history_active.len() >= mj_core::history::MAX_PENDING {
186                break;
187            }
188            if self.history_active.insert(request.request_id.clone()) {
189                self.history_jobs
190                    .spawn(crate::sessionwiki::history::execute(request));
191            }
192        }
193        Ok(())
194    }
195
196    /// Apply relay pages through the exact frontier captured by the first
197    /// response, then acknowledge that frontier once. Every projection page is
198    /// independently durable; delaying the relay's GC watermark avoids one
199    /// snapshot fsync per transport-sized page without risking redelivery.
200    pub(super) async fn catch_up_fixed_frontier(&mut self) -> Result<()> {
201        let after = RelayCursor {
202            ordinal: self.materialized.applied_event_ordinal,
203            digest: self.materialized.applied_event_digest.clone(),
204        };
205        let catch_up = self
206            .client
207            .begin_catch_up(after.ordinal, &after.digest)
208            .await?;
209        let mut cursor = self.apply_event_page(catch_up.first_page).await?;
210        let mut pages_remaining = catch_up.frontier.ordinal.saturating_sub(cursor.ordinal);
211        while cursor.ordinal < catch_up.frontier.ordinal {
212            ensure!(
213                pages_remaining > 0,
214                "relay catch-up exceeded its fixed page bound"
215            );
216            pages_remaining -= 1;
217            let page = self
218                .client
219                .next_catch_up_page(&cursor, &catch_up.frontier)
220                .await?;
221            cursor = self.apply_event_page(page).await?;
222        }
223        ensure!(
224            cursor == catch_up.frontier,
225            "controller projection did not reach the captured relay frontier"
226        );
227        if cursor.ordinal > 0 {
228            let acknowledged = self
229                .client
230                .acknowledge(cursor.ordinal, &cursor.digest)
231                .await?;
232            ensure!(
233                acknowledged == cursor,
234                "relay acknowledged cursor {}:{} instead of {}:{}",
235                acknowledged.ordinal,
236                acknowledged.digest,
237                cursor.ordinal,
238                cursor.digest,
239            );
240        }
241        let mut operational = catch_up.state;
242        operational.acknowledged_through = cursor.ordinal;
243        operational.acknowledged_digest = cursor.digest;
244        self.operational = operational;
245        Ok(())
246    }
247
248    pub(super) async fn repair_projection(&mut self) -> Result<()> {
249        let session_id = self.materialized.session_id.clone();
250        let record =
251            tokio::task::spawn_blocking(move || crate::database::load_session_record(&session_id))
252                .await
253                .context("projection repair record reader failed")??
254                .context("controller session disappeared while repairing its projection")?;
255        let Some(checkpoint) = record.checkpoint.as_ref() else {
256            let replacement = MaterializedSession::empty(&self.materialized.session_id);
257            self.client
258                .attach(
259                    replacement.applied_event_ordinal,
260                    &replacement.applied_event_digest,
261                )
262                .await
263                .context("relay cannot rebuild the projection from its genesis")?;
264            save_materialized_session(&replacement)?;
265            self.window = mj_core::state::ProjectionWindow::of(&replacement);
266            self.materialized = replacement;
267            return Ok(());
268        };
269        let checkpoint_path = checkpoint.archive_path.clone();
270        let archive = tokio::task::spawn_blocking(move || {
271            verify_archive_streaming(&checkpoint_path).with_context(|| {
272                format!(
273                    "verify projection repair checkpoint {}",
274                    checkpoint_path.display()
275                )
276            })
277        })
278        .await
279        .context("projection repair archive verification task failed")??;
280        ensure!(
281            archive.archive_sha256 == checkpoint.sha256,
282            "projection repair checkpoint checksum does not match controller metadata"
283        );
284        ensure!(
285            archive.manifest.session.id == self.materialized.session_id,
286            "projection repair checkpoint belongs to session {}, not {}",
287            archive.manifest.session.id,
288            self.materialized.session_id
289        );
290        let canonical = archive.canonical_session;
291        ensure!(
292            canonical.event_frontier == checkpoint.event_frontier,
293            "projection repair checkpoint metadata frontier {} does not match archive frontier {}",
294            checkpoint.event_frontier,
295            canonical.event_frontier
296        );
297
298        // Prove that the relay recognizes this exact event-chain cursor before
299        // replacing any controller state. A matching ordinal alone is not a
300        // repair proof.
301        self.client
302            .attach(canonical.event_frontier, &canonical.event_frontier_digest)
303            .await
304            .context("relay rejected the verified checkpoint repair cursor")?;
305        let replacement =
306            materialized_session_from_canonical(&self.materialized.session_id, &canonical)?;
307        save_materialized_session(&replacement)?;
308        self.window = mj_core::state::ProjectionWindow::of(&replacement);
309        self.materialized = replacement;
310        self.window.trim(
311            &mut self.materialized,
312            crate::database::PROJECTION_TAIL_ITEMS,
313        );
314        Ok(())
315    }
316
317    /// The relay state as of the last sync, without copying the projection.
318    pub fn operational(&self) -> &RelayOperationalState {
319        &self.operational
320    }
321
322    pub fn snapshot(&self) -> ManagedSessionSnapshot {
323        ManagedSessionSnapshot {
324            window: self.window.clone(),
325            materialized: self.materialized.clone(),
326            operational: self.operational.clone(),
327            latest_credential_sync_signal: self.latest_credential_sync_signal.clone(),
328            worker_build: self.client.worker_build().map(str::to_owned),
329            subagent_requests: self.subagent_requests.clone(),
330            subagent_results: self.subagent_results.clone(),
331        }
332    }
333
334    pub async fn complete_subagent_request(
335        &mut self,
336        result: mj_core::subagent::SubagentToolResult,
337    ) -> Result<()> {
338        self.client.complete_subagent_request(result).await?;
339        (self.subagent_requests, self.subagent_results) = self.client.subagent_requests().await?;
340        Ok(())
341    }
342
343    /// Hands one command to the relay and returns the ordinal it accepted it
344    /// at, without catching the local projection up to it.
345    ///
346    /// Callers that need the projection current call [`Self::sync`] after.
347    /// Keeping the two apart matters on the prompt path: the catch-up is the
348    /// expensive half, and a caller waiting to hear that the relay took the
349    /// command should not wait for it. It also stops a failed catch-up from
350    /// looking like a failed submission to a caller that would retry.
351    pub async fn submit_accepted(
352        &mut self,
353        command_id: String,
354        command: RelayCommand,
355    ) -> Result<u64> {
356        self.client.submit(command_id, command).await
357    }
358
359    pub async fn submit(&mut self, command_id: String, command: RelayCommand) -> Result<u64> {
360        let ordinal = self.submit_accepted(command_id, command).await?;
361        self.sync_in_place().await?;
362        Ok(ordinal)
363    }
364
365    pub async fn respond_elicitation(
366        &mut self,
367        elicitation_id: String,
368        response: ElicitationResponse,
369    ) -> Result<()> {
370        self.client
371            .respond_elicitation(elicitation_id, response)
372            .await?;
373        self.sync_in_place().await?;
374        Ok(())
375    }
376
377    pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
378        self.client.stop_background_task(background_task_id).await?;
379        self.sync_in_place().await?;
380        Ok(())
381    }
382
383    /// Persist relay-private context for the next real prompt. It never
384    /// contributes an event to the canonical projection.
385    pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
386        self.client.install_prompt_context(text).await
387    }
388
389    /// Apply one relay transport page in bounded durable chunks. A transport
390    /// page can contain thousands of events, but SQLite has one global writer;
391    /// regularly releasing it lets other session actors keep their views
392    /// current. The relay GC watermark advances only after the complete page.
393    pub(super) async fn apply_event_page(&mut self, page: RelayEventPage) -> Result<RelayCursor> {
394        for event in &page.events {
395            if let mj_core::relay::RelayObservation::CommandQueued {
396                command: RelayCommand::Prompt { prompt },
397                ..
398            } = &event.observation
399            {
400                for reference in mj_core::attachment::references(prompt)? {
401                    if let Err(error) = self.client.cache_attachment(&reference).await {
402                        // History remains readable even if a blob was lost. A
403                        // later submission still verifies every image before
404                        // admission, and must report missing data to the user.
405                        tracing::warn!(
406                            session_id = %self.materialized.session_id,
407                            attachment = %reference.sha256,
408                            %error,
409                            "could not cache image attachment during replay"
410                        );
411                    }
412                }
413            }
414        }
415
416        let RelayEventPage {
417            events,
418            through_ordinal,
419            through_digest,
420        } = page;
421        let event_count = events.len();
422        let transaction_count = event_count.div_ceil(PROJECTION_TRANSACTION_EVENT_BUDGET);
423        let started = Instant::now();
424        for events in events.chunks(PROJECTION_TRANSACTION_EVENT_BUDGET) {
425            let session_id = self.materialized.session_id.clone();
426            let events = events.to_vec();
427            let projection = self.materialized.clone();
428            let mut window = self.window.clone();
429            // Projection is CPU work and its durable page uses synchronous
430            // SQLite. Keep both off the async actor runtime so independent
431            // sessions stay responsive during each bounded catch-up chunk.
432            let (projection, window, credential_sync_signal) = tokio::task::spawn_blocking(
433                move || -> Result<(MaterializedSession, mj_core::state::ProjectionWindow, Option<CredentialSyncSignal>)> {
434                    // The in-memory projection advances on a working copy and
435                    // is published only once its page is durable.
436                    let mut projection = projection;
437                    if window.omitted_items > 0 {
438                        let mut historical = crate::database::load_projection_references(&projection, &events)?;
439                        window.omitted_items = window.omitted_items.checked_sub(historical.len())
440                            .context("projection reference count exceeds omitted history")?;
441                        historical.append(&mut projection.transcript);
442                        historical.sort_by(|a,b| (a.position, &a.stable_id).cmp(&(b.position, &b.stable_id)));
443                        projection.transcript = historical;
444                    }
445                    let mut projection_index = ProjectionIndex::new(&projection);
446                    let mut credential_sync_signal = None;
447                    let mut prepared = Vec::with_capacity(events.len());
448                    for event in &events {
449                        let mutation =
450                            project_relay_event_indexed(&projection, &projection_index, event)?
451                                .mutation;
452                        prepared.push((
453                            event.ordinal,
454                            event.previous_digest.clone(),
455                            event.digest.clone(),
456                            mutation.clone(),
457                        ));
458                        apply_committed_projection_event_indexed(
459                            &mut projection,
460                            &mut projection_index,
461                            event,
462                            mutation,
463                        )?;
464                        if let Some(reason) = relay_event_credential_sync_reason(event) {
465                            credential_sync_signal = Some(CredentialSyncSignal {
466                                ordinal: event.ordinal,
467                                reason,
468                            });
469                        }
470                    }
471                    drop(projection_index);
472                    apply_projection_page(&session_id, move |committed| {
473                        for (ordinal, previous_digest, digest, mutation) in prepared {
474                            match committed.apply(ordinal, &previous_digest, &digest, &mutation)? {
475                                ProjectionApplyOutcome::Applied => {}
476                                ProjectionApplyOutcome::AlreadyApplied => {
477                                    return Err(ProjectionAdvancedError {
478                                        event_ordinal: ordinal,
479                                    }
480                                    .into());
481                                }
482                            }
483                        }
484                        window.trim(&mut projection, crate::database::PROJECTION_TAIL_ITEMS);
485                        Ok((projection, window, credential_sync_signal))
486                    })
487                },
488            )
489            .await
490            .context("relay projection page task failed")??;
491            self.materialized = projection;
492            self.window = window;
493            if let Some(signal) = credential_sync_signal {
494                self.latest_credential_sync_signal = Some(signal);
495            }
496        }
497        tracing::debug!(target: "mj_controller::latency", session_id = %self.materialized.session_id,
498            through_ordinal, event_count, transaction_count, elapsed_ms = started.elapsed().as_secs_f64() * 1000.0,
499            "projection page committed");
500        if transaction_count > 1 {
501            tracing::debug!(
502                session_id = self.materialized.session_id,
503                event_count,
504                transaction_count,
505                elapsed_ms = started.elapsed().as_millis(),
506                "applied a large relay page in bounded projection transactions"
507            );
508        }
509        let delivered_through = self.materialized.applied_event_ordinal;
510        ensure!(
511            delivered_through == through_ordinal,
512            "relay page claimed frontier {} but delivered through {delivered_through}",
513            through_ordinal
514        );
515        ensure!(
516            self.materialized.applied_event_digest == through_digest,
517            "relay page digest does not match its claimed frontier"
518        );
519        Ok(RelayCursor {
520            ordinal: delivered_through,
521            digest: self.materialized.applied_event_digest.clone(),
522        })
523    }
524
525    /// Reconcile this worker's project-memory replica at an explicit durable
526    /// boundary. Normal relay attachment and polling must never perform this
527    /// filesystem work: a degraded target could otherwise turn reconnects
528    /// into an unbounded queue of timed-out snapshot writes.
529    pub async fn sync_project_memory(&mut self) -> Result<()> {
530        let Some(target) = self.project_memory.clone() else {
531            return Ok(());
532        };
533        if !self.client.supports_project_memory_sync() {
534            tracing::warn!(
535                session_id = self.materialized.session_id,
536                "worker protocol predates project-memory synchronization; preserving memory through checkpoints only"
537            );
538            self.project_memory = None;
539            return Ok(());
540        }
541        let (baseline, replica) = match self.client.project_memory_snapshot().await {
542            Ok(snapshot) => snapshot,
543            Err(error)
544                if error
545                    .downcast_ref::<RelayRejected>()
546                    .is_some_and(|rejected| {
547                        rejected.0.code == mj_core::relay::RelayErrorCode::InvalidState
548                    }) =>
549            {
550                tracing::warn!(
551                    session_id = self.materialized.session_id,
552                    "worker has no project-memory endpoint; preserving memory through checkpoints only"
553                );
554                self.project_memory = None;
555                return Ok(());
556            }
557            Err(error) => return Err(error),
558        };
559        let canonical_root = target.canonical_root;
560        let session_id = self.materialized.session_id.clone();
561        let (reconciliation, worker_install_needed) = tokio::task::spawn_blocking(move || {
562            let reconciliation = mj_core::project_memory::reconcile_into_canonical(
563                &canonical_root,
564                &baseline,
565                &replica,
566                &session_id,
567            )?;
568            let worker_install_needed =
569                reconciliation.merged != baseline || reconciliation.merged != replica;
570            Ok::<_, anyhow::Error>((reconciliation, worker_install_needed))
571        })
572        .await
573        .context("project memory reconciliation task failed")??;
574        for conflict in &reconciliation.conflicts {
575            tracing::warn!(session_id = self.materialized.session_id, %conflict, "project memory conflict preserved");
576        }
577        if worker_install_needed {
578            self.client
579                .install_project_memory_snapshot(reconciliation.merged)
580                .await?;
581        }
582        Ok(())
583    }
584}
585
586pub(super) fn relay_desynchronized(error: &anyhow::Error) -> bool {
587    error.chain().any(|cause| {
588        cause
589            .downcast_ref::<RelayRejected>()
590            .is_some_and(RelayRejected::is_desynchronized)
591    })
592}
593
594pub(super) fn projection_integrity_failure(error: &anyhow::Error) -> bool {
595    error
596        .chain()
597        .any(|cause| cause.downcast_ref::<ProjectionIntegrityError>().is_some())
598}
599
600/// A stopped actor and the manager that resolves its live replacement.
601///
602/// This fixture and its constructor are compiled unconditionally and hidden
603/// from the documentation because the chat crate's tests need them, and a
604/// `#[cfg(test)]` item is invisible to another crate.
605#[cfg(test)]
606pub(super) struct ReplacementSessionTestFixture {
607    pub(super) stopped: ManagedSessionHandle,
608    pub(super) control: SessionManagerControl,
609    pub(super) submitted: mpsc::UnboundedReceiver<RelayCommand>,
610}
611
612/// A stopped actor and a manager that resolves its live replacement. Chat
613/// tests use this hand-written actor instead of mocking the session manager
614/// protocol.
615#[cfg(test)]
616pub(super) fn replacement_session_test_fixture(
617    session_id: &str,
618    accepted_ordinal: u64,
619) -> ReplacementSessionTestFixture {
620    let (stopped_commands, stopped_commands_rx) = mpsc::channel(1);
621    drop(stopped_commands_rx);
622    let (stopped_releases, stopped_releases_rx) = mpsc::unbounded_channel();
623    drop(stopped_releases_rx);
624    let (stopped_view_tx, stopped_view) = watch::channel(ManagedSessionView::default());
625    drop(stopped_view_tx);
626    let stopped = ManagedSessionHandle {
627        session_id: session_id.to_owned(),
628        commands: stopped_commands,
629        releases: stopped_releases,
630        view: stopped_view,
631    };
632
633    let (commands, mut commands_rx) = mpsc::channel(4);
634    let (releases, _releases_rx) = mpsc::unbounded_channel();
635    let (view_tx, view) = watch::channel(ManagedSessionView::default());
636    let replacement = ManagedSessionHandle {
637        session_id: session_id.to_owned(),
638        commands,
639        releases,
640        view,
641    };
642    let actor_session_id = session_id.to_owned();
643    let (submitted_tx, submitted) = mpsc::unbounded_channel();
644    tokio::spawn(async move {
645        let _view_tx = view_tx;
646        while let Some(command) = commands_rx.recv().await {
647            match command {
648                ActorCommand::Submit { command, reply, .. } => {
649                    // Tests can drop the optional observer when they only
650                    // care about acceptance/reconnection.
651                    let _ = submitted_tx.send(command);
652                    let _ = reply.send(Ok(accepted_ordinal));
653                }
654                ActorCommand::Sync { reply } => {
655                    let _ = reply.send(Ok(()));
656                }
657                command => command.reject(&actor_session_id, "unsupported test operation"),
658            }
659        }
660    });
661
662    let (manager_commands, mut manager_commands_rx) = mpsc::channel(4);
663    let manager_replacement = replacement.clone();
664    tokio::spawn(async move {
665        while let Some(ManagerCommand::Session {
666            session_id: requested,
667            reply,
668        }) = manager_commands_rx.recv().await
669        {
670            let resolved =
671                (requested == manager_replacement.session_id).then(|| manager_replacement.clone());
672            let _ = reply.send(resolved);
673        }
674    });
675    ReplacementSessionTestFixture {
676        stopped,
677        submitted,
678        control: SessionManagerControl {
679            commands: manager_commands,
680        },
681    }
682}