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. Every other action may change freely behind a
383/// `PROTOCOL_VERSION` bump.
384#[derive(Debug, Clone, Serialize, Deserialize)]
385#[serde(rename_all = "snake_case", tag = "action", content = "arguments")]
386pub enum DaemonAction {
387    JevDecisions {
388        session_id: String,
389        decision_id: Option<String>,
390    },
391    NativeAgentHistory {
392        owner: String,
393        child: String,
394        before: Option<(u64, String)>,
395    },
396    Ping,
397    Status,
398    WebViewerAccess,
399    RecoverWebViewer(crate::web::WebViewerRecovery),
400    InspectWebListener,
401    ListWorkspaces,
402    CreateWorkspace {
403        name: String,
404    },
405    RenameWorkspace {
406        workspace_id: String,
407        name: String,
408    },
409    TouchWorkspace {
410        workspace_id: String,
411    },
412    CloseWorkspace {
413        workspace_id: String,
414    },
415    CancelWorkspaceClose {
416        workspace_id: String,
417    },
418    DeleteWorkspace {
419        workspace_id: String,
420    },
421    Attach {
422        client_id: String,
423        pid: u32,
424    },
425    Detach {
426        client_id: String,
427    },
428    PersistReadReceipt {
429        client_id: String,
430        workspace_id: String,
431        session_id: String,
432        through: u64,
433    },
434    PersistDetachedSessionState {
435        client_id: String,
436        workspace_id: String,
437        session_id: String,
438        through: u64,
439        owner_pid: u32,
440        draft: mj_core::storage::DetachedSessionDraft,
441    },
442    SaveActiveReview {
443        session_id: String,
444        review: mj_core::storage::StoredReview,
445    },
446    ClearActiveReview {
447        session_id: String,
448    },
449    SaveWorkspacePaneSizes {
450        workspace_id: String,
451        sizes: mj_core::workspace::PaneSizes,
452    },
453    SaveWorkspaceLayout {
454        workspace_id: String,
455        layout: mj_core::workspace::ConversationLayout,
456    },
457    PersistImportedSession {
458        session: Box<SessionRecord>,
459    },
460    SetSessionTitle {
461        session_id: String,
462        title: String,
463    },
464    SetSessionContainerSettings {
465        session_id: String,
466        cpus: Option<String>,
467        memory: Option<String>,
468        mounts: Vec<AdditionalMount>,
469        mount_history: Vec<PathBuf>,
470    },
471    SetSessionAcpTitle {
472        session_id: String,
473        title: Option<String>,
474    },
475    MarkSessionTargetMissing {
476        session_id: String,
477        detail: String,
478        updated_at: String,
479    },
480    CheckpointSession {
481        session_id: String,
482    },
483    /// Search the user's SessionWiki index. An empty query lists the most
484    /// recent sessions.
485    WikiSearch {
486        query: String,
487        limit: usize,
488    },
489    /// The markdown briefing for one indexed session.
490    WikiBrief {
491        wiki_id: String,
492        max_chars: usize,
493    },
494    /// The passages of one indexed session that match a query, with context.
495    WikiHits {
496        wiki_id: String,
497        query: String,
498        context_messages: usize,
499        per_message_chars: usize,
500    },
501    /// What one indexed session is, and what continuing it would mean.
502    WikiSession {
503        wiki_id: String,
504    },
505    /// Start a new session from an archived one's transcript.
506    WikiRestore(WikiRestoreRequest),
507    ScanRecovery {
508        all_instances: bool,
509    },
510    AdoptRecovery {
511        session_id: String,
512        target_id: String,
513        profile: Option<String>,
514        bundle: Option<String>,
515        all_instances: bool,
516    },
517    DestroyRecovery {
518        session_id: String,
519        target_id: String,
520        confirmation: String,
521        all_instances: bool,
522    },
523    Snapshot {
524        workspace_id: String,
525    },
526    RuntimeSnapshot {
527        workspace_id: String,
528        after_revision: u64,
529        #[serde(default)]
530        all_workspaces: bool,
531    },
532    RenameProfile {
533        old_id: String,
534        new_id: String,
535    },
536    RenameTarget {
537        old_id: String,
538        new_id: String,
539    },
540    SubmitSessionCommand {
541        #[serde(default)]
542        inherited_draft: Option<String>,
543        session_id: String,
544        command_id: String,
545        command: RelayCommand,
546    },
547    /// Deliver a prompt typed while a session was still starting, once the
548    /// daemon sees that session's harness become ready. The daemon owns the
549    /// wait, so the prompt arrives whether or not this client is still
550    /// running or still showing that session.
551    QueueStartupPrompt {
552        session_id: String,
553        text: String,
554        /// The saved draft text this prompt was typed from, if the client
555        /// also persisted it. Cleared after a successful submit so the
556        /// delivered prompt does not reappear as a draft.
557        #[serde(default)]
558        inherited_draft: Option<String>,
559    },
560    SyncSession {
561        session_id: String,
562    },
563    RespondElicitation {
564        session_id: String,
565        elicitation_id: String,
566        response: ElicitationResponse,
567    },
568    StopBackgroundTask {
569        session_id: String,
570        background_task_id: String,
571    },
572    /// Drive a session's second-opinion reviewer. The reviewer is a sidecar of
573    /// the session's worker, so it travels the session's own relay rather than
574    /// becoming a session of its own here.
575    ReviewerAction {
576        session_id: String,
577        /// Which reviewing role the action drives; absent means the default
578        /// one, which is what plan review uses.
579        #[serde(default, skip_serializing_if = "Option::is_none")]
580        role: Option<String>,
581        action: crate::session::ReviewerAction,
582    },
583    /// Review the turn this session just finished, on a surface's request.
584    StartTurnReview {
585        session_id: String,
586    },
587    /// Forward, dismiss, or cancel the open review.
588    ResolveTurnReview {
589        session_id: String,
590        resolution: Resolution,
591    },
592    SuspendSession {
593        session_id: String,
594    },
595    StartCreateSession(CreateSessionRequest),
596    WaitCreateSession {
597        session_id: String,
598    },
599    ResumeSession(ResumeSessionRequest),
600    PrepareMoveSession(MoveSelection),
601    MoveSession(MoveSessionRequest),
602    DiscardSinceCheckpoint {
603        session_id: String,
604        checkpoint: mj_core::state::CheckpointMetadata,
605    },
606    DestroyStoppedSession {
607        session_id: String,
608        /// Whether to delete the session's managed git branch as well. The
609        /// branch can hold work the user still wants, so destroying keeps it
610        /// unless the request asks for the deletion.
611        delete_branch: bool,
612    },
613    ForceDestroySession {
614        session_id: String,
615        /// See [`DaemonAction::DestroyStoppedSession`].
616        delete_branch: bool,
617    },
618    ForceDeleteWorkspace {
619        workspace_id: String,
620    },
621    CancelLifecycle {
622        session_id: String,
623    },
624    RecoverDraft {
625        draft_id: String,
626    },
627    Stop,
628}
629
630#[derive(Debug, Serialize, Deserialize)]
631#[serde(deny_unknown_fields)]
632pub struct RequestEnvelope {
633    pub protocol_version: u32,
634    pub request_id: u64,
635    pub token: String,
636    pub action: DaemonAction,
637}
638
639#[derive(Debug, Serialize, Deserialize)]
640#[serde(deny_unknown_fields)]
641pub struct ResponseEnvelope {
642    pub protocol_version: u32,
643    pub request_id: u64,
644    pub result: std::result::Result<DaemonReply, String>,
645}
646
647#[derive(Debug, Clone, Serialize, Deserialize)]
648#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
649pub enum DaemonReply {
650    JevDecisions(mj_core::jev::DecisionPage),
651    NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
652    Pong,
653    Status(DaemonStatus),
654    WebViewerAccess(crate::web::WebViewerAccess),
655    WebListeners(Vec<crate::web::WebListenerProcess>),
656    Workspaces(Vec<WorkspaceListing>),
657    Workspace(WorkspaceRecord),
658    Snapshot(WorkspaceSnapshot),
659    RuntimeSnapshot(Box<RuntimeSnapshot>),
660    RegisteredSession(Box<RegisteredSession>),
661    MovePreparation(Box<MovePreparation>),
662    MoveOutcome(MoveOutcome),
663    Ordinal(u64),
664    Text(String),
665    OptionalSessionState(Option<SessionState>),
666    Checkpoint(mj_core::state::CheckpointMetadata),
667    RecoveryScan(mj_core::state::RecoveryScan),
668    WikiRows(WikiSearchPage),
669    WikiHits(Option<WikiHitTranscript>),
670    WikiSession(Option<Box<WikiSessionInfo>>),
671    Reviewer(Box<crate::session::ReviewerOutcome>),
672    Done,
673}
674
675#[derive(Debug, Clone, Serialize, Deserialize)]
676#[serde(deny_unknown_fields)]
677pub struct DaemonStatus {
678    pub pid: u32,
679    pub started_at: String,
680    pub build_version: String,
681    pub attached_clients: usize,
682    pub phone_status: WebViewerStatus,
683}
684
685#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
686#[serde(rename_all = "snake_case", tag = "state")]
687pub enum WebViewerStatus {
688    Disabled,
689    Starting,
690    Ready {
691        viewer_url: String,
692        viewer_code: String,
693        qr_login_url: Option<String>,
694        fallback_reason: Option<String>,
695    },
696    Stopped,
697    Error {
698        message: String,
699    },
700}
701
702impl std::fmt::Debug for WebViewerStatus {
703    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
704        match self {
705            Self::Ready {
706                viewer_url,
707                viewer_code,
708                fallback_reason,
709                ..
710            } => formatter
711                .debug_struct("Ready")
712                .field("viewer_url", viewer_url)
713                .field("viewer_code", viewer_code)
714                .field("qr_login_url", &"[redacted]")
715                .field("fallback_reason", fallback_reason)
716                .finish(),
717            Self::Disabled => formatter.write_str("Disabled"),
718            Self::Starting => formatter.write_str("Starting"),
719            Self::Stopped => formatter.write_str("Stopped"),
720            Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
721        }
722    }
723}
724
725impl std::fmt::Display for WebViewerStatus {
726    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
727        match self {
728            Self::Disabled => formatter.write_str("disabled"),
729            Self::Starting => formatter.write_str("starting"),
730            Self::Stopped => formatter.write_str("stopped unexpectedly"),
731            Self::Error { message } => write!(formatter, "error: {message}"),
732            Self::Ready {
733                viewer_url,
734                viewer_code,
735                fallback_reason,
736                ..
737            } => {
738                write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
739                if let Some(reason) = fallback_reason {
740                    write!(
741                        formatter,
742                        "; local only because Tailscale HTTPS is unavailable: {reason}"
743                    )?;
744                }
745                Ok(())
746            }
747        }
748    }
749}
750
751/// Whether a non-child process has exited but not yet been reaped.
752///
753/// A zombie still answers `kill(pid, 0)`, because its process-table entry
754/// survives until its parent waits for it — so an existence probe alone calls
755/// it alive forever and anything waiting for it to leave waits forever. That is
756/// exactly the shape of `mj daemon restart` refusing to restart a daemon that
757/// had already stopped: `spawn_detached` used to leave the daemon a child of a
758/// long-lived Mjolnir process that never reaped it. It now double-forks, so the
759/// daemon is init's to reap, but any other unreaped child of this process would
760/// look the same, and the check stays cheap.
761///
762/// Treating a zombie as gone is also safe in the direction that matters: a
763/// zombie's PID cannot be reused until it is reaped, so nothing else can be
764/// occupying that number while this returns true.
765#[cfg(unix)]
766pub fn process_is_zombie(pid: u32) -> bool {
767    let pid = sysinfo::Pid::from_u32(pid);
768    let mut system = sysinfo::System::new();
769    system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
770    system
771        .process(pid)
772        .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
773}
774
775/// Wait for a process to leave, within [`STOP_TIMEOUT`].
776///
777/// The error says the process was still running rather than that it "did not
778/// stop": a daemon that is still winding down has not refused, and the two
779/// read very differently to somebody deciding whether to reach for a kill.
780pub async fn wait_for_exit(pid: u32) -> Result<()> {
781    let deadline = Instant::now() + STOP_TIMEOUT;
782    while daemon_process_is_alive(pid) {
783        ensure!(Instant::now() < deadline, "process {pid} is still running");
784        tokio::time::sleep(RETRY_DELAY).await;
785    }
786    Ok(())
787}
788
789/// Whether a daemon that Mjolnir launched in its own process group is alive.
790///
791/// Reaping is deliberately confined to this daemon-specific path. Attachment
792/// PIDs are merely observations and may alias unrelated children owned by this
793/// process, so their liveness probe below must never call `waitpid`.
794pub fn daemon_process_is_alive(pid: u32) -> bool {
795    #[cfg(unix)]
796    {
797        if pid == 0 {
798            return false;
799        }
800        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
801            return false;
802        };
803        let mut status = 0;
804        // SAFETY: `status` is writable for the call and WNOHANG never blocks.
805        // A non-child fails with ECHILD without changing any process state.
806        let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
807        if waited == raw_pid {
808            return false;
809        }
810        if waited == 0 {
811            return true;
812        }
813        let wait_error = std::io::Error::last_os_error();
814        if wait_error.raw_os_error() != Some(libc::ECHILD) {
815            return true;
816        }
817
818        #[cfg(target_os = "macos")]
819        return owned_daemon_group_is_alive(raw_pid);
820
821        #[cfg(not(target_os = "macos"))]
822        process_is_alive(pid)
823    }
824    #[cfg(not(unix))]
825    process_is_alive(pid)
826}
827
828#[cfg(target_os = "macos")]
829pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
830    // `spawn_detached` makes the daemon a process-group leader. Darwin
831    // excludes zombies from group signal probes: ESRCH means the group is gone
832    // and EPERM means only exiting members remain. The latter is safe here
833    // because this is a group we created for our own same-user child, not an
834    // arbitrary process group.
835    // SAFETY: signal 0 is only an existence probe, and the negative PID targets
836    // the daemon-owned group rather than another process.
837    if unsafe { libc::kill(-pid, 0) } == 0 {
838        return true;
839    }
840    let error = std::io::Error::last_os_error();
841    !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
842}
843
844pub fn process_is_alive(pid: u32) -> bool {
845    #[cfg(unix)]
846    {
847        if pid == 0 {
848            return false;
849        }
850        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
851            return false;
852        };
853        // SAFETY: kill(pid, 0) sends no signal and is the standard existence
854        // probe. EPERM still means the process exists.
855        let result = unsafe { libc::kill(raw_pid, 0) };
856        let exists =
857            result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
858        exists && !process_is_zombie(pid)
859    }
860    #[cfg(not(unix))]
861    {
862        let _ = pid;
863        true
864    }
865}
866
867pub fn read_metadata() -> Result<DaemonMetadata> {
868    let metadata = read_metadata_any()?;
869    ensure!(
870        metadata.protocol_version == PROTOCOL_VERSION,
871        "daemon protocol {} is incompatible with client protocol {}",
872        metadata.protocol_version,
873        PROTOCOL_VERSION
874    );
875    Ok(metadata)
876}
877
878pub fn read_metadata_any() -> Result<DaemonMetadata> {
879    let path = metadata_path();
880    let body = fs::read(&path).with_context(|| format!("read {}", path.display()))?;
881    let metadata: DaemonMetadata =
882        serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
883    Ok(metadata)
884}
885
886pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
887    let body = serde_json::to_vec(value)?;
888    write_encoded_frame(stream, &body).await
889}
890
891pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
892    ensure!(
893        body.len() <= MAX_FRAME_BYTES,
894        "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
895        body.len()
896    );
897    stream.write_u32(body.len() as u32).await?;
898    stream.write_all(body).await?;
899    stream.flush().await?;
900    Ok(())
901}
902
903pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
904    let length = stream.read_u32().await? as usize;
905    ensure!(
906        length <= MAX_FRAME_BYTES,
907        "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
908    );
909    let mut body = vec![0_u8; length];
910    stream.read_exact(&mut body).await?;
911    serde_json::from_slice(&body).context("decode daemon frame")
912}
913
914pub struct DaemonClient {
915    metadata: DaemonMetadata,
916    stream: TcpStream,
917    next_request_id: u64,
918}
919
920impl DaemonClient {
921    pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
922        let stream =
923            tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
924                .await
925                .context("time out connecting to Mjolnir daemon")??;
926        Ok(Self {
927            metadata,
928            stream,
929            next_request_id: 1,
930        })
931    }
932
933    /// Speak the daemon's advertised dialect, not this build's: management
934    /// requests must reach daemons of any protocol version, and the frozen
935    /// subset encodes identically across all of them.
936    pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
937        let protocol_version = self.metadata.protocol_version;
938        let request_id = self.next_request_id;
939        self.next_request_id += 1;
940        write_frame(
941            &mut self.stream,
942            &RequestEnvelope {
943                protocol_version,
944                request_id,
945                token: self.metadata.token.clone(),
946                action,
947            },
948        )
949        .await?;
950        let response: ResponseEnvelope = read_frame(&mut self.stream).await?;
951        ensure!(
952            response.protocol_version == protocol_version,
953            "daemon changed protocol"
954        );
955        ensure!(
956            response.request_id == request_id,
957            "daemon crossed request IDs"
958        );
959        response.result.map_err(anyhow::Error::msg)
960    }
961
962    pub async fn status(&mut self) -> Result<DaemonStatus> {
963        match self.request(DaemonAction::Status).await? {
964            DaemonReply::Status(status) => Ok(status),
965            reply => bail!("unexpected daemon status reply {reply:?}"),
966        }
967    }
968
969    pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
970        match self.request(DaemonAction::WebViewerAccess).await? {
971            DaemonReply::WebViewerAccess(access) => Ok(access),
972            reply => bail!("unexpected web viewer reply {reply:?}"),
973        }
974    }
975
976    pub async fn recover_web_viewer(
977        &mut self,
978        action: crate::web::WebViewerRecovery,
979    ) -> Result<()> {
980        match self.request(DaemonAction::RecoverWebViewer(action)).await? {
981            DaemonReply::Done => Ok(()),
982            reply => bail!("unexpected web viewer recovery reply {reply:?}"),
983        }
984    }
985
986    pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
987        match self.request(DaemonAction::InspectWebListener).await? {
988            DaemonReply::WebListeners(processes) => Ok(processes),
989            reply => bail!("unexpected listener inspection reply {reply:?}"),
990        }
991    }
992
993    pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
994        match self.request(DaemonAction::ListWorkspaces).await? {
995            DaemonReply::Workspaces(workspaces) => Ok(workspaces),
996            reply => bail!("unexpected daemon workspace reply {reply:?}"),
997        }
998    }
999
1000    pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
1001        match self
1002            .request(DaemonAction::RenameProfile { old_id, new_id })
1003            .await?
1004        {
1005            DaemonReply::Done => Ok(()),
1006            reply => bail!("unexpected rename-profile reply {reply:?}"),
1007        }
1008    }
1009
1010    pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1011        match self
1012            .request(DaemonAction::RenameTarget { old_id, new_id })
1013            .await?
1014        {
1015            DaemonReply::Done => Ok(()),
1016            reply => bail!("unexpected rename-target reply {reply:?}"),
1017        }
1018    }
1019
1020    pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1021        match self.request(DaemonAction::CreateWorkspace { name }).await? {
1022            DaemonReply::Workspace(workspace) => Ok(workspace),
1023            reply => bail!("unexpected create-workspace reply {reply:?}"),
1024        }
1025    }
1026
1027    pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1028        match self
1029            .request(DaemonAction::RenameWorkspace { workspace_id, name })
1030            .await?
1031        {
1032            DaemonReply::Done => Ok(()),
1033            reply => bail!("unexpected rename-workspace reply {reply:?}"),
1034        }
1035    }
1036
1037    pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1038        match self
1039            .request(DaemonAction::TouchWorkspace { workspace_id })
1040            .await?
1041        {
1042            DaemonReply::Done => Ok(()),
1043            reply => bail!("unexpected touch-workspace reply {reply:?}"),
1044        }
1045    }
1046
1047    pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1048        match self
1049            .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1050            .await?
1051        {
1052            DaemonReply::Done => Ok(()),
1053            reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1054        }
1055    }
1056
1057    pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1058        match self
1059            .request(DaemonAction::CloseWorkspace { workspace_id })
1060            .await?
1061        {
1062            DaemonReply::Done => Ok(()),
1063            reply => bail!("unexpected close workspace reply: {reply:?}"),
1064        }
1065    }
1066
1067    pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1068        match self
1069            .request(DaemonAction::DeleteWorkspace { workspace_id })
1070            .await?
1071        {
1072            DaemonReply::Done => Ok(()),
1073            reply => bail!("unexpected delete-workspace reply {reply:?}"),
1074        }
1075    }
1076
1077    pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1078        match self
1079            .request(DaemonAction::Attach { client_id, pid })
1080            .await?
1081        {
1082            DaemonReply::Done => Ok(()),
1083            reply => bail!("unexpected attach reply {reply:?}"),
1084        }
1085    }
1086
1087    pub async fn detach(&mut self, client_id: String) -> Result<()> {
1088        match self.request(DaemonAction::Detach { client_id }).await? {
1089            DaemonReply::Done => Ok(()),
1090            reply => bail!("unexpected detach reply {reply:?}"),
1091        }
1092    }
1093
1094    pub async fn persist_read_receipt(
1095        &mut self,
1096        client_id: String,
1097        workspace_id: String,
1098        session_id: String,
1099        through: u64,
1100    ) -> Result<u64> {
1101        match self
1102            .request(DaemonAction::PersistReadReceipt {
1103                client_id,
1104                workspace_id,
1105                session_id,
1106                through,
1107            })
1108            .await?
1109        {
1110            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1111            reply => bail!("unexpected read-receipt reply {reply:?}"),
1112        }
1113    }
1114
1115    pub async fn persist_detached_session_state(
1116        &mut self,
1117        client_id: String,
1118        workspace_id: String,
1119        session_id: String,
1120        through: u64,
1121        owner_pid: u32,
1122        draft: mj_core::storage::DetachedSessionDraft,
1123    ) -> Result<()> {
1124        match self
1125            .request(DaemonAction::PersistDetachedSessionState {
1126                client_id,
1127                workspace_id,
1128                session_id,
1129                through,
1130                owner_pid,
1131                draft,
1132            })
1133            .await?
1134        {
1135            DaemonReply::Done => Ok(()),
1136            reply => bail!("unexpected detached-session-state reply {reply:?}"),
1137        }
1138    }
1139
1140    pub async fn save_active_review(
1141        &mut self,
1142        session_id: String,
1143        review: mj_core::storage::StoredReview,
1144    ) -> Result<()> {
1145        match self
1146            .request(DaemonAction::SaveActiveReview { session_id, review })
1147            .await?
1148        {
1149            DaemonReply::Done => Ok(()),
1150            reply => bail!("unexpected save-review reply {reply:?}"),
1151        }
1152    }
1153
1154    pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1155        match self
1156            .request(DaemonAction::ClearActiveReview { session_id })
1157            .await?
1158        {
1159            DaemonReply::Done => Ok(()),
1160            reply => bail!("unexpected clear-review reply {reply:?}"),
1161        }
1162    }
1163
1164    pub async fn save_workspace_pane_sizes(
1165        &mut self,
1166        workspace_id: String,
1167        sizes: mj_core::workspace::PaneSizes,
1168    ) -> Result<()> {
1169        match self
1170            .request(DaemonAction::SaveWorkspacePaneSizes {
1171                workspace_id,
1172                sizes,
1173            })
1174            .await?
1175        {
1176            DaemonReply::Done => Ok(()),
1177            reply => bail!("unexpected pane-size save reply {reply:?}"),
1178        }
1179    }
1180
1181    pub async fn save_workspace_layout(
1182        &mut self,
1183        workspace_id: String,
1184        layout: mj_core::workspace::ConversationLayout,
1185    ) -> Result<()> {
1186        match self
1187            .request(DaemonAction::SaveWorkspaceLayout {
1188                workspace_id,
1189                layout,
1190            })
1191            .await?
1192        {
1193            DaemonReply::Done => Ok(()),
1194            reply => bail!("unexpected layout save reply {reply:?}"),
1195        }
1196    }
1197
1198    pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1199        match self
1200            .request(DaemonAction::PersistImportedSession {
1201                session: Box::new(session),
1202            })
1203            .await?
1204        {
1205            DaemonReply::Done => Ok(()),
1206            reply => bail!("unexpected imported-session reply {reply:?}"),
1207        }
1208    }
1209
1210    pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1211        match self
1212            .request(DaemonAction::SetSessionTitle { session_id, title })
1213            .await?
1214        {
1215            DaemonReply::Text(title) => Ok(title),
1216            reply => bail!("unexpected session-title reply {reply:?}"),
1217        }
1218    }
1219
1220    pub async fn set_session_container_settings(
1221        &mut self,
1222        session_id: String,
1223        cpus: Option<String>,
1224        memory: Option<String>,
1225        mounts: Vec<AdditionalMount>,
1226        mount_history: Vec<PathBuf>,
1227    ) -> Result<()> {
1228        match self
1229            .request(DaemonAction::SetSessionContainerSettings {
1230                session_id,
1231                cpus,
1232                memory,
1233                mounts,
1234                mount_history,
1235            })
1236            .await?
1237        {
1238            DaemonReply::Done => Ok(()),
1239            reply => bail!("unexpected container-settings reply {reply:?}"),
1240        }
1241    }
1242
1243    pub async fn set_session_acp_title(
1244        &mut self,
1245        session_id: String,
1246        title: Option<String>,
1247    ) -> Result<()> {
1248        match self
1249            .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1250            .await?
1251        {
1252            DaemonReply::Done => Ok(()),
1253            reply => bail!("unexpected ACP-title reply {reply:?}"),
1254        }
1255    }
1256
1257    pub async fn mark_session_target_missing(
1258        &mut self,
1259        session_id: String,
1260        detail: String,
1261        updated_at: String,
1262    ) -> Result<Option<SessionState>> {
1263        match self
1264            .request(DaemonAction::MarkSessionTargetMissing {
1265                session_id,
1266                detail,
1267                updated_at,
1268            })
1269            .await?
1270        {
1271            DaemonReply::OptionalSessionState(state) => Ok(state),
1272            reply => bail!("unexpected target-missing reply {reply:?}"),
1273        }
1274    }
1275
1276    pub async fn checkpoint_session(
1277        &mut self,
1278        session_id: String,
1279    ) -> Result<mj_core::state::CheckpointMetadata> {
1280        match self
1281            .request(DaemonAction::CheckpointSession { session_id })
1282            .await?
1283        {
1284            DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1285            reply => bail!("unexpected checkpoint reply {reply:?}"),
1286        }
1287    }
1288
1289    /// Search the user's SessionWiki index, newest first when the query is
1290    /// empty and best match first otherwise. The reply carries the state of
1291    /// the index as well as the rows, so a caller can say the first build is
1292    /// still running.
1293    pub async fn jev_decisions(
1294        &mut self,
1295        session_id: String,
1296        decision_id: Option<String>,
1297    ) -> Result<mj_core::jev::DecisionPage> {
1298        match self
1299            .request(DaemonAction::JevDecisions {
1300                session_id,
1301                decision_id,
1302            })
1303            .await?
1304        {
1305            DaemonReply::JevDecisions(page) => Ok(page),
1306            reply => bail!("unexpected Jev diagnostics reply {reply:?}"),
1307        }
1308    }
1309
1310    pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1311        match self
1312            .request(DaemonAction::WikiSearch { query, limit })
1313            .await?
1314        {
1315            DaemonReply::WikiRows(page) => Ok(page),
1316            reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1317        }
1318    }
1319
1320    /// The markdown briefing for one indexed session.
1321    pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1322        match self
1323            .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1324            .await?
1325        {
1326            DaemonReply::Text(markdown) => Ok(markdown),
1327            reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1328        }
1329    }
1330
1331    /// The passages of one indexed session that match a query, each matching
1332    /// message with `context_messages` neighbours on either side and its text
1333    /// capped at `per_message_chars`. `None` when the index holds no session
1334    /// with that id.
1335    pub async fn wiki_hits(
1336        &mut self,
1337        wiki_id: String,
1338        query: String,
1339        context_messages: usize,
1340        per_message_chars: usize,
1341    ) -> Result<Option<WikiHitTranscript>> {
1342        match self
1343            .request(DaemonAction::WikiHits {
1344                wiki_id,
1345                query,
1346                context_messages,
1347                per_message_chars,
1348            })
1349            .await?
1350        {
1351            DaemonReply::WikiHits(transcript) => Ok(transcript),
1352            reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1353        }
1354    }
1355
1356    /// What the index knows about one session, or `None` when the index holds
1357    /// no session with that id.
1358    pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1359        match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1360            DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1361            reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1362        }
1363    }
1364
1365    /// Start a new session carrying a hand-off compacted from an archived one.
1366    /// It answers like any other session start: the record exists and is
1367    /// provisioning, and the hand-off follows once the harness is ready.
1368    pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1369        match self.request(DaemonAction::WikiRestore(request)).await? {
1370            DaemonReply::RegisteredSession(registered) => Ok(*registered),
1371            reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1372        }
1373    }
1374
1375    pub async fn scan_recovery(
1376        &mut self,
1377        all_instances: bool,
1378    ) -> Result<mj_core::state::RecoveryScan> {
1379        match self
1380            .request(DaemonAction::ScanRecovery { all_instances })
1381            .await?
1382        {
1383            DaemonReply::RecoveryScan(scan) => Ok(scan),
1384            reply => bail!("unexpected recovery-scan reply {reply:?}"),
1385        }
1386    }
1387
1388    pub async fn adopt_recovery(
1389        &mut self,
1390        session_id: String,
1391        target_id: String,
1392        profile: Option<String>,
1393        bundle: Option<String>,
1394        all_instances: bool,
1395    ) -> Result<()> {
1396        match self
1397            .request(DaemonAction::AdoptRecovery {
1398                session_id,
1399                target_id,
1400                profile,
1401                bundle,
1402                all_instances,
1403            })
1404            .await?
1405        {
1406            DaemonReply::Done => Ok(()),
1407            reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1408        }
1409    }
1410
1411    pub async fn destroy_recovery(
1412        &mut self,
1413        session_id: String,
1414        target_id: String,
1415        confirmation: String,
1416        all_instances: bool,
1417    ) -> Result<()> {
1418        match self
1419            .request(DaemonAction::DestroyRecovery {
1420                session_id,
1421                target_id,
1422                confirmation,
1423                all_instances,
1424            })
1425            .await?
1426        {
1427            DaemonReply::Done => Ok(()),
1428            reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1429        }
1430    }
1431
1432    pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1433        match self
1434            .request(DaemonAction::Snapshot { workspace_id })
1435            .await?
1436        {
1437            DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1438            reply => bail!("unexpected snapshot reply {reply:?}"),
1439        }
1440    }
1441
1442    pub async fn runtime_snapshot(
1443        &mut self,
1444        workspace_id: String,
1445        after_revision: u64,
1446        all_workspaces: bool,
1447    ) -> Result<RuntimeSnapshot> {
1448        match self
1449            .request(DaemonAction::RuntimeSnapshot {
1450                workspace_id,
1451                after_revision,
1452                all_workspaces,
1453            })
1454            .await?
1455        {
1456            DaemonReply::RuntimeSnapshot(snapshot) => Ok(*snapshot),
1457            reply => bail!("unexpected runtime snapshot reply {reply:?}"),
1458        }
1459    }
1460
1461    pub async fn submit_session_command(
1462        &mut self,
1463        session_id: String,
1464        command_id: String,
1465        command: RelayCommand,
1466        inherited_draft: Option<String>,
1467    ) -> Result<u64> {
1468        match self
1469            .request(DaemonAction::SubmitSessionCommand {
1470                inherited_draft,
1471                session_id,
1472                command_id,
1473                command,
1474            })
1475            .await?
1476        {
1477            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1478            reply => bail!("unexpected session command reply {reply:?}"),
1479        }
1480    }
1481
1482    /// Hand the daemon a prompt for a session that is still starting. The
1483    /// daemon replies as soon as the prompt is queued, not when it is
1484    /// delivered; delivery failures come back as a session notice.
1485    pub async fn queue_startup_prompt(
1486        &mut self,
1487        session_id: String,
1488        text: String,
1489        inherited_draft: Option<String>,
1490    ) -> Result<()> {
1491        match self
1492            .request(DaemonAction::QueueStartupPrompt {
1493                session_id,
1494                text,
1495                inherited_draft,
1496            })
1497            .await?
1498        {
1499            DaemonReply::Done => Ok(()),
1500            reply => bail!("unexpected startup prompt reply {reply:?}"),
1501        }
1502    }
1503
1504    /// Ask the daemon to review the turn this session just finished.
1505    ///
1506    /// The refusal is a sentence for a person -- "prompts are queued", "set
1507    /// [review] profile in config.toml" -- so it travels as text rather than
1508    /// as a code every surface would have to translate.
1509    pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
1510        match self
1511            .request(DaemonAction::StartTurnReview { session_id })
1512            .await?
1513        {
1514            DaemonReply::Done => Ok(()),
1515            reply => bail!("unexpected turn-review reply {reply:?}"),
1516        }
1517    }
1518
1519    pub async fn resolve_turn_review(
1520        &mut self,
1521        session_id: String,
1522        resolution: Resolution,
1523    ) -> Result<()> {
1524        match self
1525            .request(DaemonAction::ResolveTurnReview {
1526                session_id,
1527                resolution,
1528            })
1529            .await?
1530        {
1531            DaemonReply::Done => Ok(()),
1532            reply => bail!("unexpected turn-review resolution reply {reply:?}"),
1533        }
1534    }
1535
1536    pub async fn reviewer_action(
1537        &mut self,
1538        session_id: String,
1539        role: Option<String>,
1540        action: crate::session::ReviewerAction,
1541    ) -> Result<crate::session::ReviewerOutcome> {
1542        match self
1543            .request(DaemonAction::ReviewerAction {
1544                session_id,
1545                role,
1546                action,
1547            })
1548            .await?
1549        {
1550            DaemonReply::Reviewer(outcome) => Ok(*outcome),
1551            reply => bail!("unexpected reviewer reply {reply:?}"),
1552        }
1553    }
1554
1555    pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
1556        match self
1557            .request(DaemonAction::SyncSession { session_id })
1558            .await?
1559        {
1560            DaemonReply::Done => Ok(()),
1561            reply => bail!("unexpected session sync reply {reply:?}"),
1562        }
1563    }
1564
1565    pub async fn respond_elicitation(
1566        &mut self,
1567        session_id: String,
1568        elicitation_id: String,
1569        response: ElicitationResponse,
1570    ) -> Result<()> {
1571        match self
1572            .request(DaemonAction::RespondElicitation {
1573                session_id,
1574                elicitation_id,
1575                response,
1576            })
1577            .await?
1578        {
1579            DaemonReply::Done => Ok(()),
1580            reply => bail!("unexpected elicitation reply {reply:?}"),
1581        }
1582    }
1583
1584    pub async fn native_agent_history(
1585        &mut self,
1586        owner: String,
1587        child: String,
1588        before: Option<(u64, String)>,
1589    ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
1590        match self
1591            .request(DaemonAction::NativeAgentHistory {
1592                owner,
1593                child,
1594                before,
1595            })
1596            .await?
1597        {
1598            DaemonReply::NativeAgentHistory(page) => Ok(page),
1599            reply => bail!("unexpected native agent history reply {reply:?}"),
1600        }
1601    }
1602
1603    pub async fn stop_background_task(
1604        &mut self,
1605        session_id: String,
1606        background_task_id: String,
1607    ) -> Result<()> {
1608        match self
1609            .request(DaemonAction::StopBackgroundTask {
1610                session_id,
1611                background_task_id,
1612            })
1613            .await?
1614        {
1615            DaemonReply::Done => Ok(()),
1616            reply => bail!("unexpected background task stop reply {reply:?}"),
1617        }
1618    }
1619
1620    pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
1621        match self
1622            .request(DaemonAction::SuspendSession { session_id })
1623            .await?
1624        {
1625            DaemonReply::Done => Ok(()),
1626            reply => bail!("unexpected close-session reply {reply:?}"),
1627        }
1628    }
1629
1630    pub async fn start_create_session(
1631        &mut self,
1632        request: CreateSessionRequest,
1633    ) -> Result<RegisteredSession> {
1634        match self
1635            .request(DaemonAction::StartCreateSession(request))
1636            .await?
1637        {
1638            DaemonReply::RegisteredSession(registered) => Ok(*registered),
1639            reply => bail!("unexpected start-create reply {reply:?}"),
1640        }
1641    }
1642
1643    pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
1644        match self
1645            .request(DaemonAction::WaitCreateSession { session_id })
1646            .await?
1647        {
1648            DaemonReply::Done => Ok(()),
1649            reply => bail!("unexpected wait-create reply {reply:?}"),
1650        }
1651    }
1652
1653    pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
1654        match self.request(DaemonAction::ResumeSession(request)).await? {
1655            DaemonReply::Done => Ok(()),
1656            reply => bail!("unexpected resume-session reply {reply:?}"),
1657        }
1658    }
1659
1660    pub async fn discard_since_checkpoint(
1661        &mut self,
1662        session_id: String,
1663        checkpoint: mj_core::state::CheckpointMetadata,
1664    ) -> Result<()> {
1665        match self
1666            .request(DaemonAction::DiscardSinceCheckpoint {
1667                session_id,
1668                checkpoint,
1669            })
1670            .await?
1671        {
1672            DaemonReply::Done => Ok(()),
1673            reply => bail!("unexpected force-stop reply {reply:?}"),
1674        }
1675    }
1676
1677    pub async fn destroy_stopped_session(
1678        &mut self,
1679        session_id: String,
1680        delete_branch: bool,
1681    ) -> Result<()> {
1682        match self
1683            .request(DaemonAction::DestroyStoppedSession {
1684                session_id,
1685                delete_branch,
1686            })
1687            .await?
1688        {
1689            DaemonReply::Done => Ok(()),
1690            reply => bail!("unexpected destroy-stopped reply {reply:?}"),
1691        }
1692    }
1693
1694    pub async fn force_destroy_session(
1695        &mut self,
1696        session_id: String,
1697        delete_branch: bool,
1698    ) -> Result<()> {
1699        match self
1700            .request(DaemonAction::ForceDestroySession {
1701                session_id,
1702                delete_branch,
1703            })
1704            .await?
1705        {
1706            DaemonReply::Done => Ok(()),
1707            reply => bail!("unexpected force-destroy reply {reply:?}"),
1708        }
1709    }
1710
1711    pub async fn force_delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1712        match self
1713            .request(DaemonAction::ForceDeleteWorkspace { workspace_id })
1714            .await?
1715        {
1716            DaemonReply::Done => Ok(()),
1717            reply => bail!("unexpected force-delete-workspace reply {reply:?}"),
1718        }
1719    }
1720
1721    pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
1722        match self
1723            .request(DaemonAction::CancelLifecycle { session_id })
1724            .await?
1725        {
1726            DaemonReply::Done => Ok(()),
1727            reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
1728        }
1729    }
1730
1731    pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
1732        match self
1733            .request(DaemonAction::RecoverDraft { draft_id })
1734            .await?
1735        {
1736            DaemonReply::Done => Ok(()),
1737            reply => bail!("unexpected recover-draft reply {reply:?}"),
1738        }
1739    }
1740
1741    pub async fn stop(&mut self) -> Result<()> {
1742        match self.request(DaemonAction::Stop).await? {
1743            DaemonReply::Done => Ok(()),
1744            reply => bail!("unexpected stop reply {reply:?}"),
1745        }
1746    }
1747}
1748
1749pub async fn connect_existing() -> Result<DaemonClient> {
1750    let metadata = tokio::task::spawn_blocking(read_metadata)
1751        .await
1752        .context("read daemon metadata task failed")??;
1753    DaemonClient::connect(metadata).await
1754}
1755
1756/// A handle to whatever daemon the metadata file advertises, regardless of its
1757/// protocol version. It only exposes the frozen management subset (`Ping`,
1758/// `Status`, `Stop`), which encodes identically in every protocol version.
1759pub struct ManagementClient {
1760    inner: DaemonClient,
1761}
1762
1763impl ManagementClient {
1764    pub fn new(inner: DaemonClient) -> Self {
1765        Self { inner }
1766    }
1767    pub fn protocol_version(&self) -> u32 {
1768        self.inner.metadata.protocol_version
1769    }
1770
1771    pub async fn status(&mut self) -> Result<DaemonStatus> {
1772        self.inner.status().await
1773    }
1774
1775    pub async fn stop(&mut self) -> Result<()> {
1776        tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
1777            .await
1778            .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
1779    }
1780
1781    /// Ask the daemon to stop and wait for its process to actually exit.
1782    pub async fn stop_and_wait(mut self) -> Result<()> {
1783        let pid = self.inner.metadata.pid;
1784        self.stop().await?;
1785        wait_for_exit(pid).await.with_context(|| {
1786            format!(
1787                "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
1788                STOP_TIMEOUT.as_secs()
1789            )
1790        })
1791    }
1792}
1793
1794pub async fn connect_management() -> Result<ManagementClient> {
1795    Ok(ManagementClient {
1796        inner: DaemonClient::connect(read_metadata_any()?).await?,
1797    })
1798}
1799
1800/// Refuse to speak to a daemon whose protocol this build does not know.
1801///
1802/// The two protocol numbers alone do not say which `mj` ran. The usual cause is
1803/// a second installation: a `cargo install`ed client sits earlier on PATH than
1804/// the build whose daemon is running, so every command fails here while the
1805/// other binary works, and nothing in the message says where either lives. It
1806/// therefore names both executables and both versions.
1807pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
1808    ensure!(
1809        metadata.protocol_version <= PROTOCOL_VERSION,
1810        "{}",
1811        unsupported_daemon_protocol_message(
1812            metadata.protocol_version,
1813            &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
1814        )
1815    );
1816    Ok(())
1817}
1818
1819/// The message [`ensure_supported_daemon_protocol`] fails with, given the
1820/// sentence that names both builds, so it can be read without a daemon.
1821fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
1822    format!(
1823        "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
1824         Put the daemon's directory first on PATH, or reinstall this client from that build."
1825    )
1826}
1827pub const PROTOCOL_VERSION: u32 = 32;
1828pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
1829/// How long a daemon is given to exit after it accepts a stop.
1830///
1831/// Stopping cancels a token and returns immediately; the daemon then unwinds
1832/// its session manager, its phone server and its pollers. That is normally
1833/// fast, but a daemon whose database has been migrated out from under it fails
1834/// every read while it winds down and has been observed taking over five
1835/// seconds — which the previous five-second bound missed by a fraction,
1836/// reporting a stop that had in fact worked as `did not stop` and aborting the
1837/// restart that depended on it.
1838pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
1839pub const RETRY_DELAY: Duration = Duration::from_millis(40);
1840impl DaemonClient {
1841    pub async fn prepare_move_session(
1842        &mut self,
1843        selection: MoveSelection,
1844    ) -> Result<MovePreparation> {
1845        match self
1846            .request(DaemonAction::PrepareMoveSession(selection))
1847            .await?
1848        {
1849            DaemonReply::MovePreparation(preparation) => Ok(*preparation),
1850            _ => bail!("daemon returned an unexpected move preparation reply"),
1851        }
1852    }
1853
1854    pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
1855        match self.request(DaemonAction::MoveSession(request)).await? {
1856            DaemonReply::MoveOutcome(outcome) => Ok(outcome),
1857            _ => bail!("daemon returned an unexpected move reply"),
1858        }
1859    }
1860}
1861
1862#[cfg(test)]
1863mod tests {
1864    use super::*;
1865    use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
1866    use std::path::Path;
1867
1868    #[test]
1869    fn unsupported_protocol_message_names_both_binaries_and_versions() {
1870        let message = unsupported_daemon_protocol_message(
1871            PROTOCOL_VERSION + 5,
1872            &describe_daemon_and_client_builds(
1873                4242,
1874                BuildDescription {
1875                    executable: Some(Path::new("/home/dev/mj/target/release/mj")),
1876                    version: "2.14.0",
1877                },
1878                BuildDescription {
1879                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
1880                    version: "2.9.0",
1881                },
1882            ),
1883        );
1884        assert_eq!(
1885            message,
1886            format!(
1887                "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
1888                 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
1889                 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
1890                 Put the daemon's directory first on PATH, or reinstall this client from that build.",
1891                PROTOCOL_VERSION + 5
1892            )
1893        );
1894    }
1895
1896    #[test]
1897    fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
1898        let message = unsupported_daemon_protocol_message(
1899            PROTOCOL_VERSION + 1,
1900            &describe_daemon_and_client_builds(
1901                4242,
1902                BuildDescription {
1903                    executable: None,
1904                    version: "2.14.0",
1905                },
1906                BuildDescription {
1907                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
1908                    version: "2.9.0",
1909                },
1910            ),
1911        );
1912        assert!(
1913            message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
1914            "{message}"
1915        );
1916        assert!(
1917            message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
1918            "{message}"
1919        );
1920    }
1921}