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 session_id = self.materialized.session_id.clone();
225 let record =
226 tokio::task::spawn_blocking(move || crate::database::load_session_record(&session_id))
227 .await
228 .context("projection repair record reader failed")??
229 .context("controller session disappeared while repairing its projection")?;
230 let Some(checkpoint) = record.checkpoint.as_ref() else {
231 let replacement = MaterializedSession::empty(&self.materialized.session_id);
232 self.client
233 .attach(
234 replacement.applied_event_ordinal,
235 &replacement.applied_event_digest,
236 )
237 .await
238 .context("relay cannot rebuild the projection from its genesis")?;
239 save_materialized_session(&replacement)?;
240 self.window = mj_core::state::ProjectionWindow::of(&replacement);
241 self.materialized = replacement;
242 return Ok(());
243 };
244 let checkpoint_path = checkpoint.archive_path.clone();
245 let archive = tokio::task::spawn_blocking(move || {
246 verify_archive_streaming(&checkpoint_path).with_context(|| {
247 format!(
248 "verify projection repair checkpoint {}",
249 checkpoint_path.display()
250 )
251 })
252 })
253 .await
254 .context("projection repair archive verification task failed")??;
255 ensure!(
256 archive.archive_sha256 == checkpoint.sha256,
257 "projection repair checkpoint checksum does not match controller metadata"
258 );
259 ensure!(
260 archive.manifest.session.id == self.materialized.session_id,
261 "projection repair checkpoint belongs to session {}, not {}",
262 archive.manifest.session.id,
263 self.materialized.session_id
264 );
265 let canonical = archive.canonical_session;
266 ensure!(
267 canonical.event_frontier == checkpoint.event_frontier,
268 "projection repair checkpoint metadata frontier {} does not match archive frontier {}",
269 checkpoint.event_frontier,
270 canonical.event_frontier
271 );
272
273 self.client
277 .attach(canonical.event_frontier, &canonical.event_frontier_digest)
278 .await
279 .context("relay rejected the verified checkpoint repair cursor")?;
280 let replacement =
281 materialized_session_from_canonical(&self.materialized.session_id, &canonical)?;
282 save_materialized_session(&replacement)?;
283 self.window = mj_core::state::ProjectionWindow::of(&replacement);
284 self.materialized = replacement;
285 self.window.trim(
286 &mut self.materialized,
287 crate::database::PROJECTION_TAIL_ITEMS,
288 );
289 Ok(())
290 }
291
292 pub fn operational(&self) -> &RelayOperationalState {
294 &self.operational
295 }
296
297 pub fn snapshot(&self) -> ManagedSessionSnapshot {
298 ManagedSessionSnapshot {
299 window: self.window.clone(),
300 materialized: self.materialized.clone(),
301 operational: self.operational.clone(),
302 latest_credential_sync_signal: self.latest_credential_sync_signal.clone(),
303 worker_build: self.client.worker_build().map(str::to_owned),
304 subagent_requests: self.subagent_requests.clone(),
305 subagent_results: self.subagent_results.clone(),
306 }
307 }
308
309 pub async fn complete_subagent_request(
310 &mut self,
311 result: mj_core::subagent::SubagentToolResult,
312 ) -> Result<()> {
313 self.client.complete_subagent_request(result).await?;
314 (self.subagent_requests, self.subagent_results) = self.client.subagent_requests().await?;
315 Ok(())
316 }
317
318 pub async fn submit_accepted(
327 &mut self,
328 command_id: String,
329 command: RelayCommand,
330 ) -> Result<u64> {
331 self.client.submit(command_id, command).await
332 }
333
334 pub async fn submit(&mut self, command_id: String, command: RelayCommand) -> Result<u64> {
335 let ordinal = self.submit_accepted(command_id, command).await?;
336 self.sync_in_place().await?;
337 Ok(ordinal)
338 }
339
340 pub async fn respond_elicitation(
341 &mut self,
342 elicitation_id: String,
343 response: ElicitationResponse,
344 ) -> Result<()> {
345 self.client
346 .respond_elicitation(elicitation_id, response)
347 .await?;
348 self.sync_in_place().await?;
349 Ok(())
350 }
351
352 pub async fn stop_background_task(&mut self, background_task_id: String) -> Result<()> {
353 self.client.stop_background_task(background_task_id).await?;
354 self.sync_in_place().await?;
355 Ok(())
356 }
357
358 pub async fn install_prompt_context(&mut self, text: String) -> Result<()> {
361 self.client.install_prompt_context(text).await
362 }
363
364 pub(super) async fn apply_event_page(&mut self, page: RelayEventPage) -> Result<RelayCursor> {
369 for event in &page.events {
370 if let mj_core::relay::RelayObservation::CommandQueued {
371 command: RelayCommand::Prompt { prompt },
372 ..
373 } = &event.observation
374 {
375 for reference in mj_core::attachment::references(prompt)? {
376 if let Err(error) = self.client.cache_attachment(&reference).await {
377 tracing::warn!(
381 session_id = %self.materialized.session_id,
382 attachment = %reference.sha256,
383 %error,
384 "could not cache image attachment during replay"
385 );
386 }
387 }
388 }
389 }
390
391 let RelayEventPage {
392 events,
393 through_ordinal,
394 through_digest,
395 } = page;
396 let event_count = events.len();
397 let transaction_count = event_count.div_ceil(PROJECTION_TRANSACTION_EVENT_BUDGET);
398 let started = Instant::now();
399 for events in events.chunks(PROJECTION_TRANSACTION_EVENT_BUDGET) {
400 let session_id = self.materialized.session_id.clone();
401 let events = events.to_vec();
402 let projection = self.materialized.clone();
403 let mut window = self.window.clone();
404 let (projection, window, credential_sync_signal) = tokio::task::spawn_blocking(
408 move || -> Result<(MaterializedSession, mj_core::state::ProjectionWindow, Option<CredentialSyncSignal>)> {
409 let mut projection = projection;
412 if window.omitted_items > 0 {
413 let mut historical = crate::database::load_projection_references(&projection, &events)?;
414 window.omitted_items = window.omitted_items.checked_sub(historical.len())
415 .context("projection reference count exceeds omitted history")?;
416 historical.append(&mut projection.transcript);
417 historical.sort_by(|a,b| (a.position, &a.stable_id).cmp(&(b.position, &b.stable_id)));
418 projection.transcript = historical;
419 }
420 let mut projection_index = ProjectionIndex::new(&projection);
421 let mut credential_sync_signal = None;
422 let mut prepared = Vec::with_capacity(events.len());
423 for event in &events {
424 let mutation =
425 project_relay_event_indexed(&projection, &projection_index, event)?
426 .mutation;
427 prepared.push((
428 event.ordinal,
429 event.previous_digest.clone(),
430 event.digest.clone(),
431 mutation.clone(),
432 ));
433 apply_committed_projection_event_indexed(
434 &mut projection,
435 &mut projection_index,
436 event,
437 mutation,
438 )?;
439 if let Some(reason) = relay_event_credential_sync_reason(event) {
440 credential_sync_signal = Some(CredentialSyncSignal {
441 ordinal: event.ordinal,
442 reason,
443 });
444 }
445 }
446 drop(projection_index);
447 apply_projection_page(&session_id, move |committed| {
448 for (ordinal, previous_digest, digest, mutation) in prepared {
449 match committed.apply(ordinal, &previous_digest, &digest, &mutation)? {
450 ProjectionApplyOutcome::Applied => {}
451 ProjectionApplyOutcome::AlreadyApplied => {
452 return Err(ProjectionAdvancedError {
453 event_ordinal: ordinal,
454 }
455 .into());
456 }
457 }
458 }
459 window.trim(&mut projection, crate::database::PROJECTION_TAIL_ITEMS);
460 Ok((projection, window, credential_sync_signal))
461 })
462 },
463 )
464 .await
465 .context("relay projection page task failed")??;
466 self.materialized = projection;
467 self.window = window;
468 if let Some(signal) = credential_sync_signal {
469 self.latest_credential_sync_signal = Some(signal);
470 }
471 }
472 tracing::debug!(target: "mj_controller::latency", session_id = %self.materialized.session_id,
473 through_ordinal, event_count, transaction_count, elapsed_ms = started.elapsed().as_secs_f64() * 1000.0,
474 "projection page committed");
475 if transaction_count > 1 {
476 tracing::debug!(
477 session_id = self.materialized.session_id,
478 event_count,
479 transaction_count,
480 elapsed_ms = started.elapsed().as_millis(),
481 "applied a large relay page in bounded projection transactions"
482 );
483 }
484 let delivered_through = self.materialized.applied_event_ordinal;
485 ensure!(
486 delivered_through == through_ordinal,
487 "relay page claimed frontier {} but delivered through {delivered_through}",
488 through_ordinal
489 );
490 ensure!(
491 self.materialized.applied_event_digest == through_digest,
492 "relay page digest does not match its claimed frontier"
493 );
494 Ok(RelayCursor {
495 ordinal: delivered_through,
496 digest: self.materialized.applied_event_digest.clone(),
497 })
498 }
499
500 pub async fn sync_project_memory(&mut self) -> Result<()> {
505 let Some(target) = self.project_memory.clone() else {
506 return Ok(());
507 };
508 if !self.client.supports_project_memory_sync() {
509 tracing::warn!(
510 session_id = self.materialized.session_id,
511 "worker protocol predates project-memory synchronization; preserving memory through checkpoints only"
512 );
513 self.project_memory = None;
514 return Ok(());
515 }
516 let (baseline, replica) = match self.client.project_memory_snapshot().await {
517 Ok(snapshot) => snapshot,
518 Err(error)
519 if error
520 .downcast_ref::<RelayRejected>()
521 .is_some_and(|rejected| {
522 rejected.0.code == mj_core::relay::RelayErrorCode::InvalidState
523 }) =>
524 {
525 tracing::warn!(
526 session_id = self.materialized.session_id,
527 "worker has no project-memory endpoint; preserving memory through checkpoints only"
528 );
529 self.project_memory = None;
530 return Ok(());
531 }
532 Err(error) => return Err(error),
533 };
534 let canonical_root = target.canonical_root;
535 let session_id = self.materialized.session_id.clone();
536 let (reconciliation, worker_install_needed) = tokio::task::spawn_blocking(move || {
537 let reconciliation = mj_core::project_memory::reconcile_into_canonical(
538 &canonical_root,
539 &baseline,
540 &replica,
541 &session_id,
542 )?;
543 let worker_install_needed =
544 reconciliation.merged != baseline || reconciliation.merged != replica;
545 Ok::<_, anyhow::Error>((reconciliation, worker_install_needed))
546 })
547 .await
548 .context("project memory reconciliation task failed")??;
549 for conflict in &reconciliation.conflicts {
550 tracing::warn!(session_id = self.materialized.session_id, %conflict, "project memory conflict preserved");
551 }
552 if worker_install_needed {
553 self.client
554 .install_project_memory_snapshot(reconciliation.merged)
555 .await?;
556 }
557 Ok(())
558 }
559}
560
561pub(super) fn relay_desynchronized(error: &anyhow::Error) -> bool {
562 error.chain().any(|cause| {
563 cause
564 .downcast_ref::<RelayRejected>()
565 .is_some_and(RelayRejected::is_desynchronized)
566 })
567}
568
569pub(super) fn projection_integrity_failure(error: &anyhow::Error) -> bool {
570 error
571 .chain()
572 .any(|cause| cause.downcast_ref::<ProjectionIntegrityError>().is_some())
573}
574
575#[cfg(test)]
581pub(super) struct ReplacementSessionTestFixture {
582 pub(super) stopped: ManagedSessionHandle,
583 pub(super) control: SessionManagerControl,
584 pub(super) submitted: mpsc::UnboundedReceiver<RelayCommand>,
585}
586
587#[cfg(test)]
591pub(super) fn replacement_session_test_fixture(
592 session_id: &str,
593 accepted_ordinal: u64,
594) -> ReplacementSessionTestFixture {
595 let (stopped_commands, stopped_commands_rx) = mpsc::channel(1);
596 drop(stopped_commands_rx);
597 let (stopped_releases, stopped_releases_rx) = mpsc::unbounded_channel();
598 drop(stopped_releases_rx);
599 let (stopped_view_tx, stopped_view) = watch::channel(ManagedSessionView::default());
600 drop(stopped_view_tx);
601 let stopped = ManagedSessionHandle {
602 session_id: session_id.to_owned(),
603 commands: stopped_commands,
604 releases: stopped_releases,
605 view: stopped_view,
606 };
607
608 let (commands, mut commands_rx) = mpsc::channel(4);
609 let (releases, _releases_rx) = mpsc::unbounded_channel();
610 let (view_tx, view) = watch::channel(ManagedSessionView::default());
611 let replacement = ManagedSessionHandle {
612 session_id: session_id.to_owned(),
613 commands,
614 releases,
615 view,
616 };
617 let actor_session_id = session_id.to_owned();
618 let (submitted_tx, submitted) = mpsc::unbounded_channel();
619 tokio::spawn(async move {
620 let _view_tx = view_tx;
621 while let Some(command) = commands_rx.recv().await {
622 match command {
623 ActorCommand::Submit { command, reply, .. } => {
624 let _ = submitted_tx.send(command);
627 let _ = reply.send(Ok(accepted_ordinal));
628 }
629 ActorCommand::Sync { reply } => {
630 let _ = reply.send(Ok(()));
631 }
632 command => command.reject(&actor_session_id, "unsupported test operation"),
633 }
634 }
635 });
636
637 let (manager_commands, mut manager_commands_rx) = mpsc::channel(4);
638 let manager_replacement = replacement.clone();
639 tokio::spawn(async move {
640 while let Some(ManagerCommand::Session {
641 session_id: requested,
642 reply,
643 }) = manager_commands_rx.recv().await
644 {
645 let resolved =
646 (requested == manager_replacement.session_id).then(|| manager_replacement.clone());
647 let _ = reply.send(resolved);
648 }
649 });
650 ReplacementSessionTestFixture {
651 stopped,
652 submitted,
653 control: SessionManagerControl {
654 commands: manager_commands,
655 },
656 }
657}