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