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 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
319 background_task_stop_jobs.spawn(async move {
320 let _upgrade_task = upgrade_task;
321 let result = match session_control.session(&request.session_id).await {
322 Ok(session) => session
323 .client()
324 .stop_background_task(request.background_task_id.clone())
325 .await
326 .map_err(|error| {
327 tracing::warn!(
328 session_id = %request.session_id,
329 background_task_id = %request.background_task_id,
330 %error,
331 "provider rejected background-task stop"
332 );
333 BackgroundTaskStopFailure::Provider
334 }),
335 Err(error) => {
336 tracing::warn!(
337 session_id = %request.session_id,
338 %error,
339 "could not resolve live session for background-task stop"
340 );
341 Err(BackgroundTaskStopFailure::SessionUnavailable)
342 }
343 };
344 if request.reply.send(result).is_err() {
345 tracing::debug!(
346 session_id = %request.session_id,
347 background_task_id = %request.background_task_id,
348 "background-task stop result dropped after viewer disconnected"
349 );
350 }
351 });
352 }
353 completed = background_task_stop_jobs.join_next(), if !background_task_stop_jobs.is_empty() => {
354 if let Some(Err(error)) = completed {
355 tracing::error!(%error, "background-task stop task failed unexpectedly");
356 }
357 }
358 move_reloaded = move_recovery_rx.recv() => {
359 let Some(result) = move_reloaded else {
360 failure = feed_stopped(
361 termination.is_cancelled(),
362 "the Move recovery projection stopped while the phone server was running",
363 );
364 break;
365 };
366 move_recovery_load_in_flight = false;
367 match result {
368 Ok(recoveries) => {
369 move_recoveries = recoveries;
370 revision = daemon_runtime.allocate_revision();
371 publish_snapshot!(revision);
372 }
373 Err(error) => tracing::warn!(%error, "could not refresh Move recovery projection"),
374 }
375 }
376 resolved = project_sources.jobs.join_next(), if !project_sources.jobs.is_empty() => {
377 match resolved {
378 Some(Ok(resolved)) => project_sources.complete(resolved),
379 Some(Err(error)) => {
380 failure = Some(anyhow::anyhow!("web project source task failed: {error}"));
381 break;
382 }
383 None => unreachable!("project source jobs were not empty"),
384 }
385 revision = daemon_runtime.allocate_revision();
386 publish_snapshot!(revision);
387 }
388 changed = daemon_revisions.changed() => {
389 native_agents_dirty = true;
390 if changed.is_err() {
391 failure = feed_stopped(
392 termination.is_cancelled(),
393 "the daemon stopped publishing runtime revisions to the phone server",
394 );
395 break;
396 }
397 daemon_revisions.borrow_and_update();
398 revision = daemon_runtime.allocate_revision();
399 publish_snapshot!(revision);
400 request_controller_reload(
401 &mut controller_reload_in_flight,
402 &mut controller_reload_requested,
403 &controller_reload_tx,
404 );
405 }
406 changed = workspace_updates.changed() => {
407 if changed.is_err() {
408 failure = feed_stopped(
409 termination.is_cancelled(),
410 "the daemon stopped publishing workspaces to the phone server",
411 );
412 break;
413 }
414 phone_workspaces = workspace_updates.borrow_and_update().clone();
415 revision = daemon_runtime.allocate_revision();
416 publish_snapshot!(revision);
417 }
418 update = capacity_updates_rx.recv() => {
419 let Some(update) = update else {
420 failure = feed_stopped(termination.is_cancelled(), "the capacity poller stopped while the phone server was running");
421 break;
422 };
423 if let Some(entry) = capacity_state.get_mut(&update.target_id) {
424 entry.refreshing = false;
425 entry.sampled_at_epoch_seconds = Some(update.sampled_at_epoch_seconds);
426 match update.result {
427 Ok(usage) => {
428 entry.on_demand = usage.is_none();
432 entry.usage = usage;
433 entry.failed = false;
434 }
435 Err(_) => entry.failed = true,
439 }
440 }
441 revision = daemon_runtime.allocate_revision();
442 publish_snapshot!(revision);
443 }
444 update = quota_updates_rx.recv(), if quota_updates_open => {
445 match update {
446 Some(QuotaUpdate::Report(outcome)) => {
447 if outcome.credentials_changed {
448 credential_sync_handle
449 .sync_profile_now(&outcome.report.profile_id, None);
450 }
451 quotas.insert(outcome.report.profile_id.clone(), outcome.report.clone());
452 subagent_quota_reports
453 .lock()
454 .expect("sub-agent quota reports lock poisoned")
455 .insert(outcome.report.profile_id.clone(), outcome.report);
456 revision = daemon_runtime.allocate_revision();
457 publish_snapshot!(revision);
458 }
459 Some(QuotaUpdate::Refreshing { .. } | QuotaUpdate::Finished { .. }) => {}
460 None => {
461 quota_updates_open = false;
462 tracing::warn!("quota refresher stopped while the phone server is running");
463 }
464 }
465 }
466 projected = conversation_projection_rx.recv() => {
467 let Some(projected) = projected else {
468 failure = feed_stopped(
469 termination.is_cancelled(),
470 "the browser transcript projection feed stopped",
471 );
472 break;
473 };
474 let session_id = projected.session_id.clone();
475 let session_active = controller
476 .state
477 .sessions
478 .get(&session_id)
479 .is_some_and(|session| session.state.is_active());
480 if !session_active {
481 conversation_projections.forget(&session_id);
486 }
487 if let Some((session_id, _key, transcript)) =
488 conversation_projections.finish(projected, session_active)
489 {
490 conversations.insert(session_id, transcript);
491 revision = daemon_runtime.allocate_revision();
492 conversation_tx.send_replace(conversations.clone());
493 publish_snapshot!(revision);
494 } else if !session_active && conversations.remove(&session_id).is_some() {
495 revision = daemon_runtime.allocate_revision();
500 conversation_tx.send_replace(conversations.clone());
501 publish_snapshot!(revision);
502 }
503 tokio::task::yield_now().await;
507 }
508 update = worker_updates_rx.recv() => {
509 native_agents_dirty = true;
510 let Some(update) = update else {
511 failure = feed_stopped(termination.is_cancelled(), "the session manager stopped; the phone server can no longer follow sessions");
512 break;
513 };
514 if let Some(snapshot) = update.view.snapshot.as_ref()
515 && let Some(session) = controller.state.sessions.get(&update.session_id)
516 && let Some(signal) = snapshot.latest_credential_sync_signal.clone()
517 {
518 credential_sync_signals.observe(
519 &update.session_id,
520 &session.last_profile,
521 signal,
522 );
523 }
524 schedule_due_credential_syncs(
525 &mut credential_sync_signals,
526 &credential_sync_handle,
527 Instant::now(),
528 );
529 apply_worker_record_update(&mut controller, &update);
530 if let Some(snapshot) = update.view.snapshot {
531 for request in snapshot.subagent_requests.iter().cloned() {
532 let identity = (update.session_id.clone(), request.request_id.clone());
533 if !active_subagent_requests.insert(identity.clone()) {
534 continue;
535 }
536 let backend = api_backend.clone();
537 let runtime = daemon_runtime.clone();
538 let parent_session_id = update.session_id.clone();
539 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
540 subagent_jobs.spawn(async move {
541 let _upgrade_task = upgrade_task;
542 let result = backend
543 .execute_subagent_tool(parent_session_id.clone(), request)
544 .await;
545 let outcome = async {
546 let handle = runtime
551 .workspace_session_handle(&parent_session_id)
552 .await?;
553 let mut lease = handle.lease_connection().await?;
554 lease
555 .connection_mut()
556 .complete_subagent_request(result)
557 .await?;
558 lease.release();
559 anyhow::Ok(())
560 }
561 .await;
562 (identity, outcome)
563 });
564 }
565 if let Some(relation) = controller.state.subagents.get_mut(&update.session_id)
566 && matches!(snapshot.materialized.execution, mj_core::state::MaterializedExecutionState::Idle)
567 && let Some(outcome) = snapshot.materialized.last_turn_outcome.as_ref()
568 && relation.noticed_turn != Some(outcome.completed_ordinal)
569 {
570 let turn = outcome.completed_ordinal;
571 let child_id = relation.child_session_id.clone();
572 let parent_id = relation.parent_session_id.clone();
573 let task_name = relation.task_name.clone();
574 let outcome_name = format!("{:?}", outcome.outcome).to_lowercase();
575 relation.noticed_turn = Some(turn);
576 let backend = api_backend.clone();
577 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
578 subagent_completion_jobs.spawn(async move {
579 let _upgrade_task = upgrade_task;
580 let result = async {
581 backend
582 .record_subagent_completion_notice(
583 parent_id,
584 &child_id,
585 &task_name,
586 turn,
587 &outcome_name,
588 )
589 .await?;
590 tokio::task::spawn_blocking({
591 let child_id = child_id.clone();
592 move || crate::database::mark_subagent_turn_noticed(&child_id, turn)
593 })
594 .await??;
595 anyhow::Ok(())
596 }
597 .await;
598 (child_id, turn, result)
599 });
600 }
601 if snapshot.operational.native_session_is_ready()
602 && operational.get(&update.session_id).is_none_or(|old: &mj_core::relay::RelayOperationalState| old.config_options != snapshot.operational.config_options)
603 && let Some(session) = controller.state.sessions.get(&update.session_id)
604 && matches!(session.target, Some(mj_core::state::TargetLocator::LocalBare { .. } | mj_core::state::TargetLocator::SshBare { .. } | mj_core::state::TargetLocator::AwsEc2 { .. }))
605 && let Some(build) = snapshot.worker_build.clone()
606 {
607 let profile = session.last_profile.clone();
608 let state = snapshot.operational.clone();
609 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
610 tokio::spawn(async move {
611 let _upgrade_task = upgrade_task;
612 if let Err(error) = crate::controller::profile_config::observe(profile, build, state).await {
613 tracing::warn!(%error, "could not cache observed profile choices");
614 }
615 });
616 }
617 let materialized = snapshot.materialized;
618 let operational_state = snapshot.operational;
619 materialized_activity.insert(
620 update.session_id.clone(),
621 crate::server_runtime::snapshot::MaterializedActivity {
622 last_activity_at_ms: materialized.last_activity_at_ms,
623 execution: materialized.execution,
624 },
625 );
626 let queued = queued_prompt_projection(&materialized);
627 let pending = materialized.pending_elicitations.clone();
628 let active_shells = operational_state.active_user_shells.clone();
629 let prompt_images_supported =
630 agent_accepts_prompt_images(&operational_state);
631 active_user_shells.insert(
632 update.session_id.clone(),
633 active_shells,
634 );
635 conversation_projections.enqueue(materialized);
636 queued_prompts.insert(
637 update.session_id.clone(),
638 queued,
639 );
640 pending_elicitations.insert(
641 update.session_id.clone(),
642 pending,
643 );
644 if prompt_images_supported {
645 prompt_images.insert(update.session_id.clone());
646 } else {
647 prompt_images.remove(&update.session_id);
648 }
649 operational.insert(
650 update.session_id.clone(),
651 operational_state,
652 );
653 revision = daemon_runtime.allocate_revision();
654 conversation_tx.send_replace(conversations.clone());
655 publish_snapshot!(revision);
656 }
657 tokio::task::yield_now().await;
662 }
663 completed = subagent_jobs.join_next(), if !subagent_jobs.is_empty() => {
664 match completed {
665 Some(Ok((identity, Ok(())))) => {
666 active_subagent_requests.remove(&identity);
667 }
668 Some(Ok((identity, Err(error)))) => {
669 active_subagent_requests.remove(&identity);
670 tracing::warn!(
671 parent_session_id = %identity.0,
672 request_id = %identity.1,
673 error = %format!("{error:#}"),
674 "sub-agent tool request failed"
675 );
676 }
677 Some(Err(error)) => tracing::warn!(%error, "sub-agent tool task panicked"),
678 None => {}
679 }
680 }
681 completed = subagent_completion_jobs.join_next(), if !subagent_completion_jobs.is_empty() => {
682 match completed {
683 Some(Ok((_, _, Ok(())))) => {}
684 Some(Ok((child_id, turn, Err(error)))) => {
685 if let Some(relation) = controller.state.subagents.get_mut(&child_id)
686 && relation.noticed_turn == Some(turn)
687 {
688 relation.noticed_turn = None;
689 }
690 tracing::warn!(%child_id, turn, error = %format!("{error:#}"), "could not record the sub-agent completion notice");
691 }
692 Some(Err(error)) => tracing::warn!(%error, "sub-agent completion task panicked"),
693 None => {}
694 }
695 }
696 _ = prune_tick.tick() => {
697 let archive_after_days = controller.config.sessionwiki.archive_after_days;
708 match archive_after_days {
709 Some(days) if archive_jobs.is_empty() => {
710 let runtime = daemon_runtime.clone();
711 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
712 archive_jobs.spawn(async move {
713 let _upgrade_task = upgrade_task;
714 runtime.archive_aged_sessions(days).await
715 });
716 }
717 Some(_) => tracing::debug!(
718 "the previous SessionWiki archive pass is still running; skipping this tick"
719 ),
720 None => daemon_runtime.wiki().request_sync(true),
721 }
722 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
726 tokio::spawn(async move {
727 let _upgrade_task = upgrade_task;
728 let upgrade_blocking = _upgrade_task.clone();
729 let pruned = tokio::task::spawn_blocking(move || {
730 let _upgrade_blocking = upgrade_blocking;
731 crate::database::prune_phone_client_state(client_state_retention)
732 })
733 .await;
734 match pruned {
735 Ok(Ok(0)) => {}
736 Ok(Ok(rows)) => tracing::debug!(rows, "pruned expired phone viewer state"),
737 Ok(Err(error)) => tracing::warn!(%error, "could not prune phone viewer state"),
738 Err(error) => tracing::warn!(%error, "phone viewer state pruning task failed"),
739 }
740 });
741 }
742 _ = credential_tick.tick() => {
743 schedule_due_credential_syncs(
744 &mut credential_sync_signals,
745 &credential_sync_handle,
746 Instant::now(),
747 );
748 while let Some(result) = credential_sync.try_result() {
749 crate::pollers::log_credential_sync_actions(&result);
750 let harness = controller
751 .config
752 .profiles
753 .get(&result.profile_id)
754 .map(|profile| profile.kind);
755 if let Some(notice) = credential_sync_notices.notice(&result, harness) {
756 eprintln!("Mjolnir: {notice}");
757 }
758 }
759 }
760 request = dictation_rx.recv() => {
761 let Some(request) = request else {
762 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering dictation requests");
763 break;
764 };
765 let paths = controller.state.sessions.get(&request.session_id).map(|session| {
766 crate::dictation::auth_paths(&controller.config, &session.last_profile)
767 });
768 let Ok(work) = crate::upgrade::activity("dictation") else { continue };
769 let termination = termination.clone();
770 dictation_jobs.spawn(async move {
771 let _work = work;
772 crate::dictation::execute(request, paths, termination).await
773 });
774 }
775 job = dictation_jobs.join_next(), if !dictation_jobs.is_empty() => {
776 if let Some(Err(error)) = job {
777 tracing::warn!(%error, "web dictation task failed");
778 }
779 }
780 job = archive_jobs.join_next(), if !archive_jobs.is_empty() => {
781 match job {
782 Some(Ok(Ok(0))) | None => {}
783 Some(Ok(Ok(archived))) => tracing::debug!(
784 archived, "the SessionWiki archive pass finished"
785 ),
786 Some(Ok(Err(error))) => tracing::warn!(
787 error = %format!("{error:#}"),
788 "the SessionWiki archive pass failed"
789 ),
790 Some(Err(error)) => tracing::warn!(
791 %error, "the SessionWiki archive task panicked"
792 ),
793 }
794 }
795 stored = client_state_rx.recv() => {
796 let Some(stored) = stored else {
797 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering viewer state requests");
798 break;
799 };
800 let workspace_of = |session_id: &str| {
804 controller
805 .state
806 .sessions
807 .get(session_id)
808 .map(|session| session.workspace_id.clone())
809 };
810 let bundle_of = |session_id: &str| {
811 controller
812 .state
813 .sessions
814 .get(session_id)
815 .map(|session| session.bundle_id.clone())
816 };
817 match stored {
818 crate::server::ClientStateRequest::Read { client_id, session_id, reply } => {
819 let workspace = workspace_of(&session_id);
820 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
821 tokio::spawn(async move {
822 let _upgrade_task = upgrade_task;
823 let upgrade_blocking = _upgrade_task.clone();
824 let answer = tokio::task::spawn_blocking(move || {
825 let _upgrade_blocking = upgrade_blocking;
826 let workspace = workspace.context("unknown session")?;
827 let state = crate::database::client_session_state(
828 &client_id, &workspace, &session_id,
829 )?;
830 anyhow::Ok(crate::server::ViewerClientState {
831 draft: state.draft,
832 through_event_ordinal: state.through_event_ordinal,
833 })
834 })
835 .await;
836 reply.send(flatten_stored(answer)).ok();
837 });
838 }
839 crate::server::ClientStateRequest::SaveDraft { client_id, session_id, draft, reply } => {
840 let workspace = workspace_of(&session_id);
841 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
842 tokio::spawn(async move {
843 let _upgrade_task = upgrade_task;
844 let upgrade_blocking = _upgrade_task.clone();
845 let answer = tokio::task::spawn_blocking(move || {
846 let _upgrade_blocking = upgrade_blocking;
847 let workspace = workspace.context("unknown session")?;
848 crate::database::persist_client_draft(
849 &client_id, &workspace, &session_id, &draft,
850 )
851 })
852 .await;
853 reply.send(flatten_stored(answer)).ok();
854 });
855 }
856 crate::server::ClientStateRequest::MarkWorkspaceRead { client_id, workspace_id, reply } => {
857 let sessions = controller
858 .state
859 .sessions
860 .values()
861 .filter(|session| session.workspace_id == workspace_id)
862 .map(|session| (session.id.clone(), session.viewed_through_event_ordinal))
863 .collect::<Vec<_>>();
864 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
865 tokio::spawn(async move {
866 let _upgrade_task = upgrade_task;
867 let upgrade_blocking = _upgrade_task.clone();
868 let answer = tokio::task::spawn_blocking(move || {
869 let _upgrade_blocking = upgrade_blocking;
870 for (session_id, through) in sessions {
871 crate::database::persist_read_receipt(
875 &client_id, &workspace_id, &session_id, through,
876 )
877 .ok();
878 }
879 anyhow::Ok(())
880 })
881 .await;
882 reply.send(flatten_stored(answer)).ok();
883 });
884 }
885 crate::server::ClientStateRequest::History { session_id, query, scope, reply } => {
886 let bundle = bundle_of(&session_id);
887 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
888 tokio::spawn(async move {
889 let _upgrade_task = upgrade_task;
890 let upgrade_blocking = _upgrade_task.clone();
891 let answer = tokio::task::spawn_blocking(move || {
892 let _upgrade_blocking = upgrade_blocking;
893 let bundle = bundle.context("unknown session")?;
894 let scope = match scope.as_str() {
895 "session" => crate::database::HistoryScope::Session,
896 "all" => crate::database::HistoryScope::All,
897 _ => crate::database::HistoryScope::Project,
898 };
899 let found = crate::database::search_prompts_bounded(
900 &session_id,
901 &bundle,
902 scope,
903 &query,
904 crate::server::MAX_HISTORY_MATCHES,
905 )?;
906 anyhow::Ok(crate::server::ViewerPromptHistory {
907 entries: found
908 .entries
909 .into_iter()
910 .map(|entry| entry.text)
911 .collect(),
912 truncated: found.truncated,
913 })
914 })
915 .await;
916 reply.send(flatten_stored(answer)).ok();
917 });
918 }
919 }
920 }
921 bundle = bundle_rx.recv(), if bundle_jobs.len() < MAX_CONCURRENT_BUNDLE_CREATIONS => {
922 let Some(crate::server::BundleRequest { source, reply }) = bundle else {
923 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering bundle requests");
924 break;
925 };
926 let done = bundle_done_tx.clone();
931 let daemon_runtime = daemon_runtime.clone();
932 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
933 bundle_jobs.spawn(async move {
934 let _upgrade_task = upgrade_task;
935 let result = daemon_runtime
936 .create_quick_bundle(source)
937 .await;
938 if let Err(error) = done.send(BundleCreated { result, reply }) {
939 tracing::debug!(%error, "bundle creation finished after the server stopped");
940 }
941 });
942 }
943 bundle_done = bundle_done_rx.recv() => {
944 let Some(BundleCreated { result, reply }) = bundle_done else {
945 failure = feed_stopped(termination.is_cancelled(), "the bundle creation pipeline stopped while the phone server was running");
946 break;
947 };
948 match result {
949 Ok(created) => {
950 let bundle_id = created.bundle_id;
951 let Some(bundle) = created.config.bundles.get(&bundle_id) else {
957 tracing::error!(%bundle_id, "bundle creation returned a config without its bundle");
958 if reply.send(Err(crate::server::BundleFailure::Controller)).is_err() {
959 tracing::debug!("bundle creation failure reply dropped after client disconnect");
960 }
961 continue;
962 };
963 controller
964 .config
965 .bundles
966 .insert(bundle_id.clone(), bundle.clone());
967 controller_reload_invalidated |= controller_reload_in_flight;
971 revision = daemon_runtime.allocate_revision();
972 publish_snapshot!(revision);
973 request_daemon_controller_reload(
974 daemon_runtime.clone(),
975 "new bundle publication",
976 );
977 if reply.send(Ok(bundle_id)).is_err() {
978 tracing::debug!("bundle creation reply dropped after client disconnect");
979 }
980 }
981 Err(error) => {
982 let failure = match error {
983 crate::controller::QuickBundleFailure::InvalidSource(
984 detail,
985 ) => {
986 tracing::debug!(error = %detail, "phone bundle source was invalid");
987 crate::server::BundleFailure::InvalidSource
988 }
989 crate::controller::QuickBundleFailure::Persistence(
990 detail,
991 ) => {
992 tracing::warn!(error = %detail, "phone bundle creation failed");
993 crate::server::BundleFailure::Controller
994 }
995 };
996 if reply.send(Err(failure)).is_err() {
997 tracing::debug!("bundle creation failure reply dropped after client disconnect");
998 }
999 }
1000 }
1001 }
1002 bundle_job = bundle_jobs.join_next(), if !bundle_jobs.is_empty() => {
1003 if let Some(Err(error)) = bundle_job {
1004 tracing::warn!(%error, "bundle creation task panicked");
1005 }
1006 }
1007 preflight = preflight_rx.recv(), if preflight_jobs.len() < MAX_CONCURRENT_PREFLIGHTS => {
1008 let Some(preflight) = preflight else {
1009 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering preflight requests");
1010 break;
1011 };
1012 let crate::server::NewPreflightRequest {
1016 bundle_id,
1017 target_id,
1018 project_directory,
1019 mut reply,
1020 remote_repairs,
1021 } = match preflight {
1022 crate::server::PreflightRequest::New(request) => request,
1023 crate::server::PreflightRequest::Resume(request) => {
1024 spawn_resume_preflight(
1025 &mut preflight_jobs,
1026 &controller,
1027 request,
1028 &termination,
1029 );
1030 continue;
1031 }
1032 crate::server::PreflightRequest::CompletePath(request) => {
1033 spawn_path_completion(
1034 &mut preflight_jobs,
1035 &controller.config,
1036 request,
1037 &termination,
1038 );
1039 continue;
1040 }
1041 };
1042 let config = controller.config.clone();
1047 let project_validation = project_directory.is_some();
1048 let task_termination = termination.clone();
1049 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
1050 preflight_jobs.spawn(async move {
1051 let _upgrade_task = upgrade_task;
1052 let cancelled = Arc::new(AtomicBool::new(false));
1053 let cancellation_guard = ProcessCancellationGuard(cancelled.clone());
1054 let mut blocking = tokio::task::spawn_blocking(move || {
1055 run_new_preflight_with_cancellation(
1056 config,
1057 bundle_id,
1058 target_id,
1059 project_directory,
1060 cancelled,
1061 remote_repairs,
1062 )
1063 });
1064 let answer = tokio::select! {
1065 biased;
1066 _ = task_termination.cancelled() => None,
1067 _ = reply.closed() => None,
1068 answer = &mut blocking => Some(answer),
1069 };
1070 let Some(answer) = answer else {
1071 drop(cancellation_guard);
1072 match blocking.await {
1073 Err(error) => tracing::warn!(%error, "cancelled phone preflight task failed"),
1074 Ok(Err(error)) => tracing::debug!(%error, "phone preflight cancelled"),
1075 Ok(Ok(_)) => {}
1076 }
1077 return;
1078 };
1079 let answer = match answer {
1080 Ok(Ok(answer)) => Ok(answer),
1081 Ok(Err(error)) => {
1082 tracing::debug!(
1083 error = %error,
1084 project_validation,
1085 "phone preflight check failed"
1086 );
1087 Err(if project_validation {
1088 PreflightFailure::Validation
1089 } else {
1090 PreflightFailure::InvalidRepository(format!("{error:#}"))
1091 })
1092 }
1093 Err(error) => {
1094 tracing::warn!(%error, "phone preflight task failed");
1095 Err(PreflightFailure::Controller(format!(
1096 "preflight task failed: {error}"
1097 )))
1098 }
1099 };
1100 if reply.send(answer).is_err() {
1101 tracing::debug!("phone preflight reply dropped after client disconnect");
1102 }
1103 });
1104 }
1105 preflight_job = preflight_jobs.join_next(), if !preflight_jobs.is_empty() => {
1106 if let Some(Err(error)) = preflight_job {
1107 tracing::warn!(%error, "phone preflight task panicked");
1108 }
1109 }
1110 preparation = move_preparation_rx.recv() => {
1111 let Some(MovePreparationRequest { selection, reply }) = preparation else {
1112 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering move preparation requests");
1113 break;
1114 };
1115 let done = move_prepared_tx.clone();
1120 let daemon_runtime = daemon_runtime.clone();
1121 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
1122 move_preparation_jobs.spawn(async move {
1123 let _upgrade_task = upgrade_task;
1124 let result = daemon_runtime
1125 .prepare_move_session(selection)
1126 .await
1127 .map_err(|error| format!("{error:#}"));
1128 if let Err(error) = done.send(MovePrepared { result, reply }) {
1129 tracing::debug!(%error, "move preparation finished after the server stopped");
1130 }
1131 });
1132 }
1133 prepared = move_prepared_rx.recv() => {
1134 let Some(MovePrepared { result, reply }) = prepared else {
1135 failure = feed_stopped(termination.is_cancelled(), "the move preparation pipeline stopped while the phone server was running");
1136 break;
1137 };
1138 if reply.send(result).is_err() {
1139 tracing::debug!("move preparation reply dropped after client disconnect");
1140 }
1141 }
1142 move_preparation_job = move_preparation_jobs.join_next(), if !move_preparation_jobs.is_empty() => {
1143 if let Some(Err(error)) = move_preparation_job {
1144 tracing::warn!(%error, "move preparation task failed");
1145 }
1146 }
1147 result = native_agent_jobs.join_next(), if !native_agent_jobs.is_empty() => {
1148 match result {
1149 Some(Ok(Ok(agents))) => {
1150 native_agents = agents;
1151 revision = daemon_runtime.allocate_revision();
1152 publish_snapshot!(revision);
1153 }
1154 Some(Ok(Err(error))) => tracing::error!(%error, "could not refresh native subagent identities"),
1155 Some(Err(error)) => tracing::error!(%error, "native subagent refresh task failed"),
1156 None => {}
1157 }
1158 }
1159 move_recovery_job = move_recovery_jobs.join_next(), if !move_recovery_jobs.is_empty() => {
1160 if let Some(Err(error)) = move_recovery_job {
1161 move_recovery_load_in_flight = false;
1162 tracing::warn!(%error, "Move recovery projection task failed");
1163 }
1164 }
1165 receipt = receipt_rx.recv() => {
1166 let Some(ReadReceiptRequest { client_id, session_id, through, reply }) = receipt else {
1167 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering read receipts");
1168 break;
1169 };
1170 match controller.state.sessions.get(&session_id) {
1171 None => {
1172 if reply.send(Err("unknown session".into())).is_err() {
1173 tracing::debug!(%session_id, "unknown-session read receipt reply dropped after client disconnect");
1174 }
1175 }
1176 Some(session) => {
1177 let workspace_id = session.workspace_id.clone();
1178 let done = receipt_done_tx.clone();
1179 let persisted_session_id = session_id.clone();
1180 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
1181 tokio::spawn(async move {
1182 let _upgrade_task = upgrade_task;
1183 let upgrade_blocking = _upgrade_task.clone();
1184 let joined = tokio::task::spawn_blocking(move || {
1185 let _upgrade_blocking = upgrade_blocking;
1186 crate::database::persist_read_receipt(
1187 &client_id,
1188 &workspace_id,
1189 &persisted_session_id,
1190 through,
1191 )
1192 })
1193 .await;
1194 let result = match joined {
1195 Ok(result) => result.map_err(|error| format!("{error:#}")),
1196 Err(error) => Err(format!("phone read receipt task failed: {error}")),
1197 };
1198 if let Err(error) = done.send(ReadReceiptPersisted { session_id, result, reply }) {
1199 tracing::debug!(%error, "phone read receipt finished after the server stopped");
1200 }
1201 });
1202 }
1203 }
1204 }
1205 persisted = receipt_done_rx.recv() => {
1206 let Some(ReadReceiptPersisted { session_id, result, reply }) = persisted else { continue };
1207 match result {
1208 Ok(receipt) => {
1209 let _ = receipt;
1210 if reply.send(Ok(())).is_err() {
1211 tracing::debug!(%session_id, "phone read receipt reply dropped after client disconnect");
1212 }
1213 }
1214 Err(error) => {
1215 tracing::warn!(%session_id, "could not persist a phone read receipt: {error}");
1216 if reply.send(Err(error)).is_err() {
1217 tracing::debug!(%session_id, "failed phone read receipt reply dropped after client disconnect");
1218 }
1219 }
1220 }
1221 }
1222 action = action_rx.recv() => {
1223 let Some(request) = action else {
1224 failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering actions");
1225 break;
1226 };
1227 match &request.action {
1231 ControllerAction::RefreshCapacity { target_id } => {
1232 let known = capacity_state.contains_key(target_id);
1233 let accepted = known && match capacity_triggers_tx.try_send(()) {
1236 Ok(()) | Err(tokio::sync::mpsc::error::TrySendError::Full(())) => true,
1237 Err(tokio::sync::mpsc::error::TrySendError::Closed(())) => {
1238 tracing::warn!("phone capacity refresh rejected: poller stopped");
1239 false
1240 }
1241 };
1242 if known {
1243 if let Some(entry) = capacity_state.get_mut(target_id) {
1244 entry.refreshing = accepted;
1245 if !accepted {
1246 entry.failed = true;
1247 }
1248 }
1249 revision = daemon_runtime.allocate_revision();
1250 publish_snapshot!(revision);
1251 }
1252 let outcome = match (known, accepted) {
1253 (_, true) => ActionOutcome::accepted(),
1254 (false, _) => ActionOutcome::Refused(Refusal::unusable(format!(
1255 "no target named {target_id} is configured"
1256 ))),
1257 (true, false) => ActionOutcome::Failed {
1258 reference: "capacity-poller".to_owned(),
1259 },
1260 };
1261 if request.reply.send(outcome).is_err() {
1262 tracing::debug!(%target_id, "phone capacity refresh reply dropped after client disconnect");
1263 }
1264 tokio::task::yield_now().await;
1265 continue;
1266 }
1267 ControllerAction::RefreshQuota { profile_id } => {
1268 let known = controller.config.enabled_profile(profile_id).is_some();
1269 if known {
1270 quota_batch.generation = quota_batch.generation.saturating_add(1);
1274 quota_batch.profiles = quota_refresh_profiles(&controller);
1275 quota_profiles_tx.send_replace(quota_batch.clone());
1276 }
1277 let outcome = if known {
1278 ActionOutcome::accepted()
1279 } else {
1280 ActionOutcome::Refused(Refusal::unusable(format!(
1281 "no enabled profile named {profile_id} is configured"
1282 )))
1283 };
1284 if request.reply.send(outcome).is_err() {
1285 tracing::debug!(%profile_id, "phone quota refresh reply dropped after client disconnect");
1286 }
1287 tokio::task::yield_now().await;
1288 continue;
1289 }
1290 _ => {}
1291 }
1292 if let ControllerAction::Cancel { session_id } = &request.action {
1293 let outcome = if request_phone_action_cancellation(
1294 session_id,
1295 &action_sessions,
1296 &action_cancellations,
1297 ) {
1298 daemon_runtime.cancel_lifecycle_if_active(session_id);
1299 ActionOutcome::accepted()
1300 } else {
1301 ActionOutcome::NotCancellable
1302 };
1303 if request.reply.send(outcome).is_err() {
1304 tracing::debug!(%session_id, "phone cancellation reply dropped after client disconnect");
1305 }
1306 tokio::task::yield_now().await;
1307 continue;
1308 }
1309 if let ControllerAction::Suspend { session_id } = &request.action {
1310 if !closing_actions.contains_key(session_id) {
1311 request_phone_action_cancellation(session_id, &action_sessions, &action_cancellations);
1312 }
1313 daemon_runtime.request_close(session_id);
1314 }
1315 if let ControllerAction::Destroy { session_id, .. } = &request.action {
1319 request_phone_action_cancellation(session_id, &action_sessions, &action_cancellations);
1320 daemon_runtime.request_close(session_id);
1321 }
1322 let session_id = match admit_phone_action(
1323 &request.action,
1324 action_cancellations.len(),
1325 &mut active_actions,
1326 ) {
1327 Ok(session_id) => session_id,
1328 Err(refusal) => {
1329 if request.reply.send(refusal).is_err() {
1330 tracing::debug!("phone action refusal reply dropped after client disconnect");
1331 }
1332 tokio::task::yield_now().await;
1333 continue;
1334 }
1335 };
1336 let ControllerRequest { action, reply } = request;
1337 let upgrade_work = match crate::upgrade::activity("web action") {
1338 Ok(work) => work,
1339 Err(error) => {
1340 let _ = reply.send(ActionOutcome::Refused(Refusal::unusable(error.to_string())));
1341 continue;
1342 }
1343 };
1344 let done = action_done_tx.clone();
1345 let session_control = worker_commands_tx.clone();
1346 let daemon_runtime = daemon_runtime.clone();
1347 let started = action_started_tx.clone();
1348 next_action_id = next_action_id.wrapping_add(1).max(1);
1349 let action_id = next_action_id;
1350 if let ControllerAction::Suspend { session_id } | ControllerAction::Destroy { session_id, .. } = &action { closing_actions.insert(session_id.clone(), action_id); }
1351 if let ControllerAction::New { workspace_id, .. } = &action {
1352 let workspace_id = if workspace_id.is_empty() && phone_workspaces.len() == 1 {
1353 phone_workspaces[0].id.clone()
1354 } else {
1355 workspace_id.clone()
1356 };
1357 launch_workspaces.insert(action_id, workspace_id);
1358 }
1359 let control = PhoneActionControl::for_action(&action);
1360 action_cancellations.insert(action_id, control.clone());
1361 if let Some(session_id) = &session_id {
1362 action_sessions.insert(action_id, session_id.clone());
1363 }
1364 let suspension_reply = if matches!(action, ControllerAction::Suspend { .. }) {
1365 Some(reply)
1366 } else {
1367 action_replies.accept(action_id, &action, reply);
1368 None
1369 };
1370 let notice_action = match &action {
1371 ControllerAction::SetConfig { .. } => Some("Configuration change"),
1372 ControllerAction::InterruptTurn { .. } => Some("Cancellation"),
1373 ControllerAction::CancelShell { .. } => Some("Shell cancellation"),
1374 ControllerAction::RemoveQueuedPrompt { .. } => Some("Queued prompt removal"),
1375 ControllerAction::RespondElicitation { .. } => Some("Answer"),
1376 _ => None,
1377 };
1378 let notice_sessions = session_control.clone();
1379 let failure_runtime = daemon_runtime.clone();
1380 let lifecycle_failure_prefix = match &action {
1381 ControllerAction::Suspend { .. } => Some(mj_core::state::CLOSE_FAILURE_PREFIX),
1382 ControllerAction::Destroy { .. } => Some(mj_core::state::DESTRUCTION_FAILURE_PREFIX),
1383 _ => None,
1384 };
1385 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
1386 tokio::spawn(async move {
1387 let _upgrade_task = upgrade_task;
1388 let _upgrade_work = upgrade_work;
1389 let upgrade_blocking = _upgrade_task.clone();
1390 let joined = tokio::task::spawn_blocking(move || {
1391 let _upgrade_blocking = upgrade_blocking;
1392 let mut suspension_reply = suspension_reply;
1393 let result = (|| -> Result<()> {
1394 if let ControllerAction::Suspend { session_id } = &action {
1395 mj_core::runtime::block_on(daemon_runtime.prepare_suspension(session_id))??;
1396 if let Some(reply) = suspension_reply.take() {
1397 let _ = reply.send(ActionOutcome::accepted());
1398 }
1399 }
1400 if control.cancelled.load(Ordering::Acquire) {
1401 bail!("phone action cancelled");
1402 }
1403 let mut operation_controller = Controller::load()?;
1404 let executor =
1405 CancellableProcessExecutor::new(control.cancelled.clone());
1406 mj_core::runtime::block_on(apply_phone_action(
1407 &mut operation_controller,
1408 PhoneActionServices {
1409 sessions: &session_control,
1410 daemon_runtime: &daemon_runtime,
1411 },
1412 action,
1413 &executor,
1414 action_id,
1415 &started,
1416 &control,
1417 ))?
1418 })();
1419 let result = result.map_err(|error| PhoneActionFailure::of(&error));
1420 if let Some(reply) = suspension_reply {
1421 let outcome = match &result {
1422 Ok(()) => ActionOutcome::accepted(),
1423 Err(failure) => failure.outcome(&action_reference(action_id)),
1424 };
1425 let _ = reply.send(outcome);
1426 }
1427 result
1428 })
1429 .await;
1430 let result = match joined {
1431 Ok(result) => result,
1432 Err(error) => Err(PhoneActionFailure::internal(format!(
1433 "phone action task failed: {error}"
1434 ))),
1435 };
1436 if let (Err(failure), Some(prefix), Some(id)) = (&result, lifecycle_failure_prefix, &session_id) {
1437 failure_runtime.record_lifecycle_failure(id, &action_reference(action_id), &crate::daemon::LifecycleFailure {
1438 detail: failure.detail.clone(), refusal: failure.refusal.clone(),
1439 }, prefix).await;
1440 }
1441 let notice = match (&result, notice_action, &session_id) {
1442 (Err(failure), Some(action), Some(session_id)) => {
1443 Some((session_id.clone(), failure.conversation_notice(action, action_id)))
1444 }
1445 _ => None,
1446 };
1447 if let Err(error) = done.send((action_id, session_id, result)) {
1448 tracing::debug!(action_id, %error, "phone action finished after the server stopped");
1449 }
1450 if let Some((session_id, text)) = notice {
1451 let recorded = async {
1452 notice_sessions.session(&session_id).await?
1453 .submit(new_command_id("action-notice")?, RelayCommand::RecordNotice { text }).await?;
1454 Ok::<_, anyhow::Error>(())
1455 }.await;
1456 if let Err(error) = recorded {
1457 tracing::warn!(%session_id, %error, "could not record failed session action in conversation");
1458 }
1459 }
1460 });
1461 }
1462 started = action_started_rx.recv() => {
1463 let Some(started) = started else {
1464 tokio::task::yield_now().await;
1465 continue;
1466 };
1467 let started_session_id = started.session.id.clone();
1468 let publication = if !action_cancellations.contains_key(&started.action_id) {
1469 Err("phone action completed before its provisional session was published".into())
1470 } else {
1471 track_started_phone_session(
1472 &mut controller.state,
1473 &mut active_actions,
1474 &mut action_sessions,
1475 started.action_id,
1476 started.session,
1477 )
1478 };
1479 if publication.is_ok() {
1480 revision = daemon_runtime.allocate_revision();
1481 publish_snapshot!(revision);
1482 request_daemon_controller_reload(
1483 daemon_runtime.clone(),
1484 "new session publication",
1485 );
1486 };
1487 if publication.is_err()
1488 && let Some(control) = action_cancellations.get(&started.action_id)
1489 {
1490 control.request_cancel();
1491 }
1492 action_replies.resolve(
1495 started.action_id,
1496 match &publication {
1497 Ok(()) => ActionOutcome::Accepted {
1498 session_id: Some(started_session_id),
1499 },
1500 Err(reason) => ActionOutcome::Refused(Refusal::precondition(
1505 reason.clone(),
1506 )),
1507 },
1508 );
1509 if started.published.send(publication).is_err() {
1510 tracing::debug!(action_id = started.action_id, "phone new-session publication reply dropped after client disconnect");
1511 }
1512 }
1513 completed = action_done_rx.recv() => {
1514 let Some((action_id, session_id, result)) = completed else {
1515 failure = feed_stopped(termination.is_cancelled(), "the phone action pipeline stopped reporting completions");
1516 break;
1517 };
1518 action_cancellations.remove(&action_id);
1519 let session_id = action_sessions.remove(&action_id).or(session_id);
1520 if closing_actions.values().any(|closing_id| *closing_id == action_id) && let Some(id) = &session_id { daemon_runtime.clear_close_request(id); }
1521 closing_actions.retain(|_, closing_id| *closing_id != action_id);
1522 if let Some(session_id) = &session_id && !action_sessions.values().any(|active| active == session_id) {
1523 active_actions.remove(session_id);
1524 }
1525 let reference = action_reference(action_id);
1532 action_replies.resolve(
1533 action_id,
1534 match &result {
1535 Ok(()) => ActionOutcome::Accepted { session_id: session_id.clone() },
1536 Err(failure) => failure.outcome(&reference),
1537 },
1538 );
1539 if let Some(workspace_id) = launch_workspaces.remove(&action_id)
1540 && result.is_err()
1541 && !session_id.as_ref().is_some_and(|id| closing_actions.contains_key(id))
1542 {
1543 record_launch_failure(
1544 &mut launch_failures,
1545 action_id,
1546 workspace_id,
1547 session_id.clone(),
1548 result.as_ref().err().map(|failure| failure.detail.clone()),
1549 );
1550 revision = daemon_runtime.allocate_revision();
1551 publish_snapshot!(revision);
1552 }
1553 if let Err(failure) = &result {
1554 tracing::warn!(
1555 action_id,
1556 reference,
1557 session_id = session_id.as_deref(),
1558 error = %failure.detail,
1559 "phone action failed"
1560 );
1561 }
1562 record_action_result(
1563 &mut pending_action_errors,
1564 session_id.as_deref(),
1565 &result,
1566 );
1567 if result.is_ok() && let Some(session_id) = session_id.clone() {
1570 let daemon_runtime = daemon_runtime.clone();
1571 let Ok(upgrade_task) = crate::upgrade::activity("web background operation") else { continue };
1572 tokio::spawn(async move {
1573 let _upgrade_task = upgrade_task;
1574 daemon_runtime.clear_recorded_close_failure(&session_id).await;
1575 });
1576 }
1577 request_controller_reload(
1578 &mut controller_reload_in_flight,
1579 &mut controller_reload_requested,
1580 &controller_reload_tx,
1581 );
1582 request_move_recovery_reload(
1583 &move_recovery_tx,
1584 &mut move_recovery_load_in_flight,
1585 &mut move_recovery_jobs,
1586 );
1587 request_daemon_controller_reload(
1588 daemon_runtime.clone(),
1589 "phone action completion",
1590 );
1591 }
1592 reloaded = controller_reload_rx.recv() => {
1593 let Some(ControllerReloaded { result }) = reloaded else {
1594 failure = feed_stopped(
1595 termination.is_cancelled(),
1596 "the controller reload pipeline stopped while the phone server was running",
1597 );
1598 break;
1599 };
1600 controller_reload_in_flight = false;
1601 if std::mem::take(&mut controller_reload_invalidated) {
1602 if let Err(error) = &result {
1603 tracing::warn!(%error, "superseded controller reload failed");
1604 }
1605 controller_reload_requested = false;
1606 request_controller_reload(
1607 &mut controller_reload_in_flight,
1608 &mut controller_reload_requested,
1609 &controller_reload_tx,
1610 );
1611 continue;
1612 }
1613 match result {
1614 Ok(mut reloaded) => {
1615 for (session_id, error) in &pending_action_errors {
1616 if let Some(session) = reloaded.state.sessions.get_mut(session_id)
1617 && session.last_error.is_none()
1618 {
1619 session.last_error = Some(error.clone());
1620 }
1621 }
1622 controller = reloaded;
1623 quotas.retain(|id, _| controller.config.enabled_profile(id).is_some());
1624 subagent_quota_reports
1625 .lock()
1626 .expect("sub-agent quota reports lock poisoned")
1627 .retain(|id, _| controller.config.enabled_profile(id).is_some());
1628 worker_targets_tx.send_replace(dashboard_worker_targets(&controller));
1629 publish_capacity_targets(
1630 &controller,
1631 &capacity_targets_tx,
1632 &mut capacity_state,
1633 );
1634 credential_sync_handle.set_targets(credential_sync_targets(&controller));
1635 republish_quota_profiles(
1636 &controller,
1637 &mut published_quota_profiles,
1638 &mut quota_batch,
1639 "a_profiles_tx,
1640 );
1641 profile_catalog.sync(&controller.config);
1649 queued_prompts.retain(|session_id, _| {
1650 controller.state.sessions.contains_key(session_id)
1651 });
1652 pending_elicitations.retain(|session_id, _| {
1653 controller.state.sessions.contains_key(session_id)
1654 });
1655 prompt_images.retain(|session_id| {
1656 controller.state.sessions.contains_key(session_id)
1657 });
1658 operational.retain(|session_id, _| {
1659 controller.state.sessions.contains_key(session_id)
1660 });
1661 materialized_activity.retain(|session_id, _| {
1662 controller.state.sessions.contains_key(session_id)
1663 });
1664 request_move_recovery_reload(
1665 &move_recovery_tx,
1666 &mut move_recovery_load_in_flight,
1667 &mut move_recovery_jobs,
1668 );
1669 conversations.retain(|id, _| {
1670 controller.state.sessions.get(id).is_some_and(|session| session.state.is_active())
1671 });
1672 for session_id in conversation_projections.session_ids() {
1673 if !controller
1674 .state
1675 .sessions
1676 .get(&session_id)
1677 .is_some_and(|session| session.state.is_active())
1678 {
1679 conversation_projections.forget(&session_id);
1680 }
1681 }
1682 revision = daemon_runtime.allocate_revision();
1683 conversation_tx.send_replace(conversations.clone());
1684 publish_snapshot!(revision);
1685 }
1686 Err(error) => {
1687 tracing::warn!(%error, "completed phone operation could not reload controller state");
1688 }
1689 }
1690 if controller_reload_requested {
1691 controller_reload_requested = false;
1692 controller_reload_in_flight = true;
1693 spawn_controller_reload(controller_reload_tx.clone());
1694 }
1695 }
1696 }
1697 }
1698 dictation_jobs.shutdown().await;
1700 bundle_jobs.shutdown().await;
1703 preflight_jobs.shutdown().await;
1707 move_preparation_jobs.shutdown().await;
1710 native_agent_jobs.shutdown().await;
1711 move_recovery_jobs.shutdown().await;
1712 crate::controller::profile_config::cancel_all();
1714 for control in action_cancellations.values() {
1715 control.request_cancel();
1716 }
1717 match failure {
1718 Some(failure) => Err(failure),
1719 None => Ok::<(), anyhow::Error>(()),
1720 }
1721 };
1722 let recorder = tokio::spawn(api_activity::record_activity_stream(
1726 activity_snapshots,
1727 crate::database::record_api_activities,
1728 ));
1729 let result = tokio::select! {
1730 result = serve.stopped() => result,
1731 result = control => result,
1732 };
1733 recorder.abort();
1734 if let Err(error) = recorder.await
1735 && error.is_panic()
1736 {
1737 tracing::warn!(%error, "native API activity recorder panicked");
1738 }
1739 drop(serve);
1742 conversation_projection_shutdown.cancel();
1743 renewal_cancellation.cancel();
1744 if let Some(task) = renewal_task
1745 && let Err(error) = task.await
1746 {
1747 tracing::warn!(%error, "Tailscale certificate renewal task failed");
1748 }
1749 worker_shutdown
1750 .shutdown()
1751 .await
1752 .context("shut down phone server session manager")?;
1753 result?;
1754 Ok(())
1755}
1756
1757pub(crate) struct ViewerServer(tokio::task::JoinHandle<Result<()>>);
1768
1769impl ViewerServer {
1770 pub(crate) fn spawn(
1771 server: impl std::future::Future<Output = Result<()>> + Send + 'static,
1772 ) -> Self {
1773 Self(tokio::spawn(server))
1774 }
1775
1776 pub(crate) async fn stopped(&mut self) -> Result<()> {
1779 match (&mut self.0).await {
1780 Ok(result) => result,
1781 Err(error) => Err(anyhow::Error::new(error).context("the web viewer task failed")),
1782 }
1783 }
1784}
1785
1786impl Drop for ViewerServer {
1787 fn drop(&mut self) {
1788 self.0.abort();
1789 }
1790}