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