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