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 set_subagent_admission(&mut self, open: bool) -> Result<()> {
350 self.client.set_subagent_admission(open).await
351 }
352
353 pub async fn submit_accepted(
362 &mut self,
363 command_id: String,
364 command: RelayCommand,
365 ) -> Result<u64> {
366 self.client.submit(command_id, command).await
367 }
368
369 pub async fn submit(&mut self, command_id: String, command: RelayCommand) -> Result<u64> {
370 let ordinal = self.submit_accepted(command_id, command).await?;
371 self.sync_in_place().await?;
372 Ok(ordinal)
373 }
374
375 pub async fn respond_elicitation(
376 &mut self,
377 elicitation_id: String,
378 response: ElicitationResponse,
379 ) -> Result<()> {
380 self.client
381 .respond_elicitation(elicitation_id, response)
382 .await?;
383 self.sync_in_place().await?;
384 Ok(())
385 }
386
387 pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
388 self.client.stop_background_task(background_task_id).await?;
389 self.sync_in_place().await?;
390 Ok(())
391 }
392
393 pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
396 self.client.install_prompt_context(text).await
397 }
398
399 pub(super) async fn apply_event_page(&mut self, page: RelayEventPage) -> Result<RelayCursor> {
404 for event in &page.events {
405 if let mj_core::relay::RelayObservation::CommandQueued {
406 command: RelayCommand::Prompt { prompt },
407 ..
408 } = &event.observation
409 {
410 for reference in mj_core::attachment::references(prompt)? {
411 if let Err(error) = self.client.cache_attachment(&reference).await {
412 tracing::warn!(
416 session_id = %self.materialized.session_id,
417 attachment = %reference.sha256,
418 %error,
419 "could not cache image attachment during replay"
420 );
421 }
422 }
423 }
424 }
425
426 let RelayEventPage {
427 events,
428 through_ordinal,
429 through_digest,
430 } = page;
431 let event_count = events.len();
432 let transaction_count = event_count.div_ceil(PROJECTION_TRANSACTION_EVENT_BUDGET);
433 let started = Instant::now();
434 for events in events.chunks(PROJECTION_TRANSACTION_EVENT_BUDGET) {
435 let session_id = self.materialized.session_id.clone();
436 let events = events.to_vec();
437 let projection = self.materialized.clone();
438 let mut window = self.window.clone();
439 let (projection, window, credential_sync_signal) = tokio::task::spawn_blocking(
443 move || -> Result<(MaterializedSession, mj_core::state::ProjectionWindow, Option<CredentialSyncSignal>)> {
444 let mut projection = projection;
447 if window.omitted_items > 0 {
448 let mut historical = crate::database::load_projection_references(&projection, &events)?;
449 window.omitted_items = window.omitted_items.checked_sub(historical.len())
450 .context("projection reference count exceeds omitted history")?;
451 historical.append(&mut projection.transcript);
452 historical.sort_by(|a,b| (a.position, &a.stable_id).cmp(&(b.position, &b.stable_id)));
453 projection.transcript = historical;
454 }
455 let mut projection_index = ProjectionIndex::new(&projection);
456 let mut credential_sync_signal = None;
457 let mut prepared = Vec::with_capacity(events.len());
458 for event in &events {
459 let mutation =
460 project_relay_event_indexed(&projection, &projection_index, event)?
461 .mutation;
462 prepared.push((
463 event.ordinal,
464 event.previous_digest.clone(),
465 event.digest.clone(),
466 mutation.clone(),
467 ));
468 apply_committed_projection_event_indexed(
469 &mut projection,
470 &mut projection_index,
471 event,
472 mutation,
473 )?;
474 if let Some(reason) = relay_event_credential_sync_reason(event) {
475 credential_sync_signal = Some(CredentialSyncSignal {
476 ordinal: event.ordinal,
477 reason,
478 });
479 }
480 }
481 drop(projection_index);
482 apply_projection_page(&session_id, move |committed| {
483 for (ordinal, previous_digest, digest, mutation) in prepared {
484 match committed.apply(ordinal, &previous_digest, &digest, &mutation)? {
485 ProjectionApplyOutcome::Applied => {}
486 ProjectionApplyOutcome::AlreadyApplied => {
487 return Err(ProjectionAdvancedError {
488 event_ordinal: ordinal,
489 }
490 .into());
491 }
492 }
493 }
494 window.trim(&mut projection, crate::database::PROJECTION_TAIL_ITEMS);
495 Ok((projection, window, credential_sync_signal))
496 })
497 },
498 )
499 .await
500 .context("relay projection page task failed")??;
501 self.materialized = projection;
502 self.window = window;
503 if let Some(signal) = credential_sync_signal {
504 self.latest_credential_sync_signal = Some(signal);
505 }
506 }
507 tracing::debug!(target: "mj_controller::latency", session_id = %self.materialized.session_id,
508 through_ordinal, event_count, transaction_count, elapsed_ms = started.elapsed().as_secs_f64() * 1000.0,
509 "projection page committed");
510 if transaction_count > 1 {
511 tracing::debug!(
512 session_id = self.materialized.session_id,
513 event_count,
514 transaction_count,
515 elapsed_ms = started.elapsed().as_millis(),
516 "applied a large relay page in bounded projection transactions"
517 );
518 }
519 let delivered_through = self.materialized.applied_event_ordinal;
520 ensure!(
521 delivered_through == through_ordinal,
522 "relay page claimed frontier {} but delivered through {delivered_through}",
523 through_ordinal
524 );
525 ensure!(
526 self.materialized.applied_event_digest == through_digest,
527 "relay page digest does not match its claimed frontier"
528 );
529 Ok(RelayCursor {
530 ordinal: delivered_through,
531 digest: self.materialized.applied_event_digest.clone(),
532 })
533 }
534
535 pub async fn sync_project_memory(
540 &mut self,
541 config: &mj_core::config::Config,
542 cancel: tokio_util::sync::CancellationToken,
543 ) -> Result<()> {
544 let Some(target) = self.project_memory.clone() else {
545 return Ok(());
546 };
547 if !self.client.supports_project_memory_sync() {
548 tracing::warn!(
549 session_id = self.materialized.session_id,
550 "worker protocol predates project-memory synchronization; preserving memory through checkpoints only"
551 );
552 self.project_memory = None;
553 return Ok(());
554 }
555 let (baseline, replica) = match self.client.project_memory_snapshot().await {
556 Ok(snapshot) => snapshot,
557 Err(error)
558 if error
559 .downcast_ref::<RelayRejected>()
560 .is_some_and(|rejected| {
561 rejected.0.code == mj_core::relay::RelayErrorCode::InvalidState
562 }) =>
563 {
564 tracing::warn!(
565 session_id = self.materialized.session_id,
566 "worker has no project-memory endpoint; preserving memory through checkpoints only"
567 );
568 self.project_memory = None;
569 return Ok(());
570 }
571 Err(error) => return Err(error),
572 };
573 let resolver = crate::project_memory_merge::UtilityModelConflictResolver::new(config);
574 let sync = crate::project_memory_merge::merge_and_swap_canonical(
575 &target.canonical_root,
576 &baseline,
577 &replica,
578 &resolver,
579 cancel,
580 )
581 .await?;
582 let sync = match sync {
583 crate::project_memory_merge::CanonicalSyncOutcome::Swapped(sync) => sync,
584 crate::project_memory_merge::CanonicalSyncOutcome::AttemptsExhausted { attempts } => {
585 tracing::warn!(
586 session_id = self.materialized.session_id,
587 attempts,
588 "project memory changed during every canonical compare-and-swap attempt; checkpoint will continue"
589 );
590 return Ok(());
591 }
592 };
593 for resolution in &sync.resolutions {
594 tracing::warn!(
595 session_id = self.materialized.session_id,
596 path = %resolution.conflict.path,
597 resolution = resolution.method.as_str(),
598 "project memory conflict resolved"
599 );
600 tracing::debug!(
601 session_id = self.materialized.session_id,
602 path = %resolution.conflict.path,
603 base = ?resolution.conflict.base,
604 canonical = %resolution.conflict.canonical,
605 replica = %resolution.conflict.replica,
606 replica_wins = %resolution.conflict.replica_wins,
607 "project memory conflict source text"
608 );
609 }
610 let install = crate::project_memory_merge::install_merged_tree(
611 &mut self.client,
612 &baseline,
613 &replica,
614 sync.tree,
615 )
616 .await?;
617 if install == Some(crate::project_memory_merge::WorkerInstallOutcome::ReplicaChanged) {
618 tracing::debug!(
619 session_id = self.materialized.session_id,
620 "worker project-memory replica changed after the snapshot; skipped installing the merged tree"
621 );
622 }
623 Ok(())
624 }
625}
626
627pub(super) fn relay_desynchronized(error: &anyhow::Error) -> bool {
628 error.chain().any(|cause| {
629 cause
630 .downcast_ref::<RelayRejected>()
631 .is_some_and(RelayRejected::is_desynchronized)
632 })
633}
634
635pub(super) fn projection_integrity_failure(error: &anyhow::Error) -> bool {
636 error
637 .chain()
638 .any(|cause| cause.downcast_ref::<ProjectionIntegrityError>().is_some())
639}