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 Lost,
40 Error,
41 DestroyedWithDataLoss,
42}
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
47#[serde(rename_all = "kebab-case")]
48pub enum SessionTransitionKind {
49 Starting,
50 Resuming,
51 Moving,
52 Suspending,
53 Destroying,
54}
55
56impl SessionTransitionKind {
57 pub const fn label(self) -> &'static str {
58 match self {
59 Self::Starting => "Starting",
60 Self::Resuming => "Resuming",
61 Self::Moving => "Moving",
62 Self::Suspending => "Suspending",
63 Self::Destroying => "Destroying",
64 }
65 }
66
67 pub fn for_session(state: SessionState, operation: Option<Self>) -> Option<Self> {
68 operation.or_else(|| state.transition_kind())
69 }
70}
71
72#[cfg(test)]
73mod transition_tests {
74 use super::{SessionState, SessionTransitionKind};
75
76 #[test]
77 fn operation_ownership_hides_intermediate_move_states_but_not_ordinary_live_work() {
78 for state in [
79 SessionState::Stopped,
80 SessionState::Running,
81 SessionState::Disconnected,
82 ] {
83 assert_eq!(
84 SessionTransitionKind::for_session(state, Some(SessionTransitionKind::Moving)),
85 Some(SessionTransitionKind::Moving)
86 );
87 assert_eq!(SessionTransitionKind::for_session(state, None), None);
88 }
89 assert_eq!(SessionState::Checkpointing.transition_kind(), None);
90 assert_eq!(
91 SessionState::Closing.transition_kind(),
92 Some(SessionTransitionKind::Suspending)
93 );
94 }
95}
96
97#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
99#[serde(tag = "state", rename_all = "snake_case")]
100pub enum MaterializedExecutionState {
101 #[default]
102 Idle,
103 Running {
104 started_at_ms: i64,
105 },
106 Closing,
107 Closed,
108}
109
110pub use crate::transcript::{TerminalOutputRecord, TranscriptBody, TranscriptItem};
111
112#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
117#[serde(rename_all = "snake_case")]
118pub enum QueuedCommandKind {
119 #[default]
120 Prompt,
121 SetConfig {
122 key: String,
123 value: String,
124 },
125}
126
127impl QueuedCommandKind {
128 pub fn is_prompt(&self) -> bool {
129 matches!(self, Self::Prompt)
130 }
131}
132
133pub fn config_command_text(key: &str, value: &str) -> String {
136 if key == "fast-mode" {
137 "/fast".to_owned()
138 } else {
139 format!("/{key} {value}")
140 }
141}
142
143#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
144#[serde(deny_unknown_fields)]
145pub struct MaterializedQueuedPrompt {
146 pub command_id: String,
147 #[serde(default, skip_serializing_if = "QueuedCommandKind::is_prompt")]
148 pub kind: QueuedCommandKind,
149 pub content: Vec<serde_json::Value>,
150 pub queued_at_ms: i64,
151 #[serde(default, skip_serializing_if = "Option::is_none")]
155 pub accepted_ordinal: Option<u64>,
156}
157
158#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
161#[serde(deny_unknown_fields)]
162pub struct MaterializedTurn {
163 pub command_id: String,
164 #[serde(default, skip_serializing_if = "Option::is_none")]
165 pub accepted_ordinal: Option<u64>,
166 pub turn_start_position: u64,
169 pub started_at_ms: i64,
170 #[serde(default, skip_serializing_if = "Option::is_none")]
173 pub steered_into: Option<String>,
174}
175
176impl MaterializedTurn {
177 pub fn belongs_to(&self, command_id: &str) -> bool {
179 self.command_id == command_id || self.steered_into.as_deref() == Some(command_id)
180 }
181}
182
183#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
185#[serde(tag = "kind", rename_all = "snake_case")]
186pub enum TurnOutcomeKind {
187 Completed { stop_reason: String },
189 Rejected { message: String },
191 Interrupted { message: String },
193}
194
195#[derive(Debug, Clone, Copy, PartialEq, Eq)]
196pub enum PromptCompletion {
197 InputRequired,
198 Finished,
199 Cancelled,
200 QuotaLimit,
201 Error,
202}
203
204pub fn classify_prompt_completion(stop_reason: &str) -> PromptCompletion {
206 let normalized = stop_reason
207 .chars()
208 .filter(|character| *character != '_' && *character != '-')
209 .flat_map(char::to_lowercase)
210 .collect::<String>();
211 match normalized.as_str() {
212 "endturn" => PromptCompletion::Finished,
213 "awaitinginput" => PromptCompletion::InputRequired,
214 "cancelled" | "canceled" => PromptCompletion::Cancelled,
215 "quotalimit" => PromptCompletion::QuotaLimit,
216 _ => PromptCompletion::Error,
217 }
218}
219
220#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
222#[serde(deny_unknown_fields)]
223pub struct MaterializedTurnOutcome {
224 #[serde(default, skip_serializing_if = "Option::is_none")]
225 pub diagnostic: Option<crate::diagnostic::TurnDiagnostic>,
226
227 #[serde(default, skip_serializing_if = "Option::is_none")]
228 pub usage: Option<crate::usage::TokenUsage>,
229 pub command_id: String,
230 #[serde(default, skip_serializing_if = "Option::is_none")]
231 pub accepted_ordinal: Option<u64>,
232 #[serde(default, skip_serializing_if = "Option::is_none")]
233 pub turn_start_position: Option<u64>,
234 pub completed_ordinal: u64,
235 pub completed_at_ms: i64,
236 pub outcome: TurnOutcomeKind,
237}
238
239impl MaterializedTurnOutcome {
240 pub fn interruption_ordinal(&self) -> Option<u64> {
242 (self.turn_start_position.is_some()
243 && matches!(self.outcome, TurnOutcomeKind::Interrupted { .. }))
244 .then_some(self.completed_ordinal)
245 }
246}
247
248#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
250#[serde(deny_unknown_fields)]
251pub struct MaterializedSession {
252 pub session_id: String,
253 pub applied_event_ordinal: u64,
254 pub applied_event_digest: String,
255 pub last_activity_at_ms: Option<i64>,
258 pub execution: MaterializedExecutionState,
259 #[serde(default, skip_serializing_if = "Option::is_none")]
260 pub session_title: Option<String>,
261 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
262 pub configuration: BTreeMap<String, serde_json::Value>,
263 #[serde(default, skip_serializing_if = "Vec::is_empty")]
264 pub transcript: Vec<Arc<TranscriptItem>>,
267 #[serde(default, skip_serializing_if = "Vec::is_empty")]
268 pub queued_prompts: Vec<MaterializedQueuedPrompt>,
269 #[serde(default, skip_serializing_if = "Vec::is_empty")]
272 pub pending_elicitations: Vec<crate::elicitation::ElicitationRequest>,
273 #[serde(default, skip_serializing_if = "Option::is_none")]
275 pub active_turn: Option<MaterializedTurn>,
276 #[serde(default, skip_serializing_if = "Option::is_none")]
279 pub last_turn_outcome: Option<MaterializedTurnOutcome>,
280}
281
282#[derive(Debug, Clone, PartialEq, Eq)]
285pub struct MaterializedSessionSummary {
286 pub session_id: String,
287 pub applied_event_ordinal: u64,
288 pub last_activity_at_ms: Option<i64>,
289 pub execution: MaterializedExecutionState,
290 pub session_title: Option<String>,
291 pub last_agent_message: Option<String>,
292 pub last_user_message: Option<String>,
293 pub last_agent_message_follows_last_user: bool,
296 pub agent_message_latest_content_ordinals: Vec<u64>,
297 pub interruption_event_ordinals: Vec<u64>,
298}
299
300impl MaterializedSession {
301 pub fn empty(session_id: impl Into<String>) -> Self {
302 Self {
303 session_id: session_id.into(),
304 applied_event_ordinal: 0,
305 applied_event_digest: RELAY_EVENT_GENESIS_DIGEST.into(),
306 last_activity_at_ms: None,
307 execution: MaterializedExecutionState::Idle,
308 session_title: None,
309 configuration: BTreeMap::new(),
310 transcript: Vec::new(),
311 queued_prompts: Vec::new(),
312 pending_elicitations: Vec::new(),
313 active_turn: None,
314 last_turn_outcome: None,
315 }
316 }
317
318 pub fn last_activity_at_ms(&self) -> Option<i64> {
319 self.last_activity_at_ms
320 }
321
322 pub fn resolved_title(&self) -> Option<String> {
328 self.session_title
329 .as_deref()
330 .and_then(normalize_session_title)
331 .or_else(|| {
332 self.transcript.iter().find_map(|item| {
333 let TranscriptBody::User { content } = &item.body else {
334 return None;
335 };
336 provisional_session_title(&crate::transcript::materialized_content_text(
337 content,
338 ))
339 })
340 })
341 .or_else(|| {
342 self.queued_prompts
343 .iter()
344 .filter(|prompt| prompt.kind.is_prompt())
345 .find_map(|prompt| {
346 provisional_session_title(&crate::transcript::materialized_content_text(
347 &prompt.content,
348 ))
349 })
350 })
351 }
352
353 pub fn unread_agent_messages_after(&self, viewed_through_event_ordinal: u64) -> u64 {
354 self.transcript
355 .iter()
356 .filter(|item| {
357 item.latest_content_event_ordinal
358 .is_some_and(|ordinal| ordinal > viewed_through_event_ordinal)
359 && item.is_nonempty_agent_message()
360 })
361 .count() as u64
362 }
363
364 pub fn unread_interruptions_after(&self, viewed_through_event_ordinal: u64) -> u64 {
365 self.interruption_event_ordinals()
366 .into_iter()
367 .filter(|ordinal| *ordinal > viewed_through_event_ordinal)
368 .count() as u64
369 }
370
371 pub fn interruption_event_ordinals(&self) -> Vec<u64> {
372 let mut ordinals = self
373 .transcript
374 .iter()
375 .filter(|item| item.is_work_interruption())
376 .map(|item| item.position)
377 .collect::<Vec<_>>();
378 if let Some(ordinal) = self
379 .last_turn_outcome
380 .as_ref()
381 .and_then(MaterializedTurnOutcome::interruption_ordinal)
382 {
383 ordinals.push(ordinal);
384 }
385 ordinals.sort_unstable();
386 ordinals.dedup();
387 ordinals
388 }
389
390 pub fn validate(&self) -> Result<()> {
391 validate_id("session", &self.session_id)?;
392 validate_relay_event_frontier(
393 self.applied_event_ordinal,
394 &self.applied_event_digest,
395 "materialized session event frontier",
396 )?;
397 if self
398 .session_title
399 .as_ref()
400 .is_some_and(|title| title.trim().is_empty())
401 {
402 bail!("materialized session has an empty title");
403 }
404 let mut item_ids = BTreeSet::new();
405 for item in &self.transcript {
406 item.validate(self.applied_event_ordinal)?;
407 if !item_ids.insert(item.stable_id.as_str()) {
408 bail!(
409 "materialized transcript contains duplicate item {:?}",
410 item.stable_id
411 );
412 }
413 }
414 let mut command_ids = BTreeSet::new();
415 for prompt in &self.queued_prompts {
416 if prompt.command_id.trim().is_empty() {
417 bail!("materialized prompt queue has an empty command id");
418 }
419 if !command_ids.insert(prompt.command_id.as_str()) {
420 bail!(
421 "materialized prompt queue contains duplicate command {:?}",
422 prompt.command_id
423 );
424 }
425 if let QueuedCommandKind::SetConfig { key, value } = &prompt.kind
426 && (key.trim().is_empty() || value.trim().is_empty())
427 {
428 bail!(
429 "materialized queued configuration change {:?} is incomplete",
430 prompt.command_id
431 );
432 }
433 }
434 Ok(())
435 }
436}
437
438#[derive(Debug, Clone, PartialEq)]
442pub struct ManagedSessionSnapshot {
443 pub materialized: MaterializedSession,
444 pub window: ProjectionWindow,
447 pub operational: RelayOperationalState,
448 pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
452 pub worker_build: Option<String>,
457 pub subagent_requests: Vec<crate::subagent::SubagentToolRequest>,
459 pub subagent_results: Vec<crate::subagent::SubagentToolResult>,
461}
462
463#[derive(Debug, Clone, PartialEq, Eq)]
475pub struct ProjectionWindow {
476 pub omitted_items: usize,
478 pub provisional_title: Option<String>,
480 pub latest_turn_start_position: Option<u64>,
484}
485
486impl ProjectionWindow {
487 pub fn trim(&mut self, session: &mut MaterializedSession, target: usize) {
490 let observed = Self::of(session);
491 if self.provisional_title.is_none() {
492 self.provisional_title = observed.provisional_title;
493 }
494 self.latest_turn_start_position = observed
495 .latest_turn_start_position
496 .or(self.latest_turn_start_position);
497 let mut boundary = session.transcript.len().saturating_sub(target.max(1));
498 for (index, item) in session.transcript.iter().enumerate() {
499 let mutable = match &item.body {
500 TranscriptBody::Agent { streaming, .. }
501 | TranscriptBody::Thought { streaming, .. } => *streaming,
502 TranscriptBody::Tool { call, .. } => matches!(
503 call.get("status").and_then(serde_json::Value::as_str),
504 Some("pending" | "in_progress")
505 ),
506 _ => false,
507 };
508 if mutable || Some(item.position) == self.latest_turn_start_position {
509 boundary = boundary.min(index);
510 }
511 }
512 let cut = session
513 .transcript
514 .iter()
515 .take(boundary + 1)
516 .rposition(|item| item.is_turn_start())
517 .unwrap_or(0);
518 if cut > 0 {
519 session.transcript.drain(..cut);
520 self.omitted_items += cut;
521 }
522 }
523
524 #[must_use]
526 pub fn of(session: &MaterializedSession) -> Self {
527 Self {
528 omitted_items: 0,
529 provisional_title: session.transcript.iter().find_map(|item| {
530 let TranscriptBody::User { content } = &item.body else {
531 return None;
532 };
533 provisional_session_title(&crate::transcript::materialized_content_text(content))
534 }),
535 latest_turn_start_position: session
536 .transcript
537 .iter()
538 .rev()
539 .find(|item| item.is_turn_start())
540 .map(|item| item.position),
541 }
542 }
543}
544
545impl ManagedSessionSnapshot {
546 #[must_use]
551 pub fn resolved_title(&self) -> Option<String> {
552 self.materialized
553 .session_title
554 .as_deref()
555 .and_then(normalize_session_title)
556 .or_else(|| self.window.provisional_title.clone())
557 .or_else(|| {
558 self.materialized
559 .queued_prompts
560 .iter()
561 .filter(|prompt| prompt.kind.is_prompt())
562 .find_map(|prompt| {
563 provisional_session_title(&crate::transcript::materialized_content_text(
564 &prompt.content,
565 ))
566 })
567 })
568 }
569
570 #[must_use]
575 pub fn latest_completed_turn_ordinal(&self) -> Option<u64> {
576 if self.materialized.execution != MaterializedExecutionState::Idle {
577 return None;
578 }
579 self.window.latest_turn_start_position
580 }
581}
582
583#[derive(Debug, Clone)]
585pub struct RecoveryObservation {
586 pub session: SessionRecord,
587 pub config: Config,
588 pub latest_completed_turn_ordinal: Option<u64>,
589 pub execution: MaterializedExecutionState,
590 pub checkpoint_safe: bool,
594}
595
596pub fn latest_completed_turn_ordinal(session: &MaterializedSession) -> Option<u64> {
601 if session.execution != MaterializedExecutionState::Idle {
602 return None;
603 }
604 session
605 .transcript
606 .iter()
607 .rev()
608 .find(|item| item.is_turn_start())
609 .map(|item| item.position)
610}
611
612pub fn validate_relay_event_digest(digest: &str, name: &str) -> Result<()> {
613 if digest.len() != 64
614 || !digest
615 .bytes()
616 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
617 {
618 bail!("{name} must be a lowercase SHA-256 digest");
619 }
620 Ok(())
621}
622
623pub fn validate_relay_event_frontier(ordinal: u64, digest: &str, name: &str) -> Result<()> {
624 validate_relay_event_digest(digest, name)?;
625 if (ordinal == 0) != (digest == RELAY_EVENT_GENESIS_DIGEST) {
626 bail!("{name} has inconsistent ordinal {ordinal} and digest {digest}");
627 }
628 Ok(())
629}
630
631fn is_false(value: &bool) -> bool {
632 !*value
633}
634
635impl SessionState {
636 pub const fn as_str(self) -> &'static str {
638 match self {
639 Self::Provisioning => "provisioning",
640 Self::Running => "running",
641 Self::Disconnected => "disconnected",
642 Self::Checkpointing => "checkpointing",
643 Self::Closing => "closing",
644 Self::Destroying => "destroying",
645 Self::Stopped => "stopped",
646 Self::Lost => "lost",
647 Self::Error => "error",
648 Self::DestroyedWithDataLoss => "destroyed-with-data-loss",
649 }
650 }
651
652 pub fn from_stored(value: &str) -> Option<Self> {
655 Some(match value {
656 "provisioning" => Self::Provisioning,
657 "running" => Self::Running,
658 "disconnected" => Self::Disconnected,
659 "checkpointing" => Self::Checkpointing,
660 "closing" => Self::Closing,
661 "destroying" => Self::Destroying,
662 "stopped" | "archived" => Self::Stopped,
663 "lost" => Self::Lost,
664 "error" => Self::Error,
665 "destroyed-with-data-loss" => Self::DestroyedWithDataLoss,
666 _ => return None,
667 })
668 }
669
670 pub const fn transition_kind(self) -> Option<SessionTransitionKind> {
673 match self {
674 Self::Provisioning => Some(SessionTransitionKind::Starting),
675 Self::Closing => Some(SessionTransitionKind::Suspending),
676 Self::Destroying => Some(SessionTransitionKind::Destroying),
677 _ => None,
678 }
679 }
680
681 pub const fn is_active(self) -> bool {
685 matches!(
686 self,
687 Self::Provisioning
688 | Self::Running
689 | Self::Disconnected
690 | Self::Checkpointing
691 | Self::Closing
692 | Self::Destroying
693 | Self::Error
694 )
695 }
696}
697
698#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
699#[serde(tag = "kind", rename_all = "kebab-case")]
700pub enum PodmanWorkspaceLocator {
701 #[default]
702 ContainerLayer,
703 Volume {
704 name: String,
705 },
706 HostPath {
707 path: PathBuf,
708 helper: Vec<String>,
709 resource: String,
710 },
711}
712
713#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
714#[serde(tag = "kind", rename_all = "kebab-case")]
715pub enum TargetLocator {
716 LocalBare {
717 worker_root: PathBuf,
718 },
719 LocalPodman {
720 container_id: String,
721 #[serde(default)]
722 workspace_storage: PodmanWorkspaceLocator,
723 #[serde(default, skip_serializing_if = "Option::is_none")]
727 borrowed_from: Option<String>,
728 },
729 LocalDocker {
730 container_id: String,
731 #[serde(default, skip_serializing_if = "Option::is_none")]
735 borrowed_from: Option<String>,
736 },
737 AppleContainer {
738 container_id: String,
739 #[serde(default, skip_serializing_if = "Option::is_none")]
743 borrowed_from: Option<String>,
744 },
745 AwsEc2 {
746 instance_id: String,
747 #[serde(default, skip_serializing_if = "Option::is_none")]
748 address: Option<String>,
749 },
750 SshBare {
751 host: String,
752 workspace: PathBuf,
753 #[serde(default, skip_serializing_if = "Option::is_none")]
754 worker_id: Option<String>,
755 },
756 SshPodman {
757 host: String,
758 container_id: String,
759 #[serde(default)]
760 workspace_storage: PodmanWorkspaceLocator,
761 #[serde(default, skip_serializing_if = "Option::is_none")]
765 borrowed_from: Option<String>,
766 },
767 SshDocker {
768 host: String,
769 container_id: String,
770 #[serde(default, skip_serializing_if = "Option::is_none")]
774 borrowed_from: Option<String>,
775 },
776}
777
778impl ManagedWorktreeTarget {
779 pub fn same_location(&self, other: &Self) -> bool {
787 match (self, other) {
788 (Self::Local, Self::Local) => true,
789 (
790 Self::Ssh {
791 destination,
792 ssh_args,
793 },
794 Self::Ssh {
795 destination: other_destination,
796 ssh_args: other_args,
797 },
798 ) => {
799 destination == other_destination
800 && ssh_location_option(ssh_args, 'p', "port")
801 == ssh_location_option(other_args, 'p', "port")
802 && ssh_location_option(ssh_args, 'l', "user")
803 == ssh_location_option(other_args, 'l', "user")
804 }
805 _ => false,
806 }
807 }
808}
809
810fn ssh_location_option(args: &[String], flag: char, option: &str) -> Option<String> {
814 let mut args = args.iter();
815 while let Some(argument) = args.next() {
816 let Some(rest) = argument.strip_prefix('-') else {
817 continue;
818 };
819 let mut chars = rest.chars();
820 let Some(name) = chars.next() else { continue };
821 if name != flag && name != 'o' {
822 continue;
823 }
824 let inline = chars.as_str();
825 let value = if inline.is_empty() {
826 args.next().cloned()
827 } else {
828 Some(inline.to_owned())
829 };
830 if name == flag {
831 return value;
832 }
833 if let Some(setting) = value {
834 let (key, found) = setting
835 .split_once(['=', ' ', '\t'])
836 .unwrap_or((setting.as_str(), ""));
837 if key.trim().eq_ignore_ascii_case(option) {
838 return Some(found.trim().to_owned());
839 }
840 }
841 }
842 None
843}
844
845#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
846#[serde(tag = "kind", rename_all = "kebab-case")]
847pub enum ManagedWorktreeTarget {
848 Local,
849 Ssh {
850 destination: String,
851 #[serde(default, skip_serializing_if = "Vec::is_empty")]
852 ssh_args: Vec<String>,
853 },
854}
855
856#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
858#[serde(deny_unknown_fields)]
859pub struct ManagedWorktreeOptions {
860 pub available: bool,
861 pub default_create: bool,
862}
863
864#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
865#[serde(deny_unknown_fields)]
866pub struct ManagedWorktree {
867 #[serde(default, skip_serializing_if = "ManagedCheckoutKind::is_worktree")]
869 pub kind: ManagedCheckoutKind,
870 pub source_project_directory: PathBuf,
871 pub source_repository: PathBuf,
872 pub worktree_root: PathBuf,
873 pub branch: String,
874 pub target: ManagedWorktreeTarget,
875 #[serde(default, skip_serializing_if = "Option::is_none")]
879 pub base_commit: Option<String>,
880}
881
882#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
883#[serde(rename_all = "snake_case")]
884pub enum ManagedCheckoutKind {
885 #[default]
886 Worktree,
887 Clone,
888}
889
890impl ManagedCheckoutKind {
891 fn is_worktree(&self) -> bool {
892 matches!(self, Self::Worktree)
893 }
894}
895
896impl ManagedWorktree {
897 fn validate(&self, session_id: &str, project_directory: Option<&Path>) -> Result<()> {
898 for (label, path) in [
899 ("source project directory", &self.source_project_directory),
900 ("source repository", &self.source_repository),
901 ("worktree root", &self.worktree_root),
902 ] {
903 if !path.is_absolute() || path.components().any(|part| part == Component::ParentDir) {
904 bail!("managed worktree {label} must be an absolute safe path");
905 }
906 }
907 if !self
908 .source_project_directory
909 .starts_with(&self.source_repository)
910 {
911 bail!("managed worktree source directory is outside its repository");
912 }
913 let expected_root = self
914 .source_repository
915 .join(".mj")
916 .join(match self.kind {
917 ManagedCheckoutKind::Worktree => "worktrees",
918 ManagedCheckoutKind::Clone => "clones",
919 })
920 .join(session_id);
921 if self.worktree_root != expected_root {
922 bail!("managed worktree root does not match the session-owned path");
923 }
924 if self.kind == ManagedCheckoutKind::Worktree && self.branch != format!("mj/{session_id}") {
925 bail!("managed worktree branch does not match the session id");
926 }
927 if self.kind == ManagedCheckoutKind::Clone && self.branch.trim().is_empty() {
928 bail!("managed clone has no starting branch");
929 }
930 let relative = self
931 .source_project_directory
932 .strip_prefix(&self.source_repository)
933 .expect("source relationship checked above");
934 if project_directory != Some(self.worktree_root.join(relative).as_path()) {
935 bail!("session project directory does not match its managed worktree");
936 }
937 match &self.target {
938 ManagedWorktreeTarget::Local => {}
939 ManagedWorktreeTarget::Ssh { destination, .. } if destination.trim().is_empty() => {
940 bail!("managed SSH worktree has an empty destination")
941 }
942 ManagedWorktreeTarget::Ssh { .. } => {}
943 }
944 Ok(())
945 }
946}
947
948#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
949#[serde(tag = "kind", rename_all = "kebab-case")]
950pub enum SessionResourceAllocation {
951 Container {
952 cpus: u64,
953 memory_bytes: u64,
954 },
955 AwsEc2 {
956 instance_type: String,
957 vcpus: u64,
958 memory_bytes: u64,
959 },
960}
961
962impl SessionResourceAllocation {
963 pub fn validate(&self) -> Result<()> {
964 match self {
965 Self::Container { cpus, memory_bytes } if *cpus == 0 || *memory_bytes == 0 => {
966 bail!("container resource allocation must have non-zero CPU and memory")
967 }
968 Self::AwsEc2 {
969 instance_type,
970 vcpus,
971 memory_bytes,
972 } if instance_type.trim().is_empty() || *vcpus == 0 || *memory_bytes == 0 => {
973 bail!("EC2 resource allocation must have an instance type, CPU, and memory")
974 }
975 _ => Ok(()),
976 }
977 }
978}
979
980pub fn allocation_cpus(allocation: &SessionResourceAllocation) -> u64 {
982 match allocation {
983 SessionResourceAllocation::Container { cpus, .. } => *cpus,
984 SessionResourceAllocation::AwsEc2 { vcpus, .. } => *vcpus,
985 }
986}
987
988pub fn allocation_memory(allocation: &SessionResourceAllocation) -> u64 {
990 match allocation {
991 SessionResourceAllocation::Container { memory_bytes, .. }
992 | SessionResourceAllocation::AwsEc2 { memory_bytes, .. } => *memory_bytes,
993 }
994}
995
996impl TargetLocator {
997 fn validate(&self, session_id: &str) -> Result<()> {
998 match self {
999 Self::LocalBare { worker_root } => {
1000 if !worker_root.is_absolute()
1001 || worker_root
1002 .components()
1003 .any(|part| part == Component::ParentDir)
1004 || !worker_root.ends_with(session_id)
1005 {
1006 bail!(
1007 "local bare worker root must be an absolute safe path ending in the session id"
1008 );
1009 }
1010 }
1011 Self::LocalPodman { container_id, .. }
1012 | Self::LocalDocker { container_id, .. }
1013 | Self::AppleContainer { container_id, .. }
1014 | Self::SshPodman { container_id, .. }
1015 | Self::SshDocker { container_id, .. }
1016 if container_id.trim().is_empty() =>
1017 {
1018 bail!("target locator has an empty container id")
1019 }
1020 Self::AwsEc2 { instance_id, .. } if instance_id.trim().is_empty() => {
1021 bail!("target locator has an empty AWS instance id")
1022 }
1023 Self::SshBare {
1024 host, workspace, ..
1025 } => {
1026 if host.trim().is_empty() {
1027 bail!("bare SSH target locator has an empty host");
1028 }
1029 if workspace.as_os_str().is_empty()
1030 || workspace
1031 .components()
1032 .any(|part| part == Component::ParentDir)
1033 || !workspace.ends_with(session_id)
1034 {
1035 bail!("bare SSH target locator must be a safe path ending in the session id");
1036 }
1037 }
1038 Self::SshPodman { host, .. } if host.trim().is_empty() => {
1039 bail!("SSH Podman target locator has an empty host")
1040 }
1041 Self::SshDocker { host, .. } if host.trim().is_empty() => {
1042 bail!("SSH Docker target locator has an empty host")
1043 }
1044 _ => {}
1045 }
1046 Ok(())
1047 }
1048}
1049
1050#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1051#[serde(deny_unknown_fields)]
1052pub struct CheckpointMetadata {
1053 pub archive_path: PathBuf,
1054 pub sha256: String,
1056 pub created_at: String,
1057 pub event_frontier: u64,
1058}
1059
1060#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1062#[serde(rename_all = "snake_case")]
1063pub enum PublicationState {
1064 Published,
1065 Unpublished,
1066 Unknown,
1067}
1068
1069#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1071pub struct PublicationAssessment {
1072 pub checkpoint_sha256: String,
1073 pub state: PublicationState,
1074 pub dirty: bool,
1075 pub stashed: bool,
1076 pub saved_commits: Vec<String>,
1077 pub destinations: Vec<String>,
1078 pub checked_at: String,
1079 pub reason: Option<String>,
1080}
1081
1082impl CheckpointMetadata {
1083 fn validate(&self) -> Result<()> {
1084 if self.archive_path.as_os_str().is_empty() {
1085 bail!("checkpoint archive path is empty");
1086 }
1087 if self.sha256.len() != 64
1088 || !self
1089 .sha256
1090 .bytes()
1091 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1092 {
1093 bail!("checkpoint SHA-256 must be 64 lowercase hexadecimal characters");
1094 }
1095 if self.created_at.trim().is_empty() {
1096 bail!("checkpoint timestamp is empty");
1097 }
1098 Ok(())
1099 }
1100}
1101
1102#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1106#[serde(deny_unknown_fields)]
1107pub struct SessionBuildCache {
1108 pub host: String,
1111 pub directory: PathBuf,
1112 #[serde(default, skip_serializing_if = "Option::is_none")]
1115 pub max_size: Option<String>,
1116 #[serde(default, skip_serializing_if = "Option::is_none")]
1119 pub target_root: Option<PathBuf>,
1120}
1121
1122#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1126pub struct ArchiveSpacePreview {
1127 pub sessions: usize,
1129 pub bytes: u64,
1131 pub reclaimable_sessions: usize,
1134 pub reclaimable_bytes: u64,
1135}
1136
1137#[derive(Debug, Clone, PartialEq, Eq)]
1141pub struct BuildCachePreview {
1142 pub native_mbx: Option<String>,
1144 pub directory: Option<PathBuf>,
1146 pub max_size: Option<BuildCacheLimit>,
1148 pub stats: Option<BuildCacheStats>,
1150 pub off_reason: Option<BuildCacheOff>,
1153}
1154
1155#[derive(Debug, Clone, PartialEq, Eq)]
1163pub struct BuildCacheStats {
1164 pub builds: u64,
1166 pub cached_compilations: u64,
1168 pub avoided_compiler_ns: u64,
1170 pub reflinked_bytes: u64,
1172}
1173
1174#[derive(Debug, Clone, PartialEq, Eq)]
1176pub enum BuildCacheOff {
1177 TurnedOff,
1179 Unavailable(String),
1182}
1183
1184impl std::fmt::Display for BuildCacheOff {
1185 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1186 match self {
1187 Self::TurnedOff => formatter.write_str("turned off for this machine"),
1188 Self::Unavailable(reason) => formatter.write_str(reason),
1189 }
1190 }
1191}
1192
1193#[derive(Debug, Clone, PartialEq, Eq)]
1195pub enum BuildCacheLimit {
1196 Size(String),
1198 HostConfiguration(Option<String>),
1201}
1202
1203#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1204#[serde(deny_unknown_fields)]
1205pub struct SessionRecord {
1206 pub id: String,
1207 #[serde(default = "default_session_workspace_id")]
1212 pub workspace_id: String,
1213 pub title: String,
1214 pub harness_kind: HarnessKind,
1215 pub last_profile: String,
1216 pub bundle_id: String,
1217 #[serde(default, skip_serializing_if = "Option::is_none")]
1219 pub project_directory: Option<PathBuf>,
1220 #[serde(default, skip_serializing_if = "Option::is_none")]
1222 pub managed_worktree: Option<ManagedWorktree>,
1223 #[serde(default, skip_serializing_if = "Option::is_none")]
1225 pub create_managed_worktree: Option<bool>,
1226 #[serde(default, skip_serializing_if = "Option::is_none")]
1231 pub launch_base: Option<String>,
1232 #[serde(default, skip_serializing_if = "Option::is_none")]
1234 pub launch_branch: Option<String>,
1235 #[serde(default, skip_serializing_if = "Option::is_none")]
1237 pub publication: Option<PublicationAssessment>,
1238 #[serde(default, skip_serializing_if = "Option::is_none")]
1241 pub mjolnir_subagents: Option<bool>,
1242 pub target_template_id: String,
1243 #[serde(default, skip_serializing_if = "Option::is_none")]
1244 pub resource_allocation: Option<SessionResourceAllocation>,
1245 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1246 pub additional_mounts: Vec<AdditionalMount>,
1247 #[serde(default, skip_serializing_if = "Option::is_none")]
1250 pub container_cpus: Option<String>,
1251 #[serde(default, skip_serializing_if = "Option::is_none")]
1254 pub container_memory: Option<String>,
1255 #[serde(default, skip_serializing_if = "Option::is_none")]
1261 pub container_workspace: Option<PathBuf>,
1262 #[serde(default, skip_serializing_if = "Option::is_none")]
1266 pub build_cache: Option<SessionBuildCache>,
1267 pub state: SessionState,
1268 #[serde(default, skip_serializing_if = "is_false")]
1271 pub archived: bool,
1272 #[serde(default, skip_serializing_if = "Option::is_none")]
1273 pub target: Option<TargetLocator>,
1274 #[serde(default, skip_serializing_if = "Option::is_none")]
1276 pub target_runtime: Option<TargetRuntimeSettings>,
1277 #[serde(default, skip_serializing_if = "Option::is_none")]
1278 pub native_session_id: Option<String>,
1279 #[serde(default, skip_serializing_if = "Option::is_none")]
1280 pub acp_session_title: Option<String>,
1281 #[serde(default, skip_serializing_if = "Option::is_none")]
1282 pub session_title_override: Option<String>,
1283 pub created_at: String,
1284 pub updated_at: String,
1285 #[serde(default, alias = "detached_after_event_ordinal")]
1286 pub viewed_through_event_ordinal: u64,
1287 #[serde(default, skip_serializing_if = "String::is_empty")]
1290 pub draft_input: String,
1291 #[serde(default, skip_serializing_if = "Option::is_none")]
1299 pub last_error: Option<String>,
1300 #[serde(default, skip_serializing_if = "Option::is_none")]
1301 pub last_checkpoint_error: Option<String>,
1302 #[serde(default, skip_serializing_if = "Option::is_none")]
1303 pub checkpoint: Option<CheckpointMetadata>,
1304}
1305
1306#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1307#[serde(deny_unknown_fields)]
1308pub struct HostContainerSize {
1309 pub cpus: u64,
1310 pub memory_bytes: u64,
1311}
1312
1313fn default_session_workspace_id() -> String {
1314 crate::workspace::DEFAULT_WORKSPACE_ID.to_owned()
1315}
1316
1317pub const DESTRUCTION_FAILURE_PREFIX: &str = "the destruction did not finish";
1326
1327pub const CLOSE_FAILURE_PREFIX: &str = "the suspension did not finish";
1328
1329pub fn is_public_lifecycle_error(error: &str) -> bool {
1331 error.starts_with(CLOSE_FAILURE_PREFIX)
1332 || error.starts_with(DESTRUCTION_FAILURE_PREFIX)
1333 || error.starts_with("the close did not finish")
1334}
1335
1336#[must_use]
1349pub fn target_label(config: &Config, target_id: &str, project: Option<&Path>) -> String {
1350 if !matches!(
1351 config.targets.get(target_id),
1352 Some(TargetTemplate::LocalBare | TargetTemplate::SshBare { .. })
1353 ) {
1354 return target_id.to_owned();
1355 }
1356 project.and_then(Path::file_name).map_or_else(
1357 || target_id.to_owned(),
1358 |directory| format!("{target_id}/{}", directory.to_string_lossy()),
1359 )
1360}
1361
1362impl SessionRecord {
1363 pub fn publication_state(&self) -> Option<PublicationState> {
1366 let independent_clone = self
1367 .managed_worktree
1368 .as_ref()
1369 .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Clone)
1370 || (self.managed_worktree.is_none() && self.project_directory.is_none());
1371 if !independent_clone {
1372 return None;
1373 }
1374 if self.state.is_active() {
1375 return Some(PublicationState::Unknown);
1376 }
1377 Some(
1378 self.checkpoint
1379 .as_ref()
1380 .zip(self.publication.as_ref())
1381 .filter(|(checkpoint, assessment)| {
1382 assessment.checkpoint_sha256 == checkpoint.sha256
1383 })
1384 .map_or(PublicationState::Unknown, |(_, assessment)| {
1385 if assessment.dirty || assessment.stashed {
1386 PublicationState::Unpublished
1387 } else {
1388 assessment.state
1389 }
1390 }),
1391 )
1392 }
1393
1394 pub fn target_runtime_settings<'a>(
1399 &'a self,
1400 config: &Config,
1401 ) -> Result<std::borrow::Cow<'a, TargetRuntimeSettings>> {
1402 if let Some(runtime) = &self.target_runtime {
1403 if let Some(refreshed) =
1406 config
1407 .targets
1408 .get(&self.target_template_id)
1409 .and_then(|template| {
1410 runtime.with_current_ssh_options(&TargetRuntimeSettings::from(template))
1411 })
1412 {
1413 return Ok(std::borrow::Cow::Owned(refreshed));
1414 }
1415 return Ok(std::borrow::Cow::Borrowed(runtime));
1416 }
1417 let template = config.targets.get(&self.target_template_id).ok_or_else(|| {
1418 crate::refusal::Refusal::precondition(format!(
1419 "Session {:?} has no recorded target access settings. Restore target {:?} in config.toml once, then retry.",
1420 self.id, self.target_template_id))
1421 })?;
1422 let runtime = TargetRuntimeSettings::from(template);
1423 if let Some(locator) = &self.target {
1424 crate::targets::TargetLocator::try_from(crate::targets::RecordedTarget {
1425 locator, runtime: Some(&runtime), session_id: &self.id,
1426 }).map_err(|error| crate::refusal::Refusal::precondition(format!(
1427 "Session {:?} cannot recover target {:?}: {error}. Restore its original target settings, then retry.", self.id, self.target_template_id)))?;
1428 }
1429 Ok(std::borrow::Cow::Owned(runtime))
1430 }
1431
1432 #[must_use]
1436 pub fn public_error(&self) -> Option<&str> {
1437 self.last_error
1438 .as_deref()
1439 .filter(|error| is_public_lifecycle_error(error))
1440 }
1441
1442 pub fn configuration_issue(&self, config: &Config) -> Option<String> {
1445 if !self.state.is_active() {
1446 return None;
1447 }
1448 let mut issues = Vec::new();
1449 match config.profiles.get(&self.last_profile) {
1450 None => issues.push(format!("missing profile {:?}", self.last_profile)),
1451 Some(profile) if profile.kind != self.harness_kind => issues.push(format!(
1452 "expects {:?}, but profile {:?} is {:?}",
1453 self.harness_kind, self.last_profile, profile.kind
1454 )),
1455 Some(_) => {}
1456 }
1457 if self.project_directory.is_none() && !config.bundles.contains_key(&self.bundle_id) {
1458 issues.push(format!("missing bundle {:?}", self.bundle_id));
1459 }
1460 if self.target_runtime.is_none() && !config.targets.contains_key(&self.target_template_id) {
1461 issues.push(format!(
1462 "missing target template {:?}",
1463 self.target_template_id
1464 ));
1465 }
1466 (!issues.is_empty()).then(|| format!(
1467 "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.",
1468 self.id, issues.join("; ")
1469 ))
1470 }
1471
1472 pub fn validate_configuration(&self, config: &Config) -> Result<()> {
1473 if let Some(issue) = self.configuration_issue(config) {
1474 return Err(crate::refusal::Refusal::precondition(issue).into());
1475 }
1476 Ok(())
1477 }
1478
1479 pub fn display_title(&self) -> &str {
1481 self.session_title_override
1482 .as_deref()
1483 .or(self.acp_session_title.as_deref())
1484 .unwrap_or(&self.id)
1485 }
1486
1487 pub fn listed_title(&self) -> &str {
1492 let named = self.session_title_override.is_some() || self.acp_session_title.is_some();
1493 if !named && !self.title.trim().is_empty() {
1494 return &self.title;
1495 }
1496 self.display_title()
1497 }
1498
1499 pub fn project_name(&self, config: &Config) -> String {
1504 if let Some(worktree) = &self.managed_worktree {
1505 return path_leaf(&worktree.source_repository);
1506 }
1507 if let Some(project_directory) = &self.project_directory {
1508 return path_leaf(project_directory);
1509 }
1510 self.bundle_source_name(config)
1511 }
1512
1513 pub fn project_target(&self, config: &Config, target_id: &str) -> String {
1517 let project = self
1518 .managed_worktree
1519 .as_ref()
1520 .map(|worktree| &worktree.source_project_directory)
1521 .or(self.project_directory.as_ref());
1522 target_label(config, target_id, project.map(PathBuf::as_path))
1523 }
1524
1525 pub fn project_source(&self, config: &Config) -> ProjectSourceIdentity {
1530 if let Some(worktree) = &self.managed_worktree {
1531 return ProjectSourceIdentity::path(&worktree.source_repository, None);
1532 }
1533 if let Some(project_directory) = &self.project_directory {
1534 let remote = match &self.target {
1535 Some(TargetLocator::SshBare { host, .. }) => Some(host.as_str()),
1536 _ => None,
1537 };
1538 return ProjectSourceIdentity::path(project_directory, remote);
1539 }
1540 self.bundle_source_identity(config)
1541 .unwrap_or_else(|| ProjectSourceIdentity {
1542 key: format!("bundle:{}", self.bundle_id),
1543 short: path_leaf(Path::new(&self.bundle_id)),
1544 full: self.bundle_id.clone(),
1545 })
1546 }
1547
1548 fn bundle_source_name(&self, config: &Config) -> String {
1551 self.bundle_source_identity(config)
1552 .map(|source| source.short)
1553 .unwrap_or_else(|| path_leaf(Path::new(&self.bundle_id)))
1554 }
1555
1556 fn bundle_source_identity(&self, config: &Config) -> Option<ProjectSourceIdentity> {
1559 let bundle = config.bundles.get(&self.bundle_id)?;
1560 let sources = bundle
1561 .repositories
1562 .iter()
1563 .map(repository_source_identity)
1564 .collect::<Option<Vec<_>>>()?;
1565 ProjectSourceIdentity::bundle(sources)
1566 }
1567
1568 pub fn compare_by_creation(&self, other: &Self) -> std::cmp::Ordering {
1572 self.creation_order_key().cmp(&other.creation_order_key())
1573 }
1574
1575 pub fn creation_order_key(&self) -> (bool, Option<i64>, &str) {
1577 let timestamp = created_at_seconds(&self.created_at);
1578 (timestamp.is_none(), timestamp, &self.id)
1579 }
1580
1581 fn validate(&self, map_id: &str) -> Result<()> {
1582 validate_id("session", &self.id)?;
1583 if self.id != map_id {
1584 bail!(
1585 "session map key {map_id:?} does not match record id {:?}",
1586 self.id
1587 );
1588 }
1589 validate_id("workspace", &self.workspace_id)?;
1590 validate_id("profile", &self.last_profile)?;
1591 validate_id("bundle", &self.bundle_id)?;
1592 if let Some(project_directory) = &self.project_directory
1593 && (!project_directory.is_absolute()
1594 || project_directory
1595 .components()
1596 .any(|part| part == Component::ParentDir))
1597 {
1598 bail!("session {:?} has an unsafe project directory", self.id);
1599 }
1600 if let Some(managed_worktree) = &self.managed_worktree {
1601 managed_worktree.validate(&self.id, self.project_directory.as_deref())?;
1602 }
1603 validate_id("target template", &self.target_template_id)?;
1604 if let Some(allocation) = &self.resource_allocation {
1605 allocation.validate()?;
1606 }
1607 validate_additional_mounts(&self.additional_mounts)?;
1608 if self.title.trim().is_empty() {
1609 bail!("session {:?} has an empty title", self.id);
1610 }
1611 if self
1612 .acp_session_title
1613 .as_ref()
1614 .is_some_and(|title| title.trim().is_empty())
1615 || self
1616 .session_title_override
1617 .as_ref()
1618 .is_some_and(|title| title.trim().is_empty())
1619 {
1620 bail!("session {:?} has an empty display title", self.id);
1621 }
1622 if self.created_at.trim().is_empty() || self.updated_at.trim().is_empty() {
1623 bail!("session {:?} has an empty timestamp", self.id);
1624 }
1625 if let Some(target) = &self.target {
1626 target.validate(&self.id)?;
1627 }
1628 if let Some(checkpoint) = &self.checkpoint {
1629 checkpoint.validate()?;
1630 }
1631 Ok(())
1632 }
1633}
1634
1635fn repository_source_identity(repository: &ProjectRepository) -> Option<ProjectSourceIdentity> {
1636 repository
1637 .github
1638 .as_deref()
1639 .and_then(ProjectSourceIdentity::git_remote)
1640 .or_else(|| {
1641 repository
1642 .local
1643 .as_deref()
1644 .map(|path| ProjectSourceIdentity::path(path, None))
1645 })
1646}
1647
1648#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1649pub struct ProjectSourceIdentity {
1650 pub key: String,
1651 pub short: String,
1652 pub full: String,
1653}
1654
1655impl ProjectSourceIdentity {
1656 pub fn bundle(mut sources: Vec<Self>) -> Option<Self> {
1658 if sources.is_empty() {
1659 return None;
1660 }
1661 sources.sort_by(|left, right| {
1662 left.key
1663 .cmp(&right.key)
1664 .then_with(|| left.full.cmp(&right.full))
1665 .then_with(|| left.short.cmp(&right.short))
1666 });
1667 sources.dedup_by(|left, right| left.key == right.key);
1668 if sources.len() == 1 {
1669 return sources.pop();
1670 }
1671 let keys = sources
1672 .iter()
1673 .map(|source| source.key.clone())
1674 .collect::<Vec<_>>();
1675 let key = serde_json::to_string(&keys).ok()?;
1676 Some(Self {
1677 key: format!("bundle:{key}"),
1678 short: sources
1679 .iter()
1680 .map(|source| source.short.as_str())
1681 .collect::<Vec<_>>()
1682 .join(" + "),
1683 full: sources
1684 .iter()
1685 .map(|source| source.full.as_str())
1686 .collect::<Vec<_>>()
1687 .join(" + "),
1688 })
1689 }
1690
1691 pub fn git_remote(source: &str) -> Option<Self> {
1694 if let Some(normalized) = normalize_github_source(source) {
1695 let short = normalized
1696 .rsplit_once('/')
1697 .map_or(normalized.as_str(), |(_, repository)| repository)
1698 .to_owned();
1699 return Some(Self {
1700 key: format!("github:{}", normalized.to_lowercase()),
1701 short,
1702 full: normalized,
1703 });
1704 }
1705 let normalized = source.trim().trim_end_matches('/').trim_end_matches(".git");
1706 if normalized.is_empty() {
1707 return None;
1708 }
1709 let short = normalized
1710 .rsplit(['/', ':'])
1711 .find(|part| !part.is_empty())
1712 .unwrap_or(normalized)
1713 .to_owned();
1714 Some(Self {
1715 key: format!("git:{}", normalized.to_lowercase()),
1716 short,
1717 full: normalized.to_owned(),
1718 })
1719 }
1720
1721 pub fn path(path: &Path, remote: Option<&str>) -> Self {
1723 let normalized = path.components().collect::<PathBuf>();
1724 let path_text = normalized.to_string_lossy().into_owned();
1725 let full = remote.map_or_else(|| path_text.clone(), |host| format!("{host}:{path_text}"));
1726 let key = remote.map_or_else(
1727 || format!("path:{path_text}"),
1728 |host| format!("path:{}:{path_text}", host.to_lowercase()),
1729 );
1730 Self {
1731 key,
1732 short: path_leaf(path),
1733 full,
1734 }
1735 }
1736}
1737
1738fn normalize_github_source(source: &str) -> Option<String> {
1739 let source = source.trim();
1740 let path = source
1741 .strip_prefix("https://github.com/")
1742 .or_else(|| source.strip_prefix("http://github.com/"))
1743 .or_else(|| source.strip_prefix("git@github.com:"))
1744 .or_else(|| source.strip_prefix("ssh://git@github.com/"))
1745 .or_else(|| {
1746 (!source.contains("://") && !source.contains('@') && !source.contains(':'))
1747 .then_some(source)
1748 })?
1749 .trim_end_matches(".git");
1750 let mut parts = path.split('/');
1751 let owner = parts.next()?;
1752 let repository = parts.next()?;
1753 (!owner.is_empty() && !repository.is_empty() && parts.next().is_none())
1754 .then(|| format!("{owner}/{repository}"))
1755}
1756
1757fn path_leaf(path: &Path) -> String {
1759 path.file_name()
1760 .unwrap_or(path.as_os_str())
1761 .to_string_lossy()
1762 .into_owned()
1763}
1764
1765fn created_at_seconds(timestamp: &str) -> Option<i64> {
1766 chrono::DateTime::parse_from_rfc3339(timestamp)
1767 .ok()
1768 .map(|timestamp| timestamp.timestamp())
1769}
1770
1771#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1772#[serde(deny_unknown_fields)]
1773pub struct State {
1774 pub version: u32,
1775 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1776 pub sessions: BTreeMap<String, SessionRecord>,
1777 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1780 pub subagents: BTreeMap<String, SubagentRecord>,
1781 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1783 pub mount_history: BTreeMap<String, Vec<PathBuf>>,
1784 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1786 pub container_sizes: BTreeMap<String, HostContainerSize>,
1787}
1788
1789impl Default for State {
1790 fn default() -> Self {
1791 Self {
1792 version: STATE_VERSION,
1793 sessions: BTreeMap::new(),
1794 subagents: BTreeMap::new(),
1795 mount_history: BTreeMap::new(),
1796 container_sizes: BTreeMap::new(),
1797 }
1798 }
1799}
1800
1801impl State {
1802 #[must_use]
1807 pub fn session_notice_name(&self, session_id: &str) -> String {
1808 match self.sessions.get(session_id) {
1809 Some(session) if session.listed_title() != session.id => {
1810 session.listed_title().to_owned()
1811 }
1812 _ => short_id(session_id).to_owned(),
1813 }
1814 }
1815
1816 #[must_use]
1825 pub fn project_identity_session<'a>(&'a self, session: &'a SessionRecord) -> &'a SessionRecord {
1826 self.subagents
1827 .get(&session.id)
1828 .and_then(|record| self.sessions.get(&record.parent_session_id))
1829 .unwrap_or(session)
1830 }
1831
1832 #[must_use]
1839 pub fn is_subagent_session(&self, id: &str) -> bool {
1840 self.subagents.contains_key(id) || crate::native_agent::is_view_id(id)
1841 }
1842
1843 pub fn validate(&self) -> Result<()> {
1844 if self.version != STATE_VERSION {
1845 bail!(
1846 "unsupported Mjolnir state version {}; expected {STATE_VERSION}",
1847 self.version
1848 );
1849 }
1850 for (id, session) in &self.sessions {
1851 session.validate(id)?;
1852 }
1853 for (child_id, subagent) in &self.subagents {
1854 if child_id != &subagent.child_session_id {
1855 bail!("sub-agent key {child_id:?} does not match its child session id");
1856 }
1857 if child_id == &subagent.parent_session_id {
1858 bail!("sub-agent {child_id:?} cannot be its own parent");
1859 }
1860 if !self.sessions.contains_key(child_id) {
1861 bail!("sub-agent {child_id:?} has no child session");
1862 }
1863 if !self.sessions.contains_key(&subagent.parent_session_id) {
1864 bail!(
1865 "sub-agent {child_id:?} has unknown parent {:?}",
1866 subagent.parent_session_id
1867 );
1868 }
1869 if self.subagents.contains_key(&subagent.parent_session_id) {
1870 bail!("sub-agent {child_id:?} cannot belong to another sub-agent");
1871 }
1872 if subagent.task_name.trim().is_empty()
1873 || subagent.profile_id.trim().is_empty()
1874 || subagent.request_key.trim().is_empty()
1875 {
1876 bail!("sub-agent {child_id:?} has incomplete relationship metadata");
1877 }
1878 }
1879 for (host, sources) in &self.mount_history {
1880 if host.trim().is_empty() {
1881 bail!("mount history contains an empty host key");
1882 }
1883 if sources.iter().any(|source| !source.is_absolute()) {
1884 bail!("mount history for {host:?} contains a non-absolute source path");
1885 }
1886 }
1887 for (host, size) in &self.container_sizes {
1888 if host.trim().is_empty() {
1889 bail!("container size history contains an empty host key");
1890 }
1891 if size.cpus == 0 || size.memory_bytes == 0 {
1892 bail!("container size history for {host:?} contains a zero value");
1893 }
1894 if size.cpus > i64::MAX as u64 || size.memory_bytes > i64::MAX as u64 {
1895 bail!("container size history for {host:?} exceeds SQLite integer range");
1896 }
1897 }
1898 Ok(())
1899 }
1900
1901 pub fn remember_mount_sources(&mut self, host: &str, mounts: &[AdditionalMount]) {
1902 if mounts.is_empty() {
1903 return;
1904 }
1905 let sources = self.mount_history.entry(host.to_owned()).or_default();
1906 for mount in mounts.iter().rev() {
1907 sources.retain(|source| source != &mount.source);
1908 sources.insert(0, mount.source.clone());
1909 }
1910 sources.truncate(20);
1911 }
1912
1913 pub fn remember_container_size(&mut self, host: &str, size: HostContainerSize) {
1914 self.container_sizes.insert(host.to_owned(), size);
1915 }
1916
1917 pub fn project_directories(&self, host: &str) -> &[PathBuf] {
1918 self.mount_history
1919 .get(&project_history_key(host))
1920 .map(Vec::as_slice)
1921 .unwrap_or_default()
1922 }
1923
1924 pub fn remember_project_directory(&mut self, host: &str, directory: &Path) {
1925 let key = project_history_key(host);
1926 let directories = self.mount_history.entry(key).or_default();
1927 directories.retain(|existing| existing != directory);
1928 directories.insert(0, directory.to_path_buf());
1929 directories.truncate(20);
1930 }
1931
1932 pub fn destroy_stopped_session(&mut self, session_id: &str) -> Result<SessionRecord> {
1933 let session = self
1934 .sessions
1935 .get(session_id)
1936 .with_context(|| format!("unknown session {session_id}"))?;
1937 if session.state.is_active() {
1938 bail!("refusing to destroy active session {session_id}");
1939 }
1940 Ok(self
1941 .sessions
1942 .remove(session_id)
1943 .expect("session checked above"))
1944 }
1945
1946 pub fn destroy_session_force(&mut self, session_id: &str) -> Result<SessionRecord> {
1952 self.sessions
1953 .get(session_id)
1954 .with_context(|| format!("unknown session {session_id}"))?;
1955 Ok(self
1956 .sessions
1957 .remove(session_id)
1958 .expect("session checked above"))
1959 }
1960
1961 pub fn validate_setup_update(&self, before: &Config, after: &Config) -> Result<()> {
1964 for session in self
1965 .sessions
1966 .values()
1967 .filter(|session| session.state.is_active())
1968 {
1969 let protected = if let Some(profile) = before.profiles.get(&session.last_profile) {
1970 let mut comparable = profile.clone();
1971 if let Some(updated) = after.profiles.get(&session.last_profile) {
1972 comparable.enabled = updated.enabled;
1973 }
1974 profile.kind == session.harness_kind
1976 && after.profiles.get(&session.last_profile) != Some(&comparable)
1977 } else {
1978 false
1979 };
1980 let bundle_changed = session.project_directory.is_none()
1981 && before
1982 .bundles
1983 .get(&session.bundle_id)
1984 .is_some_and(|bundle| after.bundles.get(&session.bundle_id) != Some(bundle));
1985 let target_changed =
1989 before
1990 .targets
1991 .get(&session.target_template_id)
1992 .is_some_and(|target| {
1993 after
1994 .targets
1995 .get(&session.target_template_id)
1996 .map(TargetTemplate::without_launch_only_settings)
1997 != Some(target.without_launch_only_settings())
1998 });
1999 if protected || bundle_changed || target_changed {
2000 let mut used = Vec::new();
2004 if protected {
2005 used.push(format!("agent profile {:?}", session.last_profile));
2006 }
2007 if bundle_changed {
2008 used.push(format!("project {:?}", session.project_name(before)));
2009 }
2010 if target_changed {
2011 used.push(format!("runtime {:?}", session.target_template_id));
2012 }
2013 let used = match used.as_slice() {
2014 [only] => only.clone(),
2015 [rest @ .., last] => format!("{} and {last}", rest.join(", ")),
2016 [] => unreachable!("something changed"),
2017 };
2018 let title = session.display_title();
2019 let named = if title == session.id {
2020 format!(
2021 "a running session in project {:?}",
2022 session.project_name(before)
2023 )
2024 } else {
2025 format!("the running session {title:?}")
2026 };
2027 bail!(
2028 "Setup would change the {used} that {named} uses. Save the new settings under a new name, or stop the session first."
2029 );
2030 }
2031 }
2032 Ok(())
2033 }
2034
2035 pub fn validate_against_config(&self, config: &Config) -> Result<()> {
2037 self.validate()?;
2038 config.validate()?;
2039 for session in self.sessions.values() {
2040 session.validate_configuration(config)?;
2041 }
2042 Ok(())
2043 }
2044}
2045
2046fn project_history_key(host: &str) -> String {
2047 format!("project:{host}")
2048}
2049
2050pub fn new_session_id() -> Result<String> {
2052 let mut random = [0u8; 16];
2053 getrandom::fill(&mut random)
2054 .map_err(|error| anyhow::anyhow!("generate Mjolnir session id: {error}"))?;
2055 Ok(crate::hex::lower_hex(random))
2056}
2057
2058pub fn harness_session_title(events: &[SequencedEvent]) -> Option<String> {
2060 events.iter().rev().find_map(|event| {
2061 let WorkerEvent::Adapter { payload, .. } = &event.event else {
2062 return None;
2063 };
2064 let crate::acp::RuntimeEvent::SessionUpdate { update } =
2065 serde_json::from_value(payload.clone()).ok()?
2066 else {
2067 return None;
2068 };
2069 let kind = update
2070 .get("sessionUpdate")
2071 .and_then(serde_json::Value::as_str)?;
2072 let title = match kind {
2073 "session_info_update" | "session_title" => {
2074 update.get("title").and_then(serde_json::Value::as_str)
2075 }
2076 _ => None,
2077 }?;
2078 normalize_session_title(title)
2079 })
2080}
2081
2082pub fn default_session_title(
2087 project_directory: Option<&Path>,
2088 bundle_id: &str,
2089 profile_id: &str,
2090) -> String {
2091 let project = project_directory.and_then(Path::file_name).map_or_else(
2092 || bundle_id.to_owned(),
2093 |name| name.to_string_lossy().into_owned(),
2094 );
2095 format!("{project} via {profile_id}")
2096}
2097
2098pub fn normalize_session_title(title: &str) -> Option<String> {
2099 let normalized = crate::relay::strip_hidden_prompt_context(title)
2100 .split_whitespace()
2101 .collect::<Vec<_>>()
2102 .join(" ");
2103 (!normalized.is_empty()).then_some(normalized)
2104}
2105
2106pub fn provisional_session_title(prompt: &str) -> Option<String> {
2112 const MAX_TITLE_CHARS: usize = 64;
2113
2114 let normalized = normalize_session_title(prompt)?;
2115 if normalized.chars().count() <= MAX_TITLE_CHARS {
2116 return Some(normalized);
2117 }
2118
2119 let mut truncated = normalized
2120 .chars()
2121 .take(MAX_TITLE_CHARS - 1)
2122 .collect::<String>();
2123 if let Some(boundary) = truncated.rfind(char::is_whitespace) {
2124 truncated.truncate(boundary);
2125 }
2126 truncated.push('…');
2127 Some(truncated)
2128}
2129
2130pub fn short_id(id: &str) -> &str {
2131 id.get(..8).unwrap_or(id)
2132}
2133
2134#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
2135pub struct RecoveryCandidate {
2136 pub session_id: String,
2137 pub target_template_id: String,
2138 pub locator: TargetLocator,
2139 pub ownership: Option<crate::worker_launch::WorkerOwnership>,
2140 #[serde(default)]
2143 pub instance_id: Option<String>,
2144 #[serde(default, skip_serializing_if = "Option::is_none")]
2149 pub tracked_session: Option<SessionState>,
2150}
2151
2152#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
2153pub struct RecoveryScan {
2154 pub candidates: Vec<RecoveryCandidate>,
2155 pub warnings: Vec<String>,
2156 #[serde(default)]
2158 pub instance_id: String,
2159 #[serde(default)]
2162 pub hidden_other_instances: usize,
2163}
2164
2165#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2166#[serde(deny_unknown_fields)]
2167pub struct ResumeRepositorySourceReceipt {
2168 pub session_id: String,
2169 pub bundle_id: String,
2170 pub checkpoint_sha256: String,
2171 pub repositories: Vec<crate::config::ProjectRepository>,
2172}
2173
2174#[cfg(test)]
2175mod tests;