Skip to main content

mj_controller/server_runtime/
run.rs

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    // `list_profiles` is called in the middle of a model's turn, so the
16    // catalogue discovers the profiles' capabilities in the background and
17    // the call only waits for what is already under way.
18    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        &quota_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    // What the capacity poller last said, per probe target. The projection is
60    // built from this on every publish rather than being accumulated, so a
61    // target that disappears from the configuration disappears from the page.
62    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        &quotas,
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    // Captured before `options` is moved into the server.
113    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    // Keep signed-in phones and logged-out identities consistent across
132    // restarts. Loading the signing key and revocations runs off this loop.
133    let cookie_key_path = crate::server::cookie_key_path();
134    options.load_cookie_credentials(cookie_key_path).await?;
135    // The documented `/api/v1` surface authenticates with a persisted bearer
136    // token and drives sessions through the daemon-side backend.
137    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    // Publish a pin only for a certificate the operator configured. A
155    // Tailscale certificate chains to a public CA and renews in place, so a
156    // pin taken now would go stale while ordinary verification keeps working.
157    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        // Stored viewer state expires with the authentication that created it.
213        // The sweep is hourly rather than on every request, because it is
214        // housekeeping and nothing waits for it.
215        let mut prune_tick = tokio::time::interval(prune_tick_interval());
216        // The first index build can take many minutes on a large corpus, and
217        // nothing else would start it until the first hourly tick. Ask for it
218        // here, explicitly, rather than leaning on the interval's immediate
219        // first tick: that tick may run the archive job instead.
220        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        // A feed that ends is not a reason to exit quietly: the phone server
269        // exists to follow sessions, so losing that feed is a named failure
270        // rather than a silent success.
271        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                    &quotas,
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                                // A fleet with nothing running reports no
446                                // figures, and that is an answer rather than a
447                                // failure.
448                                entry.on_demand = usage.is_none();
449                                entry.usage = usage;
450                                entry.failed = false;
451                            }
452                            // The last good reading stays on screen beside the
453                            // failure: one failed probe is not a reason to
454                            // forget what the machine was doing.
455                            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                        // Invalidate a late result before it can be applied or
499                        // launch another queued projection. A session that is
500                        // resumed later gets a new generation from its next
501                        // worker snapshot.
502                        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                        // Controller reload normally removes inactive rows
513                        // first, but this also covers a worker result racing
514                        // that reload and keeps the viewer from seeing a
515                        // conversation for a dead session.
516                        revision = daemon_runtime.allocate_revision();
517                        conversation_tx.send_replace(conversations.clone());
518                        publish_snapshot!(revision);
519                    }
520                    // SessionManagerUpdates has a synchronous pending fast
521                    // path. Yield after each completion so a hot stream of
522                    // updates cannot monopolize this runtime worker.
523                    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                                    // The result reaches the model as the tool
564                                    // call's own answer: completing the request
565                                    // unblocks the worker socket the harness is
566                                    // waiting on. It is not injected as a turn.
567                                    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                            // The notice names the child as every listing does.
591                            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                                    // A child that owes its report is reminded
618                                    // first; the parent hears about it when the
619                                    // reminder turn ends.
620                                    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                                    // The parent has been told this turn
644                                    // ended and nothing else is queued, so
645                                    // the child gives its processes back until
646                                    // the parent sends it more input (#1161).
647                                    // The park checks again, under the
648                                    // child's lifecycle, that it is idle.
649                                    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                    // The session update receiver can return pending entries
715                    // without touching Tokio's budgeted receive operation.
716                    // Give HTTP/TLS tasks a scheduling opportunity after each
717                    // update even when the worker is publishing continuously.
718                    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                    // The hourly full SessionWiki sync. Session closes drive
755                    // bounded syncs; this one also reconciles sessions deleted
756                    // outside the daemon and picks up any bounded run that
757                    // failed.
758                    //
759                    // With `archive_after_days` set, the same tick runs the
760                    // archive job instead, because that job starts with the
761                    // full sync itself. It runs as a background task, and only
762                    // when no earlier one is still running: a pass over a large
763                    // corpus can outlast the tick.
764                    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                    // Only rows whose client id names a phone are considered:
780                    // a terminal client's place in a conversation is not the
781                    // phone's to expire.
782                    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                    // Every one of these touches SQLite, so each runs on its
864                    // own task. A composer autosaving on a debounce must never
865                    // be able to stall the loop that follows sessions.
866                    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                                        // A receipt that would move backwards
935                                        // is not an error; it is a session this
936                                        // viewer had already read past.
937                                        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                    // Repository canonicalization and config persistence both
990                    // touch the filesystem. Keep them off this loop, and
991                    // report a panic as a failed request rather than dropping
992                    // the browser's reply.
993                    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                            // Other config sections may have changed while
1016                            // this request was in flight. Publish only the
1017                            // bundle this transaction created; a full fresh
1018                            // config is requested below and must not make a
1019                            // later completion hide another completed bundle.
1020                            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                            // A reload started before this save may still be
1032                            // queued. It must not hide a bundle after we have
1033                            // acknowledged it as available to the browser.
1034                            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                    // A resume preflight asks a different question about the
1077                    // same disk, so it runs on the same supervised task set
1078                    // and under the same cap as a new-session preflight.
1079                    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                    // Reading a working tree's status or validating a project
1111                    // directory touches the disk, so it runs on its own task
1112                    // rather than on the loop that has to stay responsive to
1113                    // every other feed.
1114                    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                    // Preparation can inspect archives, target prerequisites,
1184                    // and harness capabilities. Keep it supervised and away
1185                    // from this feed loop so another browser can still read
1186                    // snapshots while a move form is open.
1187                    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                    // A refresh nudges a poller this loop owns. It takes no
1296                    // session slot and starts no lifecycle work, so it is
1297                    // answered here rather than admitted as an action.
1298                    match &request.action {
1299                        ControllerAction::RefreshCapacity { target_id } => {
1300                            let known = capacity_state.contains_key(target_id);
1301                            // One queued nudge refreshes every target. Do not
1302                            // block the consumer while readings wait for it.
1303                            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                                // The refresher works from a generation-stamped
1339                                // batch, so a new generation is how one is asked
1340                                // for again rather than a per-profile trigger.
1341                                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                    // A force close runs even while a graceful close for the
1384                    // same session is still in flight; that stuck close is
1385                    // exactly what it is meant to take over.
1386                    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                    // The phone asked for a session, and now there is one to
1562                    // point at: that is what its request was waiting for.
1563                    action_replies.resolve(
1564                        started.action_id,
1565                        match &publication {
1566                            Ok(()) => ActionOutcome::Accepted {
1567                                session_id: Some(started_session_id),
1568                            },
1569                            // Publication fails only for reasons this loop
1570                            // words itself -- a race for the new session, or
1571                            // an action that ended first -- so the text is
1572                            // already safe and specific enough to send.
1573                            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                    // A `new` that failed before publishing a session never
1595                    // reached the arm that answers it, so its phone is still
1596                    // waiting for a reply it can act on. A failure that named
1597                    // a reason the caller can fix answers with that reason;
1598                    // every other one answers generically and points at the
1599                    // log entry below.
1600                    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                    // A session that is taking work again has recovered from
1637                    // whatever its last close did, so it stops reporting it.
1638                    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                                &quota_profiles_tx,
1709                            );
1710                            // A changed profile set or sub-agent policy makes
1711                            // the catalogue's answers wrong, so it drops them,
1712                            // adopts the configuration it is given here, and
1713                            // discovers the new one in the background. A
1714                            // `list_profiles` call that arrives first waits on
1715                            // that pass's discoveries rather than starting its
1716                            // own.
1717                            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        // Stop provider requests before the HTTP request channels disappear.
1768        dictation_jobs.shutdown().await;
1769        // Bundle jobs are supervised so shutdown never leaves a detached
1770        // request task behind holding the config mutation lock.
1771        bundle_jobs.shutdown().await;
1772        // Preflight jobs own cancellation guards for their blocking Git
1773        // probes. Aborting them here signals those probes before the server's
1774        // request channels disappear.
1775        preflight_jobs.shutdown().await;
1776        // Preparation tasks may be inspecting an archive or probing a target;
1777        // abort and drain them before the HTTP server's channels disappear.
1778        move_preparation_jobs.shutdown().await;
1779        native_agent_jobs.shutdown().await;
1780        move_recovery_jobs.shutdown().await;
1781        // Every exit stops in-flight work, whether it was asked for or forced.
1782        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    // The recorder runs beside the server, never as an arm of this select:
1792    // a recording failure must not end `run_server` and take the API down
1793    // with it (issue 1117).
1794    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    // Dropping the handle aborts the server task, which is what dropping the
1809    // server future used to do when this `select!` owned it directly.
1810    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
1826/// The viewer's HTTP server, running on a task of its own.
1827///
1828/// The listener must not share a task with the control loop above: whatever
1829/// the loop is doing during one of its turns, a server polled by the same
1830/// `select!` cannot accept a connection until that turn ends. A cheap read
1831/// such as `GET /api/v1/sessions` then waits for unrelated work — the stall
1832/// reported in issue 1061, where a list request timed out at ten seconds
1833/// while a session was provisioning and answered instantly on the next try.
1834/// On its own task the listener is scheduled independently, so a slow turn in
1835/// the control loop can only make an answer stale, never late.
1836pub(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    /// Resolves when the server stops on its own, with whatever it stopped
1846    /// for. A panicked server is a failure rather than a silent exit.
1847    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}