1use 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#[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 pub archived: bool,
70 pub native_id: Option<String>,
73 pub snippet: Option<String>,
75 pub hel_session_id: Option<String>,
77 #[serde(default)]
82 pub target: Option<String>,
83 #[serde(default)]
86 pub profile: Option<String>,
87 #[serde(default)]
91 pub harness: Option<String>,
92}
93
94#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
96#[serde(rename_all = "snake_case")]
97pub enum WikiIndexState {
98 Ready,
100 #[default]
103 Indexing,
104 VersionMismatch,
108}
109
110#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
112#[serde(deny_unknown_fields)]
113pub struct WikiStatus {
114 pub state: WikiIndexState,
115 pub topping_up: bool,
117}
118
119#[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#[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#[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#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
148#[serde(deny_unknown_fields)]
149pub struct WikiHitBlock {
150 pub role: String,
152 pub text: String,
155 pub hits: Vec<(usize, usize)>,
158 pub omitted_before: usize,
161 pub truncated: bool,
163}
164
165#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
167#[serde(deny_unknown_fields)]
168pub struct WikiHitTranscript {
169 pub blocks: Vec<WikiHitBlock>,
170 pub omitted_after: usize,
172}
173
174#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
176#[serde(rename_all = "snake_case")]
177pub enum WikiSessionStatus {
178 Mine,
180 Archived,
183 Native,
185}
186
187#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
192#[serde(deny_unknown_fields)]
193pub struct WikiSessionInfo {
194 pub wiki_id: String,
196 pub tool: String,
198 pub path: PathBuf,
200 pub status: WikiSessionStatus,
201 pub mjolnir_session_id: Option<String>,
203 pub profile_id: Option<String>,
205 pub target_template_id: Option<String>,
207 pub harness: Option<mj_core::config::HarnessKind>,
210 pub title: String,
211 pub project: String,
212 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
217 pub nothing_to_restore: bool,
218}
219
220#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
222#[serde(deny_unknown_fields)]
223pub struct WikiRestoreRequest {
224 pub wiki_id: String,
226 pub workspace_id: String,
227 pub profile_id: String,
228 pub target_template_id: String,
229 #[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#[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 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 #[serde(default, skip_serializing_if = "Option::is_none")]
347 pub at: Option<String>,
348 #[serde(default, skip_serializing_if = "Option::is_none")]
351 pub branch: Option<String>,
352 #[serde(default, skip_serializing_if = "Option::is_none")]
355 pub base: Option<String>,
356 #[serde(
358 default,
359 alias = "mjolnir_subagents",
360 deserialize_with = "mj_core::subagent::deserialize_optional_policy"
361 )]
362 pub subagents: Option<mj_core::subagent::SubagentPolicy>,
363 #[serde(default)]
364 pub initial_prompt: Option<String>,
365 pub workspace_id: String,
366 pub profile_id: String,
367 pub bundle_id: String,
368 pub project_directory: Option<PathBuf>,
369 pub target_template_id: String,
370 pub additional_mounts: Vec<AdditionalMount>,
371 pub resource_allocation: Option<SessionResourceAllocation>,
372 pub title: String,
373 pub session_title_override: Option<String>,
374}
375
376#[derive(Debug, Clone, Serialize, Deserialize)]
377#[serde(deny_unknown_fields)]
378pub struct RegisteredSession {
379 pub session: SessionRecord,
380 pub remembered_container_size: Option<(String, HostContainerSize)>,
381}
382
383#[derive(Debug, Clone, Serialize, Deserialize)]
384#[serde(deny_unknown_fields)]
385pub struct DraftPreview {
386 pub id: String,
387 pub session_id: Option<String>,
388 pub source: String,
389 pub owner_pid: Option<u32>,
390 pub saved_at: String,
391}
392
393#[derive(Debug, Clone, Serialize, Deserialize)]
402#[serde(rename_all = "snake_case", tag = "action", content = "arguments")]
403pub enum DaemonAction {
404 NativeAgentHistory {
405 owner: String,
406 child: String,
407 before: Option<(u64, String)>,
408 },
409 Ping,
410 PrepareUpgrade,
413 UpgradeBlockers,
422 Status,
423 WebViewerAccess,
424 RecoverWebViewer(crate::web::WebViewerRecovery),
425 InspectWebListener,
426 ListWorkspaces,
427 CreateWorkspace {
428 name: String,
429 },
430 RenameWorkspace {
431 workspace_id: String,
432 name: String,
433 },
434 TouchWorkspace {
435 workspace_id: String,
436 },
437 CloseWorkspace {
438 workspace_id: String,
439 },
440 CancelWorkspaceClose {
441 workspace_id: String,
442 },
443 DeleteWorkspace {
444 workspace_id: String,
445 },
446 Attach {
447 client_id: String,
448 pid: u32,
449 },
450 Detach {
451 client_id: String,
452 },
453 PersistReadReceipt {
454 client_id: String,
455 workspace_id: String,
456 session_id: String,
457 through: u64,
458 },
459 PersistDetachedSessionState {
460 client_id: String,
461 workspace_id: String,
462 session_id: String,
463 through: u64,
464 owner_pid: u32,
465 draft: mj_core::storage::DetachedSessionDraft,
466 },
467 SaveActiveReview {
468 session_id: String,
469 review: mj_core::storage::StoredReview,
470 },
471 ClearActiveReview {
472 session_id: String,
473 },
474 SaveWorkspacePaneSizes {
475 workspace_id: String,
476 sizes: mj_core::workspace::PaneSizes,
477 },
478 SaveWorkspaceLayout {
479 workspace_id: String,
480 layout: mj_core::workspace::ConversationLayout,
481 },
482 PersistImportedSession {
483 session: Box<SessionRecord>,
484 },
485 SetSessionTitle {
486 session_id: String,
487 title: String,
488 },
489 SetSessionWorkspace {
490 session_id: String,
491 workspace_id: String,
492 },
493 SetSessionContainerSettings {
494 session_id: String,
495 cpus: Option<String>,
496 memory: Option<String>,
497 mounts: Vec<AdditionalMount>,
498 mount_history: Vec<PathBuf>,
499 },
500 SetSessionAcpTitle {
501 session_id: String,
502 title: Option<String>,
503 },
504 MarkSessionTargetMissing {
505 session_id: String,
506 detail: String,
507 updated_at: String,
508 },
509 CheckpointSession {
510 session_id: String,
511 },
512 ProjectCatalog {
515 refresh: bool,
516 retry: bool,
517 },
518 CreateProject {
519 sources: Vec<String>,
520 },
521 WikiSearch {
522 query: String,
523 limit: usize,
524 },
525 SessionTextSearch {
528 query: String,
529 },
530 WikiBrief {
532 wiki_id: String,
533 max_chars: usize,
534 },
535 WikiHits {
537 wiki_id: String,
538 query: String,
539 context_messages: usize,
540 per_message_chars: usize,
541 },
542 WikiSession {
544 wiki_id: String,
545 },
546 WikiRestore(WikiRestoreRequest),
548 ScanRecovery {
549 all_instances: bool,
550 },
551 AdoptRecovery {
552 session_id: String,
553 target_id: String,
554 profile: Option<String>,
555 bundle: Option<String>,
556 all_instances: bool,
557 },
558 DestroyRecovery {
559 session_id: String,
560 target_id: String,
561 confirmation: String,
562 all_instances: bool,
563 },
564 Snapshot {
565 workspace_id: String,
566 },
567 RuntimeChanges {
568 cursor: Option<crate::runtime_feed::RuntimeCursor>,
569 wait: bool,
570 },
571 ResumeCandidates,
574 SessionRecord {
577 session_id: String,
578 },
579 GoStartupSession {
583 workspace_id: String,
584 last_session_id: Option<String>,
585 },
586 RefreshQuota,
590 SubagentOptions {
592 profile: String,
593 model: Option<String>,
594 #[serde(default, skip_serializing_if = "Option::is_none")]
596 config: Option<Box<Config>>,
597 },
598 RenameProfile {
599 old_id: String,
600 new_id: String,
601 },
602 RenameTarget {
603 old_id: String,
604 new_id: String,
605 },
606 SubmitSessionCommand {
607 #[serde(default)]
608 inherited_draft: Option<String>,
609 session_id: String,
610 command_id: String,
611 command: RelayCommand,
612 },
613 QueueStartupPrompt {
618 session_id: String,
619 text: String,
620 #[serde(default)]
624 inherited_draft: Option<String>,
625 },
626 WithdrawStartupPrompt {
632 session_id: String,
633 text: String,
634 },
635 SyncSession {
636 session_id: String,
637 },
638 RespondElicitation {
639 session_id: String,
640 elicitation_id: String,
641 response: ElicitationResponse,
642 },
643 StopBackgroundTask {
644 session_id: String,
645 background_task_id: String,
646 },
647 ReviewerAction {
651 session_id: String,
652 #[serde(default, skip_serializing_if = "Option::is_none")]
655 role: Option<String>,
656 action: crate::session::ReviewerAction,
657 },
658 StartTurnReview {
660 session_id: String,
661 },
662 ResolveTurnReview {
664 session_id: String,
665 resolution: Resolution,
666 },
667 SuspendSession {
668 session_id: String,
669 #[serde(default)]
670 acknowledge_unpublished_work: bool,
671 },
672 StartCreateSession(CreateSessionRequest),
673 WaitCreateSession {
674 session_id: String,
675 },
676 ResumeSession(ResumeSessionRequest),
677 PrepareMoveSession(MoveSelection),
678 MoveSession(MoveSessionRequest),
679 MoveSources {
680 session_id: String,
681 cleanup_operation_id: Option<String>,
682 },
683 DiscardSinceCheckpoint {
684 session_id: String,
685 checkpoint: mj_core::state::CheckpointMetadata,
686 },
687 DestroyStoppedSession {
688 session_id: String,
689 delete_branch: bool,
693 },
694 ForceDestroySession {
695 session_id: String,
696 delete_branch: bool,
698 },
699 CancelLifecycle {
700 session_id: String,
701 },
702 RecoverDraft {
703 draft_id: String,
704 },
705 Stop,
706}
707
708#[derive(Debug, Serialize, Deserialize)]
709#[serde(deny_unknown_fields)]
710pub struct RequestEnvelope {
711 pub protocol_version: u32,
712 pub request_id: u64,
713 pub token: String,
714 pub action: DaemonAction,
715}
716
717#[derive(Debug)]
721pub struct DaemonRefusal(pub String);
722
723impl std::fmt::Display for DaemonRefusal {
724 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
725 f.write_str(&self.0)
726 }
727}
728
729impl std::error::Error for DaemonRefusal {}
730
731impl DaemonRefusal {
732 #[must_use]
735 pub fn delivery_unconfirmed(&self) -> bool {
736 self.0
737 .starts_with(&crate::session::DeliveryUnconfirmed.to_string())
738 }
739}
740
741#[derive(Debug, Serialize, Deserialize)]
742#[serde(deny_unknown_fields)]
743pub struct ResponseEnvelope {
744 pub protocol_version: u32,
745 pub request_id: u64,
746 pub result: std::result::Result<DaemonReply, String>,
747}
748
749#[derive(Debug, Clone, Serialize, Deserialize)]
750#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
751pub enum DaemonReply {
752 NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
753 Pong,
754 UpgradePending,
755 UpgradeBlockers(Vec<String>),
758 Status(DaemonStatus),
759 WebViewerAccess(crate::web::WebViewerAccess),
760 WebListeners(Vec<crate::web::WebListenerProcess>),
761 SubagentOptions(mj_core::subagent::SubagentOptions),
762 Workspaces(Vec<WorkspaceListing>),
763 Workspace(WorkspaceRecord),
764 Snapshot(WorkspaceSnapshot),
765 RuntimeChanges(Box<crate::runtime_feed::RuntimeFrame>),
766 ResumeCandidates(Box<ResumeCandidates>),
767 GoStartupSession(Option<Box<SessionRecord>>),
768 SessionRecord(Option<Box<SessionRecord>>),
769 ReplyChunk {
772 bytes: Vec<u8>,
773 finished: bool,
774 },
775 RegisteredSession(Box<RegisteredSession>),
776 MovePreparation(Box<MovePreparation>),
777 MoveOutcome(MoveOutcome),
778 MoveSources(Vec<mj_core::move_workspace::RetainedMoveSource>),
779 Ordinal(u64),
780 Text(String),
781 OptionalSessionState(Option<SessionState>),
782 Checkpoint(mj_core::state::CheckpointMetadata),
783 RecoveryScan(mj_core::state::RecoveryScan),
784 WikiRows(WikiSearchPage),
785 ProjectCatalog(mj_core::project_catalog::ProjectCatalogView),
786 ProjectCreated(mj_core::project_catalog::SavedProject),
787 SessionTextMatches(Vec<SessionTextMatch>),
788 WikiHits(Option<WikiHitTranscript>),
789 WikiSession(Option<Box<WikiSessionInfo>>),
790 Reviewer(Box<crate::session::ReviewerOutcome>),
791 PromptWithdrawn(bool),
793 Done,
794}
795
796impl DaemonReply {
797 pub fn is_chunked(&self) -> bool {
800 matches!(self, Self::RuntimeChanges(_) | Self::ResumeCandidates(_))
801 }
802}
803
804#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
806#[serde(deny_unknown_fields)]
807pub struct ResumeCandidates {
808 pub candidates: Vec<ResumeCandidate>,
810 pub moves: Vec<MoveOperation>,
812 pub adopted_native_sessions: Vec<(mj_core::config::HarnessKind, String)>,
815 pub local_checkout_roots: Vec<PathBuf>,
818}
819
820#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
825#[serde(deny_unknown_fields)]
826pub struct ResumeCandidate {
827 pub session_id: String,
828 pub state: SessionState,
829 pub last_profile: String,
830 pub title: String,
832 pub origin: String,
834 pub project: String,
835 pub has_checkpoint: bool,
836 pub last_activity_ms: Option<i64>,
839 pub publication: Option<PublicationState>,
840 pub checkpoint_archive: Option<PathBuf>,
843 pub worktree_checkout: bool,
845}
846
847impl ResumeCandidate {
848 pub const TITLE_CHARS: usize = 200;
850
851 pub fn of(record: &SessionRecord, config: &Config) -> Self {
854 let timestamp_ms = |timestamp: &str| {
855 chrono::DateTime::parse_from_rfc3339(timestamp)
856 .ok()
857 .map(|parsed| parsed.timestamp_millis())
858 };
859 Self {
860 session_id: record.id.clone(),
861 state: record.state,
862 last_profile: record.last_profile.clone(),
863 title: record
864 .listed_title()
865 .chars()
866 .take(Self::TITLE_CHARS)
867 .collect(),
868 origin: record.project_target(config, &record.target_template_id),
869 project: record.project_name(config),
870 has_checkpoint: record.checkpoint.is_some(),
871 last_activity_ms: record
872 .checkpoint
873 .as_ref()
874 .and_then(|checkpoint| timestamp_ms(&checkpoint.created_at))
875 .or_else(|| timestamp_ms(&record.updated_at)),
876 publication: record.publication_state(),
877 checkpoint_archive: record
878 .checkpoint
879 .as_ref()
880 .filter(|_| record.state == SessionState::Stopped)
881 .map(|checkpoint| checkpoint.archive_path.clone()),
882 worktree_checkout: record
883 .managed_worktree
884 .as_ref()
885 .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Worktree),
886 }
887 }
888}
889
890#[derive(Debug, Clone, Serialize, Deserialize)]
891#[serde(deny_unknown_fields)]
892pub struct DaemonStatus {
893 pub pid: u32,
894 pub started_at: String,
895 pub build_version: String,
896 pub attached_clients: usize,
897 pub phone_status: WebViewerStatus,
898}
899
900#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
901#[serde(rename_all = "snake_case", tag = "state")]
902pub enum WebViewerStatus {
903 Disabled,
904 Starting,
905 Ready {
906 viewer_url: String,
907 viewer_code: String,
908 qr_login_url: Option<String>,
909 fallback_reason: Option<String>,
910 },
911 Stopped,
912 Error {
913 message: String,
914 },
915}
916
917impl std::fmt::Debug for WebViewerStatus {
918 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
919 match self {
920 Self::Ready {
921 viewer_url,
922 viewer_code,
923 fallback_reason,
924 ..
925 } => formatter
926 .debug_struct("Ready")
927 .field("viewer_url", viewer_url)
928 .field("viewer_code", viewer_code)
929 .field("qr_login_url", &"[redacted]")
930 .field("fallback_reason", fallback_reason)
931 .finish(),
932 Self::Disabled => formatter.write_str("Disabled"),
933 Self::Starting => formatter.write_str("Starting"),
934 Self::Stopped => formatter.write_str("Stopped"),
935 Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
936 }
937 }
938}
939
940impl std::fmt::Display for WebViewerStatus {
941 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
942 match self {
943 Self::Disabled => formatter.write_str("disabled"),
944 Self::Starting => formatter.write_str("starting"),
945 Self::Stopped => formatter.write_str("stopped unexpectedly"),
946 Self::Error { message } => write!(formatter, "error: {message}"),
947 Self::Ready {
948 viewer_url,
949 viewer_code,
950 fallback_reason,
951 ..
952 } => {
953 write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
954 if let Some(reason) = fallback_reason {
955 write!(
956 formatter,
957 "; local only because Tailscale HTTPS is unavailable: {reason}"
958 )?;
959 }
960 Ok(())
961 }
962 }
963 }
964}
965
966#[cfg(target_os = "linux")]
985pub fn process_is_zombie(pid: u32) -> bool {
986 let Ok(stat) = std::fs::read(format!("/proc/{pid}/stat")) else {
987 return false;
988 };
989 stat.iter()
992 .rposition(|byte| *byte == b')')
993 .and_then(|end| stat.get(end + 2))
994 .is_some_and(|state| *state == b'Z')
995}
996
997#[cfg(target_os = "macos")]
1003pub fn process_is_zombie(pid: u32) -> bool {
1004 let Ok(raw_pid) = libc::c_int::try_from(pid) else {
1005 return false;
1006 };
1007 let size = std::mem::size_of::<libc::proc_bsdinfo>();
1008 let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
1010 let written = unsafe {
1012 libc::proc_pidinfo(
1013 raw_pid,
1014 libc::PROC_PIDTBSDINFO,
1015 1,
1016 (&mut info as *mut libc::proc_bsdinfo).cast(),
1017 size as libc::c_int,
1018 )
1019 };
1020 usize::try_from(written).is_ok_and(|written| written == size) && info.pbi_status == libc::SZOMB
1021}
1022
1023#[cfg(all(unix, not(target_os = "linux"), not(target_os = "macos")))]
1024pub fn process_is_zombie(pid: u32) -> bool {
1025 let pid = sysinfo::Pid::from_u32(pid);
1026 let mut system = sysinfo::System::new();
1027 system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
1028 system
1029 .process(pid)
1030 .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
1031}
1032
1033pub async fn wait_for_exit(pid: u32) -> Result<()> {
1039 wait_for_exit_within(pid, STOP_TIMEOUT).await
1040}
1041
1042async fn wait_for_exit_within(pid: u32, timeout: Duration) -> Result<()> {
1043 let deadline = Instant::now() + timeout;
1044 while daemon_process_is_alive(pid) {
1045 ensure!(Instant::now() < deadline, "process {pid} is still running");
1046 tokio::time::sleep(RETRY_DELAY).await;
1047 }
1048 Ok(())
1049}
1050
1051pub fn daemon_process_is_alive(pid: u32) -> bool {
1057 #[cfg(unix)]
1058 {
1059 if pid == 0 {
1060 return false;
1061 }
1062 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
1063 return false;
1064 };
1065 let mut status = 0;
1066 let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
1069 if waited == raw_pid {
1070 return false;
1071 }
1072 if waited == 0 {
1073 return true;
1074 }
1075 let wait_error = std::io::Error::last_os_error();
1076 if wait_error.raw_os_error() != Some(libc::ECHILD) {
1077 return true;
1078 }
1079
1080 #[cfg(target_os = "macos")]
1081 return owned_daemon_group_is_alive(raw_pid);
1082
1083 #[cfg(not(target_os = "macos"))]
1084 process_is_alive(pid)
1085 }
1086 #[cfg(not(unix))]
1087 process_is_alive(pid)
1088}
1089
1090#[cfg(target_os = "macos")]
1091pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
1092 if unsafe { libc::kill(-pid, 0) } == 0 {
1100 return true;
1101 }
1102 let error = std::io::Error::last_os_error();
1103 !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
1104}
1105
1106pub fn process_is_alive(pid: u32) -> bool {
1107 #[cfg(unix)]
1108 {
1109 if pid == 0 {
1110 return false;
1111 }
1112 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
1113 return false;
1114 };
1115 let result = unsafe { libc::kill(raw_pid, 0) };
1118 let exists =
1119 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
1120 exists && !process_is_zombie(pid)
1121 }
1122 #[cfg(not(unix))]
1123 {
1124 let process_id = sysinfo::Pid::from_u32(pid);
1125 let mut system = sysinfo::System::new();
1126 system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[process_id]), true);
1127 system.process(process_id).is_some()
1128 }
1129}
1130
1131pub fn read_metadata() -> Result<DaemonMetadata> {
1132 let metadata = read_metadata_any()?;
1133 ensure!(
1134 metadata.protocol_version == PROTOCOL_VERSION,
1135 "daemon protocol {} is incompatible with client protocol {}",
1136 metadata.protocol_version,
1137 PROTOCOL_VERSION
1138 );
1139 Ok(metadata)
1140}
1141
1142#[derive(Debug, Clone, PartialEq, Eq)]
1148pub struct DaemonNotRunning {
1149 pub metadata_path: std::path::PathBuf,
1151}
1152
1153pub const DAEMON_NOT_RUNNING_MESSAGE: &str = "the Mjolnir daemon is not running";
1155
1156impl std::fmt::Display for DaemonNotRunning {
1157 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1158 formatter.write_str(DAEMON_NOT_RUNNING_MESSAGE)
1159 }
1160}
1161
1162impl std::error::Error for DaemonNotRunning {}
1163
1164pub fn daemon_not_running(error: &anyhow::Error) -> Option<&DaemonNotRunning> {
1166 error
1167 .chain()
1168 .find_map(|cause| cause.downcast_ref::<DaemonNotRunning>())
1169}
1170
1171pub fn read_metadata_any() -> Result<DaemonMetadata> {
1172 read_metadata_at(&metadata_path())
1173}
1174
1175fn read_metadata_at(path: &std::path::Path) -> Result<DaemonMetadata> {
1176 let body = match fs::read(path) {
1177 Ok(body) => body,
1178 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1179 return Err(DaemonNotRunning {
1180 metadata_path: path.to_owned(),
1181 }
1182 .into());
1183 }
1184 Err(error) => return Err(error).with_context(|| format!("read {}", path.display())),
1185 };
1186 let metadata: DaemonMetadata =
1187 serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1188 Ok(metadata)
1189}
1190
1191pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
1192 let body = serde_json::to_vec(value)?;
1193 write_encoded_frame(stream, &body).await
1194}
1195
1196pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
1197 ensure!(
1198 body.len() <= MAX_FRAME_BYTES,
1199 "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
1200 body.len()
1201 );
1202 stream.write_u32(body.len() as u32).await?;
1203 stream.write_all(body).await?;
1204 stream.flush().await?;
1205 Ok(())
1206}
1207
1208pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
1209 let length = stream.read_u32().await? as usize;
1210 ensure!(
1211 length <= MAX_FRAME_BYTES,
1212 "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
1213 );
1214 let mut body = vec![0_u8; length];
1215 stream.read_exact(&mut body).await?;
1216 serde_json::from_slice(&body).context("decode daemon frame")
1217}
1218
1219pub async fn read_response(stream: &mut TcpStream) -> Result<ResponseEnvelope> {
1222 let mut response: ResponseEnvelope = read_frame(stream).await?;
1223 if !matches!(response.result, Ok(DaemonReply::ReplyChunk { .. })) {
1224 return Ok(response);
1225 }
1226 let protocol_version = response.protocol_version;
1227 let request_id = response.request_id;
1228 let mut body = Vec::new();
1229 loop {
1230 ensure!(
1231 response.protocol_version == protocol_version && response.request_id == request_id,
1232 "daemon crossed reply fragment identities"
1233 );
1234 let Ok(DaemonReply::ReplyChunk { bytes, finished }) = response.result else {
1235 bail!("daemon interrupted a chunked reply");
1236 };
1237 ensure!(!bytes.is_empty(), "empty reply fragment");
1238 body.extend(bytes);
1239 if finished {
1240 break;
1241 }
1242 response = read_frame(stream).await?;
1243 }
1244 let reply: DaemonReply = tokio::task::spawn_blocking(move || serde_json::from_slice(&body))
1245 .await
1246 .context("reply decoder task failed")??;
1247 ensure!(
1248 reply.is_chunked(),
1249 "daemon chunked a reply that is never chunked"
1250 );
1251 Ok(ResponseEnvelope {
1252 protocol_version,
1253 request_id,
1254 result: Ok(reply),
1255 })
1256}
1257
1258pub struct DaemonClient {
1259 metadata: DaemonMetadata,
1260 stream: TcpStream,
1261 next_request_id: u64,
1262}
1263
1264impl DaemonClient {
1265 pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
1266 let stream =
1267 tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
1268 .await
1269 .context("time out connecting to Mjolnir daemon")??;
1270 Ok(Self {
1271 metadata,
1272 stream,
1273 next_request_id: 1,
1274 })
1275 }
1276
1277 pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1281 self.request_with_reconnect(action, || async {
1282 loop {
1283 if let Ok(client) = connect_existing().await {
1284 return Ok(client);
1285 }
1286 tokio::time::sleep(Duration::from_millis(250)).await;
1287 }
1288 })
1289 .await
1290 }
1291
1292 async fn request_with_reconnect<F, Fut>(
1293 &mut self,
1294 action: DaemonAction,
1295 mut reconnect: F,
1296 ) -> Result<DaemonReply>
1297 where
1298 F: FnMut() -> Fut,
1299 Fut: Future<Output = Result<Self>>,
1300 {
1301 loop {
1302 let response = self.request_once(action.clone()).await?;
1303 if matches!(response, DaemonReply::UpgradePending)
1304 && !matches!(action, DaemonAction::PrepareUpgrade)
1305 {
1306 tokio::time::sleep(Duration::from_millis(100)).await;
1309 *self = reconnect().await?;
1310 } else {
1311 return Ok(response);
1312 }
1313 }
1314 }
1315
1316 async fn request_once(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1317 let protocol_version = self.metadata.protocol_version;
1318 let request_id = self.next_request_id;
1319 self.next_request_id += 1;
1320 write_frame(
1321 &mut self.stream,
1322 &RequestEnvelope {
1323 protocol_version,
1324 request_id,
1325 token: self.metadata.token.clone(),
1326 action,
1327 },
1328 )
1329 .await?;
1330 let response = read_response(&mut self.stream).await?;
1331 ensure!(
1332 response.protocol_version == protocol_version,
1333 "daemon changed protocol"
1334 );
1335 ensure!(
1336 response.request_id == request_id,
1337 "daemon crossed request IDs"
1338 );
1339 response
1340 .result
1341 .map_err(|message| anyhow::Error::new(DaemonRefusal(message)))
1342 }
1343
1344 pub async fn status(&mut self) -> Result<DaemonStatus> {
1345 match self.request(DaemonAction::Status).await? {
1346 DaemonReply::Status(status) => Ok(status),
1347 reply => bail!("unexpected daemon status reply {reply:?}"),
1348 }
1349 }
1350
1351 pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
1352 match self.request(DaemonAction::WebViewerAccess).await? {
1353 DaemonReply::WebViewerAccess(access) => Ok(access),
1354 reply => bail!("unexpected web viewer reply {reply:?}"),
1355 }
1356 }
1357
1358 pub async fn subagent_options(
1359 &mut self,
1360 profile: String,
1361 model: Option<String>,
1362 config: Option<Config>,
1363 ) -> Result<mj_core::subagent::SubagentOptions> {
1364 match self
1365 .request(DaemonAction::SubagentOptions {
1366 profile,
1367 model,
1368 config: config.map(Box::new),
1369 })
1370 .await?
1371 {
1372 DaemonReply::SubagentOptions(options) => Ok(options),
1373 reply => bail!("unexpected subagent options reply {reply:?}"),
1374 }
1375 }
1376
1377 pub async fn recover_web_viewer(
1378 &mut self,
1379 action: crate::web::WebViewerRecovery,
1380 ) -> Result<()> {
1381 match self.request(DaemonAction::RecoverWebViewer(action)).await? {
1382 DaemonReply::Done => Ok(()),
1383 reply => bail!("unexpected web viewer recovery reply {reply:?}"),
1384 }
1385 }
1386
1387 pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
1388 match self.request(DaemonAction::InspectWebListener).await? {
1389 DaemonReply::WebListeners(processes) => Ok(processes),
1390 reply => bail!("unexpected listener inspection reply {reply:?}"),
1391 }
1392 }
1393
1394 pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
1395 match self.request(DaemonAction::ListWorkspaces).await? {
1396 DaemonReply::Workspaces(workspaces) => Ok(workspaces),
1397 reply => bail!("unexpected daemon workspace reply {reply:?}"),
1398 }
1399 }
1400
1401 pub async fn refresh_quota(&mut self) -> Result<()> {
1404 match self.request(DaemonAction::RefreshQuota).await? {
1405 DaemonReply::Done => Ok(()),
1406 reply => bail!("unexpected refresh-quota reply {reply:?}"),
1407 }
1408 }
1409
1410 pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
1411 match self
1412 .request(DaemonAction::RenameProfile { old_id, new_id })
1413 .await?
1414 {
1415 DaemonReply::Done => Ok(()),
1416 reply => bail!("unexpected rename-profile reply {reply:?}"),
1417 }
1418 }
1419
1420 pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1421 match self
1422 .request(DaemonAction::RenameTarget { old_id, new_id })
1423 .await?
1424 {
1425 DaemonReply::Done => Ok(()),
1426 reply => bail!("unexpected rename-target reply {reply:?}"),
1427 }
1428 }
1429
1430 pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1431 match self.request(DaemonAction::CreateWorkspace { name }).await? {
1432 DaemonReply::Workspace(workspace) => Ok(workspace),
1433 reply => bail!("unexpected create-workspace reply {reply:?}"),
1434 }
1435 }
1436
1437 pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1438 match self
1439 .request(DaemonAction::RenameWorkspace { workspace_id, name })
1440 .await?
1441 {
1442 DaemonReply::Done => Ok(()),
1443 reply => bail!("unexpected rename-workspace reply {reply:?}"),
1444 }
1445 }
1446
1447 pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1448 match self
1449 .request(DaemonAction::TouchWorkspace { workspace_id })
1450 .await?
1451 {
1452 DaemonReply::Done => Ok(()),
1453 reply => bail!("unexpected touch-workspace reply {reply:?}"),
1454 }
1455 }
1456
1457 pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1458 match self
1459 .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1460 .await?
1461 {
1462 DaemonReply::Done => Ok(()),
1463 reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1464 }
1465 }
1466
1467 pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1468 match self
1469 .request(DaemonAction::CloseWorkspace { workspace_id })
1470 .await?
1471 {
1472 DaemonReply::Done => Ok(()),
1473 reply => bail!("unexpected close workspace reply: {reply:?}"),
1474 }
1475 }
1476
1477 pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1478 match self
1479 .request(DaemonAction::DeleteWorkspace { workspace_id })
1480 .await?
1481 {
1482 DaemonReply::Done => Ok(()),
1483 reply => bail!("unexpected delete-workspace reply {reply:?}"),
1484 }
1485 }
1486
1487 pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1488 match self
1489 .request(DaemonAction::Attach { client_id, pid })
1490 .await?
1491 {
1492 DaemonReply::Done => Ok(()),
1493 reply => bail!("unexpected attach reply {reply:?}"),
1494 }
1495 }
1496
1497 pub async fn detach(&mut self, client_id: String) -> Result<()> {
1498 match self.request(DaemonAction::Detach { client_id }).await? {
1499 DaemonReply::Done => Ok(()),
1500 reply => bail!("unexpected detach reply {reply:?}"),
1501 }
1502 }
1503
1504 pub async fn persist_read_receipt(
1505 &mut self,
1506 client_id: String,
1507 workspace_id: String,
1508 session_id: String,
1509 through: u64,
1510 ) -> Result<u64> {
1511 match self
1512 .request(DaemonAction::PersistReadReceipt {
1513 client_id,
1514 workspace_id,
1515 session_id,
1516 through,
1517 })
1518 .await?
1519 {
1520 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1521 reply => bail!("unexpected read-receipt reply {reply:?}"),
1522 }
1523 }
1524
1525 pub async fn persist_detached_session_state(
1526 &mut self,
1527 client_id: String,
1528 workspace_id: String,
1529 session_id: String,
1530 through: u64,
1531 owner_pid: u32,
1532 draft: mj_core::storage::DetachedSessionDraft,
1533 ) -> Result<()> {
1534 match self
1535 .request(DaemonAction::PersistDetachedSessionState {
1536 client_id,
1537 workspace_id,
1538 session_id,
1539 through,
1540 owner_pid,
1541 draft,
1542 })
1543 .await?
1544 {
1545 DaemonReply::Done => Ok(()),
1546 reply => bail!("unexpected detached-session-state reply {reply:?}"),
1547 }
1548 }
1549
1550 pub async fn save_active_review(
1551 &mut self,
1552 session_id: String,
1553 review: mj_core::storage::StoredReview,
1554 ) -> Result<()> {
1555 match self
1556 .request(DaemonAction::SaveActiveReview { session_id, review })
1557 .await?
1558 {
1559 DaemonReply::Done => Ok(()),
1560 reply => bail!("unexpected save-review reply {reply:?}"),
1561 }
1562 }
1563
1564 pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1565 match self
1566 .request(DaemonAction::ClearActiveReview { session_id })
1567 .await?
1568 {
1569 DaemonReply::Done => Ok(()),
1570 reply => bail!("unexpected clear-review reply {reply:?}"),
1571 }
1572 }
1573
1574 pub async fn save_workspace_pane_sizes(
1575 &mut self,
1576 workspace_id: String,
1577 sizes: mj_core::workspace::PaneSizes,
1578 ) -> Result<()> {
1579 match self
1580 .request(DaemonAction::SaveWorkspacePaneSizes {
1581 workspace_id,
1582 sizes,
1583 })
1584 .await?
1585 {
1586 DaemonReply::Done => Ok(()),
1587 reply => bail!("unexpected pane-size save reply {reply:?}"),
1588 }
1589 }
1590
1591 pub async fn save_workspace_layout(
1592 &mut self,
1593 workspace_id: String,
1594 layout: mj_core::workspace::ConversationLayout,
1595 ) -> Result<()> {
1596 match self
1597 .request(DaemonAction::SaveWorkspaceLayout {
1598 workspace_id,
1599 layout,
1600 })
1601 .await?
1602 {
1603 DaemonReply::Done => Ok(()),
1604 reply => bail!("unexpected layout save reply {reply:?}"),
1605 }
1606 }
1607
1608 pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1609 match self
1610 .request(DaemonAction::PersistImportedSession {
1611 session: Box::new(session),
1612 })
1613 .await?
1614 {
1615 DaemonReply::Done => Ok(()),
1616 reply => bail!("unexpected imported-session reply {reply:?}"),
1617 }
1618 }
1619
1620 pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1621 match self
1622 .request(DaemonAction::SetSessionTitle { session_id, title })
1623 .await?
1624 {
1625 DaemonReply::Text(title) => Ok(title),
1626 reply => bail!("unexpected session-title reply {reply:?}"),
1627 }
1628 }
1629
1630 pub async fn set_session_workspace(
1631 &mut self,
1632 session_id: String,
1633 workspace_id: String,
1634 ) -> Result<()> {
1635 match self
1636 .request(DaemonAction::SetSessionWorkspace {
1637 session_id,
1638 workspace_id,
1639 })
1640 .await?
1641 {
1642 DaemonReply::Done => Ok(()),
1643 reply => bail!("unexpected session-workspace reply {reply:?}"),
1644 }
1645 }
1646
1647 pub async fn set_session_container_settings(
1648 &mut self,
1649 session_id: String,
1650 cpus: Option<String>,
1651 memory: Option<String>,
1652 mounts: Vec<AdditionalMount>,
1653 mount_history: Vec<PathBuf>,
1654 ) -> Result<()> {
1655 match self
1656 .request(DaemonAction::SetSessionContainerSettings {
1657 session_id,
1658 cpus,
1659 memory,
1660 mounts,
1661 mount_history,
1662 })
1663 .await?
1664 {
1665 DaemonReply::Done => Ok(()),
1666 reply => bail!("unexpected container-settings reply {reply:?}"),
1667 }
1668 }
1669
1670 pub async fn set_session_acp_title(
1671 &mut self,
1672 session_id: String,
1673 title: Option<String>,
1674 ) -> Result<()> {
1675 match self
1676 .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1677 .await?
1678 {
1679 DaemonReply::Done => Ok(()),
1680 reply => bail!("unexpected ACP-title reply {reply:?}"),
1681 }
1682 }
1683
1684 pub async fn mark_session_target_missing(
1685 &mut self,
1686 session_id: String,
1687 detail: String,
1688 updated_at: String,
1689 ) -> Result<Option<SessionState>> {
1690 match self
1691 .request(DaemonAction::MarkSessionTargetMissing {
1692 session_id,
1693 detail,
1694 updated_at,
1695 })
1696 .await?
1697 {
1698 DaemonReply::OptionalSessionState(state) => Ok(state),
1699 reply => bail!("unexpected target-missing reply {reply:?}"),
1700 }
1701 }
1702
1703 pub async fn checkpoint_session(
1704 &mut self,
1705 session_id: String,
1706 ) -> Result<mj_core::state::CheckpointMetadata> {
1707 match self
1708 .request(DaemonAction::CheckpointSession { session_id })
1709 .await?
1710 {
1711 DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1712 reply => bail!("unexpected checkpoint reply {reply:?}"),
1713 }
1714 }
1715
1716 pub async fn project_catalog(
1721 &mut self,
1722 refresh: bool,
1723 retry: bool,
1724 ) -> Result<mj_core::project_catalog::ProjectCatalogView> {
1725 match self
1726 .request(DaemonAction::ProjectCatalog { refresh, retry })
1727 .await?
1728 {
1729 DaemonReply::ProjectCatalog(view) => Ok(view),
1730 other => bail!("unexpected project catalog reply: {other:?}"),
1731 }
1732 }
1733
1734 pub async fn create_project(
1735 &mut self,
1736 sources: Vec<String>,
1737 ) -> Result<mj_core::project_catalog::SavedProject> {
1738 match self
1739 .request(DaemonAction::CreateProject { sources })
1740 .await?
1741 {
1742 DaemonReply::ProjectCreated(project) => Ok(project),
1743 other => bail!("unexpected project creation reply: {other:?}"),
1744 }
1745 }
1746
1747 pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1748 match self
1749 .request(DaemonAction::WikiSearch { query, limit })
1750 .await?
1751 {
1752 DaemonReply::WikiRows(page) => Ok(page),
1753 reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1754 }
1755 }
1756
1757 pub async fn session_text_search(&mut self, query: String) -> Result<Vec<SessionTextMatch>> {
1761 match self
1762 .request(DaemonAction::SessionTextSearch { query })
1763 .await?
1764 {
1765 DaemonReply::SessionTextMatches(matches) => Ok(matches),
1766 reply => bail!("unexpected session text search reply {reply:?}"),
1767 }
1768 }
1769
1770 pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1772 match self
1773 .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1774 .await?
1775 {
1776 DaemonReply::Text(markdown) => Ok(markdown),
1777 reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1778 }
1779 }
1780
1781 pub async fn wiki_hits(
1786 &mut self,
1787 wiki_id: String,
1788 query: String,
1789 context_messages: usize,
1790 per_message_chars: usize,
1791 ) -> Result<Option<WikiHitTranscript>> {
1792 match self
1793 .request(DaemonAction::WikiHits {
1794 wiki_id,
1795 query,
1796 context_messages,
1797 per_message_chars,
1798 })
1799 .await?
1800 {
1801 DaemonReply::WikiHits(transcript) => Ok(transcript),
1802 reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1803 }
1804 }
1805
1806 pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1809 match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1810 DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1811 reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1812 }
1813 }
1814
1815 pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1819 match self.request(DaemonAction::WikiRestore(request)).await? {
1820 DaemonReply::RegisteredSession(registered) => Ok(*registered),
1821 reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1822 }
1823 }
1824
1825 pub async fn scan_recovery(
1826 &mut self,
1827 all_instances: bool,
1828 ) -> Result<mj_core::state::RecoveryScan> {
1829 match self
1830 .request(DaemonAction::ScanRecovery { all_instances })
1831 .await?
1832 {
1833 DaemonReply::RecoveryScan(scan) => Ok(scan),
1834 reply => bail!("unexpected recovery-scan reply {reply:?}"),
1835 }
1836 }
1837
1838 pub async fn adopt_recovery(
1839 &mut self,
1840 session_id: String,
1841 target_id: String,
1842 profile: Option<String>,
1843 bundle: Option<String>,
1844 all_instances: bool,
1845 ) -> Result<()> {
1846 match self
1847 .request(DaemonAction::AdoptRecovery {
1848 session_id,
1849 target_id,
1850 profile,
1851 bundle,
1852 all_instances,
1853 })
1854 .await?
1855 {
1856 DaemonReply::Done => Ok(()),
1857 reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1858 }
1859 }
1860
1861 pub async fn destroy_recovery(
1862 &mut self,
1863 session_id: String,
1864 target_id: String,
1865 confirmation: String,
1866 all_instances: bool,
1867 ) -> Result<()> {
1868 match self
1869 .request(DaemonAction::DestroyRecovery {
1870 session_id,
1871 target_id,
1872 confirmation,
1873 all_instances,
1874 })
1875 .await?
1876 {
1877 DaemonReply::Done => Ok(()),
1878 reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1879 }
1880 }
1881
1882 pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1883 match self
1884 .request(DaemonAction::Snapshot { workspace_id })
1885 .await?
1886 {
1887 DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1888 reply => bail!("unexpected snapshot reply {reply:?}"),
1889 }
1890 }
1891
1892 pub async fn runtime_changes(
1893 &mut self,
1894 cursor: Option<crate::runtime_feed::RuntimeCursor>,
1895 wait: bool,
1896 ) -> Result<crate::runtime_feed::RuntimeFrame> {
1897 match self
1898 .request(DaemonAction::RuntimeChanges { cursor, wait })
1899 .await?
1900 {
1901 DaemonReply::RuntimeChanges(frame) => Ok(*frame),
1902 reply => bail!("unexpected runtime changes reply {reply:?}"),
1903 }
1904 }
1905
1906 pub async fn resume_candidates(&mut self) -> Result<ResumeCandidates> {
1907 match self.request(DaemonAction::ResumeCandidates).await? {
1908 DaemonReply::ResumeCandidates(candidates) => Ok(*candidates),
1909 reply => bail!("unexpected resume candidates reply {reply:?}"),
1910 }
1911 }
1912
1913 pub async fn session_record(&mut self, session_id: String) -> Result<Option<SessionRecord>> {
1914 match self
1915 .request(DaemonAction::SessionRecord { session_id })
1916 .await?
1917 {
1918 DaemonReply::SessionRecord(record) => Ok(record.map(|record| *record)),
1919 reply => bail!("unexpected session record reply {reply:?}"),
1920 }
1921 }
1922
1923 pub async fn go_startup_session(
1924 &mut self,
1925 workspace_id: String,
1926 last_session_id: Option<String>,
1927 ) -> Result<Option<SessionRecord>> {
1928 match self
1929 .request(DaemonAction::GoStartupSession {
1930 workspace_id,
1931 last_session_id,
1932 })
1933 .await?
1934 {
1935 DaemonReply::GoStartupSession(record) => Ok(record.map(|record| *record)),
1936 reply => bail!("unexpected go startup session reply {reply:?}"),
1937 }
1938 }
1939
1940 pub async fn submit_session_command(
1941 &mut self,
1942 session_id: String,
1943 command_id: String,
1944 command: RelayCommand,
1945 inherited_draft: Option<String>,
1946 ) -> Result<u64> {
1947 match self
1948 .request(DaemonAction::SubmitSessionCommand {
1949 inherited_draft,
1950 session_id,
1951 command_id,
1952 command,
1953 })
1954 .await?
1955 {
1956 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1957 reply => bail!("unexpected session command reply {reply:?}"),
1958 }
1959 }
1960
1961 pub async fn queue_startup_prompt(
1965 &mut self,
1966 session_id: String,
1967 text: String,
1968 inherited_draft: Option<String>,
1969 ) -> Result<()> {
1970 match self
1971 .request(DaemonAction::QueueStartupPrompt {
1972 session_id,
1973 text,
1974 inherited_draft,
1975 })
1976 .await?
1977 {
1978 DaemonReply::Done => Ok(()),
1979 reply => bail!("unexpected startup prompt reply {reply:?}"),
1980 }
1981 }
1982
1983 pub async fn withdraw_startup_prompt(
1988 &mut self,
1989 session_id: String,
1990 text: String,
1991 ) -> Result<bool> {
1992 match self
1993 .request(DaemonAction::WithdrawStartupPrompt { session_id, text })
1994 .await?
1995 {
1996 DaemonReply::PromptWithdrawn(withdrawn) => Ok(withdrawn),
1997 reply => bail!("unexpected startup prompt withdrawal reply {reply:?}"),
1998 }
1999 }
2000
2001 pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
2007 match self
2008 .request(DaemonAction::StartTurnReview { session_id })
2009 .await?
2010 {
2011 DaemonReply::Done => Ok(()),
2012 reply => bail!("unexpected turn-review reply {reply:?}"),
2013 }
2014 }
2015
2016 pub async fn resolve_turn_review(
2017 &mut self,
2018 session_id: String,
2019 resolution: Resolution,
2020 ) -> Result<()> {
2021 match self
2022 .request(DaemonAction::ResolveTurnReview {
2023 session_id,
2024 resolution,
2025 })
2026 .await?
2027 {
2028 DaemonReply::Done => Ok(()),
2029 reply => bail!("unexpected turn-review resolution reply {reply:?}"),
2030 }
2031 }
2032
2033 pub async fn reviewer_action(
2034 &mut self,
2035 session_id: String,
2036 role: Option<String>,
2037 action: crate::session::ReviewerAction,
2038 ) -> Result<crate::session::ReviewerOutcome> {
2039 match self
2040 .request(DaemonAction::ReviewerAction {
2041 session_id,
2042 role,
2043 action,
2044 })
2045 .await?
2046 {
2047 DaemonReply::Reviewer(outcome) => Ok(*outcome),
2048 reply => bail!("unexpected reviewer reply {reply:?}"),
2049 }
2050 }
2051
2052 pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
2053 match self
2054 .request(DaemonAction::SyncSession { session_id })
2055 .await?
2056 {
2057 DaemonReply::Done => Ok(()),
2058 reply => bail!("unexpected session sync reply {reply:?}"),
2059 }
2060 }
2061
2062 pub async fn respond_elicitation(
2063 &mut self,
2064 session_id: String,
2065 elicitation_id: String,
2066 response: ElicitationResponse,
2067 ) -> Result<()> {
2068 match self
2069 .request(DaemonAction::RespondElicitation {
2070 session_id,
2071 elicitation_id,
2072 response,
2073 })
2074 .await?
2075 {
2076 DaemonReply::Done => Ok(()),
2077 reply => bail!("unexpected elicitation reply {reply:?}"),
2078 }
2079 }
2080
2081 pub async fn native_agent_history(
2082 &mut self,
2083 owner: String,
2084 child: String,
2085 before: Option<(u64, String)>,
2086 ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
2087 match self
2088 .request(DaemonAction::NativeAgentHistory {
2089 owner,
2090 child,
2091 before,
2092 })
2093 .await?
2094 {
2095 DaemonReply::NativeAgentHistory(page) => Ok(page),
2096 reply => bail!("unexpected native agent history reply {reply:?}"),
2097 }
2098 }
2099
2100 pub async fn stop_background_task(
2101 &mut self,
2102 session_id: String,
2103 background_task_id: String,
2104 ) -> Result<()> {
2105 match self
2106 .request(DaemonAction::StopBackgroundTask {
2107 session_id,
2108 background_task_id,
2109 })
2110 .await?
2111 {
2112 DaemonReply::Done => Ok(()),
2113 reply => bail!("unexpected background task stop reply {reply:?}"),
2114 }
2115 }
2116
2117 pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
2118 self.suspend_session_with_ack(session_id, false).await
2119 }
2120
2121 pub async fn suspend_session_with_ack(
2122 &mut self,
2123 session_id: String,
2124 acknowledge_unpublished_work: bool,
2125 ) -> Result<()> {
2126 match self
2127 .request(DaemonAction::SuspendSession {
2128 session_id,
2129 acknowledge_unpublished_work,
2130 })
2131 .await?
2132 {
2133 DaemonReply::Done => Ok(()),
2134 reply => bail!("unexpected close-session reply {reply:?}"),
2135 }
2136 }
2137
2138 pub async fn start_create_session(
2139 &mut self,
2140 request: CreateSessionRequest,
2141 ) -> Result<RegisteredSession> {
2142 match self
2143 .request(DaemonAction::StartCreateSession(request))
2144 .await?
2145 {
2146 DaemonReply::RegisteredSession(registered) => Ok(*registered),
2147 reply => bail!("unexpected start-create reply {reply:?}"),
2148 }
2149 }
2150
2151 pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
2152 match self
2153 .request(DaemonAction::WaitCreateSession { session_id })
2154 .await?
2155 {
2156 DaemonReply::Done => Ok(()),
2157 reply => bail!("unexpected wait-create reply {reply:?}"),
2158 }
2159 }
2160
2161 pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
2162 match self.request(DaemonAction::ResumeSession(request)).await? {
2163 DaemonReply::Done => Ok(()),
2164 reply => bail!("unexpected resume-session reply {reply:?}"),
2165 }
2166 }
2167
2168 pub async fn discard_since_checkpoint(
2169 &mut self,
2170 session_id: String,
2171 checkpoint: mj_core::state::CheckpointMetadata,
2172 ) -> Result<()> {
2173 match self
2174 .request(DaemonAction::DiscardSinceCheckpoint {
2175 session_id,
2176 checkpoint,
2177 })
2178 .await?
2179 {
2180 DaemonReply::Done => Ok(()),
2181 reply => bail!("unexpected force-stop reply {reply:?}"),
2182 }
2183 }
2184
2185 pub async fn destroy_stopped_session(
2186 &mut self,
2187 session_id: String,
2188 delete_branch: bool,
2189 ) -> Result<()> {
2190 match self
2191 .request(DaemonAction::DestroyStoppedSession {
2192 session_id,
2193 delete_branch,
2194 })
2195 .await?
2196 {
2197 DaemonReply::Done => Ok(()),
2198 reply => bail!("unexpected destroy-stopped reply {reply:?}"),
2199 }
2200 }
2201
2202 pub async fn force_destroy_session(
2203 &mut self,
2204 session_id: String,
2205 delete_branch: bool,
2206 ) -> Result<()> {
2207 match self
2208 .request(DaemonAction::ForceDestroySession {
2209 session_id,
2210 delete_branch,
2211 })
2212 .await?
2213 {
2214 DaemonReply::Done => Ok(()),
2215 reply => bail!("unexpected force-destroy reply {reply:?}"),
2216 }
2217 }
2218
2219 pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
2220 match self
2221 .request(DaemonAction::CancelLifecycle { session_id })
2222 .await?
2223 {
2224 DaemonReply::Done => Ok(()),
2225 reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
2226 }
2227 }
2228
2229 pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
2230 match self
2231 .request(DaemonAction::RecoverDraft { draft_id })
2232 .await?
2233 {
2234 DaemonReply::Done => Ok(()),
2235 reply => bail!("unexpected recover-draft reply {reply:?}"),
2236 }
2237 }
2238
2239 pub async fn stop(&mut self) -> Result<()> {
2240 match self.request(DaemonAction::Stop).await? {
2241 DaemonReply::Done => Ok(()),
2242 reply => bail!("unexpected stop reply {reply:?}"),
2243 }
2244 }
2245}
2246
2247pub async fn connect_existing() -> Result<DaemonClient> {
2248 let metadata = tokio::task::spawn_blocking(read_metadata)
2249 .await
2250 .context("read daemon metadata task failed")??;
2251 DaemonClient::connect(metadata).await
2252}
2253
2254pub struct ManagementClient {
2258 inner: DaemonClient,
2259}
2260
2261impl ManagementClient {
2262 pub fn new(inner: DaemonClient) -> Self {
2263 Self { inner }
2264 }
2265 pub fn protocol_version(&self) -> u32 {
2266 self.inner.metadata.protocol_version
2267 }
2268
2269 pub async fn status(&mut self) -> Result<DaemonStatus> {
2270 self.inner.status().await
2271 }
2272
2273 pub async fn stop(&mut self) -> Result<()> {
2274 tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
2275 .await
2276 .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
2277 }
2278
2279 pub async fn stop_and_wait(mut self) -> Result<()> {
2281 let pid = self.inner.metadata.pid;
2282 self.stop().await?;
2283 wait_for_exit_within(pid, STOP_DRAIN_TIMEOUT)
2284 .await
2285 .with_context(|| {
2286 format!(
2287 "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
2288 STOP_DRAIN_TIMEOUT.as_secs()
2289 )
2290 })
2291 }
2292}
2293
2294pub async fn connect_management() -> Result<ManagementClient> {
2295 Ok(ManagementClient {
2296 inner: DaemonClient::connect(read_metadata_any()?).await?,
2297 })
2298}
2299
2300pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
2308 ensure!(
2309 metadata.protocol_version <= PROTOCOL_VERSION,
2310 "{}",
2311 unsupported_daemon_protocol_message(
2312 metadata.protocol_version,
2313 &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
2314 )
2315 );
2316 Ok(())
2317}
2318
2319fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
2322 format!(
2323 "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
2324 Put the daemon's directory first on PATH, or reinstall this client from that build."
2325 )
2326}
2327pub const PROTOCOL_VERSION: u32 = 51;
2329pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
2330pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
2340pub const STOP_DRAIN_TIMEOUT: Duration = Duration::from_secs(90);
2343pub const RETRY_DELAY: Duration = Duration::from_millis(40);
2344impl DaemonClient {
2345 pub async fn prepare_move_session(
2346 &mut self,
2347 selection: MoveSelection,
2348 ) -> Result<MovePreparation> {
2349 match self
2350 .request(DaemonAction::PrepareMoveSession(selection))
2351 .await?
2352 {
2353 DaemonReply::MovePreparation(preparation) => Ok(*preparation),
2354 _ => bail!("daemon returned an unexpected move preparation reply"),
2355 }
2356 }
2357
2358 pub async fn move_sources(
2359 &mut self,
2360 session_id: String,
2361 cleanup_operation_id: Option<String>,
2362 ) -> Result<Vec<mj_core::move_workspace::RetainedMoveSource>> {
2363 match self
2364 .request(DaemonAction::MoveSources {
2365 session_id,
2366 cleanup_operation_id,
2367 })
2368 .await?
2369 {
2370 DaemonReply::MoveSources(sources) => Ok(sources),
2371 reply => bail!("unexpected Move sources reply: {reply:?}"),
2372 }
2373 }
2374
2375 pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
2376 match self.request(DaemonAction::MoveSession(request)).await? {
2377 DaemonReply::MoveOutcome(outcome) => Ok(outcome),
2378 _ => bail!("daemon returned an unexpected move reply"),
2379 }
2380 }
2381}
2382
2383#[cfg(test)]
2384mod tests {
2385 use super::*;
2386 use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
2387 use std::path::Path;
2388
2389 #[cfg(unix)]
2390 #[test]
2391 fn an_exited_unreaped_process_is_a_zombie_and_not_alive() {
2392 use std::time::{Duration, Instant};
2396
2397 let mut running = std::process::Command::new("sleep")
2398 .arg("30")
2399 .spawn()
2400 .unwrap();
2401 assert!(process_is_alive(running.id()));
2402 assert!(!process_is_zombie(running.id()));
2403 running.kill().unwrap();
2404 running.wait().unwrap();
2405
2406 let mut exited = std::process::Command::new("sh")
2407 .args(["-c", "exit 0"])
2408 .spawn()
2409 .unwrap();
2410 let pid = exited.id();
2411 let deadline = Instant::now() + Duration::from_secs(10);
2412 while !process_is_zombie(pid) {
2413 assert!(
2414 Instant::now() < deadline,
2415 "an exited, unreaped child was never reported as a zombie"
2416 );
2417 std::thread::sleep(Duration::from_millis(20));
2418 }
2419 assert!(!process_is_alive(pid));
2420 exited.wait().unwrap();
2421 }
2422
2423 #[test]
2424 fn a_missing_endpoint_file_reads_as_a_stopped_daemon() {
2425 let directory = tempfile::tempdir().unwrap();
2426 let path = directory.path().join("daemon.json");
2427 let error = read_metadata_at(&path).unwrap_err();
2428 let stopped = daemon_not_running(&error).expect("stopped, not an I/O failure");
2429 assert_eq!(stopped.metadata_path, path);
2430 assert_eq!(format!("{error:#}"), "the Mjolnir daemon is not running");
2431
2432 std::fs::write(&path, b"not json").unwrap();
2433 let error = read_metadata_at(&path).unwrap_err();
2434 assert!(daemon_not_running(&error).is_none());
2435 }
2436
2437 #[tokio::test]
2438 async fn subagent_discovery_uses_the_daemon_and_preserves_choices_and_failures() {
2439 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2440 let metadata = DaemonMetadata {
2441 protocol_version: PROTOCOL_VERSION,
2442 pid: std::process::id(),
2443 address: listener.local_addr().unwrap(),
2444 token: "subagent-discovery-test".into(),
2445 started_at: "test".into(),
2446 build_version: env!("CARGO_PKG_VERSION").into(),
2447 };
2448 let options: mj_core::subagent::SubagentOptions =
2449 serde_json::from_value(serde_json::json!({
2450 "models": [{"value": "gpt-6-luna", "name": "Luna", "description": null}],
2451 "efforts": [{"value": "high", "name": "High", "description": null}],
2452 "unavailable": ["another-profile: authentication failed"]
2453 }))
2454 .unwrap();
2455 let expected = serde_json::to_value(&options).unwrap();
2456 let mut draft = Config::default();
2457 draft.profiles.insert(
2458 "codex".into(),
2459 serde_json::from_value(serde_json::json!({
2460 "kind": "codex", "home": "/unsaved/profile",
2461 "subagents": {"mode": "single_model", "model": "gpt-6-luna", "effort": "high"}
2462 }))
2463 .unwrap(),
2464 );
2465 let expected_draft = serde_json::to_value(&draft).unwrap();
2466 let server = tokio::spawn(async move {
2467 let (mut stream, _) = listener.accept().await.unwrap();
2468 for (index, model) in [
2469 None,
2470 Some("gpt-6-luna".to_owned()),
2471 Some("gpt-6-luna".to_owned()),
2472 ]
2473 .into_iter()
2474 .enumerate()
2475 {
2476 let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2477 assert_eq!(request.token, "subagent-discovery-test");
2478 let DaemonAction::SubagentOptions {
2479 profile,
2480 model: requested,
2481 config,
2482 } = request.action
2483 else {
2484 panic!("unexpected discovery request")
2485 };
2486 assert_eq!(profile, "codex");
2487 assert_eq!(requested, model);
2488 assert_eq!(
2489 serde_json::to_value(config).unwrap(),
2490 if index == 1 {
2491 expected_draft.clone()
2492 } else {
2493 serde_json::Value::Null
2494 }
2495 );
2496 write_frame(
2497 &mut stream,
2498 &ResponseEnvelope {
2499 protocol_version: request.protocol_version,
2500 request_id: request.request_id,
2501 result: if index != 2 {
2502 Ok(DaemonReply::SubagentOptions(options.clone()))
2503 } else {
2504 Err("profile discovery cancelled".into())
2505 },
2506 },
2507 )
2508 .await
2509 .unwrap();
2510 }
2511 });
2512 let mut client = DaemonClient::connect(metadata).await.unwrap();
2513 let discovered = client
2514 .subagent_options("codex".into(), None, None)
2515 .await
2516 .unwrap();
2517 assert_eq!(serde_json::to_value(discovered).unwrap(), expected);
2518 let discovered = client
2519 .subagent_options("codex".into(), Some("gpt-6-luna".into()), Some(draft))
2520 .await
2521 .unwrap();
2522 assert_eq!(serde_json::to_value(discovered).unwrap(), expected);
2523 let error = client
2524 .subagent_options("codex".into(), Some("gpt-6-luna".into()), None)
2525 .await
2526 .unwrap_err();
2527 assert_eq!(error.to_string(), "profile discovery cancelled");
2528 server.await.unwrap();
2529 }
2530
2531 #[tokio::test]
2532 async fn upgrade_refusal_retries_the_identical_command_but_lost_acknowledgements_do_not() {
2533 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2534 let metadata = DaemonMetadata {
2535 protocol_version: PROTOCOL_VERSION,
2536 pid: std::process::id(),
2537 address: listener.local_addr().unwrap(),
2538 token: "test-token".into(),
2539 started_at: "test".into(),
2540 build_version: env!("CARGO_PKG_VERSION").into(),
2541 };
2542 let action = DaemonAction::SubmitSessionCommand {
2543 inherited_draft: None,
2544 session_id: "test-session".into(),
2545 command_id: "steer-command".into(),
2546 command: RelayCommand::Steer {
2547 active_prompt_id: "active-command".into(),
2548 queued_prompt_id: "queued-command".into(),
2549 },
2550 };
2551 let expected = serde_json::to_value(&action).unwrap();
2552 let server = tokio::spawn(async move {
2553 for reply in [
2554 Some(DaemonReply::UpgradePending),
2555 Some(DaemonReply::Done),
2556 None,
2557 ] {
2558 let (mut stream, _) = listener.accept().await.unwrap();
2559 let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2560 assert_eq!(serde_json::to_value(&request.action).unwrap(), expected);
2561 if let Some(reply) = reply {
2562 write_frame(
2563 &mut stream,
2564 &ResponseEnvelope {
2565 protocol_version: request.protocol_version,
2566 request_id: request.request_id,
2567 result: Ok(reply),
2568 },
2569 )
2570 .await
2571 .unwrap();
2572 }
2573 }
2574 });
2575 let mut client = DaemonClient::connect(metadata.clone()).await.unwrap();
2576 let reply = client
2577 .request_with_reconnect(action.clone(), || DaemonClient::connect(metadata.clone()))
2578 .await
2579 .unwrap();
2580 assert!(matches!(reply, DaemonReply::Done));
2581 let mut client = DaemonClient::connect(metadata).await.unwrap();
2582 assert!(
2583 client
2584 .request_with_reconnect(action, || async {
2585 panic!("an ambiguous acknowledgement must not replay a mutation");
2586 })
2587 .await
2588 .is_err()
2589 );
2590 server.await.unwrap();
2591 }
2592
2593 #[test]
2594 fn unsupported_protocol_message_names_both_binaries_and_versions() {
2595 let message = unsupported_daemon_protocol_message(
2596 PROTOCOL_VERSION + 5,
2597 &describe_daemon_and_client_builds(
2598 4242,
2599 BuildDescription {
2600 executable: Some(Path::new("/home/dev/mj/target/release/mj")),
2601 version: "2.14.0",
2602 },
2603 BuildDescription {
2604 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2605 version: "2.9.0",
2606 },
2607 ),
2608 );
2609 assert_eq!(
2610 message,
2611 format!(
2612 "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
2613 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
2614 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
2615 Put the daemon's directory first on PATH, or reinstall this client from that build.",
2616 PROTOCOL_VERSION + 5
2617 )
2618 );
2619 }
2620
2621 #[test]
2622 fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
2623 let message = unsupported_daemon_protocol_message(
2624 PROTOCOL_VERSION + 1,
2625 &describe_daemon_and_client_builds(
2626 4242,
2627 BuildDescription {
2628 executable: None,
2629 version: "2.14.0",
2630 },
2631 BuildDescription {
2632 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2633 version: "2.9.0",
2634 },
2635 ),
2636 );
2637 assert!(
2638 message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
2639 "{message}"
2640 );
2641 assert!(
2642 message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
2643 "{message}"
2644 );
2645 }
2646}