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