1use std::collections::{BTreeMap, BTreeSet};
4use std::path::{Component, Path, PathBuf};
5use std::sync::Arc;
6
7use anyhow::{Context, Result, bail};
8use serde::{Deserialize, Serialize};
9
10use crate::config::{Config, HarnessKind, ProjectRepository, TargetTemplate, validate_id};
11use crate::credentials::CredentialSyncSignal;
12use crate::relay::{
13 RELAY_EVENT_GENESIS_DIGEST, RelayOperationalState, SequencedEvent, WorkerEvent,
14};
15use crate::subagent::SubagentRecord;
16use crate::targets::{AdditionalMount, validate_additional_mounts};
17
18pub const STATE_VERSION: u32 = 1;
19
20mod target_runtime;
21pub use target_runtime::{TargetConnection, TargetRuntimeSettings};
22
23mod session_move;
24pub use session_move::*;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
27#[serde(rename_all = "kebab-case")]
28pub enum SessionState {
29 Provisioning,
30 Running,
31 Disconnected,
32 Checkpointing,
33 Closing,
34 Destroying,
35 #[serde(alias = "archived")]
38 Stopped,
39 Parked,
46 Lost,
47 Error,
48 DestroyedWithDataLoss,
49}
50
51#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
54#[serde(rename_all = "kebab-case")]
55pub enum SessionTransitionKind {
56 Starting,
57 Resuming,
58 Moving,
59 Suspending,
60 Destroying,
61 Stopping,
63}
64
65impl SessionTransitionKind {
66 pub const fn label(self) -> &'static str {
67 match self {
68 Self::Starting => "Starting",
69 Self::Resuming => "Resuming",
70 Self::Moving => "Moving",
71 Self::Suspending => "Suspending",
72 Self::Destroying => "Destroying",
73 Self::Stopping => "Stopping",
74 }
75 }
76
77 pub fn for_session(state: SessionState, operation: Option<Self>) -> Option<Self> {
78 operation.or_else(|| state.transition_kind())
79 }
80}
81
82#[cfg(test)]
83mod transition_tests {
84 use super::{SessionState, SessionTransitionKind};
85
86 #[test]
87 fn operation_ownership_hides_intermediate_move_states_but_not_ordinary_live_work() {
88 for state in [
89 SessionState::Stopped,
90 SessionState::Running,
91 SessionState::Disconnected,
92 ] {
93 assert_eq!(
94 SessionTransitionKind::for_session(state, Some(SessionTransitionKind::Moving)),
95 Some(SessionTransitionKind::Moving)
96 );
97 assert_eq!(SessionTransitionKind::for_session(state, None), None);
98 }
99 assert_eq!(SessionState::Checkpointing.transition_kind(), None);
100 assert_eq!(
101 SessionState::Closing.transition_kind(),
102 Some(SessionTransitionKind::Suspending)
103 );
104 }
105}
106
107#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
109#[serde(tag = "state", rename_all = "snake_case")]
110pub enum MaterializedExecutionState {
111 #[default]
112 Idle,
113 Running {
114 started_at_ms: i64,
115 },
116 Closing,
117 Closed,
118}
119
120pub use crate::transcript::{TerminalOutputRecord, TranscriptBody, TranscriptItem};
121
122#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
127#[serde(rename_all = "snake_case")]
128pub enum QueuedCommandKind {
129 #[default]
130 Prompt,
131 SetConfig {
132 key: String,
133 value: String,
134 },
135}
136
137impl QueuedCommandKind {
138 pub fn is_prompt(&self) -> bool {
139 matches!(self, Self::Prompt)
140 }
141}
142
143pub fn config_command_text(key: &str, value: &str) -> String {
146 if key == "fast-mode" {
147 "/fast".to_owned()
148 } else {
149 format!("/{key} {value}")
150 }
151}
152
153#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
154#[serde(deny_unknown_fields)]
155pub struct MaterializedQueuedPrompt {
156 pub command_id: String,
157 #[serde(default, skip_serializing_if = "QueuedCommandKind::is_prompt")]
158 pub kind: QueuedCommandKind,
159 pub content: Vec<serde_json::Value>,
160 pub queued_at_ms: i64,
161 #[serde(default, skip_serializing_if = "Option::is_none")]
165 pub accepted_ordinal: Option<u64>,
166}
167
168#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
171#[serde(deny_unknown_fields)]
172pub struct MaterializedTurn {
173 pub command_id: String,
174 #[serde(default, skip_serializing_if = "Option::is_none")]
175 pub accepted_ordinal: Option<u64>,
176 pub turn_start_position: u64,
179 pub started_at_ms: i64,
180 #[serde(default, skip_serializing_if = "Option::is_none")]
183 pub steered_into: Option<String>,
184}
185
186impl MaterializedTurn {
187 pub fn belongs_to(&self, command_id: &str) -> bool {
189 self.command_id == command_id || self.steered_into.as_deref() == Some(command_id)
190 }
191}
192
193#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
195#[serde(tag = "kind", rename_all = "snake_case")]
196pub enum TurnOutcomeKind {
197 Completed { stop_reason: String },
199 Rejected { message: String },
201 Interrupted { message: String },
203}
204
205impl std::fmt::Display for TurnOutcomeKind {
208 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
209 match self {
210 Self::Completed { stop_reason } => match classify_prompt_completion(stop_reason) {
211 PromptCompletion::Finished => formatter.write_str("completed, end of turn"),
212 PromptCompletion::InputRequired => {
213 formatter.write_str("completed, waiting for input")
214 }
215 PromptCompletion::Cancelled => formatter.write_str("interrupted"),
217 PromptCompletion::QuotaLimit => formatter.write_str("failed: quota limit reached"),
218 PromptCompletion::Error => {
219 write!(formatter, "failed: {}", stop_reason_words(stop_reason))
220 }
221 },
222 Self::Rejected { message } => write!(
223 formatter,
224 "failed: {}",
225 message.lines().next().unwrap_or_default().trim()
226 ),
227 Self::Interrupted { .. } => formatter.write_str("interrupted"),
228 }
229 }
230}
231
232fn stop_reason_words(stop_reason: &str) -> String {
235 let mut words = String::new();
236 let mut previous_lower = false;
237 for character in stop_reason.trim().chars() {
238 if character == '_' || character == '-' || character.is_whitespace() {
239 if !words.ends_with(' ') && !words.is_empty() {
240 words.push(' ');
241 }
242 previous_lower = false;
243 continue;
244 }
245 if character.is_uppercase() && previous_lower {
246 words.push(' ');
247 }
248 previous_lower = character.is_lowercase() || character.is_ascii_digit();
249 words.extend(character.to_lowercase());
250 }
251 match words.trim_end() {
252 "" => "no reason given".to_owned(),
253 words => words.to_owned(),
254 }
255}
256
257#[derive(Debug, Clone, Copy, PartialEq, Eq)]
258pub enum PromptCompletion {
259 InputRequired,
260 Finished,
261 Cancelled,
262 QuotaLimit,
263 Error,
264}
265
266pub fn classify_prompt_completion(stop_reason: &str) -> PromptCompletion {
268 let normalized = stop_reason
269 .chars()
270 .filter(|character| *character != '_' && *character != '-')
271 .flat_map(char::to_lowercase)
272 .collect::<String>();
273 match normalized.as_str() {
274 "endturn" => PromptCompletion::Finished,
275 "awaitinginput" => PromptCompletion::InputRequired,
276 "cancelled" | "canceled" => PromptCompletion::Cancelled,
277 "quotalimit" => PromptCompletion::QuotaLimit,
278 _ => PromptCompletion::Error,
279 }
280}
281
282#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
284#[serde(deny_unknown_fields)]
285pub struct MaterializedTurnOutcome {
286 #[serde(default, skip_serializing_if = "Option::is_none")]
287 pub diagnostic: Option<crate::diagnostic::TurnDiagnostic>,
288
289 #[serde(default, skip_serializing_if = "Option::is_none")]
290 pub usage: Option<crate::usage::TokenUsage>,
291 pub command_id: String,
292 #[serde(default, skip_serializing_if = "Option::is_none")]
293 pub accepted_ordinal: Option<u64>,
294 #[serde(default, skip_serializing_if = "Option::is_none")]
295 pub turn_start_position: Option<u64>,
296 pub completed_ordinal: u64,
297 pub completed_at_ms: i64,
298 pub outcome: TurnOutcomeKind,
299}
300
301impl MaterializedTurnOutcome {
302 pub fn interruption_ordinal(&self) -> Option<u64> {
304 (self.turn_start_position.is_some()
305 && matches!(self.outcome, TurnOutcomeKind::Interrupted { .. }))
306 .then_some(self.completed_ordinal)
307 }
308}
309
310#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
312#[serde(deny_unknown_fields)]
313pub struct MaterializedSession {
314 pub session_id: String,
315 pub applied_event_ordinal: u64,
316 pub applied_event_digest: String,
317 pub last_activity_at_ms: Option<i64>,
320 pub execution: MaterializedExecutionState,
321 #[serde(default, skip_serializing_if = "Option::is_none")]
322 pub session_title: Option<String>,
323 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
324 pub configuration: BTreeMap<String, serde_json::Value>,
325 #[serde(default, skip_serializing_if = "Vec::is_empty")]
326 pub transcript: Vec<Arc<TranscriptItem>>,
329 #[serde(default, skip_serializing_if = "Vec::is_empty")]
330 pub queued_prompts: Vec<MaterializedQueuedPrompt>,
331 #[serde(default, skip_serializing_if = "Vec::is_empty")]
334 pub pending_elicitations: Vec<crate::elicitation::ElicitationRequest>,
335 #[serde(default, skip_serializing_if = "Option::is_none")]
337 pub active_turn: Option<MaterializedTurn>,
338 #[serde(default, skip_serializing_if = "Option::is_none")]
341 pub last_turn_outcome: Option<MaterializedTurnOutcome>,
342}
343
344#[derive(Debug, Clone, PartialEq, Eq)]
347pub struct MaterializedSessionSummary {
348 pub session_id: String,
349 pub applied_event_ordinal: u64,
350 pub last_activity_at_ms: Option<i64>,
351 pub execution: MaterializedExecutionState,
352 pub session_title: Option<String>,
353 pub last_agent_message: Option<String>,
354 pub last_user_message: Option<String>,
355 pub last_agent_message_follows_last_user: bool,
358 pub agent_message_latest_content_ordinals: Vec<u64>,
359 pub interruption_event_ordinals: Vec<u64>,
360}
361
362impl MaterializedSession {
363 pub fn empty(session_id: impl Into<String>) -> Self {
364 Self {
365 session_id: session_id.into(),
366 applied_event_ordinal: 0,
367 applied_event_digest: RELAY_EVENT_GENESIS_DIGEST.into(),
368 last_activity_at_ms: None,
369 execution: MaterializedExecutionState::Idle,
370 session_title: None,
371 configuration: BTreeMap::new(),
372 transcript: Vec::new(),
373 queued_prompts: Vec::new(),
374 pending_elicitations: Vec::new(),
375 active_turn: None,
376 last_turn_outcome: None,
377 }
378 }
379
380 pub fn last_activity_at_ms(&self) -> Option<i64> {
381 self.last_activity_at_ms
382 }
383
384 pub fn resolved_title(&self) -> Option<String> {
390 self.session_title
391 .as_deref()
392 .and_then(normalize_session_title)
393 .or_else(|| {
394 self.transcript.iter().find_map(|item| {
395 let TranscriptBody::User { content } = &item.body else {
396 return None;
397 };
398 provisional_session_title(&crate::transcript::materialized_content_text(
399 content,
400 ))
401 })
402 })
403 .or_else(|| {
404 self.queued_prompts
405 .iter()
406 .filter(|prompt| prompt.kind.is_prompt())
407 .find_map(|prompt| {
408 provisional_session_title(&crate::transcript::materialized_content_text(
409 &prompt.content,
410 ))
411 })
412 })
413 }
414
415 pub fn unread_agent_messages_after(&self, viewed_through_event_ordinal: u64) -> u64 {
416 self.transcript
417 .iter()
418 .filter(|item| {
419 item.latest_content_event_ordinal
420 .is_some_and(|ordinal| ordinal > viewed_through_event_ordinal)
421 && item.is_nonempty_agent_message()
422 })
423 .count() as u64
424 }
425
426 pub fn unread_interruptions_after(&self, viewed_through_event_ordinal: u64) -> u64 {
427 self.interruption_event_ordinals()
428 .into_iter()
429 .filter(|ordinal| *ordinal > viewed_through_event_ordinal)
430 .count() as u64
431 }
432
433 pub fn interruption_event_ordinals(&self) -> Vec<u64> {
434 let mut ordinals = self
435 .transcript
436 .iter()
437 .filter(|item| item.is_work_interruption())
438 .map(|item| item.position)
439 .collect::<Vec<_>>();
440 if let Some(ordinal) = self
441 .last_turn_outcome
442 .as_ref()
443 .and_then(MaterializedTurnOutcome::interruption_ordinal)
444 {
445 ordinals.push(ordinal);
446 }
447 ordinals.sort_unstable();
448 ordinals.dedup();
449 ordinals
450 }
451
452 pub fn validate(&self) -> Result<()> {
453 validate_id("session", &self.session_id)?;
454 validate_relay_event_frontier(
455 self.applied_event_ordinal,
456 &self.applied_event_digest,
457 "materialized session event frontier",
458 )?;
459 if self
460 .session_title
461 .as_ref()
462 .is_some_and(|title| title.trim().is_empty())
463 {
464 bail!("materialized session has an empty title");
465 }
466 let mut item_ids = BTreeSet::new();
467 for item in &self.transcript {
468 item.validate(self.applied_event_ordinal)?;
469 if !item_ids.insert(item.stable_id.as_str()) {
470 bail!(
471 "materialized transcript contains duplicate item {:?}",
472 item.stable_id
473 );
474 }
475 }
476 let mut command_ids = BTreeSet::new();
477 for prompt in &self.queued_prompts {
478 if prompt.command_id.trim().is_empty() {
479 bail!("materialized prompt queue has an empty command id");
480 }
481 if !command_ids.insert(prompt.command_id.as_str()) {
482 bail!(
483 "materialized prompt queue contains duplicate command {:?}",
484 prompt.command_id
485 );
486 }
487 if let QueuedCommandKind::SetConfig { key, value } = &prompt.kind
488 && (key.trim().is_empty() || value.trim().is_empty())
489 {
490 bail!(
491 "materialized queued configuration change {:?} is incomplete",
492 prompt.command_id
493 );
494 }
495 }
496 Ok(())
497 }
498}
499
500#[derive(Debug, Clone, PartialEq)]
504pub struct ManagedSessionSnapshot {
505 pub materialized: MaterializedSession,
506 pub window: ProjectionWindow,
509 pub operational: RelayOperationalState,
510 pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
514 pub worker_build: Option<String>,
519 pub subagent_requests: Vec<crate::subagent::SubagentToolRequest>,
521 pub subagent_results: Vec<crate::subagent::SubagentToolResult>,
523}
524
525#[derive(Debug, Clone, PartialEq, Eq)]
537pub struct ProjectionWindow {
538 pub omitted_items: usize,
540 pub provisional_title: Option<String>,
542 pub latest_turn_start_position: Option<u64>,
546}
547
548impl ProjectionWindow {
549 pub fn trim(&mut self, session: &mut MaterializedSession, target: usize) {
552 let observed = Self::of(session);
553 if self.provisional_title.is_none() {
554 self.provisional_title = observed.provisional_title;
555 }
556 self.latest_turn_start_position = observed
557 .latest_turn_start_position
558 .or(self.latest_turn_start_position);
559 let mut boundary = session.transcript.len().saturating_sub(target.max(1));
560 for (index, item) in session.transcript.iter().enumerate() {
561 let mutable = match &item.body {
562 TranscriptBody::Agent { streaming, .. }
563 | TranscriptBody::Thought { streaming, .. } => *streaming,
564 TranscriptBody::Tool { call, .. } => matches!(
565 call.get("status").and_then(serde_json::Value::as_str),
566 Some("pending" | "in_progress")
567 ),
568 _ => false,
569 };
570 if mutable || Some(item.position) == self.latest_turn_start_position {
571 boundary = boundary.min(index);
572 }
573 }
574 let cut = session
575 .transcript
576 .iter()
577 .take(boundary + 1)
578 .rposition(|item| item.is_turn_start())
579 .unwrap_or(0);
580 if cut > 0 {
581 session.transcript.drain(..cut);
582 self.omitted_items += cut;
583 }
584 }
585
586 #[must_use]
588 pub fn of(session: &MaterializedSession) -> Self {
589 Self {
590 omitted_items: 0,
591 provisional_title: session.transcript.iter().find_map(|item| {
592 let TranscriptBody::User { content } = &item.body else {
593 return None;
594 };
595 provisional_session_title(&crate::transcript::materialized_content_text(content))
596 }),
597 latest_turn_start_position: session
598 .transcript
599 .iter()
600 .rev()
601 .find(|item| item.is_turn_start())
602 .map(|item| item.position),
603 }
604 }
605}
606
607impl ManagedSessionSnapshot {
608 #[must_use]
613 pub fn resolved_title(&self) -> Option<String> {
614 self.materialized
615 .session_title
616 .as_deref()
617 .and_then(normalize_session_title)
618 .or_else(|| self.window.provisional_title.clone())
619 .or_else(|| {
620 self.materialized
621 .queued_prompts
622 .iter()
623 .filter(|prompt| prompt.kind.is_prompt())
624 .find_map(|prompt| {
625 provisional_session_title(&crate::transcript::materialized_content_text(
626 &prompt.content,
627 ))
628 })
629 })
630 }
631
632 #[must_use]
637 pub fn latest_completed_turn_ordinal(&self) -> Option<u64> {
638 if self.materialized.execution != MaterializedExecutionState::Idle {
639 return None;
640 }
641 self.window.latest_turn_start_position
642 }
643}
644
645#[derive(Debug, Clone)]
647pub struct RecoveryObservation {
648 pub session: SessionRecord,
649 pub config: Config,
650 pub latest_completed_turn_ordinal: Option<u64>,
651 pub execution: MaterializedExecutionState,
652 pub checkpoint_safe: bool,
656}
657
658pub fn latest_completed_turn_ordinal(session: &MaterializedSession) -> Option<u64> {
663 if session.execution != MaterializedExecutionState::Idle {
664 return None;
665 }
666 session
667 .transcript
668 .iter()
669 .rev()
670 .find(|item| item.is_turn_start())
671 .map(|item| item.position)
672}
673
674pub fn validate_relay_event_digest(digest: &str, name: &str) -> Result<()> {
675 if digest.len() != 64
676 || !digest
677 .bytes()
678 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
679 {
680 bail!("{name} must be a lowercase SHA-256 digest");
681 }
682 Ok(())
683}
684
685pub fn validate_relay_event_frontier(ordinal: u64, digest: &str, name: &str) -> Result<()> {
686 validate_relay_event_digest(digest, name)?;
687 if (ordinal == 0) != (digest == RELAY_EVENT_GENESIS_DIGEST) {
688 bail!("{name} has inconsistent ordinal {ordinal} and digest {digest}");
689 }
690 Ok(())
691}
692
693fn is_false(value: &bool) -> bool {
694 !*value
695}
696
697impl SessionState {
698 pub const fn as_str(self) -> &'static str {
700 match self {
701 Self::Provisioning => "provisioning",
702 Self::Running => "running",
703 Self::Disconnected => "disconnected",
704 Self::Checkpointing => "checkpointing",
705 Self::Closing => "closing",
706 Self::Destroying => "destroying",
707 Self::Stopped => "stopped",
708 Self::Parked => "parked",
709 Self::Lost => "lost",
710 Self::Error => "error",
711 Self::DestroyedWithDataLoss => "destroyed-with-data-loss",
712 }
713 }
714
715 pub fn from_stored(value: &str) -> Option<Self> {
718 Some(match value {
719 "provisioning" => Self::Provisioning,
720 "running" => Self::Running,
721 "disconnected" => Self::Disconnected,
722 "checkpointing" => Self::Checkpointing,
723 "closing" => Self::Closing,
724 "destroying" => Self::Destroying,
725 "stopped" | "archived" => Self::Stopped,
726 "parked" => Self::Parked,
727 "lost" => Self::Lost,
728 "error" => Self::Error,
729 "destroyed-with-data-loss" => Self::DestroyedWithDataLoss,
730 _ => return None,
731 })
732 }
733
734 pub const fn transition_kind(self) -> Option<SessionTransitionKind> {
737 match self {
738 Self::Provisioning => Some(SessionTransitionKind::Starting),
739 Self::Closing => Some(SessionTransitionKind::Suspending),
740 Self::Destroying => Some(SessionTransitionKind::Destroying),
741 _ => None,
742 }
743 }
744
745 pub const fn is_active(self) -> bool {
752 matches!(
753 self,
754 Self::Provisioning
755 | Self::Running
756 | Self::Disconnected
757 | Self::Checkpointing
758 | Self::Closing
759 | Self::Destroying
760 | Self::Parked
761 | Self::Error
762 )
763 }
764
765 pub const fn has_live_worker(self) -> bool {
769 matches!(
770 self,
771 Self::Provisioning
772 | Self::Running
773 | Self::Disconnected
774 | Self::Checkpointing
775 | Self::Closing
776 | Self::Destroying
777 )
778 }
779}
780
781#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
782#[serde(tag = "kind", rename_all = "kebab-case")]
783pub enum PodmanWorkspaceLocator {
784 #[default]
785 ContainerLayer,
786 Volume {
787 name: String,
788 },
789 HostPath {
790 path: PathBuf,
791 helper: Vec<String>,
792 resource: String,
793 },
794}
795
796#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
797#[serde(tag = "kind", rename_all = "kebab-case")]
798pub enum TargetLocator {
799 LocalBare {
800 worker_root: PathBuf,
801 },
802 LocalPodman {
803 container_id: String,
804 #[serde(default)]
805 workspace_storage: PodmanWorkspaceLocator,
806 #[serde(default, skip_serializing_if = "Option::is_none")]
810 borrowed_from: Option<String>,
811 },
812 LocalDocker {
813 container_id: String,
814 #[serde(default, skip_serializing_if = "Option::is_none")]
818 borrowed_from: Option<String>,
819 },
820 AppleContainer {
821 container_id: String,
822 #[serde(default, skip_serializing_if = "Option::is_none")]
826 borrowed_from: Option<String>,
827 },
828 AwsEc2 {
829 instance_id: String,
830 #[serde(default, skip_serializing_if = "Option::is_none")]
831 address: Option<String>,
832 },
833 SshBare {
834 host: String,
835 workspace: PathBuf,
836 #[serde(default, skip_serializing_if = "Option::is_none")]
837 worker_id: Option<String>,
838 },
839 SshPodman {
840 host: String,
841 container_id: String,
842 #[serde(default)]
843 workspace_storage: PodmanWorkspaceLocator,
844 #[serde(default, skip_serializing_if = "Option::is_none")]
848 borrowed_from: Option<String>,
849 },
850 SshDocker {
851 host: String,
852 container_id: String,
853 #[serde(default, skip_serializing_if = "Option::is_none")]
857 borrowed_from: Option<String>,
858 },
859}
860
861impl ManagedWorktreeTarget {
862 pub fn same_location(&self, other: &Self) -> bool {
870 match (self, other) {
871 (Self::Local, Self::Local) => true,
872 (
873 Self::Ssh {
874 destination,
875 ssh_args,
876 },
877 Self::Ssh {
878 destination: other_destination,
879 ssh_args: other_args,
880 },
881 ) => {
882 destination == other_destination
883 && ssh_location_option(ssh_args, 'p', "port")
884 == ssh_location_option(other_args, 'p', "port")
885 && ssh_location_option(ssh_args, 'l', "user")
886 == ssh_location_option(other_args, 'l', "user")
887 }
888 _ => false,
889 }
890 }
891}
892
893fn ssh_location_option(args: &[String], flag: char, option: &str) -> Option<String> {
897 let mut args = args.iter();
898 while let Some(argument) = args.next() {
899 let Some(rest) = argument.strip_prefix('-') else {
900 continue;
901 };
902 let mut chars = rest.chars();
903 let Some(name) = chars.next() else { continue };
904 if name != flag && name != 'o' {
905 continue;
906 }
907 let inline = chars.as_str();
908 let value = if inline.is_empty() {
909 args.next().cloned()
910 } else {
911 Some(inline.to_owned())
912 };
913 if name == flag {
914 return value;
915 }
916 if let Some(setting) = value {
917 let (key, found) = setting
918 .split_once(['=', ' ', '\t'])
919 .unwrap_or((setting.as_str(), ""));
920 if key.trim().eq_ignore_ascii_case(option) {
921 return Some(found.trim().to_owned());
922 }
923 }
924 }
925 None
926}
927
928#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
929#[serde(tag = "kind", rename_all = "kebab-case")]
930pub enum ManagedWorktreeTarget {
931 Local,
932 Ssh {
933 destination: String,
934 #[serde(default, skip_serializing_if = "Vec::is_empty")]
935 ssh_args: Vec<String>,
936 },
937}
938
939#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
941#[serde(deny_unknown_fields)]
942pub struct ManagedWorktreeOptions {
943 pub available: bool,
944 pub default_create: bool,
945}
946
947#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
948#[serde(deny_unknown_fields)]
949pub struct ManagedWorktree {
950 #[serde(default, skip_serializing_if = "ManagedCheckoutKind::is_worktree")]
952 pub kind: ManagedCheckoutKind,
953 pub source_project_directory: PathBuf,
954 pub source_repository: PathBuf,
955 pub worktree_root: PathBuf,
956 pub branch: String,
957 pub target: ManagedWorktreeTarget,
958 #[serde(default, skip_serializing_if = "Option::is_none")]
962 pub base_commit: Option<String>,
963}
964
965#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
966#[serde(rename_all = "snake_case")]
967pub enum ManagedCheckoutKind {
968 #[default]
969 Worktree,
970 Clone,
971}
972
973impl ManagedCheckoutKind {
974 fn is_worktree(&self) -> bool {
975 matches!(self, Self::Worktree)
976 }
977}
978
979impl ManagedWorktree {
980 fn validate(&self, session_id: &str, project_directory: Option<&Path>) -> Result<()> {
981 for (label, path) in [
982 ("source project directory", &self.source_project_directory),
983 ("source repository", &self.source_repository),
984 ("worktree root", &self.worktree_root),
985 ] {
986 if !path.is_absolute() || path.components().any(|part| part == Component::ParentDir) {
987 bail!("managed worktree {label} must be an absolute safe path");
988 }
989 }
990 if !self
991 .source_project_directory
992 .starts_with(&self.source_repository)
993 {
994 bail!("managed worktree source directory is outside its repository");
995 }
996 let expected_root = self
997 .source_repository
998 .join(".mj")
999 .join(match self.kind {
1000 ManagedCheckoutKind::Worktree => "worktrees",
1001 ManagedCheckoutKind::Clone => "clones",
1002 })
1003 .join(session_id);
1004 if self.worktree_root != expected_root {
1005 bail!("managed worktree root does not match the session-owned path");
1006 }
1007 if self.kind == ManagedCheckoutKind::Worktree && self.branch != format!("mj/{session_id}") {
1008 bail!("managed worktree branch does not match the session id");
1009 }
1010 if self.kind == ManagedCheckoutKind::Clone && self.branch.trim().is_empty() {
1011 bail!("managed clone has no starting branch");
1012 }
1013 let relative = self
1014 .source_project_directory
1015 .strip_prefix(&self.source_repository)
1016 .expect("source relationship checked above");
1017 if project_directory != Some(self.worktree_root.join(relative).as_path()) {
1018 bail!("session project directory does not match its managed worktree");
1019 }
1020 match &self.target {
1021 ManagedWorktreeTarget::Local => {}
1022 ManagedWorktreeTarget::Ssh { destination, .. } if destination.trim().is_empty() => {
1023 bail!("managed SSH worktree has an empty destination")
1024 }
1025 ManagedWorktreeTarget::Ssh { .. } => {}
1026 }
1027 Ok(())
1028 }
1029}
1030
1031#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1032#[serde(tag = "kind", rename_all = "kebab-case")]
1033pub enum SessionResourceAllocation {
1034 Container {
1035 cpus: u64,
1036 memory_bytes: u64,
1037 },
1038 AwsEc2 {
1039 instance_type: String,
1040 vcpus: u64,
1041 memory_bytes: u64,
1042 },
1043}
1044
1045impl SessionResourceAllocation {
1046 pub fn validate(&self) -> Result<()> {
1047 match self {
1048 Self::Container { cpus, memory_bytes } if *cpus == 0 || *memory_bytes == 0 => {
1049 bail!("container resource allocation must have non-zero CPU and memory")
1050 }
1051 Self::AwsEc2 {
1052 instance_type,
1053 vcpus,
1054 memory_bytes,
1055 } if instance_type.trim().is_empty() || *vcpus == 0 || *memory_bytes == 0 => {
1056 bail!("EC2 resource allocation must have an instance type, CPU, and memory")
1057 }
1058 _ => Ok(()),
1059 }
1060 }
1061}
1062
1063pub fn allocation_cpus(allocation: &SessionResourceAllocation) -> u64 {
1065 match allocation {
1066 SessionResourceAllocation::Container { cpus, .. } => *cpus,
1067 SessionResourceAllocation::AwsEc2 { vcpus, .. } => *vcpus,
1068 }
1069}
1070
1071pub fn allocation_memory(allocation: &SessionResourceAllocation) -> u64 {
1073 match allocation {
1074 SessionResourceAllocation::Container { memory_bytes, .. }
1075 | SessionResourceAllocation::AwsEc2 { memory_bytes, .. } => *memory_bytes,
1076 }
1077}
1078
1079impl TargetLocator {
1080 fn validate(&self, session_id: &str) -> Result<()> {
1081 match self {
1082 Self::LocalBare { worker_root } => {
1083 if !worker_root.is_absolute()
1084 || worker_root
1085 .components()
1086 .any(|part| part == Component::ParentDir)
1087 || !worker_root.ends_with(session_id)
1088 {
1089 bail!(
1090 "local bare worker root must be an absolute safe path ending in the session id"
1091 );
1092 }
1093 }
1094 Self::LocalPodman { container_id, .. }
1095 | Self::LocalDocker { container_id, .. }
1096 | Self::AppleContainer { container_id, .. }
1097 | Self::SshPodman { container_id, .. }
1098 | Self::SshDocker { container_id, .. }
1099 if container_id.trim().is_empty() =>
1100 {
1101 bail!("target locator has an empty container id")
1102 }
1103 Self::AwsEc2 { instance_id, .. } if instance_id.trim().is_empty() => {
1104 bail!("target locator has an empty AWS instance id")
1105 }
1106 Self::SshBare {
1107 host, workspace, ..
1108 } => {
1109 if host.trim().is_empty() {
1110 bail!("bare SSH target locator has an empty host");
1111 }
1112 if workspace.as_os_str().is_empty()
1113 || workspace
1114 .components()
1115 .any(|part| part == Component::ParentDir)
1116 || !workspace.ends_with(session_id)
1117 {
1118 bail!("bare SSH target locator must be a safe path ending in the session id");
1119 }
1120 }
1121 Self::SshPodman { host, .. } if host.trim().is_empty() => {
1122 bail!("SSH Podman target locator has an empty host")
1123 }
1124 Self::SshDocker { host, .. } if host.trim().is_empty() => {
1125 bail!("SSH Docker target locator has an empty host")
1126 }
1127 _ => {}
1128 }
1129 Ok(())
1130 }
1131}
1132
1133#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1134#[serde(deny_unknown_fields)]
1135pub struct CheckpointMetadata {
1136 pub archive_path: PathBuf,
1137 pub sha256: String,
1139 pub created_at: String,
1140 pub event_frontier: u64,
1141}
1142
1143#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1145#[serde(rename_all = "snake_case")]
1146pub enum PublicationState {
1147 Published,
1148 Unpublished,
1149 Unknown,
1150}
1151
1152#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1154pub struct PublicationAssessment {
1155 pub checkpoint_sha256: String,
1156 pub state: PublicationState,
1157 pub dirty: bool,
1158 pub stashed: bool,
1159 pub saved_commits: Vec<String>,
1160 pub destinations: Vec<String>,
1161 pub checked_at: String,
1162 pub reason: Option<String>,
1163}
1164
1165impl CheckpointMetadata {
1166 fn validate(&self) -> Result<()> {
1167 if self.archive_path.as_os_str().is_empty() {
1168 bail!("checkpoint archive path is empty");
1169 }
1170 if self.sha256.len() != 64
1171 || !self
1172 .sha256
1173 .bytes()
1174 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1175 {
1176 bail!("checkpoint SHA-256 must be 64 lowercase hexadecimal characters");
1177 }
1178 if self.created_at.trim().is_empty() {
1179 bail!("checkpoint timestamp is empty");
1180 }
1181 Ok(())
1182 }
1183}
1184
1185#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1189#[serde(deny_unknown_fields)]
1190pub struct SessionBuildCache {
1191 pub host: String,
1194 pub directory: PathBuf,
1195 #[serde(default, skip_serializing_if = "Option::is_none")]
1198 pub max_size: Option<String>,
1199 #[serde(default, skip_serializing_if = "Option::is_none")]
1202 pub target_root: Option<PathBuf>,
1203}
1204
1205#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1209pub struct ArchiveSpacePreview {
1210 pub sessions: usize,
1212 pub bytes: u64,
1214 pub reclaimable_sessions: usize,
1217 pub reclaimable_bytes: u64,
1218}
1219
1220#[derive(Debug, Clone, PartialEq, Eq)]
1224pub struct BuildCachePreview {
1225 pub native_mbx: Option<String>,
1227 pub directory: Option<PathBuf>,
1229 pub max_size: Option<BuildCacheLimit>,
1231 pub stats: Option<BuildCacheStats>,
1233 pub off_reason: Option<BuildCacheOff>,
1236}
1237
1238#[derive(Debug, Clone, PartialEq, Eq)]
1246pub struct BuildCacheStats {
1247 pub builds: u64,
1249 pub cached_compilations: u64,
1251 pub avoided_compiler_ns: u64,
1253 pub reflinked_bytes: u64,
1255}
1256
1257#[derive(Debug, Clone, PartialEq, Eq)]
1259pub enum BuildCacheOff {
1260 TurnedOff,
1262 Unavailable(String),
1265}
1266
1267impl std::fmt::Display for BuildCacheOff {
1268 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1269 match self {
1270 Self::TurnedOff => formatter.write_str("turned off for this machine"),
1271 Self::Unavailable(reason) => formatter.write_str(reason),
1272 }
1273 }
1274}
1275
1276#[derive(Debug, Clone, PartialEq, Eq)]
1278pub enum BuildCacheLimit {
1279 Size(String),
1281 HostConfiguration(Option<String>),
1284}
1285
1286#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1287#[serde(deny_unknown_fields)]
1288pub struct SessionRecord {
1289 pub id: String,
1290 #[serde(default = "default_session_workspace_id")]
1295 pub workspace_id: String,
1296 pub title: String,
1297 pub harness_kind: HarnessKind,
1298 pub last_profile: String,
1299 pub bundle_id: String,
1300 #[serde(default, skip_serializing_if = "Option::is_none")]
1302 pub project_directory: Option<PathBuf>,
1303 #[serde(default, skip_serializing_if = "Option::is_none")]
1305 pub managed_worktree: Option<ManagedWorktree>,
1306 #[serde(default, skip_serializing_if = "Option::is_none")]
1308 pub create_managed_worktree: Option<bool>,
1309 #[serde(default, skip_serializing_if = "Option::is_none")]
1314 pub launch_base: Option<String>,
1315 #[serde(default, skip_serializing_if = "Option::is_none")]
1317 pub launch_branch: Option<String>,
1318 #[serde(default, skip_serializing_if = "Option::is_none")]
1321 pub checkout: Option<crate::remote_git::ExactCheckout>,
1322 #[serde(default, skip_serializing_if = "Option::is_none")]
1323 pub expected_runtime_identity: Option<String>,
1324 #[serde(default, skip_serializing_if = "Option::is_none")]
1326 pub publication: Option<PublicationAssessment>,
1327 #[serde(default, skip_serializing_if = "Option::is_none")]
1330 pub mjolnir_subagents: Option<bool>,
1331 pub target_template_id: String,
1332 #[serde(default, skip_serializing_if = "Option::is_none")]
1333 pub resource_allocation: Option<SessionResourceAllocation>,
1334 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1335 pub additional_mounts: Vec<AdditionalMount>,
1336 #[serde(default, skip_serializing_if = "Option::is_none")]
1339 pub container_cpus: Option<String>,
1340 #[serde(default, skip_serializing_if = "Option::is_none")]
1343 pub container_memory: Option<String>,
1344 #[serde(default, skip_serializing_if = "Option::is_none")]
1350 pub container_workspace: Option<PathBuf>,
1351 #[serde(default, skip_serializing_if = "Option::is_none")]
1355 pub build_cache: Option<SessionBuildCache>,
1356 pub state: SessionState,
1357 #[serde(default, skip_serializing_if = "is_false")]
1360 pub archived: bool,
1361 #[serde(default, skip_serializing_if = "Option::is_none")]
1362 pub target: Option<TargetLocator>,
1363 #[serde(default, skip_serializing_if = "Option::is_none")]
1365 pub target_runtime: Option<TargetRuntimeSettings>,
1366 #[serde(default, skip_serializing_if = "Option::is_none")]
1367 pub native_session_id: Option<String>,
1368 #[serde(default, skip_serializing_if = "Option::is_none")]
1369 pub acp_session_title: Option<String>,
1370 #[serde(default, skip_serializing_if = "Option::is_none")]
1371 pub session_title_override: Option<String>,
1372 pub created_at: String,
1373 pub updated_at: String,
1374 #[serde(default, alias = "detached_after_event_ordinal")]
1375 pub viewed_through_event_ordinal: u64,
1376 #[serde(default, skip_serializing_if = "String::is_empty")]
1379 pub draft_input: String,
1380 #[serde(default, skip_serializing_if = "Option::is_none")]
1388 pub last_error: Option<String>,
1389 #[serde(default, skip_serializing_if = "Option::is_none")]
1390 pub last_checkpoint_error: Option<String>,
1391 #[serde(default, skip_serializing_if = "Option::is_none")]
1392 pub checkpoint: Option<CheckpointMetadata>,
1393}
1394
1395#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1396#[serde(deny_unknown_fields)]
1397pub struct HostContainerSize {
1398 pub cpus: u64,
1399 pub memory_bytes: u64,
1400}
1401
1402fn default_session_workspace_id() -> String {
1403 crate::workspace::DEFAULT_WORKSPACE_ID.to_owned()
1404}
1405
1406pub const DESTRUCTION_FAILURE_PREFIX: &str = "the destruction did not finish";
1415
1416pub const CLOSE_FAILURE_PREFIX: &str = "the suspension did not finish";
1417
1418pub fn is_public_lifecycle_error(error: &str) -> bool {
1420 error.starts_with(CLOSE_FAILURE_PREFIX)
1421 || error.starts_with(DESTRUCTION_FAILURE_PREFIX)
1422 || error.starts_with("the close did not finish")
1423}
1424
1425#[must_use]
1438pub fn target_label(config: &Config, target_id: &str, project: Option<&Path>) -> String {
1439 if !matches!(
1440 config.targets.get(target_id),
1441 Some(TargetTemplate::LocalBare | TargetTemplate::SshBare { .. })
1442 ) {
1443 return target_id.to_owned();
1444 }
1445 project.and_then(Path::file_name).map_or_else(
1446 || target_id.to_owned(),
1447 |directory| format!("{target_id}/{}", directory.to_string_lossy()),
1448 )
1449}
1450
1451impl SessionRecord {
1452 pub fn publication_state(&self) -> Option<PublicationState> {
1455 let independent_clone = self
1456 .managed_worktree
1457 .as_ref()
1458 .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Clone)
1459 || (self.managed_worktree.is_none() && self.project_directory.is_none());
1460 if !independent_clone {
1461 return None;
1462 }
1463 if self.state.is_active() {
1464 return Some(PublicationState::Unknown);
1465 }
1466 Some(
1467 self.checkpoint
1468 .as_ref()
1469 .zip(self.publication.as_ref())
1470 .filter(|(checkpoint, assessment)| {
1471 assessment.checkpoint_sha256 == checkpoint.sha256
1472 })
1473 .map_or(PublicationState::Unknown, |(_, assessment)| {
1474 if assessment.dirty || assessment.stashed {
1475 PublicationState::Unpublished
1476 } else {
1477 assessment.state
1478 }
1479 }),
1480 )
1481 }
1482
1483 pub fn target_runtime_settings<'a>(
1488 &'a self,
1489 config: &Config,
1490 ) -> Result<std::borrow::Cow<'a, TargetRuntimeSettings>> {
1491 if let Some(runtime) = &self.target_runtime {
1492 if let Some(refreshed) =
1495 config
1496 .targets
1497 .get(&self.target_template_id)
1498 .and_then(|template| {
1499 runtime.with_current_ssh_options(&TargetRuntimeSettings::from(template))
1500 })
1501 {
1502 return Ok(std::borrow::Cow::Owned(refreshed));
1503 }
1504 return Ok(std::borrow::Cow::Borrowed(runtime));
1505 }
1506 let template = config.targets.get(&self.target_template_id).ok_or_else(|| {
1507 crate::refusal::Refusal::precondition(format!(
1508 "Session {:?} has no recorded target access settings. Restore target {:?} in config.toml once, then retry.",
1509 self.id, self.target_template_id))
1510 })?;
1511 let runtime = TargetRuntimeSettings::from(template);
1512 if let Some(locator) = &self.target {
1513 crate::targets::TargetLocator::try_from(crate::targets::RecordedTarget {
1514 locator, runtime: Some(&runtime), session_id: &self.id,
1515 }).map_err(|error| crate::refusal::Refusal::precondition(format!(
1516 "Session {:?} cannot recover target {:?}: {error}. Restore its original target settings, then retry.", self.id, self.target_template_id)))?;
1517 }
1518 Ok(std::borrow::Cow::Owned(runtime))
1519 }
1520
1521 #[must_use]
1525 pub fn public_error(&self) -> Option<&str> {
1526 self.last_error
1527 .as_deref()
1528 .filter(|error| is_public_lifecycle_error(error))
1529 }
1530
1531 pub fn configuration_issue(&self, config: &Config) -> Option<String> {
1534 if !self.state.is_active() {
1535 return None;
1536 }
1537 let mut issues = Vec::new();
1538 match config.profiles.get(&self.last_profile) {
1539 None => issues.push(format!("missing profile {:?}", self.last_profile)),
1540 Some(profile) if profile.kind != self.harness_kind => issues.push(format!(
1541 "expects {:?}, but profile {:?} is {:?}",
1542 self.harness_kind, self.last_profile, profile.kind
1543 )),
1544 Some(_) => {}
1545 }
1546 if self.project_directory.is_none() && !config.bundles.contains_key(&self.bundle_id) {
1547 issues.push(format!("missing bundle {:?}", self.bundle_id));
1548 }
1549 if self.target_runtime.is_none() && !config.targets.contains_key(&self.target_template_id) {
1550 issues.push(format!(
1551 "missing target template {:?}",
1552 self.target_template_id
1553 ));
1554 }
1555 (!issues.is_empty()).then(|| format!(
1556 "Session {:?} needs configuration repair: {}. Restore these entries in config.toml, then retry. Run mj setup to rediscover installed profiles and targets; existing sessions are preserved.",
1557 self.id, issues.join("; ")
1558 ))
1559 }
1560
1561 pub fn validate_configuration(&self, config: &Config) -> Result<()> {
1562 if let Some(issue) = self.configuration_issue(config) {
1563 return Err(crate::refusal::Refusal::precondition(issue).into());
1564 }
1565 Ok(())
1566 }
1567
1568 pub fn display_title(&self) -> &str {
1570 self.session_title_override
1571 .as_deref()
1572 .or(self.acp_session_title.as_deref())
1573 .unwrap_or(&self.id)
1574 }
1575
1576 pub fn listed_title(&self) -> &str {
1581 let named = self.session_title_override.is_some() || self.acp_session_title.is_some();
1582 if !named && !self.title.trim().is_empty() {
1583 return &self.title;
1584 }
1585 self.display_title()
1586 }
1587
1588 pub fn project_name(&self, config: &Config) -> String {
1593 if let Some(worktree) = &self.managed_worktree {
1594 return path_leaf(&worktree.source_repository);
1595 }
1596 if let Some(project_directory) = &self.project_directory {
1597 return path_leaf(project_directory);
1598 }
1599 self.bundle_source_name(config)
1600 }
1601
1602 pub fn project_target(&self, config: &Config, target_id: &str) -> String {
1606 let project = self
1607 .managed_worktree
1608 .as_ref()
1609 .map(|worktree| &worktree.source_project_directory)
1610 .or(self.project_directory.as_ref());
1611 target_label(config, target_id, project.map(PathBuf::as_path))
1612 }
1613
1614 pub fn project_source(&self, config: &Config) -> ProjectSourceIdentity {
1619 if let Some(worktree) = &self.managed_worktree {
1620 return ProjectSourceIdentity::path(&worktree.source_repository, None);
1621 }
1622 if let Some(project_directory) = &self.project_directory {
1623 let remote = match &self.target {
1624 Some(TargetLocator::SshBare { host, .. }) => Some(host.as_str()),
1625 _ => None,
1626 };
1627 return ProjectSourceIdentity::path(project_directory, remote);
1628 }
1629 self.bundle_source_identity(config)
1630 .unwrap_or_else(|| ProjectSourceIdentity {
1631 key: format!("bundle:{}", self.bundle_id),
1632 short: path_leaf(Path::new(&self.bundle_id)),
1633 full: self.bundle_id.clone(),
1634 })
1635 }
1636
1637 fn bundle_source_name(&self, config: &Config) -> String {
1640 self.bundle_source_identity(config)
1641 .map(|source| source.short)
1642 .unwrap_or_else(|| path_leaf(Path::new(&self.bundle_id)))
1643 }
1644
1645 fn bundle_source_identity(&self, config: &Config) -> Option<ProjectSourceIdentity> {
1648 let bundle = config.bundles.get(&self.bundle_id)?;
1649 let sources = bundle
1650 .repositories
1651 .iter()
1652 .map(repository_source_identity)
1653 .collect::<Option<Vec<_>>>()?;
1654 ProjectSourceIdentity::bundle(sources)
1655 }
1656
1657 pub fn compare_by_creation(&self, other: &Self) -> std::cmp::Ordering {
1661 self.creation_order_key().cmp(&other.creation_order_key())
1662 }
1663
1664 pub fn creation_order_key(&self) -> (bool, Option<i64>, &str) {
1666 let timestamp = created_at_seconds(&self.created_at);
1667 (timestamp.is_none(), timestamp, &self.id)
1668 }
1669
1670 fn validate(&self, map_id: &str) -> Result<()> {
1671 validate_id("session", &self.id)?;
1672 if self.id != map_id {
1673 bail!(
1674 "session map key {map_id:?} does not match record id {:?}",
1675 self.id
1676 );
1677 }
1678 validate_id("workspace", &self.workspace_id)?;
1679 validate_id("profile", &self.last_profile)?;
1680 validate_id("bundle", &self.bundle_id)?;
1681 if let Some(project_directory) = &self.project_directory
1682 && (!project_directory.is_absolute()
1683 || project_directory
1684 .components()
1685 .any(|part| part == Component::ParentDir))
1686 {
1687 bail!("session {:?} has an unsafe project directory", self.id);
1688 }
1689 if let Some(managed_worktree) = &self.managed_worktree {
1690 managed_worktree.validate(&self.id, self.project_directory.as_deref())?;
1691 }
1692 validate_id("target template", &self.target_template_id)?;
1693 if let Some(allocation) = &self.resource_allocation {
1694 allocation.validate()?;
1695 }
1696 validate_additional_mounts(&self.additional_mounts)?;
1697 if self.title.trim().is_empty() {
1698 bail!("session {:?} has an empty title", self.id);
1699 }
1700 if self
1701 .acp_session_title
1702 .as_ref()
1703 .is_some_and(|title| title.trim().is_empty())
1704 || self
1705 .session_title_override
1706 .as_ref()
1707 .is_some_and(|title| title.trim().is_empty())
1708 {
1709 bail!("session {:?} has an empty display title", self.id);
1710 }
1711 if self.created_at.trim().is_empty() || self.updated_at.trim().is_empty() {
1712 bail!("session {:?} has an empty timestamp", self.id);
1713 }
1714 if let Some(target) = &self.target {
1715 target.validate(&self.id)?;
1716 }
1717 if let Some(checkpoint) = &self.checkpoint {
1718 checkpoint.validate()?;
1719 }
1720 Ok(())
1721 }
1722}
1723
1724fn repository_source_identity(repository: &ProjectRepository) -> Option<ProjectSourceIdentity> {
1725 repository
1726 .github
1727 .as_deref()
1728 .and_then(ProjectSourceIdentity::git_remote)
1729 .or_else(|| {
1730 repository
1731 .local
1732 .as_deref()
1733 .map(|path| ProjectSourceIdentity::path(path, None))
1734 })
1735}
1736
1737#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1738pub struct ProjectSourceIdentity {
1739 pub key: String,
1740 pub short: String,
1741 pub full: String,
1742}
1743
1744impl ProjectSourceIdentity {
1745 pub fn bundle(mut sources: Vec<Self>) -> Option<Self> {
1747 if sources.is_empty() {
1748 return None;
1749 }
1750 sources.sort_by(|left, right| {
1751 left.key
1752 .cmp(&right.key)
1753 .then_with(|| left.full.cmp(&right.full))
1754 .then_with(|| left.short.cmp(&right.short))
1755 });
1756 sources.dedup_by(|left, right| left.key == right.key);
1757 if sources.len() == 1 {
1758 return sources.pop();
1759 }
1760 let keys = sources
1761 .iter()
1762 .map(|source| source.key.clone())
1763 .collect::<Vec<_>>();
1764 let key = serde_json::to_string(&keys).ok()?;
1765 Some(Self {
1766 key: format!("bundle:{key}"),
1767 short: sources
1768 .iter()
1769 .map(|source| source.short.as_str())
1770 .collect::<Vec<_>>()
1771 .join(" + "),
1772 full: sources
1773 .iter()
1774 .map(|source| source.full.as_str())
1775 .collect::<Vec<_>>()
1776 .join(" + "),
1777 })
1778 }
1779
1780 pub fn git_remote(source: &str) -> Option<Self> {
1783 if let Some(normalized) = normalize_github_source(source) {
1784 let short = normalized
1785 .rsplit_once('/')
1786 .map_or(normalized.as_str(), |(_, repository)| repository)
1787 .to_owned();
1788 return Some(Self {
1789 key: format!("github:{}", normalized.to_lowercase()),
1790 short,
1791 full: normalized,
1792 });
1793 }
1794 let normalized = source.trim().trim_end_matches('/').trim_end_matches(".git");
1795 if normalized.is_empty() {
1796 return None;
1797 }
1798 let short = normalized
1799 .rsplit(['/', ':'])
1800 .find(|part| !part.is_empty())
1801 .unwrap_or(normalized)
1802 .to_owned();
1803 Some(Self {
1804 key: format!("git:{}", normalized.to_lowercase()),
1805 short,
1806 full: normalized.to_owned(),
1807 })
1808 }
1809
1810 pub fn path(path: &Path, remote: Option<&str>) -> Self {
1812 let normalized = path.components().collect::<PathBuf>();
1813 let path_text = normalized.to_string_lossy().into_owned();
1814 let full = remote.map_or_else(|| path_text.clone(), |host| format!("{host}:{path_text}"));
1815 let key = remote.map_or_else(
1816 || format!("path:{path_text}"),
1817 |host| format!("path:{}:{path_text}", host.to_lowercase()),
1818 );
1819 Self {
1820 key,
1821 short: path_leaf(path),
1822 full,
1823 }
1824 }
1825}
1826
1827fn normalize_github_source(source: &str) -> Option<String> {
1828 let source = source.trim();
1829 let path = source
1830 .strip_prefix("https://github.com/")
1831 .or_else(|| source.strip_prefix("http://github.com/"))
1832 .or_else(|| source.strip_prefix("git@github.com:"))
1833 .or_else(|| source.strip_prefix("ssh://git@github.com/"))
1834 .or_else(|| {
1835 (!source.contains("://") && !source.contains('@') && !source.contains(':'))
1836 .then_some(source)
1837 })?
1838 .trim_end_matches(".git");
1839 let mut parts = path.split('/');
1840 let owner = parts.next()?;
1841 let repository = parts.next()?;
1842 (!owner.is_empty() && !repository.is_empty() && parts.next().is_none())
1843 .then(|| format!("{owner}/{repository}"))
1844}
1845
1846fn path_leaf(path: &Path) -> String {
1848 path.file_name()
1849 .unwrap_or(path.as_os_str())
1850 .to_string_lossy()
1851 .into_owned()
1852}
1853
1854fn created_at_seconds(timestamp: &str) -> Option<i64> {
1855 chrono::DateTime::parse_from_rfc3339(timestamp)
1856 .ok()
1857 .map(|timestamp| timestamp.timestamp())
1858}
1859
1860#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1861#[serde(deny_unknown_fields)]
1862pub struct State {
1863 pub version: u32,
1864 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1865 pub sessions: BTreeMap<String, SessionRecord>,
1866 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1869 pub subagents: BTreeMap<String, SubagentRecord>,
1870 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1872 pub mount_history: BTreeMap<String, Vec<PathBuf>>,
1873 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1875 pub container_sizes: BTreeMap<String, HostContainerSize>,
1876}
1877
1878impl Default for State {
1879 fn default() -> Self {
1880 Self {
1881 version: STATE_VERSION,
1882 sessions: BTreeMap::new(),
1883 subagents: BTreeMap::new(),
1884 mount_history: BTreeMap::new(),
1885 container_sizes: BTreeMap::new(),
1886 }
1887 }
1888}
1889
1890impl State {
1891 #[must_use]
1896 pub fn session_notice_name(&self, session_id: &str) -> String {
1897 match self.sessions.get(session_id) {
1898 Some(session) if session.listed_title() != session.id => {
1899 session.listed_title().to_owned()
1900 }
1901 _ => short_id(session_id).to_owned(),
1902 }
1903 }
1904
1905 #[must_use]
1914 pub fn project_identity_session<'a>(&'a self, session: &'a SessionRecord) -> &'a SessionRecord {
1915 self.subagents
1916 .get(&session.id)
1917 .and_then(|record| self.sessions.get(&record.parent_session_id))
1918 .unwrap_or(session)
1919 }
1920
1921 #[must_use]
1928 pub fn is_subagent_session(&self, id: &str) -> bool {
1929 self.subagents.contains_key(id) || crate::native_agent::is_view_id(id)
1930 }
1931
1932 pub fn validate(&self) -> Result<()> {
1933 if self.version != STATE_VERSION {
1934 bail!(
1935 "unsupported Mjolnir state version {}; expected {STATE_VERSION}",
1936 self.version
1937 );
1938 }
1939 for (id, session) in &self.sessions {
1940 session.validate(id)?;
1941 }
1942 for (child_id, subagent) in &self.subagents {
1943 if child_id != &subagent.child_session_id {
1944 bail!("sub-agent key {child_id:?} does not match its child session id");
1945 }
1946 if child_id == &subagent.parent_session_id {
1947 bail!("sub-agent {child_id:?} cannot be its own parent");
1948 }
1949 if !self.sessions.contains_key(child_id) {
1950 bail!("sub-agent {child_id:?} has no child session");
1951 }
1952 if !self.sessions.contains_key(&subagent.parent_session_id) {
1953 bail!(
1954 "sub-agent {child_id:?} has unknown parent {:?}",
1955 subagent.parent_session_id
1956 );
1957 }
1958 if self.subagents.contains_key(&subagent.parent_session_id) {
1959 bail!("sub-agent {child_id:?} cannot belong to another sub-agent");
1960 }
1961 if subagent.task_name.trim().is_empty()
1962 || subagent.profile_id.trim().is_empty()
1963 || subagent.request_key.trim().is_empty()
1964 {
1965 bail!("sub-agent {child_id:?} has incomplete relationship metadata");
1966 }
1967 }
1968 for (host, sources) in &self.mount_history {
1969 if host.trim().is_empty() {
1970 bail!("mount history contains an empty host key");
1971 }
1972 if sources.iter().any(|source| !source.is_absolute()) {
1973 bail!("mount history for {host:?} contains a non-absolute source path");
1974 }
1975 }
1976 for (host, size) in &self.container_sizes {
1977 if host.trim().is_empty() {
1978 bail!("container size history contains an empty host key");
1979 }
1980 if size.cpus == 0 || size.memory_bytes == 0 {
1981 bail!("container size history for {host:?} contains a zero value");
1982 }
1983 if size.cpus > i64::MAX as u64 || size.memory_bytes > i64::MAX as u64 {
1984 bail!("container size history for {host:?} exceeds SQLite integer range");
1985 }
1986 }
1987 Ok(())
1988 }
1989
1990 pub fn remember_mount_sources(&mut self, host: &str, mounts: &[AdditionalMount]) {
1991 if mounts.is_empty() {
1992 return;
1993 }
1994 let sources = self.mount_history.entry(host.to_owned()).or_default();
1995 for mount in mounts.iter().rev() {
1996 sources.retain(|source| source != &mount.source);
1997 sources.insert(0, mount.source.clone());
1998 }
1999 sources.truncate(20);
2000 }
2001
2002 pub fn remember_container_size(&mut self, host: &str, size: HostContainerSize) {
2003 self.container_sizes.insert(host.to_owned(), size);
2004 }
2005
2006 pub fn project_directories(&self, host: &str) -> &[PathBuf] {
2007 self.mount_history
2008 .get(&project_history_key(host))
2009 .map(Vec::as_slice)
2010 .unwrap_or_default()
2011 }
2012
2013 pub fn remember_project_directory(&mut self, host: &str, directory: &Path) {
2014 let key = project_history_key(host);
2015 let directories = self.mount_history.entry(key).or_default();
2016 directories.retain(|existing| existing != directory);
2017 directories.insert(0, directory.to_path_buf());
2018 directories.truncate(20);
2019 }
2020
2021 pub fn destroy_stopped_session(&mut self, session_id: &str) -> Result<SessionRecord> {
2022 let session = self
2023 .sessions
2024 .get(session_id)
2025 .with_context(|| format!("unknown session {session_id}"))?;
2026 if session.state.is_active() {
2027 bail!("refusing to destroy active session {session_id}");
2028 }
2029 Ok(self
2030 .sessions
2031 .remove(session_id)
2032 .expect("session checked above"))
2033 }
2034
2035 pub fn destroy_session_force(&mut self, session_id: &str) -> Result<SessionRecord> {
2041 self.sessions
2042 .get(session_id)
2043 .with_context(|| format!("unknown session {session_id}"))?;
2044 Ok(self
2045 .sessions
2046 .remove(session_id)
2047 .expect("session checked above"))
2048 }
2049
2050 pub fn validate_setup_update(&self, before: &Config, after: &Config) -> Result<()> {
2053 for session in self
2054 .sessions
2055 .values()
2056 .filter(|session| session.state.is_active())
2057 {
2058 let protected = if let Some(profile) = before.profiles.get(&session.last_profile) {
2059 let mut comparable = profile.clone();
2060 if let Some(updated) = after.profiles.get(&session.last_profile) {
2061 comparable.enabled = updated.enabled;
2062 }
2063 profile.kind == session.harness_kind
2065 && after.profiles.get(&session.last_profile) != Some(&comparable)
2066 } else {
2067 false
2068 };
2069 let bundle_changed = session.project_directory.is_none()
2070 && before
2071 .bundles
2072 .get(&session.bundle_id)
2073 .is_some_and(|bundle| after.bundles.get(&session.bundle_id) != Some(bundle));
2074 let target_changed =
2078 before
2079 .targets
2080 .get(&session.target_template_id)
2081 .is_some_and(|target| {
2082 after
2083 .targets
2084 .get(&session.target_template_id)
2085 .map(TargetTemplate::without_launch_only_settings)
2086 != Some(target.without_launch_only_settings())
2087 });
2088 if protected || bundle_changed || target_changed {
2089 let mut used = Vec::new();
2093 if protected {
2094 used.push(format!("agent profile {:?}", session.last_profile));
2095 }
2096 if bundle_changed {
2097 used.push(format!("project {:?}", session.project_name(before)));
2098 }
2099 if target_changed {
2100 used.push(format!("runtime {:?}", session.target_template_id));
2101 }
2102 let used = match used.as_slice() {
2103 [only] => only.clone(),
2104 [rest @ .., last] => format!("{} and {last}", rest.join(", ")),
2105 [] => unreachable!("something changed"),
2106 };
2107 let title = session.display_title();
2108 let named = if title == session.id {
2109 format!(
2110 "a running session in project {:?}",
2111 session.project_name(before)
2112 )
2113 } else {
2114 format!("the running session {title:?}")
2115 };
2116 bail!(
2117 "Setup would change the {used} that {named} uses. Save the new settings under a new name, or stop the session first."
2118 );
2119 }
2120 }
2121 Ok(())
2122 }
2123
2124 pub fn validate_against_config(&self, config: &Config) -> Result<()> {
2126 self.validate()?;
2127 config.validate()?;
2128 for session in self.sessions.values() {
2129 session.validate_configuration(config)?;
2130 }
2131 Ok(())
2132 }
2133}
2134
2135fn project_history_key(host: &str) -> String {
2136 format!("project:{host}")
2137}
2138
2139pub fn new_session_id() -> Result<String> {
2141 let mut random = [0u8; 16];
2142 getrandom::fill(&mut random)
2143 .map_err(|error| anyhow::anyhow!("generate Mjolnir session id: {error}"))?;
2144 Ok(crate::hex::lower_hex(random))
2145}
2146
2147pub fn harness_session_title(events: &[SequencedEvent]) -> Option<String> {
2149 events.iter().rev().find_map(|event| {
2150 let WorkerEvent::Adapter { payload, .. } = &event.event else {
2151 return None;
2152 };
2153 let crate::acp::RuntimeEvent::SessionUpdate { update } =
2154 serde_json::from_value(payload.clone()).ok()?
2155 else {
2156 return None;
2157 };
2158 let kind = update
2159 .get("sessionUpdate")
2160 .and_then(serde_json::Value::as_str)?;
2161 let title = match kind {
2162 "session_info_update" | "session_title" => {
2163 update.get("title").and_then(serde_json::Value::as_str)
2164 }
2165 _ => None,
2166 }?;
2167 normalize_session_title(title)
2168 })
2169}
2170
2171pub fn default_session_title(
2176 project_directory: Option<&Path>,
2177 bundle_id: &str,
2178 profile_id: &str,
2179) -> String {
2180 let project = project_directory.and_then(Path::file_name).map_or_else(
2181 || bundle_id.to_owned(),
2182 |name| name.to_string_lossy().into_owned(),
2183 );
2184 format!("{project} via {profile_id}")
2185}
2186
2187pub fn normalize_session_title(title: &str) -> Option<String> {
2188 let normalized = crate::relay::strip_hidden_prompt_context(title)
2189 .split_whitespace()
2190 .collect::<Vec<_>>()
2191 .join(" ");
2192 (!normalized.is_empty()).then_some(normalized)
2193}
2194
2195pub fn provisional_session_title(prompt: &str) -> Option<String> {
2201 const MAX_TITLE_CHARS: usize = 64;
2202
2203 let normalized = normalize_session_title(prompt)?;
2204 if normalized.chars().count() <= MAX_TITLE_CHARS {
2205 return Some(normalized);
2206 }
2207
2208 let mut truncated = normalized
2209 .chars()
2210 .take(MAX_TITLE_CHARS - 1)
2211 .collect::<String>();
2212 if let Some(boundary) = truncated.rfind(char::is_whitespace) {
2213 truncated.truncate(boundary);
2214 }
2215 truncated.push('…');
2216 Some(truncated)
2217}
2218
2219pub fn short_id(id: &str) -> &str {
2220 id.get(..8).unwrap_or(id)
2221}
2222
2223#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
2224pub struct RecoveryCandidate {
2225 pub session_id: String,
2226 pub target_template_id: String,
2227 pub locator: TargetLocator,
2228 pub ownership: Option<crate::worker_launch::WorkerOwnership>,
2229 #[serde(default)]
2232 pub instance_id: Option<String>,
2233 #[serde(default, skip_serializing_if = "Option::is_none")]
2238 pub tracked_session: Option<SessionState>,
2239}
2240
2241#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
2242pub struct RecoveryScan {
2243 pub candidates: Vec<RecoveryCandidate>,
2244 pub warnings: Vec<String>,
2245 #[serde(default)]
2247 pub instance_id: String,
2248 #[serde(default)]
2251 pub hidden_other_instances: usize,
2252}
2253
2254#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2255#[serde(deny_unknown_fields)]
2256pub struct ResumeRepositorySourceReceipt {
2257 pub session_id: String,
2258 pub bundle_id: String,
2259 pub checkpoint_sha256: String,
2260 pub repositories: Vec<crate::config::ProjectRepository>,
2261}
2262
2263#[cfg(test)]
2264mod tests;