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