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