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