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