1use super::*;
2
3pub struct StandaloneSession {
4 pub(super) client: RelayClient,
5 pub(super) materialized: MaterializedSession,
6 pub(super) window: mj_core::state::ProjectionWindow,
7 pub(super) operational: RelayOperationalState,
8 pub(super) latest_credential_sync_signal: Option<CredentialSyncSignal>,
9 pub(super) project_memory: Option<ProjectMemorySyncTarget>,
10 pub(super) subagent_requests: Vec<mj_core::subagent::SubagentToolRequest>,
11 pub(super) subagent_results: Vec<mj_core::subagent::SubagentToolResult>,
12 history_jobs: tokio::task::JoinSet<mj_core::history::HistoryResult>,
13 history_active: std::collections::BTreeSet<String>,
14}
15
16impl StandaloneSession {
17 pub fn set_project_memory_target(&mut self, target: Option<ProjectMemorySyncTarget>) {
18 self.project_memory = target;
19 }
20
21 pub async fn connect(target: &RelaySessionTarget) -> Result<Self> {
22 let mut client = RelayClient::connect(&target.spec, &target.session_id).await?;
26 let operational = client.status().await?;
27 let (materialized, window) = load_projection(&target.session_id).await?;
28 let mut connection = Self {
29 client,
30 materialized,
31 window,
32 operational,
33 latest_credential_sync_signal: None,
34 project_memory: target.project_memory.clone(),
35 subagent_requests: Vec::new(),
36 subagent_results: Vec::new(),
37 history_jobs: tokio::task::JoinSet::new(),
38 history_active: Default::default(),
39 };
40 connection.sync_in_place().await?;
41 Ok(connection)
42 }
43
44 pub async fn connect_command(spec: &CommandSpec, session_id: &str) -> Result<Self> {
45 Self::connect(&RelaySessionTarget {
46 session_id: session_id.to_owned(),
47 spec: spec.clone(),
48 worker_recovery: None,
49 project_memory: None,
50 })
51 .await
52 }
53
54 pub fn protocol_version(&self) -> u32 {
58 self.client.protocol_version()
59 }
60
61 pub async fn reserve_idle(&mut self, command_id: String) -> Result<bool> {
62 self.client.reserve_idle(command_id).await
63 }
64
65 pub(super) async fn detach(self) -> Result<()> {
66 self.client.detach().await
67 }
68
69 pub async fn sync(&mut self) -> Result<ManagedSessionSnapshot> {
70 self.sync_in_place().await?;
71 Ok(self.snapshot())
72 }
73
74 pub(super) fn follows_closely(&self) -> bool {
84 const FRESH_CLASSIFICATION_MS: i64 = 60_000;
85 let now = mj_core::clock::epoch_millis();
86 !self.history_active.is_empty()
87 || self.operational.has_work_in_flight()
88 || self
89 .operational
90 .assessment
91 .as_ref()
92 .is_some_and(|assessment| {
93 assessment.current()
94 && assessment.needs_classification(now)
95 && now.saturating_sub(assessment.completed_at_ms) < FRESH_CLASSIFICATION_MS
96 })
97 }
98
99 pub(super) async fn sync_in_place(&mut self) -> Result<bool> {
100 self.sync_history().await?;
101 let original_ordinal = self.materialized.applied_event_ordinal;
102 let original_digest = self.materialized.applied_event_digest.clone();
103 let original_operational = self.operational.clone();
104 let mut repaired = false;
105 let mut repaired_frontiers = std::collections::HashSet::new();
106 loop {
107 let after_ordinal = self.materialized.applied_event_ordinal;
108 match self.catch_up_fixed_frontier().await {
109 Ok(()) => break,
110 Err(error) if error.downcast_ref::<ProjectionAdvancedError>().is_some() => {
111 let (durable, window) = load_projection(&self.materialized.session_id).await?;
112 if durable.applied_event_ordinal <= after_ordinal {
113 return Err(error);
114 }
115 self.materialized = durable;
116 self.window = window;
117 continue;
118 }
119 Err(error) if relay_desynchronized(&error) => {
120 self.repair_projection()
121 .await
122 .with_context(|| {
123 format!(
124 "controller projection for {} cannot catch up from ordinal {after_ordinal}: {error:#}",
125 self.materialized.session_id
126 )
127 })?;
128 repaired = true;
129 let frontier = self.materialized.applied_event_ordinal;
136 if !repaired_frontiers.insert(frontier) {
137 bail!(
138 "controller projection for {} cannot catch up: relay history is \
139 unreadable and rebuilding from checkpoint frontier {frontier} does \
140 not get past it",
141 self.materialized.session_id
142 );
143 }
144 continue;
145 }
146 Err(error) => return Err(error),
147 }
148 }
149 let previous_requests = self.subagent_requests.clone();
150 let previous_results = self.subagent_results.clone();
151 (self.subagent_requests, self.subagent_results) = self.client.subagent_requests().await?;
152 let changed = repaired
153 || self.materialized.applied_event_ordinal != original_ordinal
154 || self.materialized.applied_event_digest != original_digest
155 || self.operational != original_operational
156 || self.subagent_requests != previous_requests
157 || self.subagent_results != previous_results;
158 Ok(changed)
159 }
160
161 async fn sync_history(&mut self) -> Result<()> {
164 while let Some(completed) = self.history_jobs.try_join_next() {
165 match completed {
166 Ok(result) => {
167 self.history_active.remove(&result.request_id);
168 self.client.complete_history_request(result).await?;
169 }
170 Err(error) => {
171 tracing::error!(%error, "history task failed; pending requests will retry");
172 self.history_jobs.abort_all();
173 while let Some(result) = self.history_jobs.join_next().await {
174 if let Err(error) = result
175 && !error.is_cancelled()
176 {
177 tracing::error!(%error, "history task failed during cleanup");
178 }
179 }
180 self.history_active.clear();
181 }
182 }
183 }
184 for request in self.client.history_requests().await? {
185 if self.history_active.len() >= mj_core::history::MAX_PENDING {
186 break;
187 }
188 if self.history_active.insert(request.request_id.clone()) {
189 self.history_jobs
190 .spawn(crate::sessionwiki::history::execute(request));
191 }
192 }
193 Ok(())
194 }
195
196 pub(super) async fn catch_up_fixed_frontier(&mut self) -> Result<()> {
201 let after = RelayCursor {
202 ordinal: self.materialized.applied_event_ordinal,
203 digest: self.materialized.applied_event_digest.clone(),
204 };
205 let catch_up = self
206 .client
207 .begin_catch_up(after.ordinal, &after.digest)
208 .await?;
209 let mut cursor = self.apply_event_page(catch_up.first_page).await?;
210 let mut pages_remaining = catch_up.frontier.ordinal.saturating_sub(cursor.ordinal);
211 while cursor.ordinal < catch_up.frontier.ordinal {
212 ensure!(
213 pages_remaining > 0,
214 "relay catch-up exceeded its fixed page bound"
215 );
216 pages_remaining -= 1;
217 let page = self
218 .client
219 .next_catch_up_page(&cursor, &catch_up.frontier)
220 .await?;
221 cursor = self.apply_event_page(page).await?;
222 }
223 ensure!(
224 cursor == catch_up.frontier,
225 "controller projection did not reach the captured relay frontier"
226 );
227 if cursor.ordinal > 0 {
228 let acknowledged = self
229 .client
230 .acknowledge(cursor.ordinal, &cursor.digest)
231 .await?;
232 ensure!(
233 acknowledged == cursor,
234 "relay acknowledged cursor {}:{} instead of {}:{}",
235 acknowledged.ordinal,
236 acknowledged.digest,
237 cursor.ordinal,
238 cursor.digest,
239 );
240 }
241 let mut operational = catch_up.state;
242 operational.acknowledged_through = cursor.ordinal;
243 operational.acknowledged_digest = cursor.digest;
244 self.operational = operational;
245 Ok(())
246 }
247
248 pub(super) async fn repair_projection(&mut self) -> Result<()> {
249 let session_id = self.materialized.session_id.clone();
250 let record =
251 tokio::task::spawn_blocking(move || crate::database::load_session_record(&session_id))
252 .await
253 .context("projection repair record reader failed")??
254 .context("controller session disappeared while repairing its projection")?;
255 let Some(checkpoint) = record.checkpoint.as_ref() else {
256 let replacement = MaterializedSession::empty(&self.materialized.session_id);
257 self.client
258 .attach(
259 replacement.applied_event_ordinal,
260 &replacement.applied_event_digest,
261 )
262 .await
263 .context("relay cannot rebuild the projection from its genesis")?;
264 save_materialized_session(&replacement)?;
265 self.window = mj_core::state::ProjectionWindow::of(&replacement);
266 self.materialized = replacement;
267 return Ok(());
268 };
269 let checkpoint_path = checkpoint.archive_path.clone();
270 let archive = tokio::task::spawn_blocking(move || {
271 verify_archive_streaming(&checkpoint_path).with_context(|| {
272 format!(
273 "verify projection repair checkpoint {}",
274 checkpoint_path.display()
275 )
276 })
277 })
278 .await
279 .context("projection repair archive verification task failed")??;
280 ensure!(
281 archive.archive_sha256 == checkpoint.sha256,
282 "projection repair checkpoint checksum does not match controller metadata"
283 );
284 ensure!(
285 archive.manifest.session.id == self.materialized.session_id,
286 "projection repair checkpoint belongs to session {}, not {}",
287 archive.manifest.session.id,
288 self.materialized.session_id
289 );
290 let canonical = archive.canonical_session;
291 ensure!(
292 canonical.event_frontier == checkpoint.event_frontier,
293 "projection repair checkpoint metadata frontier {} does not match archive frontier {}",
294 checkpoint.event_frontier,
295 canonical.event_frontier
296 );
297
298 self.client
302 .attach(canonical.event_frontier, &canonical.event_frontier_digest)
303 .await
304 .context("relay rejected the verified checkpoint repair cursor")?;
305 let replacement =
306 materialized_session_from_canonical(&self.materialized.session_id, &canonical)?;
307 save_materialized_session(&replacement)?;
308 self.window = mj_core::state::ProjectionWindow::of(&replacement);
309 self.materialized = replacement;
310 self.window.trim(
311 &mut self.materialized,
312 crate::database::PROJECTION_TAIL_ITEMS,
313 );
314 Ok(())
315 }
316
317 pub fn operational(&self) -> &RelayOperationalState {
319 &self.operational
320 }
321
322 pub fn snapshot(&self) -> ManagedSessionSnapshot {
323 ManagedSessionSnapshot {
324 window: self.window.clone(),
325 materialized: self.materialized.clone(),
326 operational: self.operational.clone(),
327 latest_credential_sync_signal: self.latest_credential_sync_signal.clone(),
328 worker_build: self.client.worker_build().map(str::to_owned),
329 subagent_requests: self.subagent_requests.clone(),
330 subagent_results: self.subagent_results.clone(),
331 }
332 }
333
334 pub async fn complete_subagent_request(
335 &mut self,
336 result: mj_core::subagent::SubagentToolResult,
337 ) -> Result<()> {
338 self.client.complete_subagent_request(result).await?;
339 (self.subagent_requests, self.subagent_results) = self.client.subagent_requests().await?;
340 Ok(())
341 }
342
343 pub async fn submit_accepted(
352 &mut self,
353 command_id: String,
354 command: RelayCommand,
355 ) -> Result<u64> {
356 self.client.submit(command_id, command).await
357 }
358
359 pub async fn submit(&mut self, command_id: String, command: RelayCommand) -> Result<u64> {
360 let ordinal = self.submit_accepted(command_id, command).await?;
361 self.sync_in_place().await?;
362 Ok(ordinal)
363 }
364
365 pub async fn respond_elicitation(
366 &mut self,
367 elicitation_id: String,
368 response: ElicitationResponse,
369 ) -> Result<()> {
370 self.client
371 .respond_elicitation(elicitation_id, response)
372 .await?;
373 self.sync_in_place().await?;
374 Ok(())
375 }
376
377 pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
378 self.client.stop_background_task(background_task_id).await?;
379 self.sync_in_place().await?;
380 Ok(())
381 }
382
383 pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
386 self.client.install_prompt_context(text).await
387 }
388
389 pub(super) async fn apply_event_page(&mut self, page: RelayEventPage) -> Result<RelayCursor> {
394 for event in &page.events {
395 if let mj_core::relay::RelayObservation::CommandQueued {
396 command: RelayCommand::Prompt { prompt },
397 ..
398 } = &event.observation
399 {
400 for reference in mj_core::attachment::references(prompt)? {
401 if let Err(error) = self.client.cache_attachment(&reference).await {
402 tracing::warn!(
406 session_id = %self.materialized.session_id,
407 attachment = %reference.sha256,
408 %error,
409 "could not cache image attachment during replay"
410 );
411 }
412 }
413 }
414 }
415
416 let RelayEventPage {
417 events,
418 through_ordinal,
419 through_digest,
420 } = page;
421 let event_count = events.len();
422 let transaction_count = event_count.div_ceil(PROJECTION_TRANSACTION_EVENT_BUDGET);
423 let started = Instant::now();
424 for events in events.chunks(PROJECTION_TRANSACTION_EVENT_BUDGET) {
425 let session_id = self.materialized.session_id.clone();
426 let events = events.to_vec();
427 let projection = self.materialized.clone();
428 let mut window = self.window.clone();
429 let (projection, window, credential_sync_signal) = tokio::task::spawn_blocking(
433 move || -> Result<(MaterializedSession, mj_core::state::ProjectionWindow, Option<CredentialSyncSignal>)> {
434 let mut projection = projection;
437 if window.omitted_items > 0 {
438 let mut historical = crate::database::load_projection_references(&projection, &events)?;
439 window.omitted_items = window.omitted_items.checked_sub(historical.len())
440 .context("projection reference count exceeds omitted history")?;
441 historical.append(&mut projection.transcript);
442 historical.sort_by(|a,b| (a.position, &a.stable_id).cmp(&(b.position, &b.stable_id)));
443 projection.transcript = historical;
444 }
445 let mut projection_index = ProjectionIndex::new(&projection);
446 let mut credential_sync_signal = None;
447 let mut prepared = Vec::with_capacity(events.len());
448 for event in &events {
449 let mutation =
450 project_relay_event_indexed(&projection, &projection_index, event)?
451 .mutation;
452 prepared.push((
453 event.ordinal,
454 event.previous_digest.clone(),
455 event.digest.clone(),
456 mutation.clone(),
457 ));
458 apply_committed_projection_event_indexed(
459 &mut projection,
460 &mut projection_index,
461 event,
462 mutation,
463 )?;
464 if let Some(reason) = relay_event_credential_sync_reason(event) {
465 credential_sync_signal = Some(CredentialSyncSignal {
466 ordinal: event.ordinal,
467 reason,
468 });
469 }
470 }
471 drop(projection_index);
472 apply_projection_page(&session_id, move |committed| {
473 for (ordinal, previous_digest, digest, mutation) in prepared {
474 match committed.apply(ordinal, &previous_digest, &digest, &mutation)? {
475 ProjectionApplyOutcome::Applied => {}
476 ProjectionApplyOutcome::AlreadyApplied => {
477 return Err(ProjectionAdvancedError {
478 event_ordinal: ordinal,
479 }
480 .into());
481 }
482 }
483 }
484 window.trim(&mut projection, crate::database::PROJECTION_TAIL_ITEMS);
485 Ok((projection, window, credential_sync_signal))
486 })
487 },
488 )
489 .await
490 .context("relay projection page task failed")??;
491 self.materialized = projection;
492 self.window = window;
493 if let Some(signal) = credential_sync_signal {
494 self.latest_credential_sync_signal = Some(signal);
495 }
496 }
497 tracing::debug!(target: "mj_controller::latency", session_id = %self.materialized.session_id,
498 through_ordinal, event_count, transaction_count, elapsed_ms = started.elapsed().as_secs_f64() * 1000.0,
499 "projection page committed");
500 if transaction_count > 1 {
501 tracing::debug!(
502 session_id = self.materialized.session_id,
503 event_count,
504 transaction_count,
505 elapsed_ms = started.elapsed().as_millis(),
506 "applied a large relay page in bounded projection transactions"
507 );
508 }
509 let delivered_through = self.materialized.applied_event_ordinal;
510 ensure!(
511 delivered_through == through_ordinal,
512 "relay page claimed frontier {} but delivered through {delivered_through}",
513 through_ordinal
514 );
515 ensure!(
516 self.materialized.applied_event_digest == through_digest,
517 "relay page digest does not match its claimed frontier"
518 );
519 Ok(RelayCursor {
520 ordinal: delivered_through,
521 digest: self.materialized.applied_event_digest.clone(),
522 })
523 }
524
525 pub async fn sync_project_memory(&mut self) -> Result<()> {
530 let Some(target) = self.project_memory.clone() else {
531 return Ok(());
532 };
533 if !self.client.supports_project_memory_sync() {
534 tracing::warn!(
535 session_id = self.materialized.session_id,
536 "worker protocol predates project-memory synchronization; preserving memory through checkpoints only"
537 );
538 self.project_memory = None;
539 return Ok(());
540 }
541 let (baseline, replica) = match self.client.project_memory_snapshot().await {
542 Ok(snapshot) => snapshot,
543 Err(error)
544 if error
545 .downcast_ref::<RelayRejected>()
546 .is_some_and(|rejected| {
547 rejected.0.code == mj_core::relay::RelayErrorCode::InvalidState
548 }) =>
549 {
550 tracing::warn!(
551 session_id = self.materialized.session_id,
552 "worker has no project-memory endpoint; preserving memory through checkpoints only"
553 );
554 self.project_memory = None;
555 return Ok(());
556 }
557 Err(error) => return Err(error),
558 };
559 let canonical_root = target.canonical_root;
560 let session_id = self.materialized.session_id.clone();
561 let (reconciliation, worker_install_needed) = tokio::task::spawn_blocking(move || {
562 let reconciliation = mj_core::project_memory::reconcile_into_canonical(
563 &canonical_root,
564 &baseline,
565 &replica,
566 &session_id,
567 )?;
568 let worker_install_needed =
569 reconciliation.merged != baseline || reconciliation.merged != replica;
570 Ok::<_, anyhow::Error>((reconciliation, worker_install_needed))
571 })
572 .await
573 .context("project memory reconciliation task failed")??;
574 for conflict in &reconciliation.conflicts {
575 tracing::warn!(session_id = self.materialized.session_id, %conflict, "project memory conflict preserved");
576 }
577 if worker_install_needed {
578 self.client
579 .install_project_memory_snapshot(reconciliation.merged)
580 .await?;
581 }
582 Ok(())
583 }
584}
585
586pub(super) fn relay_desynchronized(error: &anyhow::Error) -> bool {
587 error.chain().any(|cause| {
588 cause
589 .downcast_ref::<RelayRejected>()
590 .is_some_and(RelayRejected::is_desynchronized)
591 })
592}
593
594pub(super) fn projection_integrity_failure(error: &anyhow::Error) -> bool {
595 error
596 .chain()
597 .any(|cause| cause.downcast_ref::<ProjectionIntegrityError>().is_some())
598}
599
600#[cfg(test)]
606pub(super) struct ReplacementSessionTestFixture {
607 pub(super) stopped: ManagedSessionHandle,
608 pub(super) control: SessionManagerControl,
609 pub(super) submitted: mpsc::UnboundedReceiver<RelayCommand>,
610}
611
612#[cfg(test)]
616pub(super) fn replacement_session_test_fixture(
617 session_id: &str,
618 accepted_ordinal: u64,
619) -> ReplacementSessionTestFixture {
620 let (stopped_commands, stopped_commands_rx) = mpsc::channel(1);
621 drop(stopped_commands_rx);
622 let (stopped_releases, stopped_releases_rx) = mpsc::unbounded_channel();
623 drop(stopped_releases_rx);
624 let (stopped_view_tx, stopped_view) = watch::channel(ManagedSessionView::default());
625 drop(stopped_view_tx);
626 let stopped = ManagedSessionHandle {
627 session_id: session_id.to_owned(),
628 commands: stopped_commands,
629 releases: stopped_releases,
630 view: stopped_view,
631 };
632
633 let (commands, mut commands_rx) = mpsc::channel(4);
634 let (releases, _releases_rx) = mpsc::unbounded_channel();
635 let (view_tx, view) = watch::channel(ManagedSessionView::default());
636 let replacement = ManagedSessionHandle {
637 session_id: session_id.to_owned(),
638 commands,
639 releases,
640 view,
641 };
642 let actor_session_id = session_id.to_owned();
643 let (submitted_tx, submitted) = mpsc::unbounded_channel();
644 tokio::spawn(async move {
645 let _view_tx = view_tx;
646 while let Some(command) = commands_rx.recv().await {
647 match command {
648 ActorCommand::Submit { command, reply, .. } => {
649 let _ = submitted_tx.send(command);
652 let _ = reply.send(Ok(accepted_ordinal));
653 }
654 ActorCommand::Sync { reply } => {
655 let _ = reply.send(Ok(()));
656 }
657 command => command.reject(&actor_session_id, "unsupported test operation"),
658 }
659 }
660 });
661
662 let (manager_commands, mut manager_commands_rx) = mpsc::channel(4);
663 let manager_replacement = replacement.clone();
664 tokio::spawn(async move {
665 while let Some(ManagerCommand::Session {
666 session_id: requested,
667 reply,
668 }) = manager_commands_rx.recv().await
669 {
670 let resolved =
671 (requested == manager_replacement.session_id).then(|| manager_replacement.clone());
672 let _ = reply.send(resolved);
673 }
674 });
675 ReplacementSessionTestFixture {
676 stopped,
677 submitted,
678 control: SessionManagerControl {
679 commands: manager_commands,
680 },
681 }
682}