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