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