Skip to main content

mj_client/
daemon.rs

1//! Authenticated local daemon protocol and client transport.
2use crate::executable::describe_running_daemon_and_client_builds;
3use crate::review::RuntimeReviewView;
4use crate::session::{ManagedSessionView, ViewError};
5use anyhow::{Context, Result, bail, ensure};
6use mj_core::config::{Config, data_dir};
7use mj_core::credentials::CredentialSyncSignal;
8use mj_core::elicitation::ElicitationResponse;
9use mj_core::relay::{RelayCommand, RelayOperationalState};
10use mj_core::review::driver::Resolution;
11use mj_core::state::*;
12use mj_core::targets::{AdditionalMount, ProvisionStage};
13use mj_core::workspace::WorkspaceRecord;
14use serde::{Deserialize, Serialize};
15use std::collections::BTreeMap;
16use std::fs;
17use std::net::SocketAddr;
18use std::path::PathBuf;
19use std::time::{Duration, Instant};
20use tokio::io::{AsyncReadExt, AsyncWriteExt};
21use tokio::net::TcpStream;
22pub fn metadata_path() -> PathBuf {
23    data_dir().join("daemon.json")
24}
25
26#[derive(Debug, Clone, Serialize, Deserialize)]
27#[serde(deny_unknown_fields)]
28pub struct DaemonMetadata {
29    pub protocol_version: u32,
30    pub pid: u32,
31    pub address: SocketAddr,
32    pub token: String,
33    pub started_at: String,
34    pub build_version: String,
35}
36
37#[derive(Debug, Clone, Serialize, Deserialize)]
38#[serde(deny_unknown_fields)]
39pub struct WorkspaceListing {
40    pub workspace: WorkspaceRecord,
41}
42
43#[derive(Debug, Clone, Serialize, Deserialize)]
44#[serde(deny_unknown_fields)]
45pub struct SessionPreview {
46    pub id: String,
47    pub title: String,
48    pub project: String,
49    pub harness: String,
50    pub state: String,
51    pub active: bool,
52    pub updated_at: String,
53}
54
55/// One session in the user's SessionWiki index, as a control surface shows it.
56///
57/// It is the daemon's own shape rather than SessionWiki's row: it carries the
58/// search snippet that found the row and, for a Mjolnir session this daemon
59/// still holds, the id that resumes it instead of restoring it.
60#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
61#[serde(deny_unknown_fields)]
62pub struct WikiRow {
63    pub id: String,
64    pub tool: String,
65    pub project: String,
66    pub title: String,
67    pub started: Option<String>,
68    pub msgs: i64,
69    pub preview: Option<String>,
70    /// The tool deleted its own copy and SessionWiki kept the transcript.
71    pub archived: bool,
72    /// The session id the tool that ran it knows it by, when its stored path
73    /// carries one. It is what matches a row against an import scan.
74    pub native_id: Option<String>,
75    /// The matching text, when this row came from a search.
76    pub snippet: Option<String>,
77    /// The live Mjolnir session this row describes, when this daemon has it.
78    pub hel_session_id: Option<String>,
79    /// The Mjolnir target template the session ran under, when the index
80    /// carries it. Only Mjolnir's own rows have one: the daemon writes it into
81    /// the index as an `mj-target:` tag while the session is still known, so an
82    /// archived row can still say where it ran.
83    #[serde(default)]
84    pub target: Option<String>,
85    /// The Mjolnir harness profile the session last ran under, from the index's
86    /// `mj-profile:` tag. Only Mjolnir's own rows have one.
87    #[serde(default)]
88    pub profile: Option<String>,
89    /// The harness kind the session ran (`codex`, `claude`, `kimi`, `grok`,
90    /// `muse`), from the index's `mj-harness:` tag. It stays meaningful after
91    /// the profile id has been removed from the configuration.
92    #[serde(default)]
93    pub harness: Option<String>,
94}
95
96/// How far along the daemon's SessionWiki index is when a search answers.
97#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
98#[serde(rename_all = "snake_case")]
99pub enum WikiIndexState {
100    /// The index has completed one full build; its answers are complete.
101    Ready,
102    /// The first full build has not finished yet, so a search can miss
103    /// sessions that exist. This is the state a fresh index starts in.
104    #[default]
105    Indexing,
106    /// The index file on disk was written by a different SessionWiki schema
107    /// version. Mjolnir will not open it, because opening it would drop and
108    /// rebuild the user's whole cache.
109    VersionMismatch,
110}
111
112/// What a search says about the index it answered from.
113#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
114#[serde(deny_unknown_fields)]
115pub struct WikiStatus {
116    pub state: WikiIndexState,
117    /// A sync is running now, so repeating the query may return more.
118    pub topping_up: bool,
119}
120
121/// One page of search results with the state of the index behind them.
122#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
123#[serde(deny_unknown_fields)]
124pub struct WikiSearchPage {
125    pub rows: Vec<WikiRow>,
126    pub status: WikiStatus,
127}
128
129/// One message of an indexed transcript, reduced to what a search preview
130/// shows: the text around the query's matches, with the matches located in it.
131#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
132#[serde(deny_unknown_fields)]
133pub struct WikiHitBlock {
134    /// `user`, `assistant` or `tool`.
135    pub role: String,
136    /// The message text, redacted, and windowed to the caller's per-message
137    /// budget when the message is longer than that.
138    pub text: String,
139    /// Byte ranges of the matches inside `text`, on character boundaries, in
140    /// order. A context message has none.
141    pub hits: Vec<(usize, usize)>,
142    /// Messages between the previous block and this one that no group covered.
143    /// Non-zero only on the first block of a group.
144    pub omitted_before: usize,
145    /// `text` is a window of the message rather than the whole of it.
146    pub truncated: bool,
147}
148
149/// The matching passages of one indexed transcript.
150#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
151#[serde(deny_unknown_fields)]
152pub struct WikiHitTranscript {
153    pub blocks: Vec<WikiHitBlock>,
154    /// Messages after the last block that no group covered.
155    pub omitted_after: usize,
156}
157
158/// What one indexed session is, as far as continuing it is concerned.
159#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
160#[serde(rename_all = "snake_case")]
161pub enum WikiSessionStatus {
162    /// A Mjolnir session this daemon still has a record of.
163    Mine,
164    /// A Mjolnir session whose record the archive job destroyed; only the
165    /// indexed transcript is left.
166    Archived,
167    /// Another tool's session, which Mjolnir would have to import.
168    Native,
169}
170
171/// One row of the SessionWiki index, with what Mjolnir knows about it.
172///
173/// This is what `mj resume --wiki` branches on and what `mj sessions
174/// --session` reports when the id names no Mjolnir session.
175#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
176#[serde(deny_unknown_fields)]
177pub struct WikiSessionInfo {
178    /// The SessionWiki id of the row.
179    pub wiki_id: String,
180    /// The tool that produced the session, as SessionWiki names it.
181    pub tool: String,
182    /// The transcript path the index stores for the row.
183    pub path: PathBuf,
184    pub status: WikiSessionStatus,
185    /// The Mjolnir session id, for a row Mjolnir itself published.
186    pub mjolnir_session_id: Option<String>,
187    /// The profile the session last ran under, from the index's own tags.
188    pub profile_id: Option<String>,
189    /// The target template the session ran on, from the index's own tags.
190    pub target_template_id: Option<String>,
191    /// The harness that drove the session, from the tags for a Mjolnir row and
192    /// from the tool name for a native one.
193    pub harness: Option<mj_core::config::HarnessKind>,
194    pub title: String,
195    pub project: String,
196}
197
198/// Start a new session carrying a compacted hand-off from an archived one.
199#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
200#[serde(deny_unknown_fields)]
201pub struct WikiRestoreRequest {
202    /// The SessionWiki session id to restore from.
203    pub wiki_id: String,
204    pub workspace_id: String,
205    pub profile_id: String,
206    pub target_template_id: String,
207    /// Where the new session opens. None takes the project the archived
208    /// session ran in, when that directory still exists.
209    #[serde(default)]
210    pub project_directory: Option<PathBuf>,
211    #[serde(default)]
212    pub additional_mounts: Vec<AdditionalMount>,
213    #[serde(default)]
214    pub resource_allocation: Option<SessionResourceAllocation>,
215}
216
217#[derive(Debug, Clone, Serialize, Deserialize)]
218#[serde(deny_unknown_fields)]
219pub struct WorkspaceSnapshot {
220    pub workspace: WorkspaceRecord,
221    pub sessions: Vec<SessionPreview>,
222    pub drafts: Vec<DraftPreview>,
223}
224
225#[derive(Debug, Clone, Serialize, Deserialize)]
226#[serde(deny_unknown_fields)]
227pub struct RuntimeSessionView {
228    pub session_id: String,
229    pub projection_ordinal: u64,
230    pub projection_digest: String,
231    pub operational: Option<RelayOperationalState>,
232    pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
233    pub connected: bool,
234    pub error: Option<ViewError>,
235}
236
237impl RuntimeSessionView {
238    pub fn from_managed(session_id: String, view: ManagedSessionView) -> Self {
239        let (projection_ordinal, projection_digest, operational, signal) =
240            view.snapshot
241                .map_or((0, String::new(), None, None), |snapshot| {
242                    (
243                        snapshot.materialized.applied_event_ordinal,
244                        snapshot.materialized.applied_event_digest,
245                        Some(snapshot.operational),
246                        snapshot.latest_credential_sync_signal,
247                    )
248                });
249        Self {
250            session_id,
251            projection_ordinal,
252            projection_digest,
253            operational,
254            latest_credential_sync_signal: signal,
255            connected: view.connected,
256            error: view.error,
257        }
258    }
259}
260
261/// Something the daemon did on its own that a surface should report once.
262///
263/// Background work has no lifecycle entry to hang a message on, so notices
264/// travel with the snapshot and carry an id: a surface reports the ones newer
265/// than the last it saw and nothing else, however often it polls.
266#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
267#[serde(deny_unknown_fields)]
268pub struct RuntimeNotice {
269    pub id: u64,
270    pub session_id: String,
271    pub text: String,
272}
273
274#[derive(Debug, Clone, Serialize, Deserialize)]
275#[serde(deny_unknown_fields)]
276pub struct RuntimeSnapshot {
277    #[serde(default)]
278    pub native_agents: Vec<mj_core::native_agent::NativeAgentSummary>,
279    #[serde(default)]
280    pub workspace_names: BTreeMap<String, String>,
281    #[serde(default)]
282    pub moves: Vec<mj_core::state::MoveOperation>,
283    pub revision: u64,
284    pub config: Config,
285    pub records: Vec<SessionRecord>,
286    pub sessions: Vec<RuntimeSessionView>,
287    pub lifecycles: Vec<RuntimeLifecycleView>,
288    /// Reviews the daemon is running, so every surface renders the same one.
289    #[serde(default)]
290    pub reviews: Vec<RuntimeReviewView>,
291    /// Recent background events for this workspace's sessions, oldest first.
292    #[serde(default)]
293    pub notices: Vec<RuntimeNotice>,
294    /// Parent/child relations for the sessions in `records`, so a surface can
295    /// keep a daemon-created child out of the real workspace without a full
296    /// state reload.
297    #[serde(default)]
298    pub subagents: Vec<mj_core::subagent::SubagentRecord>,
299}
300
301#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
302#[serde(rename_all = "snake_case")]
303pub enum RuntimeLifecycleKind {
304    Create,
305    Suspend,
306    Resume,
307    Move,
308    ForceStop,
309    DestroyStopped,
310    ForceDestroy,
311    Cleanup,
312}
313
314#[derive(Debug, Clone, Serialize, Deserialize)]
315#[serde(deny_unknown_fields)]
316pub struct RuntimeLifecycleView {
317    pub operation_id: String,
318    pub cancellable: bool,
319    pub session_id: String,
320    pub kind: RuntimeLifecycleKind,
321    pub started_at_epoch_seconds: u64,
322    pub active_stages: Vec<(ProvisionStage, u64)>,
323    pub resume_destination: Option<(String, String)>,
324    pub notice: Option<String>,
325}
326
327#[derive(Debug, Clone, Serialize, Deserialize)]
328#[serde(deny_unknown_fields)]
329pub struct ResumeSessionRequest {
330    pub session_id: String,
331    pub workspace_id: String,
332    pub profile_id: String,
333    pub target_template_id: String,
334    pub additional_mounts: Option<Vec<AdditionalMount>>,
335    pub resource_allocation: Option<SessionResourceAllocation>,
336    pub discard_queue: bool,
337    pub repository_preflight: Option<ResumeRepositorySourceReceipt>,
338}
339
340#[derive(Debug, Clone, Serialize, Deserialize)]
341#[serde(deny_unknown_fields)]
342pub struct CreateSessionRequest {
343    #[serde(default)]
344    pub create_managed_worktree: Option<bool>,
345    /// None follows the global `[subagents] enabled` setting at launch time.
346    #[serde(default)]
347    pub mjolnir_subagents: Option<bool>,
348    #[serde(default)]
349    pub initial_prompt: Option<String>,
350    pub workspace_id: String,
351    pub profile_id: String,
352    pub bundle_id: String,
353    pub project_directory: Option<PathBuf>,
354    pub target_template_id: String,
355    pub additional_mounts: Vec<AdditionalMount>,
356    pub resource_allocation: Option<SessionResourceAllocation>,
357    pub title: String,
358    pub session_title_override: Option<String>,
359}
360
361#[derive(Debug, Clone, Serialize, Deserialize)]
362#[serde(deny_unknown_fields)]
363pub struct RegisteredSession {
364    pub session: SessionRecord,
365    pub remembered_container_size: Option<(String, HostContainerSize)>,
366}
367
368#[derive(Debug, Clone, Serialize, Deserialize)]
369#[serde(deny_unknown_fields)]
370pub struct DraftPreview {
371    pub id: String,
372    pub session_id: Option<String>,
373    pub source: String,
374    pub owner_pid: Option<u32>,
375    pub saved_at: String,
376}
377
378/// `Ping`, `Status`, and `Stop` form the frozen management subset: their wire
379/// encoding — together with `RequestEnvelope`, `ResponseEnvelope`,
380/// `DaemonStatus`, and `WebViewerStatus` — must never change shape, because
381/// clients and daemons of *any* protocol version rely on them to identify,
382/// stop, and replace each other. `PrepareUpgrade` and its `Done` /
383/// `UpgradePending` replies are also frozen from protocol 33 onward: they
384/// drain accepted work before closing admission atomically. Every other
385/// action may change freely behind a `PROTOCOL_VERSION` bump.
386#[derive(Debug, Clone, Serialize, Deserialize)]
387#[serde(rename_all = "snake_case", tag = "action", content = "arguments")]
388pub enum DaemonAction {
389    NativeAgentHistory {
390        owner: String,
391        child: String,
392        before: Option<(u64, String)>,
393    },
394    Ping,
395    /// Frozen upgrade handshake, available from daemon protocol 33 onward.
396    /// Busy replies leave every operation and control channel running.
397    PrepareUpgrade,
398    Status,
399    WebViewerAccess,
400    RecoverWebViewer(crate::web::WebViewerRecovery),
401    InspectWebListener,
402    ListWorkspaces,
403    CreateWorkspace {
404        name: String,
405    },
406    RenameWorkspace {
407        workspace_id: String,
408        name: String,
409    },
410    TouchWorkspace {
411        workspace_id: String,
412    },
413    CloseWorkspace {
414        workspace_id: String,
415    },
416    CancelWorkspaceClose {
417        workspace_id: String,
418    },
419    DeleteWorkspace {
420        workspace_id: String,
421    },
422    Attach {
423        client_id: String,
424        pid: u32,
425    },
426    Detach {
427        client_id: String,
428    },
429    PersistReadReceipt {
430        client_id: String,
431        workspace_id: String,
432        session_id: String,
433        through: u64,
434    },
435    PersistDetachedSessionState {
436        client_id: String,
437        workspace_id: String,
438        session_id: String,
439        through: u64,
440        owner_pid: u32,
441        draft: mj_core::storage::DetachedSessionDraft,
442    },
443    SaveActiveReview {
444        session_id: String,
445        review: mj_core::storage::StoredReview,
446    },
447    ClearActiveReview {
448        session_id: String,
449    },
450    SaveWorkspacePaneSizes {
451        workspace_id: String,
452        sizes: mj_core::workspace::PaneSizes,
453    },
454    SaveWorkspaceLayout {
455        workspace_id: String,
456        layout: mj_core::workspace::ConversationLayout,
457    },
458    PersistImportedSession {
459        session: Box<SessionRecord>,
460    },
461    SetSessionTitle {
462        session_id: String,
463        title: String,
464    },
465    SetSessionContainerSettings {
466        session_id: String,
467        cpus: Option<String>,
468        memory: Option<String>,
469        mounts: Vec<AdditionalMount>,
470        mount_history: Vec<PathBuf>,
471    },
472    SetSessionAcpTitle {
473        session_id: String,
474        title: Option<String>,
475    },
476    MarkSessionTargetMissing {
477        session_id: String,
478        detail: String,
479        updated_at: String,
480    },
481    CheckpointSession {
482        session_id: String,
483    },
484    /// Search the user's SessionWiki index. An empty query lists the most
485    /// recent sessions.
486    WikiSearch {
487        query: String,
488        limit: usize,
489    },
490    /// The markdown briefing for one indexed session.
491    WikiBrief {
492        wiki_id: String,
493        max_chars: usize,
494    },
495    /// The passages of one indexed session that match a query, with context.
496    WikiHits {
497        wiki_id: String,
498        query: String,
499        context_messages: usize,
500        per_message_chars: usize,
501    },
502    /// What one indexed session is, and what continuing it would mean.
503    WikiSession {
504        wiki_id: String,
505    },
506    /// Start a new session from an archived one's transcript.
507    WikiRestore(WikiRestoreRequest),
508    ScanRecovery {
509        all_instances: bool,
510    },
511    AdoptRecovery {
512        session_id: String,
513        target_id: String,
514        profile: Option<String>,
515        bundle: Option<String>,
516        all_instances: bool,
517    },
518    DestroyRecovery {
519        session_id: String,
520        target_id: String,
521        confirmation: String,
522        all_instances: bool,
523    },
524    Snapshot {
525        workspace_id: String,
526    },
527    RuntimeSnapshot {
528        workspace_id: String,
529        after_revision: u64,
530        #[serde(default)]
531        all_workspaces: bool,
532    },
533    RenameProfile {
534        old_id: String,
535        new_id: String,
536    },
537    RenameTarget {
538        old_id: String,
539        new_id: String,
540    },
541    SubmitSessionCommand {
542        #[serde(default)]
543        inherited_draft: Option<String>,
544        session_id: String,
545        command_id: String,
546        command: RelayCommand,
547    },
548    /// Deliver a prompt typed while a session was still starting, once the
549    /// daemon sees that session's harness become ready. The daemon owns the
550    /// wait, so the prompt arrives whether or not this client is still
551    /// running or still showing that session.
552    QueueStartupPrompt {
553        session_id: String,
554        text: String,
555        /// The saved draft text this prompt was typed from, if the client
556        /// also persisted it. Cleared after a successful submit so the
557        /// delivered prompt does not reappear as a draft.
558        #[serde(default)]
559        inherited_draft: Option<String>,
560    },
561    SyncSession {
562        session_id: String,
563    },
564    RespondElicitation {
565        session_id: String,
566        elicitation_id: String,
567        response: ElicitationResponse,
568    },
569    StopBackgroundTask {
570        session_id: String,
571        background_task_id: String,
572    },
573    /// Drive a session's second-opinion reviewer. The reviewer is a sidecar of
574    /// the session's worker, so it travels the session's own relay rather than
575    /// becoming a session of its own here.
576    ReviewerAction {
577        session_id: String,
578        /// Which reviewing role the action drives; absent means the default
579        /// one, which is what plan review uses.
580        #[serde(default, skip_serializing_if = "Option::is_none")]
581        role: Option<String>,
582        action: crate::session::ReviewerAction,
583    },
584    /// Review the turn this session just finished, on a surface's request.
585    StartTurnReview {
586        session_id: String,
587    },
588    /// Forward, dismiss, or cancel the open review.
589    ResolveTurnReview {
590        session_id: String,
591        resolution: Resolution,
592    },
593    SuspendSession {
594        session_id: String,
595    },
596    StartCreateSession(CreateSessionRequest),
597    WaitCreateSession {
598        session_id: String,
599    },
600    ResumeSession(ResumeSessionRequest),
601    PrepareMoveSession(MoveSelection),
602    MoveSession(MoveSessionRequest),
603    DiscardSinceCheckpoint {
604        session_id: String,
605        checkpoint: mj_core::state::CheckpointMetadata,
606    },
607    DestroyStoppedSession {
608        session_id: String,
609        /// Whether to delete the session's managed git branch as well. The
610        /// branch can hold work the user still wants, so destroying keeps it
611        /// unless the request asks for the deletion.
612        delete_branch: bool,
613    },
614    ForceDestroySession {
615        session_id: String,
616        /// See [`DaemonAction::DestroyStoppedSession`].
617        delete_branch: bool,
618    },
619    ForceDeleteWorkspace {
620        workspace_id: String,
621    },
622    CancelLifecycle {
623        session_id: String,
624    },
625    RecoverDraft {
626        draft_id: String,
627    },
628    Stop,
629}
630
631#[derive(Debug, Serialize, Deserialize)]
632#[serde(deny_unknown_fields)]
633pub struct RequestEnvelope {
634    pub protocol_version: u32,
635    pub request_id: u64,
636    pub token: String,
637    pub action: DaemonAction,
638}
639
640#[derive(Debug, Serialize, Deserialize)]
641#[serde(deny_unknown_fields)]
642pub struct ResponseEnvelope {
643    pub protocol_version: u32,
644    pub request_id: u64,
645    pub result: std::result::Result<DaemonReply, String>,
646}
647
648#[derive(Debug, Clone, Serialize, Deserialize)]
649#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
650pub enum DaemonReply {
651    NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
652    Pong,
653    UpgradePending,
654    Status(DaemonStatus),
655    WebViewerAccess(crate::web::WebViewerAccess),
656    WebListeners(Vec<crate::web::WebListenerProcess>),
657    Workspaces(Vec<WorkspaceListing>),
658    Workspace(WorkspaceRecord),
659    Snapshot(WorkspaceSnapshot),
660    RuntimeSnapshot(Box<RuntimeSnapshot>),
661    RegisteredSession(Box<RegisteredSession>),
662    MovePreparation(Box<MovePreparation>),
663    MoveOutcome(MoveOutcome),
664    Ordinal(u64),
665    Text(String),
666    OptionalSessionState(Option<SessionState>),
667    Checkpoint(mj_core::state::CheckpointMetadata),
668    RecoveryScan(mj_core::state::RecoveryScan),
669    WikiRows(WikiSearchPage),
670    WikiHits(Option<WikiHitTranscript>),
671    WikiSession(Option<Box<WikiSessionInfo>>),
672    Reviewer(Box<crate::session::ReviewerOutcome>),
673    Done,
674}
675
676#[derive(Debug, Clone, Serialize, Deserialize)]
677#[serde(deny_unknown_fields)]
678pub struct DaemonStatus {
679    pub pid: u32,
680    pub started_at: String,
681    pub build_version: String,
682    pub attached_clients: usize,
683    pub phone_status: WebViewerStatus,
684}
685
686#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
687#[serde(rename_all = "snake_case", tag = "state")]
688pub enum WebViewerStatus {
689    Disabled,
690    Starting,
691    Ready {
692        viewer_url: String,
693        viewer_code: String,
694        qr_login_url: Option<String>,
695        fallback_reason: Option<String>,
696    },
697    Stopped,
698    Error {
699        message: String,
700    },
701}
702
703impl std::fmt::Debug for WebViewerStatus {
704    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
705        match self {
706            Self::Ready {
707                viewer_url,
708                viewer_code,
709                fallback_reason,
710                ..
711            } => formatter
712                .debug_struct("Ready")
713                .field("viewer_url", viewer_url)
714                .field("viewer_code", viewer_code)
715                .field("qr_login_url", &"[redacted]")
716                .field("fallback_reason", fallback_reason)
717                .finish(),
718            Self::Disabled => formatter.write_str("Disabled"),
719            Self::Starting => formatter.write_str("Starting"),
720            Self::Stopped => formatter.write_str("Stopped"),
721            Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
722        }
723    }
724}
725
726impl std::fmt::Display for WebViewerStatus {
727    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
728        match self {
729            Self::Disabled => formatter.write_str("disabled"),
730            Self::Starting => formatter.write_str("starting"),
731            Self::Stopped => formatter.write_str("stopped unexpectedly"),
732            Self::Error { message } => write!(formatter, "error: {message}"),
733            Self::Ready {
734                viewer_url,
735                viewer_code,
736                fallback_reason,
737                ..
738            } => {
739                write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
740                if let Some(reason) = fallback_reason {
741                    write!(
742                        formatter,
743                        "; local only because Tailscale HTTPS is unavailable: {reason}"
744                    )?;
745                }
746                Ok(())
747            }
748        }
749    }
750}
751
752/// Whether a non-child process has exited but not yet been reaped.
753///
754/// A zombie still answers `kill(pid, 0)`, because its process-table entry
755/// survives until its parent waits for it — so an existence probe alone calls
756/// it alive forever and anything waiting for it to leave waits forever. That is
757/// exactly the shape of `mj daemon restart` refusing to restart a daemon that
758/// had already stopped: `spawn_detached` used to leave the daemon a child of a
759/// long-lived Mjolnir process that never reaped it. It now double-forks, so the
760/// daemon is init's to reap, but any other unreaped child of this process would
761/// look the same, and the check stays cheap.
762///
763/// Treating a zombie as gone is also safe in the direction that matters: a
764/// zombie's PID cannot be reused until it is reaped, so nothing else can be
765/// occupying that number while this returns true.
766#[cfg(unix)]
767pub fn process_is_zombie(pid: u32) -> bool {
768    let pid = sysinfo::Pid::from_u32(pid);
769    let mut system = sysinfo::System::new();
770    system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
771    system
772        .process(pid)
773        .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
774}
775
776/// Wait for a process to leave, within [`STOP_TIMEOUT`].
777///
778/// The error says the process was still running rather than that it "did not
779/// stop": a daemon that is still winding down has not refused, and the two
780/// read very differently to somebody deciding whether to reach for a kill.
781pub async fn wait_for_exit(pid: u32) -> Result<()> {
782    let deadline = Instant::now() + STOP_TIMEOUT;
783    while daemon_process_is_alive(pid) {
784        ensure!(Instant::now() < deadline, "process {pid} is still running");
785        tokio::time::sleep(RETRY_DELAY).await;
786    }
787    Ok(())
788}
789
790/// Whether a daemon that Mjolnir launched in its own process group is alive.
791///
792/// Reaping is deliberately confined to this daemon-specific path. Attachment
793/// PIDs are merely observations and may alias unrelated children owned by this
794/// process, so their liveness probe below must never call `waitpid`.
795pub fn daemon_process_is_alive(pid: u32) -> bool {
796    #[cfg(unix)]
797    {
798        if pid == 0 {
799            return false;
800        }
801        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
802            return false;
803        };
804        let mut status = 0;
805        // SAFETY: `status` is writable for the call and WNOHANG never blocks.
806        // A non-child fails with ECHILD without changing any process state.
807        let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
808        if waited == raw_pid {
809            return false;
810        }
811        if waited == 0 {
812            return true;
813        }
814        let wait_error = std::io::Error::last_os_error();
815        if wait_error.raw_os_error() != Some(libc::ECHILD) {
816            return true;
817        }
818
819        #[cfg(target_os = "macos")]
820        return owned_daemon_group_is_alive(raw_pid);
821
822        #[cfg(not(target_os = "macos"))]
823        process_is_alive(pid)
824    }
825    #[cfg(not(unix))]
826    process_is_alive(pid)
827}
828
829#[cfg(target_os = "macos")]
830pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
831    // `spawn_detached` makes the daemon a process-group leader. Darwin
832    // excludes zombies from group signal probes: ESRCH means the group is gone
833    // and EPERM means only exiting members remain. The latter is safe here
834    // because this is a group we created for our own same-user child, not an
835    // arbitrary process group.
836    // SAFETY: signal 0 is only an existence probe, and the negative PID targets
837    // the daemon-owned group rather than another process.
838    if unsafe { libc::kill(-pid, 0) } == 0 {
839        return true;
840    }
841    let error = std::io::Error::last_os_error();
842    !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
843}
844
845pub fn process_is_alive(pid: u32) -> bool {
846    #[cfg(unix)]
847    {
848        if pid == 0 {
849            return false;
850        }
851        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
852            return false;
853        };
854        // SAFETY: kill(pid, 0) sends no signal and is the standard existence
855        // probe. EPERM still means the process exists.
856        let result = unsafe { libc::kill(raw_pid, 0) };
857        let exists =
858            result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
859        exists && !process_is_zombie(pid)
860    }
861    #[cfg(not(unix))]
862    {
863        let process_id = sysinfo::Pid::from_u32(pid);
864        let mut system = sysinfo::System::new();
865        system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[process_id]), true);
866        system.process(process_id).is_some()
867    }
868}
869
870pub fn read_metadata() -> Result<DaemonMetadata> {
871    let metadata = read_metadata_any()?;
872    ensure!(
873        metadata.protocol_version == PROTOCOL_VERSION,
874        "daemon protocol {} is incompatible with client protocol {}",
875        metadata.protocol_version,
876        PROTOCOL_VERSION
877    );
878    Ok(metadata)
879}
880
881pub fn read_metadata_any() -> Result<DaemonMetadata> {
882    let path = metadata_path();
883    let body = fs::read(&path).with_context(|| format!("read {}", path.display()))?;
884    let metadata: DaemonMetadata =
885        serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
886    Ok(metadata)
887}
888
889pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
890    let body = serde_json::to_vec(value)?;
891    write_encoded_frame(stream, &body).await
892}
893
894pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
895    ensure!(
896        body.len() <= MAX_FRAME_BYTES,
897        "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
898        body.len()
899    );
900    stream.write_u32(body.len() as u32).await?;
901    stream.write_all(body).await?;
902    stream.flush().await?;
903    Ok(())
904}
905
906pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
907    let length = stream.read_u32().await? as usize;
908    ensure!(
909        length <= MAX_FRAME_BYTES,
910        "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
911    );
912    let mut body = vec![0_u8; length];
913    stream.read_exact(&mut body).await?;
914    serde_json::from_slice(&body).context("decode daemon frame")
915}
916
917pub struct DaemonClient {
918    metadata: DaemonMetadata,
919    stream: TcpStream,
920    next_request_id: u64,
921}
922
923impl DaemonClient {
924    pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
925        let stream =
926            tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
927                .await
928                .context("time out connecting to Mjolnir daemon")??;
929        Ok(Self {
930            metadata,
931            stream,
932            next_request_id: 1,
933        })
934    }
935
936    /// Speak the daemon's advertised dialect, not this build's: management
937    /// requests must reach daemons of any protocol version, and the frozen
938    /// subset encodes identically across all of them.
939    pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
940        self.request_with_reconnect(action, || async {
941            loop {
942                if let Ok(client) = connect_existing().await {
943                    return Ok(client);
944                }
945                tokio::time::sleep(Duration::from_millis(250)).await;
946            }
947        })
948        .await
949    }
950
951    async fn request_with_reconnect<F, Fut>(
952        &mut self,
953        action: DaemonAction,
954        mut reconnect: F,
955    ) -> Result<DaemonReply>
956    where
957        F: FnMut() -> Fut,
958        Fut: Future<Output = Result<Self>>,
959    {
960        loop {
961            let response = self.request_once(action.clone()).await?;
962            if matches!(response, DaemonReply::UpgradePending)
963                && !matches!(action, DaemonAction::PrepareUpgrade)
964            {
965                // Only an explicit refusal guarantees non-admission. Never
966                // replay arbitrary mutations after a lost acknowledgement.
967                tokio::time::sleep(Duration::from_millis(100)).await;
968                *self = reconnect().await?;
969            } else {
970                return Ok(response);
971            }
972        }
973    }
974
975    async fn request_once(&mut self, action: DaemonAction) -> Result<DaemonReply> {
976        let protocol_version = self.metadata.protocol_version;
977        let request_id = self.next_request_id;
978        self.next_request_id += 1;
979        write_frame(
980            &mut self.stream,
981            &RequestEnvelope {
982                protocol_version,
983                request_id,
984                token: self.metadata.token.clone(),
985                action,
986            },
987        )
988        .await?;
989        let response: ResponseEnvelope = read_frame(&mut self.stream).await?;
990        ensure!(
991            response.protocol_version == protocol_version,
992            "daemon changed protocol"
993        );
994        ensure!(
995            response.request_id == request_id,
996            "daemon crossed request IDs"
997        );
998        response.result.map_err(anyhow::Error::msg)
999    }
1000
1001    pub async fn status(&mut self) -> Result<DaemonStatus> {
1002        match self.request(DaemonAction::Status).await? {
1003            DaemonReply::Status(status) => Ok(status),
1004            reply => bail!("unexpected daemon status reply {reply:?}"),
1005        }
1006    }
1007
1008    pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
1009        match self.request(DaemonAction::WebViewerAccess).await? {
1010            DaemonReply::WebViewerAccess(access) => Ok(access),
1011            reply => bail!("unexpected web viewer reply {reply:?}"),
1012        }
1013    }
1014
1015    pub async fn recover_web_viewer(
1016        &mut self,
1017        action: crate::web::WebViewerRecovery,
1018    ) -> Result<()> {
1019        match self.request(DaemonAction::RecoverWebViewer(action)).await? {
1020            DaemonReply::Done => Ok(()),
1021            reply => bail!("unexpected web viewer recovery reply {reply:?}"),
1022        }
1023    }
1024
1025    pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
1026        match self.request(DaemonAction::InspectWebListener).await? {
1027            DaemonReply::WebListeners(processes) => Ok(processes),
1028            reply => bail!("unexpected listener inspection reply {reply:?}"),
1029        }
1030    }
1031
1032    pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
1033        match self.request(DaemonAction::ListWorkspaces).await? {
1034            DaemonReply::Workspaces(workspaces) => Ok(workspaces),
1035            reply => bail!("unexpected daemon workspace reply {reply:?}"),
1036        }
1037    }
1038
1039    pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
1040        match self
1041            .request(DaemonAction::RenameProfile { old_id, new_id })
1042            .await?
1043        {
1044            DaemonReply::Done => Ok(()),
1045            reply => bail!("unexpected rename-profile reply {reply:?}"),
1046        }
1047    }
1048
1049    pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1050        match self
1051            .request(DaemonAction::RenameTarget { old_id, new_id })
1052            .await?
1053        {
1054            DaemonReply::Done => Ok(()),
1055            reply => bail!("unexpected rename-target reply {reply:?}"),
1056        }
1057    }
1058
1059    pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1060        match self.request(DaemonAction::CreateWorkspace { name }).await? {
1061            DaemonReply::Workspace(workspace) => Ok(workspace),
1062            reply => bail!("unexpected create-workspace reply {reply:?}"),
1063        }
1064    }
1065
1066    pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1067        match self
1068            .request(DaemonAction::RenameWorkspace { workspace_id, name })
1069            .await?
1070        {
1071            DaemonReply::Done => Ok(()),
1072            reply => bail!("unexpected rename-workspace reply {reply:?}"),
1073        }
1074    }
1075
1076    pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1077        match self
1078            .request(DaemonAction::TouchWorkspace { workspace_id })
1079            .await?
1080        {
1081            DaemonReply::Done => Ok(()),
1082            reply => bail!("unexpected touch-workspace reply {reply:?}"),
1083        }
1084    }
1085
1086    pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1087        match self
1088            .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1089            .await?
1090        {
1091            DaemonReply::Done => Ok(()),
1092            reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1093        }
1094    }
1095
1096    pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1097        match self
1098            .request(DaemonAction::CloseWorkspace { workspace_id })
1099            .await?
1100        {
1101            DaemonReply::Done => Ok(()),
1102            reply => bail!("unexpected close workspace reply: {reply:?}"),
1103        }
1104    }
1105
1106    pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1107        match self
1108            .request(DaemonAction::DeleteWorkspace { workspace_id })
1109            .await?
1110        {
1111            DaemonReply::Done => Ok(()),
1112            reply => bail!("unexpected delete-workspace reply {reply:?}"),
1113        }
1114    }
1115
1116    pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1117        match self
1118            .request(DaemonAction::Attach { client_id, pid })
1119            .await?
1120        {
1121            DaemonReply::Done => Ok(()),
1122            reply => bail!("unexpected attach reply {reply:?}"),
1123        }
1124    }
1125
1126    pub async fn detach(&mut self, client_id: String) -> Result<()> {
1127        match self.request(DaemonAction::Detach { client_id }).await? {
1128            DaemonReply::Done => Ok(()),
1129            reply => bail!("unexpected detach reply {reply:?}"),
1130        }
1131    }
1132
1133    pub async fn persist_read_receipt(
1134        &mut self,
1135        client_id: String,
1136        workspace_id: String,
1137        session_id: String,
1138        through: u64,
1139    ) -> Result<u64> {
1140        match self
1141            .request(DaemonAction::PersistReadReceipt {
1142                client_id,
1143                workspace_id,
1144                session_id,
1145                through,
1146            })
1147            .await?
1148        {
1149            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1150            reply => bail!("unexpected read-receipt reply {reply:?}"),
1151        }
1152    }
1153
1154    pub async fn persist_detached_session_state(
1155        &mut self,
1156        client_id: String,
1157        workspace_id: String,
1158        session_id: String,
1159        through: u64,
1160        owner_pid: u32,
1161        draft: mj_core::storage::DetachedSessionDraft,
1162    ) -> Result<()> {
1163        match self
1164            .request(DaemonAction::PersistDetachedSessionState {
1165                client_id,
1166                workspace_id,
1167                session_id,
1168                through,
1169                owner_pid,
1170                draft,
1171            })
1172            .await?
1173        {
1174            DaemonReply::Done => Ok(()),
1175            reply => bail!("unexpected detached-session-state reply {reply:?}"),
1176        }
1177    }
1178
1179    pub async fn save_active_review(
1180        &mut self,
1181        session_id: String,
1182        review: mj_core::storage::StoredReview,
1183    ) -> Result<()> {
1184        match self
1185            .request(DaemonAction::SaveActiveReview { session_id, review })
1186            .await?
1187        {
1188            DaemonReply::Done => Ok(()),
1189            reply => bail!("unexpected save-review reply {reply:?}"),
1190        }
1191    }
1192
1193    pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1194        match self
1195            .request(DaemonAction::ClearActiveReview { session_id })
1196            .await?
1197        {
1198            DaemonReply::Done => Ok(()),
1199            reply => bail!("unexpected clear-review reply {reply:?}"),
1200        }
1201    }
1202
1203    pub async fn save_workspace_pane_sizes(
1204        &mut self,
1205        workspace_id: String,
1206        sizes: mj_core::workspace::PaneSizes,
1207    ) -> Result<()> {
1208        match self
1209            .request(DaemonAction::SaveWorkspacePaneSizes {
1210                workspace_id,
1211                sizes,
1212            })
1213            .await?
1214        {
1215            DaemonReply::Done => Ok(()),
1216            reply => bail!("unexpected pane-size save reply {reply:?}"),
1217        }
1218    }
1219
1220    pub async fn save_workspace_layout(
1221        &mut self,
1222        workspace_id: String,
1223        layout: mj_core::workspace::ConversationLayout,
1224    ) -> Result<()> {
1225        match self
1226            .request(DaemonAction::SaveWorkspaceLayout {
1227                workspace_id,
1228                layout,
1229            })
1230            .await?
1231        {
1232            DaemonReply::Done => Ok(()),
1233            reply => bail!("unexpected layout save reply {reply:?}"),
1234        }
1235    }
1236
1237    pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1238        match self
1239            .request(DaemonAction::PersistImportedSession {
1240                session: Box::new(session),
1241            })
1242            .await?
1243        {
1244            DaemonReply::Done => Ok(()),
1245            reply => bail!("unexpected imported-session reply {reply:?}"),
1246        }
1247    }
1248
1249    pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1250        match self
1251            .request(DaemonAction::SetSessionTitle { session_id, title })
1252            .await?
1253        {
1254            DaemonReply::Text(title) => Ok(title),
1255            reply => bail!("unexpected session-title reply {reply:?}"),
1256        }
1257    }
1258
1259    pub async fn set_session_container_settings(
1260        &mut self,
1261        session_id: String,
1262        cpus: Option<String>,
1263        memory: Option<String>,
1264        mounts: Vec<AdditionalMount>,
1265        mount_history: Vec<PathBuf>,
1266    ) -> Result<()> {
1267        match self
1268            .request(DaemonAction::SetSessionContainerSettings {
1269                session_id,
1270                cpus,
1271                memory,
1272                mounts,
1273                mount_history,
1274            })
1275            .await?
1276        {
1277            DaemonReply::Done => Ok(()),
1278            reply => bail!("unexpected container-settings reply {reply:?}"),
1279        }
1280    }
1281
1282    pub async fn set_session_acp_title(
1283        &mut self,
1284        session_id: String,
1285        title: Option<String>,
1286    ) -> Result<()> {
1287        match self
1288            .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1289            .await?
1290        {
1291            DaemonReply::Done => Ok(()),
1292            reply => bail!("unexpected ACP-title reply {reply:?}"),
1293        }
1294    }
1295
1296    pub async fn mark_session_target_missing(
1297        &mut self,
1298        session_id: String,
1299        detail: String,
1300        updated_at: String,
1301    ) -> Result<Option<SessionState>> {
1302        match self
1303            .request(DaemonAction::MarkSessionTargetMissing {
1304                session_id,
1305                detail,
1306                updated_at,
1307            })
1308            .await?
1309        {
1310            DaemonReply::OptionalSessionState(state) => Ok(state),
1311            reply => bail!("unexpected target-missing reply {reply:?}"),
1312        }
1313    }
1314
1315    pub async fn checkpoint_session(
1316        &mut self,
1317        session_id: String,
1318    ) -> Result<mj_core::state::CheckpointMetadata> {
1319        match self
1320            .request(DaemonAction::CheckpointSession { session_id })
1321            .await?
1322        {
1323            DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1324            reply => bail!("unexpected checkpoint reply {reply:?}"),
1325        }
1326    }
1327
1328    /// Search the user's SessionWiki index, newest first when the query is
1329    /// empty and best match first otherwise. The reply carries the state of
1330    /// the index as well as the rows, so a caller can say the first build is
1331    /// still running.
1332    pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1333        match self
1334            .request(DaemonAction::WikiSearch { query, limit })
1335            .await?
1336        {
1337            DaemonReply::WikiRows(page) => Ok(page),
1338            reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1339        }
1340    }
1341
1342    /// The markdown briefing for one indexed session.
1343    pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1344        match self
1345            .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1346            .await?
1347        {
1348            DaemonReply::Text(markdown) => Ok(markdown),
1349            reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1350        }
1351    }
1352
1353    /// The passages of one indexed session that match a query, each matching
1354    /// message with `context_messages` neighbours on either side and its text
1355    /// capped at `per_message_chars`. `None` when the index holds no session
1356    /// with that id.
1357    pub async fn wiki_hits(
1358        &mut self,
1359        wiki_id: String,
1360        query: String,
1361        context_messages: usize,
1362        per_message_chars: usize,
1363    ) -> Result<Option<WikiHitTranscript>> {
1364        match self
1365            .request(DaemonAction::WikiHits {
1366                wiki_id,
1367                query,
1368                context_messages,
1369                per_message_chars,
1370            })
1371            .await?
1372        {
1373            DaemonReply::WikiHits(transcript) => Ok(transcript),
1374            reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1375        }
1376    }
1377
1378    /// What the index knows about one session, or `None` when the index holds
1379    /// no session with that id.
1380    pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1381        match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1382            DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1383            reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1384        }
1385    }
1386
1387    /// Start a new session carrying a hand-off compacted from an archived one.
1388    /// It answers like any other session start: the record exists and is
1389    /// provisioning, and the hand-off follows once the harness is ready.
1390    pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1391        match self.request(DaemonAction::WikiRestore(request)).await? {
1392            DaemonReply::RegisteredSession(registered) => Ok(*registered),
1393            reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1394        }
1395    }
1396
1397    pub async fn scan_recovery(
1398        &mut self,
1399        all_instances: bool,
1400    ) -> Result<mj_core::state::RecoveryScan> {
1401        match self
1402            .request(DaemonAction::ScanRecovery { all_instances })
1403            .await?
1404        {
1405            DaemonReply::RecoveryScan(scan) => Ok(scan),
1406            reply => bail!("unexpected recovery-scan reply {reply:?}"),
1407        }
1408    }
1409
1410    pub async fn adopt_recovery(
1411        &mut self,
1412        session_id: String,
1413        target_id: String,
1414        profile: Option<String>,
1415        bundle: Option<String>,
1416        all_instances: bool,
1417    ) -> Result<()> {
1418        match self
1419            .request(DaemonAction::AdoptRecovery {
1420                session_id,
1421                target_id,
1422                profile,
1423                bundle,
1424                all_instances,
1425            })
1426            .await?
1427        {
1428            DaemonReply::Done => Ok(()),
1429            reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1430        }
1431    }
1432
1433    pub async fn destroy_recovery(
1434        &mut self,
1435        session_id: String,
1436        target_id: String,
1437        confirmation: String,
1438        all_instances: bool,
1439    ) -> Result<()> {
1440        match self
1441            .request(DaemonAction::DestroyRecovery {
1442                session_id,
1443                target_id,
1444                confirmation,
1445                all_instances,
1446            })
1447            .await?
1448        {
1449            DaemonReply::Done => Ok(()),
1450            reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1451        }
1452    }
1453
1454    pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1455        match self
1456            .request(DaemonAction::Snapshot { workspace_id })
1457            .await?
1458        {
1459            DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1460            reply => bail!("unexpected snapshot reply {reply:?}"),
1461        }
1462    }
1463
1464    pub async fn runtime_snapshot(
1465        &mut self,
1466        workspace_id: String,
1467        after_revision: u64,
1468        all_workspaces: bool,
1469    ) -> Result<RuntimeSnapshot> {
1470        match self
1471            .request(DaemonAction::RuntimeSnapshot {
1472                workspace_id,
1473                after_revision,
1474                all_workspaces,
1475            })
1476            .await?
1477        {
1478            DaemonReply::RuntimeSnapshot(snapshot) => Ok(*snapshot),
1479            reply => bail!("unexpected runtime snapshot reply {reply:?}"),
1480        }
1481    }
1482
1483    pub async fn submit_session_command(
1484        &mut self,
1485        session_id: String,
1486        command_id: String,
1487        command: RelayCommand,
1488        inherited_draft: Option<String>,
1489    ) -> Result<u64> {
1490        match self
1491            .request(DaemonAction::SubmitSessionCommand {
1492                inherited_draft,
1493                session_id,
1494                command_id,
1495                command,
1496            })
1497            .await?
1498        {
1499            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1500            reply => bail!("unexpected session command reply {reply:?}"),
1501        }
1502    }
1503
1504    /// Hand the daemon a prompt for a session that is still starting. The
1505    /// daemon replies as soon as the prompt is queued, not when it is
1506    /// delivered; delivery failures come back as a session notice.
1507    pub async fn queue_startup_prompt(
1508        &mut self,
1509        session_id: String,
1510        text: String,
1511        inherited_draft: Option<String>,
1512    ) -> Result<()> {
1513        match self
1514            .request(DaemonAction::QueueStartupPrompt {
1515                session_id,
1516                text,
1517                inherited_draft,
1518            })
1519            .await?
1520        {
1521            DaemonReply::Done => Ok(()),
1522            reply => bail!("unexpected startup prompt reply {reply:?}"),
1523        }
1524    }
1525
1526    /// Ask the daemon to review the turn this session just finished.
1527    ///
1528    /// The refusal is a sentence for a person -- "prompts are queued", "set
1529    /// [review] profile in config.toml" -- so it travels as text rather than
1530    /// as a code every surface would have to translate.
1531    pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
1532        match self
1533            .request(DaemonAction::StartTurnReview { session_id })
1534            .await?
1535        {
1536            DaemonReply::Done => Ok(()),
1537            reply => bail!("unexpected turn-review reply {reply:?}"),
1538        }
1539    }
1540
1541    pub async fn resolve_turn_review(
1542        &mut self,
1543        session_id: String,
1544        resolution: Resolution,
1545    ) -> Result<()> {
1546        match self
1547            .request(DaemonAction::ResolveTurnReview {
1548                session_id,
1549                resolution,
1550            })
1551            .await?
1552        {
1553            DaemonReply::Done => Ok(()),
1554            reply => bail!("unexpected turn-review resolution reply {reply:?}"),
1555        }
1556    }
1557
1558    pub async fn reviewer_action(
1559        &mut self,
1560        session_id: String,
1561        role: Option<String>,
1562        action: crate::session::ReviewerAction,
1563    ) -> Result<crate::session::ReviewerOutcome> {
1564        match self
1565            .request(DaemonAction::ReviewerAction {
1566                session_id,
1567                role,
1568                action,
1569            })
1570            .await?
1571        {
1572            DaemonReply::Reviewer(outcome) => Ok(*outcome),
1573            reply => bail!("unexpected reviewer reply {reply:?}"),
1574        }
1575    }
1576
1577    pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
1578        match self
1579            .request(DaemonAction::SyncSession { session_id })
1580            .await?
1581        {
1582            DaemonReply::Done => Ok(()),
1583            reply => bail!("unexpected session sync reply {reply:?}"),
1584        }
1585    }
1586
1587    pub async fn respond_elicitation(
1588        &mut self,
1589        session_id: String,
1590        elicitation_id: String,
1591        response: ElicitationResponse,
1592    ) -> Result<()> {
1593        match self
1594            .request(DaemonAction::RespondElicitation {
1595                session_id,
1596                elicitation_id,
1597                response,
1598            })
1599            .await?
1600        {
1601            DaemonReply::Done => Ok(()),
1602            reply => bail!("unexpected elicitation reply {reply:?}"),
1603        }
1604    }
1605
1606    pub async fn native_agent_history(
1607        &mut self,
1608        owner: String,
1609        child: String,
1610        before: Option<(u64, String)>,
1611    ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
1612        match self
1613            .request(DaemonAction::NativeAgentHistory {
1614                owner,
1615                child,
1616                before,
1617            })
1618            .await?
1619        {
1620            DaemonReply::NativeAgentHistory(page) => Ok(page),
1621            reply => bail!("unexpected native agent history reply {reply:?}"),
1622        }
1623    }
1624
1625    pub async fn stop_background_task(
1626        &mut self,
1627        session_id: String,
1628        background_task_id: String,
1629    ) -> Result<()> {
1630        match self
1631            .request(DaemonAction::StopBackgroundTask {
1632                session_id,
1633                background_task_id,
1634            })
1635            .await?
1636        {
1637            DaemonReply::Done => Ok(()),
1638            reply => bail!("unexpected background task stop reply {reply:?}"),
1639        }
1640    }
1641
1642    pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
1643        match self
1644            .request(DaemonAction::SuspendSession { session_id })
1645            .await?
1646        {
1647            DaemonReply::Done => Ok(()),
1648            reply => bail!("unexpected close-session reply {reply:?}"),
1649        }
1650    }
1651
1652    pub async fn start_create_session(
1653        &mut self,
1654        request: CreateSessionRequest,
1655    ) -> Result<RegisteredSession> {
1656        match self
1657            .request(DaemonAction::StartCreateSession(request))
1658            .await?
1659        {
1660            DaemonReply::RegisteredSession(registered) => Ok(*registered),
1661            reply => bail!("unexpected start-create reply {reply:?}"),
1662        }
1663    }
1664
1665    pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
1666        match self
1667            .request(DaemonAction::WaitCreateSession { session_id })
1668            .await?
1669        {
1670            DaemonReply::Done => Ok(()),
1671            reply => bail!("unexpected wait-create reply {reply:?}"),
1672        }
1673    }
1674
1675    pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
1676        match self.request(DaemonAction::ResumeSession(request)).await? {
1677            DaemonReply::Done => Ok(()),
1678            reply => bail!("unexpected resume-session reply {reply:?}"),
1679        }
1680    }
1681
1682    pub async fn discard_since_checkpoint(
1683        &mut self,
1684        session_id: String,
1685        checkpoint: mj_core::state::CheckpointMetadata,
1686    ) -> Result<()> {
1687        match self
1688            .request(DaemonAction::DiscardSinceCheckpoint {
1689                session_id,
1690                checkpoint,
1691            })
1692            .await?
1693        {
1694            DaemonReply::Done => Ok(()),
1695            reply => bail!("unexpected force-stop reply {reply:?}"),
1696        }
1697    }
1698
1699    pub async fn destroy_stopped_session(
1700        &mut self,
1701        session_id: String,
1702        delete_branch: bool,
1703    ) -> Result<()> {
1704        match self
1705            .request(DaemonAction::DestroyStoppedSession {
1706                session_id,
1707                delete_branch,
1708            })
1709            .await?
1710        {
1711            DaemonReply::Done => Ok(()),
1712            reply => bail!("unexpected destroy-stopped reply {reply:?}"),
1713        }
1714    }
1715
1716    pub async fn force_destroy_session(
1717        &mut self,
1718        session_id: String,
1719        delete_branch: bool,
1720    ) -> Result<()> {
1721        match self
1722            .request(DaemonAction::ForceDestroySession {
1723                session_id,
1724                delete_branch,
1725            })
1726            .await?
1727        {
1728            DaemonReply::Done => Ok(()),
1729            reply => bail!("unexpected force-destroy reply {reply:?}"),
1730        }
1731    }
1732
1733    pub async fn force_delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1734        match self
1735            .request(DaemonAction::ForceDeleteWorkspace { workspace_id })
1736            .await?
1737        {
1738            DaemonReply::Done => Ok(()),
1739            reply => bail!("unexpected force-delete-workspace reply {reply:?}"),
1740        }
1741    }
1742
1743    pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
1744        match self
1745            .request(DaemonAction::CancelLifecycle { session_id })
1746            .await?
1747        {
1748            DaemonReply::Done => Ok(()),
1749            reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
1750        }
1751    }
1752
1753    pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
1754        match self
1755            .request(DaemonAction::RecoverDraft { draft_id })
1756            .await?
1757        {
1758            DaemonReply::Done => Ok(()),
1759            reply => bail!("unexpected recover-draft reply {reply:?}"),
1760        }
1761    }
1762
1763    pub async fn stop(&mut self) -> Result<()> {
1764        match self.request(DaemonAction::Stop).await? {
1765            DaemonReply::Done => Ok(()),
1766            reply => bail!("unexpected stop reply {reply:?}"),
1767        }
1768    }
1769}
1770
1771pub async fn connect_existing() -> Result<DaemonClient> {
1772    let metadata = tokio::task::spawn_blocking(read_metadata)
1773        .await
1774        .context("read daemon metadata task failed")??;
1775    DaemonClient::connect(metadata).await
1776}
1777
1778/// A handle to whatever daemon the metadata file advertises, regardless of its
1779/// protocol version. It only exposes the frozen management subset (`Ping`,
1780/// `Status`, `Stop`), which encodes identically in every protocol version.
1781pub struct ManagementClient {
1782    inner: DaemonClient,
1783}
1784
1785impl ManagementClient {
1786    pub fn new(inner: DaemonClient) -> Self {
1787        Self { inner }
1788    }
1789    pub fn protocol_version(&self) -> u32 {
1790        self.inner.metadata.protocol_version
1791    }
1792
1793    pub async fn status(&mut self) -> Result<DaemonStatus> {
1794        self.inner.status().await
1795    }
1796
1797    pub async fn stop(&mut self) -> Result<()> {
1798        tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
1799            .await
1800            .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
1801    }
1802
1803    /// Ask the daemon to stop and wait for its process to actually exit.
1804    pub async fn stop_and_wait(mut self) -> Result<()> {
1805        let pid = self.inner.metadata.pid;
1806        self.stop().await?;
1807        wait_for_exit(pid).await.with_context(|| {
1808            format!(
1809                "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
1810                STOP_TIMEOUT.as_secs()
1811            )
1812        })
1813    }
1814}
1815
1816pub async fn connect_management() -> Result<ManagementClient> {
1817    Ok(ManagementClient {
1818        inner: DaemonClient::connect(read_metadata_any()?).await?,
1819    })
1820}
1821
1822/// Refuse to speak to a daemon whose protocol this build does not know.
1823///
1824/// The two protocol numbers alone do not say which `mj` ran. The usual cause is
1825/// a second installation: a `cargo install`ed client sits earlier on PATH than
1826/// the build whose daemon is running, so every command fails here while the
1827/// other binary works, and nothing in the message says where either lives. It
1828/// therefore names both executables and both versions.
1829pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
1830    ensure!(
1831        metadata.protocol_version <= PROTOCOL_VERSION,
1832        "{}",
1833        unsupported_daemon_protocol_message(
1834            metadata.protocol_version,
1835            &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
1836        )
1837    );
1838    Ok(())
1839}
1840
1841/// The message [`ensure_supported_daemon_protocol`] fails with, given the
1842/// sentence that names both builds, so it can be read without a daemon.
1843fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
1844    format!(
1845        "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
1846         Put the daemon's directory first on PATH, or reinstall this client from that build."
1847    )
1848}
1849pub const PROTOCOL_VERSION: u32 = 33;
1850pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
1851/// How long a daemon is given to exit after it accepts a stop.
1852///
1853/// Stopping cancels a token and returns immediately; the daemon then unwinds
1854/// its session manager, its phone server and its pollers. That is normally
1855/// fast, but a daemon whose database has been migrated out from under it fails
1856/// every read while it winds down and has been observed taking over five
1857/// seconds — which the previous five-second bound missed by a fraction,
1858/// reporting a stop that had in fact worked as `did not stop` and aborting the
1859/// restart that depended on it.
1860pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
1861pub const RETRY_DELAY: Duration = Duration::from_millis(40);
1862impl DaemonClient {
1863    pub async fn prepare_move_session(
1864        &mut self,
1865        selection: MoveSelection,
1866    ) -> Result<MovePreparation> {
1867        match self
1868            .request(DaemonAction::PrepareMoveSession(selection))
1869            .await?
1870        {
1871            DaemonReply::MovePreparation(preparation) => Ok(*preparation),
1872            _ => bail!("daemon returned an unexpected move preparation reply"),
1873        }
1874    }
1875
1876    pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
1877        match self.request(DaemonAction::MoveSession(request)).await? {
1878            DaemonReply::MoveOutcome(outcome) => Ok(outcome),
1879            _ => bail!("daemon returned an unexpected move reply"),
1880        }
1881    }
1882}
1883
1884#[cfg(test)]
1885mod tests {
1886    use super::*;
1887    use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
1888    use std::path::Path;
1889
1890    #[tokio::test]
1891    async fn upgrade_refusal_retries_the_identical_command_but_lost_acknowledgements_do_not() {
1892        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1893        let metadata = DaemonMetadata {
1894            protocol_version: PROTOCOL_VERSION,
1895            pid: std::process::id(),
1896            address: listener.local_addr().unwrap(),
1897            token: "test-token".into(),
1898            started_at: "test".into(),
1899            build_version: env!("CARGO_PKG_VERSION").into(),
1900        };
1901        let action = DaemonAction::SubmitSessionCommand {
1902            inherited_draft: None,
1903            session_id: "test-session".into(),
1904            command_id: "steer-command".into(),
1905            command: RelayCommand::Steer {
1906                active_prompt_id: "active-command".into(),
1907                queued_prompt_id: "queued-command".into(),
1908            },
1909        };
1910        let expected = serde_json::to_value(&action).unwrap();
1911        let server = tokio::spawn(async move {
1912            for reply in [
1913                Some(DaemonReply::UpgradePending),
1914                Some(DaemonReply::Done),
1915                None,
1916            ] {
1917                let (mut stream, _) = listener.accept().await.unwrap();
1918                let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
1919                assert_eq!(serde_json::to_value(&request.action).unwrap(), expected);
1920                if let Some(reply) = reply {
1921                    write_frame(
1922                        &mut stream,
1923                        &ResponseEnvelope {
1924                            protocol_version: request.protocol_version,
1925                            request_id: request.request_id,
1926                            result: Ok(reply),
1927                        },
1928                    )
1929                    .await
1930                    .unwrap();
1931                }
1932            }
1933        });
1934        let mut client = DaemonClient::connect(metadata.clone()).await.unwrap();
1935        let reply = client
1936            .request_with_reconnect(action.clone(), || DaemonClient::connect(metadata.clone()))
1937            .await
1938            .unwrap();
1939        assert!(matches!(reply, DaemonReply::Done));
1940        let mut client = DaemonClient::connect(metadata).await.unwrap();
1941        assert!(
1942            client
1943                .request_with_reconnect(action, || async {
1944                    panic!("an ambiguous acknowledgement must not replay a mutation");
1945                })
1946                .await
1947                .is_err()
1948        );
1949        server.await.unwrap();
1950    }
1951
1952    #[test]
1953    fn unsupported_protocol_message_names_both_binaries_and_versions() {
1954        let message = unsupported_daemon_protocol_message(
1955            PROTOCOL_VERSION + 5,
1956            &describe_daemon_and_client_builds(
1957                4242,
1958                BuildDescription {
1959                    executable: Some(Path::new("/home/dev/mj/target/release/mj")),
1960                    version: "2.14.0",
1961                },
1962                BuildDescription {
1963                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
1964                    version: "2.9.0",
1965                },
1966            ),
1967        );
1968        assert_eq!(
1969            message,
1970            format!(
1971                "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
1972                 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
1973                 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
1974                 Put the daemon's directory first on PATH, or reinstall this client from that build.",
1975                PROTOCOL_VERSION + 5
1976            )
1977        );
1978    }
1979
1980    #[test]
1981    fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
1982        let message = unsupported_daemon_protocol_message(
1983            PROTOCOL_VERSION + 1,
1984            &describe_daemon_and_client_builds(
1985                4242,
1986                BuildDescription {
1987                    executable: None,
1988                    version: "2.14.0",
1989                },
1990                BuildDescription {
1991                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
1992                    version: "2.9.0",
1993                },
1994            ),
1995        );
1996        assert!(
1997            message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
1998            "{message}"
1999        );
2000        assert!(
2001            message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
2002            "{message}"
2003        );
2004    }
2005}