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, skip_serializing_if = "Option::is_none")]
361 pub checkout: Option<mj_core::remote_git::ExactCheckout>,
362 #[serde(default, skip_serializing_if = "Option::is_none")]
363 pub expected_runtime_identity: Option<String>,
364 #[serde(default)]
366 pub mjolnir_subagents: Option<bool>,
367 #[serde(default)]
368 pub initial_prompt: Option<String>,
369 pub workspace_id: String,
370 pub profile_id: String,
371 pub bundle_id: String,
372 pub project_directory: Option<PathBuf>,
373 pub target_template_id: String,
374 pub additional_mounts: Vec<AdditionalMount>,
375 pub resource_allocation: Option<SessionResourceAllocation>,
376 pub title: String,
377 pub session_title_override: Option<String>,
378}
379
380#[derive(Debug, Clone, Serialize, Deserialize)]
381#[serde(deny_unknown_fields)]
382pub struct RegisteredSession {
383 pub session: SessionRecord,
384 pub remembered_container_size: Option<(String, HostContainerSize)>,
385}
386
387#[derive(Debug, Clone, Serialize, Deserialize)]
388#[serde(deny_unknown_fields)]
389pub struct DraftPreview {
390 pub id: String,
391 pub session_id: Option<String>,
392 pub source: String,
393 pub owner_pid: Option<u32>,
394 pub saved_at: String,
395}
396
397#[derive(Debug, Clone, Serialize, Deserialize)]
406#[serde(rename_all = "snake_case", tag = "action", content = "arguments")]
407pub enum DaemonAction {
408 NativeAgentHistory {
409 owner: String,
410 child: String,
411 before: Option<(u64, String)>,
412 },
413 Ping,
414 PrepareUpgrade,
417 UpgradeBlockers,
426 Status,
427 WebViewerAccess,
428 RecoverWebViewer(crate::web::WebViewerRecovery),
429 InspectWebListener,
430 ListWorkspaces,
431 CreateWorkspace {
432 name: String,
433 },
434 RenameWorkspace {
435 workspace_id: String,
436 name: String,
437 },
438 TouchWorkspace {
439 workspace_id: String,
440 },
441 CloseWorkspace {
442 workspace_id: String,
443 },
444 CancelWorkspaceClose {
445 workspace_id: String,
446 },
447 DeleteWorkspace {
448 workspace_id: String,
449 },
450 Attach {
451 client_id: String,
452 pid: u32,
453 },
454 Detach {
455 client_id: String,
456 },
457 PersistReadReceipt {
458 client_id: String,
459 workspace_id: String,
460 session_id: String,
461 through: u64,
462 },
463 PersistDetachedSessionState {
464 client_id: String,
465 workspace_id: String,
466 session_id: String,
467 through: u64,
468 owner_pid: u32,
469 draft: mj_core::storage::DetachedSessionDraft,
470 },
471 SaveActiveReview {
472 session_id: String,
473 review: mj_core::storage::StoredReview,
474 },
475 ClearActiveReview {
476 session_id: String,
477 },
478 SaveWorkspacePaneSizes {
479 workspace_id: String,
480 sizes: mj_core::workspace::PaneSizes,
481 },
482 SaveWorkspaceLayout {
483 workspace_id: String,
484 layout: mj_core::workspace::ConversationLayout,
485 },
486 PersistImportedSession {
487 session: Box<SessionRecord>,
488 },
489 SetSessionTitle {
490 session_id: String,
491 title: String,
492 },
493 SetSessionContainerSettings {
494 session_id: String,
495 cpus: Option<String>,
496 memory: Option<String>,
497 mounts: Vec<AdditionalMount>,
498 mount_history: Vec<PathBuf>,
499 },
500 SetSessionAcpTitle {
501 session_id: String,
502 title: Option<String>,
503 },
504 MarkSessionTargetMissing {
505 session_id: String,
506 detail: String,
507 updated_at: String,
508 },
509 CheckpointSession {
510 session_id: String,
511 },
512 WikiSearch {
515 query: String,
516 limit: usize,
517 },
518 WikiBrief {
520 wiki_id: String,
521 max_chars: usize,
522 },
523 WikiHits {
525 wiki_id: String,
526 query: String,
527 context_messages: usize,
528 per_message_chars: usize,
529 },
530 WikiSession {
532 wiki_id: String,
533 },
534 WikiRestore(WikiRestoreRequest),
536 ScanRecovery {
537 all_instances: bool,
538 },
539 AdoptRecovery {
540 session_id: String,
541 target_id: String,
542 profile: Option<String>,
543 bundle: Option<String>,
544 all_instances: bool,
545 },
546 DestroyRecovery {
547 session_id: String,
548 target_id: String,
549 confirmation: String,
550 all_instances: bool,
551 },
552 Snapshot {
553 workspace_id: String,
554 },
555 RuntimeSnapshot {
556 workspace_id: String,
557 after_revision: u64,
558 #[serde(default)]
559 all_workspaces: bool,
560 },
561 RenameProfile {
562 old_id: String,
563 new_id: String,
564 },
565 RenameTarget {
566 old_id: String,
567 new_id: String,
568 },
569 SubmitSessionCommand {
570 #[serde(default)]
571 inherited_draft: Option<String>,
572 session_id: String,
573 command_id: String,
574 command: RelayCommand,
575 },
576 QueueStartupPrompt {
581 session_id: String,
582 text: String,
583 #[serde(default)]
587 inherited_draft: Option<String>,
588 },
589 SyncSession {
590 session_id: String,
591 },
592 RespondElicitation {
593 session_id: String,
594 elicitation_id: String,
595 response: ElicitationResponse,
596 },
597 StopBackgroundTask {
598 session_id: String,
599 background_task_id: String,
600 },
601 ReviewerAction {
605 session_id: String,
606 #[serde(default, skip_serializing_if = "Option::is_none")]
609 role: Option<String>,
610 action: crate::session::ReviewerAction,
611 },
612 StartTurnReview {
614 session_id: String,
615 },
616 ResolveTurnReview {
618 session_id: String,
619 resolution: Resolution,
620 },
621 SuspendSession {
622 session_id: String,
623 #[serde(default)]
624 acknowledge_unpublished_work: bool,
625 },
626 StartCreateSession(CreateSessionRequest),
627 WaitCreateSession {
628 session_id: String,
629 },
630 ResumeSession(ResumeSessionRequest),
631 PrepareMoveSession(MoveSelection),
632 MoveSession(MoveSessionRequest),
633 DiscardSinceCheckpoint {
634 session_id: String,
635 checkpoint: mj_core::state::CheckpointMetadata,
636 },
637 DestroyStoppedSession {
638 session_id: String,
639 delete_branch: bool,
643 },
644 ForceDestroySession {
645 session_id: String,
646 delete_branch: bool,
648 },
649 CancelLifecycle {
650 session_id: String,
651 },
652 RecoverDraft {
653 draft_id: String,
654 },
655 Stop,
656}
657
658#[derive(Debug, Serialize, Deserialize)]
659#[serde(deny_unknown_fields)]
660pub struct RequestEnvelope {
661 pub protocol_version: u32,
662 pub request_id: u64,
663 pub token: String,
664 pub action: DaemonAction,
665}
666
667#[derive(Debug)]
671pub struct DaemonRefusal(pub String);
672
673impl std::fmt::Display for DaemonRefusal {
674 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
675 f.write_str(&self.0)
676 }
677}
678
679impl std::error::Error for DaemonRefusal {}
680
681impl DaemonRefusal {
682 #[must_use]
685 pub fn delivery_unconfirmed(&self) -> bool {
686 self.0
687 .starts_with(&crate::session::DeliveryUnconfirmed.to_string())
688 }
689}
690
691#[derive(Debug, Serialize, Deserialize)]
692#[serde(deny_unknown_fields)]
693pub struct ResponseEnvelope {
694 pub protocol_version: u32,
695 pub request_id: u64,
696 pub result: std::result::Result<DaemonReply, String>,
697}
698
699#[derive(Debug, Clone, Serialize, Deserialize)]
700#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
701pub enum DaemonReply {
702 NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
703 Pong,
704 UpgradePending,
705 UpgradeBlockers(Vec<String>),
708 Status(DaemonStatus),
709 WebViewerAccess(crate::web::WebViewerAccess),
710 WebListeners(Vec<crate::web::WebListenerProcess>),
711 Workspaces(Vec<WorkspaceListing>),
712 Workspace(WorkspaceRecord),
713 Snapshot(WorkspaceSnapshot),
714 RuntimeSnapshot(Box<RuntimeSnapshot>),
715 RegisteredSession(Box<RegisteredSession>),
716 MovePreparation(Box<MovePreparation>),
717 MoveOutcome(MoveOutcome),
718 Ordinal(u64),
719 Text(String),
720 OptionalSessionState(Option<SessionState>),
721 Checkpoint(mj_core::state::CheckpointMetadata),
722 RecoveryScan(mj_core::state::RecoveryScan),
723 WikiRows(WikiSearchPage),
724 WikiHits(Option<WikiHitTranscript>),
725 WikiSession(Option<Box<WikiSessionInfo>>),
726 Reviewer(Box<crate::session::ReviewerOutcome>),
727 Done,
728}
729
730#[derive(Debug, Clone, Serialize, Deserialize)]
731#[serde(deny_unknown_fields)]
732pub struct DaemonStatus {
733 pub pid: u32,
734 pub started_at: String,
735 pub build_version: String,
736 pub attached_clients: usize,
737 pub phone_status: WebViewerStatus,
738}
739
740#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
741#[serde(rename_all = "snake_case", tag = "state")]
742pub enum WebViewerStatus {
743 Disabled,
744 Starting,
745 Ready {
746 viewer_url: String,
747 viewer_code: String,
748 qr_login_url: Option<String>,
749 fallback_reason: Option<String>,
750 },
751 Stopped,
752 Error {
753 message: String,
754 },
755}
756
757impl std::fmt::Debug for WebViewerStatus {
758 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
759 match self {
760 Self::Ready {
761 viewer_url,
762 viewer_code,
763 fallback_reason,
764 ..
765 } => formatter
766 .debug_struct("Ready")
767 .field("viewer_url", viewer_url)
768 .field("viewer_code", viewer_code)
769 .field("qr_login_url", &"[redacted]")
770 .field("fallback_reason", fallback_reason)
771 .finish(),
772 Self::Disabled => formatter.write_str("Disabled"),
773 Self::Starting => formatter.write_str("Starting"),
774 Self::Stopped => formatter.write_str("Stopped"),
775 Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
776 }
777 }
778}
779
780impl std::fmt::Display for WebViewerStatus {
781 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
782 match self {
783 Self::Disabled => formatter.write_str("disabled"),
784 Self::Starting => formatter.write_str("starting"),
785 Self::Stopped => formatter.write_str("stopped unexpectedly"),
786 Self::Error { message } => write!(formatter, "error: {message}"),
787 Self::Ready {
788 viewer_url,
789 viewer_code,
790 fallback_reason,
791 ..
792 } => {
793 write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
794 if let Some(reason) = fallback_reason {
795 write!(
796 formatter,
797 "; local only because Tailscale HTTPS is unavailable: {reason}"
798 )?;
799 }
800 Ok(())
801 }
802 }
803 }
804}
805
806#[cfg(target_os = "linux")]
825pub fn process_is_zombie(pid: u32) -> bool {
826 let Ok(stat) = std::fs::read(format!("/proc/{pid}/stat")) else {
827 return false;
828 };
829 stat.iter()
832 .rposition(|byte| *byte == b')')
833 .and_then(|end| stat.get(end + 2))
834 .is_some_and(|state| *state == b'Z')
835}
836
837#[cfg(target_os = "macos")]
843pub fn process_is_zombie(pid: u32) -> bool {
844 let Ok(raw_pid) = libc::c_int::try_from(pid) else {
845 return false;
846 };
847 let size = std::mem::size_of::<libc::proc_bsdinfo>();
848 let mut info: libc::proc_bsdinfo = unsafe { std::mem::zeroed() };
850 let written = unsafe {
852 libc::proc_pidinfo(
853 raw_pid,
854 libc::PROC_PIDTBSDINFO,
855 1,
856 (&mut info as *mut libc::proc_bsdinfo).cast(),
857 size as libc::c_int,
858 )
859 };
860 usize::try_from(written).is_ok_and(|written| written == size) && info.pbi_status == libc::SZOMB
861}
862
863#[cfg(all(unix, not(target_os = "linux"), not(target_os = "macos")))]
864pub fn process_is_zombie(pid: u32) -> bool {
865 let pid = sysinfo::Pid::from_u32(pid);
866 let mut system = sysinfo::System::new();
867 system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
868 system
869 .process(pid)
870 .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
871}
872
873pub async fn wait_for_exit(pid: u32) -> Result<()> {
879 let deadline = Instant::now() + STOP_TIMEOUT;
880 while daemon_process_is_alive(pid) {
881 ensure!(Instant::now() < deadline, "process {pid} is still running");
882 tokio::time::sleep(RETRY_DELAY).await;
883 }
884 Ok(())
885}
886
887pub fn daemon_process_is_alive(pid: u32) -> bool {
893 #[cfg(unix)]
894 {
895 if pid == 0 {
896 return false;
897 }
898 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
899 return false;
900 };
901 let mut status = 0;
902 let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
905 if waited == raw_pid {
906 return false;
907 }
908 if waited == 0 {
909 return true;
910 }
911 let wait_error = std::io::Error::last_os_error();
912 if wait_error.raw_os_error() != Some(libc::ECHILD) {
913 return true;
914 }
915
916 #[cfg(target_os = "macos")]
917 return owned_daemon_group_is_alive(raw_pid);
918
919 #[cfg(not(target_os = "macos"))]
920 process_is_alive(pid)
921 }
922 #[cfg(not(unix))]
923 process_is_alive(pid)
924}
925
926#[cfg(target_os = "macos")]
927pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
928 if unsafe { libc::kill(-pid, 0) } == 0 {
936 return true;
937 }
938 let error = std::io::Error::last_os_error();
939 !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
940}
941
942pub fn process_is_alive(pid: u32) -> bool {
943 #[cfg(unix)]
944 {
945 if pid == 0 {
946 return false;
947 }
948 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
949 return false;
950 };
951 let result = unsafe { libc::kill(raw_pid, 0) };
954 let exists =
955 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
956 exists && !process_is_zombie(pid)
957 }
958 #[cfg(not(unix))]
959 {
960 let process_id = sysinfo::Pid::from_u32(pid);
961 let mut system = sysinfo::System::new();
962 system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[process_id]), true);
963 system.process(process_id).is_some()
964 }
965}
966
967pub fn read_metadata() -> Result<DaemonMetadata> {
968 let metadata = read_metadata_any()?;
969 ensure!(
970 metadata.protocol_version == PROTOCOL_VERSION,
971 "daemon protocol {} is incompatible with client protocol {}",
972 metadata.protocol_version,
973 PROTOCOL_VERSION
974 );
975 Ok(metadata)
976}
977
978#[derive(Debug, Clone, PartialEq, Eq)]
984pub struct DaemonNotRunning {
985 pub metadata_path: std::path::PathBuf,
987}
988
989impl std::fmt::Display for DaemonNotRunning {
990 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
991 formatter.write_str("the Mjolnir daemon is not running")
992 }
993}
994
995impl std::error::Error for DaemonNotRunning {}
996
997pub fn daemon_not_running(error: &anyhow::Error) -> Option<&DaemonNotRunning> {
999 error
1000 .chain()
1001 .find_map(|cause| cause.downcast_ref::<DaemonNotRunning>())
1002}
1003
1004pub fn read_metadata_any() -> Result<DaemonMetadata> {
1005 read_metadata_at(&metadata_path())
1006}
1007
1008fn read_metadata_at(path: &std::path::Path) -> Result<DaemonMetadata> {
1009 let body = match fs::read(path) {
1010 Ok(body) => body,
1011 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1012 return Err(DaemonNotRunning {
1013 metadata_path: path.to_owned(),
1014 }
1015 .into());
1016 }
1017 Err(error) => return Err(error).with_context(|| format!("read {}", path.display())),
1018 };
1019 let metadata: DaemonMetadata =
1020 serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
1021 Ok(metadata)
1022}
1023
1024pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
1025 let body = serde_json::to_vec(value)?;
1026 write_encoded_frame(stream, &body).await
1027}
1028
1029pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
1030 ensure!(
1031 body.len() <= MAX_FRAME_BYTES,
1032 "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
1033 body.len()
1034 );
1035 stream.write_u32(body.len() as u32).await?;
1036 stream.write_all(body).await?;
1037 stream.flush().await?;
1038 Ok(())
1039}
1040
1041pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
1042 let length = stream.read_u32().await? as usize;
1043 ensure!(
1044 length <= MAX_FRAME_BYTES,
1045 "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
1046 );
1047 let mut body = vec![0_u8; length];
1048 stream.read_exact(&mut body).await?;
1049 serde_json::from_slice(&body).context("decode daemon frame")
1050}
1051
1052pub struct DaemonClient {
1053 metadata: DaemonMetadata,
1054 stream: TcpStream,
1055 next_request_id: u64,
1056}
1057
1058impl DaemonClient {
1059 pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
1060 let stream =
1061 tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
1062 .await
1063 .context("time out connecting to Mjolnir daemon")??;
1064 Ok(Self {
1065 metadata,
1066 stream,
1067 next_request_id: 1,
1068 })
1069 }
1070
1071 pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1075 self.request_with_reconnect(action, || async {
1076 loop {
1077 if let Ok(client) = connect_existing().await {
1078 return Ok(client);
1079 }
1080 tokio::time::sleep(Duration::from_millis(250)).await;
1081 }
1082 })
1083 .await
1084 }
1085
1086 async fn request_with_reconnect<F, Fut>(
1087 &mut self,
1088 action: DaemonAction,
1089 mut reconnect: F,
1090 ) -> Result<DaemonReply>
1091 where
1092 F: FnMut() -> Fut,
1093 Fut: Future<Output = Result<Self>>,
1094 {
1095 loop {
1096 let response = self.request_once(action.clone()).await?;
1097 if matches!(response, DaemonReply::UpgradePending)
1098 && !matches!(action, DaemonAction::PrepareUpgrade)
1099 {
1100 tokio::time::sleep(Duration::from_millis(100)).await;
1103 *self = reconnect().await?;
1104 } else {
1105 return Ok(response);
1106 }
1107 }
1108 }
1109
1110 async fn request_once(&mut self, action: DaemonAction) -> Result<DaemonReply> {
1111 let protocol_version = self.metadata.protocol_version;
1112 let request_id = self.next_request_id;
1113 self.next_request_id += 1;
1114 write_frame(
1115 &mut self.stream,
1116 &RequestEnvelope {
1117 protocol_version,
1118 request_id,
1119 token: self.metadata.token.clone(),
1120 action,
1121 },
1122 )
1123 .await?;
1124 let response: ResponseEnvelope = read_frame(&mut self.stream).await?;
1125 ensure!(
1126 response.protocol_version == protocol_version,
1127 "daemon changed protocol"
1128 );
1129 ensure!(
1130 response.request_id == request_id,
1131 "daemon crossed request IDs"
1132 );
1133 response
1134 .result
1135 .map_err(|message| anyhow::Error::new(DaemonRefusal(message)))
1136 }
1137
1138 pub async fn status(&mut self) -> Result<DaemonStatus> {
1139 match self.request(DaemonAction::Status).await? {
1140 DaemonReply::Status(status) => Ok(status),
1141 reply => bail!("unexpected daemon status reply {reply:?}"),
1142 }
1143 }
1144
1145 pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
1146 match self.request(DaemonAction::WebViewerAccess).await? {
1147 DaemonReply::WebViewerAccess(access) => Ok(access),
1148 reply => bail!("unexpected web viewer reply {reply:?}"),
1149 }
1150 }
1151
1152 pub async fn recover_web_viewer(
1153 &mut self,
1154 action: crate::web::WebViewerRecovery,
1155 ) -> Result<()> {
1156 match self.request(DaemonAction::RecoverWebViewer(action)).await? {
1157 DaemonReply::Done => Ok(()),
1158 reply => bail!("unexpected web viewer recovery reply {reply:?}"),
1159 }
1160 }
1161
1162 pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
1163 match self.request(DaemonAction::InspectWebListener).await? {
1164 DaemonReply::WebListeners(processes) => Ok(processes),
1165 reply => bail!("unexpected listener inspection reply {reply:?}"),
1166 }
1167 }
1168
1169 pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
1170 match self.request(DaemonAction::ListWorkspaces).await? {
1171 DaemonReply::Workspaces(workspaces) => Ok(workspaces),
1172 reply => bail!("unexpected daemon workspace reply {reply:?}"),
1173 }
1174 }
1175
1176 pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
1177 match self
1178 .request(DaemonAction::RenameProfile { old_id, new_id })
1179 .await?
1180 {
1181 DaemonReply::Done => Ok(()),
1182 reply => bail!("unexpected rename-profile reply {reply:?}"),
1183 }
1184 }
1185
1186 pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1187 match self
1188 .request(DaemonAction::RenameTarget { old_id, new_id })
1189 .await?
1190 {
1191 DaemonReply::Done => Ok(()),
1192 reply => bail!("unexpected rename-target reply {reply:?}"),
1193 }
1194 }
1195
1196 pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1197 match self.request(DaemonAction::CreateWorkspace { name }).await? {
1198 DaemonReply::Workspace(workspace) => Ok(workspace),
1199 reply => bail!("unexpected create-workspace reply {reply:?}"),
1200 }
1201 }
1202
1203 pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1204 match self
1205 .request(DaemonAction::RenameWorkspace { workspace_id, name })
1206 .await?
1207 {
1208 DaemonReply::Done => Ok(()),
1209 reply => bail!("unexpected rename-workspace reply {reply:?}"),
1210 }
1211 }
1212
1213 pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1214 match self
1215 .request(DaemonAction::TouchWorkspace { workspace_id })
1216 .await?
1217 {
1218 DaemonReply::Done => Ok(()),
1219 reply => bail!("unexpected touch-workspace reply {reply:?}"),
1220 }
1221 }
1222
1223 pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1224 match self
1225 .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1226 .await?
1227 {
1228 DaemonReply::Done => Ok(()),
1229 reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1230 }
1231 }
1232
1233 pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1234 match self
1235 .request(DaemonAction::CloseWorkspace { workspace_id })
1236 .await?
1237 {
1238 DaemonReply::Done => Ok(()),
1239 reply => bail!("unexpected close workspace reply: {reply:?}"),
1240 }
1241 }
1242
1243 pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1244 match self
1245 .request(DaemonAction::DeleteWorkspace { workspace_id })
1246 .await?
1247 {
1248 DaemonReply::Done => Ok(()),
1249 reply => bail!("unexpected delete-workspace reply {reply:?}"),
1250 }
1251 }
1252
1253 pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1254 match self
1255 .request(DaemonAction::Attach { client_id, pid })
1256 .await?
1257 {
1258 DaemonReply::Done => Ok(()),
1259 reply => bail!("unexpected attach reply {reply:?}"),
1260 }
1261 }
1262
1263 pub async fn detach(&mut self, client_id: String) -> Result<()> {
1264 match self.request(DaemonAction::Detach { client_id }).await? {
1265 DaemonReply::Done => Ok(()),
1266 reply => bail!("unexpected detach reply {reply:?}"),
1267 }
1268 }
1269
1270 pub async fn persist_read_receipt(
1271 &mut self,
1272 client_id: String,
1273 workspace_id: String,
1274 session_id: String,
1275 through: u64,
1276 ) -> Result<u64> {
1277 match self
1278 .request(DaemonAction::PersistReadReceipt {
1279 client_id,
1280 workspace_id,
1281 session_id,
1282 through,
1283 })
1284 .await?
1285 {
1286 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1287 reply => bail!("unexpected read-receipt reply {reply:?}"),
1288 }
1289 }
1290
1291 pub async fn persist_detached_session_state(
1292 &mut self,
1293 client_id: String,
1294 workspace_id: String,
1295 session_id: String,
1296 through: u64,
1297 owner_pid: u32,
1298 draft: mj_core::storage::DetachedSessionDraft,
1299 ) -> Result<()> {
1300 match self
1301 .request(DaemonAction::PersistDetachedSessionState {
1302 client_id,
1303 workspace_id,
1304 session_id,
1305 through,
1306 owner_pid,
1307 draft,
1308 })
1309 .await?
1310 {
1311 DaemonReply::Done => Ok(()),
1312 reply => bail!("unexpected detached-session-state reply {reply:?}"),
1313 }
1314 }
1315
1316 pub async fn save_active_review(
1317 &mut self,
1318 session_id: String,
1319 review: mj_core::storage::StoredReview,
1320 ) -> Result<()> {
1321 match self
1322 .request(DaemonAction::SaveActiveReview { session_id, review })
1323 .await?
1324 {
1325 DaemonReply::Done => Ok(()),
1326 reply => bail!("unexpected save-review reply {reply:?}"),
1327 }
1328 }
1329
1330 pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1331 match self
1332 .request(DaemonAction::ClearActiveReview { session_id })
1333 .await?
1334 {
1335 DaemonReply::Done => Ok(()),
1336 reply => bail!("unexpected clear-review reply {reply:?}"),
1337 }
1338 }
1339
1340 pub async fn save_workspace_pane_sizes(
1341 &mut self,
1342 workspace_id: String,
1343 sizes: mj_core::workspace::PaneSizes,
1344 ) -> Result<()> {
1345 match self
1346 .request(DaemonAction::SaveWorkspacePaneSizes {
1347 workspace_id,
1348 sizes,
1349 })
1350 .await?
1351 {
1352 DaemonReply::Done => Ok(()),
1353 reply => bail!("unexpected pane-size save reply {reply:?}"),
1354 }
1355 }
1356
1357 pub async fn save_workspace_layout(
1358 &mut self,
1359 workspace_id: String,
1360 layout: mj_core::workspace::ConversationLayout,
1361 ) -> Result<()> {
1362 match self
1363 .request(DaemonAction::SaveWorkspaceLayout {
1364 workspace_id,
1365 layout,
1366 })
1367 .await?
1368 {
1369 DaemonReply::Done => Ok(()),
1370 reply => bail!("unexpected layout save reply {reply:?}"),
1371 }
1372 }
1373
1374 pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1375 match self
1376 .request(DaemonAction::PersistImportedSession {
1377 session: Box::new(session),
1378 })
1379 .await?
1380 {
1381 DaemonReply::Done => Ok(()),
1382 reply => bail!("unexpected imported-session reply {reply:?}"),
1383 }
1384 }
1385
1386 pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1387 match self
1388 .request(DaemonAction::SetSessionTitle { session_id, title })
1389 .await?
1390 {
1391 DaemonReply::Text(title) => Ok(title),
1392 reply => bail!("unexpected session-title reply {reply:?}"),
1393 }
1394 }
1395
1396 pub async fn set_session_container_settings(
1397 &mut self,
1398 session_id: String,
1399 cpus: Option<String>,
1400 memory: Option<String>,
1401 mounts: Vec<AdditionalMount>,
1402 mount_history: Vec<PathBuf>,
1403 ) -> Result<()> {
1404 match self
1405 .request(DaemonAction::SetSessionContainerSettings {
1406 session_id,
1407 cpus,
1408 memory,
1409 mounts,
1410 mount_history,
1411 })
1412 .await?
1413 {
1414 DaemonReply::Done => Ok(()),
1415 reply => bail!("unexpected container-settings reply {reply:?}"),
1416 }
1417 }
1418
1419 pub async fn set_session_acp_title(
1420 &mut self,
1421 session_id: String,
1422 title: Option<String>,
1423 ) -> Result<()> {
1424 match self
1425 .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1426 .await?
1427 {
1428 DaemonReply::Done => Ok(()),
1429 reply => bail!("unexpected ACP-title reply {reply:?}"),
1430 }
1431 }
1432
1433 pub async fn mark_session_target_missing(
1434 &mut self,
1435 session_id: String,
1436 detail: String,
1437 updated_at: String,
1438 ) -> Result<Option<SessionState>> {
1439 match self
1440 .request(DaemonAction::MarkSessionTargetMissing {
1441 session_id,
1442 detail,
1443 updated_at,
1444 })
1445 .await?
1446 {
1447 DaemonReply::OptionalSessionState(state) => Ok(state),
1448 reply => bail!("unexpected target-missing reply {reply:?}"),
1449 }
1450 }
1451
1452 pub async fn checkpoint_session(
1453 &mut self,
1454 session_id: String,
1455 ) -> Result<mj_core::state::CheckpointMetadata> {
1456 match self
1457 .request(DaemonAction::CheckpointSession { session_id })
1458 .await?
1459 {
1460 DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1461 reply => bail!("unexpected checkpoint reply {reply:?}"),
1462 }
1463 }
1464
1465 pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1470 match self
1471 .request(DaemonAction::WikiSearch { query, limit })
1472 .await?
1473 {
1474 DaemonReply::WikiRows(page) => Ok(page),
1475 reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1476 }
1477 }
1478
1479 pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1481 match self
1482 .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1483 .await?
1484 {
1485 DaemonReply::Text(markdown) => Ok(markdown),
1486 reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1487 }
1488 }
1489
1490 pub async fn wiki_hits(
1495 &mut self,
1496 wiki_id: String,
1497 query: String,
1498 context_messages: usize,
1499 per_message_chars: usize,
1500 ) -> Result<Option<WikiHitTranscript>> {
1501 match self
1502 .request(DaemonAction::WikiHits {
1503 wiki_id,
1504 query,
1505 context_messages,
1506 per_message_chars,
1507 })
1508 .await?
1509 {
1510 DaemonReply::WikiHits(transcript) => Ok(transcript),
1511 reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1512 }
1513 }
1514
1515 pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1518 match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1519 DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1520 reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1521 }
1522 }
1523
1524 pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1528 match self.request(DaemonAction::WikiRestore(request)).await? {
1529 DaemonReply::RegisteredSession(registered) => Ok(*registered),
1530 reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1531 }
1532 }
1533
1534 pub async fn scan_recovery(
1535 &mut self,
1536 all_instances: bool,
1537 ) -> Result<mj_core::state::RecoveryScan> {
1538 match self
1539 .request(DaemonAction::ScanRecovery { all_instances })
1540 .await?
1541 {
1542 DaemonReply::RecoveryScan(scan) => Ok(scan),
1543 reply => bail!("unexpected recovery-scan reply {reply:?}"),
1544 }
1545 }
1546
1547 pub async fn adopt_recovery(
1548 &mut self,
1549 session_id: String,
1550 target_id: String,
1551 profile: Option<String>,
1552 bundle: Option<String>,
1553 all_instances: bool,
1554 ) -> Result<()> {
1555 match self
1556 .request(DaemonAction::AdoptRecovery {
1557 session_id,
1558 target_id,
1559 profile,
1560 bundle,
1561 all_instances,
1562 })
1563 .await?
1564 {
1565 DaemonReply::Done => Ok(()),
1566 reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1567 }
1568 }
1569
1570 pub async fn destroy_recovery(
1571 &mut self,
1572 session_id: String,
1573 target_id: String,
1574 confirmation: String,
1575 all_instances: bool,
1576 ) -> Result<()> {
1577 match self
1578 .request(DaemonAction::DestroyRecovery {
1579 session_id,
1580 target_id,
1581 confirmation,
1582 all_instances,
1583 })
1584 .await?
1585 {
1586 DaemonReply::Done => Ok(()),
1587 reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1588 }
1589 }
1590
1591 pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1592 match self
1593 .request(DaemonAction::Snapshot { workspace_id })
1594 .await?
1595 {
1596 DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1597 reply => bail!("unexpected snapshot reply {reply:?}"),
1598 }
1599 }
1600
1601 pub async fn runtime_snapshot(
1602 &mut self,
1603 workspace_id: String,
1604 after_revision: u64,
1605 all_workspaces: bool,
1606 ) -> Result<RuntimeSnapshot> {
1607 match self
1608 .request(DaemonAction::RuntimeSnapshot {
1609 workspace_id,
1610 after_revision,
1611 all_workspaces,
1612 })
1613 .await?
1614 {
1615 DaemonReply::RuntimeSnapshot(snapshot) => Ok(*snapshot),
1616 reply => bail!("unexpected runtime snapshot reply {reply:?}"),
1617 }
1618 }
1619
1620 pub async fn submit_session_command(
1621 &mut self,
1622 session_id: String,
1623 command_id: String,
1624 command: RelayCommand,
1625 inherited_draft: Option<String>,
1626 ) -> Result<u64> {
1627 match self
1628 .request(DaemonAction::SubmitSessionCommand {
1629 inherited_draft,
1630 session_id,
1631 command_id,
1632 command,
1633 })
1634 .await?
1635 {
1636 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1637 reply => bail!("unexpected session command reply {reply:?}"),
1638 }
1639 }
1640
1641 pub async fn queue_startup_prompt(
1645 &mut self,
1646 session_id: String,
1647 text: String,
1648 inherited_draft: Option<String>,
1649 ) -> Result<()> {
1650 match self
1651 .request(DaemonAction::QueueStartupPrompt {
1652 session_id,
1653 text,
1654 inherited_draft,
1655 })
1656 .await?
1657 {
1658 DaemonReply::Done => Ok(()),
1659 reply => bail!("unexpected startup prompt reply {reply:?}"),
1660 }
1661 }
1662
1663 pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
1669 match self
1670 .request(DaemonAction::StartTurnReview { session_id })
1671 .await?
1672 {
1673 DaemonReply::Done => Ok(()),
1674 reply => bail!("unexpected turn-review reply {reply:?}"),
1675 }
1676 }
1677
1678 pub async fn resolve_turn_review(
1679 &mut self,
1680 session_id: String,
1681 resolution: Resolution,
1682 ) -> Result<()> {
1683 match self
1684 .request(DaemonAction::ResolveTurnReview {
1685 session_id,
1686 resolution,
1687 })
1688 .await?
1689 {
1690 DaemonReply::Done => Ok(()),
1691 reply => bail!("unexpected turn-review resolution reply {reply:?}"),
1692 }
1693 }
1694
1695 pub async fn reviewer_action(
1696 &mut self,
1697 session_id: String,
1698 role: Option<String>,
1699 action: crate::session::ReviewerAction,
1700 ) -> Result<crate::session::ReviewerOutcome> {
1701 match self
1702 .request(DaemonAction::ReviewerAction {
1703 session_id,
1704 role,
1705 action,
1706 })
1707 .await?
1708 {
1709 DaemonReply::Reviewer(outcome) => Ok(*outcome),
1710 reply => bail!("unexpected reviewer reply {reply:?}"),
1711 }
1712 }
1713
1714 pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
1715 match self
1716 .request(DaemonAction::SyncSession { session_id })
1717 .await?
1718 {
1719 DaemonReply::Done => Ok(()),
1720 reply => bail!("unexpected session sync reply {reply:?}"),
1721 }
1722 }
1723
1724 pub async fn respond_elicitation(
1725 &mut self,
1726 session_id: String,
1727 elicitation_id: String,
1728 response: ElicitationResponse,
1729 ) -> Result<()> {
1730 match self
1731 .request(DaemonAction::RespondElicitation {
1732 session_id,
1733 elicitation_id,
1734 response,
1735 })
1736 .await?
1737 {
1738 DaemonReply::Done => Ok(()),
1739 reply => bail!("unexpected elicitation reply {reply:?}"),
1740 }
1741 }
1742
1743 pub async fn native_agent_history(
1744 &mut self,
1745 owner: String,
1746 child: String,
1747 before: Option<(u64, String)>,
1748 ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
1749 match self
1750 .request(DaemonAction::NativeAgentHistory {
1751 owner,
1752 child,
1753 before,
1754 })
1755 .await?
1756 {
1757 DaemonReply::NativeAgentHistory(page) => Ok(page),
1758 reply => bail!("unexpected native agent history reply {reply:?}"),
1759 }
1760 }
1761
1762 pub async fn stop_background_task(
1763 &mut self,
1764 session_id: String,
1765 background_task_id: String,
1766 ) -> Result<()> {
1767 match self
1768 .request(DaemonAction::StopBackgroundTask {
1769 session_id,
1770 background_task_id,
1771 })
1772 .await?
1773 {
1774 DaemonReply::Done => Ok(()),
1775 reply => bail!("unexpected background task stop reply {reply:?}"),
1776 }
1777 }
1778
1779 pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
1780 self.suspend_session_with_ack(session_id, false).await
1781 }
1782
1783 pub async fn suspend_session_with_ack(
1784 &mut self,
1785 session_id: String,
1786 acknowledge_unpublished_work: bool,
1787 ) -> Result<()> {
1788 match self
1789 .request(DaemonAction::SuspendSession {
1790 session_id,
1791 acknowledge_unpublished_work,
1792 })
1793 .await?
1794 {
1795 DaemonReply::Done => Ok(()),
1796 reply => bail!("unexpected close-session reply {reply:?}"),
1797 }
1798 }
1799
1800 pub async fn start_create_session(
1801 &mut self,
1802 request: CreateSessionRequest,
1803 ) -> Result<RegisteredSession> {
1804 match self
1805 .request(DaemonAction::StartCreateSession(request))
1806 .await?
1807 {
1808 DaemonReply::RegisteredSession(registered) => Ok(*registered),
1809 reply => bail!("unexpected start-create reply {reply:?}"),
1810 }
1811 }
1812
1813 pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
1814 match self
1815 .request(DaemonAction::WaitCreateSession { session_id })
1816 .await?
1817 {
1818 DaemonReply::Done => Ok(()),
1819 reply => bail!("unexpected wait-create reply {reply:?}"),
1820 }
1821 }
1822
1823 pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
1824 match self.request(DaemonAction::ResumeSession(request)).await? {
1825 DaemonReply::Done => Ok(()),
1826 reply => bail!("unexpected resume-session reply {reply:?}"),
1827 }
1828 }
1829
1830 pub async fn discard_since_checkpoint(
1831 &mut self,
1832 session_id: String,
1833 checkpoint: mj_core::state::CheckpointMetadata,
1834 ) -> Result<()> {
1835 match self
1836 .request(DaemonAction::DiscardSinceCheckpoint {
1837 session_id,
1838 checkpoint,
1839 })
1840 .await?
1841 {
1842 DaemonReply::Done => Ok(()),
1843 reply => bail!("unexpected force-stop reply {reply:?}"),
1844 }
1845 }
1846
1847 pub async fn destroy_stopped_session(
1848 &mut self,
1849 session_id: String,
1850 delete_branch: bool,
1851 ) -> Result<()> {
1852 match self
1853 .request(DaemonAction::DestroyStoppedSession {
1854 session_id,
1855 delete_branch,
1856 })
1857 .await?
1858 {
1859 DaemonReply::Done => Ok(()),
1860 reply => bail!("unexpected destroy-stopped reply {reply:?}"),
1861 }
1862 }
1863
1864 pub async fn force_destroy_session(
1865 &mut self,
1866 session_id: String,
1867 delete_branch: bool,
1868 ) -> Result<()> {
1869 match self
1870 .request(DaemonAction::ForceDestroySession {
1871 session_id,
1872 delete_branch,
1873 })
1874 .await?
1875 {
1876 DaemonReply::Done => Ok(()),
1877 reply => bail!("unexpected force-destroy reply {reply:?}"),
1878 }
1879 }
1880
1881 pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
1882 match self
1883 .request(DaemonAction::CancelLifecycle { session_id })
1884 .await?
1885 {
1886 DaemonReply::Done => Ok(()),
1887 reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
1888 }
1889 }
1890
1891 pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
1892 match self
1893 .request(DaemonAction::RecoverDraft { draft_id })
1894 .await?
1895 {
1896 DaemonReply::Done => Ok(()),
1897 reply => bail!("unexpected recover-draft reply {reply:?}"),
1898 }
1899 }
1900
1901 pub async fn stop(&mut self) -> Result<()> {
1902 match self.request(DaemonAction::Stop).await? {
1903 DaemonReply::Done => Ok(()),
1904 reply => bail!("unexpected stop reply {reply:?}"),
1905 }
1906 }
1907}
1908
1909pub async fn connect_existing() -> Result<DaemonClient> {
1910 let metadata = tokio::task::spawn_blocking(read_metadata)
1911 .await
1912 .context("read daemon metadata task failed")??;
1913 DaemonClient::connect(metadata).await
1914}
1915
1916pub struct ManagementClient {
1920 inner: DaemonClient,
1921}
1922
1923impl ManagementClient {
1924 pub fn new(inner: DaemonClient) -> Self {
1925 Self { inner }
1926 }
1927 pub fn protocol_version(&self) -> u32 {
1928 self.inner.metadata.protocol_version
1929 }
1930
1931 pub async fn status(&mut self) -> Result<DaemonStatus> {
1932 self.inner.status().await
1933 }
1934
1935 pub async fn stop(&mut self) -> Result<()> {
1936 tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
1937 .await
1938 .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
1939 }
1940
1941 pub async fn stop_and_wait(mut self) -> Result<()> {
1943 let pid = self.inner.metadata.pid;
1944 self.stop().await?;
1945 wait_for_exit(pid).await.with_context(|| {
1946 format!(
1947 "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
1948 STOP_TIMEOUT.as_secs()
1949 )
1950 })
1951 }
1952}
1953
1954pub async fn connect_management() -> Result<ManagementClient> {
1955 Ok(ManagementClient {
1956 inner: DaemonClient::connect(read_metadata_any()?).await?,
1957 })
1958}
1959
1960pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
1968 ensure!(
1969 metadata.protocol_version <= PROTOCOL_VERSION,
1970 "{}",
1971 unsupported_daemon_protocol_message(
1972 metadata.protocol_version,
1973 &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
1974 )
1975 );
1976 Ok(())
1977}
1978
1979fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
1982 format!(
1983 "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
1984 Put the daemon's directory first on PATH, or reinstall this client from that build."
1985 )
1986}
1987pub const PROTOCOL_VERSION: u32 = 38;
1988pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
1989pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
1999pub const RETRY_DELAY: Duration = Duration::from_millis(40);
2000impl DaemonClient {
2001 pub async fn prepare_move_session(
2002 &mut self,
2003 selection: MoveSelection,
2004 ) -> Result<MovePreparation> {
2005 match self
2006 .request(DaemonAction::PrepareMoveSession(selection))
2007 .await?
2008 {
2009 DaemonReply::MovePreparation(preparation) => Ok(*preparation),
2010 _ => bail!("daemon returned an unexpected move preparation reply"),
2011 }
2012 }
2013
2014 pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
2015 match self.request(DaemonAction::MoveSession(request)).await? {
2016 DaemonReply::MoveOutcome(outcome) => Ok(outcome),
2017 _ => bail!("daemon returned an unexpected move reply"),
2018 }
2019 }
2020}
2021
2022#[cfg(test)]
2023mod tests {
2024 use super::*;
2025 use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
2026 use std::path::Path;
2027
2028 #[cfg(unix)]
2029 #[test]
2030 fn an_exited_unreaped_process_is_a_zombie_and_not_alive() {
2031 use std::time::{Duration, Instant};
2035
2036 let mut running = std::process::Command::new("sleep")
2037 .arg("30")
2038 .spawn()
2039 .unwrap();
2040 assert!(process_is_alive(running.id()));
2041 assert!(!process_is_zombie(running.id()));
2042 running.kill().unwrap();
2043 running.wait().unwrap();
2044
2045 let mut exited = std::process::Command::new("sh")
2046 .args(["-c", "exit 0"])
2047 .spawn()
2048 .unwrap();
2049 let pid = exited.id();
2050 let deadline = Instant::now() + Duration::from_secs(10);
2051 while !process_is_zombie(pid) {
2052 assert!(
2053 Instant::now() < deadline,
2054 "an exited, unreaped child was never reported as a zombie"
2055 );
2056 std::thread::sleep(Duration::from_millis(20));
2057 }
2058 assert!(!process_is_alive(pid));
2059 exited.wait().unwrap();
2060 }
2061
2062 #[test]
2063 fn a_missing_endpoint_file_reads_as_a_stopped_daemon() {
2064 let directory = tempfile::tempdir().unwrap();
2065 let path = directory.path().join("daemon.json");
2066 let error = read_metadata_at(&path).unwrap_err();
2067 let stopped = daemon_not_running(&error).expect("stopped, not an I/O failure");
2068 assert_eq!(stopped.metadata_path, path);
2069 assert_eq!(format!("{error:#}"), "the Mjolnir daemon is not running");
2070
2071 std::fs::write(&path, b"not json").unwrap();
2072 let error = read_metadata_at(&path).unwrap_err();
2073 assert!(daemon_not_running(&error).is_none());
2074 }
2075
2076 #[tokio::test]
2077 async fn upgrade_refusal_retries_the_identical_command_but_lost_acknowledgements_do_not() {
2078 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2079 let metadata = DaemonMetadata {
2080 protocol_version: PROTOCOL_VERSION,
2081 pid: std::process::id(),
2082 address: listener.local_addr().unwrap(),
2083 token: "test-token".into(),
2084 started_at: "test".into(),
2085 build_version: env!("CARGO_PKG_VERSION").into(),
2086 };
2087 let action = DaemonAction::SubmitSessionCommand {
2088 inherited_draft: None,
2089 session_id: "test-session".into(),
2090 command_id: "steer-command".into(),
2091 command: RelayCommand::Steer {
2092 active_prompt_id: "active-command".into(),
2093 queued_prompt_id: "queued-command".into(),
2094 },
2095 };
2096 let expected = serde_json::to_value(&action).unwrap();
2097 let server = tokio::spawn(async move {
2098 for reply in [
2099 Some(DaemonReply::UpgradePending),
2100 Some(DaemonReply::Done),
2101 None,
2102 ] {
2103 let (mut stream, _) = listener.accept().await.unwrap();
2104 let request: RequestEnvelope = read_frame(&mut stream).await.unwrap();
2105 assert_eq!(serde_json::to_value(&request.action).unwrap(), expected);
2106 if let Some(reply) = reply {
2107 write_frame(
2108 &mut stream,
2109 &ResponseEnvelope {
2110 protocol_version: request.protocol_version,
2111 request_id: request.request_id,
2112 result: Ok(reply),
2113 },
2114 )
2115 .await
2116 .unwrap();
2117 }
2118 }
2119 });
2120 let mut client = DaemonClient::connect(metadata.clone()).await.unwrap();
2121 let reply = client
2122 .request_with_reconnect(action.clone(), || DaemonClient::connect(metadata.clone()))
2123 .await
2124 .unwrap();
2125 assert!(matches!(reply, DaemonReply::Done));
2126 let mut client = DaemonClient::connect(metadata).await.unwrap();
2127 assert!(
2128 client
2129 .request_with_reconnect(action, || async {
2130 panic!("an ambiguous acknowledgement must not replay a mutation");
2131 })
2132 .await
2133 .is_err()
2134 );
2135 server.await.unwrap();
2136 }
2137
2138 #[test]
2139 fn unsupported_protocol_message_names_both_binaries_and_versions() {
2140 let message = unsupported_daemon_protocol_message(
2141 PROTOCOL_VERSION + 5,
2142 &describe_daemon_and_client_builds(
2143 4242,
2144 BuildDescription {
2145 executable: Some(Path::new("/home/dev/mj/target/release/mj")),
2146 version: "2.14.0",
2147 },
2148 BuildDescription {
2149 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2150 version: "2.9.0",
2151 },
2152 ),
2153 );
2154 assert_eq!(
2155 message,
2156 format!(
2157 "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
2158 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
2159 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
2160 Put the daemon's directory first on PATH, or reinstall this client from that build.",
2161 PROTOCOL_VERSION + 5
2162 )
2163 );
2164 }
2165
2166 #[test]
2167 fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
2168 let message = unsupported_daemon_protocol_message(
2169 PROTOCOL_VERSION + 1,
2170 &describe_daemon_and_client_builds(
2171 4242,
2172 BuildDescription {
2173 executable: None,
2174 version: "2.14.0",
2175 },
2176 BuildDescription {
2177 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
2178 version: "2.9.0",
2179 },
2180 ),
2181 );
2182 assert!(
2183 message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
2184 "{message}"
2185 );
2186 assert!(
2187 message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
2188 "{message}"
2189 );
2190 }
2191}