1use crate::executable::describe_running_daemon_and_client_builds;
3use crate::review::RuntimeReviewView;
4use crate::session::{ManagedSessionView, ViewError};
5use anyhow::{Context, Result, bail, ensure};
6use mj_core::config::{Config, data_dir};
7use mj_core::credentials::CredentialSyncSignal;
8use mj_core::elicitation::ElicitationResponse;
9use mj_core::relay::{RelayCommand, RelayOperationalState};
10use mj_core::review::driver::Resolution;
11use mj_core::state::*;
12use mj_core::targets::{AdditionalMount, ProvisionStage};
13use mj_core::workspace::WorkspaceRecord;
14use serde::{Deserialize, Serialize};
15use std::collections::BTreeMap;
16use std::fs;
17use std::net::SocketAddr;
18use std::path::PathBuf;
19use std::time::{Duration, Instant};
20use tokio::io::{AsyncReadExt, AsyncWriteExt};
21use tokio::net::TcpStream;
22pub fn metadata_path() -> PathBuf {
23 data_dir().join("daemon.json")
24}
25
26#[derive(Debug, Clone, Serialize, Deserialize)]
27#[serde(deny_unknown_fields)]
28pub struct DaemonMetadata {
29 pub protocol_version: u32,
30 pub pid: u32,
31 pub address: SocketAddr,
32 pub token: String,
33 pub started_at: String,
34 pub build_version: String,
35}
36
37#[derive(Debug, Clone, Serialize, Deserialize)]
38#[serde(deny_unknown_fields)]
39pub struct WorkspaceListing {
40 pub workspace: WorkspaceRecord,
41}
42
43#[derive(Debug, Clone, Serialize, Deserialize)]
44#[serde(deny_unknown_fields)]
45pub struct SessionPreview {
46 pub id: String,
47 pub title: String,
48 pub project: String,
49 pub harness: String,
50 pub state: String,
51 pub active: bool,
52 pub updated_at: String,
53}
54
55#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
61#[serde(deny_unknown_fields)]
62pub struct WikiRow {
63 pub id: String,
64 pub tool: String,
65 pub project: String,
66 pub title: String,
67 pub started: Option<String>,
68 pub msgs: i64,
69 pub preview: Option<String>,
70 pub archived: bool,
72 pub native_id: Option<String>,
75 pub snippet: Option<String>,
77 pub hel_session_id: Option<String>,
79 #[serde(default)]
84 pub target: Option<String>,
85 #[serde(default)]
88 pub profile: Option<String>,
89 #[serde(default)]
93 pub harness: Option<String>,
94}
95
96#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
98#[serde(rename_all = "snake_case")]
99pub enum WikiIndexState {
100 Ready,
102 #[default]
105 Indexing,
106 VersionMismatch,
110}
111
112#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
114#[serde(deny_unknown_fields)]
115pub struct WikiStatus {
116 pub state: WikiIndexState,
117 pub topping_up: bool,
119}
120
121#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
123#[serde(deny_unknown_fields)]
124pub struct WikiSearchPage {
125 pub rows: Vec<WikiRow>,
126 pub status: WikiStatus,
127}
128
129#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
132#[serde(deny_unknown_fields)]
133pub struct WikiHitBlock {
134 pub role: String,
136 pub text: String,
139 pub hits: Vec<(usize, usize)>,
142 pub omitted_before: usize,
145 pub truncated: bool,
147}
148
149#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
151#[serde(deny_unknown_fields)]
152pub struct WikiHitTranscript {
153 pub blocks: Vec<WikiHitBlock>,
154 pub omitted_after: usize,
156}
157
158#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
160#[serde(rename_all = "snake_case")]
161pub enum WikiSessionStatus {
162 Mine,
164 Archived,
167 Native,
169}
170
171#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
176#[serde(deny_unknown_fields)]
177pub struct WikiSessionInfo {
178 pub wiki_id: String,
180 pub tool: String,
182 pub path: PathBuf,
184 pub status: WikiSessionStatus,
185 pub mjolnir_session_id: Option<String>,
187 pub profile_id: Option<String>,
189 pub target_template_id: Option<String>,
191 pub harness: Option<mj_core::config::HarnessKind>,
194 pub title: String,
195 pub project: String,
196 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
201 pub nothing_to_restore: bool,
202}
203
204#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
206#[serde(deny_unknown_fields)]
207pub struct WikiRestoreRequest {
208 pub wiki_id: String,
210 pub workspace_id: String,
211 pub profile_id: String,
212 pub target_template_id: String,
213 #[serde(default)]
216 pub project_directory: Option<PathBuf>,
217 #[serde(default)]
218 pub additional_mounts: Vec<AdditionalMount>,
219 #[serde(default)]
220 pub resource_allocation: Option<SessionResourceAllocation>,
221}
222
223#[derive(Debug, Clone, Serialize, Deserialize)]
224#[serde(deny_unknown_fields)]
225pub struct WorkspaceSnapshot {
226 pub workspace: WorkspaceRecord,
227 pub sessions: Vec<SessionPreview>,
228 pub drafts: Vec<DraftPreview>,
229}
230
231#[derive(Debug, Clone, Serialize, Deserialize)]
232#[serde(deny_unknown_fields)]
233pub struct RuntimeSessionView {
234 pub session_id: String,
235 pub projection_ordinal: u64,
236 pub projection_digest: String,
237 pub operational: Option<RelayOperationalState>,
238 pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
239 pub connected: bool,
240 pub error: Option<ViewError>,
241}
242
243impl RuntimeSessionView {
244 pub fn from_managed(session_id: String, view: ManagedSessionView) -> Self {
245 let (projection_ordinal, projection_digest, operational, signal) =
246 view.snapshot
247 .map_or((0, String::new(), None, None), |snapshot| {
248 (
249 snapshot.materialized.applied_event_ordinal,
250 snapshot.materialized.applied_event_digest,
251 Some(snapshot.operational),
252 snapshot.latest_credential_sync_signal,
253 )
254 });
255 Self {
256 session_id,
257 projection_ordinal,
258 projection_digest,
259 operational,
260 latest_credential_sync_signal: signal,
261 connected: view.connected,
262 error: view.error,
263 }
264 }
265}
266
267#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
273#[serde(deny_unknown_fields)]
274pub struct RuntimeNotice {
275 pub id: u64,
276 pub session_id: String,
277 pub text: String,
278}
279
280#[derive(Debug, Clone, Serialize, Deserialize)]
281#[serde(deny_unknown_fields)]
282pub struct RuntimeSnapshot {
283 #[serde(default)]
284 pub native_agents: Vec<mj_core::native_agent::NativeAgentSummary>,
285 #[serde(default)]
286 pub workspace_names: BTreeMap<String, String>,
287 #[serde(default)]
288 pub moves: Vec<mj_core::state::MoveOperation>,
289 pub revision: u64,
290 pub config: Config,
291 pub records: Vec<SessionRecord>,
292 pub sessions: Vec<RuntimeSessionView>,
293 pub lifecycles: Vec<RuntimeLifecycleView>,
294 #[serde(default)]
296 pub reviews: Vec<RuntimeReviewView>,
297 #[serde(default)]
299 pub notices: Vec<RuntimeNotice>,
300 #[serde(default)]
304 pub subagents: Vec<mj_core::subagent::SubagentRecord>,
305}
306
307#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
308#[serde(rename_all = "snake_case")]
309pub enum RuntimeLifecycleKind {
310 Create,
311 Suspend,
312 Resume,
313 Move,
314 ForceStop,
315 DestroyStopped,
316 ForceDestroy,
317 StopSubagent,
321 Cleanup,
322}
323
324#[derive(Debug, Clone, Serialize, Deserialize)]
325#[serde(deny_unknown_fields)]
326pub struct RuntimeLifecycleView {
327 pub operation_id: String,
328 pub cancellable: bool,
329 pub session_id: String,
330 pub kind: RuntimeLifecycleKind,
331 pub started_at_epoch_seconds: u64,
332 pub active_stages: Vec<(ProvisionStage, u64)>,
333 pub resume_destination: Option<(String, String)>,
334 pub notice: Option<String>,
335}
336
337#[derive(Debug, Clone, Serialize, Deserialize)]
338#[serde(deny_unknown_fields)]
339pub struct ResumeSessionRequest {
340 pub session_id: String,
341 pub workspace_id: String,
342 pub profile_id: String,
343 pub target_template_id: String,
344 pub additional_mounts: Option<Vec<AdditionalMount>>,
345 pub resource_allocation: Option<SessionResourceAllocation>,
346 pub discard_queue: bool,
347 pub repository_preflight: Option<ResumeRepositorySourceReceipt>,
348}
349
350#[derive(Debug, Clone, Serialize, Deserialize)]
351#[serde(deny_unknown_fields)]
352pub struct CreateSessionRequest {
353 #[serde(default)]
354 pub create_managed_worktree: Option<bool>,
355 #[serde(default)]
357 pub launch_base: Option<String>,
358 #[serde(default)]
359 pub launch_branch: Option<String>,
360 #[serde(default)]
362 pub mjolnir_subagents: Option<bool>,
363 #[serde(default)]
364 pub initial_prompt: Option<String>,
365 pub workspace_id: String,
366 pub profile_id: String,
367 pub bundle_id: String,
368 pub project_directory: Option<PathBuf>,
369 pub target_template_id: String,
370 pub additional_mounts: Vec<AdditionalMount>,
371 pub resource_allocation: Option<SessionResourceAllocation>,
372 pub title: String,
373 pub session_title_override: Option<String>,
374}
375
376#[derive(Debug, Clone, Serialize, Deserialize)]
377#[serde(deny_unknown_fields)]
378pub struct RegisteredSession {
379 pub session: SessionRecord,
380 pub remembered_container_size: Option<(String, HostContainerSize)>,
381}
382
383#[derive(Debug, Clone, Serialize, Deserialize)]
384#[serde(deny_unknown_fields)]
385pub struct DraftPreview {
386 pub id: String,
387 pub session_id: Option<String>,
388 pub source: String,
389 pub owner_pid: Option<u32>,
390 pub saved_at: String,
391}
392
393#[derive(Debug, Clone, Serialize, Deserialize)]
402#[serde(rename_all = "snake_case", tag = "action", content = "arguments")]
403pub enum DaemonAction {
404 NativeAgentHistory {
405 owner: String,
406 child: String,
407 before: Option<(u64, String)>,
408 },
409 Ping,
410 PrepareUpgrade,
413 UpgradeBlockers,
422 Status,
423 WebViewerAccess,
424 RecoverWebViewer(crate::web::WebViewerRecovery),
425 InspectWebListener,
426 ListWorkspaces,
427 CreateWorkspace {
428 name: String,
429 },
430 RenameWorkspace {
431 workspace_id: String,
432 name: String,
433 },
434 TouchWorkspace {
435 workspace_id: String,
436 },
437 CloseWorkspace {
438 workspace_id: String,
439 },
440 CancelWorkspaceClose {
441 workspace_id: String,
442 },
443 DeleteWorkspace {
444 workspace_id: String,
445 },
446 Attach {
447 client_id: String,
448 pid: u32,
449 },
450 Detach {
451 client_id: String,
452 },
453 PersistReadReceipt {
454 client_id: String,
455 workspace_id: String,
456 session_id: String,
457 through: u64,
458 },
459 PersistDetachedSessionState {
460 client_id: String,
461 workspace_id: String,
462 session_id: String,
463 through: u64,
464 owner_pid: u32,
465 draft: mj_core::storage::DetachedSessionDraft,
466 },
467 SaveActiveReview {
468 session_id: String,
469 review: mj_core::storage::StoredReview,
470 },
471 ClearActiveReview {
472 session_id: String,
473 },
474 SaveWorkspacePaneSizes {
475 workspace_id: String,
476 sizes: mj_core::workspace::PaneSizes,
477 },
478 SaveWorkspaceLayout {
479 workspace_id: String,
480 layout: mj_core::workspace::ConversationLayout,
481 },
482 PersistImportedSession {
483 session: Box<SessionRecord>,
484 },
485 SetSessionTitle {
486 session_id: String,
487 title: String,
488 },
489 SetSessionContainerSettings {
490 session_id: String,
491 cpus: Option<String>,
492 memory: Option<String>,
493 mounts: Vec<AdditionalMount>,
494 mount_history: Vec<PathBuf>,
495 },
496 SetSessionAcpTitle {
497 session_id: String,
498 title: Option<String>,
499 },
500 MarkSessionTargetMissing {
501 session_id: String,
502 detail: String,
503 updated_at: String,
504 },
505 CheckpointSession {
506 session_id: String,
507 },
508 WikiSearch {
511 query: String,
512 limit: usize,
513 },
514 WikiBrief {
516 wiki_id: String,
517 max_chars: usize,
518 },
519 WikiHits {
521 wiki_id: String,
522 query: String,
523 context_messages: usize,
524 per_message_chars: usize,
525 },
526 WikiSession {
528 wiki_id: String,
529 },
530 WikiRestore(WikiRestoreRequest),
532 ScanRecovery {
533 all_instances: bool,
534 },
535 AdoptRecovery {
536 session_id: String,
537 target_id: String,
538 profile: Option<String>,
539 bundle: Option<String>,
540 all_instances: bool,
541 },
542 DestroyRecovery {
543 session_id: String,
544 target_id: String,
545 confirmation: String,
546 all_instances: bool,
547 },
548 Snapshot {
549 workspace_id: String,
550 },
551 RuntimeSnapshot {
552 workspace_id: String,
553 after_revision: u64,
554 #[serde(default)]
555 all_workspaces: bool,
556 },
557 RenameProfile {
558 old_id: String,
559 new_id: String,
560 },
561 RenameTarget {
562 old_id: String,
563 new_id: String,
564 },
565 SubmitSessionCommand {
566 #[serde(default)]
567 inherited_draft: Option<String>,
568 session_id: String,
569 command_id: String,
570 command: RelayCommand,
571 },
572 QueueStartupPrompt {
577 session_id: String,
578 text: String,
579 #[serde(default)]
583 inherited_draft: Option<String>,
584 },
585 SyncSession {
586 session_id: String,
587 },
588 RespondElicitation {
589 session_id: String,
590 elicitation_id: String,
591 response: ElicitationResponse,
592 },
593 StopBackgroundTask {
594 session_id: String,
595 background_task_id: String,
596 },
597 ReviewerAction {
601 session_id: String,
602 #[serde(default, skip_serializing_if = "Option::is_none")]
605 role: Option<String>,
606 action: crate::session::ReviewerAction,
607 },
608 StartTurnReview {
610 session_id: String,
611 },
612 ResolveTurnReview {
614 session_id: String,
615 resolution: Resolution,
616 },
617 SuspendSession {
618 session_id: String,
619 #[serde(default)]
620 acknowledge_unpublished_work: bool,
621 },
622 StartCreateSession(CreateSessionRequest),
623 WaitCreateSession {
624 session_id: String,
625 },
626 ResumeSession(ResumeSessionRequest),
627 PrepareMoveSession(MoveSelection),
628 MoveSession(MoveSessionRequest),
629 DiscardSinceCheckpoint {
630 session_id: String,
631 checkpoint: mj_core::state::CheckpointMetadata,
632 },
633 DestroyStoppedSession {
634 session_id: String,
635 delete_branch: bool,
639 },
640 ForceDestroySession {
641 session_id: String,
642 delete_branch: bool,
644 },
645 CancelLifecycle {
646 session_id: String,
647 },
648 RecoverDraft {
649 draft_id: String,
650 },
651 Stop,
652}
653
654#[derive(Debug, Serialize, Deserialize)]
655#[serde(deny_unknown_fields)]
656pub struct RequestEnvelope {
657 pub protocol_version: u32,
658 pub request_id: u64,
659 pub token: String,
660 pub action: DaemonAction,
661}
662
663#[derive(Debug)]
667pub struct DaemonRefusal(pub String);
668
669impl std::fmt::Display for DaemonRefusal {
670 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
671 f.write_str(&self.0)
672 }
673}
674
675impl std::error::Error for DaemonRefusal {}
676
677impl DaemonRefusal {
678 #[must_use]
681 pub fn delivery_unconfirmed(&self) -> bool {
682 self.0
683 .starts_with(&crate::session::DeliveryUnconfirmed.to_string())
684 }
685}
686
687#[derive(Debug, Serialize, Deserialize)]
688#[serde(deny_unknown_fields)]
689pub struct ResponseEnvelope {
690 pub protocol_version: u32,
691 pub request_id: u64,
692 pub result: std::result::Result<DaemonReply, String>,
693}
694
695#[derive(Debug, Clone, Serialize, Deserialize)]
696#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
697pub enum DaemonReply {
698 NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
699 Pong,
700 UpgradePending,
701 UpgradeBlockers(Vec<String>),
704 Status(DaemonStatus),
705 WebViewerAccess(crate::web::WebViewerAccess),
706 WebListeners(Vec<crate::web::WebListenerProcess>),
707 Workspaces(Vec<WorkspaceListing>),
708 Workspace(WorkspaceRecord),
709 Snapshot(WorkspaceSnapshot),
710 RuntimeSnapshot(Box<RuntimeSnapshot>),
711 RegisteredSession(Box<RegisteredSession>),
712 MovePreparation(Box<MovePreparation>),
713 MoveOutcome(MoveOutcome),
714 Ordinal(u64),
715 Text(String),
716 OptionalSessionState(Option<SessionState>),
717 Checkpoint(mj_core::state::CheckpointMetadata),
718 RecoveryScan(mj_core::state::RecoveryScan),
719 WikiRows(WikiSearchPage),
720 WikiHits(Option<WikiHitTranscript>),
721 WikiSession(Option<Box<WikiSessionInfo>>),
722 Reviewer(Box<crate::session::ReviewerOutcome>),
723 Done,
724}
725
726#[derive(Debug, Clone, Serialize, Deserialize)]
727#[serde(deny_unknown_fields)]
728pub struct DaemonStatus {
729 pub pid: u32,
730 pub started_at: String,
731 pub build_version: String,
732 pub attached_clients: usize,
733 pub phone_status: WebViewerStatus,
734}
735
736#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
737#[serde(rename_all = "snake_case", tag = "state")]
738pub enum WebViewerStatus {
739 Disabled,
740 Starting,
741 Ready {
742 viewer_url: String,
743 viewer_code: String,
744 qr_login_url: Option<String>,
745 fallback_reason: Option<String>,
746 },
747 Stopped,
748 Error {
749 message: String,
750 },
751}
752
753impl std::fmt::Debug for WebViewerStatus {
754 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
755 match self {
756 Self::Ready {
757 viewer_url,
758 viewer_code,
759 fallback_reason,
760 ..
761 } => formatter
762 .debug_struct("Ready")
763 .field("viewer_url", viewer_url)
764 .field("viewer_code", viewer_code)
765 .field("qr_login_url", &"[redacted]")
766 .field("fallback_reason", fallback_reason)
767 .finish(),
768 Self::Disabled => formatter.write_str("Disabled"),
769 Self::Starting => formatter.write_str("Starting"),
770 Self::Stopped => formatter.write_str("Stopped"),
771 Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
772 }
773 }
774}
775
776impl std::fmt::Display for WebViewerStatus {
777 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
778 match self {
779 Self::Disabled => formatter.write_str("disabled"),
780 Self::Starting => formatter.write_str("starting"),
781 Self::Stopped => formatter.write_str("stopped unexpectedly"),
782 Self::Error { message } => write!(formatter, "error: {message}"),
783 Self::Ready {
784 viewer_url,
785 viewer_code,
786 fallback_reason,
787 ..
788 } => {
789 write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
790 if let Some(reason) = fallback_reason {
791 write!(
792 formatter,
793 "; local only because Tailscale HTTPS is unavailable: {reason}"
794 )?;
795 }
796 Ok(())
797 }
798 }
799 }
800}
801
802#[cfg(target_os = "linux")]
821pub fn process_is_zombie(pid: u32) -> bool {
822 let Ok(stat) = std::fs::read(format!("/proc/{pid}/stat")) else {
823 return false;
824 };
825 stat.iter()
828 .rposition(|byte| *byte == b')')
829 .and_then(|end| stat.get(end + 2))
830 .is_some_and(|state| *state == b'Z')
831}
832
833#[cfg(target_os = "macos")]
839pub fn process_is_zombie(pid: u32) -> bool {
840 let Ok(raw_pid) = libc::c_int::try_from(pid) else {
841 return false;
842 };
843 let size = std::mem::size_of::<libc::proc_bsdinfo>();
844 let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
846 let written = unsafe {
848 libc::proc_pidinfo(
849 raw_pid,
850 libc::PROC_PIDTBSDINFO,
851 1,
852 (&mut info as *mut libc::proc_bsdinfo).cast(),
853 size as libc::c_int,
854 )
855 };
856 usize::try_from(written).is_ok_and(|written| written == size) && info.pbi_status == libc::SZOMB
857}
858
859#[cfg(all(unix, not(target_os = "linux"), not(target_os = "macos")))]
860pub fn process_is_zombie(pid: u32) -> bool {
861 let pid = sysinfo::Pid::from_u32(pid);
862 let mut system = sysinfo::System::new();
863 system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
864 system
865 .process(pid)
866 .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
867}
868
869pub async fn wait_for_exit(pid: u32) -> Result<()> {
875 let deadline = Instant::now() + STOP_TIMEOUT;
876 while daemon_process_is_alive(pid) {
877 ensure!(Instant::now() < deadline, "process {pid} is still running");
878 tokio::time::sleep(RETRY_DELAY).await;
879 }
880 Ok(())
881}
882
883pub fn daemon_process_is_alive(pid: u32) -> bool {
889 #[cfg(unix)]
890 {
891 if pid == 0 {
892 return false;
893 }
894 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
895 return false;
896 };
897 let mut status = 0;
898 let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
901 if waited == raw_pid {
902 return false;
903 }
904 if waited == 0 {
905 return true;
906 }
907 let wait_error = std::io::Error::last_os_error();
908 if wait_error.raw_os_error() != Some(libc::ECHILD) {
909 return true;
910 }
911
912 #[cfg(target_os = "macos")]
913 return owned_daemon_group_is_alive(raw_pid);
914
915 #[cfg(not(target_os = "macos"))]
916 process_is_alive(pid)
917 }
918 #[cfg(not(unix))]
919 process_is_alive(pid)
920}
921
922#[cfg(target_os = "macos")]
923pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
924 if unsafe { libc::kill(-pid, 0) } == 0 {
932 return true;
933 }
934 let error = std::io::Error::last_os_error();
935 !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
936}
937
938pub fn process_is_alive(pid: u32) -> bool {
939 #[cfg(unix)]
940 {
941 if pid == 0 {
942 return false;
943 }
944 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
945 return false;
946 };
947 let result = unsafe { libc::kill(raw_pid, 0) };
950 let exists =
951 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
952 exists && !process_is_zombie(pid)
953 }
954 #[cfg(not(unix))]
955 {
956 let process_id = sysinfo::Pid::from_u32(pid);
957 let mut system = sysinfo::System::new();
958 system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[process_id]), true);
959 system.process(process_id).is_some()
960 }
961}
962
963pub fn read_metadata() -> Result<DaemonMetadata> {
964 let metadata = read_metadata_any()?;
965 ensure!(
966 metadata.protocol_version == PROTOCOL_VERSION,
967 "daemon protocol {} is incompatible with client protocol {}",
968 metadata.protocol_version,
969 PROTOCOL_VERSION
970 );
971 Ok(metadata)
972}
973
974#[derive(Debug, Clone, PartialEq, Eq)]
980pub struct DaemonNotRunning {
981 pub metadata_path: std::path::PathBuf,
983}
984
985impl std::fmt::Display for DaemonNotRunning {
986 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
987 formatter.write_str("the Mjolnir daemon is not running")
988 }
989}
990
991impl std::error::Error for DaemonNotRunning {}
992
993pub fn daemon_not_running(error: &anyhow::Error) -> Option<&DaemonNotRunning> {
995 error
996 .chain()
997 .find_map(|cause| cause.downcast_ref::<DaemonNotRunning>())
998}
999
1000pub fn read_metadata_any() -> Result<DaemonMetadata> {
1001 read_metadata_at(&metadata_path())
1002}
1003
1004fn read_metadata_at(path: &std::path::Path) -> Result<DaemonMetadata> {
1005 let body = match fs::read(path) {
1006 Ok(body) => body,
1007 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1008 return Err(DaemonNotRunning {
1009 metadata_path: path.to_owned(),
1010 }
1011 .into());
1012 }
1013 Err(error) => return Err(error).with_context(|| format!("read {}", path.display())),
1014 };
1015 let metadata: DaemonMetadata =
1016 serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1017 Ok(metadata)
1018}
1019
1020pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
1021 let body = serde_json::to_vec(value)?;
1022 write_encoded_frame(stream, &body).await
1023}
1024
1025pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
1026 ensure!(
1027 body.len() <= MAX_FRAME_BYTES,
1028 "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
1029 body.len()
1030 );
1031 stream.write_u32(body.len() as u32).await?;
1032 stream.write_all(body).await?;
1033 stream.flush().await?;
1034 Ok(())
1035}
1036
1037pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
1038 let length = stream.read_u32().await? as usize;
1039 ensure!(
1040 length <= MAX_FRAME_BYTES,
1041 "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
1042 );
1043 let mut body = vec![0_u8; length];
1044 stream.read_exact(&mut body).await?;
1045 serde_json::from_slice(&body).context("decode daemon frame")
1046}
1047
1048pub struct DaemonClient {
1049 metadata: DaemonMetadata,
1050 stream: TcpStream,
1051 next_request_id: u64,
1052}
1053
1054impl DaemonClient {
1055 pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
1056 let stream =
1057 tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
1058 .await
1059 .context("time out connecting to Mjolnir daemon")??;
1060 Ok(Self {
1061 metadata,
1062 stream,
1063 next_request_id: 1,
1064 })
1065 }
1066
1067 pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1071 self.request_with_reconnect(action, || async {
1072 loop {
1073 if let Ok(client) = connect_existing().await {
1074 return Ok(client);
1075 }
1076 tokio::time::sleep(Duration::from_millis(250)).await;
1077 }
1078 })
1079 .await
1080 }
1081
1082 async fn request_with_reconnect<F, Fut>(
1083 &mut self,
1084 action: DaemonAction,
1085 mut reconnect: F,
1086 ) -> Result<DaemonReply>
1087 where
1088 F: FnMut() -> Fut,
1089 Fut: Future<Output = Result<Self>>,
1090 {
1091 loop {
1092 let response = self.request_once(action.clone()).await?;
1093 if matches!(response, DaemonReply::UpgradePending)
1094 && !matches!(action, DaemonAction::PrepareUpgrade)
1095 {
1096 tokio::time::sleep(Duration::from_millis(100)).await;
1099 *self = reconnect().await?;
1100 } else {
1101 return Ok(response);
1102 }
1103 }
1104 }
1105
1106 async fn request_once(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1107 let protocol_version = self.metadata.protocol_version;
1108 let request_id = self.next_request_id;
1109 self.next_request_id += 1;
1110 write_frame(
1111 &mut self.stream,
1112 &RequestEnvelope {
1113 protocol_version,
1114 request_id,
1115 token: self.metadata.token.clone(),
1116 action,
1117 },
1118 )
1119 .await?;
1120 let response: ResponseEnvelope = read_frame(&mut self.stream).await?;
1121 ensure!(
1122 response.protocol_version == protocol_version,
1123 "daemon changed protocol"
1124 );
1125 ensure!(
1126 response.request_id == request_id,
1127 "daemon crossed request IDs"
1128 );
1129 response
1130 .result
1131 .map_err(|message| anyhow::Error::new(DaemonRefusal(message)))
1132 }
1133
1134 pub async fn status(&mut self) -> Result<DaemonStatus> {
1135 match self.request(DaemonAction::Status).await? {
1136 DaemonReply::Status(status) => Ok(status),
1137 reply => bail!("unexpected daemon status reply {reply:?}"),
1138 }
1139 }
1140
1141 pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
1142 match self.request(DaemonAction::WebViewerAccess).await? {
1143 DaemonReply::WebViewerAccess(access) => Ok(access),
1144 reply => bail!("unexpected web viewer reply {reply:?}"),
1145 }
1146 }
1147
1148 pub async fn recover_web_viewer(
1149 &mut self,
1150 action: crate::web::WebViewerRecovery,
1151 ) -> Result<()> {
1152 match self.request(DaemonAction::RecoverWebViewer(action)).await? {
1153 DaemonReply::Done => Ok(()),
1154 reply => bail!("unexpected web viewer recovery reply {reply:?}"),
1155 }
1156 }
1157
1158 pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
1159 match self.request(DaemonAction::InspectWebListener).await? {
1160 DaemonReply::WebListeners(processes) => Ok(processes),
1161 reply => bail!("unexpected listener inspection reply {reply:?}"),
1162 }
1163 }
1164
1165 pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
1166 match self.request(DaemonAction::ListWorkspaces).await? {
1167 DaemonReply::Workspaces(workspaces) => Ok(workspaces),
1168 reply => bail!("unexpected daemon workspace reply {reply:?}"),
1169 }
1170 }
1171
1172 pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
1173 match self
1174 .request(DaemonAction::RenameProfile { old_id, new_id })
1175 .await?
1176 {
1177 DaemonReply::Done => Ok(()),
1178 reply => bail!("unexpected rename-profile reply {reply:?}"),
1179 }
1180 }
1181
1182 pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1183 match self
1184 .request(DaemonAction::RenameTarget { old_id, new_id })
1185 .await?
1186 {
1187 DaemonReply::Done => Ok(()),
1188 reply => bail!("unexpected rename-target reply {reply:?}"),
1189 }
1190 }
1191
1192 pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1193 match self.request(DaemonAction::CreateWorkspace { name }).await? {
1194 DaemonReply::Workspace(workspace) => Ok(workspace),
1195 reply => bail!("unexpected create-workspace reply {reply:?}"),
1196 }
1197 }
1198
1199 pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1200 match self
1201 .request(DaemonAction::RenameWorkspace { workspace_id, name })
1202 .await?
1203 {
1204 DaemonReply::Done => Ok(()),
1205 reply => bail!("unexpected rename-workspace reply {reply:?}"),
1206 }
1207 }
1208
1209 pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1210 match self
1211 .request(DaemonAction::TouchWorkspace { workspace_id })
1212 .await?
1213 {
1214 DaemonReply::Done => Ok(()),
1215 reply => bail!("unexpected touch-workspace reply {reply:?}"),
1216 }
1217 }
1218
1219 pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1220 match self
1221 .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1222 .await?
1223 {
1224 DaemonReply::Done => Ok(()),
1225 reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1226 }
1227 }
1228
1229 pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1230 match self
1231 .request(DaemonAction::CloseWorkspace { workspace_id })
1232 .await?
1233 {
1234 DaemonReply::Done => Ok(()),
1235 reply => bail!("unexpected close workspace reply: {reply:?}"),
1236 }
1237 }
1238
1239 pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1240 match self
1241 .request(DaemonAction::DeleteWorkspace { workspace_id })
1242 .await?
1243 {
1244 DaemonReply::Done => Ok(()),
1245 reply => bail!("unexpected delete-workspace reply {reply:?}"),
1246 }
1247 }
1248
1249 pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1250 match self
1251 .request(DaemonAction::Attach { client_id, pid })
1252 .await?
1253 {
1254 DaemonReply::Done => Ok(()),
1255 reply => bail!("unexpected attach reply {reply:?}"),
1256 }
1257 }
1258
1259 pub async fn detach(&mut self, client_id: String) -> Result<()> {
1260 match self.request(DaemonAction::Detach { client_id }).await? {
1261 DaemonReply::Done => Ok(()),
1262 reply => bail!("unexpected detach reply {reply:?}"),
1263 }
1264 }
1265
1266 pub async fn persist_read_receipt(
1267 &mut self,
1268 client_id: String,
1269 workspace_id: String,
1270 session_id: String,
1271 through: u64,
1272 ) -> Result<u64> {
1273 match self
1274 .request(DaemonAction::PersistReadReceipt {
1275 client_id,
1276 workspace_id,
1277 session_id,
1278 through,
1279 })
1280 .await?
1281 {
1282 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1283 reply => bail!("unexpected read-receipt reply {reply:?}"),
1284 }
1285 }
1286
1287 pub async fn persist_detached_session_state(
1288 &mut self,
1289 client_id: String,
1290 workspace_id: String,
1291 session_id: String,
1292 through: u64,
1293 owner_pid: u32,
1294 draft: mj_core::storage::DetachedSessionDraft,
1295 ) -> Result<()> {
1296 match self
1297 .request(DaemonAction::PersistDetachedSessionState {
1298 client_id,
1299 workspace_id,
1300 session_id,
1301 through,
1302 owner_pid,
1303 draft,
1304 })
1305 .await?
1306 {
1307 DaemonReply::Done => Ok(()),
1308 reply => bail!("unexpected detached-session-state reply {reply:?}"),
1309 }
1310 }
1311
1312 pub async fn save_active_review(
1313 &mut self,
1314 session_id: String,
1315 review: mj_core::storage::StoredReview,
1316 ) -> Result<()> {
1317 match self
1318 .request(DaemonAction::SaveActiveReview { session_id, review })
1319 .await?
1320 {
1321 DaemonReply::Done => Ok(()),
1322 reply => bail!("unexpected save-review reply {reply:?}"),
1323 }
1324 }
1325
1326 pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1327 match self
1328 .request(DaemonAction::ClearActiveReview { session_id })
1329 .await?
1330 {
1331 DaemonReply::Done => Ok(()),
1332 reply => bail!("unexpected clear-review reply {reply:?}"),
1333 }
1334 }
1335
1336 pub async fn save_workspace_pane_sizes(
1337 &mut self,
1338 workspace_id: String,
1339 sizes: mj_core::workspace::PaneSizes,
1340 ) -> Result<()> {
1341 match self
1342 .request(DaemonAction::SaveWorkspacePaneSizes {
1343 workspace_id,
1344 sizes,
1345 })
1346 .await?
1347 {
1348 DaemonReply::Done => Ok(()),
1349 reply => bail!("unexpected pane-size save reply {reply:?}"),
1350 }
1351 }
1352
1353 pub async fn save_workspace_layout(
1354 &mut self,
1355 workspace_id: String,
1356 layout: mj_core::workspace::ConversationLayout,
1357 ) -> Result<()> {
1358 match self
1359 .request(DaemonAction::SaveWorkspaceLayout {
1360 workspace_id,
1361 layout,
1362 })
1363 .await?
1364 {
1365 DaemonReply::Done => Ok(()),
1366 reply => bail!("unexpected layout save reply {reply:?}"),
1367 }
1368 }
1369
1370 pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1371 match self
1372 .request(DaemonAction::PersistImportedSession {
1373 session: Box::new(session),
1374 })
1375 .await?
1376 {
1377 DaemonReply::Done => Ok(()),
1378 reply => bail!("unexpected imported-session reply {reply:?}"),
1379 }
1380 }
1381
1382 pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1383 match self
1384 .request(DaemonAction::SetSessionTitle { session_id, title })
1385 .await?
1386 {
1387 DaemonReply::Text(title) => Ok(title),
1388 reply => bail!("unexpected session-title reply {reply:?}"),
1389 }
1390 }
1391
1392 pub async fn set_session_container_settings(
1393 &mut self,
1394 session_id: String,
1395 cpus: Option<String>,
1396 memory: Option<String>,
1397 mounts: Vec<AdditionalMount>,
1398 mount_history: Vec<PathBuf>,
1399 ) -> Result<()> {
1400 match self
1401 .request(DaemonAction::SetSessionContainerSettings {
1402 session_id,
1403 cpus,
1404 memory,
1405 mounts,
1406 mount_history,
1407 })
1408 .await?
1409 {
1410 DaemonReply::Done => Ok(()),
1411 reply => bail!("unexpected container-settings reply {reply:?}"),
1412 }
1413 }
1414
1415 pub async fn set_session_acp_title(
1416 &mut self,
1417 session_id: String,
1418 title: Option<String>,
1419 ) -> Result<()> {
1420 match self
1421 .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1422 .await?
1423 {
1424 DaemonReply::Done => Ok(()),
1425 reply => bail!("unexpected ACP-title reply {reply:?}"),
1426 }
1427 }
1428
1429 pub async fn mark_session_target_missing(
1430 &mut self,
1431 session_id: String,
1432 detail: String,
1433 updated_at: String,
1434 ) -> Result<Option<SessionState>> {
1435 match self
1436 .request(DaemonAction::MarkSessionTargetMissing {
1437 session_id,
1438 detail,
1439 updated_at,
1440 })
1441 .await?
1442 {
1443 DaemonReply::OptionalSessionState(state) => Ok(state),
1444 reply => bail!("unexpected target-missing reply {reply:?}"),
1445 }
1446 }
1447
1448 pub async fn checkpoint_session(
1449 &mut self,
1450 session_id: String,
1451 ) -> Result<mj_core::state::CheckpointMetadata> {
1452 match self
1453 .request(DaemonAction::CheckpointSession { session_id })
1454 .await?
1455 {
1456 DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1457 reply => bail!("unexpected checkpoint reply {reply:?}"),
1458 }
1459 }
1460
1461 pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1466 match self
1467 .request(DaemonAction::WikiSearch { query, limit })
1468 .await?
1469 {
1470 DaemonReply::WikiRows(page) => Ok(page),
1471 reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1472 }
1473 }
1474
1475 pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1477 match self
1478 .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1479 .await?
1480 {
1481 DaemonReply::Text(markdown) => Ok(markdown),
1482 reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1483 }
1484 }
1485
1486 pub async fn wiki_hits(
1491 &mut self,
1492 wiki_id: String,
1493 query: String,
1494 context_messages: usize,
1495 per_message_chars: usize,
1496 ) -> Result<Option<WikiHitTranscript>> {
1497 match self
1498 .request(DaemonAction::WikiHits {
1499 wiki_id,
1500 query,
1501 context_messages,
1502 per_message_chars,
1503 })
1504 .await?
1505 {
1506 DaemonReply::WikiHits(transcript) => Ok(transcript),
1507 reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1508 }
1509 }
1510
1511 pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1514 match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1515 DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1516 reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1517 }
1518 }
1519
1520 pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1524 match self.request(DaemonAction::WikiRestore(request)).await? {
1525 DaemonReply::RegisteredSession(registered) => Ok(*registered),
1526 reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1527 }
1528 }
1529
1530 pub async fn scan_recovery(
1531 &mut self,
1532 all_instances: bool,
1533 ) -> Result<mj_core::state::RecoveryScan> {
1534 match self
1535 .request(DaemonAction::ScanRecovery { all_instances })
1536 .await?
1537 {
1538 DaemonReply::RecoveryScan(scan) => Ok(scan),
1539 reply => bail!("unexpected recovery-scan reply {reply:?}"),
1540 }
1541 }
1542
1543 pub async fn adopt_recovery(
1544 &mut self,
1545 session_id: String,
1546 target_id: String,
1547 profile: Option<String>,
1548 bundle: Option<String>,
1549 all_instances: bool,
1550 ) -> Result<()> {
1551 match self
1552 .request(DaemonAction::AdoptRecovery {
1553 session_id,
1554 target_id,
1555 profile,
1556 bundle,
1557 all_instances,
1558 })
1559 .await?
1560 {
1561 DaemonReply::Done => Ok(()),
1562 reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1563 }
1564 }
1565
1566 pub async fn destroy_recovery(
1567 &mut self,
1568 session_id: String,
1569 target_id: String,
1570 confirmation: String,
1571 all_instances: bool,
1572 ) -> Result<()> {
1573 match self
1574 .request(DaemonAction::DestroyRecovery {
1575 session_id,
1576 target_id,
1577 confirmation,
1578 all_instances,
1579 })
1580 .await?
1581 {
1582 DaemonReply::Done => Ok(()),
1583 reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1584 }
1585 }
1586
1587 pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1588 match self
1589 .request(DaemonAction::Snapshot { workspace_id })
1590 .await?
1591 {
1592 DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1593 reply => bail!("unexpected snapshot reply {reply:?}"),
1594 }
1595 }
1596
1597 pub async fn runtime_snapshot(
1598 &mut self,
1599 workspace_id: String,
1600 after_revision: u64,
1601 all_workspaces: bool,
1602 ) -> Result<RuntimeSnapshot> {
1603 match self
1604 .request(DaemonAction::RuntimeSnapshot {
1605 workspace_id,
1606 after_revision,
1607 all_workspaces,
1608 })
1609 .await?
1610 {
1611 DaemonReply::RuntimeSnapshot(snapshot) => Ok(*snapshot),
1612 reply => bail!("unexpected runtime snapshot reply {reply:?}"),
1613 }
1614 }
1615
1616 pub async fn submit_session_command(
1617 &mut self,
1618 session_id: String,
1619 command_id: String,
1620 command: RelayCommand,
1621 inherited_draft: Option<String>,
1622 ) -> Result<u64> {
1623 match self
1624 .request(DaemonAction::SubmitSessionCommand {
1625 inherited_draft,
1626 session_id,
1627 command_id,
1628 command,
1629 })
1630 .await?
1631 {
1632 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1633 reply => bail!("unexpected session command reply {reply:?}"),
1634 }
1635 }
1636
1637 pub async fn queue_startup_prompt(
1641 &mut self,
1642 session_id: String,
1643 text: String,
1644 inherited_draft: Option<String>,
1645 ) -> Result<()> {
1646 match self
1647 .request(DaemonAction::QueueStartupPrompt {
1648 session_id,
1649 text,
1650 inherited_draft,
1651 })
1652 .await?
1653 {
1654 DaemonReply::Done => Ok(()),
1655 reply => bail!("unexpected startup prompt reply {reply:?}"),
1656 }
1657 }
1658
1659 pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
1665 match self
1666 .request(DaemonAction::StartTurnReview { session_id })
1667 .await?
1668 {
1669 DaemonReply::Done => Ok(()),
1670 reply => bail!("unexpected turn-review reply {reply:?}"),
1671 }
1672 }
1673
1674 pub async fn resolve_turn_review(
1675 &mut self,
1676 session_id: String,
1677 resolution: Resolution,
1678 ) -> Result<()> {
1679 match self
1680 .request(DaemonAction::ResolveTurnReview {
1681 session_id,
1682 resolution,
1683 })
1684 .await?
1685 {
1686 DaemonReply::Done => Ok(()),
1687 reply => bail!("unexpected turn-review resolution reply {reply:?}"),
1688 }
1689 }
1690
1691 pub async fn reviewer_action(
1692 &mut self,
1693 session_id: String,
1694 role: Option<String>,
1695 action: crate::session::ReviewerAction,
1696 ) -> Result<crate::session::ReviewerOutcome> {
1697 match self
1698 .request(DaemonAction::ReviewerAction {
1699 session_id,
1700 role,
1701 action,
1702 })
1703 .await?
1704 {
1705 DaemonReply::Reviewer(outcome) => Ok(*outcome),
1706 reply => bail!("unexpected reviewer reply {reply:?}"),
1707 }
1708 }
1709
1710 pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
1711 match self
1712 .request(DaemonAction::SyncSession { session_id })
1713 .await?
1714 {
1715 DaemonReply::Done => Ok(()),
1716 reply => bail!("unexpected session sync reply {reply:?}"),
1717 }
1718 }
1719
1720 pub async fn respond_elicitation(
1721 &mut self,
1722 session_id: String,
1723 elicitation_id: String,
1724 response: ElicitationResponse,
1725 ) -> Result<()> {
1726 match self
1727 .request(DaemonAction::RespondElicitation {
1728 session_id,
1729 elicitation_id,
1730 response,
1731 })
1732 .await?
1733 {
1734 DaemonReply::Done => Ok(()),
1735 reply => bail!("unexpected elicitation reply {reply:?}"),
1736 }
1737 }
1738
1739 pub async fn native_agent_history(
1740 &mut self,
1741 owner: String,
1742 child: String,
1743 before: Option<(u64, String)>,
1744 ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
1745 match self
1746 .request(DaemonAction::NativeAgentHistory {
1747 owner,
1748 child,
1749 before,
1750 })
1751 .await?
1752 {
1753 DaemonReply::NativeAgentHistory(page) => Ok(page),
1754 reply => bail!("unexpected native agent history reply {reply:?}"),
1755 }
1756 }
1757
1758 pub async fn stop_background_task(
1759 &mut self,
1760 session_id: String,
1761 background_task_id: String,
1762 ) -> Result<()> {
1763 match self
1764 .request(DaemonAction::StopBackgroundTask {
1765 session_id,
1766 background_task_id,
1767 })
1768 .await?
1769 {
1770 DaemonReply::Done => Ok(()),
1771 reply => bail!("unexpected background task stop reply {reply:?}"),
1772 }
1773 }
1774
1775 pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
1776 self.suspend_session_with_ack(session_id, false).await
1777 }
1778
1779 pub async fn suspend_session_with_ack(
1780 &mut self,
1781 session_id: String,
1782 acknowledge_unpublished_work: bool,
1783 ) -> Result<()> {
1784 match self
1785 .request(DaemonAction::SuspendSession {
1786 session_id,
1787 acknowledge_unpublished_work,
1788 })
1789 .await?
1790 {
1791 DaemonReply::Done => Ok(()),
1792 reply => bail!("unexpected close-session reply {reply:?}"),
1793 }
1794 }
1795
1796 pub async fn start_create_session(
1797 &mut self,
1798 request: CreateSessionRequest,
1799 ) -> Result<RegisteredSession> {
1800 match self
1801 .request(DaemonAction::StartCreateSession(request))
1802 .await?
1803 {
1804 DaemonReply::RegisteredSession(registered) => Ok(*registered),
1805 reply => bail!("unexpected start-create reply {reply:?}"),
1806 }
1807 }
1808
1809 pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
1810 match self
1811 .request(DaemonAction::WaitCreateSession { session_id })
1812 .await?
1813 {
1814 DaemonReply::Done => Ok(()),
1815 reply => bail!("unexpected wait-create reply {reply:?}"),
1816 }
1817 }
1818
1819 pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
1820 match self.request(DaemonAction::ResumeSession(request)).await? {
1821 DaemonReply::Done => Ok(()),
1822 reply => bail!("unexpected resume-session reply {reply:?}"),
1823 }
1824 }
1825
1826 pub async fn discard_since_checkpoint(
1827 &mut self,
1828 session_id: String,
1829 checkpoint: mj_core::state::CheckpointMetadata,
1830 ) -> Result<()> {
1831 match self
1832 .request(DaemonAction::DiscardSinceCheckpoint {
1833 session_id,
1834 checkpoint,
1835 })
1836 .await?
1837 {
1838 DaemonReply::Done => Ok(()),
1839 reply => bail!("unexpected force-stop reply {reply:?}"),
1840 }
1841 }
1842
1843 pub async fn destroy_stopped_session(
1844 &mut self,
1845 session_id: String,
1846 delete_branch: bool,
1847 ) -> Result<()> {
1848 match self
1849 .request(DaemonAction::DestroyStoppedSession {
1850 session_id,
1851 delete_branch,
1852 })
1853 .await?
1854 {
1855 DaemonReply::Done => Ok(()),
1856 reply => bail!("unexpected destroy-stopped reply {reply:?}"),
1857 }
1858 }
1859
1860 pub async fn force_destroy_session(
1861 &mut self,
1862 session_id: String,
1863 delete_branch: bool,
1864 ) -> Result<()> {
1865 match self
1866 .request(DaemonAction::ForceDestroySession {
1867 session_id,
1868 delete_branch,
1869 })
1870 .await?
1871 {
1872 DaemonReply::Done => Ok(()),
1873 reply => bail!("unexpected force-destroy reply {reply:?}"),
1874 }
1875 }
1876
1877 pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
1878 match self
1879 .request(DaemonAction::CancelLifecycle { session_id })
1880 .await?
1881 {
1882 DaemonReply::Done => Ok(()),
1883 reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
1884 }
1885 }
1886
1887 pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
1888 match self
1889 .request(DaemonAction::RecoverDraft { draft_id })
1890 .await?
1891 {
1892 DaemonReply::Done => Ok(()),
1893 reply => bail!("unexpected recover-draft reply {reply:?}"),
1894 }
1895 }
1896
1897 pub async fn stop(&mut self) -> Result<()> {
1898 match self.request(DaemonAction::Stop).await? {
1899 DaemonReply::Done => Ok(()),
1900 reply => bail!("unexpected stop reply {reply:?}"),
1901 }
1902 }
1903}
1904
1905pub async fn connect_existing() -> Result<DaemonClient> {
1906 let metadata = tokio::task::spawn_blocking(read_metadata)
1907 .await
1908 .context("read daemon metadata task failed")??;
1909 DaemonClient::connect(metadata).await
1910}
1911
1912pub struct ManagementClient {
1916 inner: DaemonClient,
1917}
1918
1919impl ManagementClient {
1920 pub fn new(inner: DaemonClient) -> Self {
1921 Self { inner }
1922 }
1923 pub fn protocol_version(&self) -> u32 {
1924 self.inner.metadata.protocol_version
1925 }
1926
1927 pub async fn status(&mut self) -> Result<DaemonStatus> {
1928 self.inner.status().await
1929 }
1930
1931 pub async fn stop(&mut self) -> Result<()> {
1932 tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
1933 .await
1934 .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
1935 }
1936
1937 pub async fn stop_and_wait(mut self) -> Result<()> {
1939 let pid = self.inner.metadata.pid;
1940 self.stop().await?;
1941 wait_for_exit(pid).await.with_context(|| {
1942 format!(
1943 "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
1944 STOP_TIMEOUT.as_secs()
1945 )
1946 })
1947 }
1948}
1949
1950pub async fn connect_management() -> Result<ManagementClient> {
1951 Ok(ManagementClient {
1952 inner: DaemonClient::connect(read_metadata_any()?).await?,
1953 })
1954}
1955
1956pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
1964 ensure!(
1965 metadata.protocol_version <= PROTOCOL_VERSION,
1966 "{}",
1967 unsupported_daemon_protocol_message(
1968 metadata.protocol_version,
1969 &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
1970 )
1971 );
1972 Ok(())
1973}
1974
1975fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
1978 format!(
1979 "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
1980 Put the daemon's directory first on PATH, or reinstall this client from that build."
1981 )
1982}
1983pub const PROTOCOL_VERSION: u32 = 36;
1984pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
1985pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
1995pub const RETRY_DELAY: Duration = Duration::from_millis(40);
1996impl DaemonClient {
1997 pub async fn prepare_move_session(
1998 &mut self,
1999 selection: MoveSelection,
2000 ) -> Result<MovePreparation> {
2001 match self
2002 .request(DaemonAction::PrepareMoveSession(selection))
2003 .await?
2004 {
2005 DaemonReply::MovePreparation(preparation) => Ok(*preparation),
2006 _ => bail!("daemon returned an unexpected move preparation reply"),
2007 }
2008 }
2009
2010 pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
2011 match self.request(DaemonAction::MoveSession(request)).await? {
2012 DaemonReply::MoveOutcome(outcome) => Ok(outcome),
2013 _ => bail!("daemon returned an unexpected move reply"),
2014 }
2015 }
2016}
2017
2018#[cfg(test)]
2019mod tests {
2020 use super::*;
2021 use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
2022 use std::path::Path;
2023
2024 #[cfg(unix)]
2025 #[test]
2026 fn an_exited_unreaped_process_is_a_zombie_and_not_alive() {
2027 use std::time::{Duration, Instant};
2031
2032 let mut running = std::process::Command::new("sleep")
2033 .arg("30")
2034 .spawn()
2035 .unwrap();
2036 assert!(process_is_alive(running.id()));
2037 assert!(!process_is_zombie(running.id()));
2038 running.kill().unwrap();
2039 running.wait().unwrap();
2040
2041 let mut exited = std::process::Command::new("sh")
2042 .args(["-c", "exit 0"])
2043 .spawn()
2044 .unwrap();
2045 let pid = exited.id();
2046 let deadline = Instant::now() + Duration::from_secs(10);
2047 while !process_is_zombie(pid) {
2048 assert!(
2049 Instant::now() < deadline,
2050 "an exited, unreaped child was never reported as a zombie"
2051 );
2052 std::thread::sleep(Duration::from_millis(20));
2053 }
2054 assert!(!process_is_alive(pid));
2055 exited.wait().unwrap();
2056 }
2057
2058 #[test]
2059 fn a_missing_endpoint_file_reads_as_a_stopped_daemon() {
2060 let directory = tempfile::tempdir().unwrap();
2061 let path = directory.path().join("daemon.json");
2062 let error = read_metadata_at(&path).unwrap_err();
2063 let stopped = daemon_not_running(&error).expect("stopped, not an I/O failure");
2064 assert_eq!(stopped.metadata_path, path);
2065 assert_eq!(format!("{error:#}"), "the Mjolnir daemon is not running");
2066
2067 std::fs::write(&path, b"not json").unwrap();
2068 let error = read_metadata_at(&path).unwrap_err();
2069 assert!(daemon_not_running(&error).is_none());
2070 }
2071
2072 #[tokio::test]
2073 async fn upgrade_refusal_retries_the_identical_command_but_lost_acknowledgements_do_not() {
2074 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2075 let metadata = DaemonMetadata {
2076 protocol_version: PROTOCOL_VERSION,
2077 pid: std::process::id(),
2078 address: listener.local_addr().unwrap(),
2079 token: "test-token".into(),
2080 started_at: "test".into(),
2081 build_version: env!("CARGO_PKG_VERSION").into(),
2082 };
2083 let action = DaemonAction::SubmitSessionCommand {
2084 inherited_draft: None,
2085 session_id: "test-session".into(),
2086 command_id: "steer-command".into(),
2087 command: RelayCommand::Steer {
2088 active_prompt_id: "active-command".into(),
2089 queued_prompt_id: "queued-command".into(),
2090 },
2091 };
2092 let expected = serde_json::to_value(&action).unwrap();
2093 let server = tokio::spawn(async move {
2094 for reply in [
2095 Some(DaemonReply::UpgradePending),
2096 Some(DaemonReply::Done),
2097 None,
2098 ] {
2099 let (mut stream, _) = listener.accept().await.unwrap();
2100 let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2101 assert_eq!(serde_json::to_value(&request.action).unwrap(), expected);
2102 if let Some(reply) = reply {
2103 write_frame(
2104 &mut stream,
2105 &ResponseEnvelope {
2106 protocol_version: request.protocol_version,
2107 request_id: request.request_id,
2108 result: Ok(reply),
2109 },
2110 )
2111 .await
2112 .unwrap();
2113 }
2114 }
2115 });
2116 let mut client = DaemonClient::connect(metadata.clone()).await.unwrap();
2117 let reply = client
2118 .request_with_reconnect(action.clone(), || DaemonClient::connect(metadata.clone()))
2119 .await
2120 .unwrap();
2121 assert!(matches!(reply, DaemonReply::Done));
2122 let mut client = DaemonClient::connect(metadata).await.unwrap();
2123 assert!(
2124 client
2125 .request_with_reconnect(action, || async {
2126 panic!("an ambiguous acknowledgement must not replay a mutation");
2127 })
2128 .await
2129 .is_err()
2130 );
2131 server.await.unwrap();
2132 }
2133
2134 #[test]
2135 fn unsupported_protocol_message_names_both_binaries_and_versions() {
2136 let message = unsupported_daemon_protocol_message(
2137 PROTOCOL_VERSION + 5,
2138 &describe_daemon_and_client_builds(
2139 4242,
2140 BuildDescription {
2141 executable: Some(Path::new("/home/dev/mj/target/release/mj")),
2142 version: "2.14.0",
2143 },
2144 BuildDescription {
2145 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2146 version: "2.9.0",
2147 },
2148 ),
2149 );
2150 assert_eq!(
2151 message,
2152 format!(
2153 "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
2154 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
2155 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
2156 Put the daemon's directory first on PATH, or reinstall this client from that build.",
2157 PROTOCOL_VERSION + 5
2158 )
2159 );
2160 }
2161
2162 #[test]
2163 fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
2164 let message = unsupported_daemon_protocol_message(
2165 PROTOCOL_VERSION + 1,
2166 &describe_daemon_and_client_builds(
2167 4242,
2168 BuildDescription {
2169 executable: None,
2170 version: "2.14.0",
2171 },
2172 BuildDescription {
2173 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2174 version: "2.9.0",
2175 },
2176 ),
2177 );
2178 assert!(
2179 message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
2180 "{message}"
2181 );
2182 assert!(
2183 message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
2184 "{message}"
2185 );
2186 }
2187}