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 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 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 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 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 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 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 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 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 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 pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
392 self.client.install_prompt_context(text).await
393 }
394
395 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 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 let (projection, window, credential_sync_signal) = tokio::task::spawn_blocking(
439 move || -> Result<(MaterializedSession, mj_core::state::ProjectionWindow, Option<CredentialSyncSignal>)> {
440 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 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#[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#[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 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}