Skip to main content

mj_controller/server_runtime/
run.rs

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