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