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