Skip to main content

mj_controller/
server.rs

1//! Daemon-owned, phone-oriented control surface for Hel.
2//!
3//! The server deliberately owns no controller business logic. It publishes a
4//! redacted projection of controller state and forwards validated, typed
5//! actions through a channel supplied by the controller.
6
7use std::collections::BTreeMap;
8use std::convert::Infallible;
9use std::net::SocketAddr;
10use std::path::{Component, PathBuf};
11use std::sync::{Arc, Mutex};
12use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
13
14use anyhow::{Context, Result as AnyResult};
15use axum::body::{Body, Bytes, to_bytes};
16use axum::extract::{DefaultBodyLimit, Path, Query, Request, State};
17use axum::http::header::{
18    CACHE_CONTROL, CONTENT_SECURITY_POLICY as CONTENT_SECURITY_POLICY_HEADER, CONTENT_TYPE, COOKIE,
19    HeaderValue, LOCATION, REFERRER_POLICY, SET_COOKIE, X_CONTENT_TYPE_OPTIONS,
20};
21use axum::http::{HeaderMap, Response, StatusCode};
22use axum::middleware::Next;
23use axum::response::IntoResponse;
24use axum::response::sse::{Event, KeepAlive, Sse};
25use axum::routing::{get, post, put};
26use axum::{Json, Router};
27use base64::Engine as _;
28use hmac::{Hmac, KeyInit, Mac};
29use serde::{Deserialize, Serialize};
30use sha2::Sha256;
31use tokio::sync::Semaphore;
32use tokio::sync::{mpsc, watch};
33use tokio_stream::wrappers::ReceiverStream;
34use tokio_util::sync::CancellationToken;
35
36use mj_core::attachment::{AttachmentRef, AttachmentStore, MAX_IMAGE_BYTES, MAX_IMAGES};
37use mj_core::config::{Config, TargetTemplate, project_history_host, validate_id};
38use mj_core::elicitation::{ElicitationRequest, ElicitationResponse, MAX_ELICITATION_BYTES};
39use mj_core::path_completion::{CompletionHost, CompletionKind, PathCompletion};
40use mj_core::refusal::{Refusal, RefusalKind};
41use mj_core::state::{
42    MoveOperation, MovePhase, MovePreparation, MoveSelection, MoveSessionRequest,
43    ProjectSourceIdentity, SessionResourceAllocation, SessionState, SessionTransitionKind,
44    State as AppState,
45};
46
47use crate::targets::AdditionalMount;
48
49use crate::dictation::{
50    DictationError, DictationOperation, DictationRequest, DictationResponse, MAX_AUDIO_BYTES,
51    validate_wav,
52};
53use crate::image::optimize_image;
54
55pub mod api;
56
57pub use api::{
58    ApiFailure, ApiSession, PromptRequest, PromptResponse, SessionListResponse,
59    StartSessionRequest, StartSessionResponse, SubagentBackend, WaitOutcome, WaitRequest,
60    WaitResponse, api_token_path, load_or_create_api_token, map_stop_reason, resolve_wait,
61};
62
63pub use mj_client::web::{
64    BrowserDiffStat, BrowserTranscript, BrowserTranscriptEntry, WebListenerProcess,
65    WebViewerAccess, WebViewerRecovery,
66};
67
68// Keep all control surfaces on the same queue vocabulary. The resume flow
69// used to define a private copy here, which made a move request impossible to
70// pass through the web and daemon boundaries without lossy conversion.
71pub use mj_core::state::ResumeQueueDisposition;
72
73/// Select the process-wide rustls provider before any TLS configuration is built.
74///
75/// Dependency feature unification can enable both rustls providers. Rustls
76/// deliberately refuses to guess in that case, so each executable that links
77/// the controller installs the ring provider at process startup. A provider
78/// installed even earlier is already sufficient and remains in place.
79pub fn install_rustls_crypto_provider() {
80    let _ = rustls::crypto::ring::default_provider().install_default();
81}
82
83pub const COOKIE_NAME: &str = "hel_viewer_session";
84const DEFAULT_SESSION_TTL: Duration = Duration::from_secs(30 * 24 * 60 * 60);
85const EPHEMERAL_SESSION_TTL: Duration = Duration::from_secs(24 * 60 * 60);
86const MAX_BODY_BYTES: usize = 128 * 1024;
87const MAX_CODE_FAILURES: u32 = 5;
88const CODE_LOCKOUT_BASE: Duration = Duration::from_secs(30);
89const CODE_LOCKOUT_CAP: Duration = Duration::from_secs(60 * 60);
90const MAX_TITLE_CHARS: usize = 120;
91const MAX_PROMPT_CHARS: usize = 64 * 1024;
92/// How many repositories one dirty-worktree acknowledgement may name. A bundle
93/// with more repositories than this than has bigger problems than the phone.
94const MAX_DIRTY_ACKNOWLEDGEMENTS: usize = 32;
95/// The largest draft a phone may store. A composer is for a prompt, and a
96/// prompt this size has other problems; the bound exists so one viewer cannot
97/// fill the daemon's database with text it never sent.
98const MAX_DRAFT_BYTES: usize = 64 * 1024;
99/// How many prompt-history matches one search returns. Public because the
100/// controller loop performs the search and must use the same bound the phone
101/// was promised.
102pub const MAX_HISTORY_MATCHES: usize = 40;
103/// Image prompts need far more room than any other phone request. Browser
104/// uploads are base64-encoded, so two ordinary photographs already exceed the
105/// general body limit even when each one fits it. The larger bound therefore
106/// stays scoped to the action route that carries prompts.
107const MAX_PROMPT_BODY_BYTES: usize = 32 * 1024 * 1024;
108/// A browser uploads one source image at a time. The image optimizer has its
109/// own decoded-allocation bound; this is the HTTP envelope bound before that
110/// work starts.
111const MAX_ATTACHMENT_UPLOAD_BYTES: usize = 64 * 1024 * 1024;
112/// Keep two browser uploads/transcriptions in flight. The permit is acquired
113/// before reading the request body, so an overloaded client is rejected
114/// without accepting megabytes that cannot be processed yet.
115const MAX_CONCURRENT_DICTATIONS: usize = 2;
116/// Keep this in sync with the prompt admission bound and the browser composer.
117pub const MAX_PROMPT_IMAGES: usize = MAX_IMAGES;
118const COOKIE_KEY_BYTES: usize = 32;
119const COOKIE_KEY_FILE: &str = "phone-cookie-key";
120
121/// How long stored viewer state outlives its last use.
122///
123/// It matches the session cookie's own lifetime: state keyed to an identity
124/// that can no longer authenticate has nothing left to belong to.
125pub const fn default_session_ttl() -> Duration {
126    DEFAULT_SESSION_TTL
127}
128
129pub fn cookie_key_path() -> PathBuf {
130    mj_core::config::data_dir().join(COOKIE_KEY_FILE)
131}
132
133/// Load the phone cookie signing key, creating it on first use.
134///
135/// This signing key keeps a signed-in phone signed in across daemon restarts.
136/// Logged-out identities are tracked separately. Deleting the key is
137/// therefore the explicit sign-everyone-out gesture: the next start writes a
138/// new key and every outstanding cookie stops validating. A missing file is
139/// ordinary first use; an unreadable or too-short one is replaced loudly,
140/// because refusing to start would be a worse answer than asking phones to
141/// enter the viewer code again.
142pub fn load_or_create_cookie_key(path: &std::path::Path) -> AnyResult<Vec<u8>> {
143    match std::fs::read(path) {
144        Ok(key) if key.len() >= COOKIE_KEY_BYTES => return Ok(key),
145        Ok(key) => tracing::warn!(
146            path = %path.display(),
147            bytes = key.len(),
148            "phone cookie key is shorter than {COOKIE_KEY_BYTES} bytes; generating a new key signs every phone out"
149        ),
150        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
151        Err(error) => tracing::warn!(
152            path = %path.display(),
153            "could not read the phone cookie key ({error}); generating a new key signs every phone out"
154        ),
155    }
156    let key = generate_cookie_key()?;
157    mj_core::config::atomic_write(path, &key)
158        .with_context(|| format!("persist Mjolnir phone cookie key {}", path.display()))?;
159    Ok(key.to_vec())
160}
161
162/// Options for the daemon's phone service.
163///
164/// `ServerOptions::new` generates both the six-digit viewer code and an
165/// ephemeral cookie key. A caller that wants cookies to survive server
166/// restarts loads its key and logout records with `load_cookie_credentials`.
167/// `set_cookie_key` installs a key for callers managing storage themselves. The
168/// key and viewer code are intentionally omitted from `Debug` output.
169#[derive(Clone)]
170pub struct ServerOptions {
171    pub bind: SocketAddr,
172    pub snapshot_rx: watch::Receiver<ViewerSnapshot>,
173    pub conversation_rx: watch::Receiver<BTreeMap<String, BrowserTranscript>>,
174    pub action_tx: mpsc::Sender<ControllerRequest>,
175    pub bundle_tx: mpsc::Sender<BundleRequest>,
176    pub receipt_tx: mpsc::Sender<ReadReceiptRequest>,
177    pub preflight_tx: mpsc::Sender<PreflightRequest>,
178    pub move_preparation_tx: mpsc::Sender<MovePreparationRequest>,
179    pub client_state_tx: mpsc::Sender<ClientStateRequest>,
180    pub dictation_tx: mpsc::Sender<DictationRequest>,
181    /// Dedicated bounded path for stopping one live background task. This is
182    /// deliberately separate from [`ControllerAction`]: stopping a provider
183    /// task does not occupy the controller's action admission slot.
184    background_task_stop_tx: mpsc::Sender<BackgroundTaskStopRequest>,
185    pub shutdown: CancellationToken,
186    pub session_ttl: Duration,
187    /// Keep this enabled for direct HTTPS or an HTTPS reverse proxy. It may be
188    /// disabled only for an explicitly trusted HTTP development endpoint.
189    pub secure_cookie: bool,
190    tls_config: Option<axum_server::tls_rustls::RustlsConfig>,
191    viewer_code: String,
192    login_token: String,
193    cookie_key: Vec<u8>,
194    viewer_revocations: Arc<ViewerRevocations>,
195    api_token: String,
196    subagent: Option<Arc<dyn api::SubagentBackend>>,
197    /// Where the remembered fast-start preferences live. Injected rather than
198    /// resolved from the process environment inside a handler, so the
199    /// documented API answers from its own state and tests never read the
200    /// developer's real configuration.
201    preferences_path: PathBuf,
202    /// Checks the engine behind each local container target for the launch
203    /// options route. Tests replace it so the answer does not depend on the
204    /// engines the test host has.
205    engine_checks: Arc<api::LocalEngineChecks>,
206}
207
208/// Typed request channels served by the authenticated HTTP surface.
209pub struct ServerRequests {
210    pub action_tx: mpsc::Sender<ControllerRequest>,
211    pub bundle_tx: mpsc::Sender<BundleRequest>,
212    pub receipt_tx: mpsc::Sender<ReadReceiptRequest>,
213    pub preflight_tx: mpsc::Sender<PreflightRequest>,
214    pub move_preparation_tx: mpsc::Sender<MovePreparationRequest>,
215    pub client_state_tx: mpsc::Sender<ClientStateRequest>,
216    pub dictation_tx: mpsc::Sender<DictationRequest>,
217}
218
219impl ServerOptions {
220    pub fn new(
221        bind: SocketAddr,
222        snapshot_rx: watch::Receiver<ViewerSnapshot>,
223        conversation_rx: watch::Receiver<BTreeMap<String, BrowserTranscript>>,
224        requests: ServerRequests,
225    ) -> AnyResult<Self> {
226        let cookie_key = generate_cookie_key()?.to_vec();
227        Ok(Self {
228            bind,
229            snapshot_rx,
230            conversation_rx,
231            action_tx: requests.action_tx,
232            bundle_tx: requests.bundle_tx,
233            receipt_tx: requests.receipt_tx,
234            preflight_tx: requests.preflight_tx,
235            move_preparation_tx: requests.move_preparation_tx,
236            client_state_tx: requests.client_state_tx,
237            dictation_tx: requests.dictation_tx,
238            background_task_stop_tx: mpsc::channel(1).0,
239            shutdown: CancellationToken::new(),
240            session_ttl: DEFAULT_SESSION_TTL,
241            secure_cookie: true,
242            tls_config: None,
243            viewer_code: generate_viewer_code()?,
244            login_token: derive_login_token(&cookie_key),
245            cookie_key,
246            viewer_revocations: Arc::new(ViewerRevocations::default()),
247            // An empty token authenticates nothing: the daemon installs the
248            // persisted one, and a server without it serves the viewer only.
249            api_token: String::new(),
250            subagent: None,
251            preferences_path: mj_core::go::GoPreferences::path(),
252            engine_checks: Arc::new(api::LocalEngineChecks::on_this_host()),
253        })
254    }
255
256    pub fn viewer_code(&self) -> &str {
257        &self.viewer_code
258    }
259
260    pub fn login_token(&self) -> &str {
261        &self.login_token
262    }
263
264    /// Serve HTTPS directly using the supplied Rustls configuration. Hel's
265    /// CLI can load its persisted certificate (including a Tailscale-issued
266    /// certificate) and pass it here without coupling this module to disk.
267    pub fn set_tls_config(&mut self, config: axum_server::tls_rustls::RustlsConfig) {
268        self.tls_config = Some(config);
269        self.secure_cookie = true;
270    }
271
272    /// Install a persisted signing key. Rotating this value signs every phone
273    /// out and revokes QR login URLs. This does not load durable logout records;
274    /// use `load_cookie_credentials` for the daemon's persisted credentials.
275    pub fn set_cookie_key(&mut self, key: Vec<u8>) -> AnyResult<()> {
276        anyhow::ensure!(
277            key.len() >= COOKIE_KEY_BYTES,
278            "cookie signing key must be at least {COOKIE_KEY_BYTES} bytes"
279        );
280        self.login_token = derive_login_token(&key);
281        self.cookie_key = key;
282        Ok(())
283    }
284
285    /// Load the signing key and durable logout records off the async runtime.
286    pub async fn load_cookie_credentials(&mut self, path: PathBuf) -> AnyResult<()> {
287        let (key, revocations) = tokio::task::spawn_blocking(move || {
288            let key = load_or_create_cookie_key(&path)?;
289            let revocations =
290                ViewerRevocations::load(path.with_file_name("phone-cookie-revocations.json"))?;
291            Ok::<_, anyhow::Error>((key, revocations))
292        })
293        .await
294        .context("load viewer credentials task")??;
295        self.set_cookie_key(key)?;
296        self.viewer_revocations = Arc::new(revocations);
297        Ok(())
298    }
299
300    /// Install the controller's bounded background-task stop path.
301    pub fn set_background_task_stop_tx(&mut self, tx: mpsc::Sender<BackgroundTaskStopRequest>) {
302        self.background_task_stop_tx = tx;
303    }
304
305    /// Install the persisted bearer token for the `/api/v1` routes. Rotating
306    /// it revokes every client that still holds the old one.
307    pub fn set_api_token(&mut self, token: String) {
308        self.api_token = token;
309    }
310
311    /// Install the daemon-side backend the `/api/v1` routes drive sessions
312    /// through. Without it those routes answer 503.
313    pub fn set_subagent_backend(&mut self, backend: Arc<dyn api::SubagentBackend>) {
314        self.subagent = Some(backend);
315    }
316
317    /// Read the remembered default from a specific preferences file instead
318    /// of the user's own. Tests point this at a temporary path.
319    pub fn set_preferences_path(&mut self, path: PathBuf) {
320        self.preferences_path = path;
321    }
322
323    /// Answer local engine checks with `probe` instead of this host's engines.
324    #[cfg(test)]
325    fn set_engine_probe(&mut self, probe: api::EngineProbe) {
326        self.engine_checks = Arc::new(api::LocalEngineChecks::new(probe));
327    }
328
329    #[cfg(test)]
330    fn with_test_credentials(mut self, code: &str, key: &[u8]) -> Self {
331        self.viewer_code = code.to_string();
332        self.login_token = "test-login-token".into();
333        self.cookie_key = key.to_vec();
334        self.secure_cookie = false;
335        self.api_token = "test-api-token".into();
336        self
337    }
338}
339
340/// Run the phone server until its shutdown token is cancelled.
341///
342/// This binds only the requested listener. It does not daemonize, provision a
343/// target, or keep sessions alive: controller availability is required, just
344/// like MJ's explicit remote-viewer model.
345pub async fn run_server(options: ServerOptions) -> AnyResult<()> {
346    let listener = tokio::net::TcpListener::bind(options.bind)
347        .await
348        .with_context(|| format!("bind web viewer to {}", options.bind))?;
349    run_server_on_listener(options, listener).await
350}
351
352/// Serve a reserved socket so readiness and advertised ports reflect a real listener.
353pub async fn run_server_on_listener(
354    options: ServerOptions,
355    listener: tokio::net::TcpListener,
356) -> AnyResult<()> {
357    let mut options = options;
358    let bind = listener.local_addr().context("read web viewer address")?;
359    let shutdown = options.shutdown.clone();
360    let viewer_code = options.viewer_code.clone();
361    let tls_config = options.tls_config.take();
362    let app = router(options);
363    println!("Mjolnir viewer code: {viewer_code}");
364    let listener = listener.into_std().context("prepare web viewer listener")?;
365    let handle = axum_server::Handle::new();
366    let shutdown_handle = handle.clone();
367    let serve = async move {
368        if let Some(tls_config) = tls_config {
369            axum_server::from_tcp_rustls(listener, tls_config)
370                .handle(handle)
371                .serve(app.into_make_service())
372                .await
373        } else {
374            axum_server::from_tcp(listener)
375                .handle(handle)
376                .serve(app.into_make_service())
377                .await
378        }
379    };
380    tokio::pin!(serve);
381    tokio::select! {
382        result = &mut serve => result,
383        _ = shutdown.cancelled() => {
384            shutdown_handle.graceful_shutdown(Some(Duration::from_secs(2)));
385            serve.await
386        }
387    }
388    .with_context(|| format!("serve web viewer on {bind}"))
389}
390
391mod viewer_types;
392pub use viewer_types::*;
393mod actions;
394pub use actions::*;
395mod routes;
396use routes::*;
397mod handlers;
398use handlers::*;
399mod validation;
400use validation::*;
401mod errors;
402use errors::*;
403mod auth;
404pub use auth::*;
405mod assets;
406use assets::*;
407mod config_view;
408pub use config_view::*;
409
410#[cfg(test)]
411mod tests;