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 StartupCleanup,
36 Running,
37 Disconnected,
38 Checkpointing,
39 Closing,
40 Destroying,
41 #[serde(alias = "archived")]
44 Stopped,
45 Parked,
52 Lost,
53 Error,
54 DestroyedWithDataLoss,
55}
56
57#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
60#[serde(rename_all = "kebab-case")]
61pub enum SessionTransitionKind {
62 Starting,
63 Resuming,
64 Moving,
65 Suspending,
66 Destroying,
67 Stopping,
69}
70
71impl SessionTransitionKind {
72 pub const fn label(self) -> &'static str {
73 match self {
74 Self::Starting => "Starting",
75 Self::Resuming => "Resuming",
76 Self::Moving => "Moving",
77 Self::Suspending => "Suspending",
78 Self::Destroying => "Destroying",
79 Self::Stopping => "Stopping",
80 }
81 }
82
83 pub fn for_session(state: SessionState, operation: Option<Self>) -> Option<Self> {
84 operation.or_else(|| state.transition_kind())
85 }
86}
87
88#[cfg(test)]
89mod transition_tests {
90 use super::{SessionState, SessionTransitionKind};
91
92 #[test]
93 fn operation_ownership_hides_intermediate_move_states_but_not_ordinary_live_work() {
94 for state in [
95 SessionState::Stopped,
96 SessionState::Running,
97 SessionState::Disconnected,
98 ] {
99 assert_eq!(
100 SessionTransitionKind::for_session(state, Some(SessionTransitionKind::Moving)),
101 Some(SessionTransitionKind::Moving)
102 );
103 assert_eq!(SessionTransitionKind::for_session(state, None), None);
104 }
105 assert_eq!(SessionState::Checkpointing.transition_kind(), None);
106 assert_eq!(
107 SessionState::Closing.transition_kind(),
108 Some(SessionTransitionKind::Suspending)
109 );
110 }
111}
112
113#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
115#[serde(tag = "state", rename_all = "snake_case")]
116pub enum MaterializedExecutionState {
117 #[default]
118 Idle,
119 Running {
120 started_at_ms: i64,
121 },
122 Closing,
123 Closed,
124}
125
126pub use crate::transcript::{TerminalOutputRecord, TranscriptBody, TranscriptItem};
127
128#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
133#[serde(rename_all = "snake_case")]
134pub enum QueuedCommandKind {
135 #[default]
136 Prompt,
137 SetConfig {
138 key: String,
139 value: String,
140 },
141}
142
143impl QueuedCommandKind {
144 pub fn is_prompt(&self) -> bool {
145 matches!(self, Self::Prompt)
146 }
147}
148
149pub fn config_command_text(key: &str, value: &str) -> String {
152 if key == "fast-mode" {
153 "/fast".to_owned()
154 } else {
155 format!("/{key} {value}")
156 }
157}
158
159#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
160#[serde(deny_unknown_fields)]
161pub struct MaterializedQueuedPrompt {
162 pub command_id: String,
163 #[serde(default, skip_serializing_if = "QueuedCommandKind::is_prompt")]
164 pub kind: QueuedCommandKind,
165 pub content: Vec<serde_json::Value>,
166 pub queued_at_ms: i64,
167 #[serde(default, skip_serializing_if = "Option::is_none")]
171 pub accepted_ordinal: Option<u64>,
172}
173
174#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
177#[serde(deny_unknown_fields)]
178pub struct MaterializedTurn {
179 pub command_id: String,
180 #[serde(default, skip_serializing_if = "Option::is_none")]
181 pub accepted_ordinal: Option<u64>,
182 pub turn_start_position: u64,
185 pub started_at_ms: i64,
186 #[serde(default, skip_serializing_if = "Option::is_none")]
189 pub steered_into: Option<String>,
190}
191
192impl MaterializedTurn {
193 pub fn belongs_to(&self, command_id: &str) -> bool {
195 self.command_id == command_id || self.steered_into.as_deref() == Some(command_id)
196 }
197}
198
199#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
201#[serde(tag = "kind", rename_all = "snake_case")]
202pub enum TurnOutcomeKind {
203 Completed { stop_reason: String },
205 Rejected {
207 message: String,
208 #[serde(default, skip_serializing_if = "Option::is_none")]
209 reason: Option<crate::event_outcome::OutcomeReason>,
210 },
211 Interrupted {
213 message: String,
214 #[serde(default, skip_serializing_if = "Option::is_none")]
215 reason: Option<crate::event_outcome::OutcomeReason>,
216 },
217}
218
219impl std::fmt::Display for TurnOutcomeKind {
222 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
223 use crate::event_outcome::{OutcomeReason, TurnResultKind};
224 let result = self.result();
225 match result.kind {
226 TurnResultKind::Completed => formatter.write_str("completed, end of turn"),
227 TurnResultKind::InputRequired => formatter.write_str("completed, waiting for input"),
228 TurnResultKind::Cancelled | TurnResultKind::Interrupted => {
229 formatter.write_str("interrupted")
230 }
231 TurnResultKind::Failed if result.reason == Some(OutcomeReason::QuotaLimit) => {
232 formatter.write_str("failed: quota limit reached")
233 }
234 TurnResultKind::Failed | TurnResultKind::Rejected => {
235 let message = result.message.as_deref().unwrap_or("unknown failure");
236 if result.stop_reason.is_some() {
237 write!(formatter, "failed: {}", stop_reason_words(message))
238 } else {
239 write!(
240 formatter,
241 "failed: {}",
242 message.lines().next().unwrap_or_default().trim()
243 )
244 }
245 }
246 }
247 }
248}
249
250fn stop_reason_words(stop_reason: &str) -> String {
253 let mut words = String::new();
254 let mut previous_lower = false;
255 for character in stop_reason.trim().chars() {
256 if character == '_' || character == '-' || character.is_whitespace() {
257 if !words.ends_with(' ') && !words.is_empty() {
258 words.push(' ');
259 }
260 previous_lower = false;
261 continue;
262 }
263 if character.is_uppercase() && previous_lower {
264 words.push(' ');
265 }
266 previous_lower = character.is_lowercase() || character.is_ascii_digit();
267 words.extend(character.to_lowercase());
268 }
269 match words.trim_end() {
270 "" => "no reason given".to_owned(),
271 words => words.to_owned(),
272 }
273}
274
275#[derive(Debug, Clone, Copy, PartialEq, Eq)]
276pub enum PromptCompletion {
277 InputRequired,
278 Finished,
279 Cancelled,
280 QuotaLimit,
281 Error,
282}
283
284pub fn classify_prompt_completion(stop_reason: &str) -> PromptCompletion {
286 let normalized = stop_reason
287 .chars()
288 .filter(|character| *character != '_' && *character != '-')
289 .flat_map(char::to_lowercase)
290 .collect::<String>();
291 match normalized.as_str() {
292 "endturn" => PromptCompletion::Finished,
293 "awaitinginput" => PromptCompletion::InputRequired,
294 "cancelled" | "canceled" => PromptCompletion::Cancelled,
295 "quotalimit" => PromptCompletion::QuotaLimit,
296 _ => PromptCompletion::Error,
297 }
298}
299
300#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
302#[serde(deny_unknown_fields)]
303pub struct MaterializedTurnOutcome {
304 #[serde(default, skip_serializing_if = "Option::is_none")]
305 pub diagnostic: Option<crate::diagnostic::TurnDiagnostic>,
306
307 #[serde(default, skip_serializing_if = "Option::is_none")]
308 pub usage: Option<crate::usage::TokenUsage>,
309 pub command_id: String,
310 #[serde(default, skip_serializing_if = "Option::is_none")]
311 pub accepted_ordinal: Option<u64>,
312 #[serde(default, skip_serializing_if = "Option::is_none")]
313 pub turn_start_position: Option<u64>,
314 pub completed_ordinal: u64,
315 pub completed_at_ms: i64,
316 pub outcome: TurnOutcomeKind,
317}
318
319impl MaterializedTurnOutcome {
320 pub fn interruption_ordinal(&self) -> Option<u64> {
322 (self.turn_start_position.is_some()
323 && matches!(self.outcome, TurnOutcomeKind::Interrupted { .. }))
324 .then_some(self.completed_ordinal)
325 }
326}
327
328#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
330#[serde(deny_unknown_fields)]
331pub struct MaterializedSession {
332 pub session_id: String,
333 pub applied_event_ordinal: u64,
334 pub applied_event_digest: String,
335 pub last_activity_at_ms: Option<i64>,
338 pub execution: MaterializedExecutionState,
339 #[serde(default, skip_serializing_if = "Option::is_none")]
340 pub session_title: Option<String>,
341 #[serde(default, skip_serializing_if = "SessionConfiguration::is_empty")]
342 pub configuration: SessionConfiguration,
343 #[serde(default, skip_serializing_if = "Vec::is_empty")]
344 pub transcript: Vec<Arc<TranscriptItem>>,
347 #[serde(default, skip_serializing_if = "Vec::is_empty")]
348 pub queued_prompts: Vec<MaterializedQueuedPrompt>,
349 #[serde(default, skip_serializing_if = "Vec::is_empty")]
352 pub pending_elicitations: Vec<crate::elicitation::ElicitationRequest>,
353 #[serde(default, skip_serializing_if = "Option::is_none")]
355 pub active_turn: Option<MaterializedTurn>,
356 #[serde(default, skip_serializing_if = "Option::is_none")]
359 pub last_turn_outcome: Option<MaterializedTurnOutcome>,
360}
361
362#[derive(Debug, Clone, PartialEq, Eq)]
365pub struct MaterializedSessionSummary {
366 pub session_id: String,
367 pub applied_event_ordinal: u64,
368 pub last_activity_at_ms: Option<i64>,
369 pub execution: MaterializedExecutionState,
370 pub session_title: Option<String>,
371 pub last_agent_message: Option<String>,
372 pub last_user_message: Option<String>,
373 pub last_agent_message_follows_last_user: bool,
376 pub agent_message_latest_content_ordinals: Vec<u64>,
377 pub interruption_event_ordinals: Vec<u64>,
378}
379
380impl MaterializedSession {
381 pub fn empty(session_id: impl Into<String>) -> Self {
382 Self {
383 session_id: session_id.into(),
384 applied_event_ordinal: 0,
385 applied_event_digest: RELAY_EVENT_GENESIS_DIGEST.into(),
386 last_activity_at_ms: None,
387 execution: MaterializedExecutionState::Idle,
388 session_title: None,
389 configuration: SessionConfiguration::default(),
390 transcript: Vec::new(),
391 queued_prompts: Vec::new(),
392 pending_elicitations: Vec::new(),
393 active_turn: None,
394 last_turn_outcome: None,
395 }
396 }
397
398 pub fn last_activity_at_ms(&self) -> Option<i64> {
399 self.last_activity_at_ms
400 }
401
402 pub fn resolved_title(&self) -> Option<String> {
408 self.session_title
409 .as_deref()
410 .and_then(normalize_session_title)
411 .or_else(|| {
412 self.transcript.iter().find_map(|item| {
413 let TranscriptBody::User { content } = &item.body else {
414 return None;
415 };
416 provisional_session_title(&crate::transcript::materialized_content_text(
417 content,
418 ))
419 })
420 })
421 .or_else(|| {
422 self.queued_prompts
423 .iter()
424 .filter(|prompt| prompt.kind.is_prompt())
425 .find_map(|prompt| {
426 provisional_session_title(&crate::transcript::materialized_content_text(
427 &prompt.content,
428 ))
429 })
430 })
431 }
432
433 pub fn unread_agent_messages_after(&self, viewed_through_event_ordinal: u64) -> u64 {
434 self.transcript
435 .iter()
436 .filter(|item| {
437 item.latest_content_event_ordinal
438 .is_some_and(|ordinal| ordinal > viewed_through_event_ordinal)
439 && item.is_nonempty_agent_message()
440 })
441 .count() as u64
442 }
443
444 pub fn unread_interruptions_after(&self, viewed_through_event_ordinal: u64) -> u64 {
445 self.interruption_event_ordinals()
446 .into_iter()
447 .filter(|ordinal| *ordinal > viewed_through_event_ordinal)
448 .count() as u64
449 }
450
451 pub fn interruption_event_ordinals(&self) -> Vec<u64> {
452 let mut ordinals = self
453 .transcript
454 .iter()
455 .filter(|item| item.is_work_interruption())
456 .map(|item| item.position)
457 .collect::<Vec<_>>();
458 if let Some(ordinal) = self
459 .last_turn_outcome
460 .as_ref()
461 .and_then(MaterializedTurnOutcome::interruption_ordinal)
462 {
463 ordinals.push(ordinal);
464 }
465 ordinals.sort_unstable();
466 ordinals.dedup();
467 ordinals
468 }
469
470 pub fn validate(&self) -> Result<()> {
471 validate_id("session", &self.session_id)?;
472 validate_relay_event_frontier(
473 self.applied_event_ordinal,
474 &self.applied_event_digest,
475 "materialized session event frontier",
476 )?;
477 if self
478 .session_title
479 .as_ref()
480 .is_some_and(|title| title.trim().is_empty())
481 {
482 bail!("materialized session has an empty title");
483 }
484 let mut item_ids = BTreeSet::new();
485 for item in &self.transcript {
486 item.validate(self.applied_event_ordinal)?;
487 if !item_ids.insert(item.stable_id.as_str()) {
488 bail!(
489 "materialized transcript contains duplicate item {:?}",
490 item.stable_id
491 );
492 }
493 }
494 let mut command_ids = BTreeSet::new();
495 for prompt in &self.queued_prompts {
496 if prompt.command_id.trim().is_empty() {
497 bail!("materialized prompt queue has an empty command id");
498 }
499 if !command_ids.insert(prompt.command_id.as_str()) {
500 bail!(
501 "materialized prompt queue contains duplicate command {:?}",
502 prompt.command_id
503 );
504 }
505 if let QueuedCommandKind::SetConfig { key, value } = &prompt.kind
506 && (key.trim().is_empty() || value.trim().is_empty())
507 {
508 bail!(
509 "materialized queued configuration change {:?} is incomplete",
510 prompt.command_id
511 );
512 }
513 }
514 Ok(())
515 }
516}
517
518#[derive(Debug, Clone, PartialEq)]
522pub struct ManagedSessionSnapshot {
523 pub materialized: MaterializedSession,
524 pub window: ProjectionWindow,
527 pub operational: RelayOperationalState,
528 pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
532 pub worker_build: Option<String>,
537 pub subagent_requests: Vec<crate::subagent::SubagentToolRequest>,
539 pub subagent_results: Vec<crate::subagent::SubagentToolResult>,
541}
542
543#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
555pub struct ProjectionWindow {
556 pub omitted_items: usize,
558 pub provisional_title: Option<String>,
560 pub latest_turn_start_position: Option<u64>,
564}
565
566impl ProjectionWindow {
567 pub fn trim(&mut self, session: &mut MaterializedSession, target: usize) {
570 let observed = Self::of(session);
571 if self.provisional_title.is_none() {
572 self.provisional_title = observed.provisional_title;
573 }
574 self.latest_turn_start_position = observed
575 .latest_turn_start_position
576 .or(self.latest_turn_start_position);
577 let mut boundary = session.transcript.len().saturating_sub(target.max(1));
578 for (index, item) in session.transcript.iter().enumerate() {
579 let mutable = match &item.body {
580 TranscriptBody::Agent { streaming, .. }
581 | TranscriptBody::Thought { streaming, .. } => *streaming,
582 TranscriptBody::Tool { call, .. } => matches!(
583 call.get("status").and_then(serde_json::Value::as_str),
584 Some("pending" | "in_progress")
585 ),
586 _ => false,
587 };
588 if mutable || Some(item.position) == self.latest_turn_start_position {
589 boundary = boundary.min(index);
590 }
591 }
592 let cut = session
593 .transcript
594 .iter()
595 .take(boundary + 1)
596 .rposition(|item| item.is_turn_start())
597 .unwrap_or(0);
598 if cut > 0 {
599 session.transcript.drain(..cut);
600 self.omitted_items += cut;
601 }
602 }
603
604 #[must_use]
606 pub fn of(session: &MaterializedSession) -> Self {
607 Self {
608 omitted_items: 0,
609 provisional_title: session.transcript.iter().find_map(|item| {
610 let TranscriptBody::User { content } = &item.body else {
611 return None;
612 };
613 provisional_session_title(&crate::transcript::materialized_content_text(content))
614 }),
615 latest_turn_start_position: session
616 .transcript
617 .iter()
618 .rev()
619 .find(|item| item.is_turn_start())
620 .map(|item| item.position),
621 }
622 }
623}
624
625impl ManagedSessionSnapshot {
626 #[must_use]
631 pub fn resolved_title(&self) -> Option<String> {
632 self.materialized
633 .session_title
634 .as_deref()
635 .and_then(normalize_session_title)
636 .or_else(|| self.window.provisional_title.clone())
637 .or_else(|| {
638 self.materialized
639 .queued_prompts
640 .iter()
641 .filter(|prompt| prompt.kind.is_prompt())
642 .find_map(|prompt| {
643 provisional_session_title(&crate::transcript::materialized_content_text(
644 &prompt.content,
645 ))
646 })
647 })
648 }
649
650 #[must_use]
655 pub fn latest_completed_turn_ordinal(&self) -> Option<u64> {
656 if self.materialized.execution != MaterializedExecutionState::Idle {
657 return None;
658 }
659 self.window.latest_turn_start_position
660 }
661}
662
663#[derive(Debug, Clone)]
665pub struct RecoveryObservation {
666 pub session: SessionRecord,
667 pub config: Config,
668 pub latest_completed_turn_ordinal: Option<u64>,
669 pub checkpoint_wait: Option<crate::activity::CheckpointWait>,
673}
674
675pub fn latest_completed_turn_ordinal(session: &MaterializedSession) -> Option<u64> {
680 if session.execution != MaterializedExecutionState::Idle {
681 return None;
682 }
683 session
684 .transcript
685 .iter()
686 .rev()
687 .find(|item| item.is_turn_start())
688 .map(|item| item.position)
689}
690
691pub fn validate_relay_event_digest(digest: &str, name: &str) -> Result<()> {
692 if digest.len() != 64
693 || !digest
694 .bytes()
695 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
696 {
697 bail!("{name} must be a lowercase SHA-256 digest");
698 }
699 Ok(())
700}
701
702pub fn validate_relay_event_frontier(ordinal: u64, digest: &str, name: &str) -> Result<()> {
703 validate_relay_event_digest(digest, name)?;
704 if (ordinal == 0) != (digest == RELAY_EVENT_GENESIS_DIGEST) {
705 bail!("{name} has inconsistent ordinal {ordinal} and digest {digest}");
706 }
707 Ok(())
708}
709
710fn is_false(value: &bool) -> bool {
711 !*value
712}
713
714impl SessionState {
715 pub const fn as_str(self) -> &'static str {
717 match self {
718 Self::Provisioning => "provisioning",
719 Self::StartupCleanup => "startup-cleanup",
720 Self::Running => "running",
721 Self::Disconnected => "disconnected",
722 Self::Checkpointing => "checkpointing",
723 Self::Closing => "closing",
724 Self::Destroying => "destroying",
725 Self::Stopped => "stopped",
726 Self::Parked => "parked",
727 Self::Lost => "lost",
728 Self::Error => "error",
729 Self::DestroyedWithDataLoss => "destroyed-with-data-loss",
730 }
731 }
732
733 pub fn from_stored(value: &str) -> Option<Self> {
736 Some(match value {
737 "provisioning" => Self::Provisioning,
738 "startup-cleanup" => Self::StartupCleanup,
739 "running" => Self::Running,
740 "disconnected" => Self::Disconnected,
741 "checkpointing" => Self::Checkpointing,
742 "closing" => Self::Closing,
743 "destroying" => Self::Destroying,
744 "stopped" | "archived" => Self::Stopped,
745 "parked" => Self::Parked,
746 "lost" => Self::Lost,
747 "error" => Self::Error,
748 "destroyed-with-data-loss" => Self::DestroyedWithDataLoss,
749 _ => return None,
750 })
751 }
752
753 pub const fn transition_kind(self) -> Option<SessionTransitionKind> {
756 match self {
757 Self::Provisioning => Some(SessionTransitionKind::Starting),
758 Self::StartupCleanup => Some(SessionTransitionKind::Stopping),
759 Self::Closing => Some(SessionTransitionKind::Suspending),
760 Self::Destroying => Some(SessionTransitionKind::Destroying),
761 _ => None,
762 }
763 }
764
765 pub const fn is_active(self) -> bool {
772 matches!(
773 self,
774 Self::Provisioning
775 | Self::StartupCleanup
776 | Self::Running
777 | Self::Disconnected
778 | Self::Checkpointing
779 | Self::Closing
780 | Self::Destroying
781 | Self::Parked
782 | Self::Error
783 )
784 }
785
786 pub const fn has_live_worker(self) -> bool {
790 matches!(
791 self,
792 Self::Provisioning
793 | Self::StartupCleanup
794 | Self::Running
795 | Self::Disconnected
796 | Self::Checkpointing
797 | Self::Closing
798 | Self::Destroying
799 )
800 }
801}
802
803#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
804#[serde(tag = "kind", rename_all = "kebab-case")]
805pub enum PodmanWorkspaceLocator {
806 #[default]
807 ContainerLayer,
808 Volume {
809 name: String,
810 },
811 HostPath {
812 path: PathBuf,
813 helper: Vec<String>,
814 resource: String,
815 },
816}
817
818#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
819#[serde(tag = "kind", rename_all = "kebab-case")]
820pub enum TargetLocator {
821 LocalBare {
822 worker_root: PathBuf,
823 },
824 LocalPodman {
825 container_id: String,
826 #[serde(default)]
827 workspace_storage: PodmanWorkspaceLocator,
828 #[serde(default, skip_serializing_if = "Option::is_none")]
832 borrowed_from: Option<String>,
833 },
834 LocalDocker {
835 container_id: String,
836 #[serde(default, skip_serializing_if = "Option::is_none")]
840 borrowed_from: Option<String>,
841 },
842 AppleContainer {
843 container_id: String,
844 #[serde(default, skip_serializing_if = "Option::is_none")]
848 borrowed_from: Option<String>,
849 },
850 AwsEc2 {
851 instance_id: String,
852 #[serde(default, skip_serializing_if = "Option::is_none")]
853 address: Option<String>,
854 },
855 SshBare {
856 host: String,
857 workspace: PathBuf,
858 #[serde(default, skip_serializing_if = "Option::is_none")]
859 worker_id: Option<String>,
860 },
861 SshPodman {
862 host: String,
863 container_id: String,
864 #[serde(default)]
865 workspace_storage: PodmanWorkspaceLocator,
866 #[serde(default, skip_serializing_if = "Option::is_none")]
870 borrowed_from: Option<String>,
871 },
872 SshDocker {
873 host: String,
874 container_id: String,
875 #[serde(default, skip_serializing_if = "Option::is_none")]
879 borrowed_from: Option<String>,
880 },
881}
882
883impl ManagedWorktreeTarget {
884 pub fn same_location(&self, other: &Self) -> bool {
892 match (self, other) {
893 (Self::Local, Self::Local) => true,
894 (
895 Self::Ssh {
896 destination,
897 ssh_args,
898 },
899 Self::Ssh {
900 destination: other_destination,
901 ssh_args: other_args,
902 },
903 ) => {
904 destination == other_destination
905 && ssh_location_option(ssh_args, 'p', "port")
906 == ssh_location_option(other_args, 'p', "port")
907 && ssh_location_option(ssh_args, 'l', "user")
908 == ssh_location_option(other_args, 'l', "user")
909 }
910 _ => false,
911 }
912 }
913}
914
915fn ssh_location_option(args: &[String], flag: char, option: &str) -> Option<String> {
919 let mut args = args.iter();
920 while let Some(argument) = args.next() {
921 let Some(rest) = argument.strip_prefix('-') else {
922 continue;
923 };
924 let mut chars = rest.chars();
925 let Some(name) = chars.next() else { continue };
926 if name != flag && name != 'o' {
927 continue;
928 }
929 let inline = chars.as_str();
930 let value = if inline.is_empty() {
931 args.next().cloned()
932 } else {
933 Some(inline.to_owned())
934 };
935 if name == flag {
936 return value;
937 }
938 if let Some(setting) = value {
939 let (key, found) = setting
940 .split_once(['=', ' ', '\t'])
941 .unwrap_or((setting.as_str(), ""));
942 if key.trim().eq_ignore_ascii_case(option) {
943 return Some(found.trim().to_owned());
944 }
945 }
946 }
947 None
948}
949
950#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
951#[serde(tag = "kind", rename_all = "kebab-case")]
952pub enum ManagedWorktreeTarget {
953 Local,
954 Ssh {
955 destination: String,
956 #[serde(default, skip_serializing_if = "Vec::is_empty")]
957 ssh_args: Vec<String>,
958 },
959}
960
961#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
963#[serde(deny_unknown_fields)]
964pub struct ManagedWorktreeOptions {
965 pub available: bool,
966 pub default_create: bool,
967}
968
969#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
970#[serde(deny_unknown_fields)]
971pub struct ManagedWorktree {
972 #[serde(default, skip_serializing_if = "ManagedCheckoutKind::is_worktree")]
974 pub kind: ManagedCheckoutKind,
975 pub source_project_directory: PathBuf,
976 pub source_repository: PathBuf,
977 pub worktree_root: PathBuf,
978 pub branch: String,
979 pub target: ManagedWorktreeTarget,
980 #[serde(default, skip_serializing_if = "Option::is_none")]
984 pub base_commit: Option<String>,
985}
986
987#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
988#[serde(rename_all = "snake_case")]
989pub enum ManagedCheckoutKind {
990 #[default]
991 Worktree,
992 Clone,
993}
994
995impl ManagedCheckoutKind {
996 fn is_worktree(&self) -> bool {
997 matches!(self, Self::Worktree)
998 }
999}
1000
1001impl ManagedWorktree {
1002 fn validate(&self, session_id: &str, project_directory: Option<&Path>) -> Result<()> {
1003 for (label, path) in [
1004 ("source project directory", &self.source_project_directory),
1005 ("source repository", &self.source_repository),
1006 ("worktree root", &self.worktree_root),
1007 ] {
1008 if !path.is_absolute() || path.components().any(|part| part == Component::ParentDir) {
1009 bail!("managed worktree {label} must be an absolute safe path");
1010 }
1011 }
1012 if !self
1013 .source_project_directory
1014 .starts_with(&self.source_repository)
1015 {
1016 bail!("managed worktree source directory is outside its repository");
1017 }
1018 let expected_root = self
1019 .source_repository
1020 .join(".mj")
1021 .join(match self.kind {
1022 ManagedCheckoutKind::Worktree => "worktrees",
1023 ManagedCheckoutKind::Clone => "clones",
1024 })
1025 .join(session_id);
1026 if self.worktree_root != expected_root {
1027 bail!("managed worktree root does not match the session-owned path");
1028 }
1029 if self.kind == ManagedCheckoutKind::Worktree && self.branch != format!("mj/{session_id}") {
1030 bail!("managed worktree branch does not match the session id");
1031 }
1032 if self.kind == ManagedCheckoutKind::Clone && self.branch.trim().is_empty() {
1033 bail!("managed clone has no starting branch");
1034 }
1035 let relative = self
1036 .source_project_directory
1037 .strip_prefix(&self.source_repository)
1038 .expect("source relationship checked above");
1039 if project_directory != Some(self.worktree_root.join(relative).as_path()) {
1040 bail!("session project directory does not match its managed worktree");
1041 }
1042 match &self.target {
1043 ManagedWorktreeTarget::Local => {}
1044 ManagedWorktreeTarget::Ssh { destination, .. } if destination.trim().is_empty() => {
1045 bail!("managed SSH worktree has an empty destination")
1046 }
1047 ManagedWorktreeTarget::Ssh { .. } => {}
1048 }
1049 Ok(())
1050 }
1051}
1052
1053#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1054#[serde(tag = "kind", rename_all = "kebab-case")]
1055pub enum SessionResourceAllocation {
1056 Container {
1057 cpus: u64,
1058 memory_bytes: u64,
1059 },
1060 AwsEc2 {
1061 instance_type: String,
1062 vcpus: u64,
1063 memory_bytes: u64,
1064 },
1065}
1066
1067impl SessionResourceAllocation {
1068 pub fn validate(&self) -> Result<()> {
1069 match self {
1070 Self::Container { cpus, memory_bytes } if *cpus == 0 || *memory_bytes == 0 => {
1071 bail!("container resource allocation must have non-zero CPU and memory")
1072 }
1073 Self::AwsEc2 {
1074 instance_type,
1075 vcpus,
1076 memory_bytes,
1077 } if instance_type.trim().is_empty() || *vcpus == 0 || *memory_bytes == 0 => {
1078 bail!("EC2 resource allocation must have an instance type, CPU, and memory")
1079 }
1080 _ => Ok(()),
1081 }
1082 }
1083}
1084
1085pub fn allocation_cpus(allocation: &SessionResourceAllocation) -> u64 {
1087 match allocation {
1088 SessionResourceAllocation::Container { cpus, .. } => *cpus,
1089 SessionResourceAllocation::AwsEc2 { vcpus, .. } => *vcpus,
1090 }
1091}
1092
1093pub fn allocation_memory(allocation: &SessionResourceAllocation) -> u64 {
1095 match allocation {
1096 SessionResourceAllocation::Container { memory_bytes, .. }
1097 | SessionResourceAllocation::AwsEc2 { memory_bytes, .. } => *memory_bytes,
1098 }
1099}
1100
1101impl TargetLocator {
1102 pub const fn skills_scope(&self) -> crate::skills::SkillsScope {
1104 match self {
1105 Self::LocalBare { .. } => crate::skills::SkillsScope::Localhost,
1106 _ => crate::skills::SkillsScope::Isolated,
1107 }
1108 }
1109
1110 fn validate(&self, session_id: &str) -> Result<()> {
1111 match self {
1112 Self::LocalBare { worker_root } => {
1113 if !worker_root.is_absolute()
1114 || worker_root
1115 .components()
1116 .any(|part| part == Component::ParentDir)
1117 || !worker_root.ends_with(session_id)
1118 {
1119 bail!(
1120 "local bare worker root must be an absolute safe path ending in the session id"
1121 );
1122 }
1123 }
1124 Self::LocalPodman { container_id, .. }
1125 | Self::LocalDocker { container_id, .. }
1126 | Self::AppleContainer { container_id, .. }
1127 | Self::SshPodman { container_id, .. }
1128 | Self::SshDocker { container_id, .. }
1129 if container_id.trim().is_empty() =>
1130 {
1131 bail!("target locator has an empty container id")
1132 }
1133 Self::AwsEc2 { instance_id, .. } if instance_id.trim().is_empty() => {
1134 bail!("target locator has an empty AWS instance id")
1135 }
1136 Self::SshBare {
1137 host,
1138 workspace,
1139 worker_id,
1140 } => {
1141 if host.trim().is_empty() {
1142 bail!("bare SSH target locator has an empty host");
1143 }
1144 let unsafe_path = workspace.as_os_str().is_empty()
1145 || workspace
1146 .components()
1147 .any(|part| part == Component::ParentDir);
1148 match worker_id {
1149 Some(worker_id) => {
1152 if worker_id != session_id {
1153 bail!(
1154 "bare SSH target locator's worker identity does not match the session id"
1155 );
1156 }
1157 if unsafe_path {
1158 bail!("bare SSH target locator must have a safe workspace path");
1159 }
1160 }
1161 None => {
1162 if unsafe_path || !workspace.ends_with(session_id) {
1163 bail!(
1164 "bare SSH target locator must be a safe path ending in the session id"
1165 );
1166 }
1167 }
1168 }
1169 }
1170 Self::SshPodman { host, .. } if host.trim().is_empty() => {
1171 bail!("SSH Podman target locator has an empty host")
1172 }
1173 Self::SshDocker { host, .. } if host.trim().is_empty() => {
1174 bail!("SSH Docker target locator has an empty host")
1175 }
1176 _ => {}
1177 }
1178 Ok(())
1179 }
1180}
1181
1182#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1183#[serde(deny_unknown_fields)]
1184pub struct CheckpointMetadata {
1185 pub archive_path: PathBuf,
1186 pub sha256: String,
1188 pub created_at: String,
1189 pub event_frontier: u64,
1190}
1191
1192#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1194#[serde(rename_all = "snake_case")]
1195pub enum PublicationState {
1196 Published,
1197 Unpublished,
1198 Unknown,
1199}
1200
1201#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1203pub struct PublicationAssessment {
1204 pub checkpoint_sha256: String,
1205 pub state: PublicationState,
1206 pub dirty: bool,
1207 pub stashed: bool,
1208 pub saved_commits: Vec<String>,
1209 pub destinations: Vec<String>,
1210 pub checked_at: String,
1211 pub reason: Option<String>,
1212}
1213
1214impl CheckpointMetadata {
1215 fn validate(&self) -> Result<()> {
1216 if self.archive_path.as_os_str().is_empty() {
1217 bail!("checkpoint archive path is empty");
1218 }
1219 if self.sha256.len() != 64
1220 || !self
1221 .sha256
1222 .bytes()
1223 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1224 {
1225 bail!("checkpoint SHA-256 must be 64 lowercase hexadecimal characters");
1226 }
1227 if self.created_at.trim().is_empty() {
1228 bail!("checkpoint timestamp is empty");
1229 }
1230 Ok(())
1231 }
1232}
1233
1234#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1238#[serde(deny_unknown_fields)]
1239pub struct SessionBuildCache {
1240 pub host: String,
1243 pub directory: PathBuf,
1244 #[serde(default, skip_serializing_if = "Option::is_none")]
1247 pub max_size: Option<String>,
1248 #[serde(default, skip_serializing_if = "Option::is_none")]
1251 pub target_root: Option<PathBuf>,
1252}
1253
1254#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1258pub struct ArchiveSpacePreview {
1259 pub sessions: usize,
1261 pub bytes: u64,
1263 pub reclaimable_sessions: usize,
1266 pub reclaimable_bytes: u64,
1267}
1268
1269#[derive(Debug, Clone, PartialEq, Eq)]
1273pub struct BuildCachePreview {
1274 pub native_mbx: Option<String>,
1276 pub directory: Option<PathBuf>,
1278 pub max_total_size: Option<BuildCacheLimit>,
1280 pub user_managed: bool,
1282 pub application: BuildCacheApplication,
1283 pub budget_note: Option<String>,
1285 pub stats: Option<BuildCacheStats>,
1287 pub off_reason: Option<BuildCacheOff>,
1290}
1291
1292#[derive(Debug, Clone, PartialEq, Eq)]
1300pub struct BuildCacheStats {
1301 pub builds: u64,
1303 pub cached_compilations: u64,
1305 pub avoided_compiler_ns: u64,
1307 pub reflinked_bytes: u64,
1309}
1310
1311#[derive(Debug, Clone, PartialEq, Eq)]
1313pub enum BuildCacheOff {
1314 TurnedOff,
1316 Unavailable(String),
1319}
1320
1321impl std::fmt::Display for BuildCacheOff {
1322 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1323 match self {
1324 Self::TurnedOff => formatter.write_str("turned off for this machine"),
1325 Self::Unavailable(reason) => formatter.write_str(reason),
1326 }
1327 }
1328}
1329
1330#[derive(Debug, Clone, PartialEq, Eq)]
1332pub enum BuildCacheLimit {
1333 Size(String),
1335 HostConfiguration(Option<String>),
1338 MjDefault(String),
1340 MbxDefault(Option<String>),
1342}
1343
1344#[derive(Debug, Clone, Default, PartialEq, Eq)]
1345pub enum BuildCacheApplication {
1346 #[default]
1347 Pending,
1348 Applied,
1349 Failed(String),
1350}
1351
1352#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1353#[serde(deny_unknown_fields)]
1354pub struct SessionRecord {
1355 pub id: String,
1356 #[serde(default = "default_session_workspace_id")]
1361 pub workspace_id: String,
1362 pub title: String,
1363 pub harness_kind: HarnessKind,
1364 pub last_profile: String,
1365 pub bundle_id: String,
1366 #[serde(default, skip_serializing_if = "Option::is_none")]
1368 pub project: Option<crate::repository::ProjectBundleSnapshot>,
1369 #[serde(default, skip_serializing_if = "Option::is_none")]
1371 pub project_directory: Option<PathBuf>,
1372 #[serde(default, skip_serializing_if = "Option::is_none")]
1374 pub managed_worktree: Option<ManagedWorktree>,
1375 #[serde(default, skip_serializing_if = "Option::is_none")]
1377 pub create_managed_worktree: Option<bool>,
1378 #[serde(default, skip_serializing_if = "Option::is_none")]
1385 pub launch_base: Option<String>,
1386 #[serde(default, skip_serializing_if = "Option::is_none")]
1390 pub launch_branch: Option<String>,
1391 #[serde(default, skip_serializing_if = "Option::is_none")]
1395 pub checkout: Option<crate::remote_git::ExactCheckout>,
1396 #[serde(default, skip_serializing_if = "Option::is_none")]
1398 pub publication: Option<PublicationAssessment>,
1399 #[serde(
1401 default,
1402 alias = "mjolnir_subagents",
1403 deserialize_with = "crate::subagent::deserialize_optional_policy"
1404 )]
1405 pub subagents: Option<crate::subagent::SubagentPolicy>,
1406 pub target_template_id: String,
1407 #[serde(default, skip_serializing_if = "Option::is_none")]
1408 pub resource_allocation: Option<SessionResourceAllocation>,
1409 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1410 pub additional_mounts: Vec<AdditionalMount>,
1411 #[serde(default, skip_serializing_if = "Option::is_none")]
1414 pub container_cpus: Option<String>,
1415 #[serde(default, skip_serializing_if = "Option::is_none")]
1418 pub container_memory: Option<String>,
1419 #[serde(default, skip_serializing_if = "Option::is_none")]
1425 pub container_workspace: Option<PathBuf>,
1426 #[serde(default, skip_serializing_if = "Option::is_none")]
1430 pub build_cache: Option<SessionBuildCache>,
1431 pub state: SessionState,
1432 #[serde(default, skip_serializing_if = "is_false")]
1435 pub archived: bool,
1436 #[serde(default, skip_serializing_if = "Option::is_none")]
1437 pub target: Option<TargetLocator>,
1438 #[serde(default, skip_serializing_if = "Option::is_none")]
1440 pub target_runtime: Option<TargetRuntimeSettings>,
1441 #[serde(default, skip_serializing_if = "Option::is_none")]
1442 pub native_session_id: Option<String>,
1443 #[serde(default, skip_serializing_if = "Option::is_none")]
1444 pub acp_session_title: Option<String>,
1445 #[serde(default, skip_serializing_if = "Option::is_none")]
1446 pub session_title_override: Option<String>,
1447 pub created_at: String,
1448 pub updated_at: String,
1449 #[serde(default, alias = "detached_after_event_ordinal")]
1450 pub viewed_through_event_ordinal: u64,
1451 #[serde(default, skip_serializing_if = "String::is_empty")]
1454 pub draft_input: String,
1455 #[serde(default, skip_serializing_if = "Option::is_none")]
1463 pub last_error: Option<String>,
1464 #[serde(default, skip_serializing_if = "Option::is_none")]
1465 pub last_checkpoint_error: Option<String>,
1466 #[serde(default, skip_serializing_if = "Option::is_none")]
1467 pub checkpoint: Option<CheckpointMetadata>,
1468}
1469
1470#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1471#[serde(deny_unknown_fields)]
1472pub struct HostContainerSize {
1473 pub cpus: u64,
1474 pub memory_bytes: u64,
1475}
1476
1477fn default_session_workspace_id() -> String {
1478 crate::workspace::DEFAULT_WORKSPACE_ID.to_owned()
1479}
1480
1481pub const DESTRUCTION_FAILURE_PREFIX: &str = "the destruction did not finish";
1490
1491pub const CLOSE_FAILURE_PREFIX: &str = "the suspension did not finish";
1492
1493pub const MOVE_FAILURE_PREFIX: &str = "the move did not finish";
1494
1495pub fn is_public_lifecycle_error(error: &str) -> bool {
1497 error.starts_with(CLOSE_FAILURE_PREFIX)
1498 || error.starts_with(MOVE_FAILURE_PREFIX)
1499 || error.starts_with(DESTRUCTION_FAILURE_PREFIX)
1500 || error.starts_with("the close did not finish")
1501}
1502
1503#[must_use]
1516pub fn target_label(config: &Config, target_id: &str, project: Option<&Path>) -> String {
1517 if !matches!(
1518 config.targets.get(target_id),
1519 Some(TargetTemplate::LocalBare | TargetTemplate::SshBare { .. })
1520 ) {
1521 return target_id.to_owned();
1522 }
1523 project.and_then(Path::file_name).map_or_else(
1524 || target_id.to_owned(),
1525 |directory| format!("{target_id}/{}", directory.to_string_lossy()),
1526 )
1527}
1528
1529#[derive(Debug, Clone, Default, PartialEq, Eq)]
1531pub struct StartSelection {
1532 pub at: Option<String>,
1534 pub branch: Option<String>,
1536 pub base: Option<String>,
1538}
1539
1540impl SessionRecord {
1541 pub fn start_selection(&self) -> StartSelection {
1545 let at = self
1546 .checkout
1547 .as_ref()
1548 .map(|checkout| checkout.commit.clone());
1549 StartSelection {
1550 branch: self
1551 .checkout
1552 .as_ref()
1553 .and_then(|checkout| checkout.branch.clone())
1554 .or_else(|| self.launch_branch.clone()),
1555 base: self.launch_base.clone().or_else(|| at.clone()),
1556 at,
1557 }
1558 }
1559
1560 pub fn publication_state(&self) -> Option<PublicationState> {
1563 let independent_clone = self
1564 .managed_worktree
1565 .as_ref()
1566 .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Clone)
1567 || (self.managed_worktree.is_none() && self.project_directory.is_none());
1568 if !independent_clone {
1569 return None;
1570 }
1571 if self.state.is_active() {
1572 return Some(PublicationState::Unknown);
1573 }
1574 Some(
1575 self.checkpoint
1576 .as_ref()
1577 .zip(self.publication.as_ref())
1578 .filter(|(checkpoint, assessment)| {
1579 assessment.checkpoint_sha256 == checkpoint.sha256
1580 })
1581 .map_or(PublicationState::Unknown, |(_, assessment)| {
1582 if assessment.dirty || assessment.stashed {
1583 PublicationState::Unpublished
1584 } else {
1585 assessment.state
1586 }
1587 }),
1588 )
1589 }
1590
1591 pub fn target_runtime_settings<'a>(
1596 &'a self,
1597 config: &Config,
1598 ) -> Result<std::borrow::Cow<'a, TargetRuntimeSettings>> {
1599 if let Some(runtime) = &self.target_runtime {
1600 if let Some(refreshed) =
1603 config
1604 .targets
1605 .get(&self.target_template_id)
1606 .and_then(|template| {
1607 runtime.with_current_ssh_options(&TargetRuntimeSettings::from(template))
1608 })
1609 {
1610 return Ok(std::borrow::Cow::Owned(refreshed));
1611 }
1612 return Ok(std::borrow::Cow::Borrowed(runtime));
1613 }
1614 let template = config.targets.get(&self.target_template_id).ok_or_else(|| {
1615 crate::refusal::Refusal::precondition(format!(
1616 "Session {:?} has no recorded target access settings. Restore target {:?} in config.toml once, then retry.",
1617 self.id, self.target_template_id))
1618 })?;
1619 let runtime = TargetRuntimeSettings::from(template);
1620 if let Some(locator) = &self.target {
1621 crate::targets::TargetLocator::try_from(crate::targets::RecordedTarget {
1622 locator, runtime: Some(&runtime), session_id: &self.id,
1623 }).map_err(|error| crate::refusal::Refusal::precondition(format!(
1624 "Session {:?} cannot recover target {:?}: {error}. Restore its original target settings, then retry.", self.id, self.target_template_id)))?;
1625 }
1626 Ok(std::borrow::Cow::Owned(runtime))
1627 }
1628
1629 #[must_use]
1633 pub fn public_error(&self) -> Option<&str> {
1634 self.last_error
1635 .as_deref()
1636 .filter(|error| is_public_lifecycle_error(error))
1637 }
1638
1639 pub fn configuration_issue(&self, config: &Config) -> Option<String> {
1642 if !self.state.is_active() {
1643 return None;
1644 }
1645 let mut issues = Vec::new();
1646 match config.profiles.get(&self.last_profile) {
1647 None => issues.push(format!("missing profile {:?}", self.last_profile)),
1648 Some(profile) if profile.kind != self.harness_kind => issues.push(format!(
1649 "expects {:?}, but profile {:?} is {:?}",
1650 self.harness_kind, self.last_profile, profile.kind
1651 )),
1652 Some(_) => {}
1653 }
1654 if self.project_directory.is_none() && self.project_bundle(config).is_none() {
1655 issues.push(format!("missing bundle {:?}", self.bundle_id));
1656 }
1657 if self.target_runtime.is_none() && !config.targets.contains_key(&self.target_template_id) {
1658 issues.push(format!(
1659 "missing target template {:?}",
1660 self.target_template_id
1661 ));
1662 }
1663 (!issues.is_empty()).then(|| format!(
1664 "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.",
1665 self.id, issues.join("; ")
1666 ))
1667 }
1668
1669 pub fn validate_configuration(&self, config: &Config) -> Result<()> {
1670 if let Some(issue) = self.configuration_issue(config) {
1671 return Err(crate::refusal::Refusal::precondition(issue).into());
1672 }
1673 Ok(())
1674 }
1675
1676 pub fn display_title(&self) -> &str {
1678 self.session_title_override
1679 .as_deref()
1680 .or(self.acp_session_title.as_deref())
1681 .unwrap_or(&self.id)
1682 }
1683
1684 pub fn listed_title(&self) -> &str {
1689 let named = self.session_title_override.is_some() || self.acp_session_title.is_some();
1690 if !named && !self.title.trim().is_empty() {
1691 return &self.title;
1692 }
1693 self.display_title()
1694 }
1695
1696 pub fn project_name(&self, config: &Config) -> String {
1701 if let Some(project) = &self.project {
1702 return project.name();
1703 }
1704 if let Some(worktree) = &self.managed_worktree {
1705 return path_leaf(&worktree.source_repository);
1706 }
1707 if let Some(project_directory) = &self.project_directory {
1708 return path_leaf(project_directory);
1709 }
1710 self.bundle_source_name(config)
1711 }
1712
1713 pub fn project_target(&self, config: &Config, target_id: &str) -> String {
1717 let project = self
1718 .managed_worktree
1719 .as_ref()
1720 .map(|worktree| &worktree.source_project_directory)
1721 .or(self.project_directory.as_ref());
1722 target_label(config, target_id, project.map(PathBuf::as_path))
1723 }
1724
1725 pub fn project_source(&self, config: &Config) -> ProjectSourceIdentity {
1730 if let Some(project) = &self.project {
1731 return ProjectSourceIdentity {
1732 key: project
1733 .source_key()
1734 .expect("accepted project has complete identities"),
1735 short: project.name(),
1736 full: project
1737 .identities
1738 .values()
1739 .map(crate::repository::RepositoryIdentity::key)
1740 .collect::<Vec<_>>()
1741 .join(" + "),
1742 };
1743 }
1744 if let Some(worktree) = &self.managed_worktree {
1745 return ProjectSourceIdentity::path(&worktree.source_repository, None);
1746 }
1747 if let Some(project_directory) = &self.project_directory {
1748 let remote = match &self.target {
1749 Some(TargetLocator::SshBare { host, .. }) => Some(host.as_str()),
1750 _ => None,
1751 };
1752 return ProjectSourceIdentity::path(project_directory, remote);
1753 }
1754 self.bundle_source_identity(config)
1755 .unwrap_or_else(|| ProjectSourceIdentity {
1756 key: format!("bundle:{}", self.bundle_id),
1757 short: path_leaf(Path::new(&self.bundle_id)),
1758 full: self.bundle_id.clone(),
1759 })
1760 }
1761
1762 fn bundle_source_name(&self, config: &Config) -> String {
1765 self.bundle_source_identity(config)
1766 .map(|source| source.short)
1767 .unwrap_or_else(|| path_leaf(Path::new(&self.bundle_id)))
1768 }
1769
1770 fn bundle_source_identity(&self, config: &Config) -> Option<ProjectSourceIdentity> {
1773 let bundle = self.project_bundle(config)?;
1774 let sources = bundle
1775 .repositories
1776 .iter()
1777 .map(repository_source_identity)
1778 .collect::<Option<Vec<_>>>()?;
1779 ProjectSourceIdentity::bundle(sources)
1780 }
1781
1782 pub fn project_bundle<'a>(
1784 &'a self,
1785 config: &'a Config,
1786 ) -> Option<&'a crate::config::ProjectBundle> {
1787 self.project
1788 .as_ref()
1789 .map(|project| &project.bundle)
1790 .or_else(|| config.bundles.get(&self.bundle_id))
1791 }
1792
1793 pub fn compare_by_creation(&self, other: &Self) -> std::cmp::Ordering {
1797 self.creation_order_key().cmp(&other.creation_order_key())
1798 }
1799
1800 pub fn creation_order_key(&self) -> (bool, Option<i64>, &str) {
1802 let timestamp = created_at_seconds(&self.created_at);
1803 (timestamp.is_none(), timestamp, &self.id)
1804 }
1805
1806 fn validate(&self, map_id: &str) -> Result<()> {
1807 validate_id("session", &self.id)?;
1808 if self.id != map_id {
1809 bail!(
1810 "session map key {map_id:?} does not match record id {:?}",
1811 self.id
1812 );
1813 }
1814 if let Some(project) = &self.project {
1815 project.key()?;
1816 }
1817 validate_id("workspace", &self.workspace_id)?;
1818 validate_id("profile", &self.last_profile)?;
1819 validate_id("bundle", &self.bundle_id)?;
1820 if let Some(project_directory) = &self.project_directory
1821 && (!project_directory.is_absolute()
1822 || project_directory
1823 .components()
1824 .any(|part| part == Component::ParentDir))
1825 {
1826 bail!("session {:?} has an unsafe project directory", self.id);
1827 }
1828 if let Some(managed_worktree) = &self.managed_worktree {
1829 managed_worktree.validate(&self.id, self.project_directory.as_deref())?;
1830 }
1831 validate_id("target template", &self.target_template_id)?;
1832 if let Some(allocation) = &self.resource_allocation {
1833 allocation.validate()?;
1834 }
1835 validate_additional_mounts(&self.additional_mounts)?;
1836 if self.title.trim().is_empty() {
1837 bail!("session {:?} has an empty title", self.id);
1838 }
1839 if self
1840 .acp_session_title
1841 .as_ref()
1842 .is_some_and(|title| title.trim().is_empty())
1843 || self
1844 .session_title_override
1845 .as_ref()
1846 .is_some_and(|title| title.trim().is_empty())
1847 {
1848 bail!("session {:?} has an empty display title", self.id);
1849 }
1850 if self.created_at.trim().is_empty() || self.updated_at.trim().is_empty() {
1851 bail!("session {:?} has an empty timestamp", self.id);
1852 }
1853 if let Some(target) = &self.target {
1854 target.validate(&self.id)?;
1855 }
1856 if let Some(checkpoint) = &self.checkpoint {
1857 checkpoint.validate()?;
1858 }
1859 Ok(())
1860 }
1861}
1862
1863fn repository_source_identity(repository: &ProjectRepository) -> Option<ProjectSourceIdentity> {
1864 repository
1865 .github
1866 .as_deref()
1867 .and_then(ProjectSourceIdentity::git_remote)
1868 .or_else(|| {
1869 repository
1870 .local
1871 .as_deref()
1872 .map(|path| ProjectSourceIdentity::path(path, None))
1873 })
1874}
1875
1876#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1877pub struct ProjectSourceIdentity {
1878 pub key: String,
1879 pub short: String,
1880 pub full: String,
1881}
1882
1883impl ProjectSourceIdentity {
1884 pub fn bundle(mut sources: Vec<Self>) -> Option<Self> {
1886 if sources.is_empty() {
1887 return None;
1888 }
1889 sources.sort_by(|left, right| {
1890 left.key
1891 .cmp(&right.key)
1892 .then_with(|| left.full.cmp(&right.full))
1893 .then_with(|| left.short.cmp(&right.short))
1894 });
1895 sources.dedup_by(|left, right| left.key == right.key);
1896 if sources.len() == 1 {
1897 return sources.pop();
1898 }
1899 let keys = sources
1900 .iter()
1901 .map(|source| source.key.clone())
1902 .collect::<Vec<_>>();
1903 let key = serde_json::to_string(&keys).ok()?;
1904 Some(Self {
1905 key: format!("bundle:{key}"),
1906 short: sources
1907 .iter()
1908 .map(|source| source.short.as_str())
1909 .collect::<Vec<_>>()
1910 .join(" + "),
1911 full: sources
1912 .iter()
1913 .map(|source| source.full.as_str())
1914 .collect::<Vec<_>>()
1915 .join(" + "),
1916 })
1917 }
1918
1919 pub fn git_remote(source: &str) -> Option<Self> {
1922 let identity = crate::repository::RepositoryIdentity::from_remote(source)?;
1923 let full = crate::repository::RepositoryIdentity::remote_label(source)?;
1924 Some(Self {
1925 key: identity.key(),
1926 short: full.rsplit(['/', ':']).next()?.to_owned(),
1927 full,
1928 })
1929 }
1930
1931 pub fn path(path: &Path, remote: Option<&str>) -> Self {
1933 let normalized = path.components().collect::<PathBuf>();
1934 let path_text = normalized.to_string_lossy().into_owned();
1935 let full = remote.map_or_else(|| path_text.clone(), |host| format!("{host}:{path_text}"));
1936 let key = remote.map_or_else(
1937 || format!("path:{path_text}"),
1938 |host| format!("path:{}:{path_text}", host.to_lowercase()),
1939 );
1940 Self {
1941 key,
1942 short: path_leaf(path),
1943 full,
1944 }
1945 }
1946}
1947
1948fn path_leaf(path: &Path) -> String {
1950 path.file_name()
1951 .unwrap_or(path.as_os_str())
1952 .to_string_lossy()
1953 .into_owned()
1954}
1955
1956fn created_at_seconds(timestamp: &str) -> Option<i64> {
1957 chrono::DateTime::parse_from_rfc3339(timestamp)
1958 .ok()
1959 .map(|timestamp| timestamp.timestamp())
1960}
1961
1962#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1963#[serde(deny_unknown_fields)]
1964pub struct State {
1965 #[serde(default)]
1966 pub last_subagent_policy: crate::subagent::SubagentPolicy,
1967 pub version: u32,
1968 #[serde(default, skip_serializing_if = "SnapshotMap::is_empty")]
1969 pub sessions: SnapshotMap<String, SessionRecord>,
1970 #[serde(default, skip_serializing_if = "SnapshotMap::is_empty")]
1973 pub subagents: SnapshotMap<String, SubagentRecord>,
1974 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1976 pub mount_history: BTreeMap<String, Vec<PathBuf>>,
1977 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1979 pub container_sizes: BTreeMap<String, HostContainerSize>,
1980}
1981
1982impl Default for State {
1983 fn default() -> Self {
1984 Self {
1985 version: STATE_VERSION,
1986 last_subagent_policy: Default::default(),
1987 sessions: SnapshotMap::new(),
1988 subagents: SnapshotMap::new(),
1989 mount_history: BTreeMap::new(),
1990 container_sizes: BTreeMap::new(),
1991 }
1992 }
1993}
1994
1995pub fn session_is_live(record: &SessionRecord, in_operation: bool, parent_live: bool) -> bool {
2003 record.state != SessionState::Stopped || in_operation || parent_live
2004}
2005
2006pub fn live_session_ids(
2010 sessions: &SnapshotMap<String, SessionRecord>,
2011 subagents: &SnapshotMap<String, SubagentRecord>,
2012 operations: &BTreeSet<String>,
2013) -> BTreeSet<String> {
2014 let mut live = sessions
2015 .iter()
2016 .filter(|(id, record)| session_is_live(record, operations.contains(*id), false))
2017 .map(|(id, _)| id.clone())
2018 .collect::<BTreeSet<_>>();
2019 loop {
2020 let joined = subagents
2021 .iter()
2022 .filter(|(child, relation)| {
2023 !live.contains(*child)
2024 && sessions.contains_key(*child)
2025 && live.contains(&relation.parent_session_id)
2026 })
2027 .map(|(child, _)| child.clone())
2028 .collect::<Vec<_>>();
2029 if joined.is_empty() {
2030 return live;
2031 }
2032 live.extend(joined);
2033 }
2034}
2035
2036impl State {
2037 #[must_use]
2042 pub fn session_notice_name(&self, session_id: &str) -> String {
2043 match self.sessions.get(session_id) {
2044 Some(session) if session.listed_title() != session.id => {
2045 session.listed_title().to_owned()
2046 }
2047 _ => short_id(session_id).to_owned(),
2048 }
2049 }
2050
2051 #[must_use]
2060 pub fn project_identity_session<'a>(&'a self, session: &'a SessionRecord) -> &'a SessionRecord {
2061 self.subagents
2062 .get(&session.id)
2063 .and_then(|record| self.sessions.get(&record.parent_session_id))
2064 .unwrap_or(session)
2065 }
2066
2067 #[must_use]
2074 pub fn is_subagent_session(&self, id: &str) -> bool {
2075 self.subagents.contains_key(id) || crate::native_agent::is_view_id(id)
2076 }
2077
2078 pub fn validate(&self) -> Result<()> {
2079 if self.version != STATE_VERSION {
2080 bail!(
2081 "unsupported Mjolnir state version {}; expected {STATE_VERSION}",
2082 self.version
2083 );
2084 }
2085 for (id, session) in &self.sessions {
2086 session.validate(id)?;
2087 }
2088 for child_id in self.subagents.keys() {
2089 self.validate_subagent(child_id)?;
2090 }
2091 for (host, sources) in &self.mount_history {
2092 if host.trim().is_empty() {
2093 bail!("mount history contains an empty host key");
2094 }
2095 if sources.iter().any(|source| !source.is_absolute()) {
2096 bail!("mount history for {host:?} contains a non-absolute source path");
2097 }
2098 }
2099 for (host, size) in &self.container_sizes {
2100 if host.trim().is_empty() {
2101 bail!("container size history contains an empty host key");
2102 }
2103 if size.cpus == 0 || size.memory_bytes == 0 {
2104 bail!("container size history for {host:?} contains a zero value");
2105 }
2106 if size.cpus > i64::MAX as u64 || size.memory_bytes > i64::MAX as u64 {
2107 bail!("container size history for {host:?} exceeds SQLite integer range");
2108 }
2109 }
2110 Ok(())
2111 }
2112
2113 pub fn validate_subagent(&self, child_id: &str) -> Result<()> {
2115 let Some(subagent) = self.subagents.get(child_id) else {
2116 return Ok(());
2117 };
2118 if child_id != subagent.child_session_id {
2119 bail!("sub-agent key {child_id:?} does not match its child session id");
2120 }
2121 if child_id == subagent.parent_session_id {
2122 bail!("sub-agent {child_id:?} cannot be its own parent");
2123 }
2124 if !self.sessions.contains_key(child_id) {
2125 bail!("sub-agent {child_id:?} has no child session");
2126 }
2127 if !self.sessions.contains_key(&subagent.parent_session_id) {
2128 bail!(
2129 "sub-agent {child_id:?} has unknown parent {:?}",
2130 subagent.parent_session_id
2131 );
2132 }
2133 if self.subagents.contains_key(&subagent.parent_session_id) {
2134 bail!("sub-agent {child_id:?} cannot belong to another sub-agent");
2135 }
2136 if subagent.task_name.trim().is_empty()
2137 || subagent.profile_id.trim().is_empty()
2138 || subagent.request_key.trim().is_empty()
2139 {
2140 bail!("sub-agent {child_id:?} has incomplete relationship metadata");
2141 }
2142 Ok(())
2143 }
2144
2145 pub fn remember_mount_sources(&mut self, host: &str, mounts: &[AdditionalMount]) {
2146 if mounts.is_empty() {
2147 return;
2148 }
2149 let sources = self.mount_history.entry(host.to_owned()).or_default();
2150 for mount in mounts.iter().rev() {
2151 sources.retain(|source| source != &mount.source);
2152 sources.insert(0, mount.source.clone());
2153 }
2154 sources.truncate(20);
2155 }
2156
2157 pub fn remember_container_size(&mut self, host: &str, size: HostContainerSize) {
2158 self.container_sizes.insert(host.to_owned(), size);
2159 }
2160
2161 pub fn project_directories(&self, host: &str) -> &[PathBuf] {
2162 self.mount_history
2163 .get(&project_history_key(host))
2164 .map(Vec::as_slice)
2165 .unwrap_or_default()
2166 }
2167
2168 pub fn remember_project_directory(&mut self, host: &str, directory: &Path) {
2169 let key = project_history_key(host);
2170 let directories = self.mount_history.entry(key).or_default();
2171 directories.retain(|existing| existing != directory);
2172 directories.insert(0, directory.to_path_buf());
2173 directories.truncate(20);
2174 }
2175
2176 pub fn destroy_stopped_session(&mut self, session_id: &str) -> Result<SessionRecord> {
2177 let session = self
2178 .sessions
2179 .get(session_id)
2180 .with_context(|| format!("unknown session {session_id}"))?;
2181 if session.state.is_active() {
2182 bail!("refusing to destroy active session {session_id}");
2183 }
2184 Ok(self
2185 .sessions
2186 .remove(session_id)
2187 .expect("session checked above"))
2188 }
2189
2190 pub fn destroy_session_force(&mut self, session_id: &str) -> Result<SessionRecord> {
2196 self.sessions
2197 .get(session_id)
2198 .with_context(|| format!("unknown session {session_id}"))?;
2199 Ok(self
2200 .sessions
2201 .remove(session_id)
2202 .expect("session checked above"))
2203 }
2204
2205 pub fn bundle_users(&self, bundle_id: &str) -> Vec<&SessionRecord> {
2210 self.sessions
2211 .values()
2212 .filter(|session| {
2213 session.bundle_id == bundle_id
2214 && session.project_directory.is_none()
2215 && session.state != SessionState::DestroyedWithDataLoss
2216 })
2217 .collect()
2218 }
2219
2220 pub fn bundle_removal_refusal(&self, bundle_id: &str) -> Option<String> {
2223 let users = self.bundle_users(bundle_id);
2224 if users.is_empty() {
2225 return None;
2226 }
2227 let mut names = users
2228 .iter()
2229 .take(3)
2230 .map(|session| format!("{:?}", session.listed_title()))
2231 .collect::<Vec<_>>();
2232 if users.len() > 3 {
2233 names.push(format!("{} more", users.len() - 3));
2234 }
2235 Some(format!(
2236 "Project {bundle_id:?} is used by {}: {}. Destroy those sessions before removing it.",
2237 if users.len() == 1 {
2238 "a session"
2239 } else {
2240 "sessions"
2241 },
2242 names.join(", ")
2243 ))
2244 }
2245
2246 pub fn validate_setup_update(&self, before: &Config, after: &Config) -> Result<()> {
2249 for session in self
2250 .sessions
2251 .values()
2252 .filter(|session| session.state.is_active())
2253 {
2254 let protected = if let Some(profile) = before.profiles.get(&session.last_profile) {
2255 let mut comparable = profile.clone();
2256 if let Some(updated) = after.profiles.get(&session.last_profile) {
2257 comparable.enabled = updated.enabled;
2258 comparable.subagents = updated.subagents.clone();
2259 }
2260 profile.kind == session.harness_kind
2262 && after.profiles.get(&session.last_profile) != Some(&comparable)
2263 } else {
2264 false
2265 };
2266 let bundle_changed = session.project_directory.is_none()
2267 && before
2268 .bundles
2269 .get(&session.bundle_id)
2270 .is_some_and(|bundle| after.bundles.get(&session.bundle_id) != Some(bundle));
2271 let target_changed =
2275 before
2276 .targets
2277 .get(&session.target_template_id)
2278 .is_some_and(|target| {
2279 after
2280 .targets
2281 .get(&session.target_template_id)
2282 .map(TargetTemplate::without_launch_only_settings)
2283 != Some(target.without_launch_only_settings())
2284 });
2285 if protected || bundle_changed || target_changed {
2286 let mut used = Vec::new();
2290 if protected {
2291 used.push(format!("agent profile {:?}", session.last_profile));
2292 }
2293 if bundle_changed {
2294 used.push(format!("project {:?}", session.project_name(before)));
2295 }
2296 if target_changed {
2297 used.push(format!("runtime {:?}", session.target_template_id));
2298 }
2299 let used = match used.as_slice() {
2300 [only] => only.clone(),
2301 [rest @ .., last] => format!("{} and {last}", rest.join(", ")),
2302 [] => unreachable!("something changed"),
2303 };
2304 let title = session.display_title();
2305 let named = if title == session.id {
2306 format!(
2307 "a running session in project {:?}",
2308 session.project_name(before)
2309 )
2310 } else {
2311 format!("the running session {title:?}")
2312 };
2313 bail!(
2314 "Setup would change the {used} that {named} uses. Save the new settings under a new name, or stop the session first."
2315 );
2316 }
2317 }
2318 Ok(())
2319 }
2320
2321 pub fn validate_against_config(&self, config: &Config) -> Result<()> {
2323 self.validate()?;
2324 config.validate()?;
2325 for session in self.sessions.values() {
2326 session.validate_configuration(config)?;
2327 }
2328 Ok(())
2329 }
2330}
2331
2332fn project_history_key(host: &str) -> String {
2333 format!("project:{host}")
2334}
2335
2336pub fn new_session_id() -> Result<String> {
2338 let mut random = [0u8; 16];
2339 getrandom::fill(&mut random)
2340 .map_err(|error| anyhow::anyhow!("generate Mjolnir session id: {error}"))?;
2341 Ok(crate::hex::lower_hex(random))
2342}
2343
2344pub fn harness_session_title(events: &[SequencedEvent]) -> Option<String> {
2346 events.iter().rev().find_map(|event| {
2347 let WorkerEvent::Adapter { payload, .. } = &event.event else {
2348 return None;
2349 };
2350 let crate::acp::RuntimeEvent::SessionUpdate { update } =
2351 serde_json::from_value(payload.clone()).ok()?
2352 else {
2353 return None;
2354 };
2355 let kind = update
2356 .get("sessionUpdate")
2357 .and_then(serde_json::Value::as_str)?;
2358 let title = match kind {
2359 "session_info_update" | "session_title" => {
2360 update.get("title").and_then(serde_json::Value::as_str)
2361 }
2362 _ => None,
2363 }?;
2364 normalize_session_title(title)
2365 })
2366}
2367
2368pub fn default_session_title(
2373 project_directory: Option<&Path>,
2374 bundle_id: &str,
2375 profile_id: &str,
2376) -> String {
2377 let project = project_directory.and_then(Path::file_name).map_or_else(
2378 || bundle_id.to_owned(),
2379 |name| name.to_string_lossy().into_owned(),
2380 );
2381 format!("{project} via {profile_id}")
2382}
2383
2384pub const MAX_SESSION_TITLE_CHARS: usize = 256;
2385
2386pub fn normalize_session_title(title: &str) -> Option<String> {
2388 let normalized = crate::relay::strip_hidden_prompt_context(title)
2389 .split_whitespace()
2390 .collect::<Vec<_>>()
2391 .join(" ");
2392 (!normalized.is_empty()).then(|| truncate_session_title(normalized, MAX_SESSION_TITLE_CHARS))
2393}
2394
2395fn truncate_session_title(title: String, maximum_chars: usize) -> String {
2396 if title.chars().count() <= maximum_chars {
2397 return title;
2398 }
2399
2400 let mut truncated = title.chars().take(maximum_chars - 1).collect::<String>();
2401 if let Some(boundary) = truncated.rfind(char::is_whitespace) {
2402 truncated.truncate(boundary);
2403 }
2404 truncated.push('…');
2405 truncated
2406}
2407
2408pub fn provisional_session_title(prompt: &str) -> Option<String> {
2414 const MAX_TITLE_CHARS: usize = 64;
2415
2416 let normalized = normalize_session_title(prompt)?;
2417 Some(truncate_session_title(normalized, MAX_TITLE_CHARS))
2418}
2419
2420pub fn short_id(id: &str) -> &str {
2421 id.get(..8).unwrap_or(id)
2422}
2423
2424#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
2425pub struct RecoveryCandidate {
2426 pub session_id: String,
2427 pub target_template_id: String,
2428 pub locator: TargetLocator,
2429 pub ownership: Option<crate::worker_launch::WorkerOwnership>,
2430 #[serde(default)]
2433 pub instance_id: Option<String>,
2434 #[serde(default, skip_serializing_if = "Option::is_none")]
2439 pub tracked_session: Option<SessionState>,
2440}
2441
2442#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
2443pub struct RecoveryScan {
2444 pub candidates: Vec<RecoveryCandidate>,
2445 pub warnings: Vec<String>,
2446 #[serde(default)]
2448 pub instance_id: String,
2449 #[serde(default)]
2452 pub hidden_other_instances: usize,
2453}
2454
2455#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2456#[serde(deny_unknown_fields)]
2457pub struct ResumeRepositorySourceReceipt {
2458 pub session_id: String,
2459 pub bundle_id: String,
2460 pub checkpoint_sha256: String,
2461 pub repositories: Vec<crate::config::ProjectRepository>,
2462}
2463
2464#[cfg(test)]
2465mod tests;