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}
197
198#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
200#[serde(deny_unknown_fields)]
201pub struct WikiRestoreRequest {
202 pub wiki_id: String,
204 pub workspace_id: String,
205 pub profile_id: String,
206 pub target_template_id: String,
207 #[serde(default)]
210 pub project_directory: Option<PathBuf>,
211 #[serde(default)]
212 pub additional_mounts: Vec<AdditionalMount>,
213 #[serde(default)]
214 pub resource_allocation: Option<SessionResourceAllocation>,
215}
216
217#[derive(Debug, Clone, Serialize, Deserialize)]
218#[serde(deny_unknown_fields)]
219pub struct WorkspaceSnapshot {
220 pub workspace: WorkspaceRecord,
221 pub sessions: Vec<SessionPreview>,
222 pub drafts: Vec<DraftPreview>,
223}
224
225#[derive(Debug, Clone, Serialize, Deserialize)]
226#[serde(deny_unknown_fields)]
227pub struct RuntimeSessionView {
228 pub session_id: String,
229 pub projection_ordinal: u64,
230 pub projection_digest: String,
231 pub operational: Option<RelayOperationalState>,
232 pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
233 pub connected: bool,
234 pub error: Option<ViewError>,
235}
236
237impl RuntimeSessionView {
238 pub fn from_managed(session_id: String, view: ManagedSessionView) -> Self {
239 let (projection_ordinal, projection_digest, operational, signal) =
240 view.snapshot
241 .map_or((0, String::new(), None, None), |snapshot| {
242 (
243 snapshot.materialized.applied_event_ordinal,
244 snapshot.materialized.applied_event_digest,
245 Some(snapshot.operational),
246 snapshot.latest_credential_sync_signal,
247 )
248 });
249 Self {
250 session_id,
251 projection_ordinal,
252 projection_digest,
253 operational,
254 latest_credential_sync_signal: signal,
255 connected: view.connected,
256 error: view.error,
257 }
258 }
259}
260
261#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
267#[serde(deny_unknown_fields)]
268pub struct RuntimeNotice {
269 pub id: u64,
270 pub session_id: String,
271 pub text: String,
272}
273
274#[derive(Debug, Clone, Serialize, Deserialize)]
275#[serde(deny_unknown_fields)]
276pub struct RuntimeSnapshot {
277 #[serde(default)]
278 pub native_agents: Vec<mj_core::native_agent::NativeAgentSummary>,
279 #[serde(default)]
280 pub workspace_names: BTreeMap<String, String>,
281 #[serde(default)]
282 pub moves: Vec<mj_core::state::MoveOperation>,
283 pub revision: u64,
284 pub config: Config,
285 pub records: Vec<SessionRecord>,
286 pub sessions: Vec<RuntimeSessionView>,
287 pub lifecycles: Vec<RuntimeLifecycleView>,
288 #[serde(default)]
290 pub reviews: Vec<RuntimeReviewView>,
291 #[serde(default)]
293 pub notices: Vec<RuntimeNotice>,
294 #[serde(default)]
298 pub subagents: Vec<mj_core::subagent::SubagentRecord>,
299}
300
301#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
302#[serde(rename_all = "snake_case")]
303pub enum RuntimeLifecycleKind {
304 Create,
305 Suspend,
306 Resume,
307 Move,
308 ForceStop,
309 DestroyStopped,
310 ForceDestroy,
311 Cleanup,
312}
313
314#[derive(Debug, Clone, Serialize, Deserialize)]
315#[serde(deny_unknown_fields)]
316pub struct RuntimeLifecycleView {
317 pub operation_id: String,
318 pub cancellable: bool,
319 pub session_id: String,
320 pub kind: RuntimeLifecycleKind,
321 pub started_at_epoch_seconds: u64,
322 pub active_stages: Vec<(ProvisionStage, u64)>,
323 pub resume_destination: Option<(String, String)>,
324 pub notice: Option<String>,
325}
326
327#[derive(Debug, Clone, Serialize, Deserialize)]
328#[serde(deny_unknown_fields)]
329pub struct ResumeSessionRequest {
330 pub session_id: String,
331 pub workspace_id: String,
332 pub profile_id: String,
333 pub target_template_id: String,
334 pub additional_mounts: Option<Vec<AdditionalMount>>,
335 pub resource_allocation: Option<SessionResourceAllocation>,
336 pub discard_queue: bool,
337 pub repository_preflight: Option<ResumeRepositorySourceReceipt>,
338}
339
340#[derive(Debug, Clone, Serialize, Deserialize)]
341#[serde(deny_unknown_fields)]
342pub struct CreateSessionRequest {
343 #[serde(default)]
344 pub create_managed_worktree: Option<bool>,
345 #[serde(default)]
347 pub mjolnir_subagents: Option<bool>,
348 #[serde(default)]
349 pub initial_prompt: Option<String>,
350 pub workspace_id: String,
351 pub profile_id: String,
352 pub bundle_id: String,
353 pub project_directory: Option<PathBuf>,
354 pub target_template_id: String,
355 pub additional_mounts: Vec<AdditionalMount>,
356 pub resource_allocation: Option<SessionResourceAllocation>,
357 pub title: String,
358 pub session_title_override: Option<String>,
359}
360
361#[derive(Debug, Clone, Serialize, Deserialize)]
362#[serde(deny_unknown_fields)]
363pub struct RegisteredSession {
364 pub session: SessionRecord,
365 pub remembered_container_size: Option<(String, HostContainerSize)>,
366}
367
368#[derive(Debug, Clone, Serialize, Deserialize)]
369#[serde(deny_unknown_fields)]
370pub struct DraftPreview {
371 pub id: String,
372 pub session_id: Option<String>,
373 pub source: String,
374 pub owner_pid: Option<u32>,
375 pub saved_at: String,
376}
377
378#[derive(Debug, Clone, Serialize, Deserialize)]
385#[serde(rename_all = "snake_case", tag = "action", content = "arguments")]
386pub enum DaemonAction {
387 NativeAgentHistory {
388 owner: String,
389 child: String,
390 before: Option<(u64, String)>,
391 },
392 Ping,
393 Status,
394 WebViewerAccess,
395 RecoverWebViewer(crate::web::WebViewerRecovery),
396 InspectWebListener,
397 ListWorkspaces,
398 CreateWorkspace {
399 name: String,
400 },
401 RenameWorkspace {
402 workspace_id: String,
403 name: String,
404 },
405 TouchWorkspace {
406 workspace_id: String,
407 },
408 CloseWorkspace {
409 workspace_id: String,
410 },
411 CancelWorkspaceClose {
412 workspace_id: String,
413 },
414 DeleteWorkspace {
415 workspace_id: String,
416 },
417 Attach {
418 client_id: String,
419 pid: u32,
420 },
421 Detach {
422 client_id: String,
423 },
424 PersistReadReceipt {
425 client_id: String,
426 workspace_id: String,
427 session_id: String,
428 through: u64,
429 },
430 PersistDetachedSessionState {
431 client_id: String,
432 workspace_id: String,
433 session_id: String,
434 through: u64,
435 owner_pid: u32,
436 draft: mj_core::storage::DetachedSessionDraft,
437 },
438 SaveActiveReview {
439 session_id: String,
440 review: mj_core::storage::StoredReview,
441 },
442 ClearActiveReview {
443 session_id: String,
444 },
445 SaveWorkspacePaneSizes {
446 workspace_id: String,
447 sizes: mj_core::workspace::PaneSizes,
448 },
449 SaveWorkspaceLayout {
450 workspace_id: String,
451 layout: mj_core::workspace::ConversationLayout,
452 },
453 PersistImportedSession {
454 session: Box<SessionRecord>,
455 },
456 SetSessionTitle {
457 session_id: String,
458 title: String,
459 },
460 SetSessionContainerSettings {
461 session_id: String,
462 cpus: Option<String>,
463 memory: Option<String>,
464 mounts: Vec<AdditionalMount>,
465 mount_history: Vec<PathBuf>,
466 },
467 SetSessionAcpTitle {
468 session_id: String,
469 title: Option<String>,
470 },
471 MarkSessionTargetMissing {
472 session_id: String,
473 detail: String,
474 updated_at: String,
475 },
476 CheckpointSession {
477 session_id: String,
478 },
479 WikiSearch {
482 query: String,
483 limit: usize,
484 },
485 WikiBrief {
487 wiki_id: String,
488 max_chars: usize,
489 },
490 WikiHits {
492 wiki_id: String,
493 query: String,
494 context_messages: usize,
495 per_message_chars: usize,
496 },
497 WikiSession {
499 wiki_id: String,
500 },
501 WikiRestore(WikiRestoreRequest),
503 ScanRecovery {
504 all_instances: bool,
505 },
506 AdoptRecovery {
507 session_id: String,
508 target_id: String,
509 profile: Option<String>,
510 bundle: Option<String>,
511 all_instances: bool,
512 },
513 DestroyRecovery {
514 session_id: String,
515 target_id: String,
516 confirmation: String,
517 all_instances: bool,
518 },
519 Snapshot {
520 workspace_id: String,
521 },
522 RuntimeSnapshot {
523 workspace_id: String,
524 after_revision: u64,
525 #[serde(default)]
526 all_workspaces: bool,
527 },
528 RenameProfile {
529 old_id: String,
530 new_id: String,
531 },
532 RenameTarget {
533 old_id: String,
534 new_id: String,
535 },
536 SubmitSessionCommand {
537 #[serde(default)]
538 inherited_draft: Option<String>,
539 session_id: String,
540 command_id: String,
541 command: RelayCommand,
542 },
543 QueueStartupPrompt {
548 session_id: String,
549 text: String,
550 #[serde(default)]
554 inherited_draft: Option<String>,
555 },
556 SyncSession {
557 session_id: String,
558 },
559 RespondElicitation {
560 session_id: String,
561 elicitation_id: String,
562 response: ElicitationResponse,
563 },
564 StopBackgroundTask {
565 session_id: String,
566 background_task_id: String,
567 },
568 ReviewerAction {
572 session_id: String,
573 #[serde(default, skip_serializing_if = "Option::is_none")]
576 role: Option<String>,
577 action: crate::session::ReviewerAction,
578 },
579 StartTurnReview {
581 session_id: String,
582 },
583 ResolveTurnReview {
585 session_id: String,
586 resolution: Resolution,
587 },
588 SuspendSession {
589 session_id: String,
590 },
591 StartCreateSession(CreateSessionRequest),
592 WaitCreateSession {
593 session_id: String,
594 },
595 ResumeSession(ResumeSessionRequest),
596 PrepareMoveSession(MoveSelection),
597 MoveSession(MoveSessionRequest),
598 DiscardSinceCheckpoint {
599 session_id: String,
600 checkpoint: mj_core::state::CheckpointMetadata,
601 },
602 DestroyStoppedSession {
603 session_id: String,
604 delete_branch: bool,
608 },
609 ForceDestroySession {
610 session_id: String,
611 delete_branch: bool,
613 },
614 ForceDeleteWorkspace {
615 workspace_id: String,
616 },
617 CancelLifecycle {
618 session_id: String,
619 },
620 RecoverDraft {
621 draft_id: String,
622 },
623 Stop,
624}
625
626#[derive(Debug, Serialize, Deserialize)]
627#[serde(deny_unknown_fields)]
628pub struct RequestEnvelope {
629 pub protocol_version: u32,
630 pub request_id: u64,
631 pub token: String,
632 pub action: DaemonAction,
633}
634
635#[derive(Debug, Serialize, Deserialize)]
636#[serde(deny_unknown_fields)]
637pub struct ResponseEnvelope {
638 pub protocol_version: u32,
639 pub request_id: u64,
640 pub result: std::result::Result<DaemonReply, String>,
641}
642
643#[derive(Debug, Clone, Serialize, Deserialize)]
644#[serde(rename_all = "snake_case", tag = "reply", content = "value")]
645pub enum DaemonReply {
646 NativeAgentHistory(mj_core::native_agent::NativeAgentHistoryPage),
647 Pong,
648 Status(DaemonStatus),
649 WebViewerAccess(crate::web::WebViewerAccess),
650 WebListeners(Vec<crate::web::WebListenerProcess>),
651 Workspaces(Vec<WorkspaceListing>),
652 Workspace(WorkspaceRecord),
653 Snapshot(WorkspaceSnapshot),
654 RuntimeSnapshot(Box<RuntimeSnapshot>),
655 RegisteredSession(Box<RegisteredSession>),
656 MovePreparation(Box<MovePreparation>),
657 MoveOutcome(MoveOutcome),
658 Ordinal(u64),
659 Text(String),
660 OptionalSessionState(Option<SessionState>),
661 Checkpoint(mj_core::state::CheckpointMetadata),
662 RecoveryScan(mj_core::state::RecoveryScan),
663 WikiRows(WikiSearchPage),
664 WikiHits(Option<WikiHitTranscript>),
665 WikiSession(Option<Box<WikiSessionInfo>>),
666 Reviewer(Box<crate::session::ReviewerOutcome>),
667 Done,
668}
669
670#[derive(Debug, Clone, Serialize, Deserialize)]
671#[serde(deny_unknown_fields)]
672pub struct DaemonStatus {
673 pub pid: u32,
674 pub started_at: String,
675 pub build_version: String,
676 pub attached_clients: usize,
677 pub phone_status: WebViewerStatus,
678}
679
680#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
681#[serde(rename_all = "snake_case", tag = "state")]
682pub enum WebViewerStatus {
683 Disabled,
684 Starting,
685 Ready {
686 viewer_url: String,
687 viewer_code: String,
688 qr_login_url: Option<String>,
689 fallback_reason: Option<String>,
690 },
691 Stopped,
692 Error {
693 message: String,
694 },
695}
696
697impl std::fmt::Debug for WebViewerStatus {
698 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
699 match self {
700 Self::Ready {
701 viewer_url,
702 viewer_code,
703 fallback_reason,
704 ..
705 } => formatter
706 .debug_struct("Ready")
707 .field("viewer_url", viewer_url)
708 .field("viewer_code", viewer_code)
709 .field("qr_login_url", &"[redacted]")
710 .field("fallback_reason", fallback_reason)
711 .finish(),
712 Self::Disabled => formatter.write_str("Disabled"),
713 Self::Starting => formatter.write_str("Starting"),
714 Self::Stopped => formatter.write_str("Stopped"),
715 Self::Error { message } => formatter.debug_tuple("Error").field(message).finish(),
716 }
717 }
718}
719
720impl std::fmt::Display for WebViewerStatus {
721 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
722 match self {
723 Self::Disabled => formatter.write_str("disabled"),
724 Self::Starting => formatter.write_str("starting"),
725 Self::Stopped => formatter.write_str("stopped unexpectedly"),
726 Self::Error { message } => write!(formatter, "error: {message}"),
727 Self::Ready {
728 viewer_url,
729 viewer_code,
730 fallback_reason,
731 ..
732 } => {
733 write!(formatter, "{viewer_url}; viewer code {viewer_code}")?;
734 if let Some(reason) = fallback_reason {
735 write!(
736 formatter,
737 "; local only because Tailscale HTTPS is unavailable: {reason}"
738 )?;
739 }
740 Ok(())
741 }
742 }
743 }
744}
745
746#[cfg(unix)]
761pub fn process_is_zombie(pid: u32) -> bool {
762 let pid = sysinfo::Pid::from_u32(pid);
763 let mut system = sysinfo::System::new();
764 system.refresh_processes(sysinfo::ProcessesToUpdate::Some(&[pid]), true);
765 system
766 .process(pid)
767 .is_some_and(|process| process.status() == sysinfo::ProcessStatus::Zombie)
768}
769
770pub async fn wait_for_exit(pid: u32) -> Result<()> {
776 let deadline = Instant::now() + STOP_TIMEOUT;
777 while daemon_process_is_alive(pid) {
778 ensure!(Instant::now() < deadline, "process {pid} is still running");
779 tokio::time::sleep(RETRY_DELAY).await;
780 }
781 Ok(())
782}
783
784pub fn daemon_process_is_alive(pid: u32) -> bool {
790 #[cfg(unix)]
791 {
792 if pid == 0 {
793 return false;
794 }
795 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
796 return false;
797 };
798 let mut status = 0;
799 let waited = unsafe { libc::waitpid(raw_pid, &mut status, libc::WNOHANG) };
802 if waited == raw_pid {
803 return false;
804 }
805 if waited == 0 {
806 return true;
807 }
808 let wait_error = std::io::Error::last_os_error();
809 if wait_error.raw_os_error() != Some(libc::ECHILD) {
810 return true;
811 }
812
813 #[cfg(target_os = "macos")]
814 return owned_daemon_group_is_alive(raw_pid);
815
816 #[cfg(not(target_os = "macos"))]
817 process_is_alive(pid)
818 }
819 #[cfg(not(unix))]
820 process_is_alive(pid)
821}
822
823#[cfg(target_os = "macos")]
824pub fn owned_daemon_group_is_alive(pid: libc::pid_t) -> bool {
825 if unsafe { libc::kill(-pid, 0) } == 0 {
833 return true;
834 }
835 let error = std::io::Error::last_os_error();
836 !matches!(error.raw_os_error(), Some(libc::ESRCH) | Some(libc::EPERM))
837}
838
839pub fn process_is_alive(pid: u32) -> bool {
840 #[cfg(unix)]
841 {
842 if pid == 0 {
843 return false;
844 }
845 let Ok(raw_pid) = libc::pid_t::try_from(pid) else {
846 return false;
847 };
848 let result = unsafe { libc::kill(raw_pid, 0) };
851 let exists =
852 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM);
853 exists && !process_is_zombie(pid)
854 }
855 #[cfg(not(unix))]
856 {
857 let _ = pid;
858 true
859 }
860}
861
862pub fn read_metadata() -> Result<DaemonMetadata> {
863 let metadata = read_metadata_any()?;
864 ensure!(
865 metadata.protocol_version == PROTOCOL_VERSION,
866 "daemon protocol {} is incompatible with client protocol {}",
867 metadata.protocol_version,
868 PROTOCOL_VERSION
869 );
870 Ok(metadata)
871}
872
873pub fn read_metadata_any() -> Result<DaemonMetadata> {
874 let path = metadata_path();
875 let body = fs::read(&path).with_context(|| format!("read {}", path.display()))?;
876 let metadata: DaemonMetadata =
877 serde_json::from_slice(&body).with_context(|| format!("parse {}", path.display()))?;
878 Ok(metadata)
879}
880
881pub async fn write_frame<T: Serialize>(stream: &mut TcpStream, value: &T) -> Result<()> {
882 let body = serde_json::to_vec(value)?;
883 write_encoded_frame(stream, &body).await
884}
885
886pub async fn write_encoded_frame(stream: &mut TcpStream, body: &[u8]) -> Result<()> {
887 ensure!(
888 body.len() <= MAX_FRAME_BYTES,
889 "daemon frame is too large: {} bytes exceeds {MAX_FRAME_BYTES}",
890 body.len()
891 );
892 stream.write_u32(body.len() as u32).await?;
893 stream.write_all(body).await?;
894 stream.flush().await?;
895 Ok(())
896}
897
898pub async fn read_frame<T: for<'de> Deserialize<'de>>(stream: &mut TcpStream) -> Result<T> {
899 let length = stream.read_u32().await? as usize;
900 ensure!(
901 length <= MAX_FRAME_BYTES,
902 "daemon frame exceeds {MAX_FRAME_BYTES} bytes"
903 );
904 let mut body = vec![0_u8; length];
905 stream.read_exact(&mut body).await?;
906 serde_json::from_slice(&body).context("decode daemon frame")
907}
908
909pub struct DaemonClient {
910 metadata: DaemonMetadata,
911 stream: TcpStream,
912 next_request_id: u64,
913}
914
915impl DaemonClient {
916 pub async fn connect(metadata: DaemonMetadata) -> Result<Self> {
917 let stream =
918 tokio::time::timeout(Duration::from_secs(1), TcpStream::connect(metadata.address))
919 .await
920 .context("time out connecting to Mjolnir daemon")??;
921 Ok(Self {
922 metadata,
923 stream,
924 next_request_id: 1,
925 })
926 }
927
928 pub async fn request(&mut self, action: DaemonAction) -> Result<DaemonReply> {
932 let protocol_version = self.metadata.protocol_version;
933 let request_id = self.next_request_id;
934 self.next_request_id += 1;
935 write_frame(
936 &mut self.stream,
937 &RequestEnvelope {
938 protocol_version,
939 request_id,
940 token: self.metadata.token.clone(),
941 action,
942 },
943 )
944 .await?;
945 let response: ResponseEnvelope = read_frame(&mut self.stream).await?;
946 ensure!(
947 response.protocol_version == protocol_version,
948 "daemon changed protocol"
949 );
950 ensure!(
951 response.request_id == request_id,
952 "daemon crossed request IDs"
953 );
954 response.result.map_err(anyhow::Error::msg)
955 }
956
957 pub async fn status(&mut self) -> Result<DaemonStatus> {
958 match self.request(DaemonAction::Status).await? {
959 DaemonReply::Status(status) => Ok(status),
960 reply => bail!("unexpected daemon status reply {reply:?}"),
961 }
962 }
963
964 pub async fn web_access(&mut self) -> Result<crate::web::WebViewerAccess> {
965 match self.request(DaemonAction::WebViewerAccess).await? {
966 DaemonReply::WebViewerAccess(access) => Ok(access),
967 reply => bail!("unexpected web viewer reply {reply:?}"),
968 }
969 }
970
971 pub async fn recover_web_viewer(
972 &mut self,
973 action: crate::web::WebViewerRecovery,
974 ) -> Result<()> {
975 match self.request(DaemonAction::RecoverWebViewer(action)).await? {
976 DaemonReply::Done => Ok(()),
977 reply => bail!("unexpected web viewer recovery reply {reply:?}"),
978 }
979 }
980
981 pub async fn inspect_web_listener(&mut self) -> Result<Vec<crate::web::WebListenerProcess>> {
982 match self.request(DaemonAction::InspectWebListener).await? {
983 DaemonReply::WebListeners(processes) => Ok(processes),
984 reply => bail!("unexpected listener inspection reply {reply:?}"),
985 }
986 }
987
988 pub async fn list_workspaces(&mut self) -> Result<Vec<WorkspaceListing>> {
989 match self.request(DaemonAction::ListWorkspaces).await? {
990 DaemonReply::Workspaces(workspaces) => Ok(workspaces),
991 reply => bail!("unexpected daemon workspace reply {reply:?}"),
992 }
993 }
994
995 pub async fn rename_profile(&mut self, old_id: String, new_id: String) -> Result<()> {
996 match self
997 .request(DaemonAction::RenameProfile { old_id, new_id })
998 .await?
999 {
1000 DaemonReply::Done => Ok(()),
1001 reply => bail!("unexpected rename-profile reply {reply:?}"),
1002 }
1003 }
1004
1005 pub async fn rename_target(&mut self, old_id: String, new_id: String) -> Result<()> {
1006 match self
1007 .request(DaemonAction::RenameTarget { old_id, new_id })
1008 .await?
1009 {
1010 DaemonReply::Done => Ok(()),
1011 reply => bail!("unexpected rename-target reply {reply:?}"),
1012 }
1013 }
1014
1015 pub async fn create_workspace(&mut self, name: String) -> Result<WorkspaceRecord> {
1016 match self.request(DaemonAction::CreateWorkspace { name }).await? {
1017 DaemonReply::Workspace(workspace) => Ok(workspace),
1018 reply => bail!("unexpected create-workspace reply {reply:?}"),
1019 }
1020 }
1021
1022 pub async fn rename_workspace(&mut self, workspace_id: String, name: String) -> Result<()> {
1023 match self
1024 .request(DaemonAction::RenameWorkspace { workspace_id, name })
1025 .await?
1026 {
1027 DaemonReply::Done => Ok(()),
1028 reply => bail!("unexpected rename-workspace reply {reply:?}"),
1029 }
1030 }
1031
1032 pub async fn touch_workspace(&mut self, workspace_id: String) -> Result<()> {
1033 match self
1034 .request(DaemonAction::TouchWorkspace { workspace_id })
1035 .await?
1036 {
1037 DaemonReply::Done => Ok(()),
1038 reply => bail!("unexpected touch-workspace reply {reply:?}"),
1039 }
1040 }
1041
1042 pub async fn cancel_workspace_close(&mut self, workspace_id: String) -> Result<()> {
1043 match self
1044 .request(DaemonAction::CancelWorkspaceClose { workspace_id })
1045 .await?
1046 {
1047 DaemonReply::Done => Ok(()),
1048 reply => bail!("unexpected cancel workspace close reply: {reply:?}"),
1049 }
1050 }
1051
1052 pub async fn close_workspace(&mut self, workspace_id: String) -> Result<()> {
1053 match self
1054 .request(DaemonAction::CloseWorkspace { workspace_id })
1055 .await?
1056 {
1057 DaemonReply::Done => Ok(()),
1058 reply => bail!("unexpected close workspace reply: {reply:?}"),
1059 }
1060 }
1061
1062 pub async fn delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1063 match self
1064 .request(DaemonAction::DeleteWorkspace { workspace_id })
1065 .await?
1066 {
1067 DaemonReply::Done => Ok(()),
1068 reply => bail!("unexpected delete-workspace reply {reply:?}"),
1069 }
1070 }
1071
1072 pub async fn attach(&mut self, client_id: String, pid: u32) -> Result<()> {
1073 match self
1074 .request(DaemonAction::Attach { client_id, pid })
1075 .await?
1076 {
1077 DaemonReply::Done => Ok(()),
1078 reply => bail!("unexpected attach reply {reply:?}"),
1079 }
1080 }
1081
1082 pub async fn detach(&mut self, client_id: String) -> Result<()> {
1083 match self.request(DaemonAction::Detach { client_id }).await? {
1084 DaemonReply::Done => Ok(()),
1085 reply => bail!("unexpected detach reply {reply:?}"),
1086 }
1087 }
1088
1089 pub async fn persist_read_receipt(
1090 &mut self,
1091 client_id: String,
1092 workspace_id: String,
1093 session_id: String,
1094 through: u64,
1095 ) -> Result<u64> {
1096 match self
1097 .request(DaemonAction::PersistReadReceipt {
1098 client_id,
1099 workspace_id,
1100 session_id,
1101 through,
1102 })
1103 .await?
1104 {
1105 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1106 reply => bail!("unexpected read-receipt reply {reply:?}"),
1107 }
1108 }
1109
1110 pub async fn persist_detached_session_state(
1111 &mut self,
1112 client_id: String,
1113 workspace_id: String,
1114 session_id: String,
1115 through: u64,
1116 owner_pid: u32,
1117 draft: mj_core::storage::DetachedSessionDraft,
1118 ) -> Result<()> {
1119 match self
1120 .request(DaemonAction::PersistDetachedSessionState {
1121 client_id,
1122 workspace_id,
1123 session_id,
1124 through,
1125 owner_pid,
1126 draft,
1127 })
1128 .await?
1129 {
1130 DaemonReply::Done => Ok(()),
1131 reply => bail!("unexpected detached-session-state reply {reply:?}"),
1132 }
1133 }
1134
1135 pub async fn save_active_review(
1136 &mut self,
1137 session_id: String,
1138 review: mj_core::storage::StoredReview,
1139 ) -> Result<()> {
1140 match self
1141 .request(DaemonAction::SaveActiveReview { session_id, review })
1142 .await?
1143 {
1144 DaemonReply::Done => Ok(()),
1145 reply => bail!("unexpected save-review reply {reply:?}"),
1146 }
1147 }
1148
1149 pub async fn clear_active_review(&mut self, session_id: String) -> Result<()> {
1150 match self
1151 .request(DaemonAction::ClearActiveReview { session_id })
1152 .await?
1153 {
1154 DaemonReply::Done => Ok(()),
1155 reply => bail!("unexpected clear-review reply {reply:?}"),
1156 }
1157 }
1158
1159 pub async fn save_workspace_pane_sizes(
1160 &mut self,
1161 workspace_id: String,
1162 sizes: mj_core::workspace::PaneSizes,
1163 ) -> Result<()> {
1164 match self
1165 .request(DaemonAction::SaveWorkspacePaneSizes {
1166 workspace_id,
1167 sizes,
1168 })
1169 .await?
1170 {
1171 DaemonReply::Done => Ok(()),
1172 reply => bail!("unexpected pane-size save reply {reply:?}"),
1173 }
1174 }
1175
1176 pub async fn save_workspace_layout(
1177 &mut self,
1178 workspace_id: String,
1179 layout: mj_core::workspace::ConversationLayout,
1180 ) -> Result<()> {
1181 match self
1182 .request(DaemonAction::SaveWorkspaceLayout {
1183 workspace_id,
1184 layout,
1185 })
1186 .await?
1187 {
1188 DaemonReply::Done => Ok(()),
1189 reply => bail!("unexpected layout save reply {reply:?}"),
1190 }
1191 }
1192
1193 pub async fn persist_imported_session(&mut self, session: SessionRecord) -> Result<()> {
1194 match self
1195 .request(DaemonAction::PersistImportedSession {
1196 session: Box::new(session),
1197 })
1198 .await?
1199 {
1200 DaemonReply::Done => Ok(()),
1201 reply => bail!("unexpected imported-session reply {reply:?}"),
1202 }
1203 }
1204
1205 pub async fn set_session_title(&mut self, session_id: String, title: String) -> Result<String> {
1206 match self
1207 .request(DaemonAction::SetSessionTitle { session_id, title })
1208 .await?
1209 {
1210 DaemonReply::Text(title) => Ok(title),
1211 reply => bail!("unexpected session-title reply {reply:?}"),
1212 }
1213 }
1214
1215 pub async fn set_session_container_settings(
1216 &mut self,
1217 session_id: String,
1218 cpus: Option<String>,
1219 memory: Option<String>,
1220 mounts: Vec<AdditionalMount>,
1221 mount_history: Vec<PathBuf>,
1222 ) -> Result<()> {
1223 match self
1224 .request(DaemonAction::SetSessionContainerSettings {
1225 session_id,
1226 cpus,
1227 memory,
1228 mounts,
1229 mount_history,
1230 })
1231 .await?
1232 {
1233 DaemonReply::Done => Ok(()),
1234 reply => bail!("unexpected container-settings reply {reply:?}"),
1235 }
1236 }
1237
1238 pub async fn set_session_acp_title(
1239 &mut self,
1240 session_id: String,
1241 title: Option<String>,
1242 ) -> Result<()> {
1243 match self
1244 .request(DaemonAction::SetSessionAcpTitle { session_id, title })
1245 .await?
1246 {
1247 DaemonReply::Done => Ok(()),
1248 reply => bail!("unexpected ACP-title reply {reply:?}"),
1249 }
1250 }
1251
1252 pub async fn mark_session_target_missing(
1253 &mut self,
1254 session_id: String,
1255 detail: String,
1256 updated_at: String,
1257 ) -> Result<Option<SessionState>> {
1258 match self
1259 .request(DaemonAction::MarkSessionTargetMissing {
1260 session_id,
1261 detail,
1262 updated_at,
1263 })
1264 .await?
1265 {
1266 DaemonReply::OptionalSessionState(state) => Ok(state),
1267 reply => bail!("unexpected target-missing reply {reply:?}"),
1268 }
1269 }
1270
1271 pub async fn checkpoint_session(
1272 &mut self,
1273 session_id: String,
1274 ) -> Result<mj_core::state::CheckpointMetadata> {
1275 match self
1276 .request(DaemonAction::CheckpointSession { session_id })
1277 .await?
1278 {
1279 DaemonReply::Checkpoint(checkpoint) => Ok(checkpoint),
1280 reply => bail!("unexpected checkpoint reply {reply:?}"),
1281 }
1282 }
1283
1284 pub async fn wiki_search(&mut self, query: String, limit: usize) -> Result<WikiSearchPage> {
1289 match self
1290 .request(DaemonAction::WikiSearch { query, limit })
1291 .await?
1292 {
1293 DaemonReply::WikiRows(page) => Ok(page),
1294 reply => bail!("unexpected SessionWiki search reply {reply:?}"),
1295 }
1296 }
1297
1298 pub async fn wiki_brief(&mut self, wiki_id: String, max_chars: usize) -> Result<String> {
1300 match self
1301 .request(DaemonAction::WikiBrief { wiki_id, max_chars })
1302 .await?
1303 {
1304 DaemonReply::Text(markdown) => Ok(markdown),
1305 reply => bail!("unexpected SessionWiki brief reply {reply:?}"),
1306 }
1307 }
1308
1309 pub async fn wiki_hits(
1314 &mut self,
1315 wiki_id: String,
1316 query: String,
1317 context_messages: usize,
1318 per_message_chars: usize,
1319 ) -> Result<Option<WikiHitTranscript>> {
1320 match self
1321 .request(DaemonAction::WikiHits {
1322 wiki_id,
1323 query,
1324 context_messages,
1325 per_message_chars,
1326 })
1327 .await?
1328 {
1329 DaemonReply::WikiHits(transcript) => Ok(transcript),
1330 reply => bail!("unexpected SessionWiki hits reply {reply:?}"),
1331 }
1332 }
1333
1334 pub async fn wiki_session(&mut self, wiki_id: String) -> Result<Option<WikiSessionInfo>> {
1337 match self.request(DaemonAction::WikiSession { wiki_id }).await? {
1338 DaemonReply::WikiSession(info) => Ok(info.map(|info| *info)),
1339 reply => bail!("unexpected SessionWiki session reply {reply:?}"),
1340 }
1341 }
1342
1343 pub async fn wiki_restore(&mut self, request: WikiRestoreRequest) -> Result<RegisteredSession> {
1347 match self.request(DaemonAction::WikiRestore(request)).await? {
1348 DaemonReply::RegisteredSession(registered) => Ok(*registered),
1349 reply => bail!("unexpected SessionWiki restore reply {reply:?}"),
1350 }
1351 }
1352
1353 pub async fn scan_recovery(
1354 &mut self,
1355 all_instances: bool,
1356 ) -> Result<mj_core::state::RecoveryScan> {
1357 match self
1358 .request(DaemonAction::ScanRecovery { all_instances })
1359 .await?
1360 {
1361 DaemonReply::RecoveryScan(scan) => Ok(scan),
1362 reply => bail!("unexpected recovery-scan reply {reply:?}"),
1363 }
1364 }
1365
1366 pub async fn adopt_recovery(
1367 &mut self,
1368 session_id: String,
1369 target_id: String,
1370 profile: Option<String>,
1371 bundle: Option<String>,
1372 all_instances: bool,
1373 ) -> Result<()> {
1374 match self
1375 .request(DaemonAction::AdoptRecovery {
1376 session_id,
1377 target_id,
1378 profile,
1379 bundle,
1380 all_instances,
1381 })
1382 .await?
1383 {
1384 DaemonReply::Done => Ok(()),
1385 reply => bail!("unexpected recovery-adopt reply {reply:?}"),
1386 }
1387 }
1388
1389 pub async fn destroy_recovery(
1390 &mut self,
1391 session_id: String,
1392 target_id: String,
1393 confirmation: String,
1394 all_instances: bool,
1395 ) -> Result<()> {
1396 match self
1397 .request(DaemonAction::DestroyRecovery {
1398 session_id,
1399 target_id,
1400 confirmation,
1401 all_instances,
1402 })
1403 .await?
1404 {
1405 DaemonReply::Done => Ok(()),
1406 reply => bail!("unexpected recovery-destroy reply {reply:?}"),
1407 }
1408 }
1409
1410 pub async fn snapshot(&mut self, workspace_id: String) -> Result<WorkspaceSnapshot> {
1411 match self
1412 .request(DaemonAction::Snapshot { workspace_id })
1413 .await?
1414 {
1415 DaemonReply::Snapshot(snapshot) => Ok(snapshot),
1416 reply => bail!("unexpected snapshot reply {reply:?}"),
1417 }
1418 }
1419
1420 pub async fn runtime_snapshot(
1421 &mut self,
1422 workspace_id: String,
1423 after_revision: u64,
1424 all_workspaces: bool,
1425 ) -> Result<RuntimeSnapshot> {
1426 match self
1427 .request(DaemonAction::RuntimeSnapshot {
1428 workspace_id,
1429 after_revision,
1430 all_workspaces,
1431 })
1432 .await?
1433 {
1434 DaemonReply::RuntimeSnapshot(snapshot) => Ok(*snapshot),
1435 reply => bail!("unexpected runtime snapshot reply {reply:?}"),
1436 }
1437 }
1438
1439 pub async fn submit_session_command(
1440 &mut self,
1441 session_id: String,
1442 command_id: String,
1443 command: RelayCommand,
1444 inherited_draft: Option<String>,
1445 ) -> Result<u64> {
1446 match self
1447 .request(DaemonAction::SubmitSessionCommand {
1448 inherited_draft,
1449 session_id,
1450 command_id,
1451 command,
1452 })
1453 .await?
1454 {
1455 DaemonReply::Ordinal(ordinal) => Ok(ordinal),
1456 reply => bail!("unexpected session command reply {reply:?}"),
1457 }
1458 }
1459
1460 pub async fn queue_startup_prompt(
1464 &mut self,
1465 session_id: String,
1466 text: String,
1467 inherited_draft: Option<String>,
1468 ) -> Result<()> {
1469 match self
1470 .request(DaemonAction::QueueStartupPrompt {
1471 session_id,
1472 text,
1473 inherited_draft,
1474 })
1475 .await?
1476 {
1477 DaemonReply::Done => Ok(()),
1478 reply => bail!("unexpected startup prompt reply {reply:?}"),
1479 }
1480 }
1481
1482 pub async fn start_turn_review(&mut self, session_id: String) -> Result<()> {
1488 match self
1489 .request(DaemonAction::StartTurnReview { session_id })
1490 .await?
1491 {
1492 DaemonReply::Done => Ok(()),
1493 reply => bail!("unexpected turn-review reply {reply:?}"),
1494 }
1495 }
1496
1497 pub async fn resolve_turn_review(
1498 &mut self,
1499 session_id: String,
1500 resolution: Resolution,
1501 ) -> Result<()> {
1502 match self
1503 .request(DaemonAction::ResolveTurnReview {
1504 session_id,
1505 resolution,
1506 })
1507 .await?
1508 {
1509 DaemonReply::Done => Ok(()),
1510 reply => bail!("unexpected turn-review resolution reply {reply:?}"),
1511 }
1512 }
1513
1514 pub async fn reviewer_action(
1515 &mut self,
1516 session_id: String,
1517 role: Option<String>,
1518 action: crate::session::ReviewerAction,
1519 ) -> Result<crate::session::ReviewerOutcome> {
1520 match self
1521 .request(DaemonAction::ReviewerAction {
1522 session_id,
1523 role,
1524 action,
1525 })
1526 .await?
1527 {
1528 DaemonReply::Reviewer(outcome) => Ok(*outcome),
1529 reply => bail!("unexpected reviewer reply {reply:?}"),
1530 }
1531 }
1532
1533 pub async fn sync_session(&mut self, session_id: String) -> Result<()> {
1534 match self
1535 .request(DaemonAction::SyncSession { session_id })
1536 .await?
1537 {
1538 DaemonReply::Done => Ok(()),
1539 reply => bail!("unexpected session sync reply {reply:?}"),
1540 }
1541 }
1542
1543 pub async fn respond_elicitation(
1544 &mut self,
1545 session_id: String,
1546 elicitation_id: String,
1547 response: ElicitationResponse,
1548 ) -> Result<()> {
1549 match self
1550 .request(DaemonAction::RespondElicitation {
1551 session_id,
1552 elicitation_id,
1553 response,
1554 })
1555 .await?
1556 {
1557 DaemonReply::Done => Ok(()),
1558 reply => bail!("unexpected elicitation reply {reply:?}"),
1559 }
1560 }
1561
1562 pub async fn native_agent_history(
1563 &mut self,
1564 owner: String,
1565 child: String,
1566 before: Option<(u64, String)>,
1567 ) -> Result<mj_core::native_agent::NativeAgentHistoryPage> {
1568 match self
1569 .request(DaemonAction::NativeAgentHistory {
1570 owner,
1571 child,
1572 before,
1573 })
1574 .await?
1575 {
1576 DaemonReply::NativeAgentHistory(page) => Ok(page),
1577 reply => bail!("unexpected native agent history reply {reply:?}"),
1578 }
1579 }
1580
1581 pub async fn stop_background_task(
1582 &mut self,
1583 session_id: String,
1584 background_task_id: String,
1585 ) -> Result<()> {
1586 match self
1587 .request(DaemonAction::StopBackgroundTask {
1588 session_id,
1589 background_task_id,
1590 })
1591 .await?
1592 {
1593 DaemonReply::Done => Ok(()),
1594 reply => bail!("unexpected background task stop reply {reply:?}"),
1595 }
1596 }
1597
1598 pub async fn suspend_session(&mut self, session_id: String) -> Result<()> {
1599 match self
1600 .request(DaemonAction::SuspendSession { session_id })
1601 .await?
1602 {
1603 DaemonReply::Done => Ok(()),
1604 reply => bail!("unexpected close-session reply {reply:?}"),
1605 }
1606 }
1607
1608 pub async fn start_create_session(
1609 &mut self,
1610 request: CreateSessionRequest,
1611 ) -> Result<RegisteredSession> {
1612 match self
1613 .request(DaemonAction::StartCreateSession(request))
1614 .await?
1615 {
1616 DaemonReply::RegisteredSession(registered) => Ok(*registered),
1617 reply => bail!("unexpected start-create reply {reply:?}"),
1618 }
1619 }
1620
1621 pub async fn wait_create_session(&mut self, session_id: String) -> Result<()> {
1622 match self
1623 .request(DaemonAction::WaitCreateSession { session_id })
1624 .await?
1625 {
1626 DaemonReply::Done => Ok(()),
1627 reply => bail!("unexpected wait-create reply {reply:?}"),
1628 }
1629 }
1630
1631 pub async fn resume_session(&mut self, request: ResumeSessionRequest) -> Result<()> {
1632 match self.request(DaemonAction::ResumeSession(request)).await? {
1633 DaemonReply::Done => Ok(()),
1634 reply => bail!("unexpected resume-session reply {reply:?}"),
1635 }
1636 }
1637
1638 pub async fn discard_since_checkpoint(
1639 &mut self,
1640 session_id: String,
1641 checkpoint: mj_core::state::CheckpointMetadata,
1642 ) -> Result<()> {
1643 match self
1644 .request(DaemonAction::DiscardSinceCheckpoint {
1645 session_id,
1646 checkpoint,
1647 })
1648 .await?
1649 {
1650 DaemonReply::Done => Ok(()),
1651 reply => bail!("unexpected force-stop reply {reply:?}"),
1652 }
1653 }
1654
1655 pub async fn destroy_stopped_session(
1656 &mut self,
1657 session_id: String,
1658 delete_branch: bool,
1659 ) -> Result<()> {
1660 match self
1661 .request(DaemonAction::DestroyStoppedSession {
1662 session_id,
1663 delete_branch,
1664 })
1665 .await?
1666 {
1667 DaemonReply::Done => Ok(()),
1668 reply => bail!("unexpected destroy-stopped reply {reply:?}"),
1669 }
1670 }
1671
1672 pub async fn force_destroy_session(
1673 &mut self,
1674 session_id: String,
1675 delete_branch: bool,
1676 ) -> Result<()> {
1677 match self
1678 .request(DaemonAction::ForceDestroySession {
1679 session_id,
1680 delete_branch,
1681 })
1682 .await?
1683 {
1684 DaemonReply::Done => Ok(()),
1685 reply => bail!("unexpected force-destroy reply {reply:?}"),
1686 }
1687 }
1688
1689 pub async fn force_delete_workspace(&mut self, workspace_id: String) -> Result<()> {
1690 match self
1691 .request(DaemonAction::ForceDeleteWorkspace { workspace_id })
1692 .await?
1693 {
1694 DaemonReply::Done => Ok(()),
1695 reply => bail!("unexpected force-delete-workspace reply {reply:?}"),
1696 }
1697 }
1698
1699 pub async fn cancel_lifecycle(&mut self, session_id: String) -> Result<()> {
1700 match self
1701 .request(DaemonAction::CancelLifecycle { session_id })
1702 .await?
1703 {
1704 DaemonReply::Done => Ok(()),
1705 reply => bail!("unexpected cancel-lifecycle reply {reply:?}"),
1706 }
1707 }
1708
1709 pub async fn recover_draft(&mut self, draft_id: String) -> Result<()> {
1710 match self
1711 .request(DaemonAction::RecoverDraft { draft_id })
1712 .await?
1713 {
1714 DaemonReply::Done => Ok(()),
1715 reply => bail!("unexpected recover-draft reply {reply:?}"),
1716 }
1717 }
1718
1719 pub async fn stop(&mut self) -> Result<()> {
1720 match self.request(DaemonAction::Stop).await? {
1721 DaemonReply::Done => Ok(()),
1722 reply => bail!("unexpected stop reply {reply:?}"),
1723 }
1724 }
1725}
1726
1727pub async fn connect_existing() -> Result<DaemonClient> {
1728 let metadata = tokio::task::spawn_blocking(read_metadata)
1729 .await
1730 .context("read daemon metadata task failed")??;
1731 DaemonClient::connect(metadata).await
1732}
1733
1734pub struct ManagementClient {
1738 inner: DaemonClient,
1739}
1740
1741impl ManagementClient {
1742 pub fn new(inner: DaemonClient) -> Self {
1743 Self { inner }
1744 }
1745 pub fn protocol_version(&self) -> u32 {
1746 self.inner.metadata.protocol_version
1747 }
1748
1749 pub async fn status(&mut self) -> Result<DaemonStatus> {
1750 self.inner.status().await
1751 }
1752
1753 pub async fn stop(&mut self) -> Result<()> {
1754 tokio::time::timeout(STOP_TIMEOUT, self.inner.stop())
1755 .await
1756 .context("Mjolnir daemon did not acknowledge the stop before the deadline")?
1757 }
1758
1759 pub async fn stop_and_wait(mut self) -> Result<()> {
1761 let pid = self.inner.metadata.pid;
1762 self.stop().await?;
1763 wait_for_exit(pid).await.with_context(|| {
1764 format!(
1765 "Mjolnir daemon {pid} accepted the stop but was still running after {}s",
1766 STOP_TIMEOUT.as_secs()
1767 )
1768 })
1769 }
1770}
1771
1772pub async fn connect_management() -> Result<ManagementClient> {
1773 Ok(ManagementClient {
1774 inner: DaemonClient::connect(read_metadata_any()?).await?,
1775 })
1776}
1777
1778pub fn ensure_supported_daemon_protocol(metadata: &DaemonMetadata) -> Result<()> {
1786 ensure!(
1787 metadata.protocol_version <= PROTOCOL_VERSION,
1788 "{}",
1789 unsupported_daemon_protocol_message(
1790 metadata.protocol_version,
1791 &describe_running_daemon_and_client_builds(metadata.pid, &metadata.build_version),
1792 )
1793 );
1794 Ok(())
1795}
1796
1797fn unsupported_daemon_protocol_message(daemon_protocol: u32, builds: &str) -> String {
1800 format!(
1801 "the daemon uses a newer protocol ({daemon_protocol}) than this client ({PROTOCOL_VERSION}). {builds}. \
1802 Put the daemon's directory first on PATH, or reinstall this client from that build."
1803 )
1804}
1805pub const PROTOCOL_VERSION: u32 = 32;
1806pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
1807pub const STOP_TIMEOUT: Duration = Duration::from_secs(30);
1817pub const RETRY_DELAY: Duration = Duration::from_millis(40);
1818impl DaemonClient {
1819 pub async fn prepare_move_session(
1820 &mut self,
1821 selection: MoveSelection,
1822 ) -> Result<MovePreparation> {
1823 match self
1824 .request(DaemonAction::PrepareMoveSession(selection))
1825 .await?
1826 {
1827 DaemonReply::MovePreparation(preparation) => Ok(*preparation),
1828 _ => bail!("daemon returned an unexpected move preparation reply"),
1829 }
1830 }
1831
1832 pub async fn move_session(&mut self, request: MoveSessionRequest) -> Result<MoveOutcome> {
1833 match self.request(DaemonAction::MoveSession(request)).await? {
1834 DaemonReply::MoveOutcome(outcome) => Ok(outcome),
1835 _ => bail!("daemon returned an unexpected move reply"),
1836 }
1837 }
1838}
1839
1840#[cfg(test)]
1841mod tests {
1842 use super::*;
1843 use crate::executable::{BuildDescription, describe_daemon_and_client_builds};
1844 use std::path::Path;
1845
1846 #[test]
1847 fn unsupported_protocol_message_names_both_binaries_and_versions() {
1848 let message = unsupported_daemon_protocol_message(
1849 PROTOCOL_VERSION + 5,
1850 &describe_daemon_and_client_builds(
1851 4242,
1852 BuildDescription {
1853 executable: Some(Path::new("/home/dev/mj/target/release/mj")),
1854 version: "2.14.0",
1855 },
1856 BuildDescription {
1857 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
1858 version: "2.9.0",
1859 },
1860 ),
1861 );
1862 assert_eq!(
1863 message,
1864 format!(
1865 "the daemon uses a newer protocol ({}) than this client ({PROTOCOL_VERSION}). \
1866 Daemon 4242 runs /home/dev/mj/target/release/mj (version 2.14.0), \
1867 while this client runs /home/dev/.cargo/bin/mj (version 2.9.0). \
1868 Put the daemon's directory first on PATH, or reinstall this client from that build.",
1869 PROTOCOL_VERSION + 5
1870 )
1871 );
1872 }
1873
1874 #[test]
1875 fn unsupported_protocol_message_still_names_the_client_when_the_daemon_file_is_unknown() {
1876 let message = unsupported_daemon_protocol_message(
1877 PROTOCOL_VERSION + 1,
1878 &describe_daemon_and_client_builds(
1879 4242,
1880 BuildDescription {
1881 executable: None,
1882 version: "2.14.0",
1883 },
1884 BuildDescription {
1885 executable: Some(Path::new("/home/dev/.cargo/bin/mj")),
1886 version: "2.9.0",
1887 },
1888 ),
1889 );
1890 assert!(
1891 message.contains("Daemon 4242 runs an unknown file (version 2.14.0)"),
1892 "{message}"
1893 );
1894 assert!(
1895 message.contains("this client runs /home/dev/.cargo/bin/mj (version 2.9.0)"),
1896 "{message}"
1897 );
1898 }
1899}