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