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