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