1use super::*;
2
3pub async fn run_server(
4 args: ServerArgs,
5 termination: tokio_util::sync::CancellationToken,
6 worker: SessionManagerChannels,
7 daemon_runtime: Arc<RuntimeState>,
8 mut workspace_updates: tokio::sync::watch::Receiver<Vec<WorkspaceRecord>>,
9) -> Result<()> {
10 let resolved = resolve_server_args(args, termination.clone()).await?;
11 let bind = resolved.bind;
12 let mut controller = Controller::load()?;
13 let profile_catalog = profile_catalog::ProfileCatalog::new(termination.child_token());
17 profile_catalog.sync(&controller.config);
18 let mut daemon_revisions = daemon_runtime.revisions();
19 daemon_revisions.borrow_and_update();
20 let mut phone_workspaces = workspace_updates.borrow_and_update().clone();
21 let mut quotas = std::collections::BTreeMap::new();
22 let subagent_quota_reports = Arc::new(std::sync::Mutex::new(quotas.clone()));
23 let (quota_profiles_tx, mut quota_updates_rx) = spawn_quota_refresher();
24 let mut quota_batch = QuotaRefreshBatch::default();
25 let mut published_quota_profiles = std::collections::BTreeMap::new();
26 republish_quota_profiles(
27 &controller,
28 &mut published_quota_profiles,
29 &mut quota_batch,
30 "a_profiles_tx,
31 );
32 let mut revision = daemon_runtime.allocate_revision();
33 let mut conversations = std::collections::BTreeMap::new();
34 let mut queued_prompts = projected_queued_prompts(&controller)?;
35 let mut active_user_shells = std::collections::BTreeMap::new();
36 let mut pending_elicitations = std::collections::BTreeMap::new();
37 let mut prompt_images = std::collections::BTreeSet::new();
38 let mut operational = std::collections::BTreeMap::new();
39 let mut native_agents =
40 load_native_agents(controller.state.sessions.keys().cloned().collect()).await?;
41 let mut materialized_activity = load_materialized_activity(&controller).await?;
42 let mut project_sources = PhoneProjectSources::default();
43 let (records, lifecycles) = daemon_runtime.session_projection();
44 controller.state.sessions = records;
45 let mut operations = lifecycles
46 .iter()
47 .map(|view| (view.session_id.clone(), viewer_operation(view)))
48 .collect::<std::collections::BTreeMap<_, _>>();
49 let mut move_recoveries = ViewerMoveRecoveries::new();
50 let (move_recovery_tx, mut move_recovery_rx) =
51 tokio::sync::mpsc::unbounded_channel::<Result<ViewerMoveRecoveries, String>>();
52 let mut move_recovery_load_in_flight = false;
53 let mut launch_failures = Vec::new();
54 let mut capacity_state: std::collections::BTreeMap<String, PhoneCapacity> =
58 std::collections::BTreeMap::new();
59 let (capacity_targets_tx, capacity_triggers_tx, mut capacity_updates_rx) =
60 crate::pollers::spawn_dashboard_capacity_poller();
61 let (snapshot_tx, snapshot_rx) = tokio::sync::watch::channel(viewer_snapshot(
62 &controller,
63 &phone_workspaces,
64 "as,
65 &PhoneSessionViews {
66 native_agents: &native_agents,
67 conversations: &conversations,
68 queued_prompts: &queued_prompts,
69 active_user_shells: &active_user_shells,
70 pending_elicitations: &pending_elicitations,
71 prompt_images: &prompt_images,
72 operational: &operational,
73 materialized_activity: &materialized_activity,
74 project_sources: &project_sources,
75 operations: &operations,
76 move_recoveries: &move_recoveries,
77 capacity: &viewer_capacity(&capacity_state),
78 launch_failures: &launch_failures,
79 reviews: &review_views(&daemon_runtime),
80 },
81 revision,
82 ));
83 let (conversation_tx, conversation_rx) = tokio::sync::watch::channel(conversations.clone());
84 let (action_tx, mut action_rx) = tokio::sync::mpsc::channel(32);
85 let (bundle_tx, mut bundle_rx) = tokio::sync::mpsc::channel(16);
86 let (receipt_tx, mut receipt_rx) = tokio::sync::mpsc::channel(32);
87 let (preflight_tx, mut preflight_rx) = tokio::sync::mpsc::channel(32);
88 let (move_preparation_tx, mut move_preparation_rx) = tokio::sync::mpsc::channel(32);
89 let (client_state_tx, mut client_state_rx) = tokio::sync::mpsc::channel(64);
90 let (dictation_tx, mut dictation_rx) =
91 tokio::sync::mpsc::channel::<crate::dictation::DictationRequest>(8);
92 let (background_task_stop_tx, mut background_task_stop_rx) =
93 tokio::sync::mpsc::channel::<BackgroundTaskStopRequest>(32);
94 let SessionManagerChannels {
95 targets: worker_targets_tx,
96 control: worker_commands_tx,
97 updates: mut worker_updates_rx,
98 shutdown: worker_shutdown,
99 } = worker;
100 worker_targets_tx.send_replace(dashboard_worker_targets(&controller));
101 publish_capacity_targets(&controller, &capacity_targets_tx, &mut capacity_state);
102 let mut credential_sync = CredentialSyncCoordinator::spawn();
103 let credential_sync_handle = credential_sync.handle();
104 credential_sync_handle.set_targets(credential_sync_targets(&controller));
105 let mut credential_sync_signals = CredentialSyncSignalTracker::default();
106 let mut credential_sync_notices = CredentialSyncNotices::default();
107 let options_session_ttl = crate::server::default_session_ttl();
109 let activity_snapshots = snapshot_rx.clone();
110 let mut options = ServerOptions::new(
111 bind,
112 snapshot_rx,
113 conversation_rx,
114 crate::server::ServerRequests {
115 action_tx,
116 bundle_tx,
117 receipt_tx,
118 preflight_tx,
119 move_preparation_tx,
120 client_state_tx,
121 dictation_tx,
122 },
123 )?;
124 options.set_background_task_stop_tx(background_task_stop_tx);
125 options.shutdown = termination.clone();
126 let cookie_key_path = crate::server::cookie_key_path();
129 options.load_cookie_credentials(cookie_key_path).await?;
130 options.set_api_token(crate::server::load_or_create_api_token(
133 &crate::server::api_token_path(),
134 )?);
135 let api_runtime = daemon_runtime.clone();
136 let api_backend = Arc::new(
137 api::ApiBackend::new(
138 worker_commands_tx.client(),
139 Arc::new(move |session_id: &str| api_runtime.session_state(session_id)),
140 daemon_runtime.clone(),
141 )
142 .with_quota_reports(subagent_quota_reports.clone())
143 .with_profile_catalog(profile_catalog.clone()),
144 );
145 options.set_subagent_backend(api_backend.clone());
146 let renewal_cancellation = termination.child_token();
147 let mut renewal_task = None;
148 if let Some((cert, key)) = resolved.tls_files {
149 let rustls = axum_server::tls_rustls::RustlsConfig::from_pem_file(cert, key)
150 .await
151 .context("load web viewer TLS certificate")?;
152 options.set_tls_config(rustls.clone());
153 if let Some(tailscale) = resolved.tailscale {
154 renewal_task = Some(spawn_tailscale_cert_renewer(
155 tailscale,
156 rustls,
157 renewal_cancellation.clone(),
158 ));
159 }
160 } else if bind.ip().is_loopback() {
161 options.secure_cookie = false;
162 } else {
163 anyhow::bail!("non-loopback web viewer requires TLS");
164 }
165 let fallback_reason = resolved.fallback_reason;
166 let qr_login_url = if fallback_reason.is_none() && resolved.viewer_url.starts_with("https://") {
167 let encoded = url::form_urlencoded::byte_serialize(options.login_token().as_bytes())
168 .collect::<String>();
169 Some(format!(
170 "{}/auth/login?token={encoded}",
171 resolved.viewer_url.trim_end_matches('/')
172 ))
173 } else {
174 None
175 };
176 let ready = crate::server::WebViewerAccess::Ready {
177 viewer_url: resolved.viewer_url,
178 viewer_code: options.viewer_code().to_owned(),
179 qr_login_url,
180 fallback_reason,
181 };
182
183 let mut serve = ViewerServer::spawn({
184 let daemon_runtime = daemon_runtime.clone();
185 async move {
186 crate::web_viewer::serve(options, ready, &daemon_runtime.web_viewer, |access| {
187 daemon_runtime.publish_web_access(access);
188 })
189 .await
190 }
191 });
192 let conversation_projection_shutdown = termination.child_token();
193 let control = async {
194 let mut credential_tick = tokio::time::interval(Duration::from_millis(250));
195 let mut prune_tick = tokio::time::interval(prune_tick_interval());
199 daemon_runtime.wiki().request_sync(true);
204 let client_state_retention = options_session_ttl;
205 let (action_done_tx, mut action_done_rx) = tokio::sync::mpsc::unbounded_channel::<(
206 u64,
207 Option<String>,
208 std::result::Result<(), PhoneActionFailure>,
209 )>();
210 let (action_started_tx, mut action_started_rx) =
211 tokio::sync::mpsc::unbounded_channel::<PhoneActionStarted>();
212 let (receipt_done_tx, mut receipt_done_rx) =
213 tokio::sync::mpsc::unbounded_channel::<ReadReceiptPersisted>();
214 let (controller_reload_tx, mut controller_reload_rx) =
215 tokio::sync::mpsc::unbounded_channel::<ControllerReloaded>();
216 let (bundle_done_tx, mut bundle_done_rx) =
217 tokio::sync::mpsc::unbounded_channel::<BundleCreated>();
218 let (move_prepared_tx, mut move_prepared_rx) =
219 tokio::sync::mpsc::unbounded_channel::<MovePrepared>();
220 let mut dictation_jobs = tokio::task::JoinSet::new();
221 let mut archive_jobs = tokio::task::JoinSet::new();
222 let mut bundle_jobs = tokio::task::JoinSet::new();
223 let mut preflight_jobs = tokio::task::JoinSet::new();
224 let mut move_preparation_jobs = tokio::task::JoinSet::new();
225 let mut move_recovery_jobs = tokio::task::JoinSet::new();
226 let mut native_agent_jobs = tokio::task::JoinSet::new();
227 let mut native_agents_dirty = true;
228 let mut background_task_stop_jobs = tokio::task::JoinSet::new();
229 let mut background_task_stop_open = true;
230 let mut controller_reload_in_flight = false;
231 let mut controller_reload_requested = false;
232 let mut controller_reload_invalidated = false;
233 let mut pending_action_errors = std::collections::BTreeMap::<String, String>::new();
234 let mut active_actions = std::collections::BTreeSet::new();
235 let mut closing_actions = std::collections::BTreeMap::<String, u64>::new();
236 let mut next_action_id = 0_u64;
237 let mut action_cancellations = std::collections::BTreeMap::<u64, PhoneActionControl>::new();
238 let mut action_sessions = std::collections::BTreeMap::<u64, String>::new();
239 let mut action_replies = PendingActionReplies::default();
240 let mut launch_workspaces = std::collections::BTreeMap::new();
241 let mut subagent_jobs = tokio::task::JoinSet::new();
242 let mut subagent_completion_jobs = tokio::task::JoinSet::new();
243 let mut active_subagent_requests = std::collections::BTreeSet::new();
244 let (conversation_projection_tx, mut conversation_projection_rx) =
245 tokio::sync::mpsc::channel(CONVERSATION_PROJECTION_CHANNEL_CAPACITY);
246 let mut conversation_projections = ConversationProjectionDispatcher::new(
247 conversation_projection_tx,
248 conversation_projection_shutdown.clone(),
249 );
250 let mut quota_updates_open = true;
251 let mut failure: Option<anyhow::Error> = None;
255 request_move_recovery_reload(
256 &move_recovery_tx,
257 &mut move_recovery_load_in_flight,
258 &mut move_recovery_jobs,
259 );
260 macro_rules! publish_snapshot {
261 ($revision:expr) => {
262 let (records, lifecycles) = daemon_runtime.session_projection();
263 controller.state.sessions = records;
264 for (session_id, error) in &pending_action_errors {
265 if let Some(session) = controller.state.sessions.get_mut(session_id)
266 && session.last_error.is_none()
267 {
268 session.last_error = Some(error.clone());
269 }
270 }
271 operations = lifecycles.iter()
272 .map(|view| (view.session_id.clone(), viewer_operation(view)))
273 .collect();
274 let snapshot = viewer_snapshot(
275 &controller,
276 &phone_workspaces,
277 "as,
278 &PhoneSessionViews {
279 native_agents: &native_agents,
280 conversations: &conversations,
281 queued_prompts: &queued_prompts,
282 active_user_shells: &active_user_shells,
283 pending_elicitations: &pending_elicitations,
284 prompt_images: &prompt_images,
285 operational: &operational,
286 materialized_activity: &materialized_activity,
287 project_sources: &project_sources,
288 operations: &operations,
289 move_recoveries: &move_recoveries,
290 capacity: &viewer_capacity(&capacity_state),
291 launch_failures: &launch_failures,
292 reviews: &review_views(&daemon_runtime),
293 },
294 $revision,
295 );
296 if let Err(error) = snapshot_tx.send(snapshot) {
297 tracing::debug!(revision = $revision, %error, "phone snapshot delivery failed; no viewer is subscribed");
298 }
299 };
300 }
301 loop {
302 if native_agents_dirty && native_agent_jobs.is_empty() {
303 native_agents_dirty = false;
304 native_agent_jobs.spawn(load_native_agents(
305 controller.state.sessions.keys().cloned().collect(),
306 ));
307 }
308 project_sources.synchronize(&controller);
309 tokio::select! {
310 _ = termination.cancelled() => break,
311 request = background_task_stop_rx.recv(), if background_task_stop_open => {
312 let Some(request) = request else {
313 background_task_stop_open = false;
314 tracing::warn!("background-task stop request feed closed while the phone server was running");
315 continue;
316 };
317 let session_control = worker_commands_tx.clone();
318 background_task_stop_jobs.spawn(async move {
319 let result = match session_control.session(&request.session_id).await {
320 Ok(session) => session
321 .client()
322 .stop_background_task(request.background_task_id.clone())
323 .await
324 .map_err(|error| {
325 tracing::warn!(
326 session_id = %request.session_id,
327 background_task_id = %request.background_task_id,
328 %error,
329 "provider rejected background-task stop"
330 );
331 BackgroundTaskStopFailure::Provider
332 }),
333 Err(error) => {
334 tracing::warn!(
335 session_id = %request.session_id,
336 %error,
337 "could not resolve live session for background-task stop"
338 );
339 Err(BackgroundTaskStopFailure::SessionUnavailable)
340 }
341 };
342 if request.reply.send(result).is_err() {
343 tracing::debug!(
344 session_id = %request.session_id,
345 background_task_id = %request.background_task_id,
346 "background-task stop result dropped after viewer disconnected"
347 );
348 }
349 });
350 }
351 completed = background_task_stop_jobs.join_next(), if !background_task_stop_jobs.is_empty() => {
352 if let Some(Err(error)) = completed {
353 tracing::error!(%error, "background-task stop task failed unexpectedly");
354 }
355 }
356 move_reloaded = move_recovery_rx.recv() => {
357 let Some(result) = move_reloaded else {
358 failure = feed_stopped(
359 termination.is_cancelled(),
360 "the Move recovery projection stopped while the phone server was running",
361 );
362 break;
363 };
364 move_recovery_load_in_flight = false;
365 match result {
366 Ok(recoveries) => {
367 move_recoveries = recoveries;
368 revision = daemon_runtime.allocate_revision();
369 publish_snapshot!(revision);
370 }
371 Err(error) => tracing::warn!(%error, "could not refresh Move recovery projection"),
372 }
373 }
374 resolved = project_sources.jobs.join_next(), if !project_sources.jobs.is_empty() => {
375 match resolved {
376 Some(Ok(resolved)) => project_sources.complete(resolved),
377 Some(Err(error)) => {
378 failure = Some(anyhow::anyhow!("web project source task failed: {error}"));
379 break;
380 }
381 None => unreachable!("project source jobs were not empty"),
382 }
383 revision = daemon_runtime.allocate_revision();
384 publish_snapshot!(revision);
385 }
386 changed = daemon_revisions.changed() => {
387 native_agents_dirty = true;
388 if changed.is_err() {
389 failure = feed_stopped(
390 termination.is_cancelled(),
391 "the daemon stopped publishing runtime revisions to the phone server",
392 );
393 break;
394 }
395 daemon_revisions.borrow_and_update();
396 revision = daemon_runtime.allocate_revision();
397 publish_snapshot!(revision);
398 request_controller_reload(
399 &mut controller_reload_in_flight,
400 &mut controller_reload_requested,
401 &controller_reload_tx,
402 );
403 }
404 changed = workspace_updates.changed() => {
405 if changed.is_err() {
406 failure = feed_stopped(
407 termination.is_cancelled(),
408 "the daemon stopped publishing workspaces to the phone server",
409 );
410 break;
411 }
412 phone_workspaces = workspace_updates.borrow_and_update().clone();
413 revision = daemon_runtime.allocate_revision();
414 publish_snapshot!(revision);
415 }
416 update = capacity_updates_rx.recv() => {
417 let Some(update) = update else {
418 failure = feed_stopped(termination.is_cancelled(), "the capacity poller stopped while the phone server was running");
419 break;
420 };
421 if let Some(entry) = capacity_state.get_mut(&update.target_id) {
422 entry.refreshing = false;
423 entry.sampled_at_epoch_seconds = Some(update.sampled_at_epoch_seconds);
424 match update.result {
425 Ok(usage) => {
426 entry.on_demand = usage.is_none();
430 entry.usage = usage;
431 entry.failed = false;
432 }
433 Err(_) => entry.failed = true,
437 }
438 }
439 revision = daemon_runtime.allocate_revision();
440 publish_snapshot!(revision);
441 }
442 update = quota_updates_rx.recv(), if quota_updates_open => {
443 match update {
444 Some(QuotaUpdate::Report(outcome)) => {
445 if outcome.credentials_changed {
446 credential_sync_handle
447 .sync_profile_now(&outcome.report.profile_id, None);
448 }
449 quotas.insert(outcome.report.profile_id.clone(), outcome.report.clone());
450 subagent_quota_reports
451 .lock()
452 .expect("sub-agent quota reports lock poisoned")
453 .insert(outcome.report.profile_id.clone(), outcome.report);
454 revision = daemon_runtime.allocate_revision();
455 publish_snapshot!(revision);
456 }
457 Some(QuotaUpdate::Refreshing { .. } | QuotaUpdate::Finished { .. }) => {}
458 None => {
459 quota_updates_open = false;
460 tracing::warn!("quota refresher stopped while the phone server is running");
461 }
462 }
463 }
464 projected = conversation_projection_rx.recv() => {
465 let Some(projected) = projected else {
466 failure = feed_stopped(
467 termination.is_cancelled(),
468 "the browser transcript projection feed stopped",
469 );
470 break;
471 };
472 let session_id = projected.session_id.clone();
473 let session_active = controller
474 .state
475 .sessions
476 .get(&session_id)
477 .is_some_and(|session| session.state.is_active());
478 if !session_active {
479 conversation_projections.forget(&session_id);
484 }
485 if let Some((session_id, _key, transcript)) =
486 conversation_projections.finish(projected, session_active)
487 {
488 conversations.insert(session_id, transcript);
489 revision = daemon_runtime.allocate_revision();
490 conversation_tx.send_replace(conversations.clone());
491 publish_snapshot!(revision);
492 } else if !session_active && conversations.remove(&session_id).is_some() {
493 revision = daemon_runtime.allocate_revision();
498 conversation_tx.send_replace(conversations.clone());
499 publish_snapshot!(revision);
500 }
501 tokio::task::yield_now().await;
505 }
506 update = worker_updates_rx.recv() => {
507 native_agents_dirty = true;
508 let Some(update) = update else {
509 failure = feed_stopped(termination.is_cancelled(), "the session manager stopped; the phone server can no longer follow sessions");
510 break;
511 };
512 if let Some(snapshot) = update.view.snapshot.as_ref()
513 && let Some(session) = controller.state.sessions.get(&update.session_id)
514 && let Some(signal) = snapshot.latest_credential_sync_signal.clone()
515 {
516 credential_sync_signals.observe(
517 &update.session_id,
518 &session.last_profile,
519 signal,
520 );
521 }
522 schedule_due_credential_syncs(
523 &mut credential_sync_signals,
524 &credential_sync_handle,
525 Instant::now(),
526 );
527 apply_worker_record_update(&mut controller, &update);
528 if let Some(snapshot) = update.view.snapshot {
529 for request in snapshot.subagent_requests.iter().cloned() {
530 let identity = (update.session_id.clone(), request.request_id.clone());
531 if !active_subagent_requests.insert(identity.clone()) {
532 continue;
533 }
534 let backend = api_backend.clone();
535 let runtime = daemon_runtime.clone();
536 let parent_session_id = update.session_id.clone();
537 subagent_jobs.spawn(async move {
538 let result = backend
539 .execute_subagent_tool(parent_session_id.clone(), request)
540 .await;
541 let outcome = async {
542 let handle = runtime
547 .workspace_session_handle(&parent_session_id)
548 .await?;
549 let mut lease = handle.lease_connection().await?;
550 lease
551 .connection_mut()
552 .complete_subagent_request(result)
553 .await?;
554 lease.release();
555 anyhow::Ok(())
556 }
557 .await;
558 (identity, outcome)
559 });
560 }
561 if let Some(relation) = controller.state.subagents.get_mut(&update.session_id)
562 && matches!(snapshot.materialized.execution, mj_core::state::MaterializedExecutionState::Idle)
563 && let Some(outcome) = snapshot.materialized.last_turn_outcome.as_ref()
564 && relation.noticed_turn != Some(outcome.completed_ordinal)
565 {
566 let turn = outcome.completed_ordinal;
567 let child_id = relation.child_session_id.clone();
568 let parent_id = relation.parent_session_id.clone();
569 let task_name = relation.task_name.clone();
570 let outcome_name = format!("{:?}", outcome.outcome).to_lowercase();
571 relation.noticed_turn = Some(turn);
572 let backend = api_backend.clone();
573 subagent_completion_jobs.spawn(async move {
574 let result = async {
575 backend
576 .record_subagent_completion_notice(
577 parent_id,
578 &child_id,
579 &task_name,
580 turn,
581 &outcome_name,
582 )
583 .await?;
584 tokio::task::spawn_blocking({
585 let child_id = child_id.clone();
586 move || crate::database::mark_subagent_turn_noticed(&child_id, turn)
587 })
588 .await??;
589 anyhow::Ok(())
590 }
591 .await;
592 (child_id, turn, result)
593 });
594 }
595 if snapshot.operational.native_session_is_ready()
596 && operational.get(&update.session_id).is_none_or(|old: &mj_core::relay::RelayOperationalState| old.config_options != snapshot.operational.config_options)
597 && let Some(session) = controller.state.sessions.get(&update.session_id)
598 && matches!(session.target, Some(mj_core::state::TargetLocator::LocalBare { .. } | mj_core::state::TargetLocator::SshBare { .. } | mj_core::state::TargetLocator::AwsEc2 { .. }))
599 && let Some(build) = snapshot.worker_build.clone()
600 {
601 let profile = session.last_profile.clone();
602 let state = snapshot.operational.clone();
603 tokio::spawn(async move {
604 if let Err(error) = crate::controller::profile_config::observe(profile, build, state).await {
605 tracing::warn!(%error, "could not cache observed profile choices");
606 }
607 });
608 }
609 let materialized = snapshot.materialized;
610 let operational_state = snapshot.operational;
611 materialized_activity.insert(
612 update.session_id.clone(),
613 crate::server_runtime::snapshot::MaterializedActivity {
614 last_activity_at_ms: materialized.last_activity_at_ms,
615 execution: materialized.execution,
616 },
617 );
618 let queued = queued_prompt_projection(&materialized);
619 let pending = materialized.pending_elicitations.clone();
620 let active_shells = operational_state.active_user_shells.clone();
621 let prompt_images_supported =
622 agent_accepts_prompt_images(&operational_state);
623 active_user_shells.insert(
624 update.session_id.clone(),
625 active_shells,
626 );
627 conversation_projections.enqueue(materialized);
628 queued_prompts.insert(
629 update.session_id.clone(),
630 queued,
631 );
632 pending_elicitations.insert(
633 update.session_id.clone(),
634 pending,
635 );
636 if prompt_images_supported {
637 prompt_images.insert(update.session_id.clone());
638 } else {
639 prompt_images.remove(&update.session_id);
640 }
641 operational.insert(
642 update.session_id.clone(),
643 operational_state,
644 );
645 revision = daemon_runtime.allocate_revision();
646 conversation_tx.send_replace(conversations.clone());
647 publish_snapshot!(revision);
648 }
649 tokio::task::yield_now().await;
654 }
655 completed = subagent_jobs.join_next(), if !subagent_jobs.is_empty() => {
656 match completed {
657 Some(Ok((identity, Ok(())))) => {
658 active_subagent_requests.remove(&identity);
659 }
660 Some(Ok((identity, Err(error)))) => {
661 active_subagent_requests.remove(&identity);
662 tracing::warn!(
663 parent_session_id = %identity.0,
664 request_id = %identity.1,
665 error = %format!("{error:#}"),
666 "sub-agent tool request failed"
667 );
668 }
669 Some(Err(error)) => tracing::warn!(%error, "sub-agent tool task panicked"),
670 None => {}
671 }
672 }
673 completed = subagent_completion_jobs.join_next(), if !subagent_completion_jobs.is_empty() => {
674 match completed {
675 Some(Ok((_, _, Ok(())))) => {}
676 Some(Ok((child_id, turn, Err(error)))) => {
677 if let Some(relation) = controller.state.subagents.get_mut(&child_id)
678 && relation.noticed_turn == Some(turn)
679 {
680 relation.noticed_turn = None;
681 }
682 tracing::warn!(%child_id, turn, error = %format!("{error:#}"), "could not record the sub-agent completion notice");
683 }
684 Some(Err(error)) => tracing::warn!(%error, "sub-agent completion task panicked"),
685 None => {}
686 }
687 }
688 _ = prune_tick.tick() => {
689 let archive_after_days = controller.config.sessionwiki.archive_after_days;
700 match archive_after_days {
701 Some(days) if archive_jobs.is_empty() => {
702 let runtime = daemon_runtime.clone();
703 archive_jobs.spawn(async move {
704 runtime.archive_aged_sessions(days).await
705 });
706 }
707 Some(_) => tracing::debug!(
708 "the previous SessionWiki archive pass is still running; skipping this tick"
709 ),
710 None => daemon_runtime.wiki().request_sync(true),
711 }
712 tokio::spawn(async move {
716 let pruned = tokio::task::spawn_blocking(move || {
717 crate::database::prune_phone_client_state(client_state_retention)
718 })
719 .await;
720 match pruned {
721 Ok(Ok(0)) => {}
722 Ok(Ok(rows)) => tracing::debug!(rows, "pruned expired phone viewer state"),
723 Ok(Err(error)) => tracing::warn!(%error, "could not prune phone viewer state"),
724 Err(error) => tracing::warn!(%error, "phone viewer state pruning task failed"),
725 }
726 });
727 }
728 _ = credential_tick.tick() => {
729 schedule_due_credential_syncs(
730 &mut credential_sync_signals,
731 &credential_sync_handle,
732 Instant::now(),
733 );
734 while let Some(result) = credential_sync.try_result() {
735 crate::pollers::log_credential_sync_actions(&result);
736 let harness = controller
737 .config
738 .profiles
739 .get(&result.profile_id)
740 .map(|profile| profile.kind);
741 if let Some(notice) = credential_sync_notices.notice(&result, harness) {
742 eprintln!("Mjolnir: {notice}");
743 }
744 }
745 }
746 request = dictation_rx.recv() => {
747 let Some(request) = request else {
748 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering dictation requests");
749 break;
750 };
751 let paths = controller.state.sessions.get(&request.session_id).map(|session| {
752 crate::dictation::auth_paths(&controller.config, &session.last_profile)
753 });
754 dictation_jobs.spawn(crate::dictation::execute(
755 request, paths, termination.clone(),
756 ));
757 }
758 job = dictation_jobs.join_next(), if !dictation_jobs.is_empty() => {
759 if let Some(Err(error)) = job {
760 tracing::warn!(%error, "web dictation task failed");
761 }
762 }
763 job = archive_jobs.join_next(), if !archive_jobs.is_empty() => {
764 match job {
765 Some(Ok(Ok(0))) | None => {}
766 Some(Ok(Ok(archived))) => tracing::debug!(
767 archived, "the SessionWiki archive pass finished"
768 ),
769 Some(Ok(Err(error))) => tracing::warn!(
770 error = %format!("{error:#}"),
771 "the SessionWiki archive pass failed"
772 ),
773 Some(Err(error)) => tracing::warn!(
774 %error, "the SessionWiki archive task panicked"
775 ),
776 }
777 }
778 stored = client_state_rx.recv() => {
779 let Some(stored) = stored else {
780 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering viewer state requests");
781 break;
782 };
783 let workspace_of = |session_id: &str| {
787 controller
788 .state
789 .sessions
790 .get(session_id)
791 .map(|session| session.workspace_id.clone())
792 };
793 let bundle_of = |session_id: &str| {
794 controller
795 .state
796 .sessions
797 .get(session_id)
798 .map(|session| session.bundle_id.clone())
799 };
800 match stored {
801 crate::server::ClientStateRequest::Read { client_id, session_id, reply } => {
802 let workspace = workspace_of(&session_id);
803 tokio::spawn(async move {
804 let answer = tokio::task::spawn_blocking(move || {
805 let workspace = workspace.context("unknown session")?;
806 let state = crate::database::client_session_state(
807 &client_id, &workspace, &session_id,
808 )?;
809 anyhow::Ok(crate::server::ViewerClientState {
810 draft: state.draft,
811 through_event_ordinal: state.through_event_ordinal,
812 })
813 })
814 .await;
815 reply.send(flatten_stored(answer)).ok();
816 });
817 }
818 crate::server::ClientStateRequest::SaveDraft { client_id, session_id, draft, reply } => {
819 let workspace = workspace_of(&session_id);
820 tokio::spawn(async move {
821 let answer = tokio::task::spawn_blocking(move || {
822 let workspace = workspace.context("unknown session")?;
823 crate::database::persist_client_draft(
824 &client_id, &workspace, &session_id, &draft,
825 )
826 })
827 .await;
828 reply.send(flatten_stored(answer)).ok();
829 });
830 }
831 crate::server::ClientStateRequest::MarkWorkspaceRead { client_id, workspace_id, reply } => {
832 let sessions = controller
833 .state
834 .sessions
835 .values()
836 .filter(|session| session.workspace_id == workspace_id)
837 .map(|session| (session.id.clone(), session.viewed_through_event_ordinal))
838 .collect::<Vec<_>>();
839 tokio::spawn(async move {
840 let answer = tokio::task::spawn_blocking(move || {
841 for (session_id, through) in sessions {
842 crate::database::persist_read_receipt(
846 &client_id, &workspace_id, &session_id, through,
847 )
848 .ok();
849 }
850 anyhow::Ok(())
851 })
852 .await;
853 reply.send(flatten_stored(answer)).ok();
854 });
855 }
856 crate::server::ClientStateRequest::History { session_id, query, scope, reply } => {
857 let bundle = bundle_of(&session_id);
858 tokio::spawn(async move {
859 let answer = tokio::task::spawn_blocking(move || {
860 let bundle = bundle.context("unknown session")?;
861 let scope = match scope.as_str() {
862 "session" => crate::database::HistoryScope::Session,
863 "all" => crate::database::HistoryScope::All,
864 _ => crate::database::HistoryScope::Project,
865 };
866 let found = crate::database::search_prompts_bounded(
867 &session_id,
868 &bundle,
869 scope,
870 &query,
871 crate::server::MAX_HISTORY_MATCHES,
872 )?;
873 anyhow::Ok(crate::server::ViewerPromptHistory {
874 entries: found
875 .entries
876 .into_iter()
877 .map(|entry| entry.text)
878 .collect(),
879 truncated: found.truncated,
880 })
881 })
882 .await;
883 reply.send(flatten_stored(answer)).ok();
884 });
885 }
886 }
887 }
888 bundle = bundle_rx.recv(), if bundle_jobs.len() < MAX_CONCURRENT_BUNDLE_CREATIONS => {
889 let Some(crate::server::BundleRequest { source, reply }) = bundle else {
890 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering bundle requests");
891 break;
892 };
893 let done = bundle_done_tx.clone();
898 let daemon_runtime = daemon_runtime.clone();
899 bundle_jobs.spawn(async move {
900 let result = daemon_runtime
901 .create_quick_bundle(source)
902 .await;
903 if let Err(error) = done.send(BundleCreated { result, reply }) {
904 tracing::debug!(%error, "bundle creation finished after the server stopped");
905 }
906 });
907 }
908 bundle_done = bundle_done_rx.recv() => {
909 let Some(BundleCreated { result, reply }) = bundle_done else {
910 failure = feed_stopped(termination.is_cancelled(), "the bundle creation pipeline stopped while the phone server was running");
911 break;
912 };
913 match result {
914 Ok(created) => {
915 let bundle_id = created.bundle_id;
916 let Some(bundle) = created.config.bundles.get(&bundle_id) else {
922 tracing::error!(%bundle_id, "bundle creation returned a config without its bundle");
923 if reply.send(Err(crate::server::BundleFailure::Controller)).is_err() {
924 tracing::debug!("bundle creation failure reply dropped after client disconnect");
925 }
926 continue;
927 };
928 controller
929 .config
930 .bundles
931 .insert(bundle_id.clone(), bundle.clone());
932 controller_reload_invalidated |= controller_reload_in_flight;
936 revision = daemon_runtime.allocate_revision();
937 publish_snapshot!(revision);
938 request_daemon_controller_reload(
939 daemon_runtime.clone(),
940 "new bundle publication",
941 );
942 if reply.send(Ok(bundle_id)).is_err() {
943 tracing::debug!("bundle creation reply dropped after client disconnect");
944 }
945 }
946 Err(error) => {
947 let failure = match error {
948 crate::controller::QuickBundleFailure::InvalidSource(
949 detail,
950 ) => {
951 tracing::debug!(error = %detail, "phone bundle source was invalid");
952 crate::server::BundleFailure::InvalidSource
953 }
954 crate::controller::QuickBundleFailure::Persistence(
955 detail,
956 ) => {
957 tracing::warn!(error = %detail, "phone bundle creation failed");
958 crate::server::BundleFailure::Controller
959 }
960 };
961 if reply.send(Err(failure)).is_err() {
962 tracing::debug!("bundle creation failure reply dropped after client disconnect");
963 }
964 }
965 }
966 }
967 bundle_job = bundle_jobs.join_next(), if !bundle_jobs.is_empty() => {
968 if let Some(Err(error)) = bundle_job {
969 tracing::warn!(%error, "bundle creation task panicked");
970 }
971 }
972 preflight = preflight_rx.recv(), if preflight_jobs.len() < MAX_CONCURRENT_PREFLIGHTS => {
973 let Some(preflight) = preflight else {
974 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering preflight requests");
975 break;
976 };
977 let crate::server::NewPreflightRequest {
981 bundle_id,
982 target_id,
983 project_directory,
984 mut reply,
985 remote_repairs,
986 } = match preflight {
987 crate::server::PreflightRequest::New(request) => request,
988 crate::server::PreflightRequest::Resume(request) => {
989 spawn_resume_preflight(
990 &mut preflight_jobs,
991 &controller,
992 request,
993 &termination,
994 );
995 continue;
996 }
997 crate::server::PreflightRequest::CompletePath(request) => {
998 spawn_path_completion(
999 &mut preflight_jobs,
1000 &controller.config,
1001 request,
1002 &termination,
1003 );
1004 continue;
1005 }
1006 };
1007 let config = controller.config.clone();
1012 let project_validation = project_directory.is_some();
1013 let task_termination = termination.clone();
1014 preflight_jobs.spawn(async move {
1015 let cancelled = Arc::new(AtomicBool::new(false));
1016 let cancellation_guard = ProcessCancellationGuard(cancelled.clone());
1017 let mut blocking = tokio::task::spawn_blocking(move || {
1018 run_new_preflight_with_cancellation(
1019 config,
1020 bundle_id,
1021 target_id,
1022 project_directory,
1023 cancelled,
1024 remote_repairs,
1025 )
1026 });
1027 let answer = tokio::select! {
1028 biased;
1029 _ = task_termination.cancelled() => None,
1030 _ = reply.closed() => None,
1031 answer = &mut blocking => Some(answer),
1032 };
1033 let Some(answer) = answer else {
1034 drop(cancellation_guard);
1035 match blocking.await {
1036 Err(error) => tracing::warn!(%error, "cancelled phone preflight task failed"),
1037 Ok(Err(error)) => tracing::debug!(%error, "phone preflight cancelled"),
1038 Ok(Ok(_)) => {}
1039 }
1040 return;
1041 };
1042 let answer = match answer {
1043 Ok(Ok(answer)) => Ok(answer),
1044 Ok(Err(error)) => {
1045 tracing::debug!(
1046 error = %error,
1047 project_validation,
1048 "phone preflight check failed"
1049 );
1050 Err(if project_validation {
1051 PreflightFailure::Validation
1052 } else {
1053 PreflightFailure::InvalidRepository(format!("{error:#}"))
1054 })
1055 }
1056 Err(error) => {
1057 tracing::warn!(%error, "phone preflight task failed");
1058 Err(PreflightFailure::Controller(format!(
1059 "preflight task failed: {error}"
1060 )))
1061 }
1062 };
1063 if reply.send(answer).is_err() {
1064 tracing::debug!("phone preflight reply dropped after client disconnect");
1065 }
1066 });
1067 }
1068 preflight_job = preflight_jobs.join_next(), if !preflight_jobs.is_empty() => {
1069 if let Some(Err(error)) = preflight_job {
1070 tracing::warn!(%error, "phone preflight task panicked");
1071 }
1072 }
1073 preparation = move_preparation_rx.recv() => {
1074 let Some(MovePreparationRequest { selection, reply }) = preparation else {
1075 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering move preparation requests");
1076 break;
1077 };
1078 let done = move_prepared_tx.clone();
1083 let daemon_runtime = daemon_runtime.clone();
1084 move_preparation_jobs.spawn(async move {
1085 let result = daemon_runtime
1086 .prepare_move_session(selection)
1087 .await
1088 .map_err(|error| format!("{error:#}"));
1089 if let Err(error) = done.send(MovePrepared { result, reply }) {
1090 tracing::debug!(%error, "move preparation finished after the server stopped");
1091 }
1092 });
1093 }
1094 prepared = move_prepared_rx.recv() => {
1095 let Some(MovePrepared { result, reply }) = prepared else {
1096 failure = feed_stopped(termination.is_cancelled(), "the move preparation pipeline stopped while the phone server was running");
1097 break;
1098 };
1099 if reply.send(result).is_err() {
1100 tracing::debug!("move preparation reply dropped after client disconnect");
1101 }
1102 }
1103 move_preparation_job = move_preparation_jobs.join_next(), if !move_preparation_jobs.is_empty() => {
1104 if let Some(Err(error)) = move_preparation_job {
1105 tracing::warn!(%error, "move preparation task failed");
1106 }
1107 }
1108 result = native_agent_jobs.join_next(), if !native_agent_jobs.is_empty() => {
1109 match result {
1110 Some(Ok(Ok(agents))) => {
1111 native_agents = agents;
1112 revision = daemon_runtime.allocate_revision();
1113 publish_snapshot!(revision);
1114 }
1115 Some(Ok(Err(error))) => tracing::error!(%error, "could not refresh native subagent identities"),
1116 Some(Err(error)) => tracing::error!(%error, "native subagent refresh task failed"),
1117 None => {}
1118 }
1119 }
1120 move_recovery_job = move_recovery_jobs.join_next(), if !move_recovery_jobs.is_empty() => {
1121 if let Some(Err(error)) = move_recovery_job {
1122 move_recovery_load_in_flight = false;
1123 tracing::warn!(%error, "Move recovery projection task failed");
1124 }
1125 }
1126 receipt = receipt_rx.recv() => {
1127 let Some(ReadReceiptRequest { client_id, session_id, through, reply }) = receipt else {
1128 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering read receipts");
1129 break;
1130 };
1131 match controller.state.sessions.get(&session_id) {
1132 None => {
1133 if reply.send(Err("unknown session".into())).is_err() {
1134 tracing::debug!(%session_id, "unknown-session read receipt reply dropped after client disconnect");
1135 }
1136 }
1137 Some(session) => {
1138 let workspace_id = session.workspace_id.clone();
1139 let done = receipt_done_tx.clone();
1140 let persisted_session_id = session_id.clone();
1141 tokio::spawn(async move {
1142 let joined = tokio::task::spawn_blocking(move || {
1143 crate::database::persist_read_receipt(
1144 &client_id,
1145 &workspace_id,
1146 &persisted_session_id,
1147 through,
1148 )
1149 })
1150 .await;
1151 let result = match joined {
1152 Ok(result) => result.map_err(|error| format!("{error:#}")),
1153 Err(error) => Err(format!("phone read receipt task failed: {error}")),
1154 };
1155 if let Err(error) = done.send(ReadReceiptPersisted { session_id, result, reply }) {
1156 tracing::debug!(%error, "phone read receipt finished after the server stopped");
1157 }
1158 });
1159 }
1160 }
1161 }
1162 persisted = receipt_done_rx.recv() => {
1163 let Some(ReadReceiptPersisted { session_id, result, reply }) = persisted else { continue };
1164 match result {
1165 Ok(receipt) => {
1166 let _ = receipt;
1167 if reply.send(Ok(())).is_err() {
1168 tracing::debug!(%session_id, "phone read receipt reply dropped after client disconnect");
1169 }
1170 }
1171 Err(error) => {
1172 tracing::warn!(%session_id, "could not persist a phone read receipt: {error}");
1173 if reply.send(Err(error)).is_err() {
1174 tracing::debug!(%session_id, "failed phone read receipt reply dropped after client disconnect");
1175 }
1176 }
1177 }
1178 }
1179 action = action_rx.recv() => {
1180 let Some(request) = action else {
1181 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering actions");
1182 break;
1183 };
1184 match &request.action {
1188 ControllerAction::RefreshCapacity { target_id } => {
1189 let known = capacity_state.contains_key(target_id);
1190 let accepted = known && match capacity_triggers_tx.try_send(()) {
1193 Ok(()) | Err(tokio::sync::mpsc::error::TrySendError::Full(())) => true,
1194 Err(tokio::sync::mpsc::error::TrySendError::Closed(())) => {
1195 tracing::warn!("phone capacity refresh rejected: poller stopped");
1196 false
1197 }
1198 };
1199 if known {
1200 if let Some(entry) = capacity_state.get_mut(target_id) {
1201 entry.refreshing = accepted;
1202 if !accepted {
1203 entry.failed = true;
1204 }
1205 }
1206 revision = daemon_runtime.allocate_revision();
1207 publish_snapshot!(revision);
1208 }
1209 let outcome = match (known, accepted) {
1210 (_, true) => ActionOutcome::accepted(),
1211 (false, _) => ActionOutcome::Refused(Refusal::unusable(format!(
1212 "no target named {target_id} is configured"
1213 ))),
1214 (true, false) => ActionOutcome::Failed {
1215 reference: "capacity-poller".to_owned(),
1216 },
1217 };
1218 if request.reply.send(outcome).is_err() {
1219 tracing::debug!(%target_id, "phone capacity refresh reply dropped after client disconnect");
1220 }
1221 tokio::task::yield_now().await;
1222 continue;
1223 }
1224 ControllerAction::RefreshQuota { profile_id } => {
1225 let known = controller.config.enabled_profile(profile_id).is_some();
1226 if known {
1227 quota_batch.generation = quota_batch.generation.saturating_add(1);
1231 quota_batch.profiles = quota_refresh_profiles(&controller);
1232 quota_profiles_tx.send_replace(quota_batch.clone());
1233 }
1234 let outcome = if known {
1235 ActionOutcome::accepted()
1236 } else {
1237 ActionOutcome::Refused(Refusal::unusable(format!(
1238 "no enabled profile named {profile_id} is configured"
1239 )))
1240 };
1241 if request.reply.send(outcome).is_err() {
1242 tracing::debug!(%profile_id, "phone quota refresh reply dropped after client disconnect");
1243 }
1244 tokio::task::yield_now().await;
1245 continue;
1246 }
1247 _ => {}
1248 }
1249 if let ControllerAction::Cancel { session_id } = &request.action {
1250 let outcome = if request_phone_action_cancellation(
1251 session_id,
1252 &action_sessions,
1253 &action_cancellations,
1254 ) {
1255 daemon_runtime.cancel_lifecycle_if_active(session_id);
1256 ActionOutcome::accepted()
1257 } else {
1258 ActionOutcome::NotCancellable
1259 };
1260 if request.reply.send(outcome).is_err() {
1261 tracing::debug!(%session_id, "phone cancellation reply dropped after client disconnect");
1262 }
1263 tokio::task::yield_now().await;
1264 continue;
1265 }
1266 if let ControllerAction::Suspend { session_id } = &request.action {
1267 if !closing_actions.contains_key(session_id) {
1268 request_phone_action_cancellation(session_id, &action_sessions, &action_cancellations);
1269 }
1270 daemon_runtime.request_close(session_id);
1271 }
1272 if let ControllerAction::Destroy { session_id, .. } = &request.action {
1276 request_phone_action_cancellation(session_id, &action_sessions, &action_cancellations);
1277 daemon_runtime.request_close(session_id);
1278 }
1279 let session_id = match admit_phone_action(
1280 &request.action,
1281 action_cancellations.len(),
1282 &mut active_actions,
1283 ) {
1284 Ok(session_id) => session_id,
1285 Err(refusal) => {
1286 if request.reply.send(refusal).is_err() {
1287 tracing::debug!("phone action refusal reply dropped after client disconnect");
1288 }
1289 tokio::task::yield_now().await;
1290 continue;
1291 }
1292 };
1293 let ControllerRequest { action, reply } = request;
1294 let done = action_done_tx.clone();
1295 let session_control = worker_commands_tx.clone();
1296 let daemon_runtime = daemon_runtime.clone();
1297 let started = action_started_tx.clone();
1298 next_action_id = next_action_id.wrapping_add(1).max(1);
1299 let action_id = next_action_id;
1300 if let ControllerAction::Suspend { session_id } | ControllerAction::Destroy { session_id, .. } = &action { closing_actions.insert(session_id.clone(), action_id); }
1301 if let ControllerAction::New { workspace_id, .. } = &action {
1302 let workspace_id = if workspace_id.is_empty() && phone_workspaces.len() == 1 {
1303 phone_workspaces[0].id.clone()
1304 } else {
1305 workspace_id.clone()
1306 };
1307 launch_workspaces.insert(action_id, workspace_id);
1308 }
1309 let control = PhoneActionControl::for_action(&action);
1310 action_cancellations.insert(action_id, control.clone());
1311 if let Some(session_id) = &session_id {
1312 action_sessions.insert(action_id, session_id.clone());
1313 }
1314 let suspension_reply = if matches!(action, ControllerAction::Suspend { .. }) {
1315 Some(reply)
1316 } else {
1317 action_replies.accept(action_id, &action, reply);
1318 None
1319 };
1320 let notice_action = match &action {
1321 ControllerAction::SetConfig { .. } => Some("Configuration change"),
1322 ControllerAction::InterruptTurn { .. } => Some("Cancellation"),
1323 ControllerAction::CancelShell { .. } => Some("Shell cancellation"),
1324 ControllerAction::RemoveQueuedPrompt { .. } => Some("Queued prompt removal"),
1325 ControllerAction::RespondElicitation { .. } => Some("Answer"),
1326 _ => None,
1327 };
1328 let notice_sessions = session_control.clone();
1329 let failure_runtime = daemon_runtime.clone();
1330 let lifecycle_failure_prefix = match &action {
1331 ControllerAction::Suspend { .. } => Some(mj_core::state::CLOSE_FAILURE_PREFIX),
1332 ControllerAction::Destroy { .. } => Some(mj_core::state::DESTRUCTION_FAILURE_PREFIX),
1333 _ => None,
1334 };
1335 tokio::spawn(async move {
1336 let joined = tokio::task::spawn_blocking(move || {
1337 let mut suspension_reply = suspension_reply;
1338 let result = (|| -> Result<()> {
1339 if let ControllerAction::Suspend { session_id } = &action {
1340 mj_core::runtime::block_on(daemon_runtime.prepare_suspension(session_id))??;
1341 if let Some(reply) = suspension_reply.take() {
1342 let _ = reply.send(ActionOutcome::accepted());
1343 }
1344 }
1345 if control.cancelled.load(Ordering::Acquire) {
1346 bail!("phone action cancelled");
1347 }
1348 let mut operation_controller = Controller::load()?;
1349 let executor =
1350 CancellableProcessExecutor::new(control.cancelled.clone());
1351 mj_core::runtime::block_on(apply_phone_action(
1352 &mut operation_controller,
1353 PhoneActionServices {
1354 sessions: &session_control,
1355 daemon_runtime: &daemon_runtime,
1356 },
1357 action,
1358 &executor,
1359 action_id,
1360 &started,
1361 &control,
1362 ))?
1363 })();
1364 let result = result.map_err(|error| PhoneActionFailure::of(&error));
1365 if let Some(reply) = suspension_reply {
1366 let outcome = match &result {
1367 Ok(()) => ActionOutcome::accepted(),
1368 Err(failure) => failure.outcome(&action_reference(action_id)),
1369 };
1370 let _ = reply.send(outcome);
1371 }
1372 result
1373 })
1374 .await;
1375 let result = match joined {
1376 Ok(result) => result,
1377 Err(error) => Err(PhoneActionFailure::internal(format!(
1378 "phone action task failed: {error}"
1379 ))),
1380 };
1381 if let (Err(failure), Some(prefix), Some(id)) = (&result, lifecycle_failure_prefix, &session_id) {
1382 failure_runtime.record_lifecycle_failure(id, &action_reference(action_id), &crate::daemon::LifecycleFailure {
1383 detail: failure.detail.clone(), refusal: failure.refusal.clone(),
1384 }, prefix).await;
1385 }
1386 let notice = match (&result, notice_action, &session_id) {
1387 (Err(failure), Some(action), Some(session_id)) => {
1388 Some((session_id.clone(), failure.conversation_notice(action, action_id)))
1389 }
1390 _ => None,
1391 };
1392 if let Err(error) = done.send((action_id, session_id, result)) {
1393 tracing::debug!(action_id, %error, "phone action finished after the server stopped");
1394 }
1395 if let Some((session_id, text)) = notice {
1396 let recorded = async {
1397 notice_sessions.session(&session_id).await?
1398 .submit(new_command_id("action-notice")?, RelayCommand::RecordNotice { text }).await?;
1399 Ok::<_, anyhow::Error>(())
1400 }.await;
1401 if let Err(error) = recorded {
1402 tracing::warn!(%session_id, %error, "could not record failed session action in conversation");
1403 }
1404 }
1405 });
1406 }
1407 started = action_started_rx.recv() => {
1408 let Some(started) = started else {
1409 tokio::task::yield_now().await;
1410 continue;
1411 };
1412 let started_session_id = started.session.id.clone();
1413 let publication = if !action_cancellations.contains_key(&started.action_id) {
1414 Err("phone action completed before its provisional session was published".into())
1415 } else {
1416 track_started_phone_session(
1417 &mut controller.state,
1418 &mut active_actions,
1419 &mut action_sessions,
1420 started.action_id,
1421 started.session,
1422 )
1423 };
1424 if publication.is_ok() {
1425 revision = daemon_runtime.allocate_revision();
1426 publish_snapshot!(revision);
1427 request_daemon_controller_reload(
1428 daemon_runtime.clone(),
1429 "new session publication",
1430 );
1431 };
1432 if publication.is_err()
1433 && let Some(control) = action_cancellations.get(&started.action_id)
1434 {
1435 control.request_cancel();
1436 }
1437 action_replies.resolve(
1440 started.action_id,
1441 match &publication {
1442 Ok(()) => ActionOutcome::Accepted {
1443 session_id: Some(started_session_id),
1444 },
1445 Err(reason) => ActionOutcome::Refused(Refusal::precondition(
1450 reason.clone(),
1451 )),
1452 },
1453 );
1454 if started.published.send(publication).is_err() {
1455 tracing::debug!(action_id = started.action_id, "phone new-session publication reply dropped after client disconnect");
1456 }
1457 }
1458 completed = action_done_rx.recv() => {
1459 let Some((action_id, session_id, result)) = completed else {
1460 failure = feed_stopped(termination.is_cancelled(), "the phone action pipeline stopped reporting completions");
1461 break;
1462 };
1463 action_cancellations.remove(&action_id);
1464 let session_id = action_sessions.remove(&action_id).or(session_id);
1465 if closing_actions.values().any(|closing_id| *closing_id == action_id) && let Some(id) = &session_id { daemon_runtime.clear_close_request(id); }
1466 closing_actions.retain(|_, closing_id| *closing_id != action_id);
1467 if let Some(session_id) = &session_id && !action_sessions.values().any(|active| active == session_id) {
1468 active_actions.remove(session_id);
1469 }
1470 let reference = action_reference(action_id);
1477 action_replies.resolve(
1478 action_id,
1479 match &result {
1480 Ok(()) => ActionOutcome::Accepted { session_id: session_id.clone() },
1481 Err(failure) => failure.outcome(&reference),
1482 },
1483 );
1484 if let Some(workspace_id) = launch_workspaces.remove(&action_id)
1485 && result.is_err()
1486 && !session_id.as_ref().is_some_and(|id| closing_actions.contains_key(id))
1487 {
1488 record_launch_failure(
1489 &mut launch_failures,
1490 action_id,
1491 workspace_id,
1492 session_id.clone(),
1493 result.as_ref().err().map(|failure| failure.detail.clone()),
1494 );
1495 revision = daemon_runtime.allocate_revision();
1496 publish_snapshot!(revision);
1497 }
1498 if let Err(failure) = &result {
1499 tracing::warn!(
1500 action_id,
1501 reference,
1502 session_id = session_id.as_deref(),
1503 error = %failure.detail,
1504 "phone action failed"
1505 );
1506 }
1507 record_action_result(
1508 &mut pending_action_errors,
1509 session_id.as_deref(),
1510 &result,
1511 );
1512 if result.is_ok() && let Some(session_id) = session_id.clone() {
1515 let daemon_runtime = daemon_runtime.clone();
1516 tokio::spawn(async move {
1517 daemon_runtime.clear_recorded_close_failure(&session_id).await;
1518 });
1519 }
1520 request_controller_reload(
1521 &mut controller_reload_in_flight,
1522 &mut controller_reload_requested,
1523 &controller_reload_tx,
1524 );
1525 request_move_recovery_reload(
1526 &move_recovery_tx,
1527 &mut move_recovery_load_in_flight,
1528 &mut move_recovery_jobs,
1529 );
1530 request_daemon_controller_reload(
1531 daemon_runtime.clone(),
1532 "phone action completion",
1533 );
1534 }
1535 reloaded = controller_reload_rx.recv() => {
1536 let Some(ControllerReloaded { result }) = reloaded else {
1537 failure = feed_stopped(
1538 termination.is_cancelled(),
1539 "the controller reload pipeline stopped while the phone server was running",
1540 );
1541 break;
1542 };
1543 controller_reload_in_flight = false;
1544 if std::mem::take(&mut controller_reload_invalidated) {
1545 if let Err(error) = &result {
1546 tracing::warn!(%error, "superseded controller reload failed");
1547 }
1548 controller_reload_requested = false;
1549 request_controller_reload(
1550 &mut controller_reload_in_flight,
1551 &mut controller_reload_requested,
1552 &controller_reload_tx,
1553 );
1554 continue;
1555 }
1556 match result {
1557 Ok(mut reloaded) => {
1558 for (session_id, error) in &pending_action_errors {
1559 if let Some(session) = reloaded.state.sessions.get_mut(session_id)
1560 && session.last_error.is_none()
1561 {
1562 session.last_error = Some(error.clone());
1563 }
1564 }
1565 controller = reloaded;
1566 quotas.retain(|id, _| controller.config.enabled_profile(id).is_some());
1567 subagent_quota_reports
1568 .lock()
1569 .expect("sub-agent quota reports lock poisoned")
1570 .retain(|id, _| controller.config.enabled_profile(id).is_some());
1571 worker_targets_tx.send_replace(dashboard_worker_targets(&controller));
1572 publish_capacity_targets(
1573 &controller,
1574 &capacity_targets_tx,
1575 &mut capacity_state,
1576 );
1577 credential_sync_handle.set_targets(credential_sync_targets(&controller));
1578 republish_quota_profiles(
1579 &controller,
1580 &mut published_quota_profiles,
1581 &mut quota_batch,
1582 "a_profiles_tx,
1583 );
1584 profile_catalog.sync(&controller.config);
1592 queued_prompts.retain(|session_id, _| {
1593 controller.state.sessions.contains_key(session_id)
1594 });
1595 pending_elicitations.retain(|session_id, _| {
1596 controller.state.sessions.contains_key(session_id)
1597 });
1598 prompt_images.retain(|session_id| {
1599 controller.state.sessions.contains_key(session_id)
1600 });
1601 operational.retain(|session_id, _| {
1602 controller.state.sessions.contains_key(session_id)
1603 });
1604 materialized_activity.retain(|session_id, _| {
1605 controller.state.sessions.contains_key(session_id)
1606 });
1607 request_move_recovery_reload(
1608 &move_recovery_tx,
1609 &mut move_recovery_load_in_flight,
1610 &mut move_recovery_jobs,
1611 );
1612 conversations.retain(|id, _| {
1613 controller.state.sessions.get(id).is_some_and(|session| session.state.is_active())
1614 });
1615 for session_id in conversation_projections.session_ids() {
1616 if !controller
1617 .state
1618 .sessions
1619 .get(&session_id)
1620 .is_some_and(|session| session.state.is_active())
1621 {
1622 conversation_projections.forget(&session_id);
1623 }
1624 }
1625 revision = daemon_runtime.allocate_revision();
1626 conversation_tx.send_replace(conversations.clone());
1627 publish_snapshot!(revision);
1628 }
1629 Err(error) => {
1630 tracing::warn!(%error, "completed phone operation could not reload controller state");
1631 }
1632 }
1633 if controller_reload_requested {
1634 controller_reload_requested = false;
1635 controller_reload_in_flight = true;
1636 spawn_controller_reload(controller_reload_tx.clone());
1637 }
1638 }
1639 }
1640 }
1641 dictation_jobs.shutdown().await;
1643 bundle_jobs.shutdown().await;
1646 preflight_jobs.shutdown().await;
1650 move_preparation_jobs.shutdown().await;
1653 native_agent_jobs.shutdown().await;
1654 move_recovery_jobs.shutdown().await;
1655 crate::controller::profile_config::cancel_all();
1657 for control in action_cancellations.values() {
1658 control.request_cancel();
1659 }
1660 match failure {
1661 Some(failure) => Err(failure),
1662 None => Ok::<(), anyhow::Error>(()),
1663 }
1664 };
1665 let recorder = tokio::spawn(api_activity::record_activity_stream(
1669 activity_snapshots,
1670 crate::database::record_api_activities,
1671 ));
1672 let result = tokio::select! {
1673 result = serve.stopped() => result,
1674 result = control => result,
1675 };
1676 recorder.abort();
1677 if let Err(error) = recorder.await
1678 && error.is_panic()
1679 {
1680 tracing::warn!(%error, "native API activity recorder panicked");
1681 }
1682 drop(serve);
1685 conversation_projection_shutdown.cancel();
1686 renewal_cancellation.cancel();
1687 if let Some(task) = renewal_task
1688 && let Err(error) = task.await
1689 {
1690 tracing::warn!(%error, "Tailscale certificate renewal task failed");
1691 }
1692 worker_shutdown
1693 .shutdown()
1694 .await
1695 .context("shut down phone server session manager")?;
1696 result?;
1697 Ok(())
1698}
1699
1700pub(crate) struct ViewerServer(tokio::task::JoinHandle<Result<()>>);
1711
1712impl ViewerServer {
1713 pub(crate) fn spawn(
1714 server: impl std::future::Future<Output = Result<()>> + Send + 'static,
1715 ) -> Self {
1716 Self(tokio::spawn(server))
1717 }
1718
1719 pub(crate) async fn stopped(&mut self) -> Result<()> {
1722 match (&mut self.0).await {
1723 Ok(result) => result,
1724 Err(error) => Err(anyhow::Error::new(error).context("the web viewer task failed")),
1725 }
1726 }
1727}
1728
1729impl Drop for ViewerServer {
1730 fn drop(&mut self) {
1731 self.0.abort();
1732 }
1733}