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