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::snapshot_map::SnapshotMap;
16use crate::subagent::SubagentRecord;
17use crate::targets::{AdditionalMount, validate_additional_mounts};
18
19pub const STATE_VERSION: u32 = 1;
20
21mod target_runtime;
22pub use target_runtime::{TargetConnection, TargetRuntimeSettings};
23
24mod session_configuration;
25pub use session_configuration::SessionConfiguration;
26
27mod session_move;
28pub use session_move::*;
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
31#[serde(rename_all = "kebab-case")]
32pub enum SessionState {
33 Provisioning,
34 Running,
35 Disconnected,
36 Checkpointing,
37 Closing,
38 Destroying,
39 #[serde(alias = "archived")]
42 Stopped,
43 Parked,
50 Lost,
51 Error,
52 DestroyedWithDataLoss,
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
58#[serde(rename_all = "kebab-case")]
59pub enum SessionTransitionKind {
60 Starting,
61 Resuming,
62 Moving,
63 Suspending,
64 Destroying,
65 Stopping,
67}
68
69impl SessionTransitionKind {
70 pub const fn label(self) -> &'static str {
71 match self {
72 Self::Starting => "Starting",
73 Self::Resuming => "Resuming",
74 Self::Moving => "Moving",
75 Self::Suspending => "Suspending",
76 Self::Destroying => "Destroying",
77 Self::Stopping => "Stopping",
78 }
79 }
80
81 pub fn for_session(state: SessionState, operation: Option<Self>) -> Option<Self> {
82 operation.or_else(|| state.transition_kind())
83 }
84}
85
86#[cfg(test)]
87mod transition_tests {
88 use super::{SessionState, SessionTransitionKind};
89
90 #[test]
91 fn operation_ownership_hides_intermediate_move_states_but_not_ordinary_live_work() {
92 for state in [
93 SessionState::Stopped,
94 SessionState::Running,
95 SessionState::Disconnected,
96 ] {
97 assert_eq!(
98 SessionTransitionKind::for_session(state, Some(SessionTransitionKind::Moving)),
99 Some(SessionTransitionKind::Moving)
100 );
101 assert_eq!(SessionTransitionKind::for_session(state, None), None);
102 }
103 assert_eq!(SessionState::Checkpointing.transition_kind(), None);
104 assert_eq!(
105 SessionState::Closing.transition_kind(),
106 Some(SessionTransitionKind::Suspending)
107 );
108 }
109}
110
111#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
113#[serde(tag = "state", rename_all = "snake_case")]
114pub enum MaterializedExecutionState {
115 #[default]
116 Idle,
117 Running {
118 started_at_ms: i64,
119 },
120 Closing,
121 Closed,
122}
123
124pub use crate::transcript::{TerminalOutputRecord, TranscriptBody, TranscriptItem};
125
126#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
131#[serde(rename_all = "snake_case")]
132pub enum QueuedCommandKind {
133 #[default]
134 Prompt,
135 SetConfig {
136 key: String,
137 value: String,
138 },
139}
140
141impl QueuedCommandKind {
142 pub fn is_prompt(&self) -> bool {
143 matches!(self, Self::Prompt)
144 }
145}
146
147pub fn config_command_text(key: &str, value: &str) -> String {
150 if key == "fast-mode" {
151 "/fast".to_owned()
152 } else {
153 format!("/{key} {value}")
154 }
155}
156
157#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
158#[serde(deny_unknown_fields)]
159pub struct MaterializedQueuedPrompt {
160 pub command_id: String,
161 #[serde(default, skip_serializing_if = "QueuedCommandKind::is_prompt")]
162 pub kind: QueuedCommandKind,
163 pub content: Vec<serde_json::Value>,
164 pub queued_at_ms: i64,
165 #[serde(default, skip_serializing_if = "Option::is_none")]
169 pub accepted_ordinal: Option<u64>,
170}
171
172#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
175#[serde(deny_unknown_fields)]
176pub struct MaterializedTurn {
177 pub command_id: String,
178 #[serde(default, skip_serializing_if = "Option::is_none")]
179 pub accepted_ordinal: Option<u64>,
180 pub turn_start_position: u64,
183 pub started_at_ms: i64,
184 #[serde(default, skip_serializing_if = "Option::is_none")]
187 pub steered_into: Option<String>,
188}
189
190impl MaterializedTurn {
191 pub fn belongs_to(&self, command_id: &str) -> bool {
193 self.command_id == command_id || self.steered_into.as_deref() == Some(command_id)
194 }
195}
196
197#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
199#[serde(tag = "kind", rename_all = "snake_case")]
200pub enum TurnOutcomeKind {
201 Completed { stop_reason: String },
203 Rejected {
205 message: String,
206 #[serde(default, skip_serializing_if = "Option::is_none")]
207 reason: Option<crate::event_outcome::OutcomeReason>,
208 },
209 Interrupted {
211 message: String,
212 #[serde(default, skip_serializing_if = "Option::is_none")]
213 reason: Option<crate::event_outcome::OutcomeReason>,
214 },
215}
216
217impl std::fmt::Display for TurnOutcomeKind {
220 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
221 use crate::event_outcome::{OutcomeReason, TurnResultKind};
222 let result = self.result();
223 match result.kind {
224 TurnResultKind::Completed => formatter.write_str("completed, end of turn"),
225 TurnResultKind::InputRequired => formatter.write_str("completed, waiting for input"),
226 TurnResultKind::Cancelled | TurnResultKind::Interrupted => {
227 formatter.write_str("interrupted")
228 }
229 TurnResultKind::Failed if result.reason == Some(OutcomeReason::QuotaLimit) => {
230 formatter.write_str("failed: quota limit reached")
231 }
232 TurnResultKind::Failed | TurnResultKind::Rejected => {
233 let message = result.message.as_deref().unwrap_or("unknown failure");
234 if result.stop_reason.is_some() {
235 write!(formatter, "failed: {}", stop_reason_words(message))
236 } else {
237 write!(
238 formatter,
239 "failed: {}",
240 message.lines().next().unwrap_or_default().trim()
241 )
242 }
243 }
244 }
245 }
246}
247
248fn stop_reason_words(stop_reason: &str) -> String {
251 let mut words = String::new();
252 let mut previous_lower = false;
253 for character in stop_reason.trim().chars() {
254 if character == '_' || character == '-' || character.is_whitespace() {
255 if !words.ends_with(' ') && !words.is_empty() {
256 words.push(' ');
257 }
258 previous_lower = false;
259 continue;
260 }
261 if character.is_uppercase() && previous_lower {
262 words.push(' ');
263 }
264 previous_lower = character.is_lowercase() || character.is_ascii_digit();
265 words.extend(character.to_lowercase());
266 }
267 match words.trim_end() {
268 "" => "no reason given".to_owned(),
269 words => words.to_owned(),
270 }
271}
272
273#[derive(Debug, Clone, Copy, PartialEq, Eq)]
274pub enum PromptCompletion {
275 InputRequired,
276 Finished,
277 Cancelled,
278 QuotaLimit,
279 Error,
280}
281
282pub fn classify_prompt_completion(stop_reason: &str) -> PromptCompletion {
284 let normalized = stop_reason
285 .chars()
286 .filter(|character| *character != '_' && *character != '-')
287 .flat_map(char::to_lowercase)
288 .collect::<String>();
289 match normalized.as_str() {
290 "endturn" => PromptCompletion::Finished,
291 "awaitinginput" => PromptCompletion::InputRequired,
292 "cancelled" | "canceled" => PromptCompletion::Cancelled,
293 "quotalimit" => PromptCompletion::QuotaLimit,
294 _ => PromptCompletion::Error,
295 }
296}
297
298#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
300#[serde(deny_unknown_fields)]
301pub struct MaterializedTurnOutcome {
302 #[serde(default, skip_serializing_if = "Option::is_none")]
303 pub diagnostic: Option<crate::diagnostic::TurnDiagnostic>,
304
305 #[serde(default, skip_serializing_if = "Option::is_none")]
306 pub usage: Option<crate::usage::TokenUsage>,
307 pub command_id: String,
308 #[serde(default, skip_serializing_if = "Option::is_none")]
309 pub accepted_ordinal: Option<u64>,
310 #[serde(default, skip_serializing_if = "Option::is_none")]
311 pub turn_start_position: Option<u64>,
312 pub completed_ordinal: u64,
313 pub completed_at_ms: i64,
314 pub outcome: TurnOutcomeKind,
315}
316
317impl MaterializedTurnOutcome {
318 pub fn interruption_ordinal(&self) -> Option<u64> {
320 (self.turn_start_position.is_some()
321 && matches!(self.outcome, TurnOutcomeKind::Interrupted { .. }))
322 .then_some(self.completed_ordinal)
323 }
324}
325
326#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
328#[serde(deny_unknown_fields)]
329pub struct MaterializedSession {
330 pub session_id: String,
331 pub applied_event_ordinal: u64,
332 pub applied_event_digest: String,
333 pub last_activity_at_ms: Option<i64>,
336 pub execution: MaterializedExecutionState,
337 #[serde(default, skip_serializing_if = "Option::is_none")]
338 pub session_title: Option<String>,
339 #[serde(default, skip_serializing_if = "SessionConfiguration::is_empty")]
340 pub configuration: SessionConfiguration,
341 #[serde(default, skip_serializing_if = "Vec::is_empty")]
342 pub transcript: Vec<Arc<TranscriptItem>>,
345 #[serde(default, skip_serializing_if = "Vec::is_empty")]
346 pub queued_prompts: Vec<MaterializedQueuedPrompt>,
347 #[serde(default, skip_serializing_if = "Vec::is_empty")]
350 pub pending_elicitations: Vec<crate::elicitation::ElicitationRequest>,
351 #[serde(default, skip_serializing_if = "Option::is_none")]
353 pub active_turn: Option<MaterializedTurn>,
354 #[serde(default, skip_serializing_if = "Option::is_none")]
357 pub last_turn_outcome: Option<MaterializedTurnOutcome>,
358}
359
360#[derive(Debug, Clone, PartialEq, Eq)]
363pub struct MaterializedSessionSummary {
364 pub session_id: String,
365 pub applied_event_ordinal: u64,
366 pub last_activity_at_ms: Option<i64>,
367 pub execution: MaterializedExecutionState,
368 pub session_title: Option<String>,
369 pub last_agent_message: Option<String>,
370 pub last_user_message: Option<String>,
371 pub last_agent_message_follows_last_user: bool,
374 pub agent_message_latest_content_ordinals: Vec<u64>,
375 pub interruption_event_ordinals: Vec<u64>,
376}
377
378impl MaterializedSession {
379 pub fn empty(session_id: impl Into<String>) -> Self {
380 Self {
381 session_id: session_id.into(),
382 applied_event_ordinal: 0,
383 applied_event_digest: RELAY_EVENT_GENESIS_DIGEST.into(),
384 last_activity_at_ms: None,
385 execution: MaterializedExecutionState::Idle,
386 session_title: None,
387 configuration: SessionConfiguration::default(),
388 transcript: Vec::new(),
389 queued_prompts: Vec::new(),
390 pending_elicitations: Vec::new(),
391 active_turn: None,
392 last_turn_outcome: None,
393 }
394 }
395
396 pub fn last_activity_at_ms(&self) -> Option<i64> {
397 self.last_activity_at_ms
398 }
399
400 pub fn resolved_title(&self) -> Option<String> {
406 self.session_title
407 .as_deref()
408 .and_then(normalize_session_title)
409 .or_else(|| {
410 self.transcript.iter().find_map(|item| {
411 let TranscriptBody::User { content } = &item.body else {
412 return None;
413 };
414 provisional_session_title(&crate::transcript::materialized_content_text(
415 content,
416 ))
417 })
418 })
419 .or_else(|| {
420 self.queued_prompts
421 .iter()
422 .filter(|prompt| prompt.kind.is_prompt())
423 .find_map(|prompt| {
424 provisional_session_title(&crate::transcript::materialized_content_text(
425 &prompt.content,
426 ))
427 })
428 })
429 }
430
431 pub fn unread_agent_messages_after(&self, viewed_through_event_ordinal: u64) -> u64 {
432 self.transcript
433 .iter()
434 .filter(|item| {
435 item.latest_content_event_ordinal
436 .is_some_and(|ordinal| ordinal > viewed_through_event_ordinal)
437 && item.is_nonempty_agent_message()
438 })
439 .count() as u64
440 }
441
442 pub fn unread_interruptions_after(&self, viewed_through_event_ordinal: u64) -> u64 {
443 self.interruption_event_ordinals()
444 .into_iter()
445 .filter(|ordinal| *ordinal > viewed_through_event_ordinal)
446 .count() as u64
447 }
448
449 pub fn interruption_event_ordinals(&self) -> Vec<u64> {
450 let mut ordinals = self
451 .transcript
452 .iter()
453 .filter(|item| item.is_work_interruption())
454 .map(|item| item.position)
455 .collect::<Vec<_>>();
456 if let Some(ordinal) = self
457 .last_turn_outcome
458 .as_ref()
459 .and_then(MaterializedTurnOutcome::interruption_ordinal)
460 {
461 ordinals.push(ordinal);
462 }
463 ordinals.sort_unstable();
464 ordinals.dedup();
465 ordinals
466 }
467
468 pub fn validate(&self) -> Result<()> {
469 validate_id("session", &self.session_id)?;
470 validate_relay_event_frontier(
471 self.applied_event_ordinal,
472 &self.applied_event_digest,
473 "materialized session event frontier",
474 )?;
475 if self
476 .session_title
477 .as_ref()
478 .is_some_and(|title| title.trim().is_empty())
479 {
480 bail!("materialized session has an empty title");
481 }
482 let mut item_ids = BTreeSet::new();
483 for item in &self.transcript {
484 item.validate(self.applied_event_ordinal)?;
485 if !item_ids.insert(item.stable_id.as_str()) {
486 bail!(
487 "materialized transcript contains duplicate item {:?}",
488 item.stable_id
489 );
490 }
491 }
492 let mut command_ids = BTreeSet::new();
493 for prompt in &self.queued_prompts {
494 if prompt.command_id.trim().is_empty() {
495 bail!("materialized prompt queue has an empty command id");
496 }
497 if !command_ids.insert(prompt.command_id.as_str()) {
498 bail!(
499 "materialized prompt queue contains duplicate command {:?}",
500 prompt.command_id
501 );
502 }
503 if let QueuedCommandKind::SetConfig { key, value } = &prompt.kind
504 && (key.trim().is_empty() || value.trim().is_empty())
505 {
506 bail!(
507 "materialized queued configuration change {:?} is incomplete",
508 prompt.command_id
509 );
510 }
511 }
512 Ok(())
513 }
514}
515
516#[derive(Debug, Clone, PartialEq)]
520pub struct ManagedSessionSnapshot {
521 pub materialized: MaterializedSession,
522 pub window: ProjectionWindow,
525 pub operational: RelayOperationalState,
526 pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
530 pub worker_build: Option<String>,
535 pub subagent_requests: Vec<crate::subagent::SubagentToolRequest>,
537 pub subagent_results: Vec<crate::subagent::SubagentToolResult>,
539}
540
541#[derive(Debug, Clone, PartialEq, Eq)]
553pub struct ProjectionWindow {
554 pub omitted_items: usize,
556 pub provisional_title: Option<String>,
558 pub latest_turn_start_position: Option<u64>,
562}
563
564impl ProjectionWindow {
565 pub fn trim(&mut self, session: &mut MaterializedSession, target: usize) {
568 let observed = Self::of(session);
569 if self.provisional_title.is_none() {
570 self.provisional_title = observed.provisional_title;
571 }
572 self.latest_turn_start_position = observed
573 .latest_turn_start_position
574 .or(self.latest_turn_start_position);
575 let mut boundary = session.transcript.len().saturating_sub(target.max(1));
576 for (index, item) in session.transcript.iter().enumerate() {
577 let mutable = match &item.body {
578 TranscriptBody::Agent { streaming, .. }
579 | TranscriptBody::Thought { streaming, .. } => *streaming,
580 TranscriptBody::Tool { call, .. } => matches!(
581 call.get("status").and_then(serde_json::Value::as_str),
582 Some("pending" | "in_progress")
583 ),
584 _ => false,
585 };
586 if mutable || Some(item.position) == self.latest_turn_start_position {
587 boundary = boundary.min(index);
588 }
589 }
590 let cut = session
591 .transcript
592 .iter()
593 .take(boundary + 1)
594 .rposition(|item| item.is_turn_start())
595 .unwrap_or(0);
596 if cut > 0 {
597 session.transcript.drain(..cut);
598 self.omitted_items += cut;
599 }
600 }
601
602 #[must_use]
604 pub fn of(session: &MaterializedSession) -> Self {
605 Self {
606 omitted_items: 0,
607 provisional_title: session.transcript.iter().find_map(|item| {
608 let TranscriptBody::User { content } = &item.body else {
609 return None;
610 };
611 provisional_session_title(&crate::transcript::materialized_content_text(content))
612 }),
613 latest_turn_start_position: session
614 .transcript
615 .iter()
616 .rev()
617 .find(|item| item.is_turn_start())
618 .map(|item| item.position),
619 }
620 }
621}
622
623impl ManagedSessionSnapshot {
624 #[must_use]
629 pub fn resolved_title(&self) -> Option<String> {
630 self.materialized
631 .session_title
632 .as_deref()
633 .and_then(normalize_session_title)
634 .or_else(|| self.window.provisional_title.clone())
635 .or_else(|| {
636 self.materialized
637 .queued_prompts
638 .iter()
639 .filter(|prompt| prompt.kind.is_prompt())
640 .find_map(|prompt| {
641 provisional_session_title(&crate::transcript::materialized_content_text(
642 &prompt.content,
643 ))
644 })
645 })
646 }
647
648 #[must_use]
653 pub fn latest_completed_turn_ordinal(&self) -> Option<u64> {
654 if self.materialized.execution != MaterializedExecutionState::Idle {
655 return None;
656 }
657 self.window.latest_turn_start_position
658 }
659}
660
661#[derive(Debug, Clone)]
663pub struct RecoveryObservation {
664 pub session: SessionRecord,
665 pub config: Config,
666 pub latest_completed_turn_ordinal: Option<u64>,
667 pub checkpoint_wait: Option<crate::activity::CheckpointWait>,
671}
672
673pub fn latest_completed_turn_ordinal(session: &MaterializedSession) -> Option<u64> {
678 if session.execution != MaterializedExecutionState::Idle {
679 return None;
680 }
681 session
682 .transcript
683 .iter()
684 .rev()
685 .find(|item| item.is_turn_start())
686 .map(|item| item.position)
687}
688
689pub fn validate_relay_event_digest(digest: &str, name: &str) -> Result<()> {
690 if digest.len() != 64
691 || !digest
692 .bytes()
693 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
694 {
695 bail!("{name} must be a lowercase SHA-256 digest");
696 }
697 Ok(())
698}
699
700pub fn validate_relay_event_frontier(ordinal: u64, digest: &str, name: &str) -> Result<()> {
701 validate_relay_event_digest(digest, name)?;
702 if (ordinal == 0) != (digest == RELAY_EVENT_GENESIS_DIGEST) {
703 bail!("{name} has inconsistent ordinal {ordinal} and digest {digest}");
704 }
705 Ok(())
706}
707
708fn is_false(value: &bool) -> bool {
709 !*value
710}
711
712impl SessionState {
713 pub const fn as_str(self) -> &'static str {
715 match self {
716 Self::Provisioning => "provisioning",
717 Self::Running => "running",
718 Self::Disconnected => "disconnected",
719 Self::Checkpointing => "checkpointing",
720 Self::Closing => "closing",
721 Self::Destroying => "destroying",
722 Self::Stopped => "stopped",
723 Self::Parked => "parked",
724 Self::Lost => "lost",
725 Self::Error => "error",
726 Self::DestroyedWithDataLoss => "destroyed-with-data-loss",
727 }
728 }
729
730 pub fn from_stored(value: &str) -> Option<Self> {
733 Some(match value {
734 "provisioning" => Self::Provisioning,
735 "running" => Self::Running,
736 "disconnected" => Self::Disconnected,
737 "checkpointing" => Self::Checkpointing,
738 "closing" => Self::Closing,
739 "destroying" => Self::Destroying,
740 "stopped" | "archived" => Self::Stopped,
741 "parked" => Self::Parked,
742 "lost" => Self::Lost,
743 "error" => Self::Error,
744 "destroyed-with-data-loss" => Self::DestroyedWithDataLoss,
745 _ => return None,
746 })
747 }
748
749 pub const fn transition_kind(self) -> Option<SessionTransitionKind> {
752 match self {
753 Self::Provisioning => Some(SessionTransitionKind::Starting),
754 Self::Closing => Some(SessionTransitionKind::Suspending),
755 Self::Destroying => Some(SessionTransitionKind::Destroying),
756 _ => None,
757 }
758 }
759
760 pub const fn is_active(self) -> bool {
767 matches!(
768 self,
769 Self::Provisioning
770 | Self::Running
771 | Self::Disconnected
772 | Self::Checkpointing
773 | Self::Closing
774 | Self::Destroying
775 | Self::Parked
776 | Self::Error
777 )
778 }
779
780 pub const fn has_live_worker(self) -> bool {
784 matches!(
785 self,
786 Self::Provisioning
787 | Self::Running
788 | Self::Disconnected
789 | Self::Checkpointing
790 | Self::Closing
791 | Self::Destroying
792 )
793 }
794}
795
796#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
797#[serde(tag = "kind", rename_all = "kebab-case")]
798pub enum PodmanWorkspaceLocator {
799 #[default]
800 ContainerLayer,
801 Volume {
802 name: String,
803 },
804 HostPath {
805 path: PathBuf,
806 helper: Vec<String>,
807 resource: String,
808 },
809}
810
811#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
812#[serde(tag = "kind", rename_all = "kebab-case")]
813pub enum TargetLocator {
814 LocalBare {
815 worker_root: PathBuf,
816 },
817 LocalPodman {
818 container_id: String,
819 #[serde(default)]
820 workspace_storage: PodmanWorkspaceLocator,
821 #[serde(default, skip_serializing_if = "Option::is_none")]
825 borrowed_from: Option<String>,
826 },
827 LocalDocker {
828 container_id: String,
829 #[serde(default, skip_serializing_if = "Option::is_none")]
833 borrowed_from: Option<String>,
834 },
835 AppleContainer {
836 container_id: String,
837 #[serde(default, skip_serializing_if = "Option::is_none")]
841 borrowed_from: Option<String>,
842 },
843 AwsEc2 {
844 instance_id: String,
845 #[serde(default, skip_serializing_if = "Option::is_none")]
846 address: Option<String>,
847 },
848 SshBare {
849 host: String,
850 workspace: PathBuf,
851 #[serde(default, skip_serializing_if = "Option::is_none")]
852 worker_id: Option<String>,
853 },
854 SshPodman {
855 host: String,
856 container_id: String,
857 #[serde(default)]
858 workspace_storage: PodmanWorkspaceLocator,
859 #[serde(default, skip_serializing_if = "Option::is_none")]
863 borrowed_from: Option<String>,
864 },
865 SshDocker {
866 host: String,
867 container_id: String,
868 #[serde(default, skip_serializing_if = "Option::is_none")]
872 borrowed_from: Option<String>,
873 },
874}
875
876impl ManagedWorktreeTarget {
877 pub fn same_location(&self, other: &Self) -> bool {
885 match (self, other) {
886 (Self::Local, Self::Local) => true,
887 (
888 Self::Ssh {
889 destination,
890 ssh_args,
891 },
892 Self::Ssh {
893 destination: other_destination,
894 ssh_args: other_args,
895 },
896 ) => {
897 destination == other_destination
898 && ssh_location_option(ssh_args, 'p', "port")
899 == ssh_location_option(other_args, 'p', "port")
900 && ssh_location_option(ssh_args, 'l', "user")
901 == ssh_location_option(other_args, 'l', "user")
902 }
903 _ => false,
904 }
905 }
906}
907
908fn ssh_location_option(args: &[String], flag: char, option: &str) -> Option<String> {
912 let mut args = args.iter();
913 while let Some(argument) = args.next() {
914 let Some(rest) = argument.strip_prefix('-') else {
915 continue;
916 };
917 let mut chars = rest.chars();
918 let Some(name) = chars.next() else { continue };
919 if name != flag && name != 'o' {
920 continue;
921 }
922 let inline = chars.as_str();
923 let value = if inline.is_empty() {
924 args.next().cloned()
925 } else {
926 Some(inline.to_owned())
927 };
928 if name == flag {
929 return value;
930 }
931 if let Some(setting) = value {
932 let (key, found) = setting
933 .split_once(['=', ' ', '\t'])
934 .unwrap_or((setting.as_str(), ""));
935 if key.trim().eq_ignore_ascii_case(option) {
936 return Some(found.trim().to_owned());
937 }
938 }
939 }
940 None
941}
942
943#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
944#[serde(tag = "kind", rename_all = "kebab-case")]
945pub enum ManagedWorktreeTarget {
946 Local,
947 Ssh {
948 destination: String,
949 #[serde(default, skip_serializing_if = "Vec::is_empty")]
950 ssh_args: Vec<String>,
951 },
952}
953
954#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
956#[serde(deny_unknown_fields)]
957pub struct ManagedWorktreeOptions {
958 pub available: bool,
959 pub default_create: bool,
960}
961
962#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
963#[serde(deny_unknown_fields)]
964pub struct ManagedWorktree {
965 #[serde(default, skip_serializing_if = "ManagedCheckoutKind::is_worktree")]
967 pub kind: ManagedCheckoutKind,
968 pub source_project_directory: PathBuf,
969 pub source_repository: PathBuf,
970 pub worktree_root: PathBuf,
971 pub branch: String,
972 pub target: ManagedWorktreeTarget,
973 #[serde(default, skip_serializing_if = "Option::is_none")]
977 pub base_commit: Option<String>,
978}
979
980#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
981#[serde(rename_all = "snake_case")]
982pub enum ManagedCheckoutKind {
983 #[default]
984 Worktree,
985 Clone,
986}
987
988impl ManagedCheckoutKind {
989 fn is_worktree(&self) -> bool {
990 matches!(self, Self::Worktree)
991 }
992}
993
994impl ManagedWorktree {
995 fn validate(&self, session_id: &str, project_directory: Option<&Path>) -> Result<()> {
996 for (label, path) in [
997 ("source project directory", &self.source_project_directory),
998 ("source repository", &self.source_repository),
999 ("worktree root", &self.worktree_root),
1000 ] {
1001 if !path.is_absolute() || path.components().any(|part| part == Component::ParentDir) {
1002 bail!("managed worktree {label} must be an absolute safe path");
1003 }
1004 }
1005 if !self
1006 .source_project_directory
1007 .starts_with(&self.source_repository)
1008 {
1009 bail!("managed worktree source directory is outside its repository");
1010 }
1011 let expected_root = self
1012 .source_repository
1013 .join(".mj")
1014 .join(match self.kind {
1015 ManagedCheckoutKind::Worktree => "worktrees",
1016 ManagedCheckoutKind::Clone => "clones",
1017 })
1018 .join(session_id);
1019 if self.worktree_root != expected_root {
1020 bail!("managed worktree root does not match the session-owned path");
1021 }
1022 if self.kind == ManagedCheckoutKind::Worktree && self.branch != format!("mj/{session_id}") {
1023 bail!("managed worktree branch does not match the session id");
1024 }
1025 if self.kind == ManagedCheckoutKind::Clone && self.branch.trim().is_empty() {
1026 bail!("managed clone has no starting branch");
1027 }
1028 let relative = self
1029 .source_project_directory
1030 .strip_prefix(&self.source_repository)
1031 .expect("source relationship checked above");
1032 if project_directory != Some(self.worktree_root.join(relative).as_path()) {
1033 bail!("session project directory does not match its managed worktree");
1034 }
1035 match &self.target {
1036 ManagedWorktreeTarget::Local => {}
1037 ManagedWorktreeTarget::Ssh { destination, .. } if destination.trim().is_empty() => {
1038 bail!("managed SSH worktree has an empty destination")
1039 }
1040 ManagedWorktreeTarget::Ssh { .. } => {}
1041 }
1042 Ok(())
1043 }
1044}
1045
1046#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1047#[serde(tag = "kind", rename_all = "kebab-case")]
1048pub enum SessionResourceAllocation {
1049 Container {
1050 cpus: u64,
1051 memory_bytes: u64,
1052 },
1053 AwsEc2 {
1054 instance_type: String,
1055 vcpus: u64,
1056 memory_bytes: u64,
1057 },
1058}
1059
1060impl SessionResourceAllocation {
1061 pub fn validate(&self) -> Result<()> {
1062 match self {
1063 Self::Container { cpus, memory_bytes } if *cpus == 0 || *memory_bytes == 0 => {
1064 bail!("container resource allocation must have non-zero CPU and memory")
1065 }
1066 Self::AwsEc2 {
1067 instance_type,
1068 vcpus,
1069 memory_bytes,
1070 } if instance_type.trim().is_empty() || *vcpus == 0 || *memory_bytes == 0 => {
1071 bail!("EC2 resource allocation must have an instance type, CPU, and memory")
1072 }
1073 _ => Ok(()),
1074 }
1075 }
1076}
1077
1078pub fn allocation_cpus(allocation: &SessionResourceAllocation) -> u64 {
1080 match allocation {
1081 SessionResourceAllocation::Container { cpus, .. } => *cpus,
1082 SessionResourceAllocation::AwsEc2 { vcpus, .. } => *vcpus,
1083 }
1084}
1085
1086pub fn allocation_memory(allocation: &SessionResourceAllocation) -> u64 {
1088 match allocation {
1089 SessionResourceAllocation::Container { memory_bytes, .. }
1090 | SessionResourceAllocation::AwsEc2 { memory_bytes, .. } => *memory_bytes,
1091 }
1092}
1093
1094impl TargetLocator {
1095 fn validate(&self, session_id: &str) -> Result<()> {
1096 match self {
1097 Self::LocalBare { worker_root } => {
1098 if !worker_root.is_absolute()
1099 || worker_root
1100 .components()
1101 .any(|part| part == Component::ParentDir)
1102 || !worker_root.ends_with(session_id)
1103 {
1104 bail!(
1105 "local bare worker root must be an absolute safe path ending in the session id"
1106 );
1107 }
1108 }
1109 Self::LocalPodman { container_id, .. }
1110 | Self::LocalDocker { container_id, .. }
1111 | Self::AppleContainer { container_id, .. }
1112 | Self::SshPodman { container_id, .. }
1113 | Self::SshDocker { container_id, .. }
1114 if container_id.trim().is_empty() =>
1115 {
1116 bail!("target locator has an empty container id")
1117 }
1118 Self::AwsEc2 { instance_id, .. } if instance_id.trim().is_empty() => {
1119 bail!("target locator has an empty AWS instance id")
1120 }
1121 Self::SshBare {
1122 host, workspace, ..
1123 } => {
1124 if host.trim().is_empty() {
1125 bail!("bare SSH target locator has an empty host");
1126 }
1127 if workspace.as_os_str().is_empty()
1128 || workspace
1129 .components()
1130 .any(|part| part == Component::ParentDir)
1131 || !workspace.ends_with(session_id)
1132 {
1133 bail!("bare SSH target locator must be a safe path ending in the session id");
1134 }
1135 }
1136 Self::SshPodman { host, .. } if host.trim().is_empty() => {
1137 bail!("SSH Podman target locator has an empty host")
1138 }
1139 Self::SshDocker { host, .. } if host.trim().is_empty() => {
1140 bail!("SSH Docker target locator has an empty host")
1141 }
1142 _ => {}
1143 }
1144 Ok(())
1145 }
1146}
1147
1148#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1149#[serde(deny_unknown_fields)]
1150pub struct CheckpointMetadata {
1151 pub archive_path: PathBuf,
1152 pub sha256: String,
1154 pub created_at: String,
1155 pub event_frontier: u64,
1156}
1157
1158#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1160#[serde(rename_all = "snake_case")]
1161pub enum PublicationState {
1162 Published,
1163 Unpublished,
1164 Unknown,
1165}
1166
1167#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1169pub struct PublicationAssessment {
1170 pub checkpoint_sha256: String,
1171 pub state: PublicationState,
1172 pub dirty: bool,
1173 pub stashed: bool,
1174 pub saved_commits: Vec<String>,
1175 pub destinations: Vec<String>,
1176 pub checked_at: String,
1177 pub reason: Option<String>,
1178}
1179
1180impl CheckpointMetadata {
1181 fn validate(&self) -> Result<()> {
1182 if self.archive_path.as_os_str().is_empty() {
1183 bail!("checkpoint archive path is empty");
1184 }
1185 if self.sha256.len() != 64
1186 || !self
1187 .sha256
1188 .bytes()
1189 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1190 {
1191 bail!("checkpoint SHA-256 must be 64 lowercase hexadecimal characters");
1192 }
1193 if self.created_at.trim().is_empty() {
1194 bail!("checkpoint timestamp is empty");
1195 }
1196 Ok(())
1197 }
1198}
1199
1200#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1204#[serde(deny_unknown_fields)]
1205pub struct SessionBuildCache {
1206 pub host: String,
1209 pub directory: PathBuf,
1210 #[serde(default, skip_serializing_if = "Option::is_none")]
1213 pub max_size: Option<String>,
1214 #[serde(default, skip_serializing_if = "Option::is_none")]
1217 pub target_root: Option<PathBuf>,
1218}
1219
1220#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1224pub struct ArchiveSpacePreview {
1225 pub sessions: usize,
1227 pub bytes: u64,
1229 pub reclaimable_sessions: usize,
1232 pub reclaimable_bytes: u64,
1233}
1234
1235#[derive(Debug, Clone, PartialEq, Eq)]
1239pub struct BuildCachePreview {
1240 pub native_mbx: Option<String>,
1242 pub directory: Option<PathBuf>,
1244 pub max_size: Option<BuildCacheLimit>,
1246 pub target_max_size: Option<BuildCacheLimit>,
1247 pub user_managed: bool,
1249 pub application: BuildCacheApplication,
1250 pub budget_note: Option<String>,
1252 pub stats: Option<BuildCacheStats>,
1254 pub off_reason: Option<BuildCacheOff>,
1257}
1258
1259#[derive(Debug, Clone, PartialEq, Eq)]
1267pub struct BuildCacheStats {
1268 pub builds: u64,
1270 pub cached_compilations: u64,
1272 pub avoided_compiler_ns: u64,
1274 pub reflinked_bytes: u64,
1276}
1277
1278#[derive(Debug, Clone, PartialEq, Eq)]
1280pub enum BuildCacheOff {
1281 TurnedOff,
1283 Unavailable(String),
1286}
1287
1288impl std::fmt::Display for BuildCacheOff {
1289 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1290 match self {
1291 Self::TurnedOff => formatter.write_str("turned off for this machine"),
1292 Self::Unavailable(reason) => formatter.write_str(reason),
1293 }
1294 }
1295}
1296
1297#[derive(Debug, Clone, PartialEq, Eq)]
1299pub enum BuildCacheLimit {
1300 Size(String),
1302 HostConfiguration(Option<String>),
1305 MjDefault(String),
1307 MbxDefault(Option<String>),
1309}
1310
1311#[derive(Debug, Clone, Default, PartialEq, Eq)]
1312pub enum BuildCacheApplication {
1313 #[default]
1314 Pending,
1315 Applied,
1316 Failed(String),
1317}
1318
1319#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1320#[serde(deny_unknown_fields)]
1321pub struct SessionRecord {
1322 pub id: String,
1323 #[serde(default = "default_session_workspace_id")]
1328 pub workspace_id: String,
1329 pub title: String,
1330 pub harness_kind: HarnessKind,
1331 pub last_profile: String,
1332 pub bundle_id: String,
1333 #[serde(default, skip_serializing_if = "Option::is_none")]
1335 pub project_directory: Option<PathBuf>,
1336 #[serde(default, skip_serializing_if = "Option::is_none")]
1338 pub managed_worktree: Option<ManagedWorktree>,
1339 #[serde(default, skip_serializing_if = "Option::is_none")]
1341 pub create_managed_worktree: Option<bool>,
1342 #[serde(default, skip_serializing_if = "Option::is_none")]
1349 pub launch_base: Option<String>,
1350 #[serde(default, skip_serializing_if = "Option::is_none")]
1354 pub launch_branch: Option<String>,
1355 #[serde(default, skip_serializing_if = "Option::is_none")]
1359 pub checkout: Option<crate::remote_git::ExactCheckout>,
1360 #[serde(default, skip_serializing_if = "Option::is_none")]
1361 pub expected_runtime_identity: Option<String>,
1362 #[serde(default, skip_serializing_if = "Option::is_none")]
1364 pub publication: Option<PublicationAssessment>,
1365 #[serde(
1367 default,
1368 alias = "mjolnir_subagents",
1369 deserialize_with = "crate::subagent::deserialize_optional_policy"
1370 )]
1371 pub subagents: Option<crate::subagent::SubagentPolicy>,
1372 pub target_template_id: String,
1373 #[serde(default, skip_serializing_if = "Option::is_none")]
1374 pub resource_allocation: Option<SessionResourceAllocation>,
1375 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1376 pub additional_mounts: Vec<AdditionalMount>,
1377 #[serde(default, skip_serializing_if = "Option::is_none")]
1380 pub container_cpus: Option<String>,
1381 #[serde(default, skip_serializing_if = "Option::is_none")]
1384 pub container_memory: Option<String>,
1385 #[serde(default, skip_serializing_if = "Option::is_none")]
1391 pub container_workspace: Option<PathBuf>,
1392 #[serde(default, skip_serializing_if = "Option::is_none")]
1396 pub build_cache: Option<SessionBuildCache>,
1397 pub state: SessionState,
1398 #[serde(default, skip_serializing_if = "is_false")]
1401 pub archived: bool,
1402 #[serde(default, skip_serializing_if = "Option::is_none")]
1403 pub target: Option<TargetLocator>,
1404 #[serde(default, skip_serializing_if = "Option::is_none")]
1406 pub target_runtime: Option<TargetRuntimeSettings>,
1407 #[serde(default, skip_serializing_if = "Option::is_none")]
1408 pub native_session_id: Option<String>,
1409 #[serde(default, skip_serializing_if = "Option::is_none")]
1410 pub acp_session_title: Option<String>,
1411 #[serde(default, skip_serializing_if = "Option::is_none")]
1412 pub session_title_override: Option<String>,
1413 pub created_at: String,
1414 pub updated_at: String,
1415 #[serde(default, alias = "detached_after_event_ordinal")]
1416 pub viewed_through_event_ordinal: u64,
1417 #[serde(default, skip_serializing_if = "String::is_empty")]
1420 pub draft_input: String,
1421 #[serde(default, skip_serializing_if = "Option::is_none")]
1429 pub last_error: Option<String>,
1430 #[serde(default, skip_serializing_if = "Option::is_none")]
1431 pub last_checkpoint_error: Option<String>,
1432 #[serde(default, skip_serializing_if = "Option::is_none")]
1433 pub checkpoint: Option<CheckpointMetadata>,
1434}
1435
1436#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1437#[serde(deny_unknown_fields)]
1438pub struct HostContainerSize {
1439 pub cpus: u64,
1440 pub memory_bytes: u64,
1441}
1442
1443fn default_session_workspace_id() -> String {
1444 crate::workspace::DEFAULT_WORKSPACE_ID.to_owned()
1445}
1446
1447pub const DESTRUCTION_FAILURE_PREFIX: &str = "the destruction did not finish";
1456
1457pub const CLOSE_FAILURE_PREFIX: &str = "the suspension did not finish";
1458
1459pub fn is_public_lifecycle_error(error: &str) -> bool {
1461 error.starts_with(CLOSE_FAILURE_PREFIX)
1462 || error.starts_with(DESTRUCTION_FAILURE_PREFIX)
1463 || error.starts_with("the close did not finish")
1464}
1465
1466#[must_use]
1479pub fn target_label(config: &Config, target_id: &str, project: Option<&Path>) -> String {
1480 if !matches!(
1481 config.targets.get(target_id),
1482 Some(TargetTemplate::LocalBare | TargetTemplate::SshBare { .. })
1483 ) {
1484 return target_id.to_owned();
1485 }
1486 project.and_then(Path::file_name).map_or_else(
1487 || target_id.to_owned(),
1488 |directory| format!("{target_id}/{}", directory.to_string_lossy()),
1489 )
1490}
1491
1492#[derive(Debug, Clone, Default, PartialEq, Eq)]
1494pub struct StartSelection {
1495 pub at: Option<String>,
1497 pub branch: Option<String>,
1499 pub base: Option<String>,
1501}
1502
1503impl SessionRecord {
1504 pub fn start_selection(&self) -> StartSelection {
1508 let at = self
1509 .checkout
1510 .as_ref()
1511 .map(|checkout| checkout.commit.clone());
1512 StartSelection {
1513 branch: self
1514 .checkout
1515 .as_ref()
1516 .and_then(|checkout| checkout.branch.clone())
1517 .or_else(|| self.launch_branch.clone()),
1518 base: self.launch_base.clone().or_else(|| at.clone()),
1519 at,
1520 }
1521 }
1522
1523 pub fn publication_state(&self) -> Option<PublicationState> {
1526 let independent_clone = self
1527 .managed_worktree
1528 .as_ref()
1529 .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Clone)
1530 || (self.managed_worktree.is_none() && self.project_directory.is_none());
1531 if !independent_clone {
1532 return None;
1533 }
1534 if self.state.is_active() {
1535 return Some(PublicationState::Unknown);
1536 }
1537 Some(
1538 self.checkpoint
1539 .as_ref()
1540 .zip(self.publication.as_ref())
1541 .filter(|(checkpoint, assessment)| {
1542 assessment.checkpoint_sha256 == checkpoint.sha256
1543 })
1544 .map_or(PublicationState::Unknown, |(_, assessment)| {
1545 if assessment.dirty || assessment.stashed {
1546 PublicationState::Unpublished
1547 } else {
1548 assessment.state
1549 }
1550 }),
1551 )
1552 }
1553
1554 pub fn target_runtime_settings<'a>(
1559 &'a self,
1560 config: &Config,
1561 ) -> Result<std::borrow::Cow<'a, TargetRuntimeSettings>> {
1562 if let Some(runtime) = &self.target_runtime {
1563 if let Some(refreshed) =
1566 config
1567 .targets
1568 .get(&self.target_template_id)
1569 .and_then(|template| {
1570 runtime.with_current_ssh_options(&TargetRuntimeSettings::from(template))
1571 })
1572 {
1573 return Ok(std::borrow::Cow::Owned(refreshed));
1574 }
1575 return Ok(std::borrow::Cow::Borrowed(runtime));
1576 }
1577 let template = config.targets.get(&self.target_template_id).ok_or_else(|| {
1578 crate::refusal::Refusal::precondition(format!(
1579 "Session {:?} has no recorded target access settings. Restore target {:?} in config.toml once, then retry.",
1580 self.id, self.target_template_id))
1581 })?;
1582 let runtime = TargetRuntimeSettings::from(template);
1583 if let Some(locator) = &self.target {
1584 crate::targets::TargetLocator::try_from(crate::targets::RecordedTarget {
1585 locator, runtime: Some(&runtime), session_id: &self.id,
1586 }).map_err(|error| crate::refusal::Refusal::precondition(format!(
1587 "Session {:?} cannot recover target {:?}: {error}. Restore its original target settings, then retry.", self.id, self.target_template_id)))?;
1588 }
1589 Ok(std::borrow::Cow::Owned(runtime))
1590 }
1591
1592 #[must_use]
1596 pub fn public_error(&self) -> Option<&str> {
1597 self.last_error
1598 .as_deref()
1599 .filter(|error| is_public_lifecycle_error(error))
1600 }
1601
1602 pub fn configuration_issue(&self, config: &Config) -> Option<String> {
1605 if !self.state.is_active() {
1606 return None;
1607 }
1608 let mut issues = Vec::new();
1609 match config.profiles.get(&self.last_profile) {
1610 None => issues.push(format!("missing profile {:?}", self.last_profile)),
1611 Some(profile) if profile.kind != self.harness_kind => issues.push(format!(
1612 "expects {:?}, but profile {:?} is {:?}",
1613 self.harness_kind, self.last_profile, profile.kind
1614 )),
1615 Some(_) => {}
1616 }
1617 if self.project_directory.is_none() && !config.bundles.contains_key(&self.bundle_id) {
1618 issues.push(format!("missing bundle {:?}", self.bundle_id));
1619 }
1620 if self.target_runtime.is_none() && !config.targets.contains_key(&self.target_template_id) {
1621 issues.push(format!(
1622 "missing target template {:?}",
1623 self.target_template_id
1624 ));
1625 }
1626 (!issues.is_empty()).then(|| format!(
1627 "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.",
1628 self.id, issues.join("; ")
1629 ))
1630 }
1631
1632 pub fn validate_configuration(&self, config: &Config) -> Result<()> {
1633 if let Some(issue) = self.configuration_issue(config) {
1634 return Err(crate::refusal::Refusal::precondition(issue).into());
1635 }
1636 Ok(())
1637 }
1638
1639 pub fn display_title(&self) -> &str {
1641 self.session_title_override
1642 .as_deref()
1643 .or(self.acp_session_title.as_deref())
1644 .unwrap_or(&self.id)
1645 }
1646
1647 pub fn listed_title(&self) -> &str {
1652 let named = self.session_title_override.is_some() || self.acp_session_title.is_some();
1653 if !named && !self.title.trim().is_empty() {
1654 return &self.title;
1655 }
1656 self.display_title()
1657 }
1658
1659 pub fn project_name(&self, config: &Config) -> String {
1664 if let Some(worktree) = &self.managed_worktree {
1665 return path_leaf(&worktree.source_repository);
1666 }
1667 if let Some(project_directory) = &self.project_directory {
1668 return path_leaf(project_directory);
1669 }
1670 self.bundle_source_name(config)
1671 }
1672
1673 pub fn project_target(&self, config: &Config, target_id: &str) -> String {
1677 let project = self
1678 .managed_worktree
1679 .as_ref()
1680 .map(|worktree| &worktree.source_project_directory)
1681 .or(self.project_directory.as_ref());
1682 target_label(config, target_id, project.map(PathBuf::as_path))
1683 }
1684
1685 pub fn project_source(&self, config: &Config) -> ProjectSourceIdentity {
1690 if let Some(worktree) = &self.managed_worktree {
1691 return ProjectSourceIdentity::path(&worktree.source_repository, None);
1692 }
1693 if let Some(project_directory) = &self.project_directory {
1694 let remote = match &self.target {
1695 Some(TargetLocator::SshBare { host, .. }) => Some(host.as_str()),
1696 _ => None,
1697 };
1698 return ProjectSourceIdentity::path(project_directory, remote);
1699 }
1700 self.bundle_source_identity(config)
1701 .unwrap_or_else(|| ProjectSourceIdentity {
1702 key: format!("bundle:{}", self.bundle_id),
1703 short: path_leaf(Path::new(&self.bundle_id)),
1704 full: self.bundle_id.clone(),
1705 })
1706 }
1707
1708 fn bundle_source_name(&self, config: &Config) -> String {
1711 self.bundle_source_identity(config)
1712 .map(|source| source.short)
1713 .unwrap_or_else(|| path_leaf(Path::new(&self.bundle_id)))
1714 }
1715
1716 fn bundle_source_identity(&self, config: &Config) -> Option<ProjectSourceIdentity> {
1719 let bundle = config.bundles.get(&self.bundle_id)?;
1720 let sources = bundle
1721 .repositories
1722 .iter()
1723 .map(repository_source_identity)
1724 .collect::<Option<Vec<_>>>()?;
1725 ProjectSourceIdentity::bundle(sources)
1726 }
1727
1728 pub fn compare_by_creation(&self, other: &Self) -> std::cmp::Ordering {
1732 self.creation_order_key().cmp(&other.creation_order_key())
1733 }
1734
1735 pub fn creation_order_key(&self) -> (bool, Option<i64>, &str) {
1737 let timestamp = created_at_seconds(&self.created_at);
1738 (timestamp.is_none(), timestamp, &self.id)
1739 }
1740
1741 fn validate(&self, map_id: &str) -> Result<()> {
1742 validate_id("session", &self.id)?;
1743 if self.id != map_id {
1744 bail!(
1745 "session map key {map_id:?} does not match record id {:?}",
1746 self.id
1747 );
1748 }
1749 validate_id("workspace", &self.workspace_id)?;
1750 validate_id("profile", &self.last_profile)?;
1751 validate_id("bundle", &self.bundle_id)?;
1752 if let Some(project_directory) = &self.project_directory
1753 && (!project_directory.is_absolute()
1754 || project_directory
1755 .components()
1756 .any(|part| part == Component::ParentDir))
1757 {
1758 bail!("session {:?} has an unsafe project directory", self.id);
1759 }
1760 if let Some(managed_worktree) = &self.managed_worktree {
1761 managed_worktree.validate(&self.id, self.project_directory.as_deref())?;
1762 }
1763 validate_id("target template", &self.target_template_id)?;
1764 if let Some(allocation) = &self.resource_allocation {
1765 allocation.validate()?;
1766 }
1767 validate_additional_mounts(&self.additional_mounts)?;
1768 if self.title.trim().is_empty() {
1769 bail!("session {:?} has an empty title", self.id);
1770 }
1771 if self
1772 .acp_session_title
1773 .as_ref()
1774 .is_some_and(|title| title.trim().is_empty())
1775 || self
1776 .session_title_override
1777 .as_ref()
1778 .is_some_and(|title| title.trim().is_empty())
1779 {
1780 bail!("session {:?} has an empty display title", self.id);
1781 }
1782 if self.created_at.trim().is_empty() || self.updated_at.trim().is_empty() {
1783 bail!("session {:?} has an empty timestamp", self.id);
1784 }
1785 if let Some(target) = &self.target {
1786 target.validate(&self.id)?;
1787 }
1788 if let Some(checkpoint) = &self.checkpoint {
1789 checkpoint.validate()?;
1790 }
1791 Ok(())
1792 }
1793}
1794
1795fn repository_source_identity(repository: &ProjectRepository) -> Option<ProjectSourceIdentity> {
1796 repository
1797 .github
1798 .as_deref()
1799 .and_then(ProjectSourceIdentity::git_remote)
1800 .or_else(|| {
1801 repository
1802 .local
1803 .as_deref()
1804 .map(|path| ProjectSourceIdentity::path(path, None))
1805 })
1806}
1807
1808#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1809pub struct ProjectSourceIdentity {
1810 pub key: String,
1811 pub short: String,
1812 pub full: String,
1813}
1814
1815impl ProjectSourceIdentity {
1816 pub fn bundle(mut sources: Vec<Self>) -> Option<Self> {
1818 if sources.is_empty() {
1819 return None;
1820 }
1821 sources.sort_by(|left, right| {
1822 left.key
1823 .cmp(&right.key)
1824 .then_with(|| left.full.cmp(&right.full))
1825 .then_with(|| left.short.cmp(&right.short))
1826 });
1827 sources.dedup_by(|left, right| left.key == right.key);
1828 if sources.len() == 1 {
1829 return sources.pop();
1830 }
1831 let keys = sources
1832 .iter()
1833 .map(|source| source.key.clone())
1834 .collect::<Vec<_>>();
1835 let key = serde_json::to_string(&keys).ok()?;
1836 Some(Self {
1837 key: format!("bundle:{key}"),
1838 short: sources
1839 .iter()
1840 .map(|source| source.short.as_str())
1841 .collect::<Vec<_>>()
1842 .join(" + "),
1843 full: sources
1844 .iter()
1845 .map(|source| source.full.as_str())
1846 .collect::<Vec<_>>()
1847 .join(" + "),
1848 })
1849 }
1850
1851 pub fn git_remote(source: &str) -> Option<Self> {
1854 if let Some(normalized) = normalize_github_source(source) {
1855 let short = normalized
1856 .rsplit_once('/')
1857 .map_or(normalized.as_str(), |(_, repository)| repository)
1858 .to_owned();
1859 return Some(Self {
1860 key: format!("github:{}", normalized.to_lowercase()),
1861 short,
1862 full: normalized,
1863 });
1864 }
1865 let normalized = source.trim().trim_end_matches('/').trim_end_matches(".git");
1866 if normalized.is_empty() {
1867 return None;
1868 }
1869 let short = normalized
1870 .rsplit(['/', ':'])
1871 .find(|part| !part.is_empty())
1872 .unwrap_or(normalized)
1873 .to_owned();
1874 Some(Self {
1875 key: format!("git:{}", normalized.to_lowercase()),
1876 short,
1877 full: normalized.to_owned(),
1878 })
1879 }
1880
1881 pub fn path(path: &Path, remote: Option<&str>) -> Self {
1883 let normalized = path.components().collect::<PathBuf>();
1884 let path_text = normalized.to_string_lossy().into_owned();
1885 let full = remote.map_or_else(|| path_text.clone(), |host| format!("{host}:{path_text}"));
1886 let key = remote.map_or_else(
1887 || format!("path:{path_text}"),
1888 |host| format!("path:{}:{path_text}", host.to_lowercase()),
1889 );
1890 Self {
1891 key,
1892 short: path_leaf(path),
1893 full,
1894 }
1895 }
1896}
1897
1898fn normalize_github_source(source: &str) -> Option<String> {
1899 let source = source.trim();
1900 let path = source
1901 .strip_prefix("https://github.com/")
1902 .or_else(|| source.strip_prefix("http://github.com/"))
1903 .or_else(|| source.strip_prefix("git@github.com:"))
1904 .or_else(|| source.strip_prefix("ssh://git@github.com/"))
1905 .or_else(|| {
1906 (!source.contains("://") && !source.contains('@') && !source.contains(':'))
1907 .then_some(source)
1908 })?
1909 .trim_end_matches(".git");
1910 let mut parts = path.split('/');
1911 let owner = parts.next()?;
1912 let repository = parts.next()?;
1913 (!owner.is_empty() && !repository.is_empty() && parts.next().is_none())
1914 .then(|| format!("{owner}/{repository}"))
1915}
1916
1917fn path_leaf(path: &Path) -> String {
1919 path.file_name()
1920 .unwrap_or(path.as_os_str())
1921 .to_string_lossy()
1922 .into_owned()
1923}
1924
1925fn created_at_seconds(timestamp: &str) -> Option<i64> {
1926 chrono::DateTime::parse_from_rfc3339(timestamp)
1927 .ok()
1928 .map(|timestamp| timestamp.timestamp())
1929}
1930
1931#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1932#[serde(deny_unknown_fields)]
1933pub struct State {
1934 #[serde(default)]
1935 pub last_subagent_policy: crate::subagent::SubagentPolicy,
1936 pub version: u32,
1937 #[serde(default, skip_serializing_if = "SnapshotMap::is_empty")]
1938 pub sessions: SnapshotMap<String, SessionRecord>,
1939 #[serde(default, skip_serializing_if = "SnapshotMap::is_empty")]
1942 pub subagents: SnapshotMap<String, SubagentRecord>,
1943 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1945 pub mount_history: BTreeMap<String, Vec<PathBuf>>,
1946 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1948 pub container_sizes: BTreeMap<String, HostContainerSize>,
1949}
1950
1951impl Default for State {
1952 fn default() -> Self {
1953 Self {
1954 version: STATE_VERSION,
1955 last_subagent_policy: Default::default(),
1956 sessions: SnapshotMap::new(),
1957 subagents: SnapshotMap::new(),
1958 mount_history: BTreeMap::new(),
1959 container_sizes: BTreeMap::new(),
1960 }
1961 }
1962}
1963
1964impl State {
1965 #[must_use]
1970 pub fn session_notice_name(&self, session_id: &str) -> String {
1971 match self.sessions.get(session_id) {
1972 Some(session) if session.listed_title() != session.id => {
1973 session.listed_title().to_owned()
1974 }
1975 _ => short_id(session_id).to_owned(),
1976 }
1977 }
1978
1979 #[must_use]
1988 pub fn project_identity_session<'a>(&'a self, session: &'a SessionRecord) -> &'a SessionRecord {
1989 self.subagents
1990 .get(&session.id)
1991 .and_then(|record| self.sessions.get(&record.parent_session_id))
1992 .unwrap_or(session)
1993 }
1994
1995 #[must_use]
2002 pub fn is_subagent_session(&self, id: &str) -> bool {
2003 self.subagents.contains_key(id) || crate::native_agent::is_view_id(id)
2004 }
2005
2006 pub fn validate(&self) -> Result<()> {
2007 if self.version != STATE_VERSION {
2008 bail!(
2009 "unsupported Mjolnir state version {}; expected {STATE_VERSION}",
2010 self.version
2011 );
2012 }
2013 for (id, session) in &self.sessions {
2014 session.validate(id)?;
2015 }
2016 for child_id in self.subagents.keys() {
2017 self.validate_subagent(child_id)?;
2018 }
2019 for (host, sources) in &self.mount_history {
2020 if host.trim().is_empty() {
2021 bail!("mount history contains an empty host key");
2022 }
2023 if sources.iter().any(|source| !source.is_absolute()) {
2024 bail!("mount history for {host:?} contains a non-absolute source path");
2025 }
2026 }
2027 for (host, size) in &self.container_sizes {
2028 if host.trim().is_empty() {
2029 bail!("container size history contains an empty host key");
2030 }
2031 if size.cpus == 0 || size.memory_bytes == 0 {
2032 bail!("container size history for {host:?} contains a zero value");
2033 }
2034 if size.cpus > i64::MAX as u64 || size.memory_bytes > i64::MAX as u64 {
2035 bail!("container size history for {host:?} exceeds SQLite integer range");
2036 }
2037 }
2038 Ok(())
2039 }
2040
2041 pub fn validate_subagent(&self, child_id: &str) -> Result<()> {
2043 let Some(subagent) = self.subagents.get(child_id) else {
2044 return Ok(());
2045 };
2046 if child_id != subagent.child_session_id {
2047 bail!("sub-agent key {child_id:?} does not match its child session id");
2048 }
2049 if child_id == subagent.parent_session_id {
2050 bail!("sub-agent {child_id:?} cannot be its own parent");
2051 }
2052 if !self.sessions.contains_key(child_id) {
2053 bail!("sub-agent {child_id:?} has no child session");
2054 }
2055 if !self.sessions.contains_key(&subagent.parent_session_id) {
2056 bail!(
2057 "sub-agent {child_id:?} has unknown parent {:?}",
2058 subagent.parent_session_id
2059 );
2060 }
2061 if self.subagents.contains_key(&subagent.parent_session_id) {
2062 bail!("sub-agent {child_id:?} cannot belong to another sub-agent");
2063 }
2064 if subagent.task_name.trim().is_empty()
2065 || subagent.profile_id.trim().is_empty()
2066 || subagent.request_key.trim().is_empty()
2067 {
2068 bail!("sub-agent {child_id:?} has incomplete relationship metadata");
2069 }
2070 Ok(())
2071 }
2072
2073 pub fn remember_mount_sources(&mut self, host: &str, mounts: &[AdditionalMount]) {
2074 if mounts.is_empty() {
2075 return;
2076 }
2077 let sources = self.mount_history.entry(host.to_owned()).or_default();
2078 for mount in mounts.iter().rev() {
2079 sources.retain(|source| source != &mount.source);
2080 sources.insert(0, mount.source.clone());
2081 }
2082 sources.truncate(20);
2083 }
2084
2085 pub fn remember_container_size(&mut self, host: &str, size: HostContainerSize) {
2086 self.container_sizes.insert(host.to_owned(), size);
2087 }
2088
2089 pub fn project_directories(&self, host: &str) -> &[PathBuf] {
2090 self.mount_history
2091 .get(&project_history_key(host))
2092 .map(Vec::as_slice)
2093 .unwrap_or_default()
2094 }
2095
2096 pub fn remember_project_directory(&mut self, host: &str, directory: &Path) {
2097 let key = project_history_key(host);
2098 let directories = self.mount_history.entry(key).or_default();
2099 directories.retain(|existing| existing != directory);
2100 directories.insert(0, directory.to_path_buf());
2101 directories.truncate(20);
2102 }
2103
2104 pub fn destroy_stopped_session(&mut self, session_id: &str) -> Result<SessionRecord> {
2105 let session = self
2106 .sessions
2107 .get(session_id)
2108 .with_context(|| format!("unknown session {session_id}"))?;
2109 if session.state.is_active() {
2110 bail!("refusing to destroy active session {session_id}");
2111 }
2112 Ok(self
2113 .sessions
2114 .remove(session_id)
2115 .expect("session checked above"))
2116 }
2117
2118 pub fn destroy_session_force(&mut self, session_id: &str) -> Result<SessionRecord> {
2124 self.sessions
2125 .get(session_id)
2126 .with_context(|| format!("unknown session {session_id}"))?;
2127 Ok(self
2128 .sessions
2129 .remove(session_id)
2130 .expect("session checked above"))
2131 }
2132
2133 pub fn bundle_users(&self, bundle_id: &str) -> Vec<&SessionRecord> {
2138 self.sessions
2139 .values()
2140 .filter(|session| {
2141 session.bundle_id == bundle_id
2142 && session.project_directory.is_none()
2143 && session.state != SessionState::DestroyedWithDataLoss
2144 })
2145 .collect()
2146 }
2147
2148 pub fn bundle_removal_refusal(&self, bundle_id: &str) -> Option<String> {
2151 let users = self.bundle_users(bundle_id);
2152 if users.is_empty() {
2153 return None;
2154 }
2155 let mut names = users
2156 .iter()
2157 .take(3)
2158 .map(|session| format!("{:?}", session.listed_title()))
2159 .collect::<Vec<_>>();
2160 if users.len() > 3 {
2161 names.push(format!("{} more", users.len() - 3));
2162 }
2163 Some(format!(
2164 "Project {bundle_id:?} is used by {}: {}. Destroy those sessions before removing it.",
2165 if users.len() == 1 {
2166 "a session"
2167 } else {
2168 "sessions"
2169 },
2170 names.join(", ")
2171 ))
2172 }
2173
2174 pub fn validate_setup_update(&self, before: &Config, after: &Config) -> Result<()> {
2177 for session in self
2178 .sessions
2179 .values()
2180 .filter(|session| session.state.is_active())
2181 {
2182 let protected = if let Some(profile) = before.profiles.get(&session.last_profile) {
2183 let mut comparable = profile.clone();
2184 if let Some(updated) = after.profiles.get(&session.last_profile) {
2185 comparable.enabled = updated.enabled;
2186 }
2187 profile.kind == session.harness_kind
2189 && after.profiles.get(&session.last_profile) != Some(&comparable)
2190 } else {
2191 false
2192 };
2193 let bundle_changed = session.project_directory.is_none()
2194 && before
2195 .bundles
2196 .get(&session.bundle_id)
2197 .is_some_and(|bundle| after.bundles.get(&session.bundle_id) != Some(bundle));
2198 let target_changed =
2202 before
2203 .targets
2204 .get(&session.target_template_id)
2205 .is_some_and(|target| {
2206 after
2207 .targets
2208 .get(&session.target_template_id)
2209 .map(TargetTemplate::without_launch_only_settings)
2210 != Some(target.without_launch_only_settings())
2211 });
2212 if protected || bundle_changed || target_changed {
2213 let mut used = Vec::new();
2217 if protected {
2218 used.push(format!("agent profile {:?}", session.last_profile));
2219 }
2220 if bundle_changed {
2221 used.push(format!("project {:?}", session.project_name(before)));
2222 }
2223 if target_changed {
2224 used.push(format!("runtime {:?}", session.target_template_id));
2225 }
2226 let used = match used.as_slice() {
2227 [only] => only.clone(),
2228 [rest @ .., last] => format!("{} and {last}", rest.join(", ")),
2229 [] => unreachable!("something changed"),
2230 };
2231 let title = session.display_title();
2232 let named = if title == session.id {
2233 format!(
2234 "a running session in project {:?}",
2235 session.project_name(before)
2236 )
2237 } else {
2238 format!("the running session {title:?}")
2239 };
2240 bail!(
2241 "Setup would change the {used} that {named} uses. Save the new settings under a new name, or stop the session first."
2242 );
2243 }
2244 }
2245 Ok(())
2246 }
2247
2248 pub fn validate_against_config(&self, config: &Config) -> Result<()> {
2250 self.validate()?;
2251 config.validate()?;
2252 for session in self.sessions.values() {
2253 session.validate_configuration(config)?;
2254 }
2255 Ok(())
2256 }
2257}
2258
2259fn project_history_key(host: &str) -> String {
2260 format!("project:{host}")
2261}
2262
2263pub fn new_session_id() -> Result<String> {
2265 let mut random = [0u8; 16];
2266 getrandom::fill(&mut random)
2267 .map_err(|error| anyhow::anyhow!("generate Mjolnir session id: {error}"))?;
2268 Ok(crate::hex::lower_hex(random))
2269}
2270
2271pub fn harness_session_title(events: &[SequencedEvent]) -> Option<String> {
2273 events.iter().rev().find_map(|event| {
2274 let WorkerEvent::Adapter { payload, .. } = &event.event else {
2275 return None;
2276 };
2277 let crate::acp::RuntimeEvent::SessionUpdate { update } =
2278 serde_json::from_value(payload.clone()).ok()?
2279 else {
2280 return None;
2281 };
2282 let kind = update
2283 .get("sessionUpdate")
2284 .and_then(serde_json::Value::as_str)?;
2285 let title = match kind {
2286 "session_info_update" | "session_title" => {
2287 update.get("title").and_then(serde_json::Value::as_str)
2288 }
2289 _ => None,
2290 }?;
2291 normalize_session_title(title)
2292 })
2293}
2294
2295pub fn default_session_title(
2300 project_directory: Option<&Path>,
2301 bundle_id: &str,
2302 profile_id: &str,
2303) -> String {
2304 let project = project_directory.and_then(Path::file_name).map_or_else(
2305 || bundle_id.to_owned(),
2306 |name| name.to_string_lossy().into_owned(),
2307 );
2308 format!("{project} via {profile_id}")
2309}
2310
2311pub fn normalize_session_title(title: &str) -> Option<String> {
2312 let normalized = crate::relay::strip_hidden_prompt_context(title)
2313 .split_whitespace()
2314 .collect::<Vec<_>>()
2315 .join(" ");
2316 (!normalized.is_empty()).then_some(normalized)
2317}
2318
2319pub fn provisional_session_title(prompt: &str) -> Option<String> {
2325 const MAX_TITLE_CHARS: usize = 64;
2326
2327 let normalized = normalize_session_title(prompt)?;
2328 if normalized.chars().count() <= MAX_TITLE_CHARS {
2329 return Some(normalized);
2330 }
2331
2332 let mut truncated = normalized
2333 .chars()
2334 .take(MAX_TITLE_CHARS - 1)
2335 .collect::<String>();
2336 if let Some(boundary) = truncated.rfind(char::is_whitespace) {
2337 truncated.truncate(boundary);
2338 }
2339 truncated.push('…');
2340 Some(truncated)
2341}
2342
2343pub fn short_id(id: &str) -> &str {
2344 id.get(..8).unwrap_or(id)
2345}
2346
2347#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
2348pub struct RecoveryCandidate {
2349 pub session_id: String,
2350 pub target_template_id: String,
2351 pub locator: TargetLocator,
2352 pub ownership: Option<crate::worker_launch::WorkerOwnership>,
2353 #[serde(default)]
2356 pub instance_id: Option<String>,
2357 #[serde(default, skip_serializing_if = "Option::is_none")]
2362 pub tracked_session: Option<SessionState>,
2363}
2364
2365#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
2366pub struct RecoveryScan {
2367 pub candidates: Vec<RecoveryCandidate>,
2368 pub warnings: Vec<String>,
2369 #[serde(default)]
2371 pub instance_id: String,
2372 #[serde(default)]
2375 pub hidden_other_instances: usize,
2376}
2377
2378#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2379#[serde(deny_unknown_fields)]
2380pub struct ResumeRepositorySourceReceipt {
2381 pub session_id: String,
2382 pub bundle_id: String,
2383 pub checkpoint_sha256: String,
2384 pub repositories: Vec<crate::config::ProjectRepository>,
2385}
2386
2387#[cfg(test)]
2388mod tests;