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