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