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::session::{ManagedSessionView, ViewError};
4use anyhow::{Context, Result, bail, ensure};
5use mj_core::config::{Config, data_dir};
6use mj_core::credentials::CredentialSyncSignal;
7use mj_core::elicitation::ElicitationResponse;
8use mj_core::relay::{RelayCommand, RelayOperationalState};
9use mj_core::review::driver::Resolution;
10use mj_core::state::*;
11use mj_core::targets::{AdditionalMount, ProvisionStage};
12use mj_core::workspace::WorkspaceRecord;
13use serde::{Deserialize, Serialize};
14use std::fs;
15use std::net::SocketAddr;
16use std::path::PathBuf;
17use std::time::{Duration, Instant};
18use tokio::io::{AsyncReadExt, AsyncWriteExt};
19use tokio::net::TcpStream;
20pub fn metadata_path() -> PathBuf {
21    data_dir().join("daemon.json")
22}
23
24#[derive(Debug, Clone, Serialize, Deserialize)]
25#[serde(deny_unknown_fields)]
26pub struct DaemonMetadata {
27    pub protocol_version: u32,
28    pub pid: u32,
29    pub address: SocketAddr,
30    pub token: String,
31    pub started_at: String,
32    pub build_version: String,
33}
34
35#[derive(Debug, Clone, Serialize, Deserialize)]
36#[serde(deny_unknown_fields)]
37pub struct WorkspaceListing {
38    pub workspace: WorkspaceRecord,
39}
40
41#[derive(Debug, Clone, Serialize, Deserialize)]
42#[serde(deny_unknown_fields)]
43pub struct SessionPreview {
44    pub id: String,
45    pub title: String,
46    pub project: String,
47    pub harness: String,
48    pub state: String,
49    pub active: bool,
50    pub updated_at: String,
51}
52
53/// One session in the user's SessionWiki index, as a control surface shows it.
54///
55/// It is the daemon's own shape rather than SessionWiki's row: it carries the
56/// search snippet that found the row and, for a Mjolnir session this daemon
57/// still holds, the id that resumes it instead of restoring it.
58#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
59#[serde(deny_unknown_fields)]
60pub struct WikiRow {
61    pub id: String,
62    pub tool: String,
63    pub project: String,
64    pub title: String,
65    pub started: Option<String>,
66    pub msgs: i64,
67    pub preview: Option<String>,
68    /// The tool deleted its own copy and SessionWiki kept the transcript.
69    pub archived: bool,
70    /// The session id the tool that ran it knows it by, when its stored path
71    /// carries one. It is what matches a row against an import scan.
72    pub native_id: Option<String>,
73    /// The matching text, when this row came from a search.
74    pub snippet: Option<String>,
75    /// The live Mjolnir session this row describes, when this daemon has it.
76    pub hel_session_id: Option<String>,
77    /// The Mjolnir target template the session ran under, when the index
78    /// carries it. Only Mjolnir's own rows have one: the daemon writes it into
79    /// the index as an `mj-target:` tag while the session is still known, so an
80    /// archived row can still say where it ran.
81    #[serde(default)]
82    pub target: Option<String>,
83    /// The Mjolnir harness profile the session last ran under, from the index's
84    /// `mj-profile:` tag. Only Mjolnir's own rows have one.
85    #[serde(default)]
86    pub profile: Option<String>,
87    /// The harness kind the session ran (`codex`, `claude`, `kimi`, `grok`,
88    /// `muse`), from the index's `mj-harness:` tag. It stays meaningful after
89    /// the profile id has been removed from the configuration.
90    #[serde(default)]
91    pub harness: Option<String>,
92}
93
94/// How far along the daemon's SessionWiki index is when a search answers.
95#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
96#[serde(rename_all = "snake_case")]
97pub enum WikiIndexState {
98    /// The index has completed one full build; its answers are complete.
99    Ready,
100    /// The first full build has not finished yet, so a search can miss
101    /// sessions that exist. This is the state a fresh index starts in.
102    #[default]
103    Indexing,
104    /// The index file on disk was written by a different SessionWiki schema
105    /// version. Mjolnir will not open it, because opening it would drop and
106    /// rebuild the user's whole cache.
107    VersionMismatch,
108}
109
110/// What a search says about the index it answered from.
111#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
112#[serde(deny_unknown_fields)]
113pub struct WikiStatus {
114    pub state: WikiIndexState,
115    /// A sync is running now, so repeating the query may return more.
116    pub topping_up: bool,
117}
118
119/// Where in a session's conversation a text search matched. A match in a
120/// user message ranks ahead of one in an agent message; a match only in tool
121/// output is not a match at all.
122#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
123#[serde(rename_all = "snake_case")]
124pub enum SessionTextMatchKind {
125    User,
126    Agent,
127}
128
129/// One live session whose conversation contains a search query.
130#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
131#[serde(deny_unknown_fields)]
132pub struct SessionTextMatch {
133    pub session_id: String,
134    pub kind: SessionTextMatchKind,
135}
136
137/// One page of search results with the state of the index behind them.
138#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
139#[serde(deny_unknown_fields)]
140pub struct WikiSearchPage {
141    pub rows: Vec<WikiRow>,
142    pub status: WikiStatus,
143}
144
145/// One message of an indexed transcript, reduced to what a search preview
146/// shows: the text around the query's matches, with the matches located in it.
147#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
148#[serde(deny_unknown_fields)]
149pub struct WikiHitBlock {
150    /// `user`, `assistant` or `tool`.
151    pub role: String,
152    /// The message text, redacted, and windowed to the caller's per-message
153    /// budget when the message is longer than that.
154    pub text: String,
155    /// Byte ranges of the matches inside `text`, on character boundaries, in
156    /// order. A context message has none.
157    pub hits: Vec<(usize, usize)>,
158    /// Messages between the previous block and this one that no group covered.
159    /// Non-zero only on the first block of a group.
160    pub omitted_before: usize,
161    /// `text` is a window of the message rather than the whole of it.
162    pub truncated: bool,
163}
164
165/// The matching passages of one indexed transcript.
166#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
167#[serde(deny_unknown_fields)]
168pub struct WikiHitTranscript {
169    pub blocks: Vec<WikiHitBlock>,
170    /// Messages after the last block that no group covered.
171    pub omitted_after: usize,
172}
173
174/// What one indexed session is, as far as continuing it is concerned.
175#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
176#[serde(rename_all = "snake_case")]
177pub enum WikiSessionStatus {
178    /// A Mjolnir session this daemon still has a record of.
179    Mine,
180    /// A Mjolnir session whose record the archive job destroyed; only the
181    /// indexed transcript is left.
182    Archived,
183    /// Another tool's session, which Mjolnir would have to import.
184    Native,
185}
186
187/// One row of the SessionWiki index, with what Mjolnir knows about it.
188///
189/// This is what `mj resume --wiki` branches on and what `mj sessions
190/// --session` reports when the id names no Mjolnir session.
191#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192#[serde(deny_unknown_fields)]
193pub struct WikiSessionInfo {
194    /// The SessionWiki id of the row.
195    pub wiki_id: String,
196    /// The tool that produced the session, as SessionWiki names it.
197    pub tool: String,
198    /// The transcript path the index stores for the row.
199    pub path: PathBuf,
200    pub status: WikiSessionStatus,
201    /// The Mjolnir session id, for a row Mjolnir itself published.
202    pub mjolnir_session_id: Option<String>,
203    /// The profile the session last ran under, from the index's own tags.
204    pub profile_id: Option<String>,
205    /// The target template the session ran on, from the index's own tags.
206    pub target_template_id: Option<String>,
207    /// The harness that drove the session, from the tags for a Mjolnir row and
208    /// from the tool name for a native one.
209    pub harness: Option<mj_core::config::HarnessKind>,
210    pub title: String,
211    pub project: String,
212    /// Set for an archived Mjolnir row whose transcript holds no prompt,
213    /// which `mj resume --wiki` has nothing to restore from. Sent only when
214    /// set, so a daemon that predates it and a client that predates it both
215    /// read every other row unchanged.
216    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
217    pub nothing_to_restore: bool,
218}
219
220/// Start a new session carrying a compacted hand-off from an archived one.
221#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
222#[serde(deny_unknown_fields)]
223pub struct WikiRestoreRequest {
224    /// The SessionWiki session id to restore from.
225    pub wiki_id: String,
226    pub workspace_id: String,
227    pub profile_id: String,
228    pub target_template_id: String,
229    /// Where the new session opens. None takes the project the archived
230    /// session ran in, when that directory still exists.
231    #[serde(default)]
232    pub project_directory: Option<PathBuf>,
233    #[serde(default)]
234    pub additional_mounts: Vec<AdditionalMount>,
235    #[serde(default)]
236    pub resource_allocation: Option<SessionResourceAllocation>,
237}
238
239#[derive(Debug, Clone, Serialize, Deserialize)]
240#[serde(deny_unknown_fields)]
241pub struct WorkspaceSnapshot {
242    pub workspace: WorkspaceRecord,
243    pub sessions: Vec<SessionPreview>,
244    pub drafts: Vec<DraftPreview>,
245}
246
247#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
248#[serde(deny_unknown_fields)]
249pub struct RuntimeSessionView {
250    pub session_id: String,
251    pub projection_ordinal: u64,
252    pub projection_digest: String,
253    pub operational: Option<RelayOperationalState>,
254    pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
255    pub connected: bool,
256    pub error: Option<ViewError>,
257}
258
259impl RuntimeSessionView {
260    pub fn from_managed(session_id: String, view: ManagedSessionView) -> Self {
261        let (projection_ordinal, projection_digest, operational, signal) =
262            view.snapshot
263                .map_or((0, String::new(), None, None), |snapshot| {
264                    (
265                        snapshot.materialized.applied_event_ordinal,
266                        snapshot.materialized.applied_event_digest,
267                        Some(snapshot.operational),
268                        snapshot.latest_credential_sync_signal,
269                    )
270                });
271        Self {
272            session_id,
273            projection_ordinal,
274            projection_digest,
275            operational,
276            latest_credential_sync_signal: signal,
277            connected: view.connected,
278            error: view.error,
279        }
280    }
281}
282
283/// Something the daemon did on its own that a surface should report once.
284///
285/// Background work has no lifecycle entry to hang a message on, so notices
286/// travel with the snapshot and carry an id: a surface reports the ones newer
287/// than the last it saw and nothing else, however often it polls.
288#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
289#[serde(deny_unknown_fields)]
290pub struct RuntimeNotice {
291    pub id: u64,
292    pub session_id: String,
293    pub text: String,
294}
295
296#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
297#[serde(rename_all = "snake_case")]
298pub enum RuntimeLifecycleKind {
299    Create,
300    Suspend,
301    Resume,
302    Move,
303    ForceStop,
304    DestroyStopped,
305    ForceDestroy,
306    /// A sub-agent stopped and removed because its parent is being
307    /// suspended. The same teardown as `ForceDestroy`, which surfaces call a
308    /// stop, as the suspend does.
309    StopSubagent,
310    Cleanup,
311}
312
313#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
314#[serde(deny_unknown_fields)]
315pub struct RuntimeLifecycleView {
316    pub operation_id: String,
317    pub cancellable: bool,
318    pub session_id: String,
319    pub kind: RuntimeLifecycleKind,
320    pub started_at_epoch_seconds: u64,
321    pub active_stages: Vec<(ProvisionStage, u64)>,
322    pub resume_destination: Option<(String, String)>,
323    pub notice: Option<String>,
324}
325
326#[derive(Debug, Clone, Serialize, Deserialize)]
327#[serde(deny_unknown_fields)]
328pub struct ResumeSessionRequest {
329    pub session_id: String,
330    pub workspace_id: String,
331    pub profile_id: String,
332    pub target_template_id: String,
333    pub additional_mounts: Option<Vec<AdditionalMount>>,
334    pub resource_allocation: Option<SessionResourceAllocation>,
335    pub discard_queue: bool,
336    pub repository_preflight: Option<ResumeRepositorySourceReceipt>,
337}
338
339#[derive(Debug, Clone, Serialize, Deserialize)]
340#[serde(deny_unknown_fields)]
341pub struct CreateSessionRequest {
342    #[serde(default)]
343    pub create_managed_worktree: Option<bool>,
344    /// Full commit object ID to start the workspace at: an exact checkout of
345    /// the bundle's primary repository, detached unless `branch` is given.
346    #[serde(default, skip_serializing_if = "Option::is_none")]
347    pub at: Option<String>,
348    /// With `at`, the new branch created there; otherwise the existing branch
349    /// to check out in a new isolated workspace.
350    #[serde(default, skip_serializing_if = "Option::is_none")]
351    pub branch: Option<String>,
352    /// Diff base, when it is not `at`. Without `at`, also the starting
353    /// revision for a raw managed worktree.
354    #[serde(default, skip_serializing_if = "Option::is_none")]
355    pub base: Option<String>,
356    /// Omitted uses the selected profile's subagent setting.
357    #[serde(
358        default,
359        alias = "mjolnir_subagents",
360        deserialize_with = "mj_core::subagent::deserialize_optional_policy"
361    )]
362    pub subagents: Option<mj_core::subagent::SubagentPolicy>,
363    #[serde(default)]
364    pub initial_prompt: Option<String>,
365    pub workspace_id: String,
366    pub profile_id: String,
367    pub bundle_id: String,
368    pub project_directory: Option<PathBuf>,
369    pub target_template_id: String,
370    pub additional_mounts: Vec<AdditionalMount>,
371    pub resource_allocation: Option<SessionResourceAllocation>,
372    pub title: String,
373    pub session_title_override: Option<String>,
374}
375
376#[derive(Debug, Clone, Serialize, Deserialize)]
377#[serde(deny_unknown_fields)]
378pub struct RegisteredSession {
379    pub session: SessionRecord,
380    pub remembered_container_size: Option<(String, HostContainerSize)>,
381}
382
383#[derive(Debug, Clone, Serialize, Deserialize)]
384#[serde(deny_unknown_fields)]
385pub struct DraftPreview {
386    pub id: String,
387    pub session_id: Option<String>,
388    pub source: String,
389    pub owner_pid: Option<u32>,
390    pub saved_at: String,
391}
392
393/// `Ping`, `Status`, and `Stop` form the frozen management subset: their wire
394/// encoding — together with `RequestEnvelope`, `ResponseEnvelope`,
395/// `DaemonStatus`, and `WebViewerStatus` — must never change shape, because
396/// clients and daemons of *any* protocol version rely on them to identify,
397/// stop, and replace each other. `PrepareUpgrade` and its `Done` /
398/// `UpgradePending` replies are also frozen from protocol 33 onward: they
399/// drain accepted work before closing admission atomically. Every other
400/// action may change freely behind a `PROTOCOL_VERSION` bump.
401#[derive(Debug, Clone, Serialize, Deserialize)]
402#[serde(rename_all = "snake_case", tag = "action", content = "arguments")]
403pub enum DaemonAction {
404    NativeAgentHistory {
405        owner: String,
406        child: String,
407        before: Option<(u64, String)>,
408    },
409    Ping,
410    /// Frozen upgrade handshake, available from daemon protocol 33 onward.
411    /// Busy replies leave every operation and control channel running.
412    PrepareUpgrade,
413    /// Names the daemon-owned work that is currently holding an automatic
414    /// handoff open, for the CLI's wait notice.
415    ///
416    /// Added after protocol 33, so it is not part of the frozen management
417    /// subset. A daemon that predates it cannot deserialize the frame: it
418    /// fails the request and closes the connection. Callers must therefore
419    /// treat any failure, including a dropped connection or an error frame,
420    /// as "unknown" and say nothing about what the upgrade is waiting for.
421    UpgradeBlockers,
422    Status,
423    WebViewerAccess,
424    RecoverWebViewer(crate::web::WebViewerRecovery),
425    InspectWebListener,
426    ListWorkspaces,
427    CreateWorkspace {
428        name: String,
429    },
430    RenameWorkspace {
431        workspace_id: String,
432        name: String,
433    },
434    TouchWorkspace {
435        workspace_id: String,
436    },
437    CloseWorkspace {
438        workspace_id: String,
439    },
440    CancelWorkspaceClose {
441        workspace_id: String,
442    },
443    DeleteWorkspace {
444        workspace_id: String,
445    },
446    Attach {
447        client_id: String,
448        pid: u32,
449    },
450    Detach {
451        client_id: String,
452    },
453    PersistReadReceipt {
454        client_id: String,
455        workspace_id: String,
456        session_id: String,
457        through: u64,
458    },
459    PersistDetachedSessionState {
460        client_id: String,
461        workspace_id: String,
462        session_id: String,
463        through: u64,
464        owner_pid: u32,
465        draft: mj_core::storage::DetachedSessionDraft,
466    },
467    SaveActiveReview {
468        session_id: String,
469        review: mj_core::storage::StoredReview,
470    },
471    ClearActiveReview {
472        session_id: String,
473    },
474    SaveWorkspacePaneSizes {
475        workspace_id: String,
476        sizes: mj_core::workspace::PaneSizes,
477    },
478    SaveWorkspaceLayout {
479        workspace_id: String,
480        layout: mj_core::workspace::ConversationLayout,
481    },
482    PersistImportedSession {
483        session: Box<SessionRecord>,
484    },
485    SetSessionTitle {
486        session_id: String,
487        title: String,
488    },
489    SetSessionWorkspace {
490        session_id: String,
491        workspace_id: String,
492    },
493    SetSessionContainerSettings {
494        session_id: String,
495        cpus: Option<String>,
496        memory: Option<String>,
497        mounts: Vec<AdditionalMount>,
498        mount_history: Vec<PathBuf>,
499    },
500    SetSessionAcpTitle {
501        session_id: String,
502        title: Option<String>,
503    },
504    MarkSessionTargetMissing {
505        session_id: String,
506        detail: String,
507        updated_at: String,
508    },
509    CheckpointSession {
510        session_id: String,
511    },
512    /// Search the user's SessionWiki index. An empty query lists the most
513    /// recent sessions.
514    ProjectCatalog {
515        refresh: bool,
516        retry: bool,
517    },
518    CreateProject {
519        sources: Vec<String>,
520    },
521    WikiSearch {
522        query: String,
523        limit: usize,
524    },
525    /// The live sessions whose user or agent messages contain a query,
526    /// ignoring case. Tool calls and tool output do not count.
527    SessionTextSearch {
528        query: String,
529    },
530    /// The markdown briefing for one indexed session.
531    WikiBrief {
532        wiki_id: String,
533        max_chars: usize,
534    },
535    /// The passages of one indexed session that match a query, with context.
536    WikiHits {
537        wiki_id: String,
538        query: String,
539        context_messages: usize,
540        per_message_chars: usize,
541    },
542    /// What one indexed session is, and what continuing it would mean.
543    WikiSession {
544        wiki_id: String,
545    },
546    /// Start a new session from an archived one's transcript.
547    WikiRestore(WikiRestoreRequest),
548    ScanRecovery {
549        all_instances: bool,
550    },
551    AdoptRecovery {
552        session_id: String,
553        target_id: String,
554        profile: Option<String>,
555        bundle: Option<String>,
556        all_instances: bool,
557    },
558    DestroyRecovery {
559        session_id: String,
560        target_id: String,
561        confirmation: String,
562        all_instances: bool,
563    },
564    Snapshot {
565        workspace_id: String,
566    },
567    RuntimeChanges {
568        cursor: Option<crate::runtime_feed::RuntimeCursor>,
569        wait: bool,
570    },
571    /// The sessions the resume dialog lists. The runtime feed carries only
572    /// live sessions, so the dialog asks for these when it opens.
573    ResumeCandidates,
574    /// One session's whole record, live or stopped. The resume wizard reads
575    /// it for the row picked in the resume dialog.
576    SessionRecord {
577        session_id: String,
578    },
579    /// The session `mj go` opens: the remembered one while it is still
580    /// eligible, otherwise the most recently updated eligible session of the
581    /// workspace, live or stopped.
582    GoStartupSession {
583        workspace_id: String,
584        last_session_id: Option<String>,
585    },
586    /// Ask the daemon's quota poller to probe now. The daemon is the only
587    /// process that probes; the result arrives in the runtime feed. Added in
588    /// protocol 43.
589    RefreshQuota,
590    /// Discover eligible subagent models in the daemon, which owns the cache.
591    SubagentOptions {
592        profile: String,
593        model: Option<String>,
594        /// Setup can probe an unsaved draft without changing live settings.
595        #[serde(default, skip_serializing_if = "Option::is_none")]
596        config: Option<Box<Config>>,
597    },
598    RenameProfile {
599        old_id: String,
600        new_id: String,
601    },
602    RenameTarget {
603        old_id: String,
604        new_id: String,
605    },
606    SubmitSessionCommand {
607        #[serde(default)]
608        inherited_draft: Option<String>,
609        session_id: String,
610        command_id: String,
611        command: RelayCommand,
612    },
613    /// Deliver a prompt typed while a session was still starting, once the
614    /// daemon sees that session's harness become ready. The daemon owns the
615    /// wait, so the prompt arrives whether or not this client is still
616    /// running or still showing that session.
617    QueueStartupPrompt {
618        session_id: String,
619        text: String,
620        /// The saved draft text this prompt was typed from, if the client
621        /// also persisted it. Cleared after a successful submit so the
622        /// delivered prompt does not reappear as a draft.
623        #[serde(default)]
624        inherited_draft: Option<String>,
625    },
626    /// Take back a prompt queued with `QueueStartupPrompt` that the daemon has
627    /// not started delivering, so the person can edit it. The daemon cancels
628    /// the newest queued prompt with this exact text and replies
629    /// `PromptWithdrawn(true)`, or `PromptWithdrawn(false)` when none is
630    /// waiting because delivery already started.
631    WithdrawStartupPrompt {
632        session_id: String,
633        text: String,
634    },
635    SyncSession {
636        session_id: String,
637    },
638    RespondElicitation {
639        session_id: String,
640        elicitation_id: String,
641        response: ElicitationResponse,
642    },
643    StopBackgroundTask {
644        session_id: String,
645        background_task_id: String,
646    },
647    /// Drive a session's second-opinion reviewer. The reviewer is a sidecar of
648    /// the session's worker, so it travels the session's own relay rather than
649    /// becoming a session of its own here.
650    ReviewerAction {
651        session_id: String,
652        /// Which reviewing role the action drives; absent means the default
653        /// one, which is what plan review uses.
654        #[serde(default, skip_serializing_if = "Option::is_none")]
655        role: Option<String>,
656        action: crate::session::ReviewerAction,
657    },
658    /// Review the turn this session just finished, on a surface's request.
659    StartTurnReview {
660        session_id: String,
661    },
662    /// Forward, dismiss, or cancel the open review.
663    ResolveTurnReview {
664        session_id: String,
665        resolution: Resolution,
666    },
667    SuspendSession {
668        session_id: String,
669        #[serde(default)]
670        acknowledge_unpublished_work: bool,
671    },
672    StartCreateSession(CreateSessionRequest),
673    WaitCreateSession {
674        session_id: String,
675    },
676    ResumeSession(ResumeSessionRequest),
677    PrepareMoveSession(MoveSelection),
678    MoveSession(MoveSessionRequest),
679    MoveSources {
680        session_id: String,
681        cleanup_operation_id: Option<String>,
682    },
683    DiscardSinceCheckpoint {
684        session_id: String,
685        checkpoint: mj_core::state::CheckpointMetadata,
686    },
687    DestroyStoppedSession {
688        session_id: String,
689        /// Whether to delete the session's managed git branch as well. The
690        /// branch can hold work the user still wants, so destroying keeps it
691        /// unless the request asks for the deletion.
692        delete_branch: bool,
693    },
694    ForceDestroySession {
695        session_id: String,
696        /// See [`DaemonAction::DestroyStoppedSession`].
697        delete_branch: bool,
698    },
699    CancelLifecycle {
700        session_id: String,
701    },
702    RecoverDraft {
703        draft_id: String,
704    },
705    Stop,
706}
707
708#[derive(Debug, Serialize, Deserialize)]
709#[serde(deny_unknown_fields)]
710pub struct RequestEnvelope {
711    pub protocol_version: u32,
712    pub request_id: u64,
713    pub token: String,
714    pub action: DaemonAction,
715}
716
717/// The daemon answered a request with an error. Unlike a lost connection,
718/// this is a complete round trip: the daemon read the request and said why it
719/// did not carry it out.
720#[derive(Debug)]
721pub struct DaemonRefusal(pub String);
722
723impl std::fmt::Display for DaemonRefusal {
724    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
725        f.write_str(&self.0)
726    }
727}
728
729impl std::error::Error for DaemonRefusal {}
730
731impl DaemonRefusal {
732    /// Whether the daemon itself could not confirm delivery. Its message is
733    /// then led by `session::DeliveryUnconfirmed`.
734    #[must_use]
735    pub fn delivery_unconfirmed(&self) -> bool {
736        self.0
737            .starts_with(&crate::session::DeliveryUnconfirmed.to_string())
738    }
739}
740
741#[derive(Debug, Serialize, Deserialize)]
742#[serde(deny_unknown_fields)]
743pub struct ResponseEnvelope {
744    pub protocol_version: u32,
745    pub request_id: u64,
746    pub result: std::result::Result<DaemonReply, String>,
747}
748
749#[derive(Debug, Clone, Serialize, Deserialize)]
750#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
751pub enum DaemonReply {
752    NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
753    Pong,
754    UpgradePending,
755    /// Labels for the work holding an automatic handoff open, empty when
756    /// nothing is. Answers `DaemonAction::UpgradeBlockers`.
757    UpgradeBlockers(Vec<String>),
758    Status(DaemonStatus),
759    WebViewerAccess(crate::web::WebViewerAccess),
760    WebListeners(Vec<crate::web::WebListenerProcess>),
761    SubagentOptions(mj_core::subagent::SubagentOptions),
762    Workspaces(Vec<WorkspaceListing>),
763    Workspace(WorkspaceRecord),
764    Snapshot(WorkspaceSnapshot),
765    RuntimeChanges(Box<crate::runtime_feed::RuntimeFrame>),
766    ResumeCandidates(Box<ResumeCandidates>),
767    GoStartupSession(Option<Box<SessionRecord>>),
768    SessionRecord(Option<Box<SessionRecord>>),
769    /// Transport fragments of one chunked reply (see
770    /// [`DaemonReply::is_chunked`]); never exposed to consumers.
771    ReplyChunk {
772        bytes: Vec<u8>,
773        finished: bool,
774    },
775    RegisteredSession(Box<RegisteredSession>),
776    MovePreparation(Box<MovePreparation>),
777    MoveOutcome(MoveOutcome),
778    MoveSources(Vec<mj_core::move_workspace::RetainedMoveSource>),
779    Ordinal(u64),
780    Text(String),
781    OptionalSessionState(Option<SessionState>),
782    Checkpoint(mj_core::state::CheckpointMetadata),
783    RecoveryScan(mj_core::state::RecoveryScan),
784    WikiRows(WikiSearchPage),
785    ProjectCatalog(mj_core::project_catalog::ProjectCatalogView),
786    ProjectCreated(mj_core::project_catalog::SavedProject),
787    SessionTextMatches(Vec<SessionTextMatch>),
788    WikiHits(Option<WikiHitTranscript>),
789    WikiSession(Option<Box<WikiSessionInfo>>),
790    Reviewer(Box<crate::session::ReviewerOutcome>),
791    /// Whether a queued startup prompt was withdrawn before delivery.
792    PromptWithdrawn(bool),
793    Done,
794}
795
796impl DaemonReply {
797    /// Replies that can outgrow one frame. The daemon sends them as
798    /// [`DaemonReply::ReplyChunk`] fragments and the client reassembles them.
799    pub fn is_chunked(&self) -> bool {
800        matches!(self, Self::RuntimeChanges(_) | Self::ResumeCandidates(_))
801    }
802}
803
804/// What the resume dialog lists, read on demand.
805#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
806#[serde(deny_unknown_fields)]
807pub struct ResumeCandidates {
808    /// One row for every inactive session that is not a sub-agent.
809    pub candidates: Vec<ResumeCandidate>,
810    /// Durable moves of those sessions, for "move needs recovery" marks.
811    pub moves: Vec<MoveOperation>,
812    /// Native sessions that some record, live or not, already holds. The
813    /// import list must not offer them again.
814    pub adopted_native_sessions: Vec<(mj_core::config::HarnessKind, String)>,
815    /// Local checkouts that Mjolnir created for its sessions. Every native
816    /// thread that ran inside one belongs to that session.
817    pub local_checkout_roots: Vec<PathBuf>,
818}
819
820/// One row of the resume dialog's Mjolnir tab: what the row shows, without
821/// the rest of the record. A session's title can be tens of kilobytes, and
822/// most of a store's history is never resumed, so the resume wizard fetches
823/// the whole record (`DaemonAction::SessionRecord`) only for the row picked.
824#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
825#[serde(deny_unknown_fields)]
826pub struct ResumeCandidate {
827    pub session_id: String,
828    pub state: SessionState,
829    pub last_profile: String,
830    /// The listed title, cut to [`ResumeCandidate::TITLE_CHARS`].
831    pub title: String,
832    /// Where the session ran, as the session list names it.
833    pub origin: String,
834    pub project: String,
835    pub has_checkpoint: bool,
836    /// The checkpoint's creation time, or the record's last update without
837    /// one; `None` when neither parses.
838    pub last_activity_ms: Option<i64>,
839    pub publication: Option<PublicationState>,
840    /// The checkpoint archive whose size the row shows. Only a stopped
841    /// record is described by its checkpoint.
842    pub checkpoint_archive: Option<PathBuf>,
843    /// Whether destroying the session can also delete its worktree branch.
844    pub worktree_checkout: bool,
845}
846
847impl ResumeCandidate {
848    /// Longer than any row draws it.
849    pub const TITLE_CHARS: usize = 200;
850
851    /// The row for `record`. The daemon and the dashboard both build rows
852    /// here, so a row reads the same whichever of them held the record.
853    pub fn of(record: &SessionRecord, config: &Config) -> Self {
854        let timestamp_ms = |timestamp: &str| {
855            chrono::DateTime::parse_from_rfc3339(timestamp)
856                .ok()
857                .map(|parsed| parsed.timestamp_millis())
858        };
859        Self {
860            session_id: record.id.clone(),
861            state: record.state,
862            last_profile: record.last_profile.clone(),
863            title: record
864                .listed_title()
865                .chars()
866                .take(Self::TITLE_CHARS)
867                .collect(),
868            origin: record.project_target(config, &record.target_template_id),
869            project: record.project_name(config),
870            has_checkpoint: record.checkpoint.is_some(),
871            last_activity_ms: record
872                .checkpoint
873                .as_ref()
874                .and_then(|checkpoint| timestamp_ms(&checkpoint.created_at))
875                .or_else(|| timestamp_ms(&record.updated_at)),
876            publication: record.publication_state(),
877            checkpoint_archive: record
878                .checkpoint
879                .as_ref()
880                .filter(|_| record.state == SessionState::Stopped)
881                .map(|checkpoint| checkpoint.archive_path.clone()),
882            worktree_checkout: record
883                .managed_worktree
884                .as_ref()
885                .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Worktree),
886        }
887    }
888}
889
890#[derive(Debug, Clone, Serialize, Deserialize)]
891#[serde(deny_unknown_fields)]
892pub struct DaemonStatus {
893    pub pid: u32,
894    pub started_at: String,
895    pub build_version: String,
896    pub attached_clients: usize,
897    pub phone_status: WebViewerStatus,
898}
899
900#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
901#[serde(rename_all = "snake_case", tag = "state")]
902pub enum WebViewerStatus {
903    Disabled,
904    Starting,
905    Ready {
906        viewer_url: String,
907        viewer_code: String,
908        qr_login_url: Option<String>,
909        fallback_reason: Option<String>,
910    },
911    Stopped,
912    Error {
913        message: String,
914    },
915}
916
917impl std::fmt::Debug for WebViewerStatus {
918    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
919        match self {
920            Self::Ready {
921                viewer_url,
922                viewer_code,
923                fallback_reason,
924                ..
925            } => formatter
926                .debug_struct("Ready")
927                .field("viewer_url", viewer_url)
928                .field("viewer_code", viewer_code)
929                .field("qr_login_url", &"[redacted]")
930                .field("fallback_reason", fallback_reason)
931                .finish(),
932            Self::Disabled => formatter.write_str("Disabled"),
933            Self::Starting => formatter.write_str("Starting"),
934            Self::Stopped => formatter.write_str("Stopped"),
935            Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
936        }
937    }
938}
939
940impl std::fmt::Display for WebViewerStatus {
941    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
942        match self {
943            Self::Disabled => formatter.write_str("disabled"),
944            Self::Starting => formatter.write_str("starting"),
945            Self::Stopped => formatter.write_str("stopped unexpectedly"),
946            Self::Error { message } => write!(formatter, "error: {message}"),
947            Self::Ready {
948                viewer_url,
949                viewer_code,
950                fallback_reason,
951                ..
952            } => {
953                write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
954                if let Some(reason) = fallback_reason {
955                    write!(
956                        formatter,
957                        "; local only because Tailscale HTTPS is unavailable: {reason}"
958                    )?;
959                }
960                Ok(())
961            }
962        }
963    }
964}
965
966/// Whether a non-child process has exited but not yet been reaped.
967///
968/// A zombie still answers `kill(pid, 0)`, because its process-table entry
969/// survives until its parent waits for it — so an existence probe alone calls
970/// it alive forever and anything waiting for it to leave waits forever. That is
971/// exactly the shape of `mj daemon restart` refusing to restart a daemon that
972/// had already stopped: `spawn_detached` used to leave the daemon a child of a
973/// long-lived Mjolnir process that never reaped it. It now double-forks, so the
974/// daemon is init's to reap, but any other unreaped child of this process would
975/// look the same, and the check stays cheap.
976///
977/// Treating a zombie as gone is also safe in the direction that matters: a
978/// zombie's PID cannot be reused until it is reaped, so nothing else can be
979/// occupying that number while this returns true.
980///
981/// On Linux this reads the one `/proc/<pid>/stat` file. `sysinfo` scans every
982/// process on the machine for a single-PID refresh, so it is not used here:
983/// this probe runs every 250 ms while an upgrade waits on an older daemon.
984#[cfg(target_os = "linux")]
985pub fn process_is_zombie(pid: u32) -> bool {
986    let Ok(stat) = std::fs::read(format!("/proc/{pid}/stat")) else {
987        return false;
988    };
989    // The command name in parentheses may itself contain ") ", so the state
990    // is the first field after the last closing parenthesis.
991    stat.iter()
992        .rposition(|byte| *byte == b')')
993        .and_then(|end| stat.get(end + 2))
994        .is_some_and(|state| *state == b'Z')
995}
996
997/// On macOS this asks the kernel for the one process's BSD info. `sysinfo`
998/// cannot be used: its lookups skip zombies, so it reports an exited, unreaped
999/// process as absent rather than as a zombie, and `kill(pid, 0)` then calls it
1000/// alive forever. A nonzero `arg` to `PROC_PIDTBSDINFO` makes the kernel
1001/// include zombies.
1002#[cfg(target_os = "macos")]
1003pub fn process_is_zombie(pid: u32) -> bool {
1004    let Ok(raw_pid) = libc::c_int::try_from(pid) else {
1005        return false;
1006    };
1007    let size = std::mem::size_of::<libc::proc_bsdinfo>();
1008    // SAFETY: proc_bsdinfo is plain old data, so all-zero is a valid value.
1009    let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
1010    // SAFETY: the buffer is a writable proc_bsdinfo of exactly `size` bytes.
1011    let written = unsafe {
1012        libc::proc_pidinfo(
1013            raw_pid,
1014            libc::PROC_PIDTBSDINFO,
1015            1,
1016            (&mut info as *mut libc::proc_bsdinfo).cast(),
1017            size as libc::c_int,
1018        )
1019    };
1020    usize::try_from(written).is_ok_and(|written| written == size) && info.pbi_status == libc::SZOMB
1021}
1022
1023#[cfg(all(unix, not(target_os = "linux"), not(target_os = "macos")))]
1024pub fn process_is_zombie(pid: u32) -> bool {
1025    let pid = sysinfo::Pid::from_u32(pid);
1026    let mut system = sysinfo::System::new();
1027    system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
1028    system
1029        .process(pid)
1030        .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
1031}
1032
1033/// Wait for a process to leave, within [`STOP_TIMEOUT`].
1034///
1035/// The error says the process was still running rather than that it "did not
1036/// stop": a daemon that is still winding down has not refused, and the two
1037/// read very differently to somebody deciding whether to reach for a kill.
1038pub async fn wait_for_exit(pid: u32) -> Result<()> {
1039    wait_for_exit_within(pid, STOP_TIMEOUT).await
1040}
1041
1042async fn wait_for_exit_within(pid: u32, timeout: Duration) -> Result<()> {
1043    let deadline = Instant::now() + timeout;
1044    while daemon_process_is_alive(pid) {
1045        ensure!(Instant::now() < deadline, "process {pid} is still running");
1046        tokio::time::sleep(RETRY_DELAY).await;
1047    }
1048    Ok(())
1049}
1050
1051/// Whether a daemon that Mjolnir launched in its own process group is alive.
1052///
1053/// Reaping is deliberately confined to this daemon-specific path. Attachment
1054/// PIDs are merely observations and may alias unrelated children owned by this
1055/// process, so their liveness probe below must never call `waitpid`.
1056pub fn daemon_process_is_alive(pid: u32) -> bool {
1057    #[cfg(unix)]
1058    {
1059        if pid == 0 {
1060            return false;
1061        }
1062        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
1063            return false;
1064        };
1065        let mut status = 0;
1066        // SAFETY: `status` is writable for the call and WNOHANG never blocks.
1067        // A non-child fails with ECHILD without changing any process state.
1068        let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
1069        if waited == raw_pid {
1070            return false;
1071        }
1072        if waited == 0 {
1073            return true;
1074        }
1075        let wait_error = std::io::Error::last_os_error();
1076        if wait_error.raw_os_error() != Some(libc::ECHILD) {
1077            return true;
1078        }
1079
1080        #[cfg(target_os = "macos")]
1081        return owned_daemon_group_is_alive(raw_pid);
1082
1083        #[cfg(not(target_os = "macos"))]
1084        process_is_alive(pid)
1085    }
1086    #[cfg(not(unix))]
1087    process_is_alive(pid)
1088}
1089
1090#[cfg(target_os = "macos")]
1091pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
1092    // `spawn_detached` makes the daemon a process-group leader. Darwin
1093    // excludes zombies from group signal probes: ESRCH means the group is gone
1094    // and EPERM means only exiting members remain. The latter is safe here
1095    // because this is a group we created for our own same-user child, not an
1096    // arbitrary process group.
1097    // SAFETY: signal 0 is only an existence probe, and the negative PID targets
1098    // the daemon-owned group rather than another process.
1099    if unsafe { libc::kill(-pid, 0) } == 0 {
1100        return true;
1101    }
1102    let error = std::io::Error::last_os_error();
1103    !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
1104}
1105
1106pub fn process_is_alive(pid: u32) -> bool {
1107    #[cfg(unix)]
1108    {
1109        if pid == 0 {
1110            return false;
1111        }
1112        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
1113            return false;
1114        };
1115        // SAFETY: kill(pid, 0) sends no signal and is the standard existence
1116        // probe. EPERM still means the process exists.
1117        let result = unsafe { libc::kill(raw_pid, 0) };
1118        let exists =
1119            result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
1120        exists && !process_is_zombie(pid)
1121    }
1122    #[cfg(not(unix))]
1123    {
1124        let process_id = sysinfo::Pid::from_u32(pid);
1125        let mut system = sysinfo::System::new();
1126        system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[process_id]), true);
1127        system.process(process_id).is_some()
1128    }
1129}
1130
1131pub fn read_metadata() -> Result<DaemonMetadata> {
1132    let metadata = read_metadata_any()?;
1133    ensure!(
1134        metadata.protocol_version == PROTOCOL_VERSION,
1135        "daemon protocol {} is incompatible with client protocol {}",
1136        metadata.protocol_version,
1137        PROTOCOL_VERSION
1138    );
1139    Ok(metadata)
1140}
1141
1142/// No daemon has published its endpoint: `daemon.json` does not exist.
1143///
1144/// This is the ordinary stopped state, not a failure. A daemon removes the
1145/// file when it stops, so callers show "not running" rather than the raw
1146/// file error.
1147#[derive(Debug, Clone, PartialEq, Eq)]
1148pub struct DaemonNotRunning {
1149    /// Where the endpoint would be, which also names the instance.
1150    pub metadata_path: std::path::PathBuf,
1151}
1152
1153/// What a caller shows for [`DaemonNotRunning`].
1154pub const DAEMON_NOT_RUNNING_MESSAGE: &str = "the Mjolnir daemon is not running";
1155
1156impl std::fmt::Display for DaemonNotRunning {
1157    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1158        formatter.write_str(DAEMON_NOT_RUNNING_MESSAGE)
1159    }
1160}
1161
1162impl std::error::Error for DaemonNotRunning {}
1163
1164/// The stopped state, if `error` reports it anywhere in its chain.
1165pub fn daemon_not_running(error: &anyhow::Error) -> Option<&DaemonNotRunning> {
1166    error
1167        .chain()
1168        .find_map(|cause| cause.downcast_ref::<DaemonNotRunning>())
1169}
1170
1171pub fn read_metadata_any() -> Result<DaemonMetadata> {
1172    read_metadata_at(&metadata_path())
1173}
1174
1175fn read_metadata_at(path: &std::path::Path) -> Result<DaemonMetadata> {
1176    let body = match fs::read(path) {
1177        Ok(body) => body,
1178        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1179            return Err(DaemonNotRunning {
1180                metadata_path: path.to_owned(),
1181            }
1182            .into());
1183        }
1184        Err(error) => return Err(error).with_context(|| format!("read {}", path.display())),
1185    };
1186    let metadata: DaemonMetadata =
1187        serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1188    Ok(metadata)
1189}
1190
1191pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
1192    let body = serde_json::to_vec(value)?;
1193    write_encoded_frame(stream, &body).await
1194}
1195
1196pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
1197    ensure!(
1198        body.len() <= MAX_FRAME_BYTES,
1199        "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
1200        body.len()
1201    );
1202    stream.write_u32(body.len() as u32).await?;
1203    stream.write_all(body).await?;
1204    stream.flush().await?;
1205    Ok(())
1206}
1207
1208pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
1209    let length = stream.read_u32().await? as usize;
1210    ensure!(
1211        length <= MAX_FRAME_BYTES,
1212        "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
1213    );
1214    let mut body = vec![0_u8; length];
1215    stream.read_exact(&mut body).await?;
1216    serde_json::from_slice(&body).context("decode daemon frame")
1217}
1218
1219/// Read one reply, reassembling a chunked one before exposing it. All
1220/// fragments belong to the same response; disconnects discard partial state.
1221pub async fn read_response(stream: &mut TcpStream) -> Result<ResponseEnvelope> {
1222    let mut response: ResponseEnvelope = read_frame(stream).await?;
1223    if !matches!(response.result, Ok(DaemonReply::ReplyChunk { .. })) {
1224        return Ok(response);
1225    }
1226    let protocol_version = response.protocol_version;
1227    let request_id = response.request_id;
1228    let mut body = Vec::new();
1229    loop {
1230        ensure!(
1231            response.protocol_version == protocol_version && response.request_id == request_id,
1232            "daemon crossed reply fragment identities"
1233        );
1234        let Ok(DaemonReply::ReplyChunk { bytes, finished }) = response.result else {
1235            bail!("daemon interrupted a chunked reply");
1236        };
1237        ensure!(!bytes.is_empty(), "empty reply fragment");
1238        body.extend(bytes);
1239        if finished {
1240            break;
1241        }
1242        response = read_frame(stream).await?;
1243    }
1244    let reply: DaemonReply = tokio::task::spawn_blocking(move || serde_json::from_slice(&body))
1245        .await
1246        .context("reply decoder task failed")??;
1247    ensure!(
1248        reply.is_chunked(),
1249        "daemon chunked a reply that is never chunked"
1250    );
1251    Ok(ResponseEnvelope {
1252        protocol_version,
1253        request_id,
1254        result: Ok(reply),
1255    })
1256}
1257
1258pub struct DaemonClient {
1259    metadata: DaemonMetadata,
1260    stream: TcpStream,
1261    next_request_id: u64,
1262}
1263
1264impl DaemonClient {
1265    pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
1266        let stream =
1267            tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
1268                .await
1269                .context("time out connecting to Mjolnir daemon")??;
1270        Ok(Self {
1271            metadata,
1272            stream,
1273            next_request_id: 1,
1274        })
1275    }
1276
1277    /// Speak the daemon's advertised dialect, not this build's: management
1278    /// requests must reach daemons of any protocol version, and the frozen
1279    /// subset encodes identically across all of them.
1280    pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1281        self.request_with_reconnect(action, || async {
1282            loop {
1283                if let Ok(client) = connect_existing().await {
1284                    return Ok(client);
1285                }
1286                tokio::time::sleep(Duration::from_millis(250)).await;
1287            }
1288        })
1289        .await
1290    }
1291
1292    async fn request_with_reconnect<F, Fut>(
1293        &mut self,
1294        action: DaemonAction,
1295        mut reconnect: F,
1296    ) -> Result<DaemonReply>
1297    where
1298        F: FnMut() -> Fut,
1299        Fut: Future<Output = Result<Self>>,
1300    {
1301        loop {
1302            let response = self.request_once(action.clone()).await?;
1303            if matches!(response, DaemonReply::UpgradePending)
1304                && !matches!(action, DaemonAction::PrepareUpgrade)
1305            {
1306                // Only an explicit refusal guarantees non-admission. Never
1307                // replay arbitrary mutations after a lost acknowledgement.
1308                tokio::time::sleep(Duration::from_millis(100)).await;
1309                *self = reconnect().await?;
1310            } else {
1311                return Ok(response);
1312            }
1313        }
1314    }
1315
1316    async fn request_once(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1317        let protocol_version = self.metadata.protocol_version;
1318        let request_id = self.next_request_id;
1319        self.next_request_id += 1;
1320        write_frame(
1321            &mut self.stream,
1322            &RequestEnvelope {
1323                protocol_version,
1324                request_id,
1325                token: self.metadata.token.clone(),
1326                action,
1327            },
1328        )
1329        .await?;
1330        let response = read_response(&mut self.stream).await?;
1331        ensure!(
1332            response.protocol_version == protocol_version,
1333            "daemon changed protocol"
1334        );
1335        ensure!(
1336            response.request_id == request_id,
1337            "daemon crossed request IDs"
1338        );
1339        response
1340            .result
1341            .map_err(|message| anyhow::Error::new(DaemonRefusal(message)))
1342    }
1343
1344    pub async fn status(&mut self) -> Result<DaemonStatus> {
1345        match self.request(DaemonAction::Status).await? {
1346            DaemonReply::Status(status) => Ok(status),
1347            reply => bail!("unexpected daemon status reply {reply:?}"),
1348        }
1349    }
1350
1351    pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
1352        match self.request(DaemonAction::WebViewerAccess).await? {
1353            DaemonReply::WebViewerAccess(access) => Ok(access),
1354            reply => bail!("unexpected web viewer reply {reply:?}"),
1355        }
1356    }
1357
1358    pub async fn subagent_options(
1359        &mut self,
1360        profile: String,
1361        model: Option<String>,
1362        config: Option<Config>,
1363    ) -> Result<mj_core::subagent::SubagentOptions> {
1364        match self
1365            .request(DaemonAction::SubagentOptions {
1366                profile,
1367                model,
1368                config: config.map(Box::new),
1369            })
1370            .await?
1371        {
1372            DaemonReply::SubagentOptions(options) => Ok(options),
1373            reply => bail!("unexpected subagent options reply {reply:?}"),
1374        }
1375    }
1376
1377    pub async fn recover_web_viewer(
1378        &mut self,
1379        action: crate::web::WebViewerRecovery,
1380    ) -> Result<()> {
1381        match self.request(DaemonAction::RecoverWebViewer(action)).await? {
1382            DaemonReply::Done => Ok(()),
1383            reply => bail!("unexpected web viewer recovery reply {reply:?}"),
1384        }
1385    }
1386
1387    pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
1388        match self.request(DaemonAction::InspectWebListener).await? {
1389            DaemonReply::WebListeners(processes) => Ok(processes),
1390            reply => bail!("unexpected listener inspection reply {reply:?}"),
1391        }
1392    }
1393
1394    pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
1395        match self.request(DaemonAction::ListWorkspaces).await? {
1396            DaemonReply::Workspaces(workspaces) => Ok(workspaces),
1397            reply => bail!("unexpected daemon workspace reply {reply:?}"),
1398        }
1399    }
1400
1401    /// Ask the daemon to probe every profile that is not on hold. Returns when
1402    /// the daemon has accepted the request, not when the probes finish.
1403    pub async fn refresh_quota(&mut self) -> Result<()> {
1404        match self.request(DaemonAction::RefreshQuota).await? {
1405            DaemonReply::Done => Ok(()),
1406            reply => bail!("unexpected refresh-quota reply {reply:?}"),
1407        }
1408    }
1409
1410    pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
1411        match self
1412            .request(DaemonAction::RenameProfile { old_id, new_id })
1413            .await?
1414        {
1415            DaemonReply::Done => Ok(()),
1416            reply => bail!("unexpected rename-profile reply {reply:?}"),
1417        }
1418    }
1419
1420    pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1421        match self
1422            .request(DaemonAction::RenameTarget { old_id, new_id })
1423            .await?
1424        {
1425            DaemonReply::Done => Ok(()),
1426            reply => bail!("unexpected rename-target reply {reply:?}"),
1427        }
1428    }
1429
1430    pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1431        match self.request(DaemonAction::CreateWorkspace { name }).await? {
1432            DaemonReply::Workspace(workspace) => Ok(workspace),
1433            reply => bail!("unexpected create-workspace reply {reply:?}"),
1434        }
1435    }
1436
1437    pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1438        match self
1439            .request(DaemonAction::RenameWorkspace { workspace_id, name })
1440            .await?
1441        {
1442            DaemonReply::Done => Ok(()),
1443            reply => bail!("unexpected rename-workspace reply {reply:?}"),
1444        }
1445    }
1446
1447    pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1448        match self
1449            .request(DaemonAction::TouchWorkspace { workspace_id })
1450            .await?
1451        {
1452            DaemonReply::Done => Ok(()),
1453            reply => bail!("unexpected touch-workspace reply {reply:?}"),
1454        }
1455    }
1456
1457    pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1458        match self
1459            .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1460            .await?
1461        {
1462            DaemonReply::Done => Ok(()),
1463            reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1464        }
1465    }
1466
1467    pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1468        match self
1469            .request(DaemonAction::CloseWorkspace { workspace_id })
1470            .await?
1471        {
1472            DaemonReply::Done => Ok(()),
1473            reply => bail!("unexpected close workspace reply: {reply:?}"),
1474        }
1475    }
1476
1477    pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1478        match self
1479            .request(DaemonAction::DeleteWorkspace { workspace_id })
1480            .await?
1481        {
1482            DaemonReply::Done => Ok(()),
1483            reply => bail!("unexpected delete-workspace reply {reply:?}"),
1484        }
1485    }
1486
1487    pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1488        match self
1489            .request(DaemonAction::Attach { client_id, pid })
1490            .await?
1491        {
1492            DaemonReply::Done => Ok(()),
1493            reply => bail!("unexpected attach reply {reply:?}"),
1494        }
1495    }
1496
1497    pub async fn detach(&mut self, client_id: String) -> Result<()> {
1498        match self.request(DaemonAction::Detach { client_id }).await? {
1499            DaemonReply::Done => Ok(()),
1500            reply => bail!("unexpected detach reply {reply:?}"),
1501        }
1502    }
1503
1504    pub async fn persist_read_receipt(
1505        &mut self,
1506        client_id: String,
1507        workspace_id: String,
1508        session_id: String,
1509        through: u64,
1510    ) -> Result<u64> {
1511        match self
1512            .request(DaemonAction::PersistReadReceipt {
1513                client_id,
1514                workspace_id,
1515                session_id,
1516                through,
1517            })
1518            .await?
1519        {
1520            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1521            reply => bail!("unexpected read-receipt reply {reply:?}"),
1522        }
1523    }
1524
1525    pub async fn persist_detached_session_state(
1526        &mut self,
1527        client_id: String,
1528        workspace_id: String,
1529        session_id: String,
1530        through: u64,
1531        owner_pid: u32,
1532        draft: mj_core::storage::DetachedSessionDraft,
1533    ) -> Result<()> {
1534        match self
1535            .request(DaemonAction::PersistDetachedSessionState {
1536                client_id,
1537                workspace_id,
1538                session_id,
1539                through,
1540                owner_pid,
1541                draft,
1542            })
1543            .await?
1544        {
1545            DaemonReply::Done => Ok(()),
1546            reply => bail!("unexpected detached-session-state reply {reply:?}"),
1547        }
1548    }
1549
1550    pub async fn save_active_review(
1551        &mut self,
1552        session_id: String,
1553        review: mj_core::storage::StoredReview,
1554    ) -> Result<()> {
1555        match self
1556            .request(DaemonAction::SaveActiveReview { session_id, review })
1557            .await?
1558        {
1559            DaemonReply::Done => Ok(()),
1560            reply => bail!("unexpected save-review reply {reply:?}"),
1561        }
1562    }
1563
1564    pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1565        match self
1566            .request(DaemonAction::ClearActiveReview { session_id })
1567            .await?
1568        {
1569            DaemonReply::Done => Ok(()),
1570            reply => bail!("unexpected clear-review reply {reply:?}"),
1571        }
1572    }
1573
1574    pub async fn save_workspace_pane_sizes(
1575        &mut self,
1576        workspace_id: String,
1577        sizes: mj_core::workspace::PaneSizes,
1578    ) -> Result<()> {
1579        match self
1580            .request(DaemonAction::SaveWorkspacePaneSizes {
1581                workspace_id,
1582                sizes,
1583            })
1584            .await?
1585        {
1586            DaemonReply::Done => Ok(()),
1587            reply => bail!("unexpected pane-size save reply {reply:?}"),
1588        }
1589    }
1590
1591    pub async fn save_workspace_layout(
1592        &mut self,
1593        workspace_id: String,
1594        layout: mj_core::workspace::ConversationLayout,
1595    ) -> Result<()> {
1596        match self
1597            .request(DaemonAction::SaveWorkspaceLayout {
1598                workspace_id,
1599                layout,
1600            })
1601            .await?
1602        {
1603            DaemonReply::Done => Ok(()),
1604            reply => bail!("unexpected layout save reply {reply:?}"),
1605        }
1606    }
1607
1608    pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1609        match self
1610            .request(DaemonAction::PersistImportedSession {
1611                session: Box::new(session),
1612            })
1613            .await?
1614        {
1615            DaemonReply::Done => Ok(()),
1616            reply => bail!("unexpected imported-session reply {reply:?}"),
1617        }
1618    }
1619
1620    pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1621        match self
1622            .request(DaemonAction::SetSessionTitle { session_id, title })
1623            .await?
1624        {
1625            DaemonReply::Text(title) => Ok(title),
1626            reply => bail!("unexpected session-title reply {reply:?}"),
1627        }
1628    }
1629
1630    pub async fn set_session_workspace(
1631        &mut self,
1632        session_id: String,
1633        workspace_id: String,
1634    ) -> Result<()> {
1635        match self
1636            .request(DaemonAction::SetSessionWorkspace {
1637                session_id,
1638                workspace_id,
1639            })
1640            .await?
1641        {
1642            DaemonReply::Done => Ok(()),
1643            reply => bail!("unexpected session-workspace reply {reply:?}"),
1644        }
1645    }
1646
1647    pub async fn set_session_container_settings(
1648        &mut self,
1649        session_id: String,
1650        cpus: Option<String>,
1651        memory: Option<String>,
1652        mounts: Vec<AdditionalMount>,
1653        mount_history: Vec<PathBuf>,
1654    ) -> Result<()> {
1655        match self
1656            .request(DaemonAction::SetSessionContainerSettings {
1657                session_id,
1658                cpus,
1659                memory,
1660                mounts,
1661                mount_history,
1662            })
1663            .await?
1664        {
1665            DaemonReply::Done => Ok(()),
1666            reply => bail!("unexpected container-settings reply {reply:?}"),
1667        }
1668    }
1669
1670    pub async fn set_session_acp_title(
1671        &mut self,
1672        session_id: String,
1673        title: Option<String>,
1674    ) -> Result<()> {
1675        match self
1676            .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1677            .await?
1678        {
1679            DaemonReply::Done => Ok(()),
1680            reply => bail!("unexpected ACP-title reply {reply:?}"),
1681        }
1682    }
1683
1684    pub async fn mark_session_target_missing(
1685        &mut self,
1686        session_id: String,
1687        detail: String,
1688        updated_at: String,
1689    ) -> Result<Option<SessionState>> {
1690        match self
1691            .request(DaemonAction::MarkSessionTargetMissing {
1692                session_id,
1693                detail,
1694                updated_at,
1695            })
1696            .await?
1697        {
1698            DaemonReply::OptionalSessionState(state) => Ok(state),
1699            reply => bail!("unexpected target-missing reply {reply:?}"),
1700        }
1701    }
1702
1703    pub async fn checkpoint_session(
1704        &mut self,
1705        session_id: String,
1706    ) -> Result<mj_core::state::CheckpointMetadata> {
1707        match self
1708            .request(DaemonAction::CheckpointSession { session_id })
1709            .await?
1710        {
1711            DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1712            reply => bail!("unexpected checkpoint reply {reply:?}"),
1713        }
1714    }
1715
1716    /// Search the user's SessionWiki index, newest first when the query is
1717    /// empty and best match first otherwise. The reply carries the state of
1718    /// the index as well as the rows, so a caller can say the first build is
1719    /// still running.
1720    pub async fn project_catalog(
1721        &mut self,
1722        refresh: bool,
1723        retry: bool,
1724    ) -> Result<mj_core::project_catalog::ProjectCatalogView> {
1725        match self
1726            .request(DaemonAction::ProjectCatalog { refresh, retry })
1727            .await?
1728        {
1729            DaemonReply::ProjectCatalog(view) => Ok(view),
1730            other => bail!("unexpected project catalog reply: {other:?}"),
1731        }
1732    }
1733
1734    pub async fn create_project(
1735        &mut self,
1736        sources: Vec<String>,
1737    ) -> Result<mj_core::project_catalog::SavedProject> {
1738        match self
1739            .request(DaemonAction::CreateProject { sources })
1740            .await?
1741        {
1742            DaemonReply::ProjectCreated(project) => Ok(project),
1743            other => bail!("unexpected project creation reply: {other:?}"),
1744        }
1745    }
1746
1747    pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1748        match self
1749            .request(DaemonAction::WikiSearch { query, limit })
1750            .await?
1751        {
1752            DaemonReply::WikiRows(page) => Ok(page),
1753            reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1754        }
1755    }
1756
1757    /// The live sessions whose user or agent messages contain `query`, and
1758    /// which of the two matched. Answers from the SessionWiki index, so a very
1759    /// new message can be missing until the next sync.
1760    pub async fn session_text_search(&mut self, query: String) -> Result<Vec<SessionTextMatch>> {
1761        match self
1762            .request(DaemonAction::SessionTextSearch { query })
1763            .await?
1764        {
1765            DaemonReply::SessionTextMatches(matches) => Ok(matches),
1766            reply => bail!("unexpected session text search reply {reply:?}"),
1767        }
1768    }
1769
1770    /// The markdown briefing for one indexed session.
1771    pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1772        match self
1773            .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1774            .await?
1775        {
1776            DaemonReply::Text(markdown) => Ok(markdown),
1777            reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1778        }
1779    }
1780
1781    /// The passages of one indexed session that match a query, each matching
1782    /// message with `context_messages` neighbours on either side and its text
1783    /// capped at `per_message_chars`. `None` when the index holds no session
1784    /// with that id.
1785    pub async fn wiki_hits(
1786        &mut self,
1787        wiki_id: String,
1788        query: String,
1789        context_messages: usize,
1790        per_message_chars: usize,
1791    ) -> Result<Option<WikiHitTranscript>> {
1792        match self
1793            .request(DaemonAction::WikiHits {
1794                wiki_id,
1795                query,
1796                context_messages,
1797                per_message_chars,
1798            })
1799            .await?
1800        {
1801            DaemonReply::WikiHits(transcript) => Ok(transcript),
1802            reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1803        }
1804    }
1805
1806    /// What the index knows about one session, or `None` when the index holds
1807    /// no session with that id.
1808    pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1809        match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1810            DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1811            reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1812        }
1813    }
1814
1815    /// Start a new session carrying a hand-off compacted from an archived one.
1816    /// It answers like any other session start: the record exists and is
1817    /// provisioning, and the hand-off follows once the harness is ready.
1818    pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1819        match self.request(DaemonAction::WikiRestore(request)).await? {
1820            DaemonReply::RegisteredSession(registered) => Ok(*registered),
1821            reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1822        }
1823    }
1824
1825    pub async fn scan_recovery(
1826        &mut self,
1827        all_instances: bool,
1828    ) -> Result<mj_core::state::RecoveryScan> {
1829        match self
1830            .request(DaemonAction::ScanRecovery { all_instances })
1831            .await?
1832        {
1833            DaemonReply::RecoveryScan(scan) => Ok(scan),
1834            reply => bail!("unexpected recovery-scan reply {reply:?}"),
1835        }
1836    }
1837
1838    pub async fn adopt_recovery(
1839        &mut self,
1840        session_id: String,
1841        target_id: String,
1842        profile: Option<String>,
1843        bundle: Option<String>,
1844        all_instances: bool,
1845    ) -> Result<()> {
1846        match self
1847            .request(DaemonAction::AdoptRecovery {
1848                session_id,
1849                target_id,
1850                profile,
1851                bundle,
1852                all_instances,
1853            })
1854            .await?
1855        {
1856            DaemonReply::Done => Ok(()),
1857            reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1858        }
1859    }
1860
1861    pub async fn destroy_recovery(
1862        &mut self,
1863        session_id: String,
1864        target_id: String,
1865        confirmation: String,
1866        all_instances: bool,
1867    ) -> Result<()> {
1868        match self
1869            .request(DaemonAction::DestroyRecovery {
1870                session_id,
1871                target_id,
1872                confirmation,
1873                all_instances,
1874            })
1875            .await?
1876        {
1877            DaemonReply::Done => Ok(()),
1878            reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1879        }
1880    }
1881
1882    pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1883        match self
1884            .request(DaemonAction::Snapshot { workspace_id })
1885            .await?
1886        {
1887            DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1888            reply => bail!("unexpected snapshot reply {reply:?}"),
1889        }
1890    }
1891
1892    pub async fn runtime_changes(
1893        &mut self,
1894        cursor: Option<crate::runtime_feed::RuntimeCursor>,
1895        wait: bool,
1896    ) -> Result<crate::runtime_feed::RuntimeFrame> {
1897        match self
1898            .request(DaemonAction::RuntimeChanges { cursor, wait })
1899            .await?
1900        {
1901            DaemonReply::RuntimeChanges(frame) => Ok(*frame),
1902            reply => bail!("unexpected runtime changes reply {reply:?}"),
1903        }
1904    }
1905
1906    pub async fn resume_candidates(&mut self) -> Result<ResumeCandidates> {
1907        match self.request(DaemonAction::ResumeCandidates).await? {
1908            DaemonReply::ResumeCandidates(candidates) => Ok(*candidates),
1909            reply => bail!("unexpected resume candidates reply {reply:?}"),
1910        }
1911    }
1912
1913    pub async fn session_record(&mut self, session_id: String) -> Result<Option<SessionRecord>> {
1914        match self
1915            .request(DaemonAction::SessionRecord { session_id })
1916            .await?
1917        {
1918            DaemonReply::SessionRecord(record) => Ok(record.map(|record| *record)),
1919            reply => bail!("unexpected session record reply {reply:?}"),
1920        }
1921    }
1922
1923    pub async fn go_startup_session(
1924        &mut self,
1925        workspace_id: String,
1926        last_session_id: Option<String>,
1927    ) -> Result<Option<SessionRecord>> {
1928        match self
1929            .request(DaemonAction::GoStartupSession {
1930                workspace_id,
1931                last_session_id,
1932            })
1933            .await?
1934        {
1935            DaemonReply::GoStartupSession(record) => Ok(record.map(|record| *record)),
1936            reply => bail!("unexpected go startup session reply {reply:?}"),
1937        }
1938    }
1939
1940    pub async fn submit_session_command(
1941        &mut self,
1942        session_id: String,
1943        command_id: String,
1944        command: RelayCommand,
1945        inherited_draft: Option<String>,
1946    ) -> Result<u64> {
1947        match self
1948            .request(DaemonAction::SubmitSessionCommand {
1949                inherited_draft,
1950                session_id,
1951                command_id,
1952                command,
1953            })
1954            .await?
1955        {
1956            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1957            reply => bail!("unexpected session command reply {reply:?}"),
1958        }
1959    }
1960
1961    /// Hand the daemon a prompt for a session that is still starting. The
1962    /// daemon replies as soon as the prompt is queued, not when it is
1963    /// delivered; delivery failures come back as a session notice.
1964    pub async fn queue_startup_prompt(
1965        &mut self,
1966        session_id: String,
1967        text: String,
1968        inherited_draft: Option<String>,
1969    ) -> Result<()> {
1970        match self
1971            .request(DaemonAction::QueueStartupPrompt {
1972                session_id,
1973                text,
1974                inherited_draft,
1975            })
1976            .await?
1977        {
1978            DaemonReply::Done => Ok(()),
1979            reply => bail!("unexpected startup prompt reply {reply:?}"),
1980        }
1981    }
1982
1983    /// Take a prompt back from the daemon's startup queue before delivery.
1984    /// `Ok(true)` means it is withdrawn and will not be sent; `Ok(false)`
1985    /// means delivery had already started, so it was sent. An error means the
1986    /// daemon could not be asked, and the prompt may still be queued.
1987    pub async fn withdraw_startup_prompt(
1988        &mut self,
1989        session_id: String,
1990        text: String,
1991    ) -> Result<bool> {
1992        match self
1993            .request(DaemonAction::WithdrawStartupPrompt { session_id, text })
1994            .await?
1995        {
1996            DaemonReply::PromptWithdrawn(withdrawn) => Ok(withdrawn),
1997            reply => bail!("unexpected startup prompt withdrawal reply {reply:?}"),
1998        }
1999    }
2000
2001    /// Ask the daemon to review the turn this session just finished.
2002    ///
2003    /// The refusal is a sentence for a person -- "prompts are queued", "set
2004    /// [review] profile in config.toml" -- so it travels as text rather than
2005    /// as a code every surface would have to translate.
2006    pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
2007        match self
2008            .request(DaemonAction::StartTurnReview { session_id })
2009            .await?
2010        {
2011            DaemonReply::Done => Ok(()),
2012            reply => bail!("unexpected turn-review reply {reply:?}"),
2013        }
2014    }
2015
2016    pub async fn resolve_turn_review(
2017        &mut self,
2018        session_id: String,
2019        resolution: Resolution,
2020    ) -> Result<()> {
2021        match self
2022            .request(DaemonAction::ResolveTurnReview {
2023                session_id,
2024                resolution,
2025            })
2026            .await?
2027        {
2028            DaemonReply::Done => Ok(()),
2029            reply => bail!("unexpected turn-review resolution reply {reply:?}"),
2030        }
2031    }
2032
2033    pub async fn reviewer_action(
2034        &mut self,
2035        session_id: String,
2036        role: Option<String>,
2037        action: crate::session::ReviewerAction,
2038    ) -> Result<crate::session::ReviewerOutcome> {
2039        match self
2040            .request(DaemonAction::ReviewerAction {
2041                session_id,
2042                role,
2043                action,
2044            })
2045            .await?
2046        {
2047            DaemonReply::Reviewer(outcome) => Ok(*outcome),
2048            reply => bail!("unexpected reviewer reply {reply:?}"),
2049        }
2050    }
2051
2052    pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
2053        match self
2054            .request(DaemonAction::SyncSession { session_id })
2055            .await?
2056        {
2057            DaemonReply::Done => Ok(()),
2058            reply => bail!("unexpected session sync reply {reply:?}"),
2059        }
2060    }
2061
2062    pub async fn respond_elicitation(
2063        &mut self,
2064        session_id: String,
2065        elicitation_id: String,
2066        response: ElicitationResponse,
2067    ) -> Result<()> {
2068        match self
2069            .request(DaemonAction::RespondElicitation {
2070                session_id,
2071                elicitation_id,
2072                response,
2073            })
2074            .await?
2075        {
2076            DaemonReply::Done => Ok(()),
2077            reply => bail!("unexpected elicitation reply {reply:?}"),
2078        }
2079    }
2080
2081    pub async fn native_agent_history(
2082        &mut self,
2083        owner: String,
2084        child: String,
2085        before: Option<(u64, String)>,
2086    ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
2087        match self
2088            .request(DaemonAction::NativeAgentHistory {
2089                owner,
2090                child,
2091                before,
2092            })
2093            .await?
2094        {
2095            DaemonReply::NativeAgentHistory(page) => Ok(page),
2096            reply => bail!("unexpected native agent history reply {reply:?}"),
2097        }
2098    }
2099
2100    pub async fn stop_background_task(
2101        &mut self,
2102        session_id: String,
2103        background_task_id: String,
2104    ) -> Result<()> {
2105        match self
2106            .request(DaemonAction::StopBackgroundTask {
2107                session_id,
2108                background_task_id,
2109            })
2110            .await?
2111        {
2112            DaemonReply::Done => Ok(()),
2113            reply => bail!("unexpected background task stop reply {reply:?}"),
2114        }
2115    }
2116
2117    pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
2118        self.suspend_session_with_ack(session_id, false).await
2119    }
2120
2121    pub async fn suspend_session_with_ack(
2122        &mut self,
2123        session_id: String,
2124        acknowledge_unpublished_work: bool,
2125    ) -> Result<()> {
2126        match self
2127            .request(DaemonAction::SuspendSession {
2128                session_id,
2129                acknowledge_unpublished_work,
2130            })
2131            .await?
2132        {
2133            DaemonReply::Done => Ok(()),
2134            reply => bail!("unexpected close-session reply {reply:?}"),
2135        }
2136    }
2137
2138    pub async fn start_create_session(
2139        &mut self,
2140        request: CreateSessionRequest,
2141    ) -> Result<RegisteredSession> {
2142        match self
2143            .request(DaemonAction::StartCreateSession(request))
2144            .await?
2145        {
2146            DaemonReply::RegisteredSession(registered) => Ok(*registered),
2147            reply => bail!("unexpected start-create reply {reply:?}"),
2148        }
2149    }
2150
2151    pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
2152        match self
2153            .request(DaemonAction::WaitCreateSession { session_id })
2154            .await?
2155        {
2156            DaemonReply::Done => Ok(()),
2157            reply => bail!("unexpected wait-create reply {reply:?}"),
2158        }
2159    }
2160
2161    pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
2162        match self.request(DaemonAction::ResumeSession(request)).await? {
2163            DaemonReply::Done => Ok(()),
2164            reply => bail!("unexpected resume-session reply {reply:?}"),
2165        }
2166    }
2167
2168    pub async fn discard_since_checkpoint(
2169        &mut self,
2170        session_id: String,
2171        checkpoint: mj_core::state::CheckpointMetadata,
2172    ) -> Result<()> {
2173        match self
2174            .request(DaemonAction::DiscardSinceCheckpoint {
2175                session_id,
2176                checkpoint,
2177            })
2178            .await?
2179        {
2180            DaemonReply::Done => Ok(()),
2181            reply => bail!("unexpected force-stop reply {reply:?}"),
2182        }
2183    }
2184
2185    pub async fn destroy_stopped_session(
2186        &mut self,
2187        session_id: String,
2188        delete_branch: bool,
2189    ) -> Result<()> {
2190        match self
2191            .request(DaemonAction::DestroyStoppedSession {
2192                session_id,
2193                delete_branch,
2194            })
2195            .await?
2196        {
2197            DaemonReply::Done => Ok(()),
2198            reply => bail!("unexpected destroy-stopped reply {reply:?}"),
2199        }
2200    }
2201
2202    pub async fn force_destroy_session(
2203        &mut self,
2204        session_id: String,
2205        delete_branch: bool,
2206    ) -> Result<()> {
2207        match self
2208            .request(DaemonAction::ForceDestroySession {
2209                session_id,
2210                delete_branch,
2211            })
2212            .await?
2213        {
2214            DaemonReply::Done => Ok(()),
2215            reply => bail!("unexpected force-destroy reply {reply:?}"),
2216        }
2217    }
2218
2219    pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
2220        match self
2221            .request(DaemonAction::CancelLifecycle { session_id })
2222            .await?
2223        {
2224            DaemonReply::Done => Ok(()),
2225            reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
2226        }
2227    }
2228
2229    pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
2230        match self
2231            .request(DaemonAction::RecoverDraft { draft_id })
2232            .await?
2233        {
2234            DaemonReply::Done => Ok(()),
2235            reply => bail!("unexpected recover-draft reply {reply:?}"),
2236        }
2237    }
2238
2239    pub async fn stop(&mut self) -> Result<()> {
2240        match self.request(DaemonAction::Stop).await? {
2241            DaemonReply::Done => Ok(()),
2242            reply => bail!("unexpected stop reply {reply:?}"),
2243        }
2244    }
2245}
2246
2247pub async fn connect_existing() -> Result<DaemonClient> {
2248    let metadata = tokio::task::spawn_blocking(read_metadata)
2249        .await
2250        .context("read daemon metadata task failed")??;
2251    DaemonClient::connect(metadata).await
2252}
2253
2254/// A handle to whatever daemon the metadata file advertises, regardless of its
2255/// protocol version. It only exposes the frozen management subset (`Ping`,
2256/// `Status`, `Stop`), which encodes identically in every protocol version.
2257pub struct ManagementClient {
2258    inner: DaemonClient,
2259}
2260
2261impl ManagementClient {
2262    pub fn new(inner: DaemonClient) -> Self {
2263        Self { inner }
2264    }
2265    pub fn protocol_version(&self) -> u32 {
2266        self.inner.metadata.protocol_version
2267    }
2268
2269    pub async fn status(&mut self) -> Result<DaemonStatus> {
2270        self.inner.status().await
2271    }
2272
2273    pub async fn stop(&mut self) -> Result<()> {
2274        tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
2275            .await
2276            .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
2277    }
2278
2279    /// Ask the daemon to stop and wait for its process to actually exit.
2280    pub async fn stop_and_wait(mut self) -> Result<()> {
2281        let pid = self.inner.metadata.pid;
2282        self.stop().await?;
2283        wait_for_exit_within(pid, STOP_DRAIN_TIMEOUT)
2284            .await
2285            .with_context(|| {
2286                format!(
2287                    "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
2288                    STOP_DRAIN_TIMEOUT.as_secs()
2289                )
2290            })
2291    }
2292}
2293
2294pub async fn connect_management() -> Result<ManagementClient> {
2295    Ok(ManagementClient {
2296        inner: DaemonClient::connect(read_metadata_any()?).await?,
2297    })
2298}
2299
2300/// Refuse to speak to a daemon whose protocol this build does not know.
2301///
2302/// The two protocol numbers alone do not say which `mj` ran. The usual cause is
2303/// a second installation: a `cargo install`ed client sits earlier on PATH than
2304/// the build whose daemon is running, so every command fails here while the
2305/// other binary works, and nothing in the message says where either lives. It
2306/// therefore names both executables and both versions.
2307pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
2308    ensure!(
2309        metadata.protocol_version <= PROTOCOL_VERSION,
2310        "{}",
2311        unsupported_daemon_protocol_message(
2312            metadata.protocol_version,
2313            &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
2314        )
2315    );
2316    Ok(())
2317}
2318
2319/// The message [`ensure_supported_daemon_protocol`] fails with, given the
2320/// sentence that names both builds, so it can be read without a daemon.
2321fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
2322    format!(
2323        "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
2324         Put the daemon's directory first on PATH, or reinstall this client from that build."
2325    )
2326}
2327// The resume dialog lists preview rows and fetches the record it resumes.
2328pub const PROTOCOL_VERSION: u32 = 51;
2329pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
2330/// How long a daemon is given to exit after it accepts a stop.
2331///
2332/// Stopping cancels a token and returns immediately; the daemon then unwinds
2333/// its session manager, its phone server and its pollers. That is normally
2334/// fast, but a daemon whose database has been migrated out from under it fails
2335/// every read while it winds down and has been observed taking over five
2336/// seconds — which the previous five-second bound missed by a fraction,
2337/// reporting a stop that had in fact worked as `did not stop` and aborting the
2338/// restart that depended on it.
2339pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
2340/// How long a requested stop may take: the daemon waits up to 60 seconds for
2341/// session destroys in flight (#1191) before it winds down.
2342pub const STOP_DRAIN_TIMEOUT: Duration = Duration::from_secs(90);
2343pub const RETRY_DELAY: Duration = Duration::from_millis(40);
2344impl DaemonClient {
2345    pub async fn prepare_move_session(
2346        &mut self,
2347        selection: MoveSelection,
2348    ) -> Result<MovePreparation> {
2349        match self
2350            .request(DaemonAction::PrepareMoveSession(selection))
2351            .await?
2352        {
2353            DaemonReply::MovePreparation(preparation) => Ok(*preparation),
2354            _ => bail!("daemon returned an unexpected move preparation reply"),
2355        }
2356    }
2357
2358    pub async fn move_sources(
2359        &mut self,
2360        session_id: String,
2361        cleanup_operation_id: Option<String>,
2362    ) -> Result<Vec<mj_core::move_workspace::RetainedMoveSource>> {
2363        match self
2364            .request(DaemonAction::MoveSources {
2365                session_id,
2366                cleanup_operation_id,
2367            })
2368            .await?
2369        {
2370            DaemonReply::MoveSources(sources) => Ok(sources),
2371            reply => bail!("unexpected Move sources reply: {reply:?}"),
2372        }
2373    }
2374
2375    pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
2376        match self.request(DaemonAction::MoveSession(request)).await? {
2377            DaemonReply::MoveOutcome(outcome) => Ok(outcome),
2378            _ => bail!("daemon returned an unexpected move reply"),
2379        }
2380    }
2381}
2382
2383#[cfg(test)]
2384mod tests {
2385    use super::*;
2386    use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
2387    use std::path::Path;
2388
2389    #[cfg(unix)]
2390    #[test]
2391    fn an_exited_unreaped_process_is_a_zombie_and_not_alive() {
2392        // An upgrade waits for the old daemon to leave. If that daemon exited
2393        // but is still a child nobody has reaped, `kill(pid, 0)` keeps finding
2394        // it, so only the zombie check lets the wait end.
2395        use std::time::{Duration, Instant};
2396
2397        let mut running = std::process::Command::new("sleep")
2398            .arg("30")
2399            .spawn()
2400            .unwrap();
2401        assert!(process_is_alive(running.id()));
2402        assert!(!process_is_zombie(running.id()));
2403        running.kill().unwrap();
2404        running.wait().unwrap();
2405
2406        let mut exited = std::process::Command::new("sh")
2407            .args(["-c", "exit 0"])
2408            .spawn()
2409            .unwrap();
2410        let pid = exited.id();
2411        let deadline = Instant::now() + Duration::from_secs(10);
2412        while !process_is_zombie(pid) {
2413            assert!(
2414                Instant::now() < deadline,
2415                "an exited, unreaped child was never reported as a zombie"
2416            );
2417            std::thread::sleep(Duration::from_millis(20));
2418        }
2419        assert!(!process_is_alive(pid));
2420        exited.wait().unwrap();
2421    }
2422
2423    #[test]
2424    fn a_missing_endpoint_file_reads_as_a_stopped_daemon() {
2425        let directory = tempfile::tempdir().unwrap();
2426        let path = directory.path().join("daemon.json");
2427        let error = read_metadata_at(&path).unwrap_err();
2428        let stopped = daemon_not_running(&error).expect("stopped, not an I/O failure");
2429        assert_eq!(stopped.metadata_path, path);
2430        assert_eq!(format!("{error:#}"), "the Mjolnir daemon is not running");
2431
2432        std::fs::write(&path, b"not json").unwrap();
2433        let error = read_metadata_at(&path).unwrap_err();
2434        assert!(daemon_not_running(&error).is_none());
2435    }
2436
2437    #[tokio::test]
2438    async fn subagent_discovery_uses_the_daemon_and_preserves_choices_and_failures() {
2439        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2440        let metadata = DaemonMetadata {
2441            protocol_version: PROTOCOL_VERSION,
2442            pid: std::process::id(),
2443            address: listener.local_addr().unwrap(),
2444            token: "subagent-discovery-test".into(),
2445            started_at: "test".into(),
2446            build_version: env!("CARGO_PKG_VERSION").into(),
2447        };
2448        let options: mj_core::subagent::SubagentOptions =
2449            serde_json::from_value(serde_json::json!({
2450                "models": [{"value": "gpt-6-luna", "name": "Luna", "description": null}],
2451                "efforts": [{"value": "high", "name": "High", "description": null}],
2452                "unavailable": ["another-profile: authentication failed"]
2453            }))
2454            .unwrap();
2455        let expected = serde_json::to_value(&options).unwrap();
2456        let mut draft = Config::default();
2457        draft.profiles.insert(
2458            "codex".into(),
2459            serde_json::from_value(serde_json::json!({
2460                "kind": "codex", "home": "/unsaved/profile",
2461                "subagents": {"mode": "single_model", "model": "gpt-6-luna", "effort": "high"}
2462            }))
2463            .unwrap(),
2464        );
2465        let expected_draft = serde_json::to_value(&draft).unwrap();
2466        let server = tokio::spawn(async move {
2467            let (mut stream, _) = listener.accept().await.unwrap();
2468            for (index, model) in [
2469                None,
2470                Some("gpt-6-luna".to_owned()),
2471                Some("gpt-6-luna".to_owned()),
2472            ]
2473            .into_iter()
2474            .enumerate()
2475            {
2476                let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2477                assert_eq!(request.token, "subagent-discovery-test");
2478                let DaemonAction::SubagentOptions {
2479                    profile,
2480                    model: requested,
2481                    config,
2482                } = request.action
2483                else {
2484                    panic!("unexpected discovery request")
2485                };
2486                assert_eq!(profile, "codex");
2487                assert_eq!(requested, model);
2488                assert_eq!(
2489                    serde_json::to_value(config).unwrap(),
2490                    if index == 1 {
2491                        expected_draft.clone()
2492                    } else {
2493                        serde_json::Value::Null
2494                    }
2495                );
2496                write_frame(
2497                    &mut stream,
2498                    &ResponseEnvelope {
2499                        protocol_version: request.protocol_version,
2500                        request_id: request.request_id,
2501                        result: if index != 2 {
2502                            Ok(DaemonReply::SubagentOptions(options.clone()))
2503                        } else {
2504                            Err("profile discovery cancelled".into())
2505                        },
2506                    },
2507                )
2508                .await
2509                .unwrap();
2510            }
2511        });
2512        let mut client = DaemonClient::connect(metadata).await.unwrap();
2513        let discovered = client
2514            .subagent_options("codex".into(), None, None)
2515            .await
2516            .unwrap();
2517        assert_eq!(serde_json::to_value(discovered).unwrap(), expected);
2518        let discovered = client
2519            .subagent_options("codex".into(), Some("gpt-6-luna".into()), Some(draft))
2520            .await
2521            .unwrap();
2522        assert_eq!(serde_json::to_value(discovered).unwrap(), expected);
2523        let error = client
2524            .subagent_options("codex".into(), Some("gpt-6-luna".into()), None)
2525            .await
2526            .unwrap_err();
2527        assert_eq!(error.to_string(), "profile discovery cancelled");
2528        server.await.unwrap();
2529    }
2530
2531    #[tokio::test]
2532    async fn upgrade_refusal_retries_the_identical_command_but_lost_acknowledgements_do_not() {
2533        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2534        let metadata = DaemonMetadata {
2535            protocol_version: PROTOCOL_VERSION,
2536            pid: std::process::id(),
2537            address: listener.local_addr().unwrap(),
2538            token: "test-token".into(),
2539            started_at: "test".into(),
2540            build_version: env!("CARGO_PKG_VERSION").into(),
2541        };
2542        let action = DaemonAction::SubmitSessionCommand {
2543            inherited_draft: None,
2544            session_id: "test-session".into(),
2545            command_id: "steer-command".into(),
2546            command: RelayCommand::Steer {
2547                active_prompt_id: "active-command".into(),
2548                queued_prompt_id: "queued-command".into(),
2549            },
2550        };
2551        let expected = serde_json::to_value(&action).unwrap();
2552        let server = tokio::spawn(async move {
2553            for reply in [
2554                Some(DaemonReply::UpgradePending),
2555                Some(DaemonReply::Done),
2556                None,
2557            ] {
2558                let (mut stream, _) = listener.accept().await.unwrap();
2559                let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2560                assert_eq!(serde_json::to_value(&request.action).unwrap(), expected);
2561                if let Some(reply) = reply {
2562                    write_frame(
2563                        &mut stream,
2564                        &ResponseEnvelope {
2565                            protocol_version: request.protocol_version,
2566                            request_id: request.request_id,
2567                            result: Ok(reply),
2568                        },
2569                    )
2570                    .await
2571                    .unwrap();
2572                }
2573            }
2574        });
2575        let mut client = DaemonClient::connect(metadata.clone()).await.unwrap();
2576        let reply = client
2577            .request_with_reconnect(action.clone(), || DaemonClient::connect(metadata.clone()))
2578            .await
2579            .unwrap();
2580        assert!(matches!(reply, DaemonReply::Done));
2581        let mut client = DaemonClient::connect(metadata).await.unwrap();
2582        assert!(
2583            client
2584                .request_with_reconnect(action, || async {
2585                    panic!("an ambiguous acknowledgement must not replay a mutation");
2586                })
2587                .await
2588                .is_err()
2589        );
2590        server.await.unwrap();
2591    }
2592
2593    #[test]
2594    fn unsupported_protocol_message_names_both_binaries_and_versions() {
2595        let message = unsupported_daemon_protocol_message(
2596            PROTOCOL_VERSION + 5,
2597            &describe_daemon_and_client_builds(
2598                4242,
2599                BuildDescription {
2600                    executable: Some(Path::new("/home/dev/mj/target/release/mj")),
2601                    version: "2.14.0",
2602                },
2603                BuildDescription {
2604                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2605                    version: "2.9.0",
2606                },
2607            ),
2608        );
2609        assert_eq!(
2610            message,
2611            format!(
2612                "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
2613                 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
2614                 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
2615                 Put the daemon's directory first on PATH, or reinstall this client from that build.",
2616                PROTOCOL_VERSION + 5
2617            )
2618        );
2619    }
2620
2621    #[test]
2622    fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
2623        let message = unsupported_daemon_protocol_message(
2624            PROTOCOL_VERSION + 1,
2625            &describe_daemon_and_client_builds(
2626                4242,
2627                BuildDescription {
2628                    executable: None,
2629                    version: "2.14.0",
2630                },
2631                BuildDescription {
2632                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2633                    version: "2.9.0",
2634                },
2635            ),
2636        );
2637        assert!(
2638            message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
2639            "{message}"
2640        );
2641        assert!(
2642            message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
2643            "{message}"
2644        );
2645    }
2646}