Skip to main content

mj_core/
state.rs

1//! Durable controller-side state for Hel-managed sessions.
2
3use std::collections::{BTreeMap, BTreeSet};
4use std::path::{Component, Path, PathBuf};
5use std::sync::Arc;
6
7use anyhow::{Context, Result, bail};
8use serde::{Deserialize, Serialize};
9
10use crate::config::{Config, HarnessKind, ProjectRepository, TargetTemplate, validate_id};
11use crate::credentials::CredentialSyncSignal;
12use crate::relay::{
13    RELAY_EVENT_GENESIS_DIGEST, RelayOperationalState, SequencedEvent, WorkerEvent,
14};
15use crate::snapshot_map::SnapshotMap;
16use crate::subagent::SubagentRecord;
17use crate::targets::{AdditionalMount, validate_additional_mounts};
18
19pub const STATE_VERSION: u32 = 1;
20
21mod target_runtime;
22pub use target_runtime::{TargetConnection, TargetRuntimeSettings};
23
24mod session_configuration;
25pub use session_configuration::SessionConfiguration;
26
27mod session_move;
28pub use session_move::*;
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
31#[serde(rename_all = "kebab-case")]
32pub enum SessionState {
33    Provisioning,
34    Running,
35    Disconnected,
36    Checkpointing,
37    Closing,
38    Destroying,
39    /// Checkpointed and torn down. Persisted as `"archived"` before the verb
40    /// was renamed, so the alias keeps those records loading.
41    #[serde(alias = "archived")]
42    Stopped,
43    /// A Mjolnir sub-agent whose turn ended and whose parent was told: its
44    /// worker process tree is stopped so it holds no processes in the
45    /// parent's container, while its record, relation, target locator and
46    /// worker root (relay journal, native session id) stay. Only a parent's
47    /// `send_input` starts it again. Nothing that connects to, reconnects,
48    /// recovers or upgrades live sessions acts on it.
49    Parked,
50    Lost,
51    Error,
52    DestroyedWithDataLoss,
53}
54
55/// A lifecycle transition temporarily replaces the conversation in control surfaces.
56/// Operation ownership takes precedence over intermediate durable session states.
57#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
58#[serde(rename_all = "kebab-case")]
59pub enum SessionTransitionKind {
60    Starting,
61    Resuming,
62    Moving,
63    Suspending,
64    Destroying,
65    /// A sub-agent stopped because its parent is being suspended.
66    Stopping,
67}
68
69impl SessionTransitionKind {
70    pub const fn label(self) -> &'static str {
71        match self {
72            Self::Starting => "Starting",
73            Self::Resuming => "Resuming",
74            Self::Moving => "Moving",
75            Self::Suspending => "Suspending",
76            Self::Destroying => "Destroying",
77            Self::Stopping => "Stopping",
78        }
79    }
80
81    pub fn for_session(state: SessionState, operation: Option<Self>) -> Option<Self> {
82        operation.or_else(|| state.transition_kind())
83    }
84}
85
86#[cfg(test)]
87mod transition_tests {
88    use super::{SessionState, SessionTransitionKind};
89
90    #[test]
91    fn operation_ownership_hides_intermediate_move_states_but_not_ordinary_live_work() {
92        for state in [
93            SessionState::Stopped,
94            SessionState::Running,
95            SessionState::Disconnected,
96        ] {
97            assert_eq!(
98                SessionTransitionKind::for_session(state, Some(SessionTransitionKind::Moving)),
99                Some(SessionTransitionKind::Moving)
100            );
101            assert_eq!(SessionTransitionKind::for_session(state, None), None);
102        }
103        assert_eq!(SessionState::Checkpointing.transition_kind(), None);
104        assert_eq!(
105            SessionState::Closing.transition_kind(),
106            Some(SessionTransitionKind::Suspending)
107        );
108    }
109}
110
111/// Controller-owned execution state derived from the relay event stream.
112#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
113#[serde(tag = "state", rename_all = "snake_case")]
114pub enum MaterializedExecutionState {
115    #[default]
116    Idle,
117    Running {
118        started_at_ms: i64,
119    },
120    Closing,
121    Closed,
122}
123
124pub use crate::transcript::{TerminalOutputRecord, TranscriptBody, TranscriptItem};
125
126/// What a durable queue entry does when its turn comes.
127///
128/// Serialized without a tag for prompts so entries written before configuration
129/// changes could be queued keep loading unchanged.
130#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
131#[serde(rename_all = "snake_case")]
132pub enum QueuedCommandKind {
133    #[default]
134    Prompt,
135    SetConfig {
136        key: String,
137        value: String,
138    },
139}
140
141impl QueuedCommandKind {
142    pub fn is_prompt(&self) -> bool {
143        matches!(self, Self::Prompt)
144    }
145}
146
147/// The composer form of a configuration change, used both as the queue entry's
148/// display text and as the text peeled back into the composer for editing.
149pub fn config_command_text(key: &str, value: &str) -> String {
150    if key == "fast-mode" {
151        "/fast".to_owned()
152    } else {
153        format!("/{key} {value}")
154    }
155}
156
157#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
158#[serde(deny_unknown_fields)]
159pub struct MaterializedQueuedPrompt {
160    pub command_id: String,
161    #[serde(default, skip_serializing_if = "QueuedCommandKind::is_prompt")]
162    pub kind: QueuedCommandKind,
163    pub content: Vec<serde_json::Value>,
164    pub queued_at_ms: i64,
165    /// Relay acceptance ordinal of the `CommandQueued` event that created this
166    /// entry. It is the turn identity the API hands back to callers, so wait
167    /// can tell one queued prompt's outcome from another's.
168    #[serde(default, skip_serializing_if = "Option::is_none")]
169    pub accepted_ordinal: Option<u64>,
170}
171
172/// The prompt currently executing, recorded when its `CommandStarted` event is
173/// projected and cleared when the command completes.
174#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
175#[serde(deny_unknown_fields)]
176pub struct MaterializedTurn {
177    pub command_id: String,
178    #[serde(default, skip_serializing_if = "Option::is_none")]
179    pub accepted_ordinal: Option<u64>,
180    /// Ordinal of the `CommandStarted` event, which is also the transcript
181    /// position of the turn's first item.
182    pub turn_start_position: u64,
183    pub started_at_ms: i64,
184    /// The relay prompt still executing this turn after a steer moved the
185    /// turn to a queued prompt. That prompt's completion ends this turn.
186    #[serde(default, skip_serializing_if = "Option::is_none")]
187    pub steered_into: Option<String>,
188}
189
190impl MaterializedTurn {
191    /// Whether the ending of relay prompt `command_id` ends this turn.
192    pub fn belongs_to(&self, command_id: &str) -> bool {
193        self.command_id == command_id || self.steered_into.as_deref() == Some(command_id)
194    }
195}
196
197/// How a prompt ended.
198#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
199#[serde(tag = "kind", rename_all = "snake_case")]
200pub enum TurnOutcomeKind {
201    /// The harness finished the turn and reported this stop reason.
202    Completed { stop_reason: String },
203    /// The relay refused the command before it ran.
204    Rejected {
205        message: String,
206        #[serde(default, skip_serializing_if = "Option::is_none")]
207        reason: Option<crate::event_outcome::OutcomeReason>,
208    },
209    /// The command was interrupted after being accepted.
210    Interrupted {
211        message: String,
212        #[serde(default, skip_serializing_if = "Option::is_none")]
213        reason: Option<crate::event_outcome::OutcomeReason>,
214    },
215}
216
217/// How a turn ended, in words a person reads: "completed, end of turn",
218/// "interrupted", or "failed: <reason>".
219impl std::fmt::Display for TurnOutcomeKind {
220    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
221        use crate::event_outcome::{OutcomeReason, TurnResultKind};
222        let result = self.result();
223        match result.kind {
224            TurnResultKind::Completed => formatter.write_str("completed, end of turn"),
225            TurnResultKind::InputRequired => formatter.write_str("completed, waiting for input"),
226            TurnResultKind::Cancelled | TurnResultKind::Interrupted => {
227                formatter.write_str("interrupted")
228            }
229            TurnResultKind::Failed if result.reason == Some(OutcomeReason::QuotaLimit) => {
230                formatter.write_str("failed: quota limit reached")
231            }
232            TurnResultKind::Failed | TurnResultKind::Rejected => {
233                let message = result.message.as_deref().unwrap_or("unknown failure");
234                if result.stop_reason.is_some() {
235                    write!(formatter, "failed: {}", stop_reason_words(message))
236                } else {
237                    write!(
238                        formatter,
239                        "failed: {}",
240                        message.lines().next().unwrap_or_default().trim()
241                    )
242                }
243            }
244        }
245    }
246}
247
248/// A stop reason as the harness spells it (`MaxTokens`, `max_turn_requests`)
249/// as lower-case words.
250fn stop_reason_words(stop_reason: &str) -> String {
251    let mut words = String::new();
252    let mut previous_lower = false;
253    for character in stop_reason.trim().chars() {
254        if character == '_' || character == '-' || character.is_whitespace() {
255            if !words.ends_with(' ') && !words.is_empty() {
256                words.push(' ');
257            }
258            previous_lower = false;
259            continue;
260        }
261        if character.is_uppercase() && previous_lower {
262            words.push(' ');
263        }
264        previous_lower = character.is_lowercase() || character.is_ascii_digit();
265        words.extend(character.to_lowercase());
266    }
267    match words.trim_end() {
268        "" => "no reason given".to_owned(),
269        words => words.to_owned(),
270    }
271}
272
273#[derive(Debug, Clone, Copy, PartialEq, Eq)]
274pub enum PromptCompletion {
275    InputRequired,
276    Finished,
277    Cancelled,
278    QuotaLimit,
279    Error,
280}
281
282/// Shared interpretation for wait responses and durable completion events.
283pub fn classify_prompt_completion(stop_reason: &str) -> PromptCompletion {
284    let normalized = stop_reason
285        .chars()
286        .filter(|character| *character != '_' && *character != '-')
287        .flat_map(char::to_lowercase)
288        .collect::<String>();
289    match normalized.as_str() {
290        "endturn" => PromptCompletion::Finished,
291        "awaitinginput" => PromptCompletion::InputRequired,
292        "cancelled" | "canceled" => PromptCompletion::Cancelled,
293        "quotalimit" => PromptCompletion::QuotaLimit,
294        _ => PromptCompletion::Error,
295    }
296}
297
298/// The most recent finished prompt on a session.
299#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
300#[serde(deny_unknown_fields)]
301pub struct MaterializedTurnOutcome {
302    #[serde(default, skip_serializing_if = "Option::is_none")]
303    pub diagnostic: Option<crate::diagnostic::TurnDiagnostic>,
304
305    #[serde(default, skip_serializing_if = "Option::is_none")]
306    pub usage: Option<crate::usage::TokenUsage>,
307    pub command_id: String,
308    #[serde(default, skip_serializing_if = "Option::is_none")]
309    pub accepted_ordinal: Option<u64>,
310    #[serde(default, skip_serializing_if = "Option::is_none")]
311    pub turn_start_position: Option<u64>,
312    pub completed_ordinal: u64,
313    pub completed_at_ms: i64,
314    pub outcome: TurnOutcomeKind,
315}
316
317impl MaterializedTurnOutcome {
318    /// Only work that actually started can have been interrupted.
319    pub fn interruption_ordinal(&self) -> Option<u64> {
320        (self.turn_start_position.is_some()
321            && matches!(self.outcome, TurnOutcomeKind::Interrupted { .. }))
322        .then_some(self.completed_ordinal)
323    }
324}
325
326/// Canonical controller projection for one logical ACP session.
327#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
328#[serde(deny_unknown_fields)]
329pub struct MaterializedSession {
330    pub session_id: String,
331    pub applied_event_ordinal: u64,
332    pub applied_event_digest: String,
333    /// Monotonic controller projection watermark derived from relay event
334    /// receipt times. It is deliberately independent of retained rows.
335    pub last_activity_at_ms: Option<i64>,
336    pub execution: MaterializedExecutionState,
337    #[serde(default, skip_serializing_if = "Option::is_none")]
338    pub session_title: Option<String>,
339    #[serde(default, skip_serializing_if = "SessionConfiguration::is_empty")]
340    pub configuration: SessionConfiguration,
341    #[serde(default, skip_serializing_if = "Vec::is_empty")]
342    /// Transcript items are shared by pointer so cloning a snapshot copies
343    /// handles rather than the whole conversation.
344    pub transcript: Vec<Arc<TranscriptItem>>,
345    #[serde(default, skip_serializing_if = "Vec::is_empty")]
346    pub queued_prompts: Vec<MaterializedQueuedPrompt>,
347    /// In-flight form requests are projected durably, but their answers are
348    /// connection-only and never enter this state.
349    #[serde(default, skip_serializing_if = "Vec::is_empty")]
350    pub pending_elicitations: Vec<crate::elicitation::ElicitationRequest>,
351    /// The prompt running right now, if any.
352    #[serde(default, skip_serializing_if = "Option::is_none")]
353    pub active_turn: Option<MaterializedTurn>,
354    /// The most recently finished prompt, kept after the session stops so a
355    /// caller can still read how the last turn ended.
356    #[serde(default, skip_serializing_if = "Option::is_none")]
357    pub last_turn_outcome: Option<MaterializedTurnOutcome>,
358}
359
360/// The small portion of a durable projection needed to populate dashboard
361/// rows before the live session delivers its full transcript snapshot.
362#[derive(Debug, Clone, PartialEq, Eq)]
363pub struct MaterializedSessionSummary {
364    pub session_id: String,
365    pub applied_event_ordinal: u64,
366    pub last_activity_at_ms: Option<i64>,
367    pub execution: MaterializedExecutionState,
368    pub session_title: Option<String>,
369    pub last_agent_message: Option<String>,
370    pub last_user_message: Option<String>,
371    /// Whether the last nonempty agent message appears after the last
372    /// nonempty user message in transcript order.
373    pub last_agent_message_follows_last_user: bool,
374    pub agent_message_latest_content_ordinals: Vec<u64>,
375    pub interruption_event_ordinals: Vec<u64>,
376}
377
378impl MaterializedSession {
379    pub fn empty(session_id: impl Into<String>) -> Self {
380        Self {
381            session_id: session_id.into(),
382            applied_event_ordinal: 0,
383            applied_event_digest: RELAY_EVENT_GENESIS_DIGEST.into(),
384            last_activity_at_ms: None,
385            execution: MaterializedExecutionState::Idle,
386            session_title: None,
387            configuration: SessionConfiguration::default(),
388            transcript: Vec::new(),
389            queued_prompts: Vec::new(),
390            pending_elicitations: Vec::new(),
391            active_turn: None,
392            last_turn_outcome: None,
393        }
394    }
395
396    pub fn last_activity_at_ms(&self) -> Option<i64> {
397        self.last_activity_at_ms
398    }
399
400    /// Resolve the title exposed by a live materialized session.
401    ///
402    /// Sessions created before provisional titles were projected can still
403    /// have an untitled transcript. Derive the same bounded fallback from
404    /// their first visible user prompt when reading them.
405    pub fn resolved_title(&self) -> Option<String> {
406        self.session_title
407            .as_deref()
408            .and_then(normalize_session_title)
409            .or_else(|| {
410                self.transcript.iter().find_map(|item| {
411                    let TranscriptBody::User { content } = &item.body else {
412                        return None;
413                    };
414                    provisional_session_title(&crate::transcript::materialized_content_text(
415                        content,
416                    ))
417                })
418            })
419            .or_else(|| {
420                self.queued_prompts
421                    .iter()
422                    .filter(|prompt| prompt.kind.is_prompt())
423                    .find_map(|prompt| {
424                        provisional_session_title(&crate::transcript::materialized_content_text(
425                            &prompt.content,
426                        ))
427                    })
428            })
429    }
430
431    pub fn unread_agent_messages_after(&self, viewed_through_event_ordinal: u64) -> u64 {
432        self.transcript
433            .iter()
434            .filter(|item| {
435                item.latest_content_event_ordinal
436                    .is_some_and(|ordinal| ordinal > viewed_through_event_ordinal)
437                    && item.is_nonempty_agent_message()
438            })
439            .count() as u64
440    }
441
442    pub fn unread_interruptions_after(&self, viewed_through_event_ordinal: u64) -> u64 {
443        self.interruption_event_ordinals()
444            .into_iter()
445            .filter(|ordinal| *ordinal > viewed_through_event_ordinal)
446            .count() as u64
447    }
448
449    pub fn interruption_event_ordinals(&self) -> Vec<u64> {
450        let mut ordinals = self
451            .transcript
452            .iter()
453            .filter(|item| item.is_work_interruption())
454            .map(|item| item.position)
455            .collect::<Vec<_>>();
456        if let Some(ordinal) = self
457            .last_turn_outcome
458            .as_ref()
459            .and_then(MaterializedTurnOutcome::interruption_ordinal)
460        {
461            ordinals.push(ordinal);
462        }
463        ordinals.sort_unstable();
464        ordinals.dedup();
465        ordinals
466    }
467
468    pub fn validate(&self) -> Result<()> {
469        validate_id("session", &self.session_id)?;
470        validate_relay_event_frontier(
471            self.applied_event_ordinal,
472            &self.applied_event_digest,
473            "materialized session event frontier",
474        )?;
475        if self
476            .session_title
477            .as_ref()
478            .is_some_and(|title| title.trim().is_empty())
479        {
480            bail!("materialized session has an empty title");
481        }
482        let mut item_ids = BTreeSet::new();
483        for item in &self.transcript {
484            item.validate(self.applied_event_ordinal)?;
485            if !item_ids.insert(item.stable_id.as_str()) {
486                bail!(
487                    "materialized transcript contains duplicate item {:?}",
488                    item.stable_id
489                );
490            }
491        }
492        let mut command_ids = BTreeSet::new();
493        for prompt in &self.queued_prompts {
494            if prompt.command_id.trim().is_empty() {
495                bail!("materialized prompt queue has an empty command id");
496            }
497            if !command_ids.insert(prompt.command_id.as_str()) {
498                bail!(
499                    "materialized prompt queue contains duplicate command {:?}",
500                    prompt.command_id
501                );
502            }
503            if let QueuedCommandKind::SetConfig { key, value } = &prompt.kind
504                && (key.trim().is_empty() || value.trim().is_empty())
505            {
506                bail!(
507                    "materialized queued configuration change {:?} is incomplete",
508                    prompt.command_id
509                );
510            }
511        }
512        Ok(())
513    }
514}
515
516/// A materialized session paired with the live worker's relay state. The
517/// session manager hands this to every reader that needs both the durable
518/// projection and the connection's operational status.
519#[derive(Debug, Clone, PartialEq)]
520pub struct ManagedSessionSnapshot {
521    pub materialized: MaterializedSession,
522    /// What `materialized.transcript` leaves out, and the facts that live
523    /// there. See [`ProjectionWindow`].
524    pub window: ProjectionWindow,
525    pub operational: RelayOperationalState,
526    /// Newest relay event observed by this live actor that asks for immediate
527    /// credential reconciliation. This is intentionally ephemeral: it avoids
528    /// retaining raw replay pages or rescanning projected history.
529    pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
530    /// Content address of the executable the connected worker is running, as
531    /// it reported in hello. `None` when the connection did not come from a
532    /// live worker or the worker predates the field; either way the worker is
533    /// not known to be the build this controller would install.
534    pub worker_build: Option<String>,
535    /// Pending parent-tool work fetched from the target worker.
536    pub subagent_requests: Vec<crate::subagent::SubagentToolRequest>,
537    /// Recently completed tool work cached by the worker for idempotent calls.
538    pub subagent_results: Vec<crate::subagent::SubagentToolResult>,
539}
540
541/// What a projection's transcript window leaves out.
542///
543/// A polled projection carries only the end of the transcript, because that is
544/// all any viewer shows and loading the rest is work proportional to history.
545/// Two facts a reader needs live outside that window: the provisional title
546/// comes from the *first* user message, and the newest turn start is outside
547/// it whenever a single turn is longer than the window. Both are read
548/// separately, with one indexed query each, rather than found by scanning.
549///
550/// A complete projection answers both by scanning what it already holds, which
551/// is what [`ProjectionWindow::of`] does.
552#[derive(Debug, Clone, PartialEq, Eq)]
553pub struct ProjectionWindow {
554    /// Transcript items before the window. Zero when the projection is whole.
555    pub omitted_items: usize,
556    /// The title derived from the first user message.
557    pub provisional_title: Option<String>,
558    /// Position of the newest turn start — a user message or the marker for a
559    /// turn the harness began on its own — whether or not it is in the
560    /// window. `None` when the session has none.
561    pub latest_turn_start_position: Option<u64>,
562}
563
564impl ProjectionWindow {
565    /// Keep complete turns around the tail target. Unsettled content can still
566    /// change after a newer turn starts, so retain its turn as well.
567    pub fn trim(&mut self, session: &mut MaterializedSession, target: usize) {
568        let observed = Self::of(session);
569        if self.provisional_title.is_none() {
570            self.provisional_title = observed.provisional_title;
571        }
572        self.latest_turn_start_position = observed
573            .latest_turn_start_position
574            .or(self.latest_turn_start_position);
575        let mut boundary = session.transcript.len().saturating_sub(target.max(1));
576        for (index, item) in session.transcript.iter().enumerate() {
577            let mutable = match &item.body {
578                TranscriptBody::Agent { streaming, .. }
579                | TranscriptBody::Thought { streaming, .. } => *streaming,
580                TranscriptBody::Tool { call, .. } => matches!(
581                    call.get("status").and_then(serde_json::Value::as_str),
582                    Some("pending" | "in_progress")
583                ),
584                _ => false,
585            };
586            if mutable || Some(item.position) == self.latest_turn_start_position {
587                boundary = boundary.min(index);
588            }
589        }
590        let cut = session
591            .transcript
592            .iter()
593            .take(boundary + 1)
594            .rposition(|item| item.is_turn_start())
595            .unwrap_or(0);
596        if cut > 0 {
597            session.transcript.drain(..cut);
598            self.omitted_items += cut;
599        }
600    }
601
602    /// The window of a projection that omits nothing.
603    #[must_use]
604    pub fn of(session: &MaterializedSession) -> Self {
605        Self {
606            omitted_items: 0,
607            provisional_title: session.transcript.iter().find_map(|item| {
608                let TranscriptBody::User { content } = &item.body else {
609                    return None;
610                };
611                provisional_session_title(&crate::transcript::materialized_content_text(content))
612            }),
613            latest_turn_start_position: session
614                .transcript
615                .iter()
616                .rev()
617                .find(|item| item.is_turn_start())
618                .map(|item| item.position),
619        }
620    }
621}
622
623impl ManagedSessionSnapshot {
624    /// The session's title, using the same precedence as
625    /// [`MaterializedSession::resolved_title`] but taking the provisional
626    /// title from the window rather than from a transcript head that a polled
627    /// projection does not carry.
628    #[must_use]
629    pub fn resolved_title(&self) -> Option<String> {
630        self.materialized
631            .session_title
632            .as_deref()
633            .and_then(normalize_session_title)
634            .or_else(|| self.window.provisional_title.clone())
635            .or_else(|| {
636                self.materialized
637                    .queued_prompts
638                    .iter()
639                    .filter(|prompt| prompt.kind.is_prompt())
640                    .find_map(|prompt| {
641                        provisional_session_title(&crate::transcript::materialized_content_text(
642                            &prompt.content,
643                        ))
644                    })
645            })
646    }
647
648    /// The position of the turn this session most recently finished, or `None`
649    /// while it is still working. Same answer as
650    /// [`latest_completed_turn_ordinal`], from a position the window carries
651    /// rather than a scan back through the transcript.
652    #[must_use]
653    pub fn latest_completed_turn_ordinal(&self) -> Option<u64> {
654        if self.materialized.execution != MaterializedExecutionState::Idle {
655            return None;
656        }
657        self.window.latest_turn_start_position
658    }
659}
660
661/// One session's activity, reported to the recovery coordinator.
662#[derive(Debug, Clone)]
663pub struct RecoveryObservation {
664    pub session: SessionRecord,
665    pub config: Config,
666    pub latest_completed_turn_ordinal: Option<u64>,
667    /// Why a routine checkpoint has to wait, or `None` when one may start now:
668    /// [`crate::activity::routine_checkpoint_wait`] on the worker's
669    /// operational state, the same answer checkpoint admission gives.
670    pub checkpoint_wait: Option<crate::activity::CheckpointWait>,
671}
672
673/// The position where the session's most recent finished turn began, or
674/// `None` while it is still working. A turn starts at a user message or at the
675/// marker for a turn the harness began on its own, so autonomous work is
676/// covered once it settles.
677pub fn latest_completed_turn_ordinal(session: &MaterializedSession) -> Option<u64> {
678    if session.execution != MaterializedExecutionState::Idle {
679        return None;
680    }
681    session
682        .transcript
683        .iter()
684        .rev()
685        .find(|item| item.is_turn_start())
686        .map(|item| item.position)
687}
688
689pub fn validate_relay_event_digest(digest: &str, name: &str) -> Result<()> {
690    if digest.len() != 64
691        || !digest
692            .bytes()
693            .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
694    {
695        bail!("{name} must be a lowercase SHA-256 digest");
696    }
697    Ok(())
698}
699
700pub fn validate_relay_event_frontier(ordinal: u64, digest: &str, name: &str) -> Result<()> {
701    validate_relay_event_digest(digest, name)?;
702    if (ordinal == 0) != (digest == RELAY_EVENT_GENESIS_DIGEST) {
703        bail!("{name} has inconsistent ordinal {ordinal} and digest {digest}");
704    }
705    Ok(())
706}
707
708fn is_false(value: &bool) -> bool {
709    !*value
710}
711
712impl SessionState {
713    /// The persisted and wire spelling, matching the serde encoding.
714    pub const fn as_str(self) -> &'static str {
715        match self {
716            Self::Provisioning => "provisioning",
717            Self::Running => "running",
718            Self::Disconnected => "disconnected",
719            Self::Checkpointing => "checkpointing",
720            Self::Closing => "closing",
721            Self::Destroying => "destroying",
722            Self::Stopped => "stopped",
723            Self::Parked => "parked",
724            Self::Lost => "lost",
725            Self::Error => "error",
726            Self::DestroyedWithDataLoss => "destroyed-with-data-loss",
727        }
728    }
729
730    /// Read a stored spelling. Rows written before the verb was renamed still
731    /// say `"archived"`.
732    pub fn from_stored(value: &str) -> Option<Self> {
733        Some(match value {
734            "provisioning" => Self::Provisioning,
735            "running" => Self::Running,
736            "disconnected" => Self::Disconnected,
737            "checkpointing" => Self::Checkpointing,
738            "closing" => Self::Closing,
739            "destroying" => Self::Destroying,
740            "stopped" | "archived" => Self::Stopped,
741            "parked" => Self::Parked,
742            "lost" => Self::Lost,
743            "error" => Self::Error,
744            "destroyed-with-data-loss" => Self::DestroyedWithDataLoss,
745            _ => return None,
746        })
747    }
748
749    /// Recovery without a live operation still hides an unfinished target transition.
750    /// Ordinary checkpoints and reconnects deliberately keep their conversation visible.
751    pub const fn transition_kind(self) -> Option<SessionTransitionKind> {
752        match self {
753            Self::Provisioning => Some(SessionTransitionKind::Starting),
754            Self::Closing => Some(SessionTransitionKind::Suspending),
755            Self::Destroying => Some(SessionTransitionKind::Destroying),
756            _ => None,
757        }
758    }
759
760    /// True while the session still belongs on the dashboard. `Closing` and
761    /// `Checkpointing` stay active on purpose: a stop that has not produced a
762    /// verified checkpoint must not make its row disappear. A `Parked`
763    /// sub-agent is active too: it is still its parent's child, still listed,
764    /// and its parent's suspend, destroy or workspace close must still end
765    /// it. Code that needs a live worker must ask [`Self::has_live_worker`].
766    pub const fn is_active(self) -> bool {
767        matches!(
768            self,
769            Self::Provisioning
770                | Self::Running
771                | Self::Disconnected
772                | Self::Checkpointing
773                | Self::Closing
774                | Self::Destroying
775                | Self::Parked
776                | Self::Error
777        )
778    }
779
780    /// True while the session may have a worker process tree on its target,
781    /// including one still being started or torn down: the states in which a
782    /// sub-agent counts against its parent's cap.
783    pub const fn has_live_worker(self) -> bool {
784        matches!(
785            self,
786            Self::Provisioning
787                | Self::Running
788                | Self::Disconnected
789                | Self::Checkpointing
790                | Self::Closing
791                | Self::Destroying
792        )
793    }
794}
795
796#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
797#[serde(tag = "kind", rename_all = "kebab-case")]
798pub enum PodmanWorkspaceLocator {
799    #[default]
800    ContainerLayer,
801    Volume {
802        name: String,
803    },
804    HostPath {
805        path: PathBuf,
806        helper: Vec<String>,
807        resource: String,
808    },
809}
810
811#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
812#[serde(tag = "kind", rename_all = "kebab-case")]
813pub enum TargetLocator {
814    LocalBare {
815        worker_root: PathBuf,
816    },
817    LocalPodman {
818        container_id: String,
819        #[serde(default)]
820        workspace_storage: PodmanWorkspaceLocator,
821        /// The session that owns the container when this locator is a
822        /// sub-agent child borrowing its parent's container; `None` when the
823        /// session owns the container itself.
824        #[serde(default, skip_serializing_if = "Option::is_none")]
825        borrowed_from: Option<String>,
826    },
827    LocalDocker {
828        container_id: String,
829        /// The session that owns the container when this locator is a
830        /// sub-agent child borrowing its parent's container; `None` when the
831        /// session owns the container itself.
832        #[serde(default, skip_serializing_if = "Option::is_none")]
833        borrowed_from: Option<String>,
834    },
835    AppleContainer {
836        container_id: String,
837        /// The session that owns the container when this locator is a
838        /// sub-agent child borrowing its parent's container; `None` when the
839        /// session owns the container itself.
840        #[serde(default, skip_serializing_if = "Option::is_none")]
841        borrowed_from: Option<String>,
842    },
843    AwsEc2 {
844        instance_id: String,
845        #[serde(default, skip_serializing_if = "Option::is_none")]
846        address: Option<String>,
847    },
848    SshBare {
849        host: String,
850        workspace: PathBuf,
851        #[serde(default, skip_serializing_if = "Option::is_none")]
852        worker_id: Option<String>,
853    },
854    SshPodman {
855        host: String,
856        container_id: String,
857        #[serde(default)]
858        workspace_storage: PodmanWorkspaceLocator,
859        /// The session that owns the container when this locator is a
860        /// sub-agent child borrowing its parent's container; `None` when the
861        /// session owns the container itself.
862        #[serde(default, skip_serializing_if = "Option::is_none")]
863        borrowed_from: Option<String>,
864    },
865    SshDocker {
866        host: String,
867        container_id: String,
868        /// The session that owns the container when this locator is a
869        /// sub-agent child borrowing its parent's container; `None` when the
870        /// session owns the container itself.
871        #[serde(default, skip_serializing_if = "Option::is_none")]
872        borrowed_from: Option<String>,
873    },
874}
875
876impl ManagedWorktreeTarget {
877    /// Whether `other` reaches the same checkout: the same kind and, over
878    /// SSH, the same destination, port, and login user.
879    ///
880    /// The other `ssh` options (keys, `ControlPath`, keepalives, host-key
881    /// policy) say how to connect, not where the worktree lives, so they
882    /// follow the machine's current configuration and never make a
883    /// suspended session unable to resume.
884    pub fn same_location(&self, other: &Self) -> bool {
885        match (self, other) {
886            (Self::Local, Self::Local) => true,
887            (
888                Self::Ssh {
889                    destination,
890                    ssh_args,
891                },
892                Self::Ssh {
893                    destination: other_destination,
894                    ssh_args: other_args,
895                },
896            ) => {
897                destination == other_destination
898                    && ssh_location_option(ssh_args, 'p', "port")
899                        == ssh_location_option(other_args, 'p', "port")
900                    && ssh_location_option(ssh_args, 'l', "user")
901                        == ssh_location_option(other_args, 'l', "user")
902            }
903            _ => false,
904        }
905    }
906}
907
908/// The value `ssh` would use for an option that has both a short flag
909/// (`-p 22`, `-p22`) and an `-o` spelling (`-o Port=22`, `-oPort 22`). OpenSSH
910/// keeps the first value it sees.
911fn ssh_location_option(args: &[String], flag: char, option: &str) -> Option<String> {
912    let mut args = args.iter();
913    while let Some(argument) = args.next() {
914        let Some(rest) = argument.strip_prefix('-') else {
915            continue;
916        };
917        let mut chars = rest.chars();
918        let Some(name) = chars.next() else { continue };
919        if name != flag && name != 'o' {
920            continue;
921        }
922        let inline = chars.as_str();
923        let value = if inline.is_empty() {
924            args.next().cloned()
925        } else {
926            Some(inline.to_owned())
927        };
928        if name == flag {
929            return value;
930        }
931        if let Some(setting) = value {
932            let (key, found) = setting
933                .split_once(['=', ' ', '\t'])
934                .unwrap_or((setting.as_str(), ""));
935            if key.trim().eq_ignore_ascii_case(option) {
936                return Some(found.trim().to_owned());
937            }
938        }
939    }
940    None
941}
942
943#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
944#[serde(tag = "kind", rename_all = "kebab-case")]
945pub enum ManagedWorktreeTarget {
946    Local,
947    Ssh {
948        destination: String,
949        #[serde(default, skip_serializing_if = "Vec::is_empty")]
950        ssh_args: Vec<String>,
951    },
952}
953
954/// Whether a selected project can create a session-owned Git checkout.
955#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
956#[serde(deny_unknown_fields)]
957pub struct ManagedWorktreeOptions {
958    pub available: bool,
959    pub default_create: bool,
960}
961
962#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
963#[serde(deny_unknown_fields)]
964pub struct ManagedWorktree {
965    /// Old records are linked worktrees. New isolated raw sessions own a clone.
966    #[serde(default, skip_serializing_if = "ManagedCheckoutKind::is_worktree")]
967    pub kind: ManagedCheckoutKind,
968    pub source_project_directory: PathBuf,
969    pub source_repository: PathBuf,
970    pub worktree_root: PathBuf,
971    pub branch: String,
972    pub target: ManagedWorktreeTarget,
973    /// The commit the session branch was created at. Recorded so an export can
974    /// diff against it in one read; sessions created before this field existed
975    /// fall back to the branch reflog, which expires.
976    #[serde(default, skip_serializing_if = "Option::is_none")]
977    pub base_commit: Option<String>,
978}
979
980#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
981#[serde(rename_all = "snake_case")]
982pub enum ManagedCheckoutKind {
983    #[default]
984    Worktree,
985    Clone,
986}
987
988impl ManagedCheckoutKind {
989    fn is_worktree(&self) -> bool {
990        matches!(self, Self::Worktree)
991    }
992}
993
994impl ManagedWorktree {
995    fn validate(&self, session_id: &str, project_directory: Option<&Path>) -> Result<()> {
996        for (label, path) in [
997            ("source project directory", &self.source_project_directory),
998            ("source repository", &self.source_repository),
999            ("worktree root", &self.worktree_root),
1000        ] {
1001            if !path.is_absolute() || path.components().any(|part| part == Component::ParentDir) {
1002                bail!("managed worktree {label} must be an absolute safe path");
1003            }
1004        }
1005        if !self
1006            .source_project_directory
1007            .starts_with(&self.source_repository)
1008        {
1009            bail!("managed worktree source directory is outside its repository");
1010        }
1011        let expected_root = self
1012            .source_repository
1013            .join(".mj")
1014            .join(match self.kind {
1015                ManagedCheckoutKind::Worktree => "worktrees",
1016                ManagedCheckoutKind::Clone => "clones",
1017            })
1018            .join(session_id);
1019        if self.worktree_root != expected_root {
1020            bail!("managed worktree root does not match the session-owned path");
1021        }
1022        if self.kind == ManagedCheckoutKind::Worktree && self.branch != format!("mj/{session_id}") {
1023            bail!("managed worktree branch does not match the session id");
1024        }
1025        if self.kind == ManagedCheckoutKind::Clone && self.branch.trim().is_empty() {
1026            bail!("managed clone has no starting branch");
1027        }
1028        let relative = self
1029            .source_project_directory
1030            .strip_prefix(&self.source_repository)
1031            .expect("source relationship checked above");
1032        if project_directory != Some(self.worktree_root.join(relative).as_path()) {
1033            bail!("session project directory does not match its managed worktree");
1034        }
1035        match &self.target {
1036            ManagedWorktreeTarget::Local => {}
1037            ManagedWorktreeTarget::Ssh { destination, .. } if destination.trim().is_empty() => {
1038                bail!("managed SSH worktree has an empty destination")
1039            }
1040            ManagedWorktreeTarget::Ssh { .. } => {}
1041        }
1042        Ok(())
1043    }
1044}
1045
1046#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1047#[serde(tag = "kind", rename_all = "kebab-case")]
1048pub enum SessionResourceAllocation {
1049    Container {
1050        cpus: u64,
1051        memory_bytes: u64,
1052    },
1053    AwsEc2 {
1054        instance_type: String,
1055        vcpus: u64,
1056        memory_bytes: u64,
1057    },
1058}
1059
1060impl SessionResourceAllocation {
1061    pub fn validate(&self) -> Result<()> {
1062        match self {
1063            Self::Container { cpus, memory_bytes } if *cpus == 0 || *memory_bytes == 0 => {
1064                bail!("container resource allocation must have non-zero CPU and memory")
1065            }
1066            Self::AwsEc2 {
1067                instance_type,
1068                vcpus,
1069                memory_bytes,
1070            } if instance_type.trim().is_empty() || *vcpus == 0 || *memory_bytes == 0 => {
1071                bail!("EC2 resource allocation must have an instance type, CPU, and memory")
1072            }
1073            _ => Ok(()),
1074        }
1075    }
1076}
1077
1078/// The CPU count an allocation grants, regardless of target kind.
1079pub fn allocation_cpus(allocation: &SessionResourceAllocation) -> u64 {
1080    match allocation {
1081        SessionResourceAllocation::Container { cpus, .. } => *cpus,
1082        SessionResourceAllocation::AwsEc2 { vcpus, .. } => *vcpus,
1083    }
1084}
1085
1086/// The memory, in bytes, an allocation grants, regardless of target kind.
1087pub fn allocation_memory(allocation: &SessionResourceAllocation) -> u64 {
1088    match allocation {
1089        SessionResourceAllocation::Container { memory_bytes, .. }
1090        | SessionResourceAllocation::AwsEc2 { memory_bytes, .. } => *memory_bytes,
1091    }
1092}
1093
1094impl TargetLocator {
1095    fn validate(&self, session_id: &str) -> Result<()> {
1096        match self {
1097            Self::LocalBare { worker_root } => {
1098                if !worker_root.is_absolute()
1099                    || worker_root
1100                        .components()
1101                        .any(|part| part == Component::ParentDir)
1102                    || !worker_root.ends_with(session_id)
1103                {
1104                    bail!(
1105                        "local bare worker root must be an absolute safe path ending in the session id"
1106                    );
1107                }
1108            }
1109            Self::LocalPodman { container_id, .. }
1110            | Self::LocalDocker { container_id, .. }
1111            | Self::AppleContainer { container_id, .. }
1112            | Self::SshPodman { container_id, .. }
1113            | Self::SshDocker { container_id, .. }
1114                if container_id.trim().is_empty() =>
1115            {
1116                bail!("target locator has an empty container id")
1117            }
1118            Self::AwsEc2 { instance_id, .. } if instance_id.trim().is_empty() => {
1119                bail!("target locator has an empty AWS instance id")
1120            }
1121            Self::SshBare {
1122                host, workspace, ..
1123            } => {
1124                if host.trim().is_empty() {
1125                    bail!("bare SSH target locator has an empty host");
1126                }
1127                if workspace.as_os_str().is_empty()
1128                    || workspace
1129                        .components()
1130                        .any(|part| part == Component::ParentDir)
1131                    || !workspace.ends_with(session_id)
1132                {
1133                    bail!("bare SSH target locator must be a safe path ending in the session id");
1134                }
1135            }
1136            Self::SshPodman { host, .. } if host.trim().is_empty() => {
1137                bail!("SSH Podman target locator has an empty host")
1138            }
1139            Self::SshDocker { host, .. } if host.trim().is_empty() => {
1140                bail!("SSH Docker target locator has an empty host")
1141            }
1142            _ => {}
1143        }
1144        Ok(())
1145    }
1146}
1147
1148#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1149#[serde(deny_unknown_fields)]
1150pub struct CheckpointMetadata {
1151    pub archive_path: PathBuf,
1152    /// Lowercase SHA-256 digest of the verified archive.
1153    pub sha256: String,
1154    pub created_at: String,
1155    pub event_frontier: u64,
1156}
1157
1158/// Whether every saved Git change has a verified durable copy outside mj.
1159#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1160#[serde(rename_all = "snake_case")]
1161pub enum PublicationState {
1162    Published,
1163    Unpublished,
1164    Unknown,
1165}
1166
1167/// Evidence for one exact checkpoint. A newer checkpoint invalidates it.
1168#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1169pub struct PublicationAssessment {
1170    pub checkpoint_sha256: String,
1171    pub state: PublicationState,
1172    pub dirty: bool,
1173    pub stashed: bool,
1174    pub saved_commits: Vec<String>,
1175    pub destinations: Vec<String>,
1176    pub checked_at: String,
1177    pub reason: Option<String>,
1178}
1179
1180impl CheckpointMetadata {
1181    fn validate(&self) -> Result<()> {
1182        if self.archive_path.as_os_str().is_empty() {
1183            bail!("checkpoint archive path is empty");
1184        }
1185        if self.sha256.len() != 64
1186            || !self
1187                .sha256
1188                .bytes()
1189                .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1190        {
1191            bail!("checkpoint SHA-256 must be 64 lowercase hexadecimal characters");
1192        }
1193        if self.created_at.trim().is_empty() {
1194            bail!("checkpoint timestamp is empty");
1195        }
1196        Ok(())
1197    }
1198}
1199
1200/// The mbx build cache a container session was provisioned with. The
1201/// directory is a host path that is mounted read-write at the same absolute
1202/// path inside the container.
1203#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1204#[serde(deny_unknown_fields)]
1205pub struct SessionBuildCache {
1206    /// The container host this cache was resolved on. A session moved to a
1207    /// different host cannot reuse it, so the decision is made again there.
1208    pub host: String,
1209    pub directory: PathBuf,
1210    /// Read compatibility for old records only. Launch never uses this value;
1211    /// budgets belong to the shared machine configuration.
1212    #[serde(default, skip_serializing_if = "Option::is_none")]
1213    pub max_size: Option<String>,
1214    /// A `[target] root` the host's mbx configuration relocates outside the
1215    /// cache directory, mounted read-write at the same path as well.
1216    #[serde(default, skip_serializing_if = "Option::is_none")]
1217    pub target_root: Option<PathBuf>,
1218}
1219
1220/// How much disk Mjolnir's own copies of sessions use, and how much an
1221/// `archive_after_days` value would free. Settings shows this on the
1222/// SessionWiki page so the effect of a value is visible before it is saved.
1223#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1224pub struct ArchiveSpacePreview {
1225    /// Every session record Mjolnir holds, archived or not.
1226    pub sessions: usize,
1227    /// What those sessions' checkpoints and attachments occupy.
1228    pub bytes: u64,
1229    /// The sessions an `archive_after_days` value would catch, and their
1230    /// share of `bytes`. Both are zero when no value is set.
1231    pub reclaimable_sessions: usize,
1232    pub reclaimable_bytes: u64,
1233}
1234
1235/// What a container target's host resolves for its blank build cache
1236/// settings right now. Settings shows this beside each "automatic" field so
1237/// the values a session would actually run with are visible before one starts.
1238#[derive(Debug, Clone, PartialEq, Eq)]
1239pub struct BuildCachePreview {
1240    /// The host's own mbx version, or `None` when it has none on `PATH`.
1241    pub native_mbx: Option<String>,
1242    /// The cache directory sessions would mount, once known.
1243    pub directory: Option<PathBuf>,
1244    /// The budget sessions would run with, once known.
1245    pub max_size: Option<BuildCacheLimit>,
1246    pub target_max_size: Option<BuildCacheLimit>,
1247    /// True when a general-purpose host installation owns the configuration.
1248    pub user_managed: bool,
1249    pub application: BuildCacheApplication,
1250    /// Other mbx limits can further reduce the space available to worktrees.
1251    pub budget_note: Option<String>,
1252    /// What the cache on that host has done so far, when it has a tally.
1253    pub stats: Option<BuildCacheStats>,
1254    /// Why sessions on this target run without a cache, or `None` when they
1255    /// share one.
1256    pub off_reason: Option<BuildCacheOff>,
1257}
1258
1259/// mbx's own running totals for one host's cache, read from the tally beside
1260/// the store.
1261///
1262/// These are machine-wide and cumulative: every session on that host and any
1263/// native builds the user ran themselves are counted together, since they
1264/// share one cache. They answer whether the cache is being used at all, not
1265/// what one session got out of it.
1266#[derive(Debug, Clone, PartialEq, Eq)]
1267pub struct BuildCacheStats {
1268    /// Builds that went through mbx. Zero means nothing has used the shim.
1269    pub builds: u64,
1270    /// Compilations answered from the cache instead of run.
1271    pub cached_compilations: u64,
1272    /// Compiler time those answers avoided, as mbx estimates it.
1273    pub avoided_compiler_ns: u64,
1274    /// Restored output bytes that were cloned rather than copied.
1275    pub reflinked_bytes: u64,
1276}
1277
1278/// Why a target's sessions run without the build cache.
1279#[derive(Debug, Clone, PartialEq, Eq)]
1280pub enum BuildCacheOff {
1281    /// This machine's own `enabled = false`.
1282    TurnedOff,
1283    /// Nothing on the machine's settings page can turn it on: the global
1284    /// switch, the host's mbx, or its filesystem. The text says which.
1285    Unavailable(String),
1286}
1287
1288impl std::fmt::Display for BuildCacheOff {
1289    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1290        match self {
1291            Self::TurnedOff => formatter.write_str("turned off for this machine"),
1292            Self::Unavailable(reason) => formatter.write_str(reason),
1293        }
1294    }
1295}
1296
1297/// Where a machine build cache budget comes from.
1298#[derive(Debug, Clone, PartialEq, Eq)]
1299pub enum BuildCacheLimit {
1300    /// An explicit mj machine setting written to shared mbx configuration.
1301    Size(String),
1302    /// The host's own `~/.config/mbx/config.toml` carries the budget. The
1303    /// total it sets, when it sets one.
1304    HostConfiguration(Option<String>),
1305    /// An automatic total initialized by mj, independent of later disk growth.
1306    MjDefault(String),
1307    /// The pinned mbx's disk-scaled default, or no combined limit.
1308    MbxDefault(Option<String>),
1309}
1310
1311#[derive(Debug, Clone, Default, PartialEq, Eq)]
1312pub enum BuildCacheApplication {
1313    #[default]
1314    Pending,
1315    Applied,
1316    Failed(String),
1317}
1318
1319#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1320#[serde(deny_unknown_fields)]
1321pub struct SessionRecord {
1322    pub id: String,
1323    /// Owning workspace while active, or the most recent workspace while inactive.
1324    ///
1325    /// Inactive histories are globally resumable, so this id may refer to a
1326    /// workspace that has since been deleted.
1327    #[serde(default = "default_session_workspace_id")]
1328    pub workspace_id: String,
1329    pub title: String,
1330    pub harness_kind: HarnessKind,
1331    pub last_profile: String,
1332    pub bundle_id: String,
1333    /// Existing project directory used directly by a local or SSH bare target.
1334    #[serde(default, skip_serializing_if = "Option::is_none")]
1335    pub project_directory: Option<PathBuf>,
1336    /// Git worktree created and owned by Hel for this raw-project session.
1337    #[serde(default, skip_serializing_if = "Option::is_none")]
1338    pub managed_worktree: Option<ManagedWorktree>,
1339    /// None preserves automatic selection; false uses the selected directory.
1340    #[serde(default, skip_serializing_if = "Option::is_none")]
1341    pub create_managed_worktree: Option<bool>,
1342    /// Diff baseline as supplied by the caller (the API's `base`). A raw
1343    /// managed worktree also starts here; a bundle checkout keeps its selected
1344    /// remote branch tip. With `checkout`, it applies to the checked-out
1345    /// repository only, and when absent that repository's base is the
1346    /// checkout commit. The resolved baseline lands in
1347    /// `managed_worktree.base_commit` or the clone's `mj.baseCommit`.
1348    #[serde(default, skip_serializing_if = "Option::is_none")]
1349    pub launch_base: Option<String>,
1350    /// Existing branch to check out in a new isolated workspace (the API's
1351    /// `branch` without `at`). Never set together with `checkout`, which
1352    /// carries its own branch.
1353    #[serde(default, skip_serializing_if = "Option::is_none")]
1354    pub launch_branch: Option<String>,
1355    /// Immutable exact starting selection for one bundle repository, built
1356    /// from the API's `at` and `branch`. Resume preserves checkpointed work
1357    /// rather than applying this selection again.
1358    #[serde(default, skip_serializing_if = "Option::is_none")]
1359    pub checkout: Option<crate::remote_git::ExactCheckout>,
1360    #[serde(default, skip_serializing_if = "Option::is_none")]
1361    pub expected_runtime_identity: Option<String>,
1362    /// Last verified publication verdict, tied to its checkpoint digest.
1363    #[serde(default, skip_serializing_if = "Option::is_none")]
1364    pub publication: Option<PublicationAssessment>,
1365    /// Stored delegation policy. Historical records without a choice use native delegation.
1366    #[serde(
1367        default,
1368        alias = "mjolnir_subagents",
1369        deserialize_with = "crate::subagent::deserialize_optional_policy"
1370    )]
1371    pub subagents: Option<crate::subagent::SubagentPolicy>,
1372    pub target_template_id: String,
1373    #[serde(default, skip_serializing_if = "Option::is_none")]
1374    pub resource_allocation: Option<SessionResourceAllocation>,
1375    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1376    pub additional_mounts: Vec<AdditionalMount>,
1377    /// Per-session container CPU limit that overrides the target template's
1378    /// value. It is applied the next time the container is created.
1379    #[serde(default, skip_serializing_if = "Option::is_none")]
1380    pub container_cpus: Option<String>,
1381    /// Per-session container memory limit that overrides the target
1382    /// template's value. It is applied the next time the container is created.
1383    #[serde(default, skip_serializing_if = "Option::is_none")]
1384    pub container_memory: Option<String>,
1385    /// In-container workspace root this session's repositories live under.
1386    /// `None` is a session whose container predates per-session workspaces and
1387    /// therefore keeps the shared legacy `/workspace`; every session created
1388    /// since records `/workspace/<session id>`, so two checkouts of one project
1389    /// on a host never share an absolute path.
1390    #[serde(default, skip_serializing_if = "Option::is_none")]
1391    pub container_workspace: Option<PathBuf>,
1392    /// The mbx build cache this session's container runs with, decided once at
1393    /// provisioning. `None` means the session runs without a build cache;
1394    /// resume, move, and sub-agent children reuse the recorded value.
1395    #[serde(default, skip_serializing_if = "Option::is_none")]
1396    pub build_cache: Option<SessionBuildCache>,
1397    pub state: SessionState,
1398    /// Legacy visibility preference, retained for record compatibility.
1399    /// Current surfaces do not hide sessions based on this flag.
1400    #[serde(default, skip_serializing_if = "is_false")]
1401    pub archived: bool,
1402    #[serde(default, skip_serializing_if = "Option::is_none")]
1403    pub target: Option<TargetLocator>,
1404    /// Connection and worker settings captured when this target was selected.
1405    #[serde(default, skip_serializing_if = "Option::is_none")]
1406    pub target_runtime: Option<TargetRuntimeSettings>,
1407    #[serde(default, skip_serializing_if = "Option::is_none")]
1408    pub native_session_id: Option<String>,
1409    #[serde(default, skip_serializing_if = "Option::is_none")]
1410    pub acp_session_title: Option<String>,
1411    #[serde(default, skip_serializing_if = "Option::is_none")]
1412    pub session_title_override: Option<String>,
1413    pub created_at: String,
1414    pub updated_at: String,
1415    #[serde(default, alias = "detached_after_event_ordinal")]
1416    pub viewed_through_event_ordinal: u64,
1417    /// Unsent chat input carried across a detach, so returning to a session
1418    /// restores what the user was typing. Empty means no draft.
1419    #[serde(default, skip_serializing_if = "String::is_empty")]
1420    pub draft_input: String,
1421    /// Why the last operation on this session failed.
1422    ///
1423    /// This is usually a raw controller error chain, which names profile
1424    /// homes, project paths and SSH hosts, so a public projection publishes it
1425    /// only for a session that is stopped or failed. The one exception is a
1426    /// sentence the controller composed for the person; see
1427    /// [`CLOSE_FAILURE_PREFIX`].
1428    #[serde(default, skip_serializing_if = "Option::is_none")]
1429    pub last_error: Option<String>,
1430    #[serde(default, skip_serializing_if = "Option::is_none")]
1431    pub last_checkpoint_error: Option<String>,
1432    #[serde(default, skip_serializing_if = "Option::is_none")]
1433    pub checkpoint: Option<CheckpointMetadata>,
1434}
1435
1436#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1437#[serde(deny_unknown_fields)]
1438pub struct HostContainerSize {
1439    pub cpus: u64,
1440    pub memory_bytes: u64,
1441}
1442
1443fn default_session_workspace_id() -> String {
1444    crate::workspace::DEFAULT_WORKSPACE_ID.to_owned()
1445}
1446
1447/// How a failed close's reason begins in [`SessionRecord::last_error`].
1448///
1449/// A close that fails is non-destructive: the session goes back to the state
1450/// it was running in. Its reason therefore has to be published for a live
1451/// session, which the raw error chains in the same field never are. The
1452/// controller composes a sentence for the person and tags it with this prefix,
1453/// and the projection reads the tag to tell the two apart. Written in one
1454/// place and read in one place, so the tag cannot drift.
1455pub const DESTRUCTION_FAILURE_PREFIX: &str = "the destruction did not finish";
1456
1457pub const CLOSE_FAILURE_PREFIX: &str = "the suspension did not finish";
1458
1459/// Recognize safe lifecycle outcomes, including records saved before the rename.
1460pub fn is_public_lifecycle_error(error: &str) -> bool {
1461    error.starts_with(CLOSE_FAILURE_PREFIX)
1462        || error.starts_with(DESTRUCTION_FAILURE_PREFIX)
1463        || error.starts_with("the close did not finish")
1464}
1465
1466/// How a target is named beside a session, wherever a surface shows one.
1467///
1468/// A bare target (`local-bare`, `ssh-bare`) opens a project directory directly,
1469/// so the target alone does not say what the session was working on and the
1470/// project's own folder name is appended. Every other kind names a provisioned
1471/// environment that already identifies itself, so the target id stands alone.
1472/// A target id the configuration no longer holds is shown verbatim, because its
1473/// kind is no longer known.
1474///
1475/// Shared so the live session summary and the Resume dialog's archived rows
1476/// cannot drift apart: [`SessionRecord::project_target`] calls this, and so
1477/// does the archived row built from the SessionWiki index.
1478#[must_use]
1479pub fn target_label(config: &Config, target_id: &str, project: Option<&Path>) -> String {
1480    if !matches!(
1481        config.targets.get(target_id),
1482        Some(TargetTemplate::LocalBare | TargetTemplate::SshBare { .. })
1483    ) {
1484        return target_id.to_owned();
1485    }
1486    project.and_then(Path::file_name).map_or_else(
1487        || target_id.to_owned(),
1488        |directory| format!("{target_id}/{}", directory.to_string_lossy()),
1489    )
1490}
1491
1492/// A session's starting selection in the terms the API and CLI use.
1493#[derive(Debug, Clone, Default, PartialEq, Eq)]
1494pub struct StartSelection {
1495    /// Commit the workspace started checked out at.
1496    pub at: Option<String>,
1497    /// Branch created at `at`, or the existing branch checked out without it.
1498    pub branch: Option<String>,
1499    /// Diff base; `at` unless the caller named another.
1500    pub base: Option<String>,
1501}
1502
1503impl SessionRecord {
1504    /// The starting selection this record stores, as `at`, `branch` and
1505    /// `base`. The record keeps the older field layout, so this is the one
1506    /// place that maps it to the public names.
1507    pub fn start_selection(&self) -> StartSelection {
1508        let at = self
1509            .checkout
1510            .as_ref()
1511            .map(|checkout| checkout.commit.clone());
1512        StartSelection {
1513            branch: self
1514                .checkout
1515                .as_ref()
1516                .and_then(|checkout| checkout.branch.clone())
1517                .or_else(|| self.launch_branch.clone()),
1518            base: self.launch_base.clone().or_else(|| at.clone()),
1519            at,
1520        }
1521    }
1522
1523    /// Cached verdict for an independent clone. Active checkouts are unknown
1524    /// until a new checkpoint binds an assessment to their exact contents.
1525    pub fn publication_state(&self) -> Option<PublicationState> {
1526        let independent_clone = self
1527            .managed_worktree
1528            .as_ref()
1529            .is_some_and(|owned| owned.kind == ManagedCheckoutKind::Clone)
1530            || (self.managed_worktree.is_none() && self.project_directory.is_none());
1531        if !independent_clone {
1532            return None;
1533        }
1534        if self.state.is_active() {
1535            return Some(PublicationState::Unknown);
1536        }
1537        Some(
1538            self.checkpoint
1539                .as_ref()
1540                .zip(self.publication.as_ref())
1541                .filter(|(checkpoint, assessment)| {
1542                    assessment.checkpoint_sha256 == checkpoint.sha256
1543                })
1544                .map_or(PublicationState::Unknown, |(_, assessment)| {
1545                    if assessment.dirty || assessment.stashed {
1546                        PublicationState::Unpublished
1547                    } else {
1548                        assessment.state
1549                    }
1550                }),
1551        )
1552    }
1553
1554    /// The target access settings this session's commands use: the ones
1555    /// recorded when its target was selected, with the machine's current ssh
1556    /// options while they still reach the same host, user, and port, or the
1557    /// configured target's when nothing was recorded.
1558    pub fn target_runtime_settings<'a>(
1559        &'a self,
1560        config: &Config,
1561    ) -> Result<std::borrow::Cow<'a, TargetRuntimeSettings>> {
1562        if let Some(runtime) = &self.target_runtime {
1563            // How to reach the target follows the machine's current ssh
1564            // options; where it is stays as recorded (launch finding R3-7).
1565            if let Some(refreshed) =
1566                config
1567                    .targets
1568                    .get(&self.target_template_id)
1569                    .and_then(|template| {
1570                        runtime.with_current_ssh_options(&TargetRuntimeSettings::from(template))
1571                    })
1572            {
1573                return Ok(std::borrow::Cow::Owned(refreshed));
1574            }
1575            return Ok(std::borrow::Cow::Borrowed(runtime));
1576        }
1577        let template = config.targets.get(&self.target_template_id).ok_or_else(|| {
1578            crate::refusal::Refusal::precondition(format!(
1579                "Session {:?} has no recorded target access settings. Restore target {:?} in config.toml once, then retry.",
1580                self.id, self.target_template_id))
1581        })?;
1582        let runtime = TargetRuntimeSettings::from(template);
1583        if let Some(locator) = &self.target {
1584            crate::targets::TargetLocator::try_from(crate::targets::RecordedTarget {
1585                locator, runtime: Some(&runtime), session_id: &self.id,
1586            }).map_err(|error| crate::refusal::Refusal::precondition(format!(
1587                "Session {:?} cannot recover target {:?}: {error}. Restore its original target settings, then retry.", self.id, self.target_template_id)))?;
1588        }
1589        Ok(std::borrow::Cow::Owned(runtime))
1590    }
1591
1592    /// The recorded failure that is safe to publish whatever state this
1593    /// session is in, because the controller wrote it for the person rather
1594    /// than copying an error chain into it.
1595    #[must_use]
1596    pub fn public_error(&self) -> Option<&str> {
1597        self.last_error
1598            .as_deref()
1599            .filter(|error| is_public_lifecycle_error(error))
1600    }
1601
1602    /// Configuration drift belongs to this session, not the entire controller.
1603    /// The diagnostic contains only public identifiers, so both UIs can show it.
1604    pub fn configuration_issue(&self, config: &Config) -> Option<String> {
1605        if !self.state.is_active() {
1606            return None;
1607        }
1608        let mut issues = Vec::new();
1609        match config.profiles.get(&self.last_profile) {
1610            None => issues.push(format!("missing profile {:?}", self.last_profile)),
1611            Some(profile) if profile.kind != self.harness_kind => issues.push(format!(
1612                "expects {:?}, but profile {:?} is {:?}",
1613                self.harness_kind, self.last_profile, profile.kind
1614            )),
1615            Some(_) => {}
1616        }
1617        if self.project_directory.is_none() && !config.bundles.contains_key(&self.bundle_id) {
1618            issues.push(format!("missing bundle {:?}", self.bundle_id));
1619        }
1620        if self.target_runtime.is_none() && !config.targets.contains_key(&self.target_template_id) {
1621            issues.push(format!(
1622                "missing target template {:?}",
1623                self.target_template_id
1624            ));
1625        }
1626        (!issues.is_empty()).then(|| format!(
1627            "Session {:?} needs configuration repair: {}. Restore these entries in config.toml, then retry. Run mj setup to rediscover installed profiles and targets; existing sessions are preserved.",
1628            self.id, issues.join("; ")
1629        ))
1630    }
1631
1632    pub fn validate_configuration(&self, config: &Config) -> Result<()> {
1633        if let Some(issue) = self.configuration_issue(config) {
1634            return Err(crate::refusal::Refusal::precondition(issue).into());
1635        }
1636        Ok(())
1637    }
1638
1639    /// User-visible session name, independent of the initial prompt stored in `title`.
1640    pub fn display_title(&self) -> &str {
1641        self.session_title_override
1642            .as_deref()
1643            .or(self.acp_session_title.as_deref())
1644            .unwrap_or(&self.id)
1645    }
1646
1647    /// The name a listing shows: the display title, except that a session
1648    /// the harness has not named yet and nobody renamed would otherwise be
1649    /// named by its id, which every listing already prints beside it. The
1650    /// title it was created with says more (launch findings F-12 and R2-8).
1651    pub fn listed_title(&self) -> &str {
1652        let named = self.session_title_override.is_some() || self.acp_session_title.is_some();
1653        if !named && !self.title.trim().is_empty() {
1654            return &self.title;
1655        }
1656        self.display_title()
1657    }
1658
1659    /// Project this session works in, as the session list and the chat header
1660    /// both name it: the source repository of a managed worktree, else the
1661    /// project directory, else the bundle's primary repository, else the
1662    /// bundle id.
1663    pub fn project_name(&self, config: &Config) -> String {
1664        if let Some(worktree) = &self.managed_worktree {
1665            return path_leaf(&worktree.source_repository);
1666        }
1667        if let Some(project_directory) = &self.project_directory {
1668            return path_leaf(project_directory);
1669        }
1670        self.bundle_source_name(config)
1671    }
1672
1673    /// Target label used by the live session summary. Bare targets identify
1674    /// the project directory they open directly; workspace targets already
1675    /// identify the provisioned environment on their own.
1676    pub fn project_target(&self, config: &Config, target_id: &str) -> String {
1677        let project = self
1678            .managed_worktree
1679            .as_ref()
1680            .map(|worktree| &worktree.source_project_directory)
1681            .or(self.project_directory.as_ref());
1682        target_label(config, target_id, project.map(PathBuf::as_path))
1683    }
1684
1685    /// Stable source identity used to group sessions. Managed worktrees point
1686    /// back at their source repository, raw sessions use their project
1687    /// directory until their Git origin is resolved, and bundle sessions use
1688    /// their complete canonical repository set when configured.
1689    pub fn project_source(&self, config: &Config) -> ProjectSourceIdentity {
1690        if let Some(worktree) = &self.managed_worktree {
1691            return ProjectSourceIdentity::path(&worktree.source_repository, None);
1692        }
1693        if let Some(project_directory) = &self.project_directory {
1694            let remote = match &self.target {
1695                Some(TargetLocator::SshBare { host, .. }) => Some(host.as_str()),
1696                _ => None,
1697            };
1698            return ProjectSourceIdentity::path(project_directory, remote);
1699        }
1700        self.bundle_source_identity(config)
1701            .unwrap_or_else(|| ProjectSourceIdentity {
1702                key: format!("bundle:{}", self.bundle_id),
1703                short: path_leaf(Path::new(&self.bundle_id)),
1704                full: self.bundle_id.clone(),
1705            })
1706    }
1707
1708    /// Resolve the display name shared by session headings, chat headers, and
1709    /// resume details for a bundle-backed session.
1710    fn bundle_source_name(&self, config: &Config) -> String {
1711        self.bundle_source_identity(config)
1712            .map(|source| source.short)
1713            .unwrap_or_else(|| path_leaf(Path::new(&self.bundle_id)))
1714    }
1715
1716    /// Resolve the canonical identity of every repository in a bundle for
1717    /// grouping and display naming.
1718    fn bundle_source_identity(&self, config: &Config) -> Option<ProjectSourceIdentity> {
1719        let bundle = config.bundles.get(&self.bundle_id)?;
1720        let sources = bundle
1721            .repositories
1722            .iter()
1723            .map(repository_source_identity)
1724            .collect::<Option<Vec<_>>>()?;
1725        ProjectSourceIdentity::bundle(sources)
1726    }
1727
1728    /// Orders two sessions the way the session list's sequence view does:
1729    /// oldest first by creation time, with the id as a stable tiebreak. A
1730    /// session whose timestamp does not parse sorts last.
1731    pub fn compare_by_creation(&self, other: &Self) -> std::cmp::Ordering {
1732        self.creation_order_key().cmp(&other.creation_order_key())
1733    }
1734
1735    /// Parse once per session when used with `sort_by_cached_key`.
1736    pub fn creation_order_key(&self) -> (bool, Option<i64>, &str) {
1737        let timestamp = created_at_seconds(&self.created_at);
1738        (timestamp.is_none(), timestamp, &self.id)
1739    }
1740
1741    fn validate(&self, map_id: &str) -> Result<()> {
1742        validate_id("session", &self.id)?;
1743        if self.id != map_id {
1744            bail!(
1745                "session map key {map_id:?} does not match record id {:?}",
1746                self.id
1747            );
1748        }
1749        validate_id("workspace", &self.workspace_id)?;
1750        validate_id("profile", &self.last_profile)?;
1751        validate_id("bundle", &self.bundle_id)?;
1752        if let Some(project_directory) = &self.project_directory
1753            && (!project_directory.is_absolute()
1754                || project_directory
1755                    .components()
1756                    .any(|part| part == Component::ParentDir))
1757        {
1758            bail!("session {:?} has an unsafe project directory", self.id);
1759        }
1760        if let Some(managed_worktree) = &self.managed_worktree {
1761            managed_worktree.validate(&self.id, self.project_directory.as_deref())?;
1762        }
1763        validate_id("target template", &self.target_template_id)?;
1764        if let Some(allocation) = &self.resource_allocation {
1765            allocation.validate()?;
1766        }
1767        validate_additional_mounts(&self.additional_mounts)?;
1768        if self.title.trim().is_empty() {
1769            bail!("session {:?} has an empty title", self.id);
1770        }
1771        if self
1772            .acp_session_title
1773            .as_ref()
1774            .is_some_and(|title| title.trim().is_empty())
1775            || self
1776                .session_title_override
1777                .as_ref()
1778                .is_some_and(|title| title.trim().is_empty())
1779        {
1780            bail!("session {:?} has an empty display title", self.id);
1781        }
1782        if self.created_at.trim().is_empty() || self.updated_at.trim().is_empty() {
1783            bail!("session {:?} has an empty timestamp", self.id);
1784        }
1785        if let Some(target) = &self.target {
1786            target.validate(&self.id)?;
1787        }
1788        if let Some(checkpoint) = &self.checkpoint {
1789            checkpoint.validate()?;
1790        }
1791        Ok(())
1792    }
1793}
1794
1795fn repository_source_identity(repository: &ProjectRepository) -> Option<ProjectSourceIdentity> {
1796    repository
1797        .github
1798        .as_deref()
1799        .and_then(ProjectSourceIdentity::git_remote)
1800        .or_else(|| {
1801            repository
1802                .local
1803                .as_deref()
1804                .map(|path| ProjectSourceIdentity::path(path, None))
1805        })
1806}
1807
1808#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1809pub struct ProjectSourceIdentity {
1810    pub key: String,
1811    pub short: String,
1812    pub full: String,
1813}
1814
1815impl ProjectSourceIdentity {
1816    /// Combine repository identities into one stable bundle identity.
1817    pub fn bundle(mut sources: Vec<Self>) -> Option<Self> {
1818        if sources.is_empty() {
1819            return None;
1820        }
1821        sources.sort_by(|left, right| {
1822            left.key
1823                .cmp(&right.key)
1824                .then_with(|| left.full.cmp(&right.full))
1825                .then_with(|| left.short.cmp(&right.short))
1826        });
1827        sources.dedup_by(|left, right| left.key == right.key);
1828        if sources.len() == 1 {
1829            return sources.pop();
1830        }
1831        let keys = sources
1832            .iter()
1833            .map(|source| source.key.clone())
1834            .collect::<Vec<_>>();
1835        let key = serde_json::to_string(&keys).ok()?;
1836        Some(Self {
1837            key: format!("bundle:{key}"),
1838            short: sources
1839                .iter()
1840                .map(|source| source.short.as_str())
1841                .collect::<Vec<_>>()
1842                .join(" + "),
1843            full: sources
1844                .iter()
1845                .map(|source| source.full.as_str())
1846                .collect::<Vec<_>>()
1847                .join(" + "),
1848        })
1849    }
1850
1851    /// Canonicalizes a Git remote so raw checkouts group as the same project
1852    /// even when their worktree paths differ.
1853    pub fn git_remote(source: &str) -> Option<Self> {
1854        if let Some(normalized) = normalize_github_source(source) {
1855            let short = normalized
1856                .rsplit_once('/')
1857                .map_or(normalized.as_str(), |(_, repository)| repository)
1858                .to_owned();
1859            return Some(Self {
1860                key: format!("github:{}", normalized.to_lowercase()),
1861                short,
1862                full: normalized,
1863            });
1864        }
1865        let normalized = source.trim().trim_end_matches('/').trim_end_matches(".git");
1866        if normalized.is_empty() {
1867            return None;
1868        }
1869        let short = normalized
1870            .rsplit(['/', ':'])
1871            .find(|part| !part.is_empty())
1872            .unwrap_or(normalized)
1873            .to_owned();
1874        Some(Self {
1875            key: format!("git:{}", normalized.to_lowercase()),
1876            short,
1877            full: normalized.to_owned(),
1878        })
1879    }
1880
1881    /// Build a local-root identity, qualified by host for remote directories.
1882    pub fn path(path: &Path, remote: Option<&str>) -> Self {
1883        let normalized = path.components().collect::<PathBuf>();
1884        let path_text = normalized.to_string_lossy().into_owned();
1885        let full = remote.map_or_else(|| path_text.clone(), |host| format!("{host}:{path_text}"));
1886        let key = remote.map_or_else(
1887            || format!("path:{path_text}"),
1888            |host| format!("path:{}:{path_text}", host.to_lowercase()),
1889        );
1890        Self {
1891            key,
1892            short: path_leaf(path),
1893            full,
1894        }
1895    }
1896}
1897
1898fn normalize_github_source(source: &str) -> Option<String> {
1899    let source = source.trim();
1900    let path = source
1901        .strip_prefix("https://github.com/")
1902        .or_else(|| source.strip_prefix("http://github.com/"))
1903        .or_else(|| source.strip_prefix("git@github.com:"))
1904        .or_else(|| source.strip_prefix("ssh://git@github.com/"))
1905        .or_else(|| {
1906            (!source.contains("://") && !source.contains('@') && !source.contains(':'))
1907                .then_some(source)
1908        })?
1909        .trim_end_matches(".git");
1910    let mut parts = path.split('/');
1911    let owner = parts.next()?;
1912    let repository = parts.next()?;
1913    (!owner.is_empty() && !repository.is_empty() && parts.next().is_none())
1914        .then(|| format!("{owner}/{repository}"))
1915}
1916
1917/// Last component of a path, falling back to the whole path when it has none.
1918fn path_leaf(path: &Path) -> String {
1919    path.file_name()
1920        .unwrap_or(path.as_os_str())
1921        .to_string_lossy()
1922        .into_owned()
1923}
1924
1925fn created_at_seconds(timestamp: &str) -> Option<i64> {
1926    chrono::DateTime::parse_from_rfc3339(timestamp)
1927        .ok()
1928        .map(|timestamp| timestamp.timestamp())
1929}
1930
1931#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1932#[serde(deny_unknown_fields)]
1933pub struct State {
1934    #[serde(default)]
1935    pub last_subagent_policy: crate::subagent::SubagentPolicy,
1936    pub version: u32,
1937    #[serde(default, skip_serializing_if = "SnapshotMap::is_empty")]
1938    pub sessions: SnapshotMap<String, SessionRecord>,
1939    /// Child sessions keyed by their session id. The relationship lives in
1940    /// controller state so every control surface sees the same session family.
1941    #[serde(default, skip_serializing_if = "SnapshotMap::is_empty")]
1942    pub subagents: SnapshotMap<String, SubagentRecord>,
1943    /// Recently used source directories, keyed by `local` or SSH host name.
1944    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1945    pub mount_history: BTreeMap<String, Vec<PathBuf>>,
1946    /// Most recently launched container size on each physical target host.
1947    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1948    pub container_sizes: BTreeMap<String, HostContainerSize>,
1949}
1950
1951impl Default for State {
1952    fn default() -> Self {
1953        Self {
1954            version: STATE_VERSION,
1955            last_subagent_policy: Default::default(),
1956            sessions: SnapshotMap::new(),
1957            subagents: SnapshotMap::new(),
1958            mount_history: BTreeMap::new(),
1959            container_sizes: BTreeMap::new(),
1960        }
1961    }
1962}
1963
1964impl State {
1965    /// How a notice names a session: the title the session list shows
1966    /// (`listed_title`, which includes the title it was created with), or its
1967    /// short id when it has no title or its record is gone (launch findings
1968    /// B-3, R5-5 and R8-3).
1969    #[must_use]
1970    pub fn session_notice_name(&self, session_id: &str) -> String {
1971        match self.sessions.get(session_id) {
1972            Some(session) if session.listed_title() != session.id => {
1973                session.listed_title().to_owned()
1974            }
1975            _ => short_id(session_id).to_owned(),
1976        }
1977    }
1978
1979    /// The session whose project identity names a row.
1980    ///
1981    /// A sub-agent child runs inside its parent's workspace and owns no
1982    /// managed worktree, so its own `project_directory` is the parent's
1983    /// worktree checkout, whose directory is named after the parent session
1984    /// id. Reading the project identity from the parent instead keeps a child
1985    /// under the same project heading and target label as the session it
1986    /// belongs to.
1987    #[must_use]
1988    pub fn project_identity_session<'a>(&'a self, session: &'a SessionRecord) -> &'a SessionRecord {
1989        self.subagents
1990            .get(&session.id)
1991            .and_then(|record| self.sessions.get(&record.parent_session_id))
1992            .unwrap_or(session)
1993    }
1994
1995    /// Whether `id` names a sub-agent rather than a session the user started:
1996    /// a Mjolnir-managed child, or the record a client builds to show a
1997    /// harness-owned child (see [`crate::native_agent::view_id`]).
1998    ///
1999    /// Every list of top-level sessions filters with this, so the lists
2000    /// cannot disagree about what a sub-agent is.
2001    #[must_use]
2002    pub fn is_subagent_session(&self, id: &str) -> bool {
2003        self.subagents.contains_key(id) || crate::native_agent::is_view_id(id)
2004    }
2005
2006    pub fn validate(&self) -> Result<()> {
2007        if self.version != STATE_VERSION {
2008            bail!(
2009                "unsupported Mjolnir state version {}; expected {STATE_VERSION}",
2010                self.version
2011            );
2012        }
2013        for (id, session) in &self.sessions {
2014            session.validate(id)?;
2015        }
2016        for child_id in self.subagents.keys() {
2017            self.validate_subagent(child_id)?;
2018        }
2019        for (host, sources) in &self.mount_history {
2020            if host.trim().is_empty() {
2021                bail!("mount history contains an empty host key");
2022            }
2023            if sources.iter().any(|source| !source.is_absolute()) {
2024                bail!("mount history for {host:?} contains a non-absolute source path");
2025            }
2026        }
2027        for (host, size) in &self.container_sizes {
2028            if host.trim().is_empty() {
2029                bail!("container size history contains an empty host key");
2030            }
2031            if size.cpus == 0 || size.memory_bytes == 0 {
2032                bail!("container size history for {host:?} contains a zero value");
2033            }
2034            if size.cpus > i64::MAX as u64 || size.memory_bytes > i64::MAX as u64 {
2035                bail!("container size history for {host:?} exceeds SQLite integer range");
2036            }
2037        }
2038        Ok(())
2039    }
2040
2041    /// Validate one relationship after an incremental committed update.
2042    pub fn validate_subagent(&self, child_id: &str) -> Result<()> {
2043        let Some(subagent) = self.subagents.get(child_id) else {
2044            return Ok(());
2045        };
2046        if child_id != subagent.child_session_id {
2047            bail!("sub-agent key {child_id:?} does not match its child session id");
2048        }
2049        if child_id == subagent.parent_session_id {
2050            bail!("sub-agent {child_id:?} cannot be its own parent");
2051        }
2052        if !self.sessions.contains_key(child_id) {
2053            bail!("sub-agent {child_id:?} has no child session");
2054        }
2055        if !self.sessions.contains_key(&subagent.parent_session_id) {
2056            bail!(
2057                "sub-agent {child_id:?} has unknown parent {:?}",
2058                subagent.parent_session_id
2059            );
2060        }
2061        if self.subagents.contains_key(&subagent.parent_session_id) {
2062            bail!("sub-agent {child_id:?} cannot belong to another sub-agent");
2063        }
2064        if subagent.task_name.trim().is_empty()
2065            || subagent.profile_id.trim().is_empty()
2066            || subagent.request_key.trim().is_empty()
2067        {
2068            bail!("sub-agent {child_id:?} has incomplete relationship metadata");
2069        }
2070        Ok(())
2071    }
2072
2073    pub fn remember_mount_sources(&mut self, host: &str, mounts: &[AdditionalMount]) {
2074        if mounts.is_empty() {
2075            return;
2076        }
2077        let sources = self.mount_history.entry(host.to_owned()).or_default();
2078        for mount in mounts.iter().rev() {
2079            sources.retain(|source| source != &mount.source);
2080            sources.insert(0, mount.source.clone());
2081        }
2082        sources.truncate(20);
2083    }
2084
2085    pub fn remember_container_size(&mut self, host: &str, size: HostContainerSize) {
2086        self.container_sizes.insert(host.to_owned(), size);
2087    }
2088
2089    pub fn project_directories(&self, host: &str) -> &[PathBuf] {
2090        self.mount_history
2091            .get(&project_history_key(host))
2092            .map(Vec::as_slice)
2093            .unwrap_or_default()
2094    }
2095
2096    pub fn remember_project_directory(&mut self, host: &str, directory: &Path) {
2097        let key = project_history_key(host);
2098        let directories = self.mount_history.entry(key).or_default();
2099        directories.retain(|existing| existing != directory);
2100        directories.insert(0, directory.to_path_buf());
2101        directories.truncate(20);
2102    }
2103
2104    pub fn destroy_stopped_session(&mut self, session_id: &str) -> Result<SessionRecord> {
2105        let session = self
2106            .sessions
2107            .get(session_id)
2108            .with_context(|| format!("unknown session {session_id}"))?;
2109        if session.state.is_active() {
2110            bail!("refusing to destroy active session {session_id}");
2111        }
2112        Ok(self
2113            .sessions
2114            .remove(session_id)
2115            .expect("session checked above"))
2116    }
2117
2118    /// Remove a session record from state regardless of its lifecycle state.
2119    ///
2120    /// Force destruction is the one caller: by the time it runs, every
2121    /// external artifact has been torn down or its loss accepted, so no state
2122    /// is refused here.
2123    pub fn destroy_session_force(&mut self, session_id: &str) -> Result<SessionRecord> {
2124        self.sessions
2125            .get(session_id)
2126            .with_context(|| format!("unknown session {session_id}"))?;
2127        Ok(self
2128            .sessions
2129            .remove(session_id)
2130            .expect("session checked above"))
2131    }
2132
2133    /// Sessions that still read `bundle_id` from the config. A suspended
2134    /// session counts: resume looks its project up again. Only a session
2135    /// opened on a plain directory, or one whose data is already gone, does
2136    /// not need it.
2137    pub fn bundle_users(&self, bundle_id: &str) -> Vec<&SessionRecord> {
2138        self.sessions
2139            .values()
2140            .filter(|session| {
2141                session.bundle_id == bundle_id
2142                    && session.project_directory.is_none()
2143                    && session.state != SessionState::DestroyedWithDataLoss
2144            })
2145            .collect()
2146    }
2147
2148    /// Why `bundle_id` cannot be removed from the config, or `None` when no
2149    /// session uses it.
2150    pub fn bundle_removal_refusal(&self, bundle_id: &str) -> Option<String> {
2151        let users = self.bundle_users(bundle_id);
2152        if users.is_empty() {
2153            return None;
2154        }
2155        let mut names = users
2156            .iter()
2157            .take(3)
2158            .map(|session| format!("{:?}", session.listed_title()))
2159            .collect::<Vec<_>>();
2160        if users.len() > 3 {
2161            names.push(format!("{} more", users.len() - 3));
2162        }
2163        Some(format!(
2164            "Project {bundle_id:?} is used by {}: {}. Destroy those sessions before removing it.",
2165            if users.len() == 1 {
2166                "a session"
2167            } else {
2168                "sessions"
2169            },
2170            names.join(", ")
2171        ))
2172    }
2173
2174    /// Setup may add replacements under new names, but must not rewrite
2175    /// dependencies still owned by active sessions.
2176    pub fn validate_setup_update(&self, before: &Config, after: &Config) -> Result<()> {
2177        for session in self
2178            .sessions
2179            .values()
2180            .filter(|session| session.state.is_active())
2181        {
2182            let protected = if let Some(profile) = before.profiles.get(&session.last_profile) {
2183                let mut comparable = profile.clone();
2184                if let Some(updated) = after.profiles.get(&session.last_profile) {
2185                    comparable.enabled = updated.enabled;
2186                }
2187                // A mismatched harness is already broken; allow repairing it.
2188                profile.kind == session.harness_kind
2189                    && after.profiles.get(&session.last_profile) != Some(&comparable)
2190            } else {
2191                false
2192            };
2193            let bundle_changed = session.project_directory.is_none()
2194                && before
2195                    .bundles
2196                    .get(&session.bundle_id)
2197                    .is_some_and(|bundle| after.bundles.get(&session.bundle_id) != Some(bundle));
2198            // Build cache settings are resolved at provisioning time and kept
2199            // on the session record, so editing them does not disturb a
2200            // running session.
2201            let target_changed =
2202                before
2203                    .targets
2204                    .get(&session.target_template_id)
2205                    .is_some_and(|target| {
2206                        after
2207                            .targets
2208                            .get(&session.target_template_id)
2209                            .map(TargetTemplate::without_launch_only_settings)
2210                            != Some(target.without_launch_only_settings())
2211                    });
2212            if protected || bundle_changed || target_changed {
2213                // Named as the screen names them: the session by its title,
2214                // the project by its name rather than the internal bundle
2215                // id, and only the parts this change touches.
2216                let mut used = Vec::new();
2217                if protected {
2218                    used.push(format!("agent profile {:?}", session.last_profile));
2219                }
2220                if bundle_changed {
2221                    used.push(format!("project {:?}", session.project_name(before)));
2222                }
2223                if target_changed {
2224                    used.push(format!("runtime {:?}", session.target_template_id));
2225                }
2226                let used = match used.as_slice() {
2227                    [only] => only.clone(),
2228                    [rest @ .., last] => format!("{} and {last}", rest.join(", ")),
2229                    [] => unreachable!("something changed"),
2230                };
2231                let title = session.display_title();
2232                let named = if title == session.id {
2233                    format!(
2234                        "a running session in project {:?}",
2235                        session.project_name(before)
2236                    )
2237                } else {
2238                    format!("the running session {title:?}")
2239                };
2240                bail!(
2241                    "Setup would change the {used} that {named} uses. Save the new settings under a new name, or stop the session first."
2242                );
2243            }
2244        }
2245        Ok(())
2246    }
2247
2248    /// Strict validation for callers that need all active references intact.
2249    pub fn validate_against_config(&self, config: &Config) -> Result<()> {
2250        self.validate()?;
2251        config.validate()?;
2252        for session in self.sessions.values() {
2253            session.validate_configuration(config)?;
2254        }
2255        Ok(())
2256    }
2257}
2258
2259fn project_history_key(host: &str) -> String {
2260    format!("project:{host}")
2261}
2262
2263/// Generate an opaque, filesystem-safe stable id for a new logical session.
2264pub fn new_session_id() -> Result<String> {
2265    let mut random = [0u8; 16];
2266    getrandom::fill(&mut random)
2267        .map_err(|error| anyhow::anyhow!("generate Mjolnir session id: {error}"))?;
2268    Ok(crate::hex::lower_hex(random))
2269}
2270
2271/// Return the newest clean ACP session title from canonical worker events.
2272pub fn harness_session_title(events: &[SequencedEvent]) -> Option<String> {
2273    events.iter().rev().find_map(|event| {
2274        let WorkerEvent::Adapter { payload, .. } = &event.event else {
2275            return None;
2276        };
2277        let crate::acp::RuntimeEvent::SessionUpdate { update } =
2278            serde_json::from_value(payload.clone()).ok()?
2279        else {
2280            return None;
2281        };
2282        let kind = update
2283            .get("sessionUpdate")
2284            .and_then(serde_json::Value::as_str)?;
2285        let title = match kind {
2286            "session_info_update" | "session_title" => {
2287                update.get("title").and_then(serde_json::Value::as_str)
2288            }
2289            _ => None,
2290        }?;
2291        normalize_session_title(title)
2292    })
2293}
2294
2295/// The title a new session gets when whoever starts it gives none: the name
2296/// of its project directory (or its bundle id) and the profile, such as
2297/// "project via fake". The dashboard, the HTTP API, and `mj new` all use it,
2298/// so a session reads the same way whichever surface started it.
2299pub fn default_session_title(
2300    project_directory: Option<&Path>,
2301    bundle_id: &str,
2302    profile_id: &str,
2303) -> String {
2304    let project = project_directory.and_then(Path::file_name).map_or_else(
2305        || bundle_id.to_owned(),
2306        |name| name.to_string_lossy().into_owned(),
2307    );
2308    format!("{project} via {profile_id}")
2309}
2310
2311pub fn normalize_session_title(title: &str) -> Option<String> {
2312    let normalized = crate::relay::strip_hidden_prompt_context(title)
2313        .split_whitespace()
2314        .collect::<Vec<_>>()
2315        .join(" ");
2316    (!normalized.is_empty()).then_some(normalized)
2317}
2318
2319/// Build the short-lived title shown before the harness supplies its own.
2320///
2321/// The first visible user prompt is immediately useful for identifying a
2322/// session, but it can be arbitrarily large. Keep this fallback bounded; a
2323/// later ACP session-info update remains authoritative and replaces it.
2324pub fn provisional_session_title(prompt: &str) -> Option<String> {
2325    const MAX_TITLE_CHARS: usize = 64;
2326
2327    let normalized = normalize_session_title(prompt)?;
2328    if normalized.chars().count() <= MAX_TITLE_CHARS {
2329        return Some(normalized);
2330    }
2331
2332    let mut truncated = normalized
2333        .chars()
2334        .take(MAX_TITLE_CHARS - 1)
2335        .collect::<String>();
2336    if let Some(boundary) = truncated.rfind(char::is_whitespace) {
2337        truncated.truncate(boundary);
2338    }
2339    truncated.push('…');
2340    Some(truncated)
2341}
2342
2343pub fn short_id(id: &str) -> &str {
2344    id.get(..8).unwrap_or(id)
2345}
2346
2347#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
2348pub struct RecoveryCandidate {
2349    pub session_id: String,
2350    pub target_template_id: String,
2351    pub locator: TargetLocator,
2352    pub ownership: Option<crate::worker_launch::WorkerOwnership>,
2353    /// Instance that created the worker, from its label or tag, else from
2354    /// the ownership marker. `None` means an older build left no stamp.
2355    #[serde(default)]
2356    pub instance_id: Option<String>,
2357    /// State of the session this resource is labelled for, when the
2358    /// controller still tracks that session. A leftover resource the session
2359    /// record no longer names can only be destroyed, never adopted, because
2360    /// the session id is already taken.
2361    #[serde(default, skip_serializing_if = "Option::is_none")]
2362    pub tracked_session: Option<SessionState>,
2363}
2364
2365#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
2366pub struct RecoveryScan {
2367    pub candidates: Vec<RecoveryCandidate>,
2368    pub warnings: Vec<String>,
2369    /// Identity of the instance that ran the scan.
2370    #[serde(default)]
2371    pub instance_id: String,
2372    /// Candidates left out because another or an unknown instance created
2373    /// them and the scan was not widened to all instances.
2374    #[serde(default)]
2375    pub hidden_other_instances: usize,
2376}
2377
2378#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2379#[serde(deny_unknown_fields)]
2380pub struct ResumeRepositorySourceReceipt {
2381    pub session_id: String,
2382    pub bundle_id: String,
2383    pub checkpoint_sha256: String,
2384    pub repositories: Vec<crate::config::ProjectRepository>,
2385}
2386
2387#[cfg(test)]
2388mod tests;