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