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