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    /// One live session's transcript tail as it was at `cursor`, so the
572    /// deltas after that cursor apply to it exactly.
573    SessionTail {
574        session_id: String,
575        cursor: crate::runtime_feed::RuntimeCursor,
576    },
577    /// The sessions the resume dialog lists. The runtime feed carries only
578    /// live sessions, so the dialog asks for these when it opens.
579    ResumeCandidates,
580    /// One session's whole record, live or stopped. The resume wizard reads
581    /// it for the row picked in the resume dialog.
582    SessionRecord {
583        session_id: String,
584    },
585    /// The session `mj go` opens: the remembered one while it is still
586    /// eligible, otherwise the most recently updated eligible session of the
587    /// workspace, live or stopped.
588    GoStartupSession {
589        workspace_id: String,
590        last_session_id: Option<String>,
591    },
592    /// Ask the daemon's quota poller to probe now. The daemon is the only
593    /// process that probes; the result arrives in the runtime feed. Added in
594    /// protocol 43.
595    RefreshQuota,
596    /// Discover eligible subagent models in the daemon, which owns the cache.
597    SubagentOptions {
598        profile: String,
599        model: Option<String>,
600        /// Setup can probe an unsaved draft without changing live settings.
601        #[serde(default, skip_serializing_if = "Option::is_none")]
602        config: Option<Box<Config>>,
603    },
604    /// Ensure valid unsaved profile definitions are hydrating. Results arrive
605    /// through runtime publications; this request does no discovery itself.
606    WarmProfileCapabilities {
607        config: Box<Config>,
608    },
609    RenameProfile {
610        old_id: String,
611        new_id: String,
612    },
613    RenameTarget {
614        old_id: String,
615        new_id: String,
616    },
617    SubmitSessionCommand {
618        #[serde(default)]
619        inherited_draft: Option<String>,
620        session_id: String,
621        command_id: String,
622        command: RelayCommand,
623    },
624    /// Deliver a prompt typed while a session was still starting, once the
625    /// daemon sees that session's harness become ready. The daemon owns the
626    /// wait, so the prompt arrives whether or not this client is still
627    /// running or still showing that session.
628    QueueStartupPrompt {
629        session_id: String,
630        text: String,
631        /// The saved draft text this prompt was typed from, if the client
632        /// also persisted it. Cleared after a successful submit so the
633        /// delivered prompt does not reappear as a draft.
634        #[serde(default)]
635        inherited_draft: Option<String>,
636    },
637    /// Take back a prompt queued with `QueueStartupPrompt` that the daemon has
638    /// not started delivering, so the person can edit it. The daemon cancels
639    /// the newest queued prompt with this exact text and replies
640    /// `PromptWithdrawn(true)`, or `PromptWithdrawn(false)` when none is
641    /// waiting because delivery already started.
642    WithdrawStartupPrompt {
643        session_id: String,
644        text: String,
645    },
646    SyncSession {
647        session_id: String,
648    },
649    RespondElicitation {
650        session_id: String,
651        elicitation_id: String,
652        response: ElicitationResponse,
653    },
654    StopBackgroundTask {
655        session_id: String,
656        background_task_id: String,
657    },
658    /// Drive a session's second-opinion reviewer. The reviewer is a sidecar of
659    /// the session's worker, so it travels the session's own relay rather than
660    /// becoming a session of its own here.
661    ReviewerAction {
662        session_id: String,
663        /// Which reviewing role the action drives; absent means the default
664        /// one, which is what plan review uses.
665        #[serde(default, skip_serializing_if = "Option::is_none")]
666        role: Option<String>,
667        action: crate::session::ReviewerAction,
668    },
669    /// Review the turn this session just finished, on a surface's request.
670    StartTurnReview {
671        session_id: String,
672    },
673    /// Forward, dismiss, or cancel the open review.
674    ResolveTurnReview {
675        session_id: String,
676        resolution: Resolution,
677    },
678    SuspendSession {
679        session_id: String,
680        #[serde(default)]
681        acknowledge_unpublished_work: bool,
682    },
683    StartCreateSession(CreateSessionRequest),
684    WaitCreateSession {
685        session_id: String,
686    },
687    ResumeSession(ResumeSessionRequest),
688    PrepareMoveSession(MoveSelection),
689    MoveSession(MoveSessionRequest),
690    MoveSources {
691        session_id: String,
692        cleanup_operation_id: Option<String>,
693    },
694    DiscardSinceCheckpoint {
695        session_id: String,
696        checkpoint: mj_core::state::CheckpointMetadata,
697    },
698    DestroyStoppedSession {
699        session_id: String,
700        /// Whether to delete the session's managed git branch as well. The
701        /// branch can hold work the user still wants, so destroying keeps it
702        /// unless the request asks for the deletion.
703        delete_branch: bool,
704    },
705    ForceDestroySession {
706        session_id: String,
707        /// See [`DaemonAction::DestroyStoppedSession`].
708        delete_branch: bool,
709    },
710    CancelLifecycle {
711        session_id: String,
712    },
713    RecoverDraft {
714        draft_id: String,
715    },
716    Stop,
717}
718
719#[derive(Debug, Serialize, Deserialize)]
720#[serde(deny_unknown_fields)]
721pub struct RequestEnvelope {
722    pub protocol_version: u32,
723    pub request_id: u64,
724    pub token: String,
725    pub action: DaemonAction,
726}
727
728/// The daemon answered a request with an error. Unlike a lost connection,
729/// this is a complete round trip: the daemon read the request and said why it
730/// did not carry it out.
731#[derive(Debug)]
732pub struct DaemonRefusal(pub String);
733
734impl std::fmt::Display for DaemonRefusal {
735    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
736        f.write_str(&self.0)
737    }
738}
739
740impl std::error::Error for DaemonRefusal {}
741
742impl DaemonRefusal {
743    /// Whether the daemon itself could not confirm delivery. Its message is
744    /// then led by `session::DeliveryUnconfirmed`.
745    #[must_use]
746    pub fn delivery_unconfirmed(&self) -> bool {
747        self.0
748            .starts_with(&crate::session::DeliveryUnconfirmed.to_string())
749    }
750}
751
752#[derive(Debug, Serialize, Deserialize)]
753#[serde(deny_unknown_fields)]
754pub struct ResponseEnvelope {
755    pub protocol_version: u32,
756    pub request_id: u64,
757    pub result: std::result::Result<DaemonReply, String>,
758}
759
760#[derive(Debug, Clone, Serialize, Deserialize)]
761#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
762pub enum DaemonReply {
763    NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
764    Pong,
765    UpgradePending,
766    /// Labels for the work holding an automatic handoff open, empty when
767    /// nothing is. Answers `DaemonAction::UpgradeBlockers`.
768    UpgradeBlockers(Vec<String>),
769    Status(DaemonStatus),
770    WebViewerAccess(crate::web::WebViewerAccess),
771    WebListeners(Vec<crate::web::WebListenerProcess>),
772    SubagentOptions(mj_core::subagent::SubagentOptions),
773    Workspaces(Vec<WorkspaceListing>),
774    Workspace(WorkspaceRecord),
775    Snapshot(WorkspaceSnapshot),
776    RuntimeChanges(Box<crate::runtime_feed::RuntimeFrame>),
777    SessionTail(Box<crate::runtime_feed::SessionTailReply>),
778    ResumeCandidates(Box<ResumeCandidates>),
779    GoStartupSession(Option<Box<SessionRecord>>),
780    SessionRecord(Option<Box<SessionRecord>>),
781    /// Transport fragments of one chunked reply (see
782    /// [`DaemonReply::is_chunked`]); never exposed to consumers.
783    ReplyChunk {
784        bytes: Vec<u8>,
785        finished: bool,
786    },
787    RegisteredSession(Box<RegisteredSession>),
788    MovePreparation(Box<MovePreparation>),
789    MoveOutcome(MoveOutcome),
790    MoveSources(Vec<mj_core::move_workspace::RetainedMoveSource>),
791    Ordinal(u64),
792    Text(String),
793    OptionalSessionState(Option<SessionState>),
794    Checkpoint(mj_core::state::CheckpointMetadata),
795    RecoveryScan(mj_core::state::RecoveryScan),
796    WikiRows(WikiSearchPage),
797    ProjectCatalog(mj_core::project_catalog::ProjectCatalogView),
798    ProjectCreated(mj_core::project_catalog::SavedProject),
799    SessionTextMatches(Vec<SessionTextMatch>),
800    WikiHits(Option<WikiHitTranscript>),
801    WikiSession(Option<Box<WikiSessionInfo>>),
802    Reviewer(Box<crate::session::ReviewerOutcome>),
803    /// Whether a queued startup prompt was withdrawn before delivery.
804    PromptWithdrawn(bool),
805    Done,
806}
807
808impl DaemonReply {
809    /// Replies that can outgrow one frame. The daemon sends them as
810    /// [`DaemonReply::ReplyChunk`] fragments and the client reassembles them.
811    pub fn is_chunked(&self) -> bool {
812        matches!(
813            self,
814            Self::RuntimeChanges(_) | Self::SessionTail(_) | Self::ResumeCandidates(_)
815        )
816    }
817}
818
819/// What the resume dialog lists, read on demand.
820#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
821#[serde(deny_unknown_fields)]
822pub struct ResumeCandidates {
823    /// One row for every inactive session that is not a sub-agent.
824    pub candidates: Vec<ResumeCandidate>,
825    /// Durable moves of those sessions, for "move needs recovery" marks.
826    pub moves: Vec<MoveOperation>,
827    /// Native sessions that some record, live or not, already holds. The
828    /// import list must not offer them again.
829    pub adopted_native_sessions: Vec<(mj_core::config::HarnessKind, String)>,
830    /// Local checkouts that Mjolnir created for its sessions. Every native
831    /// thread that ran inside one belongs to that session.
832    pub local_checkout_roots: Vec<PathBuf>,
833}
834
835/// One row of the resume dialog's Mjolnir tab: what the row shows, without
836/// the rest of the record. A session's title can be tens of kilobytes, and
837/// most of a store's history is never resumed, so the resume wizard fetches
838/// the whole record (`DaemonAction::SessionRecord`) only for the row picked.
839#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
840#[serde(deny_unknown_fields)]
841pub struct ResumeCandidate {
842    pub session_id: String,
843    pub state: SessionState,
844    pub last_profile: String,
845    /// The listed title, cut to [`ResumeCandidate::TITLE_CHARS`].
846    pub title: String,
847    /// Where the session ran, as the session list names it.
848    pub origin: String,
849    pub project: String,
850    pub has_checkpoint: bool,
851    /// The checkpoint's creation time, or the record's last update without
852    /// one; `None` when neither parses.
853    pub last_activity_ms: Option<i64>,
854    pub publication: Option<PublicationState>,
855    /// The checkpoint archive whose size the row shows. Only a stopped
856    /// record is described by its checkpoint.
857    pub checkpoint_archive: Option<PathBuf>,
858    /// Whether destroying the session can also delete its worktree branch.
859    pub worktree_checkout: bool,
860}
861
862impl ResumeCandidate {
863    /// Longer than any row draws it.
864    pub const TITLE_CHARS: usize = 200;
865
866    /// The row for `record`. The daemon and the dashboard both build rows
867    /// here, so a row reads the same whichever of them held the record.
868    pub fn of(record: &SessionRecord, config: &Config) -> Self {
869        let timestamp_ms = |timestamp: &str| {
870            chrono::DateTime::parse_from_rfc3339(timestamp)
871                .ok()
872                .map(|parsed| parsed.timestamp_millis())
873        };
874        Self {
875            session_id: record.id.clone(),
876            state: record.state,
877            last_profile: record.last_profile.clone(),
878            title: record
879                .listed_title()
880                .chars()
881                .take(Self::TITLE_CHARS)
882                .collect(),
883            origin: record.project_target(config, &record.target_template_id),
884            project: record.project_name(config),
885            has_checkpoint: record.checkpoint.is_some(),
886            last_activity_ms: record
887                .checkpoint
888                .as_ref()
889                .and_then(|checkpoint| timestamp_ms(&checkpoint.created_at))
890                .or_else(|| timestamp_ms(&record.updated_at)),
891            publication: record.publication_state(),
892            checkpoint_archive: record
893                .checkpoint
894                .as_ref()
895                .filter(|_| record.state == SessionState::Stopped)
896                .map(|checkpoint| checkpoint.archive_path.clone()),
897            worktree_checkout: record
898                .managed_worktree
899                .as_ref()
900                .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Worktree),
901        }
902    }
903}
904
905#[derive(Debug, Clone, Serialize, Deserialize)]
906#[serde(deny_unknown_fields)]
907pub struct DaemonStatus {
908    pub pid: u32,
909    pub started_at: String,
910    pub build_version: String,
911    pub attached_clients: usize,
912    pub phone_status: WebViewerStatus,
913}
914
915#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
916#[serde(rename_all = "snake_case", tag = "state")]
917pub enum WebViewerStatus {
918    Disabled,
919    Starting,
920    Ready {
921        viewer_url: String,
922        viewer_code: String,
923        qr_login_url: Option<String>,
924        fallback_reason: Option<String>,
925    },
926    Stopped,
927    Error {
928        message: String,
929    },
930}
931
932impl std::fmt::Debug for WebViewerStatus {
933    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
934        match self {
935            Self::Ready {
936                viewer_url,
937                viewer_code,
938                fallback_reason,
939                ..
940            } => formatter
941                .debug_struct("Ready")
942                .field("viewer_url", viewer_url)
943                .field("viewer_code", viewer_code)
944                .field("qr_login_url", &"[redacted]")
945                .field("fallback_reason", fallback_reason)
946                .finish(),
947            Self::Disabled => formatter.write_str("Disabled"),
948            Self::Starting => formatter.write_str("Starting"),
949            Self::Stopped => formatter.write_str("Stopped"),
950            Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
951        }
952    }
953}
954
955impl std::fmt::Display for WebViewerStatus {
956    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
957        match self {
958            Self::Disabled => formatter.write_str("disabled"),
959            Self::Starting => formatter.write_str("starting"),
960            Self::Stopped => formatter.write_str("stopped unexpectedly"),
961            Self::Error { message } => write!(formatter, "error: {message}"),
962            Self::Ready {
963                viewer_url,
964                viewer_code,
965                fallback_reason,
966                ..
967            } => {
968                write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
969                if let Some(reason) = fallback_reason {
970                    write!(
971                        formatter,
972                        "; local only because Tailscale HTTPS is unavailable: {reason}"
973                    )?;
974                }
975                Ok(())
976            }
977        }
978    }
979}
980
981/// Whether a non-child process has exited but not yet been reaped.
982///
983/// A zombie still answers `kill(pid, 0)`, because its process-table entry
984/// survives until its parent waits for it — so an existence probe alone calls
985/// it alive forever and anything waiting for it to leave waits forever. That is
986/// exactly the shape of `mj daemon restart` refusing to restart a daemon that
987/// had already stopped: `spawn_detached` used to leave the daemon a child of a
988/// long-lived Mjolnir process that never reaped it. It now double-forks, so the
989/// daemon is init's to reap, but any other unreaped child of this process would
990/// look the same, and the check stays cheap.
991///
992/// Treating a zombie as gone is also safe in the direction that matters: a
993/// zombie's PID cannot be reused until it is reaped, so nothing else can be
994/// occupying that number while this returns true.
995///
996/// On Linux this reads the one `/proc/<pid>/stat` file. `sysinfo` scans every
997/// process on the machine for a single-PID refresh, so it is not used here:
998/// this probe runs every 250 ms while an upgrade waits on an older daemon.
999#[cfg(target_os = "linux")]
1000pub fn process_is_zombie(pid: u32) -> bool {
1001    let Ok(stat) = std::fs::read(format!("/proc/{pid}/stat")) else {
1002        return false;
1003    };
1004    // The command name in parentheses may itself contain ") ", so the state
1005    // is the first field after the last closing parenthesis.
1006    stat.iter()
1007        .rposition(|byte| *byte == b')')
1008        .and_then(|end| stat.get(end + 2))
1009        .is_some_and(|state| *state == b'Z')
1010}
1011
1012/// On macOS this asks the kernel for the one process's BSD info. `sysinfo`
1013/// cannot be used: its lookups skip zombies, so it reports an exited, unreaped
1014/// process as absent rather than as a zombie, and `kill(pid, 0)` then calls it
1015/// alive forever. A nonzero `arg` to `PROC_PIDTBSDINFO` makes the kernel
1016/// include zombies.
1017#[cfg(target_os = "macos")]
1018pub fn process_is_zombie(pid: u32) -> bool {
1019    let Ok(raw_pid) = libc::c_int::try_from(pid) else {
1020        return false;
1021    };
1022    let size = std::mem::size_of::<libc::proc_bsdinfo>();
1023    // SAFETY: proc_bsdinfo is plain old data, so all-zero is a valid value.
1024    let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
1025    // SAFETY: the buffer is a writable proc_bsdinfo of exactly `size` bytes.
1026    let written = unsafe {
1027        libc::proc_pidinfo(
1028            raw_pid,
1029            libc::PROC_PIDTBSDINFO,
1030            1,
1031            (&mut info as *mut libc::proc_bsdinfo).cast(),
1032            size as libc::c_int,
1033        )
1034    };
1035    usize::try_from(written).is_ok_and(|written| written == size) && info.pbi_status == libc::SZOMB
1036}
1037
1038#[cfg(all(unix, not(target_os = "linux"), not(target_os = "macos")))]
1039pub fn process_is_zombie(pid: u32) -> bool {
1040    let pid = sysinfo::Pid::from_u32(pid);
1041    let mut system = sysinfo::System::new();
1042    system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
1043    system
1044        .process(pid)
1045        .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
1046}
1047
1048/// Wait for a process to leave, within [`STOP_TIMEOUT`].
1049///
1050/// The error says the process was still running rather than that it "did not
1051/// stop": a daemon that is still winding down has not refused, and the two
1052/// read very differently to somebody deciding whether to reach for a kill.
1053pub async fn wait_for_exit(pid: u32) -> Result<()> {
1054    wait_for_exit_within(pid, STOP_TIMEOUT).await
1055}
1056
1057async fn wait_for_exit_within(pid: u32, timeout: Duration) -> Result<()> {
1058    let deadline = Instant::now() + timeout;
1059    while daemon_process_is_alive(pid) {
1060        ensure!(Instant::now() < deadline, "process {pid} is still running");
1061        tokio::time::sleep(RETRY_DELAY).await;
1062    }
1063    Ok(())
1064}
1065
1066/// Whether a daemon that Mjolnir launched in its own process group is alive.
1067///
1068/// Reaping is deliberately confined to this daemon-specific path. Attachment
1069/// PIDs are merely observations and may alias unrelated children owned by this
1070/// process, so their liveness probe below must never call `waitpid`.
1071pub fn daemon_process_is_alive(pid: u32) -> bool {
1072    #[cfg(unix)]
1073    {
1074        if pid == 0 {
1075            return false;
1076        }
1077        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
1078            return false;
1079        };
1080        let mut status = 0;
1081        // SAFETY: `status` is writable for the call and WNOHANG never blocks.
1082        // A non-child fails with ECHILD without changing any process state.
1083        let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
1084        if waited == raw_pid {
1085            return false;
1086        }
1087        if waited == 0 {
1088            return true;
1089        }
1090        let wait_error = std::io::Error::last_os_error();
1091        if wait_error.raw_os_error() != Some(libc::ECHILD) {
1092            return true;
1093        }
1094
1095        #[cfg(target_os = "macos")]
1096        return owned_daemon_group_is_alive(raw_pid);
1097
1098        #[cfg(not(target_os = "macos"))]
1099        process_is_alive(pid)
1100    }
1101    #[cfg(not(unix))]
1102    process_is_alive(pid)
1103}
1104
1105#[cfg(target_os = "macos")]
1106pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
1107    // `spawn_detached` makes the daemon a process-group leader. Darwin
1108    // excludes zombies from group signal probes: ESRCH means the group is gone
1109    // and EPERM means only exiting members remain. The latter is safe here
1110    // because this is a group we created for our own same-user child, not an
1111    // arbitrary process group.
1112    // SAFETY: signal 0 is only an existence probe, and the negative PID targets
1113    // the daemon-owned group rather than another process.
1114    if unsafe { libc::kill(-pid, 0) } == 0 {
1115        return true;
1116    }
1117    let error = std::io::Error::last_os_error();
1118    !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
1119}
1120
1121pub fn process_is_alive(pid: u32) -> bool {
1122    #[cfg(unix)]
1123    {
1124        if pid == 0 {
1125            return false;
1126        }
1127        let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
1128            return false;
1129        };
1130        // SAFETY: kill(pid, 0) sends no signal and is the standard existence
1131        // probe. EPERM still means the process exists.
1132        let result = unsafe { libc::kill(raw_pid, 0) };
1133        let exists =
1134            result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
1135        exists && !process_is_zombie(pid)
1136    }
1137    #[cfg(not(unix))]
1138    {
1139        let process_id = sysinfo::Pid::from_u32(pid);
1140        let mut system = sysinfo::System::new();
1141        system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[process_id]), true);
1142        system.process(process_id).is_some()
1143    }
1144}
1145
1146pub fn read_metadata() -> Result<DaemonMetadata> {
1147    let metadata = read_metadata_any()?;
1148    ensure!(
1149        metadata.protocol_version == PROTOCOL_VERSION,
1150        "daemon protocol {} is incompatible with client protocol {}",
1151        metadata.protocol_version,
1152        PROTOCOL_VERSION
1153    );
1154    Ok(metadata)
1155}
1156
1157/// No daemon has published its endpoint: `daemon.json` does not exist.
1158///
1159/// This is the ordinary stopped state, not a failure. A daemon removes the
1160/// file when it stops, so callers show "not running" rather than the raw
1161/// file error.
1162#[derive(Debug, Clone, PartialEq, Eq)]
1163pub struct DaemonNotRunning {
1164    /// Where the endpoint would be, which also names the instance.
1165    pub metadata_path: std::path::PathBuf,
1166}
1167
1168/// What a caller shows for [`DaemonNotRunning`].
1169pub const DAEMON_NOT_RUNNING_MESSAGE: &str = "the Mjolnir daemon is not running";
1170
1171impl std::fmt::Display for DaemonNotRunning {
1172    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1173        formatter.write_str(DAEMON_NOT_RUNNING_MESSAGE)
1174    }
1175}
1176
1177impl std::error::Error for DaemonNotRunning {}
1178
1179/// The stopped state, if `error` reports it anywhere in its chain.
1180pub fn daemon_not_running(error: &anyhow::Error) -> Option<&DaemonNotRunning> {
1181    error
1182        .chain()
1183        .find_map(|cause| cause.downcast_ref::<DaemonNotRunning>())
1184}
1185
1186pub fn read_metadata_any() -> Result<DaemonMetadata> {
1187    read_metadata_at(&metadata_path())
1188}
1189
1190fn read_metadata_at(path: &std::path::Path) -> Result<DaemonMetadata> {
1191    let body = match fs::read(path) {
1192        Ok(body) => body,
1193        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1194            return Err(DaemonNotRunning {
1195                metadata_path: path.to_owned(),
1196            }
1197            .into());
1198        }
1199        Err(error) => return Err(error).with_context(|| format!("read {}", path.display())),
1200    };
1201    let metadata: DaemonMetadata =
1202        serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1203    Ok(metadata)
1204}
1205
1206pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
1207    let body = serde_json::to_vec(value)?;
1208    write_encoded_frame(stream, &body).await
1209}
1210
1211pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
1212    ensure!(
1213        body.len() <= MAX_FRAME_BYTES,
1214        "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
1215        body.len()
1216    );
1217    stream.write_u32(body.len() as u32).await?;
1218    stream.write_all(body).await?;
1219    stream.flush().await?;
1220    Ok(())
1221}
1222
1223pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
1224    let length = stream.read_u32().await? as usize;
1225    ensure!(
1226        length <= MAX_FRAME_BYTES,
1227        "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
1228    );
1229    let mut body = vec![0_u8; length];
1230    stream.read_exact(&mut body).await?;
1231    serde_json::from_slice(&body).context("decode daemon frame")
1232}
1233
1234/// Read one reply, reassembling a chunked one before exposing it. All
1235/// fragments belong to the same response; disconnects discard partial state.
1236pub async fn read_response(stream: &mut TcpStream) -> Result<ResponseEnvelope> {
1237    let mut response: ResponseEnvelope = read_frame(stream).await?;
1238    if !matches!(response.result, Ok(DaemonReply::ReplyChunk { .. })) {
1239        return Ok(response);
1240    }
1241    let protocol_version = response.protocol_version;
1242    let request_id = response.request_id;
1243    let mut body = Vec::new();
1244    loop {
1245        ensure!(
1246            response.protocol_version == protocol_version && response.request_id == request_id,
1247            "daemon crossed reply fragment identities"
1248        );
1249        let Ok(DaemonReply::ReplyChunk { bytes, finished }) = response.result else {
1250            bail!("daemon interrupted a chunked reply");
1251        };
1252        ensure!(!bytes.is_empty(), "empty reply fragment");
1253        body.extend(bytes);
1254        if finished {
1255            break;
1256        }
1257        response = read_frame(stream).await?;
1258    }
1259    let reply: DaemonReply = tokio::task::spawn_blocking(move || serde_json::from_slice(&body))
1260        .await
1261        .context("reply decoder task failed")??;
1262    ensure!(
1263        reply.is_chunked(),
1264        "daemon chunked a reply that is never chunked"
1265    );
1266    Ok(ResponseEnvelope {
1267        protocol_version,
1268        request_id,
1269        result: Ok(reply),
1270    })
1271}
1272
1273pub struct DaemonClient {
1274    metadata: DaemonMetadata,
1275    stream: TcpStream,
1276    next_request_id: u64,
1277}
1278
1279impl DaemonClient {
1280    pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
1281        let stream =
1282            tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
1283                .await
1284                .context("time out connecting to Mjolnir daemon")??;
1285        Ok(Self {
1286            metadata,
1287            stream,
1288            next_request_id: 1,
1289        })
1290    }
1291
1292    /// The process this client is connected to.
1293    pub fn daemon_pid(&self) -> u32 {
1294        self.metadata.pid
1295    }
1296
1297    /// Speak the daemon's advertised dialect, not this build's: management
1298    /// requests must reach daemons of any protocol version, and the frozen
1299    /// subset encodes identically across all of them.
1300    pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1301        self.request_with_reconnect(action, || async {
1302            loop {
1303                if let Ok(client) = connect_existing().await {
1304                    return Ok(client);
1305                }
1306                tokio::time::sleep(Duration::from_millis(250)).await;
1307            }
1308        })
1309        .await
1310    }
1311
1312    async fn request_with_reconnect<F, Fut>(
1313        &mut self,
1314        action: DaemonAction,
1315        mut reconnect: F,
1316    ) -> Result<DaemonReply>
1317    where
1318        F: FnMut() -> Fut,
1319        Fut: Future<Output = Result<Self>>,
1320    {
1321        loop {
1322            let response = self.request_once(action.clone()).await?;
1323            if matches!(response, DaemonReply::UpgradePending)
1324                && !matches!(action, DaemonAction::PrepareUpgrade)
1325            {
1326                // Only an explicit refusal guarantees non-admission. Never
1327                // replay arbitrary mutations after a lost acknowledgement.
1328                tokio::time::sleep(Duration::from_millis(100)).await;
1329                *self = reconnect().await?;
1330            } else {
1331                return Ok(response);
1332            }
1333        }
1334    }
1335
1336    async fn request_once(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1337        let protocol_version = self.metadata.protocol_version;
1338        let request_id = self.next_request_id;
1339        self.next_request_id += 1;
1340        write_frame(
1341            &mut self.stream,
1342            &RequestEnvelope {
1343                protocol_version,
1344                request_id,
1345                token: self.metadata.token.clone(),
1346                action,
1347            },
1348        )
1349        .await?;
1350        let response = read_response(&mut self.stream).await?;
1351        ensure!(
1352            response.protocol_version == protocol_version,
1353            "daemon changed protocol"
1354        );
1355        ensure!(
1356            response.request_id == request_id,
1357            "daemon crossed request IDs"
1358        );
1359        response
1360            .result
1361            .map_err(|message| anyhow::Error::new(DaemonRefusal(message)))
1362    }
1363
1364    pub async fn status(&mut self) -> Result<DaemonStatus> {
1365        match self.request(DaemonAction::Status).await? {
1366            DaemonReply::Status(status) => Ok(status),
1367            reply => bail!("unexpected daemon status reply {reply:?}"),
1368        }
1369    }
1370
1371    pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
1372        match self.request(DaemonAction::WebViewerAccess).await? {
1373            DaemonReply::WebViewerAccess(access) => Ok(access),
1374            reply => bail!("unexpected web viewer reply {reply:?}"),
1375        }
1376    }
1377
1378    pub async fn subagent_options(
1379        &mut self,
1380        profile: String,
1381        model: Option<String>,
1382        config: Option<Config>,
1383    ) -> Result<mj_core::subagent::SubagentOptions> {
1384        match self
1385            .request(DaemonAction::SubagentOptions {
1386                profile,
1387                model,
1388                config: config.map(Box::new),
1389            })
1390            .await?
1391        {
1392            DaemonReply::SubagentOptions(options) => Ok(options),
1393            reply => bail!("unexpected subagent options reply {reply:?}"),
1394        }
1395    }
1396
1397    pub async fn warm_profile_capabilities(&mut self, config: Config) -> Result<()> {
1398        match self
1399            .request(DaemonAction::WarmProfileCapabilities {
1400                config: Box::new(config),
1401            })
1402            .await?
1403        {
1404            DaemonReply::Done => Ok(()),
1405            reply => bail!("unexpected capability hydration reply {reply:?}"),
1406        }
1407    }
1408
1409    pub async fn recover_web_viewer(
1410        &mut self,
1411        action: crate::web::WebViewerRecovery,
1412    ) -> Result<()> {
1413        match self.request(DaemonAction::RecoverWebViewer(action)).await? {
1414            DaemonReply::Done => Ok(()),
1415            reply => bail!("unexpected web viewer recovery reply {reply:?}"),
1416        }
1417    }
1418
1419    pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
1420        match self.request(DaemonAction::InspectWebListener).await? {
1421            DaemonReply::WebListeners(processes) => Ok(processes),
1422            reply => bail!("unexpected listener inspection reply {reply:?}"),
1423        }
1424    }
1425
1426    pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
1427        match self.request(DaemonAction::ListWorkspaces).await? {
1428            DaemonReply::Workspaces(workspaces) => Ok(workspaces),
1429            reply => bail!("unexpected daemon workspace reply {reply:?}"),
1430        }
1431    }
1432
1433    /// Ask the daemon to probe every profile that is not on hold. Returns when
1434    /// the daemon has accepted the request, not when the probes finish.
1435    pub async fn refresh_quota(&mut self) -> Result<()> {
1436        match self.request(DaemonAction::RefreshQuota).await? {
1437            DaemonReply::Done => Ok(()),
1438            reply => bail!("unexpected refresh-quota reply {reply:?}"),
1439        }
1440    }
1441
1442    pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
1443        match self
1444            .request(DaemonAction::RenameProfile { old_id, new_id })
1445            .await?
1446        {
1447            DaemonReply::Done => Ok(()),
1448            reply => bail!("unexpected rename-profile reply {reply:?}"),
1449        }
1450    }
1451
1452    pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1453        match self
1454            .request(DaemonAction::RenameTarget { old_id, new_id })
1455            .await?
1456        {
1457            DaemonReply::Done => Ok(()),
1458            reply => bail!("unexpected rename-target reply {reply:?}"),
1459        }
1460    }
1461
1462    pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1463        match self.request(DaemonAction::CreateWorkspace { name }).await? {
1464            DaemonReply::Workspace(workspace) => Ok(workspace),
1465            reply => bail!("unexpected create-workspace reply {reply:?}"),
1466        }
1467    }
1468
1469    pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1470        match self
1471            .request(DaemonAction::RenameWorkspace { workspace_id, name })
1472            .await?
1473        {
1474            DaemonReply::Done => Ok(()),
1475            reply => bail!("unexpected rename-workspace reply {reply:?}"),
1476        }
1477    }
1478
1479    pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1480        match self
1481            .request(DaemonAction::TouchWorkspace { workspace_id })
1482            .await?
1483        {
1484            DaemonReply::Done => Ok(()),
1485            reply => bail!("unexpected touch-workspace reply {reply:?}"),
1486        }
1487    }
1488
1489    pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1490        match self
1491            .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1492            .await?
1493        {
1494            DaemonReply::Done => Ok(()),
1495            reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1496        }
1497    }
1498
1499    pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1500        match self
1501            .request(DaemonAction::CloseWorkspace { workspace_id })
1502            .await?
1503        {
1504            DaemonReply::Done => Ok(()),
1505            reply => bail!("unexpected close workspace reply: {reply:?}"),
1506        }
1507    }
1508
1509    pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1510        match self
1511            .request(DaemonAction::DeleteWorkspace { workspace_id })
1512            .await?
1513        {
1514            DaemonReply::Done => Ok(()),
1515            reply => bail!("unexpected delete-workspace reply {reply:?}"),
1516        }
1517    }
1518
1519    pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1520        match self
1521            .request(DaemonAction::Attach { client_id, pid })
1522            .await?
1523        {
1524            DaemonReply::Done => Ok(()),
1525            reply => bail!("unexpected attach reply {reply:?}"),
1526        }
1527    }
1528
1529    pub async fn detach(&mut self, client_id: String) -> Result<()> {
1530        match self.request(DaemonAction::Detach { client_id }).await? {
1531            DaemonReply::Done => Ok(()),
1532            reply => bail!("unexpected detach reply {reply:?}"),
1533        }
1534    }
1535
1536    pub async fn persist_read_receipt(
1537        &mut self,
1538        client_id: String,
1539        workspace_id: String,
1540        session_id: String,
1541        through: u64,
1542    ) -> Result<u64> {
1543        match self
1544            .request(DaemonAction::PersistReadReceipt {
1545                client_id,
1546                workspace_id,
1547                session_id,
1548                through,
1549            })
1550            .await?
1551        {
1552            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1553            reply => bail!("unexpected read-receipt reply {reply:?}"),
1554        }
1555    }
1556
1557    pub async fn persist_detached_session_state(
1558        &mut self,
1559        client_id: String,
1560        workspace_id: String,
1561        session_id: String,
1562        through: u64,
1563        owner_pid: u32,
1564        draft: mj_core::storage::DetachedSessionDraft,
1565    ) -> Result<()> {
1566        match self
1567            .request(DaemonAction::PersistDetachedSessionState {
1568                client_id,
1569                workspace_id,
1570                session_id,
1571                through,
1572                owner_pid,
1573                draft,
1574            })
1575            .await?
1576        {
1577            DaemonReply::Done => Ok(()),
1578            reply => bail!("unexpected detached-session-state reply {reply:?}"),
1579        }
1580    }
1581
1582    pub async fn save_active_review(
1583        &mut self,
1584        session_id: String,
1585        review: mj_core::storage::StoredReview,
1586    ) -> Result<()> {
1587        match self
1588            .request(DaemonAction::SaveActiveReview { session_id, review })
1589            .await?
1590        {
1591            DaemonReply::Done => Ok(()),
1592            reply => bail!("unexpected save-review reply {reply:?}"),
1593        }
1594    }
1595
1596    pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1597        match self
1598            .request(DaemonAction::ClearActiveReview { session_id })
1599            .await?
1600        {
1601            DaemonReply::Done => Ok(()),
1602            reply => bail!("unexpected clear-review reply {reply:?}"),
1603        }
1604    }
1605
1606    pub async fn save_workspace_pane_sizes(
1607        &mut self,
1608        workspace_id: String,
1609        sizes: mj_core::workspace::PaneSizes,
1610    ) -> Result<()> {
1611        match self
1612            .request(DaemonAction::SaveWorkspacePaneSizes {
1613                workspace_id,
1614                sizes,
1615            })
1616            .await?
1617        {
1618            DaemonReply::Done => Ok(()),
1619            reply => bail!("unexpected pane-size save reply {reply:?}"),
1620        }
1621    }
1622
1623    pub async fn save_workspace_layout(
1624        &mut self,
1625        workspace_id: String,
1626        layout: mj_core::workspace::ConversationLayout,
1627    ) -> Result<()> {
1628        match self
1629            .request(DaemonAction::SaveWorkspaceLayout {
1630                workspace_id,
1631                layout,
1632            })
1633            .await?
1634        {
1635            DaemonReply::Done => Ok(()),
1636            reply => bail!("unexpected layout save reply {reply:?}"),
1637        }
1638    }
1639
1640    pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1641        match self
1642            .request(DaemonAction::PersistImportedSession {
1643                session: Box::new(session),
1644            })
1645            .await?
1646        {
1647            DaemonReply::Done => Ok(()),
1648            reply => bail!("unexpected imported-session reply {reply:?}"),
1649        }
1650    }
1651
1652    pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1653        match self
1654            .request(DaemonAction::SetSessionTitle { session_id, title })
1655            .await?
1656        {
1657            DaemonReply::Text(title) => Ok(title),
1658            reply => bail!("unexpected session-title reply {reply:?}"),
1659        }
1660    }
1661
1662    pub async fn set_session_workspace(
1663        &mut self,
1664        session_id: String,
1665        workspace_id: String,
1666    ) -> Result<()> {
1667        match self
1668            .request(DaemonAction::SetSessionWorkspace {
1669                session_id,
1670                workspace_id,
1671            })
1672            .await?
1673        {
1674            DaemonReply::Done => Ok(()),
1675            reply => bail!("unexpected session-workspace reply {reply:?}"),
1676        }
1677    }
1678
1679    pub async fn set_session_container_settings(
1680        &mut self,
1681        session_id: String,
1682        cpus: Option<String>,
1683        memory: Option<String>,
1684        mounts: Vec<AdditionalMount>,
1685        mount_history: Vec<PathBuf>,
1686    ) -> Result<()> {
1687        match self
1688            .request(DaemonAction::SetSessionContainerSettings {
1689                session_id,
1690                cpus,
1691                memory,
1692                mounts,
1693                mount_history,
1694            })
1695            .await?
1696        {
1697            DaemonReply::Done => Ok(()),
1698            reply => bail!("unexpected container-settings reply {reply:?}"),
1699        }
1700    }
1701
1702    pub async fn set_session_acp_title(
1703        &mut self,
1704        session_id: String,
1705        title: Option<String>,
1706    ) -> Result<()> {
1707        match self
1708            .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1709            .await?
1710        {
1711            DaemonReply::Done => Ok(()),
1712            reply => bail!("unexpected ACP-title reply {reply:?}"),
1713        }
1714    }
1715
1716    pub async fn mark_session_target_missing(
1717        &mut self,
1718        session_id: String,
1719        detail: String,
1720        updated_at: String,
1721    ) -> Result<Option<SessionState>> {
1722        match self
1723            .request(DaemonAction::MarkSessionTargetMissing {
1724                session_id,
1725                detail,
1726                updated_at,
1727            })
1728            .await?
1729        {
1730            DaemonReply::OptionalSessionState(state) => Ok(state),
1731            reply => bail!("unexpected target-missing reply {reply:?}"),
1732        }
1733    }
1734
1735    pub async fn checkpoint_session(
1736        &mut self,
1737        session_id: String,
1738    ) -> Result<mj_core::state::CheckpointMetadata> {
1739        match self
1740            .request(DaemonAction::CheckpointSession { session_id })
1741            .await?
1742        {
1743            DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1744            reply => bail!("unexpected checkpoint reply {reply:?}"),
1745        }
1746    }
1747
1748    /// Search the user's SessionWiki index, newest first when the query is
1749    /// empty and best match first otherwise. The reply carries the state of
1750    /// the index as well as the rows, so a caller can say the first build is
1751    /// still running.
1752    pub async fn project_catalog(
1753        &mut self,
1754        refresh: bool,
1755        retry: bool,
1756    ) -> Result<mj_core::project_catalog::ProjectCatalogView> {
1757        match self
1758            .request(DaemonAction::ProjectCatalog { refresh, retry })
1759            .await?
1760        {
1761            DaemonReply::ProjectCatalog(view) => Ok(view),
1762            other => bail!("unexpected project catalog reply: {other:?}"),
1763        }
1764    }
1765
1766    pub async fn create_project(
1767        &mut self,
1768        sources: Vec<String>,
1769    ) -> Result<mj_core::project_catalog::SavedProject> {
1770        match self
1771            .request(DaemonAction::CreateProject { sources })
1772            .await?
1773        {
1774            DaemonReply::ProjectCreated(project) => Ok(project),
1775            other => bail!("unexpected project creation reply: {other:?}"),
1776        }
1777    }
1778
1779    pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1780        match self
1781            .request(DaemonAction::WikiSearch { query, limit })
1782            .await?
1783        {
1784            DaemonReply::WikiRows(page) => Ok(page),
1785            reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1786        }
1787    }
1788
1789    /// The live sessions whose user or agent messages contain `query`, and
1790    /// which of the two matched. Answers from the SessionWiki index, so a very
1791    /// new message can be missing until the next sync.
1792    pub async fn session_text_search(&mut self, query: String) -> Result<Vec<SessionTextMatch>> {
1793        match self
1794            .request(DaemonAction::SessionTextSearch { query })
1795            .await?
1796        {
1797            DaemonReply::SessionTextMatches(matches) => Ok(matches),
1798            reply => bail!("unexpected session text search reply {reply:?}"),
1799        }
1800    }
1801
1802    /// The markdown briefing for one indexed session.
1803    pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1804        match self
1805            .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1806            .await?
1807        {
1808            DaemonReply::Text(markdown) => Ok(markdown),
1809            reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1810        }
1811    }
1812
1813    /// The passages of one indexed session that match a query, each matching
1814    /// message with `context_messages` neighbours on either side and its text
1815    /// capped at `per_message_chars`. `None` when the index holds no session
1816    /// with that id.
1817    pub async fn wiki_hits(
1818        &mut self,
1819        wiki_id: String,
1820        query: String,
1821        context_messages: usize,
1822        per_message_chars: usize,
1823    ) -> Result<Option<WikiHitTranscript>> {
1824        match self
1825            .request(DaemonAction::WikiHits {
1826                wiki_id,
1827                query,
1828                context_messages,
1829                per_message_chars,
1830            })
1831            .await?
1832        {
1833            DaemonReply::WikiHits(transcript) => Ok(transcript),
1834            reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1835        }
1836    }
1837
1838    /// What the index knows about one session, or `None` when the index holds
1839    /// no session with that id.
1840    pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1841        match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1842            DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1843            reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1844        }
1845    }
1846
1847    /// Start a new session carrying a hand-off compacted from an archived one.
1848    /// It answers like any other session start: the record exists and is
1849    /// provisioning, and the hand-off follows once the harness is ready.
1850    pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1851        match self.request(DaemonAction::WikiRestore(request)).await? {
1852            DaemonReply::RegisteredSession(registered) => Ok(*registered),
1853            reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1854        }
1855    }
1856
1857    pub async fn scan_recovery(
1858        &mut self,
1859        all_instances: bool,
1860    ) -> Result<mj_core::state::RecoveryScan> {
1861        match self
1862            .request(DaemonAction::ScanRecovery { all_instances })
1863            .await?
1864        {
1865            DaemonReply::RecoveryScan(scan) => Ok(scan),
1866            reply => bail!("unexpected recovery-scan reply {reply:?}"),
1867        }
1868    }
1869
1870    pub async fn adopt_recovery(
1871        &mut self,
1872        session_id: String,
1873        target_id: String,
1874        profile: Option<String>,
1875        bundle: Option<String>,
1876        all_instances: bool,
1877    ) -> Result<()> {
1878        match self
1879            .request(DaemonAction::AdoptRecovery {
1880                session_id,
1881                target_id,
1882                profile,
1883                bundle,
1884                all_instances,
1885            })
1886            .await?
1887        {
1888            DaemonReply::Done => Ok(()),
1889            reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1890        }
1891    }
1892
1893    pub async fn destroy_recovery(
1894        &mut self,
1895        session_id: String,
1896        target_id: String,
1897        confirmation: String,
1898        all_instances: bool,
1899    ) -> Result<()> {
1900        match self
1901            .request(DaemonAction::DestroyRecovery {
1902                session_id,
1903                target_id,
1904                confirmation,
1905                all_instances,
1906            })
1907            .await?
1908        {
1909            DaemonReply::Done => Ok(()),
1910            reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1911        }
1912    }
1913
1914    pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1915        match self
1916            .request(DaemonAction::Snapshot { workspace_id })
1917            .await?
1918        {
1919            DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1920            reply => bail!("unexpected snapshot reply {reply:?}"),
1921        }
1922    }
1923
1924    pub async fn runtime_changes(
1925        &mut self,
1926        cursor: Option<crate::runtime_feed::RuntimeCursor>,
1927        wait: bool,
1928    ) -> Result<crate::runtime_feed::RuntimeFrame> {
1929        match self
1930            .request(DaemonAction::RuntimeChanges { cursor, wait })
1931            .await?
1932        {
1933            DaemonReply::RuntimeChanges(frame) => Ok(*frame),
1934            reply => bail!("unexpected runtime changes reply {reply:?}"),
1935        }
1936    }
1937
1938    pub async fn session_tail(
1939        &mut self,
1940        session_id: String,
1941        cursor: crate::runtime_feed::RuntimeCursor,
1942    ) -> Result<crate::runtime_feed::SessionTailReply> {
1943        match self
1944            .request(DaemonAction::SessionTail { session_id, cursor })
1945            .await?
1946        {
1947            DaemonReply::SessionTail(reply) => Ok(*reply),
1948            reply => bail!("unexpected session tail reply {reply:?}"),
1949        }
1950    }
1951
1952    pub async fn resume_candidates(&mut self) -> Result<ResumeCandidates> {
1953        match self.request(DaemonAction::ResumeCandidates).await? {
1954            DaemonReply::ResumeCandidates(candidates) => Ok(*candidates),
1955            reply => bail!("unexpected resume candidates reply {reply:?}"),
1956        }
1957    }
1958
1959    pub async fn session_record(&mut self, session_id: String) -> Result<Option<SessionRecord>> {
1960        match self
1961            .request(DaemonAction::SessionRecord { session_id })
1962            .await?
1963        {
1964            DaemonReply::SessionRecord(record) => Ok(record.map(|record| *record)),
1965            reply => bail!("unexpected session record reply {reply:?}"),
1966        }
1967    }
1968
1969    pub async fn go_startup_session(
1970        &mut self,
1971        workspace_id: String,
1972        last_session_id: Option<String>,
1973    ) -> Result<Option<SessionRecord>> {
1974        match self
1975            .request(DaemonAction::GoStartupSession {
1976                workspace_id,
1977                last_session_id,
1978            })
1979            .await?
1980        {
1981            DaemonReply::GoStartupSession(record) => Ok(record.map(|record| *record)),
1982            reply => bail!("unexpected go startup session reply {reply:?}"),
1983        }
1984    }
1985
1986    pub async fn submit_session_command(
1987        &mut self,
1988        session_id: String,
1989        command_id: String,
1990        command: RelayCommand,
1991        inherited_draft: Option<String>,
1992    ) -> Result<u64> {
1993        match self
1994            .request(DaemonAction::SubmitSessionCommand {
1995                inherited_draft,
1996                session_id,
1997                command_id,
1998                command,
1999            })
2000            .await?
2001        {
2002            DaemonReply::Ordinal(ordinal) => Ok(ordinal),
2003            reply => bail!("unexpected session command reply {reply:?}"),
2004        }
2005    }
2006
2007    /// Hand the daemon a prompt for a session that is still starting. The
2008    /// daemon replies as soon as the prompt is queued, not when it is
2009    /// delivered; delivery failures come back as a session notice.
2010    pub async fn queue_startup_prompt(
2011        &mut self,
2012        session_id: String,
2013        text: String,
2014        inherited_draft: Option<String>,
2015    ) -> Result<()> {
2016        match self
2017            .request(DaemonAction::QueueStartupPrompt {
2018                session_id,
2019                text,
2020                inherited_draft,
2021            })
2022            .await?
2023        {
2024            DaemonReply::Done => Ok(()),
2025            reply => bail!("unexpected startup prompt reply {reply:?}"),
2026        }
2027    }
2028
2029    /// Take a prompt back from the daemon's startup queue before delivery.
2030    /// `Ok(true)` means it is withdrawn and will not be sent; `Ok(false)`
2031    /// means delivery had already started, so it was sent. An error means the
2032    /// daemon could not be asked, and the prompt may still be queued.
2033    pub async fn withdraw_startup_prompt(
2034        &mut self,
2035        session_id: String,
2036        text: String,
2037    ) -> Result<bool> {
2038        match self
2039            .request(DaemonAction::WithdrawStartupPrompt { session_id, text })
2040            .await?
2041        {
2042            DaemonReply::PromptWithdrawn(withdrawn) => Ok(withdrawn),
2043            reply => bail!("unexpected startup prompt withdrawal reply {reply:?}"),
2044        }
2045    }
2046
2047    /// Ask the daemon to review the turn this session just finished.
2048    ///
2049    /// The refusal is a sentence for a person -- "prompts are queued", "set
2050    /// [review] profile in config.toml" -- so it travels as text rather than
2051    /// as a code every surface would have to translate.
2052    pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
2053        match self
2054            .request(DaemonAction::StartTurnReview { session_id })
2055            .await?
2056        {
2057            DaemonReply::Done => Ok(()),
2058            reply => bail!("unexpected turn-review reply {reply:?}"),
2059        }
2060    }
2061
2062    pub async fn resolve_turn_review(
2063        &mut self,
2064        session_id: String,
2065        resolution: Resolution,
2066    ) -> Result<()> {
2067        match self
2068            .request(DaemonAction::ResolveTurnReview {
2069                session_id,
2070                resolution,
2071            })
2072            .await?
2073        {
2074            DaemonReply::Done => Ok(()),
2075            reply => bail!("unexpected turn-review resolution reply {reply:?}"),
2076        }
2077    }
2078
2079    pub async fn reviewer_action(
2080        &mut self,
2081        session_id: String,
2082        role: Option<String>,
2083        action: crate::session::ReviewerAction,
2084    ) -> Result<crate::session::ReviewerOutcome> {
2085        match self
2086            .request(DaemonAction::ReviewerAction {
2087                session_id,
2088                role,
2089                action,
2090            })
2091            .await?
2092        {
2093            DaemonReply::Reviewer(outcome) => Ok(*outcome),
2094            reply => bail!("unexpected reviewer reply {reply:?}"),
2095        }
2096    }
2097
2098    pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
2099        match self
2100            .request(DaemonAction::SyncSession { session_id })
2101            .await?
2102        {
2103            DaemonReply::Done => Ok(()),
2104            reply => bail!("unexpected session sync reply {reply:?}"),
2105        }
2106    }
2107
2108    pub async fn respond_elicitation(
2109        &mut self,
2110        session_id: String,
2111        elicitation_id: String,
2112        response: ElicitationResponse,
2113    ) -> Result<()> {
2114        match self
2115            .request(DaemonAction::RespondElicitation {
2116                session_id,
2117                elicitation_id,
2118                response,
2119            })
2120            .await?
2121        {
2122            DaemonReply::Done => Ok(()),
2123            reply => bail!("unexpected elicitation reply {reply:?}"),
2124        }
2125    }
2126
2127    pub async fn native_agent_history(
2128        &mut self,
2129        owner: String,
2130        child: String,
2131        before: Option<(u64, String)>,
2132    ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
2133        match self
2134            .request(DaemonAction::NativeAgentHistory {
2135                owner,
2136                child,
2137                before,
2138            })
2139            .await?
2140        {
2141            DaemonReply::NativeAgentHistory(page) => Ok(page),
2142            reply => bail!("unexpected native agent history reply {reply:?}"),
2143        }
2144    }
2145
2146    pub async fn stop_background_task(
2147        &mut self,
2148        session_id: String,
2149        background_task_id: String,
2150    ) -> Result<()> {
2151        match self
2152            .request(DaemonAction::StopBackgroundTask {
2153                session_id,
2154                background_task_id,
2155            })
2156            .await?
2157        {
2158            DaemonReply::Done => Ok(()),
2159            reply => bail!("unexpected background task stop reply {reply:?}"),
2160        }
2161    }
2162
2163    pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
2164        self.suspend_session_with_ack(session_id, false).await
2165    }
2166
2167    pub async fn suspend_session_with_ack(
2168        &mut self,
2169        session_id: String,
2170        acknowledge_unpublished_work: bool,
2171    ) -> Result<()> {
2172        match self
2173            .request(DaemonAction::SuspendSession {
2174                session_id,
2175                acknowledge_unpublished_work,
2176            })
2177            .await?
2178        {
2179            DaemonReply::Done => Ok(()),
2180            reply => bail!("unexpected close-session reply {reply:?}"),
2181        }
2182    }
2183
2184    pub async fn start_create_session(
2185        &mut self,
2186        request: CreateSessionRequest,
2187    ) -> Result<RegisteredSession> {
2188        match self
2189            .request(DaemonAction::StartCreateSession(request))
2190            .await?
2191        {
2192            DaemonReply::RegisteredSession(registered) => Ok(*registered),
2193            reply => bail!("unexpected start-create reply {reply:?}"),
2194        }
2195    }
2196
2197    pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
2198        match self
2199            .request(DaemonAction::WaitCreateSession { session_id })
2200            .await?
2201        {
2202            DaemonReply::Done => Ok(()),
2203            reply => bail!("unexpected wait-create reply {reply:?}"),
2204        }
2205    }
2206
2207    pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
2208        match self.request(DaemonAction::ResumeSession(request)).await? {
2209            DaemonReply::Done => Ok(()),
2210            reply => bail!("unexpected resume-session reply {reply:?}"),
2211        }
2212    }
2213
2214    pub async fn discard_since_checkpoint(
2215        &mut self,
2216        session_id: String,
2217        checkpoint: mj_core::state::CheckpointMetadata,
2218    ) -> Result<()> {
2219        match self
2220            .request(DaemonAction::DiscardSinceCheckpoint {
2221                session_id,
2222                checkpoint,
2223            })
2224            .await?
2225        {
2226            DaemonReply::Done => Ok(()),
2227            reply => bail!("unexpected force-stop reply {reply:?}"),
2228        }
2229    }
2230
2231    pub async fn destroy_stopped_session(
2232        &mut self,
2233        session_id: String,
2234        delete_branch: bool,
2235    ) -> Result<()> {
2236        match self
2237            .request(DaemonAction::DestroyStoppedSession {
2238                session_id,
2239                delete_branch,
2240            })
2241            .await?
2242        {
2243            DaemonReply::Done => Ok(()),
2244            reply => bail!("unexpected destroy-stopped reply {reply:?}"),
2245        }
2246    }
2247
2248    pub async fn force_destroy_session(
2249        &mut self,
2250        session_id: String,
2251        delete_branch: bool,
2252    ) -> Result<()> {
2253        match self
2254            .request(DaemonAction::ForceDestroySession {
2255                session_id,
2256                delete_branch,
2257            })
2258            .await?
2259        {
2260            DaemonReply::Done => Ok(()),
2261            reply => bail!("unexpected force-destroy reply {reply:?}"),
2262        }
2263    }
2264
2265    pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
2266        match self
2267            .request(DaemonAction::CancelLifecycle { session_id })
2268            .await?
2269        {
2270            DaemonReply::Done => Ok(()),
2271            reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
2272        }
2273    }
2274
2275    pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
2276        match self
2277            .request(DaemonAction::RecoverDraft { draft_id })
2278            .await?
2279        {
2280            DaemonReply::Done => Ok(()),
2281            reply => bail!("unexpected recover-draft reply {reply:?}"),
2282        }
2283    }
2284
2285    pub async fn stop(&mut self) -> Result<()> {
2286        match self.request(DaemonAction::Stop).await? {
2287            DaemonReply::Done => Ok(()),
2288            reply => bail!("unexpected stop reply {reply:?}"),
2289        }
2290    }
2291}
2292
2293pub async fn connect_existing() -> Result<DaemonClient> {
2294    let metadata = tokio::task::spawn_blocking(read_metadata)
2295        .await
2296        .context("read daemon metadata task failed")??;
2297    DaemonClient::connect(metadata).await
2298}
2299
2300/// A handle to whatever daemon the metadata file advertises, regardless of its
2301/// protocol version. It only exposes the frozen management subset (`Ping`,
2302/// `Status`, `Stop`), which encodes identically in every protocol version.
2303pub struct ManagementClient {
2304    inner: DaemonClient,
2305}
2306
2307impl ManagementClient {
2308    pub fn new(inner: DaemonClient) -> Self {
2309        Self { inner }
2310    }
2311    pub fn protocol_version(&self) -> u32 {
2312        self.inner.metadata.protocol_version
2313    }
2314
2315    pub async fn status(&mut self) -> Result<DaemonStatus> {
2316        self.inner.status().await
2317    }
2318
2319    pub async fn stop(&mut self) -> Result<()> {
2320        tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
2321            .await
2322            .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
2323    }
2324
2325    /// Ask the daemon to stop and wait for its process to actually exit.
2326    pub async fn stop_and_wait(mut self) -> Result<()> {
2327        let pid = self.inner.metadata.pid;
2328        self.stop().await?;
2329        wait_for_exit_within(pid, STOP_DRAIN_TIMEOUT)
2330            .await
2331            .with_context(|| {
2332                format!(
2333                    "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
2334                    STOP_DRAIN_TIMEOUT.as_secs()
2335                )
2336            })
2337    }
2338}
2339
2340pub async fn connect_management() -> Result<ManagementClient> {
2341    Ok(ManagementClient {
2342        inner: DaemonClient::connect(read_metadata_any()?).await?,
2343    })
2344}
2345
2346/// Refuse to speak to a daemon whose protocol this build does not know.
2347///
2348/// The two protocol numbers alone do not say which `mj` ran. The usual cause is
2349/// a second installation: a `cargo install`ed client sits earlier on PATH than
2350/// the build whose daemon is running, so every command fails here while the
2351/// other binary works, and nothing in the message says where either lives. It
2352/// therefore names both executables and both versions.
2353pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
2354    ensure!(
2355        metadata.protocol_version <= PROTOCOL_VERSION,
2356        "{}",
2357        unsupported_daemon_protocol_message(
2358            metadata.protocol_version,
2359            &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
2360        )
2361    );
2362    Ok(())
2363}
2364
2365/// The message [`ensure_supported_daemon_protocol`] fails with, given the
2366/// sentence that names both builds, so it can be read without a daemon.
2367fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
2368    format!(
2369        "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
2370         Put the daemon's directory first on PATH, or reinstall this client from that build."
2371    )
2372}
2373// Settings can warm draft profile capabilities through the daemon-owned catalog.
2374// Runtime deltas carry transcript item changes, and a client fetches a
2375// session's tail from the daemon at its cursor instead of reading SQLite.
2376pub const PROTOCOL_VERSION: u32 = 53;
2377pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
2378/// How long a daemon is given to exit after it accepts a stop.
2379///
2380/// Stopping cancels a token and returns immediately; the daemon then unwinds
2381/// its session manager, its phone server and its pollers. That is normally
2382/// fast, but a daemon whose database has been migrated out from under it fails
2383/// every read while it winds down and has been observed taking over five
2384/// seconds — which the previous five-second bound missed by a fraction,
2385/// reporting a stop that had in fact worked as `did not stop` and aborting the
2386/// restart that depended on it.
2387pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
2388/// How long a requested stop may take: the daemon waits up to 60 seconds for
2389/// session destroys in flight (#1191) before it winds down.
2390pub const STOP_DRAIN_TIMEOUT: Duration = Duration::from_secs(90);
2391pub const RETRY_DELAY: Duration = Duration::from_millis(40);
2392impl DaemonClient {
2393    pub async fn prepare_move_session(
2394        &mut self,
2395        selection: MoveSelection,
2396    ) -> Result<MovePreparation> {
2397        match self
2398            .request(DaemonAction::PrepareMoveSession(selection))
2399            .await?
2400        {
2401            DaemonReply::MovePreparation(preparation) => Ok(*preparation),
2402            _ => bail!("daemon returned an unexpected move preparation reply"),
2403        }
2404    }
2405
2406    pub async fn move_sources(
2407        &mut self,
2408        session_id: String,
2409        cleanup_operation_id: Option<String>,
2410    ) -> Result<Vec<mj_core::move_workspace::RetainedMoveSource>> {
2411        match self
2412            .request(DaemonAction::MoveSources {
2413                session_id,
2414                cleanup_operation_id,
2415            })
2416            .await?
2417        {
2418            DaemonReply::MoveSources(sources) => Ok(sources),
2419            reply => bail!("unexpected Move sources reply: {reply:?}"),
2420        }
2421    }
2422
2423    pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
2424        match self.request(DaemonAction::MoveSession(request)).await? {
2425            DaemonReply::MoveOutcome(outcome) => Ok(outcome),
2426            _ => bail!("daemon returned an unexpected move reply"),
2427        }
2428    }
2429}
2430
2431#[cfg(test)]
2432mod tests {
2433    use super::*;
2434    use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
2435    use std::path::Path;
2436
2437    #[cfg(unix)]
2438    #[test]
2439    fn an_exited_unreaped_process_is_a_zombie_and_not_alive() {
2440        // An upgrade waits for the old daemon to leave. If that daemon exited
2441        // but is still a child nobody has reaped, `kill(pid, 0)` keeps finding
2442        // it, so only the zombie check lets the wait end.
2443        use std::time::{Duration, Instant};
2444
2445        let mut running = std::process::Command::new("sleep")
2446            .arg("30")
2447            .spawn()
2448            .unwrap();
2449        assert!(process_is_alive(running.id()));
2450        assert!(!process_is_zombie(running.id()));
2451        running.kill().unwrap();
2452        running.wait().unwrap();
2453
2454        let mut exited = std::process::Command::new("sh")
2455            .args(["-c", "exit 0"])
2456            .spawn()
2457            .unwrap();
2458        let pid = exited.id();
2459        let deadline = Instant::now() + Duration::from_secs(10);
2460        while !process_is_zombie(pid) {
2461            assert!(
2462                Instant::now() < deadline,
2463                "an exited, unreaped child was never reported as a zombie"
2464            );
2465            std::thread::sleep(Duration::from_millis(20));
2466        }
2467        assert!(!process_is_alive(pid));
2468        exited.wait().unwrap();
2469    }
2470
2471    #[test]
2472    fn a_missing_endpoint_file_reads_as_a_stopped_daemon() {
2473        let directory = tempfile::tempdir().unwrap();
2474        let path = directory.path().join("daemon.json");
2475        let error = read_metadata_at(&path).unwrap_err();
2476        let stopped = daemon_not_running(&error).expect("stopped, not an I/O failure");
2477        assert_eq!(stopped.metadata_path, path);
2478        assert_eq!(format!("{error:#}"), "the Mjolnir daemon is not running");
2479
2480        std::fs::write(&path, b"not json").unwrap();
2481        let error = read_metadata_at(&path).unwrap_err();
2482        assert!(daemon_not_running(&error).is_none());
2483    }
2484
2485    #[tokio::test]
2486    async fn subagent_discovery_uses_the_daemon_and_preserves_choices_and_failures() {
2487        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2488        let metadata = DaemonMetadata {
2489            protocol_version: PROTOCOL_VERSION,
2490            pid: std::process::id(),
2491            address: listener.local_addr().unwrap(),
2492            token: "subagent-discovery-test".into(),
2493            started_at: "test".into(),
2494            build_version: env!("CARGO_PKG_VERSION").into(),
2495        };
2496        let options: mj_core::subagent::SubagentOptions =
2497            serde_json::from_value(serde_json::json!({
2498                "models": [{"value": "gpt-6-luna", "name": "Luna", "description": null}],
2499                "efforts": [{"value": "high", "name": "High", "description": null}],
2500                "unavailable": ["another-profile: authentication failed"]
2501            }))
2502            .unwrap();
2503        let expected = serde_json::to_value(&options).unwrap();
2504        let mut draft = Config::default();
2505        draft.profiles.insert(
2506            "codex".into(),
2507            serde_json::from_value(serde_json::json!({
2508                "kind": "codex", "home": "/unsaved/profile",
2509                "subagents": {"mode": "single_model", "model": "gpt-6-luna", "effort": "high"}
2510            }))
2511            .unwrap(),
2512        );
2513        let expected_draft = serde_json::to_value(&draft).unwrap();
2514        let server = tokio::spawn(async move {
2515            let (mut stream, _) = listener.accept().await.unwrap();
2516            for (index, model) in [
2517                None,
2518                Some("gpt-6-luna".to_owned()),
2519                Some("gpt-6-luna".to_owned()),
2520            ]
2521            .into_iter()
2522            .enumerate()
2523            {
2524                let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2525                assert_eq!(request.token, "subagent-discovery-test");
2526                let DaemonAction::SubagentOptions {
2527                    profile,
2528                    model: requested,
2529                    config,
2530                } = request.action
2531                else {
2532                    panic!("unexpected discovery request")
2533                };
2534                assert_eq!(profile, "codex");
2535                assert_eq!(requested, model);
2536                assert_eq!(
2537                    serde_json::to_value(config).unwrap(),
2538                    if index == 1 {
2539                        expected_draft.clone()
2540                    } else {
2541                        serde_json::Value::Null
2542                    }
2543                );
2544                write_frame(
2545                    &mut stream,
2546                    &ResponseEnvelope {
2547                        protocol_version: request.protocol_version,
2548                        request_id: request.request_id,
2549                        result: if index != 2 {
2550                            Ok(DaemonReply::SubagentOptions(options.clone()))
2551                        } else {
2552                            Err("profile discovery cancelled".into())
2553                        },
2554                    },
2555                )
2556                .await
2557                .unwrap();
2558            }
2559        });
2560        let mut client = DaemonClient::connect(metadata).await.unwrap();
2561        let discovered = client
2562            .subagent_options("codex".into(), None, None)
2563            .await
2564            .unwrap();
2565        assert_eq!(serde_json::to_value(discovered).unwrap(), expected);
2566        let discovered = client
2567            .subagent_options("codex".into(), Some("gpt-6-luna".into()), Some(draft))
2568            .await
2569            .unwrap();
2570        assert_eq!(serde_json::to_value(discovered).unwrap(), expected);
2571        let error = client
2572            .subagent_options("codex".into(), Some("gpt-6-luna".into()), None)
2573            .await
2574            .unwrap_err();
2575        assert_eq!(error.to_string(), "profile discovery cancelled");
2576        server.await.unwrap();
2577    }
2578
2579    #[tokio::test]
2580    async fn upgrade_refusal_retries_the_identical_command_but_lost_acknowledgements_do_not() {
2581        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2582        let metadata = DaemonMetadata {
2583            protocol_version: PROTOCOL_VERSION,
2584            pid: std::process::id(),
2585            address: listener.local_addr().unwrap(),
2586            token: "test-token".into(),
2587            started_at: "test".into(),
2588            build_version: env!("CARGO_PKG_VERSION").into(),
2589        };
2590        let action = DaemonAction::SubmitSessionCommand {
2591            inherited_draft: None,
2592            session_id: "test-session".into(),
2593            command_id: "steer-command".into(),
2594            command: RelayCommand::Steer {
2595                active_prompt_id: "active-command".into(),
2596                queued_prompt_id: "queued-command".into(),
2597            },
2598        };
2599        let expected = serde_json::to_value(&action).unwrap();
2600        let server = tokio::spawn(async move {
2601            for reply in [
2602                Some(DaemonReply::UpgradePending),
2603                Some(DaemonReply::Done),
2604                None,
2605            ] {
2606                let (mut stream, _) = listener.accept().await.unwrap();
2607                let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2608                assert_eq!(serde_json::to_value(&request.action).unwrap(), expected);
2609                if let Some(reply) = reply {
2610                    write_frame(
2611                        &mut stream,
2612                        &ResponseEnvelope {
2613                            protocol_version: request.protocol_version,
2614                            request_id: request.request_id,
2615                            result: Ok(reply),
2616                        },
2617                    )
2618                    .await
2619                    .unwrap();
2620                }
2621            }
2622        });
2623        let mut client = DaemonClient::connect(metadata.clone()).await.unwrap();
2624        let reply = client
2625            .request_with_reconnect(action.clone(), || DaemonClient::connect(metadata.clone()))
2626            .await
2627            .unwrap();
2628        assert!(matches!(reply, DaemonReply::Done));
2629        let mut client = DaemonClient::connect(metadata).await.unwrap();
2630        assert!(
2631            client
2632                .request_with_reconnect(action, || async {
2633                    panic!("an ambiguous acknowledgement must not replay a mutation");
2634                })
2635                .await
2636                .is_err()
2637        );
2638        server.await.unwrap();
2639    }
2640
2641    #[test]
2642    fn unsupported_protocol_message_names_both_binaries_and_versions() {
2643        let message = unsupported_daemon_protocol_message(
2644            PROTOCOL_VERSION + 5,
2645            &describe_daemon_and_client_builds(
2646                4242,
2647                BuildDescription {
2648                    executable: Some(Path::new("/home/dev/mj/target/release/mj")),
2649                    version: "2.14.0",
2650                },
2651                BuildDescription {
2652                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2653                    version: "2.9.0",
2654                },
2655            ),
2656        );
2657        assert_eq!(
2658            message,
2659            format!(
2660                "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
2661                 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
2662                 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
2663                 Put the daemon's directory first on PATH, or reinstall this client from that build.",
2664                PROTOCOL_VERSION + 5
2665            )
2666        );
2667    }
2668
2669    #[test]
2670    fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
2671        let message = unsupported_daemon_protocol_message(
2672            PROTOCOL_VERSION + 1,
2673            &describe_daemon_and_client_builds(
2674                4242,
2675                BuildDescription {
2676                    executable: None,
2677                    version: "2.14.0",
2678                },
2679                BuildDescription {
2680                    executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2681                    version: "2.9.0",
2682                },
2683            ),
2684        );
2685        assert!(
2686            message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
2687            "{message}"
2688        );
2689        assert!(
2690            message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
2691            "{message}"
2692        );
2693    }
2694}