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