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) async fn sync_in_place(&mut self) -> Result<bool> {
75 self.sync_history().await?;
76 let original_ordinal = self.materialized.applied_event_ordinal;
77 let original_digest = self.materialized.applied_event_digest.clone();
78 let original_operational = self.operational.clone();
79 let mut repaired = false;
80 let mut repaired_frontiers = std::collections::HashSet::new();
81 loop {
82 let after_ordinal = self.materialized.applied_event_ordinal;
83 match self.catch_up_fixed_frontier().await {
84 Ok(()) => break,
85 Err(error) if error.downcast_ref::<ProjectionAdvancedError>().is_some() => {
86 let (durable, window) = load_projection(&self.materialized.session_id).await?;
87 if durable.applied_event_ordinal <= after_ordinal {
88 return Err(error);
89 }
90 self.materialized = durable;
91 self.window = window;
92 continue;
93 }
94 Err(error) if relay_desynchronized(&error) => {
95 self.repair_projection()
96 .await
97 .with_context(|| {
98 format!(
99 "controller projection for {} cannot catch up from ordinal {after_ordinal}: {error:#}",
100 self.materialized.session_id
101 )
102 })?;
103 repaired = true;
104 let frontier = self.materialized.applied_event_ordinal;
111 if !repaired_frontiers.insert(frontier) {
112 bail!(
113 "controller projection for {} cannot catch up: relay history is \
114 unreadable and rebuilding from checkpoint frontier {frontier} does \
115 not get past it",
116 self.materialized.session_id
117 );
118 }
119 continue;
120 }
121 Err(error) => return Err(error),
122 }
123 }
124 let previous_requests = self.subagent_requests.clone();
125 let previous_results = self.subagent_results.clone();
126 (self.subagent_requests, self.subagent_results) = self.client.subagent_requests().await?;
127 let changed = repaired
128 || self.materialized.applied_event_ordinal != original_ordinal
129 || self.materialized.applied_event_digest != original_digest
130 || self.operational != original_operational
131 || self.subagent_requests != previous_requests
132 || self.subagent_results != previous_results;
133 Ok(changed)
134 }
135
136 async fn sync_history(&mut self) -> Result<()> {
139 while let Some(completed) = self.history_jobs.try_join_next() {
140 match completed {
141 Ok(result) => {
142 self.history_active.remove(&result.request_id);
143 self.client.complete_history_request(result).await?;
144 }
145 Err(error) => {
146 tracing::error!(%error, "history task failed; pending requests will retry");
147 self.history_jobs.abort_all();
148 while let Some(result) = self.history_jobs.join_next().await {
149 if let Err(error) = result
150 && !error.is_cancelled()
151 {
152 tracing::error!(%error, "history task failed during cleanup");
153 }
154 }
155 self.history_active.clear();
156 }
157 }
158 }
159 for request in self.client.history_requests().await? {
160 if self.history_active.len() >= mj_core::history::MAX_PENDING {
161 break;
162 }
163 if self.history_active.insert(request.request_id.clone()) {
164 self.history_jobs
165 .spawn(crate::sessionwiki::history::execute(request));
166 }
167 }
168 Ok(())
169 }
170
171 pub(super) async fn catch_up_fixed_frontier(&mut self) -> Result<()> {
176 let after = RelayCursor {
177 ordinal: self.materialized.applied_event_ordinal,
178 digest: self.materialized.applied_event_digest.clone(),
179 };
180 let catch_up = self
181 .client
182 .begin_catch_up(after.ordinal, &after.digest)
183 .await?;
184 let mut cursor = self.apply_event_page(catch_up.first_page).await?;
185 let mut pages_remaining = catch_up.frontier.ordinal.saturating_sub(cursor.ordinal);
186 while cursor.ordinal < catch_up.frontier.ordinal {
187 ensure!(
188 pages_remaining > 0,
189 "relay catch-up exceeded its fixed page bound"
190 );
191 pages_remaining -= 1;
192 let page = self
193 .client
194 .next_catch_up_page(&cursor, &catch_up.frontier)
195 .await?;
196 cursor = self.apply_event_page(page).await?;
197 }
198 ensure!(
199 cursor == catch_up.frontier,
200 "controller projection did not reach the captured relay frontier"
201 );
202 if cursor.ordinal > 0 {
203 let acknowledged = self
204 .client
205 .acknowledge(cursor.ordinal, &cursor.digest)
206 .await?;
207 ensure!(
208 acknowledged == cursor,
209 "relay acknowledged cursor {}:{} instead of {}:{}",
210 acknowledged.ordinal,
211 acknowledged.digest,
212 cursor.ordinal,
213 cursor.digest,
214 );
215 }
216 let mut operational = catch_up.state;
217 operational.acknowledged_through = cursor.ordinal;
218 operational.acknowledged_digest = cursor.digest;
219 self.operational = operational;
220 Ok(())
221 }
222
223 pub(super) async fn repair_projection(&mut self) -> Result<()> {
224 let state = crate::database::load_state()?;
225 let record = state
226 .sessions
227 .get(&self.materialized.session_id)
228 .context("controller session disappeared while repairing its projection")?;
229 let Some(checkpoint) = record.checkpoint.as_ref() else {
230 let replacement = MaterializedSession::empty(&self.materialized.session_id);
231 self.client
232 .attach(
233 replacement.applied_event_ordinal,
234 &replacement.applied_event_digest,
235 )
236 .await
237 .context("relay cannot rebuild the projection from its genesis")?;
238 save_materialized_session(&replacement)?;
239 self.window = mj_core::state::ProjectionWindow::of(&replacement);
240 self.materialized = replacement;
241 return Ok(());
242 };
243 let checkpoint_path = checkpoint.archive_path.clone();
244 let archive = tokio::task::spawn_blocking(move || {
245 verify_archive_streaming(&checkpoint_path).with_context(|| {
246 format!(
247 "verify projection repair checkpoint {}",
248 checkpoint_path.display()
249 )
250 })
251 })
252 .await
253 .context("projection repair archive verification task failed")??;
254 ensure!(
255 archive.archive_sha256 == checkpoint.sha256,
256 "projection repair checkpoint checksum does not match controller metadata"
257 );
258 ensure!(
259 archive.manifest.session.id == self.materialized.session_id,
260 "projection repair checkpoint belongs to session {}, not {}",
261 archive.manifest.session.id,
262 self.materialized.session_id
263 );
264 let canonical = archive.canonical_session;
265 ensure!(
266 canonical.event_frontier == checkpoint.event_frontier,
267 "projection repair checkpoint metadata frontier {} does not match archive frontier {}",
268 checkpoint.event_frontier,
269 canonical.event_frontier
270 );
271
272 self.client
276 .attach(canonical.event_frontier, &canonical.event_frontier_digest)
277 .await
278 .context("relay rejected the verified checkpoint repair cursor")?;
279 let replacement =
280 materialized_session_from_canonical(&self.materialized.session_id, &canonical)?;
281 save_materialized_session(&replacement)?;
282 self.window = mj_core::state::ProjectionWindow::of(&replacement);
283 self.materialized = replacement;
284 self.window.trim(
285 &mut self.materialized,
286 crate::database::PROJECTION_TAIL_ITEMS,
287 );
288 Ok(())
289 }
290
291 pub fn operational(&self) -> &RelayOperationalState {
293 &self.operational
294 }
295
296 pub fn snapshot(&self) -> ManagedSessionSnapshot {
297 ManagedSessionSnapshot {
298 window: self.window.clone(),
299 materialized: self.materialized.clone(),
300 operational: self.operational.clone(),
301 latest_credential_sync_signal: self.latest_credential_sync_signal.clone(),
302 worker_build: self.client.worker_build().map(str::to_owned),
303 subagent_requests: self.subagent_requests.clone(),
304 subagent_results: self.subagent_results.clone(),
305 }
306 }
307
308 pub async fn complete_subagent_request(
309 &mut self,
310 result: mj_core::subagent::SubagentToolResult,
311 ) -> Result<()> {
312 self.client.complete_subagent_request(result).await?;
313 (self.subagent_requests, self.subagent_results) = self.client.subagent_requests().await?;
314 Ok(())
315 }
316
317 pub async fn submit_accepted(
326 &mut self,
327 command_id: String,
328 command: RelayCommand,
329 ) -> Result<u64> {
330 self.client.submit(command_id, command).await
331 }
332
333 pub async fn submit(&mut self, command_id: String, command: RelayCommand) -> Result<u64> {
334 let ordinal = self.submit_accepted(command_id, command).await?;
335 self.sync_in_place().await?;
336 Ok(ordinal)
337 }
338
339 pub async fn respond_elicitation(
340 &mut self,
341 elicitation_id: String,
342 response: ElicitationResponse,
343 ) -> Result<()> {
344 self.client
345 .respond_elicitation(elicitation_id, response)
346 .await?;
347 self.sync_in_place().await?;
348 Ok(())
349 }
350
351 pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
352 self.client.stop_background_task(background_task_id).await?;
353 self.sync_in_place().await?;
354 Ok(())
355 }
356
357 pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
360 self.client.install_prompt_context(text).await
361 }
362
363 pub(super) async fn apply_event_page(&mut self, page: RelayEventPage) -> Result<RelayCursor> {
368 for event in &page.events {
369 if let mj_core::relay::RelayObservation::CommandQueued {
370 command: RelayCommand::Prompt { prompt },
371 ..
372 } = &event.observation
373 {
374 for reference in mj_core::attachment::references(prompt)? {
375 if let Err(error) = self.client.cache_attachment(&reference).await {
376 tracing::warn!(
380 session_id = %self.materialized.session_id,
381 attachment = %reference.sha256,
382 %error,
383 "could not cache image attachment during replay"
384 );
385 }
386 }
387 }
388 }
389
390 let RelayEventPage {
391 events,
392 through_ordinal,
393 through_digest,
394 } = page;
395 let event_count = events.len();
396 let transaction_count = event_count.div_ceil(PROJECTION_TRANSACTION_EVENT_BUDGET);
397 let started = Instant::now();
398 for events in events.chunks(PROJECTION_TRANSACTION_EVENT_BUDGET) {
399 let session_id = self.materialized.session_id.clone();
400 let events = events.to_vec();
401 let projection = self.materialized.clone();
402 let mut window = self.window.clone();
403 let (projection, window, credential_sync_signal) = tokio::task::spawn_blocking(
407 move || -> Result<(MaterializedSession, mj_core::state::ProjectionWindow, Option<CredentialSyncSignal>)> {
408 let mut projection = projection;
411 if window.omitted_items > 0 {
412 let mut historical = crate::database::load_projection_references(&projection, &events)?;
413 window.omitted_items = window.omitted_items.checked_sub(historical.len())
414 .context("projection reference count exceeds omitted history")?;
415 historical.append(&mut projection.transcript);
416 historical.sort_by(|a,b| (a.position, &a.stable_id).cmp(&(b.position, &b.stable_id)));
417 projection.transcript = historical;
418 }
419 let mut projection_index = ProjectionIndex::new(&projection);
420 let mut credential_sync_signal = None;
421 let mut prepared = Vec::with_capacity(events.len());
422 for event in &events {
423 let mutation =
424 project_relay_event_indexed(&projection, &projection_index, event)?
425 .mutation;
426 prepared.push((
427 event.ordinal,
428 event.previous_digest.clone(),
429 event.digest.clone(),
430 mutation.clone(),
431 ));
432 apply_committed_projection_event_indexed(
433 &mut projection,
434 &mut projection_index,
435 event,
436 mutation,
437 )?;
438 if let Some(reason) = relay_event_credential_sync_reason(event) {
439 credential_sync_signal = Some(CredentialSyncSignal {
440 ordinal: event.ordinal,
441 reason,
442 });
443 }
444 }
445 drop(projection_index);
446 apply_projection_page(&session_id, move |committed| {
447 for (ordinal, previous_digest, digest, mutation) in prepared {
448 match committed.apply(ordinal, &previous_digest, &digest, &mutation)? {
449 ProjectionApplyOutcome::Applied => {}
450 ProjectionApplyOutcome::AlreadyApplied => {
451 return Err(ProjectionAdvancedError {
452 event_ordinal: ordinal,
453 }
454 .into());
455 }
456 }
457 }
458 window.trim(&mut projection, crate::database::PROJECTION_TAIL_ITEMS);
459 Ok((projection, window, credential_sync_signal))
460 })
461 },
462 )
463 .await
464 .context("relay projection page task failed")??;
465 self.materialized = projection;
466 self.window = window;
467 if let Some(signal) = credential_sync_signal {
468 self.latest_credential_sync_signal = Some(signal);
469 }
470 }
471 tracing::debug!(target: "mj_controller::latency", session_id = %self.materialized.session_id,
472 through_ordinal, event_count, transaction_count, elapsed_ms = started.elapsed().as_secs_f64() * 1000.0,
473 "projection page committed");
474 if transaction_count > 1 {
475 tracing::debug!(
476 session_id = self.materialized.session_id,
477 event_count,
478 transaction_count,
479 elapsed_ms = started.elapsed().as_millis(),
480 "applied a large relay page in bounded projection transactions"
481 );
482 }
483 let delivered_through = self.materialized.applied_event_ordinal;
484 ensure!(
485 delivered_through == through_ordinal,
486 "relay page claimed frontier {} but delivered through {delivered_through}",
487 through_ordinal
488 );
489 ensure!(
490 self.materialized.applied_event_digest == through_digest,
491 "relay page digest does not match its claimed frontier"
492 );
493 Ok(RelayCursor {
494 ordinal: delivered_through,
495 digest: self.materialized.applied_event_digest.clone(),
496 })
497 }
498
499 pub async fn sync_project_memory(&mut self) -> Result<()> {
504 let Some(target) = self.project_memory.clone() else {
505 return Ok(());
506 };
507 if !self.client.supports_project_memory_sync() {
508 tracing::warn!(
509 session_id = self.materialized.session_id,
510 "worker protocol predates project-memory synchronization; preserving memory through checkpoints only"
511 );
512 self.project_memory = None;
513 return Ok(());
514 }
515 let (baseline, replica) = match self.client.project_memory_snapshot().await {
516 Ok(snapshot) => snapshot,
517 Err(error)
518 if error
519 .downcast_ref::<RelayRejected>()
520 .is_some_and(|rejected| {
521 rejected.0.code == mj_core::relay::RelayErrorCode::InvalidState
522 }) =>
523 {
524 tracing::warn!(
525 session_id = self.materialized.session_id,
526 "worker has no project-memory endpoint; preserving memory through checkpoints only"
527 );
528 self.project_memory = None;
529 return Ok(());
530 }
531 Err(error) => return Err(error),
532 };
533 let canonical_root = target.canonical_root;
534 let session_id = self.materialized.session_id.clone();
535 let (reconciliation, worker_install_needed) = tokio::task::spawn_blocking(move || {
536 let reconciliation = mj_core::project_memory::reconcile_into_canonical(
537 &canonical_root,
538 &baseline,
539 &replica,
540 &session_id,
541 )?;
542 let worker_install_needed =
543 reconciliation.merged != baseline || reconciliation.merged != replica;
544 Ok::<_, anyhow::Error>((reconciliation, worker_install_needed))
545 })
546 .await
547 .context("project memory reconciliation task failed")??;
548 for conflict in &reconciliation.conflicts {
549 tracing::warn!(session_id = self.materialized.session_id, %conflict, "project memory conflict preserved");
550 }
551 if worker_install_needed {
552 self.client
553 .install_project_memory_snapshot(reconciliation.merged)
554 .await?;
555 }
556 Ok(())
557 }
558}
559
560pub(super) fn relay_desynchronized(error: &anyhow::Error) -> bool {
561 error.chain().any(|cause| {
562 cause
563 .downcast_ref::<RelayRejected>()
564 .is_some_and(RelayRejected::is_desynchronized)
565 })
566}
567
568pub(super) fn projection_integrity_failure(error: &anyhow::Error) -> bool {
569 error
570 .chain()
571 .any(|cause| cause.downcast_ref::<ProjectionIntegrityError>().is_some())
572}
573
574#[cfg(test)]
580pub(super) struct ReplacementSessionTestFixture {
581 pub(super) stopped: ManagedSessionHandle,
582 pub(super) control: SessionManagerControl,
583 pub(super) submitted: mpsc::UnboundedReceiver<RelayCommand>,
584}
585
586#[cfg(test)]
590pub(super) fn replacement_session_test_fixture(
591 session_id: &str,
592 accepted_ordinal: u64,
593) -> ReplacementSessionTestFixture {
594 let (stopped_commands, stopped_commands_rx) = mpsc::channel(1);
595 drop(stopped_commands_rx);
596 let (stopped_releases, stopped_releases_rx) = mpsc::unbounded_channel();
597 drop(stopped_releases_rx);
598 let (stopped_view_tx, stopped_view) = watch::channel(ManagedSessionView::default());
599 drop(stopped_view_tx);
600 let stopped = ManagedSessionHandle {
601 session_id: session_id.to_owned(),
602 commands: stopped_commands,
603 releases: stopped_releases,
604 view: stopped_view,
605 };
606
607 let (commands, mut commands_rx) = mpsc::channel(4);
608 let (releases, _releases_rx) = mpsc::unbounded_channel();
609 let (view_tx, view) = watch::channel(ManagedSessionView::default());
610 let replacement = ManagedSessionHandle {
611 session_id: session_id.to_owned(),
612 commands,
613 releases,
614 view,
615 };
616 let actor_session_id = session_id.to_owned();
617 let (submitted_tx, submitted) = mpsc::unbounded_channel();
618 tokio::spawn(async move {
619 let _view_tx = view_tx;
620 while let Some(command) = commands_rx.recv().await {
621 match command {
622 ActorCommand::Submit { command, reply, .. } => {
623 let _ = submitted_tx.send(command);
626 let _ = reply.send(Ok(accepted_ordinal));
627 }
628 ActorCommand::Sync { reply } => {
629 let _ = reply.send(Ok(()));
630 }
631 command => command.reject(&actor_session_id, "unsupported test operation"),
632 }
633 }
634 });
635
636 let (manager_commands, mut manager_commands_rx) = mpsc::channel(4);
637 let manager_replacement = replacement.clone();
638 tokio::spawn(async move {
639 while let Some(ManagerCommand::Session {
640 session_id: requested,
641 reply,
642 }) = manager_commands_rx.recv().await
643 {
644 let resolved =
645 (requested == manager_replacement.session_id).then(|| manager_replacement.clone());
646 let _ = reply.send(resolved);
647 }
648 });
649 ReplacementSessionTestFixture {
650 stopped,
651 submitted,
652 control: SessionManagerControl {
653 commands: manager_commands,
654 },
655 }
656}