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::subagent::SubagentRecord;
16use crate::targets::{AdditionalMount, validate_additional_mounts};
17
18pub const STATE_VERSION: u32 = 1;
19
20mod session_move;
21pub use session_move::*;
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
24#[serde(rename_all = "kebab-case")]
25pub enum SessionState {
26    Provisioning,
27    Running,
28    Disconnected,
29    Checkpointing,
30    Closing,
31    Destroying,
32    /// Checkpointed and torn down. Persisted as `"archived"` before the verb
33    /// was renamed, so the alias keeps those records loading.
34    #[serde(alias = "archived")]
35    Stopped,
36    Lost,
37    Error,
38    DestroyedWithDataLoss,
39}
40
41/// A lifecycle transition temporarily replaces the conversation in control surfaces.
42/// Operation ownership takes precedence over intermediate durable session states.
43#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
44#[serde(rename_all = "kebab-case")]
45pub enum SessionTransitionKind {
46    Starting,
47    Resuming,
48    Moving,
49    Suspending,
50    Destroying,
51}
52
53impl SessionTransitionKind {
54    pub const fn label(self) -> &'static str {
55        match self {
56            Self::Starting => "Starting",
57            Self::Resuming => "Resuming",
58            Self::Moving => "Moving",
59            Self::Suspending => "Suspending",
60            Self::Destroying => "Destroying",
61        }
62    }
63
64    pub fn for_session(state: SessionState, operation: Option<Self>) -> Option<Self> {
65        operation.or_else(|| state.transition_kind())
66    }
67}
68
69#[cfg(test)]
70mod transition_tests {
71    use super::{SessionState, SessionTransitionKind};
72
73    #[test]
74    fn operation_ownership_hides_intermediate_move_states_but_not_ordinary_live_work() {
75        for state in [
76            SessionState::Stopped,
77            SessionState::Running,
78            SessionState::Disconnected,
79        ] {
80            assert_eq!(
81                SessionTransitionKind::for_session(state, Some(SessionTransitionKind::Moving)),
82                Some(SessionTransitionKind::Moving)
83            );
84            assert_eq!(SessionTransitionKind::for_session(state, None), None);
85        }
86        assert_eq!(SessionState::Checkpointing.transition_kind(), None);
87        assert_eq!(
88            SessionState::Closing.transition_kind(),
89            Some(SessionTransitionKind::Suspending)
90        );
91    }
92}
93
94/// Controller-owned execution state derived from the relay event stream.
95#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
96#[serde(tag = "state", rename_all = "snake_case")]
97pub enum MaterializedExecutionState {
98    #[default]
99    Idle,
100    Running {
101        started_at_ms: i64,
102    },
103    Closing,
104    Closed,
105}
106
107pub use crate::transcript::{TerminalOutputRecord, TranscriptBody, TranscriptItem};
108
109/// What a durable queue entry does when its turn comes.
110///
111/// Serialized without a tag for prompts so entries written before configuration
112/// changes could be queued keep loading unchanged.
113#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
114#[serde(rename_all = "snake_case")]
115pub enum QueuedCommandKind {
116    #[default]
117    Prompt,
118    SetConfig {
119        key: String,
120        value: String,
121    },
122}
123
124impl QueuedCommandKind {
125    pub fn is_prompt(&self) -> bool {
126        matches!(self, Self::Prompt)
127    }
128}
129
130/// The composer form of a configuration change, used both as the queue entry's
131/// display text and as the text peeled back into the composer for editing.
132pub fn config_command_text(key: &str, value: &str) -> String {
133    if key == "fast-mode" {
134        "/fast".to_owned()
135    } else {
136        format!("/{key} {value}")
137    }
138}
139
140#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
141#[serde(deny_unknown_fields)]
142pub struct MaterializedQueuedPrompt {
143    pub command_id: String,
144    #[serde(default, skip_serializing_if = "QueuedCommandKind::is_prompt")]
145    pub kind: QueuedCommandKind,
146    pub content: Vec<serde_json::Value>,
147    pub queued_at_ms: i64,
148    /// Relay acceptance ordinal of the `CommandQueued` event that created this
149    /// entry. It is the turn identity the API hands back to callers, so wait
150    /// can tell one queued prompt's outcome from another's.
151    #[serde(default, skip_serializing_if = "Option::is_none")]
152    pub accepted_ordinal: Option<u64>,
153}
154
155/// The prompt currently executing, recorded when its `CommandStarted` event is
156/// projected and cleared when the command completes.
157#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
158#[serde(deny_unknown_fields)]
159pub struct MaterializedTurn {
160    pub command_id: String,
161    #[serde(default, skip_serializing_if = "Option::is_none")]
162    pub accepted_ordinal: Option<u64>,
163    /// Ordinal of the `CommandStarted` event, which is also the transcript
164    /// position of the turn's first item.
165    pub turn_start_position: u64,
166    pub started_at_ms: i64,
167}
168
169/// How a prompt ended.
170#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
171#[serde(tag = "kind", rename_all = "snake_case")]
172pub enum TurnOutcomeKind {
173    /// The harness finished the turn and reported this stop reason.
174    Completed { stop_reason: String },
175    /// The relay refused the command before it ran.
176    Rejected { message: String },
177    /// The command was interrupted after being accepted.
178    Interrupted { message: String },
179}
180
181#[derive(Debug, Clone, Copy, PartialEq, Eq)]
182pub enum PromptCompletion {
183    InputRequired,
184    Finished,
185    Cancelled,
186    QuotaLimit,
187    Error,
188}
189
190/// Shared interpretation for wait responses and durable completion events.
191pub fn classify_prompt_completion(stop_reason: &str) -> PromptCompletion {
192    let normalized = stop_reason
193        .chars()
194        .filter(|character| *character != '_' && *character != '-')
195        .flat_map(char::to_lowercase)
196        .collect::<String>();
197    match normalized.as_str() {
198        "endturn" => PromptCompletion::Finished,
199        "awaitinginput" => PromptCompletion::InputRequired,
200        "cancelled" | "canceled" => PromptCompletion::Cancelled,
201        "quotalimit" => PromptCompletion::QuotaLimit,
202        _ if crate::relay::is_capacity_stop_reason(stop_reason) => PromptCompletion::QuotaLimit,
203        _ => PromptCompletion::Error,
204    }
205}
206
207/// The most recent finished prompt on a session.
208#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
209#[serde(deny_unknown_fields)]
210pub struct MaterializedTurnOutcome {
211    #[serde(default, skip_serializing_if = "Option::is_none")]
212    pub diagnostic: Option<crate::diagnostic::TurnDiagnostic>,
213
214    #[serde(default, skip_serializing_if = "Option::is_none")]
215    pub usage: Option<crate::usage::TokenUsage>,
216    pub command_id: String,
217    #[serde(default, skip_serializing_if = "Option::is_none")]
218    pub accepted_ordinal: Option<u64>,
219    #[serde(default, skip_serializing_if = "Option::is_none")]
220    pub turn_start_position: Option<u64>,
221    pub completed_ordinal: u64,
222    pub completed_at_ms: i64,
223    pub outcome: TurnOutcomeKind,
224}
225
226impl MaterializedTurnOutcome {
227    /// Only work that actually started can have been interrupted.
228    pub fn interruption_ordinal(&self) -> Option<u64> {
229        (self.turn_start_position.is_some()
230            && matches!(self.outcome, TurnOutcomeKind::Interrupted { .. }))
231        .then_some(self.completed_ordinal)
232    }
233}
234
235/// Canonical controller projection for one logical ACP session.
236#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
237#[serde(deny_unknown_fields)]
238pub struct MaterializedSession {
239    pub session_id: String,
240    pub applied_event_ordinal: u64,
241    pub applied_event_digest: String,
242    /// Monotonic controller projection watermark derived from relay event
243    /// receipt times. It is deliberately independent of retained rows.
244    pub last_activity_at_ms: Option<i64>,
245    pub execution: MaterializedExecutionState,
246    #[serde(default, skip_serializing_if = "Option::is_none")]
247    pub session_title: Option<String>,
248    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
249    pub configuration: BTreeMap<String, serde_json::Value>,
250    #[serde(default, skip_serializing_if = "Vec::is_empty")]
251    /// Transcript items are shared by pointer so cloning a snapshot copies
252    /// handles rather than the whole conversation.
253    pub transcript: Vec<Arc<TranscriptItem>>,
254    #[serde(default, skip_serializing_if = "Vec::is_empty")]
255    pub queued_prompts: Vec<MaterializedQueuedPrompt>,
256    /// In-flight form requests are projected durably, but their answers are
257    /// connection-only and never enter this state.
258    #[serde(default, skip_serializing_if = "Vec::is_empty")]
259    pub pending_elicitations: Vec<crate::elicitation::ElicitationRequest>,
260    /// The prompt running right now, if any.
261    #[serde(default, skip_serializing_if = "Option::is_none")]
262    pub active_turn: Option<MaterializedTurn>,
263    /// The most recently finished prompt, kept after the session stops so a
264    /// caller can still read how the last turn ended.
265    #[serde(default, skip_serializing_if = "Option::is_none")]
266    pub last_turn_outcome: Option<MaterializedTurnOutcome>,
267}
268
269/// The small portion of a durable projection needed to populate dashboard
270/// rows before the live session delivers its full transcript snapshot.
271#[derive(Debug, Clone, PartialEq, Eq)]
272pub struct MaterializedSessionSummary {
273    pub session_id: String,
274    pub applied_event_ordinal: u64,
275    pub last_activity_at_ms: Option<i64>,
276    pub execution: MaterializedExecutionState,
277    pub session_title: Option<String>,
278    pub last_agent_message: Option<String>,
279    pub last_user_message: Option<String>,
280    /// Whether the last nonempty agent message appears after the last
281    /// nonempty user message in transcript order.
282    pub last_agent_message_follows_last_user: bool,
283    pub agent_message_latest_content_ordinals: Vec<u64>,
284    pub interruption_event_ordinals: Vec<u64>,
285}
286
287impl MaterializedSession {
288    pub fn empty(session_id: impl Into<String>) -> Self {
289        Self {
290            session_id: session_id.into(),
291            applied_event_ordinal: 0,
292            applied_event_digest: RELAY_EVENT_GENESIS_DIGEST.into(),
293            last_activity_at_ms: None,
294            execution: MaterializedExecutionState::Idle,
295            session_title: None,
296            configuration: BTreeMap::new(),
297            transcript: Vec::new(),
298            queued_prompts: Vec::new(),
299            pending_elicitations: Vec::new(),
300            active_turn: None,
301            last_turn_outcome: None,
302        }
303    }
304
305    pub fn last_activity_at_ms(&self) -> Option<i64> {
306        self.last_activity_at_ms
307    }
308
309    /// Resolve the title exposed by a live materialized session.
310    ///
311    /// Sessions created before provisional titles were projected can still
312    /// have an untitled transcript. Derive the same bounded fallback from
313    /// their first visible user prompt when reading them.
314    pub fn resolved_title(&self) -> Option<String> {
315        self.session_title
316            .as_deref()
317            .and_then(normalize_session_title)
318            .or_else(|| {
319                self.transcript.iter().find_map(|item| {
320                    let TranscriptBody::User { content } = &item.body else {
321                        return None;
322                    };
323                    provisional_session_title(&crate::transcript::materialized_content_text(
324                        content,
325                    ))
326                })
327            })
328            .or_else(|| {
329                self.queued_prompts
330                    .iter()
331                    .filter(|prompt| prompt.kind.is_prompt())
332                    .find_map(|prompt| {
333                        provisional_session_title(&crate::transcript::materialized_content_text(
334                            &prompt.content,
335                        ))
336                    })
337            })
338    }
339
340    pub fn unread_agent_messages_after(&self, viewed_through_event_ordinal: u64) -> u64 {
341        self.transcript
342            .iter()
343            .filter(|item| {
344                item.latest_content_event_ordinal
345                    .is_some_and(|ordinal| ordinal > viewed_through_event_ordinal)
346                    && item.is_nonempty_agent_message()
347            })
348            .count() as u64
349    }
350
351    pub fn unread_interruptions_after(&self, viewed_through_event_ordinal: u64) -> u64 {
352        self.interruption_event_ordinals()
353            .into_iter()
354            .filter(|ordinal| *ordinal > viewed_through_event_ordinal)
355            .count() as u64
356    }
357
358    pub fn interruption_event_ordinals(&self) -> Vec<u64> {
359        let mut ordinals = self
360            .transcript
361            .iter()
362            .filter(|item| item.is_work_interruption())
363            .map(|item| item.position)
364            .collect::<Vec<_>>();
365        if let Some(ordinal) = self
366            .last_turn_outcome
367            .as_ref()
368            .and_then(MaterializedTurnOutcome::interruption_ordinal)
369        {
370            ordinals.push(ordinal);
371        }
372        ordinals.sort_unstable();
373        ordinals.dedup();
374        ordinals
375    }
376
377    pub fn validate(&self) -> Result<()> {
378        validate_id("session", &self.session_id)?;
379        validate_relay_event_frontier(
380            self.applied_event_ordinal,
381            &self.applied_event_digest,
382            "materialized session event frontier",
383        )?;
384        if self
385            .session_title
386            .as_ref()
387            .is_some_and(|title| title.trim().is_empty())
388        {
389            bail!("materialized session has an empty title");
390        }
391        let mut item_ids = BTreeSet::new();
392        for item in &self.transcript {
393            item.validate(self.applied_event_ordinal)?;
394            if !item_ids.insert(item.stable_id.as_str()) {
395                bail!(
396                    "materialized transcript contains duplicate item {:?}",
397                    item.stable_id
398                );
399            }
400        }
401        let mut command_ids = BTreeSet::new();
402        for prompt in &self.queued_prompts {
403            if prompt.command_id.trim().is_empty() {
404                bail!("materialized prompt queue has an empty command id");
405            }
406            if !command_ids.insert(prompt.command_id.as_str()) {
407                bail!(
408                    "materialized prompt queue contains duplicate command {:?}",
409                    prompt.command_id
410                );
411            }
412            if let QueuedCommandKind::SetConfig { key, value } = &prompt.kind
413                && (key.trim().is_empty() || value.trim().is_empty())
414            {
415                bail!(
416                    "materialized queued configuration change {:?} is incomplete",
417                    prompt.command_id
418                );
419            }
420        }
421        Ok(())
422    }
423}
424
425/// A materialized session paired with the live worker's relay state. The
426/// session manager hands this to every reader that needs both the durable
427/// projection and the connection's operational status.
428#[derive(Debug, Clone, PartialEq)]
429pub struct ManagedSessionSnapshot {
430    pub materialized: MaterializedSession,
431    /// What `materialized.transcript` leaves out, and the facts that live
432    /// there. See [`ProjectionWindow`].
433    pub window: ProjectionWindow,
434    pub operational: RelayOperationalState,
435    /// Newest relay event observed by this live actor that asks for immediate
436    /// credential reconciliation. This is intentionally ephemeral: it avoids
437    /// retaining raw replay pages or rescanning projected history.
438    pub latest_credential_sync_signal: Option<CredentialSyncSignal>,
439    /// Content address of the executable the connected worker is running, as
440    /// it reported in hello. `None` when the connection did not come from a
441    /// live worker or the worker predates the field; either way the worker is
442    /// not known to be the build this controller would install.
443    pub worker_build: Option<String>,
444    /// Pending parent-tool work fetched from the target worker.
445    pub subagent_requests: Vec<crate::subagent::SubagentToolRequest>,
446    /// Recently completed tool work cached by the worker for idempotent calls.
447    pub subagent_results: Vec<crate::subagent::SubagentToolResult>,
448}
449
450/// What a projection's transcript window leaves out.
451///
452/// A polled projection carries only the end of the transcript, because that is
453/// all any viewer shows and loading the rest is work proportional to history.
454/// Two facts a reader needs live outside that window: the provisional title
455/// comes from the *first* user message, and the newest turn start is outside
456/// it whenever a single turn is longer than the window. Both are read
457/// separately, with one indexed query each, rather than found by scanning.
458///
459/// A complete projection answers both by scanning what it already holds, which
460/// is what [`ProjectionWindow::of`] does.
461#[derive(Debug, Clone, PartialEq, Eq)]
462pub struct ProjectionWindow {
463    /// Transcript items before the window. Zero when the projection is whole.
464    pub omitted_items: usize,
465    /// The title derived from the first user message.
466    pub provisional_title: Option<String>,
467    /// Position of the newest turn start — a user message or the marker for a
468    /// turn the harness began on its own — whether or not it is in the
469    /// window. `None` when the session has none.
470    pub latest_turn_start_position: Option<u64>,
471}
472
473impl ProjectionWindow {
474    /// The window of a projection that omits nothing.
475    #[must_use]
476    pub fn of(session: &MaterializedSession) -> Self {
477        Self {
478            omitted_items: 0,
479            provisional_title: session.transcript.iter().find_map(|item| {
480                let TranscriptBody::User { content } = &item.body else {
481                    return None;
482                };
483                provisional_session_title(&crate::transcript::materialized_content_text(content))
484            }),
485            latest_turn_start_position: session
486                .transcript
487                .iter()
488                .rev()
489                .find(|item| item.is_turn_start())
490                .map(|item| item.position),
491        }
492    }
493}
494
495impl ManagedSessionSnapshot {
496    /// The session's title, using the same precedence as
497    /// [`MaterializedSession::resolved_title`] but taking the provisional
498    /// title from the window rather than from a transcript head that a polled
499    /// projection does not carry.
500    #[must_use]
501    pub fn resolved_title(&self) -> Option<String> {
502        self.materialized
503            .session_title
504            .as_deref()
505            .and_then(normalize_session_title)
506            .or_else(|| self.window.provisional_title.clone())
507            .or_else(|| {
508                self.materialized
509                    .queued_prompts
510                    .iter()
511                    .filter(|prompt| prompt.kind.is_prompt())
512                    .find_map(|prompt| {
513                        provisional_session_title(&crate::transcript::materialized_content_text(
514                            &prompt.content,
515                        ))
516                    })
517            })
518    }
519
520    /// The position of the turn this session most recently finished, or `None`
521    /// while it is still working. Same answer as
522    /// [`latest_completed_turn_ordinal`], from a position the window carries
523    /// rather than a scan back through the transcript.
524    #[must_use]
525    pub fn latest_completed_turn_ordinal(&self) -> Option<u64> {
526        if self.materialized.execution != MaterializedExecutionState::Idle {
527            return None;
528        }
529        self.window.latest_turn_start_position
530    }
531}
532
533/// One session's activity, reported to the recovery coordinator.
534#[derive(Debug, Clone)]
535pub struct RecoveryObservation {
536    pub session: SessionRecord,
537    pub config: Config,
538    pub latest_completed_turn_ordinal: Option<u64>,
539    pub execution: MaterializedExecutionState,
540    /// Whether live provider-owned work permits an automatic checkpoint now.
541    /// This is separate from materialized execution because Kimi detached
542    /// agents outlive the parent turn that returned the session to `Idle`.
543    pub checkpoint_safe: bool,
544}
545
546/// The position where the session's most recent finished turn began, or
547/// `None` while it is still working. A turn starts at a user message or at the
548/// marker for a turn the harness began on its own, so autonomous work is
549/// covered once it settles.
550pub fn latest_completed_turn_ordinal(session: &MaterializedSession) -> Option<u64> {
551    if session.execution != MaterializedExecutionState::Idle {
552        return None;
553    }
554    session
555        .transcript
556        .iter()
557        .rev()
558        .find(|item| item.is_turn_start())
559        .map(|item| item.position)
560}
561
562pub fn validate_relay_event_digest(digest: &str, name: &str) -> Result<()> {
563    if digest.len() != 64
564        || !digest
565            .bytes()
566            .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
567    {
568        bail!("{name} must be a lowercase SHA-256 digest");
569    }
570    Ok(())
571}
572
573pub fn validate_relay_event_frontier(ordinal: u64, digest: &str, name: &str) -> Result<()> {
574    validate_relay_event_digest(digest, name)?;
575    if (ordinal == 0) != (digest == RELAY_EVENT_GENESIS_DIGEST) {
576        bail!("{name} has inconsistent ordinal {ordinal} and digest {digest}");
577    }
578    Ok(())
579}
580
581fn is_false(value: &bool) -> bool {
582    !*value
583}
584
585impl SessionState {
586    /// The persisted and wire spelling, matching the serde encoding.
587    pub const fn as_str(self) -> &'static str {
588        match self {
589            Self::Provisioning => "provisioning",
590            Self::Running => "running",
591            Self::Disconnected => "disconnected",
592            Self::Checkpointing => "checkpointing",
593            Self::Closing => "closing",
594            Self::Destroying => "destroying",
595            Self::Stopped => "stopped",
596            Self::Lost => "lost",
597            Self::Error => "error",
598            Self::DestroyedWithDataLoss => "destroyed-with-data-loss",
599        }
600    }
601
602    /// Read a stored spelling. Rows written before the verb was renamed still
603    /// say `"archived"`.
604    pub fn from_stored(value: &str) -> Option<Self> {
605        Some(match value {
606            "provisioning" => Self::Provisioning,
607            "running" => Self::Running,
608            "disconnected" => Self::Disconnected,
609            "checkpointing" => Self::Checkpointing,
610            "closing" => Self::Closing,
611            "destroying" => Self::Destroying,
612            "stopped" | "archived" => Self::Stopped,
613            "lost" => Self::Lost,
614            "error" => Self::Error,
615            "destroyed-with-data-loss" => Self::DestroyedWithDataLoss,
616            _ => return None,
617        })
618    }
619
620    /// Recovery without a live operation still hides an unfinished target transition.
621    /// Ordinary checkpoints and reconnects deliberately keep their conversation visible.
622    pub const fn transition_kind(self) -> Option<SessionTransitionKind> {
623        match self {
624            Self::Provisioning => Some(SessionTransitionKind::Starting),
625            Self::Closing => Some(SessionTransitionKind::Suspending),
626            Self::Destroying => Some(SessionTransitionKind::Destroying),
627            _ => None,
628        }
629    }
630
631    /// True while the session still belongs on the dashboard. `Closing` and
632    /// `Checkpointing` stay active on purpose: a stop that has not produced a
633    /// verified checkpoint must not make its row disappear.
634    pub const fn is_active(self) -> bool {
635        matches!(
636            self,
637            Self::Provisioning
638                | Self::Running
639                | Self::Disconnected
640                | Self::Checkpointing
641                | Self::Closing
642                | Self::Destroying
643                | Self::Error
644        )
645    }
646}
647
648#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
649#[serde(tag = "kind", rename_all = "kebab-case")]
650pub enum PodmanWorkspaceLocator {
651    #[default]
652    ContainerLayer,
653    Volume {
654        name: String,
655    },
656    HostPath {
657        path: PathBuf,
658        helper: Vec<String>,
659        resource: String,
660    },
661}
662
663#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
664#[serde(tag = "kind", rename_all = "kebab-case")]
665pub enum TargetLocator {
666    LocalBare {
667        worker_root: PathBuf,
668    },
669    LocalPodman {
670        container_id: String,
671        #[serde(default)]
672        workspace_storage: PodmanWorkspaceLocator,
673        /// The session that owns the container when this locator is a
674        /// sub-agent child borrowing its parent's container; `None` when the
675        /// session owns the container itself.
676        #[serde(default, skip_serializing_if = "Option::is_none")]
677        borrowed_from: Option<String>,
678    },
679    LocalDocker {
680        container_id: String,
681        /// The session that owns the container when this locator is a
682        /// sub-agent child borrowing its parent's container; `None` when the
683        /// session owns the container itself.
684        #[serde(default, skip_serializing_if = "Option::is_none")]
685        borrowed_from: Option<String>,
686    },
687    AppleContainer {
688        container_id: String,
689        /// The session that owns the container when this locator is a
690        /// sub-agent child borrowing its parent's container; `None` when the
691        /// session owns the container itself.
692        #[serde(default, skip_serializing_if = "Option::is_none")]
693        borrowed_from: Option<String>,
694    },
695    AwsEc2 {
696        instance_id: String,
697        #[serde(default, skip_serializing_if = "Option::is_none")]
698        address: Option<String>,
699    },
700    SshBare {
701        host: String,
702        workspace: PathBuf,
703        #[serde(default, skip_serializing_if = "Option::is_none")]
704        worker_id: Option<String>,
705    },
706    SshPodman {
707        host: String,
708        container_id: String,
709        #[serde(default)]
710        workspace_storage: PodmanWorkspaceLocator,
711        /// The session that owns the container when this locator is a
712        /// sub-agent child borrowing its parent's container; `None` when the
713        /// session owns the container itself.
714        #[serde(default, skip_serializing_if = "Option::is_none")]
715        borrowed_from: Option<String>,
716    },
717    SshDocker {
718        host: String,
719        container_id: String,
720        /// The session that owns the container when this locator is a
721        /// sub-agent child borrowing its parent's container; `None` when the
722        /// session owns the container itself.
723        #[serde(default, skip_serializing_if = "Option::is_none")]
724        borrowed_from: Option<String>,
725    },
726}
727
728#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
729#[serde(tag = "kind", rename_all = "kebab-case")]
730pub enum ManagedWorktreeTarget {
731    Local,
732    Ssh {
733        destination: String,
734        #[serde(default, skip_serializing_if = "Vec::is_empty")]
735        ssh_args: Vec<String>,
736    },
737}
738
739/// Whether a selected project can create a session-owned Git checkout.
740#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
741#[serde(deny_unknown_fields)]
742pub struct ManagedWorktreeOptions {
743    pub available: bool,
744    pub default_create: bool,
745}
746
747#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
748#[serde(deny_unknown_fields)]
749pub struct ManagedWorktree {
750    pub source_project_directory: PathBuf,
751    pub source_repository: PathBuf,
752    pub worktree_root: PathBuf,
753    pub branch: String,
754    pub target: ManagedWorktreeTarget,
755    /// The commit the session branch was created at. Recorded so an export can
756    /// diff against it in one read; sessions created before this field existed
757    /// fall back to the branch reflog, which expires.
758    #[serde(default, skip_serializing_if = "Option::is_none")]
759    pub base_commit: Option<String>,
760}
761
762impl ManagedWorktree {
763    fn validate(&self, session_id: &str, project_directory: Option<&Path>) -> Result<()> {
764        for (label, path) in [
765            ("source project directory", &self.source_project_directory),
766            ("source repository", &self.source_repository),
767            ("worktree root", &self.worktree_root),
768        ] {
769            if !path.is_absolute() || path.components().any(|part| part == Component::ParentDir) {
770                bail!("managed worktree {label} must be an absolute safe path");
771            }
772        }
773        if !self
774            .source_project_directory
775            .starts_with(&self.source_repository)
776        {
777            bail!("managed worktree source directory is outside its repository");
778        }
779        let expected_root = self
780            .source_repository
781            .join(".mj")
782            .join("worktrees")
783            .join(session_id);
784        if self.worktree_root != expected_root {
785            bail!("managed worktree root does not match the session-owned path");
786        }
787        if self.branch != format!("mj/{session_id}") {
788            bail!("managed worktree branch does not match the session id");
789        }
790        let relative = self
791            .source_project_directory
792            .strip_prefix(&self.source_repository)
793            .expect("source relationship checked above");
794        if project_directory != Some(self.worktree_root.join(relative).as_path()) {
795            bail!("session project directory does not match its managed worktree");
796        }
797        match &self.target {
798            ManagedWorktreeTarget::Local => {}
799            ManagedWorktreeTarget::Ssh { destination, .. } if destination.trim().is_empty() => {
800                bail!("managed SSH worktree has an empty destination")
801            }
802            ManagedWorktreeTarget::Ssh { .. } => {}
803        }
804        Ok(())
805    }
806}
807
808#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
809#[serde(tag = "kind", rename_all = "kebab-case")]
810pub enum SessionResourceAllocation {
811    Container {
812        cpus: u64,
813        memory_bytes: u64,
814    },
815    AwsEc2 {
816        instance_type: String,
817        vcpus: u64,
818        memory_bytes: u64,
819    },
820}
821
822impl SessionResourceAllocation {
823    pub fn validate(&self) -> Result<()> {
824        match self {
825            Self::Container { cpus, memory_bytes } if *cpus == 0 || *memory_bytes == 0 => {
826                bail!("container resource allocation must have non-zero CPU and memory")
827            }
828            Self::AwsEc2 {
829                instance_type,
830                vcpus,
831                memory_bytes,
832            } if instance_type.trim().is_empty() || *vcpus == 0 || *memory_bytes == 0 => {
833                bail!("EC2 resource allocation must have an instance type, CPU, and memory")
834            }
835            _ => Ok(()),
836        }
837    }
838}
839
840/// The CPU count an allocation grants, regardless of target kind.
841pub fn allocation_cpus(allocation: &SessionResourceAllocation) -> u64 {
842    match allocation {
843        SessionResourceAllocation::Container { cpus, .. } => *cpus,
844        SessionResourceAllocation::AwsEc2 { vcpus, .. } => *vcpus,
845    }
846}
847
848/// The memory, in bytes, an allocation grants, regardless of target kind.
849pub fn allocation_memory(allocation: &SessionResourceAllocation) -> u64 {
850    match allocation {
851        SessionResourceAllocation::Container { memory_bytes, .. }
852        | SessionResourceAllocation::AwsEc2 { memory_bytes, .. } => *memory_bytes,
853    }
854}
855
856impl TargetLocator {
857    fn validate(&self, session_id: &str) -> Result<()> {
858        match self {
859            Self::LocalBare { worker_root } => {
860                if !worker_root.is_absolute()
861                    || worker_root
862                        .components()
863                        .any(|part| part == Component::ParentDir)
864                    || !worker_root.ends_with(session_id)
865                {
866                    bail!(
867                        "local bare worker root must be an absolute safe path ending in the session id"
868                    );
869                }
870            }
871            Self::LocalPodman { container_id, .. }
872            | Self::LocalDocker { container_id, .. }
873            | Self::AppleContainer { container_id, .. }
874            | Self::SshPodman { container_id, .. }
875            | Self::SshDocker { container_id, .. }
876                if container_id.trim().is_empty() =>
877            {
878                bail!("target locator has an empty container id")
879            }
880            Self::AwsEc2 { instance_id, .. } if instance_id.trim().is_empty() => {
881                bail!("target locator has an empty AWS instance id")
882            }
883            Self::SshBare {
884                host, workspace, ..
885            } => {
886                if host.trim().is_empty() {
887                    bail!("bare SSH target locator has an empty host");
888                }
889                if workspace.as_os_str().is_empty()
890                    || workspace
891                        .components()
892                        .any(|part| part == Component::ParentDir)
893                    || !workspace.ends_with(session_id)
894                {
895                    bail!("bare SSH target locator must be a safe path ending in the session id");
896                }
897            }
898            Self::SshPodman { host, .. } if host.trim().is_empty() => {
899                bail!("SSH Podman target locator has an empty host")
900            }
901            Self::SshDocker { host, .. } if host.trim().is_empty() => {
902                bail!("SSH Docker target locator has an empty host")
903            }
904            _ => {}
905        }
906        Ok(())
907    }
908}
909
910#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
911#[serde(deny_unknown_fields)]
912pub struct CheckpointMetadata {
913    pub archive_path: PathBuf,
914    /// Lowercase SHA-256 digest of the verified archive.
915    pub sha256: String,
916    pub created_at: String,
917    pub event_frontier: u64,
918}
919
920impl CheckpointMetadata {
921    fn validate(&self) -> Result<()> {
922        if self.archive_path.as_os_str().is_empty() {
923            bail!("checkpoint archive path is empty");
924        }
925        if self.sha256.len() != 64
926            || !self
927                .sha256
928                .bytes()
929                .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
930        {
931            bail!("checkpoint SHA-256 must be 64 lowercase hexadecimal characters");
932        }
933        if self.created_at.trim().is_empty() {
934            bail!("checkpoint timestamp is empty");
935        }
936        Ok(())
937    }
938}
939
940/// The mbx build cache a container session was provisioned with. The
941/// directory is a host path that is mounted read-write at the same absolute
942/// path inside the container.
943#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
944#[serde(deny_unknown_fields)]
945pub struct SessionBuildCache {
946    /// The container host this cache was resolved on. A session moved to a
947    /// different host cannot reuse it, so the decision is made again there.
948    pub host: String,
949    pub directory: PathBuf,
950    /// An mbx size string passed as `MBX_GC_MAX_TOTAL_SIZE`, or `None` when
951    /// the host's own mbx configuration file already carries the budget.
952    #[serde(default, skip_serializing_if = "Option::is_none")]
953    pub max_size: Option<String>,
954    /// A `[target] root` the host's mbx configuration relocates outside the
955    /// cache directory, mounted read-write at the same path as well.
956    #[serde(default, skip_serializing_if = "Option::is_none")]
957    pub target_root: Option<PathBuf>,
958}
959
960/// How much disk Mjolnir's own copies of sessions use, and how much an
961/// `archive_after_days` value would free. Settings shows this on the
962/// SessionWiki page so the effect of a value is visible before it is saved.
963#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
964pub struct ArchiveSpacePreview {
965    /// Every session record Mjolnir holds, archived or not.
966    pub sessions: usize,
967    /// What those sessions' checkpoints and attachments occupy.
968    pub bytes: u64,
969    /// The sessions an `archive_after_days` value would catch, and their
970    /// share of `bytes`. Both are zero when no value is set.
971    pub reclaimable_sessions: usize,
972    pub reclaimable_bytes: u64,
973}
974
975/// What a container target's host resolves for its blank build cache
976/// settings right now. Settings shows this beside each "automatic" field so
977/// the values a session would actually run with are visible before one starts.
978#[derive(Debug, Clone, PartialEq, Eq)]
979pub struct BuildCachePreview {
980    /// The host's own mbx version, or `None` when it has none on `PATH`.
981    pub native_mbx: Option<String>,
982    /// The cache directory sessions would mount, once known.
983    pub directory: Option<PathBuf>,
984    /// The budget sessions would run with, once known.
985    pub max_size: Option<BuildCacheLimit>,
986    /// What the cache on that host has done so far, when it has a tally.
987    pub stats: Option<BuildCacheStats>,
988    /// Why sessions on this target run without a cache, or `None` when they
989    /// share one.
990    pub off_reason: Option<BuildCacheOff>,
991}
992
993/// mbx's own running totals for one host's cache, read from the tally beside
994/// the store.
995///
996/// These are machine-wide and cumulative: every session on that host and any
997/// native builds the user ran themselves are counted together, since they
998/// share one cache. They answer whether the cache is being used at all, not
999/// what one session got out of it.
1000#[derive(Debug, Clone, PartialEq, Eq)]
1001pub struct BuildCacheStats {
1002    /// Builds that went through mbx. Zero means nothing has used the shim.
1003    pub builds: u64,
1004    /// Compilations answered from the cache instead of run.
1005    pub cached_compilations: u64,
1006    /// Compiler time those answers avoided, as mbx estimates it.
1007    pub avoided_compiler_ns: u64,
1008    /// Restored output bytes that were cloned rather than copied.
1009    pub reflinked_bytes: u64,
1010}
1011
1012/// Why a target's sessions run without the build cache.
1013#[derive(Debug, Clone, PartialEq, Eq)]
1014pub enum BuildCacheOff {
1015    /// This machine's own `enabled = false`.
1016    TurnedOff,
1017    /// Nothing on the machine's settings page can turn it on: the global
1018    /// switch, the host's mbx, or its filesystem. The text says which.
1019    Unavailable(String),
1020}
1021
1022impl std::fmt::Display for BuildCacheOff {
1023    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1024        match self {
1025            Self::TurnedOff => formatter.write_str("turned off for this machine"),
1026            Self::Unavailable(reason) => formatter.write_str(reason),
1027        }
1028    }
1029}
1030
1031/// Where a build cache session's size budget comes from.
1032#[derive(Debug, Clone, PartialEq, Eq)]
1033pub enum BuildCacheLimit {
1034    /// An mbx size string passed as `MBX_GC_MAX_TOTAL_SIZE`.
1035    Size(String),
1036    /// The host's own `~/.config/mbx/config.toml` carries the budget. The
1037    /// total it sets, when it sets one.
1038    HostConfiguration(Option<String>),
1039}
1040
1041#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1042#[serde(deny_unknown_fields)]
1043pub struct SessionRecord {
1044    pub id: String,
1045    /// Owning workspace while active, or the most recent workspace while inactive.
1046    ///
1047    /// Inactive histories are globally resumable, so this id may refer to a
1048    /// workspace that has since been deleted.
1049    #[serde(default = "default_session_workspace_id")]
1050    pub workspace_id: String,
1051    pub title: String,
1052    pub harness_kind: HarnessKind,
1053    pub last_profile: String,
1054    pub bundle_id: String,
1055    /// Existing project directory used directly by a local or SSH bare target.
1056    #[serde(default, skip_serializing_if = "Option::is_none")]
1057    pub project_directory: Option<PathBuf>,
1058    /// Git worktree created and owned by Hel for this raw-project session.
1059    #[serde(default, skip_serializing_if = "Option::is_none")]
1060    pub managed_worktree: Option<ManagedWorktree>,
1061    /// None preserves automatic selection; false uses the selected directory.
1062    #[serde(default, skip_serializing_if = "Option::is_none")]
1063    pub create_managed_worktree: Option<bool>,
1064    /// The Git revision the session was asked to start at, as the caller typed
1065    /// it. None starts at HEAD for a managed worktree and at the remote
1066    /// default branch for a bundle session. The resolved commit lands in
1067    /// `managed_worktree.base_commit` or in the clone's `mj.baseCommit`.
1068    #[serde(default, skip_serializing_if = "Option::is_none")]
1069    pub launch_base: Option<String>,
1070    /// None follows the global `[subagents] enabled` setting at launch time;
1071    /// Some(true) and Some(false) are explicit per-session choices.
1072    #[serde(default, skip_serializing_if = "Option::is_none")]
1073    pub mjolnir_subagents: Option<bool>,
1074    pub target_template_id: String,
1075    #[serde(default, skip_serializing_if = "Option::is_none")]
1076    pub resource_allocation: Option<SessionResourceAllocation>,
1077    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1078    pub additional_mounts: Vec<AdditionalMount>,
1079    /// Per-session container CPU limit that overrides the target template's
1080    /// value. It is applied the next time the container is created.
1081    #[serde(default, skip_serializing_if = "Option::is_none")]
1082    pub container_cpus: Option<String>,
1083    /// Per-session container memory limit that overrides the target
1084    /// template's value. It is applied the next time the container is created.
1085    #[serde(default, skip_serializing_if = "Option::is_none")]
1086    pub container_memory: Option<String>,
1087    /// In-container workspace root this session's repositories live under.
1088    /// `None` is a session whose container predates per-session workspaces and
1089    /// therefore keeps the shared legacy `/workspace`; every session created
1090    /// since records `/workspace/<session id>`, so two checkouts of one project
1091    /// on a host never share an absolute path.
1092    #[serde(default, skip_serializing_if = "Option::is_none")]
1093    pub container_workspace: Option<PathBuf>,
1094    /// The mbx build cache this session's container runs with, decided once at
1095    /// provisioning. `None` means the session runs without a build cache;
1096    /// resume, move, and sub-agent children reuse the recorded value.
1097    #[serde(default, skip_serializing_if = "Option::is_none")]
1098    pub build_cache: Option<SessionBuildCache>,
1099    pub state: SessionState,
1100    /// Legacy visibility preference, retained for record compatibility.
1101    /// Current surfaces do not hide sessions based on this flag.
1102    #[serde(default, skip_serializing_if = "is_false")]
1103    pub archived: bool,
1104    #[serde(default, skip_serializing_if = "Option::is_none")]
1105    pub target: Option<TargetLocator>,
1106    #[serde(default, skip_serializing_if = "Option::is_none")]
1107    pub native_session_id: Option<String>,
1108    #[serde(default, skip_serializing_if = "Option::is_none")]
1109    pub acp_session_title: Option<String>,
1110    #[serde(default, skip_serializing_if = "Option::is_none")]
1111    pub session_title_override: Option<String>,
1112    pub created_at: String,
1113    pub updated_at: String,
1114    #[serde(default, alias = "detached_after_event_ordinal")]
1115    pub viewed_through_event_ordinal: u64,
1116    /// Unsent chat input carried across a detach, so returning to a session
1117    /// restores what the user was typing. Empty means no draft.
1118    #[serde(default, skip_serializing_if = "String::is_empty")]
1119    pub draft_input: String,
1120    /// Why the last operation on this session failed.
1121    ///
1122    /// This is usually a raw controller error chain, which names profile
1123    /// homes, project paths and SSH hosts, so a public projection publishes it
1124    /// only for a session that is stopped or failed. The one exception is a
1125    /// sentence the controller composed for the person; see
1126    /// [`CLOSE_FAILURE_PREFIX`].
1127    #[serde(default, skip_serializing_if = "Option::is_none")]
1128    pub last_error: Option<String>,
1129    #[serde(default, skip_serializing_if = "Option::is_none")]
1130    pub last_checkpoint_error: Option<String>,
1131    #[serde(default, skip_serializing_if = "Option::is_none")]
1132    pub checkpoint: Option<CheckpointMetadata>,
1133}
1134
1135#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1136#[serde(deny_unknown_fields)]
1137pub struct HostContainerSize {
1138    pub cpus: u64,
1139    pub memory_bytes: u64,
1140}
1141
1142fn default_session_workspace_id() -> String {
1143    crate::workspace::DEFAULT_WORKSPACE_ID.to_owned()
1144}
1145
1146/// How a failed close's reason begins in [`SessionRecord::last_error`].
1147///
1148/// A close that fails is non-destructive: the session goes back to the state
1149/// it was running in. Its reason therefore has to be published for a live
1150/// session, which the raw error chains in the same field never are. The
1151/// controller composes a sentence for the person and tags it with this prefix,
1152/// and the projection reads the tag to tell the two apart. Written in one
1153/// place and read in one place, so the tag cannot drift.
1154pub const DESTRUCTION_FAILURE_PREFIX: &str = "the destruction did not finish";
1155
1156pub const CLOSE_FAILURE_PREFIX: &str = "the suspension did not finish";
1157
1158/// Recognize safe lifecycle outcomes, including records saved before the rename.
1159pub fn is_public_lifecycle_error(error: &str) -> bool {
1160    error.starts_with(CLOSE_FAILURE_PREFIX)
1161        || error.starts_with(DESTRUCTION_FAILURE_PREFIX)
1162        || error.starts_with("the close did not finish")
1163}
1164
1165/// How a target is named beside a session, wherever a surface shows one.
1166///
1167/// A bare target (`local-bare`, `ssh-bare`) opens a project directory directly,
1168/// so the target alone does not say what the session was working on and the
1169/// project's own folder name is appended. Every other kind names a provisioned
1170/// environment that already identifies itself, so the target id stands alone.
1171/// A target id the configuration no longer holds is shown verbatim, because its
1172/// kind is no longer known.
1173///
1174/// Shared so the live session summary and the Resume dialog's archived rows
1175/// cannot drift apart: [`SessionRecord::project_target`] calls this, and so
1176/// does the archived row built from the SessionWiki index.
1177#[must_use]
1178pub fn target_label(config: &Config, target_id: &str, project: Option<&Path>) -> String {
1179    if !matches!(
1180        config.targets.get(target_id),
1181        Some(TargetTemplate::LocalBare | TargetTemplate::SshBare { .. })
1182    ) {
1183        return target_id.to_owned();
1184    }
1185    project.and_then(Path::file_name).map_or_else(
1186        || target_id.to_owned(),
1187        |directory| format!("{target_id}/{}", directory.to_string_lossy()),
1188    )
1189}
1190
1191impl SessionRecord {
1192    /// The recorded failure that is safe to publish whatever state this
1193    /// session is in, because the controller wrote it for the person rather
1194    /// than copying an error chain into it.
1195    #[must_use]
1196    pub fn public_error(&self) -> Option<&str> {
1197        self.last_error
1198            .as_deref()
1199            .filter(|error| is_public_lifecycle_error(error))
1200    }
1201
1202    /// Configuration drift belongs to this session, not the entire controller.
1203    /// The diagnostic contains only public identifiers, so both UIs can show it.
1204    pub fn configuration_issue(&self, config: &Config) -> Option<String> {
1205        if !self.state.is_active() {
1206            return None;
1207        }
1208        let mut issues = Vec::new();
1209        match config.profiles.get(&self.last_profile) {
1210            None => issues.push(format!("missing profile {:?}", self.last_profile)),
1211            Some(profile) if profile.kind != self.harness_kind => issues.push(format!(
1212                "expects {:?}, but profile {:?} is {:?}",
1213                self.harness_kind, self.last_profile, profile.kind
1214            )),
1215            Some(_) => {}
1216        }
1217        if self.project_directory.is_none() && !config.bundles.contains_key(&self.bundle_id) {
1218            issues.push(format!("missing bundle {:?}", self.bundle_id));
1219        }
1220        if !config.targets.contains_key(&self.target_template_id) {
1221            issues.push(format!(
1222                "missing target template {:?}",
1223                self.target_template_id
1224            ));
1225        }
1226        (!issues.is_empty()).then(|| format!(
1227            "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.",
1228            self.id, issues.join("; ")
1229        ))
1230    }
1231
1232    pub fn validate_configuration(&self, config: &Config) -> Result<()> {
1233        if let Some(issue) = self.configuration_issue(config) {
1234            bail!("{issue}");
1235        }
1236        Ok(())
1237    }
1238
1239    /// User-visible session name, independent of the initial prompt stored in `title`.
1240    pub fn display_title(&self) -> &str {
1241        self.session_title_override
1242            .as_deref()
1243            .or(self.acp_session_title.as_deref())
1244            .unwrap_or(&self.id)
1245    }
1246
1247    /// Project this session works in, as the session list and the chat header
1248    /// both name it: the source repository of a managed worktree, else the
1249    /// project directory, else the bundle's primary repository, else the
1250    /// bundle id.
1251    pub fn project_name(&self, config: &Config) -> String {
1252        if let Some(worktree) = &self.managed_worktree {
1253            return path_leaf(&worktree.source_repository);
1254        }
1255        if let Some(project_directory) = &self.project_directory {
1256            return path_leaf(project_directory);
1257        }
1258        self.bundle_source_name(config)
1259    }
1260
1261    /// Target label used by the live session summary. Bare targets identify
1262    /// the project directory they open directly; workspace targets already
1263    /// identify the provisioned environment on their own.
1264    pub fn project_target(&self, config: &Config, target_id: &str) -> String {
1265        let project = self
1266            .managed_worktree
1267            .as_ref()
1268            .map(|worktree| &worktree.source_project_directory)
1269            .or(self.project_directory.as_ref());
1270        target_label(config, target_id, project.map(PathBuf::as_path))
1271    }
1272
1273    /// Stable source identity used to group sessions. Managed worktrees point
1274    /// back at their source repository, raw sessions use their project
1275    /// directory until their Git origin is resolved, and bundle sessions use
1276    /// their complete canonical repository set when configured.
1277    pub fn project_source(&self, config: &Config) -> ProjectSourceIdentity {
1278        if let Some(worktree) = &self.managed_worktree {
1279            return ProjectSourceIdentity::path(&worktree.source_repository, None);
1280        }
1281        if let Some(project_directory) = &self.project_directory {
1282            let remote = match &self.target {
1283                Some(TargetLocator::SshBare { host, .. }) => Some(host.as_str()),
1284                _ => None,
1285            };
1286            return ProjectSourceIdentity::path(project_directory, remote);
1287        }
1288        self.bundle_source_identity(config)
1289            .unwrap_or_else(|| ProjectSourceIdentity {
1290                key: format!("bundle:{}", self.bundle_id),
1291                short: path_leaf(Path::new(&self.bundle_id)),
1292                full: self.bundle_id.clone(),
1293            })
1294    }
1295
1296    /// Resolve the display name shared by session headings, chat headers, and
1297    /// resume details for a bundle-backed session.
1298    fn bundle_source_name(&self, config: &Config) -> String {
1299        self.bundle_source_identity(config)
1300            .map(|source| source.short)
1301            .unwrap_or_else(|| path_leaf(Path::new(&self.bundle_id)))
1302    }
1303
1304    /// Resolve the canonical identity of every repository in a bundle for
1305    /// grouping and display naming.
1306    fn bundle_source_identity(&self, config: &Config) -> Option<ProjectSourceIdentity> {
1307        let bundle = config.bundles.get(&self.bundle_id)?;
1308        let sources = bundle
1309            .repositories
1310            .iter()
1311            .map(repository_source_identity)
1312            .collect::<Option<Vec<_>>>()?;
1313        ProjectSourceIdentity::bundle(sources)
1314    }
1315
1316    /// Orders two sessions the way the session list's sequence view does:
1317    /// oldest first by creation time, with the id as a stable tiebreak. A
1318    /// session whose timestamp does not parse sorts last.
1319    pub fn compare_by_creation(&self, other: &Self) -> std::cmp::Ordering {
1320        self.creation_order_key().cmp(&other.creation_order_key())
1321    }
1322
1323    /// Parse once per session when used with `sort_by_cached_key`.
1324    pub fn creation_order_key(&self) -> (bool, Option<i64>, &str) {
1325        let timestamp = created_at_seconds(&self.created_at);
1326        (timestamp.is_none(), timestamp, &self.id)
1327    }
1328
1329    fn validate(&self, map_id: &str) -> Result<()> {
1330        validate_id("session", &self.id)?;
1331        if self.id != map_id {
1332            bail!(
1333                "session map key {map_id:?} does not match record id {:?}",
1334                self.id
1335            );
1336        }
1337        validate_id("workspace", &self.workspace_id)?;
1338        validate_id("profile", &self.last_profile)?;
1339        validate_id("bundle", &self.bundle_id)?;
1340        if let Some(project_directory) = &self.project_directory
1341            && (!project_directory.is_absolute()
1342                || project_directory
1343                    .components()
1344                    .any(|part| part == Component::ParentDir))
1345        {
1346            bail!("session {:?} has an unsafe project directory", self.id);
1347        }
1348        if let Some(managed_worktree) = &self.managed_worktree {
1349            managed_worktree.validate(&self.id, self.project_directory.as_deref())?;
1350        }
1351        validate_id("target template", &self.target_template_id)?;
1352        if let Some(allocation) = &self.resource_allocation {
1353            allocation.validate()?;
1354        }
1355        validate_additional_mounts(&self.additional_mounts)?;
1356        if self.title.trim().is_empty() {
1357            bail!("session {:?} has an empty title", self.id);
1358        }
1359        if self
1360            .acp_session_title
1361            .as_ref()
1362            .is_some_and(|title| title.trim().is_empty())
1363            || self
1364                .session_title_override
1365                .as_ref()
1366                .is_some_and(|title| title.trim().is_empty())
1367        {
1368            bail!("session {:?} has an empty display title", self.id);
1369        }
1370        if self.created_at.trim().is_empty() || self.updated_at.trim().is_empty() {
1371            bail!("session {:?} has an empty timestamp", self.id);
1372        }
1373        if let Some(target) = &self.target {
1374            target.validate(&self.id)?;
1375        }
1376        if let Some(checkpoint) = &self.checkpoint {
1377            checkpoint.validate()?;
1378        }
1379        Ok(())
1380    }
1381}
1382
1383fn repository_source_identity(repository: &ProjectRepository) -> Option<ProjectSourceIdentity> {
1384    repository
1385        .github
1386        .as_deref()
1387        .and_then(ProjectSourceIdentity::git_remote)
1388        .or_else(|| {
1389            repository
1390                .local
1391                .as_deref()
1392                .map(|path| ProjectSourceIdentity::path(path, None))
1393        })
1394}
1395
1396#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1397pub struct ProjectSourceIdentity {
1398    pub key: String,
1399    pub short: String,
1400    pub full: String,
1401}
1402
1403impl ProjectSourceIdentity {
1404    /// Combine repository identities into one stable bundle identity.
1405    pub fn bundle(mut sources: Vec<Self>) -> Option<Self> {
1406        if sources.is_empty() {
1407            return None;
1408        }
1409        sources.sort_by(|left, right| {
1410            left.key
1411                .cmp(&right.key)
1412                .then_with(|| left.full.cmp(&right.full))
1413                .then_with(|| left.short.cmp(&right.short))
1414        });
1415        sources.dedup_by(|left, right| left.key == right.key);
1416        if sources.len() == 1 {
1417            return sources.pop();
1418        }
1419        let keys = sources
1420            .iter()
1421            .map(|source| source.key.clone())
1422            .collect::<Vec<_>>();
1423        let key = serde_json::to_string(&keys).ok()?;
1424        Some(Self {
1425            key: format!("bundle:{key}"),
1426            short: sources
1427                .iter()
1428                .map(|source| source.short.as_str())
1429                .collect::<Vec<_>>()
1430                .join(" + "),
1431            full: sources
1432                .iter()
1433                .map(|source| source.full.as_str())
1434                .collect::<Vec<_>>()
1435                .join(" + "),
1436        })
1437    }
1438
1439    /// Canonicalizes a Git remote so raw checkouts group as the same project
1440    /// even when their worktree paths differ.
1441    pub fn git_remote(source: &str) -> Option<Self> {
1442        if let Some(normalized) = normalize_github_source(source) {
1443            let short = normalized
1444                .rsplit_once('/')
1445                .map_or(normalized.as_str(), |(_, repository)| repository)
1446                .to_owned();
1447            return Some(Self {
1448                key: format!("github:{}", normalized.to_lowercase()),
1449                short,
1450                full: normalized,
1451            });
1452        }
1453        let normalized = source.trim().trim_end_matches('/').trim_end_matches(".git");
1454        if normalized.is_empty() {
1455            return None;
1456        }
1457        let short = normalized
1458            .rsplit(['/', ':'])
1459            .find(|part| !part.is_empty())
1460            .unwrap_or(normalized)
1461            .to_owned();
1462        Some(Self {
1463            key: format!("git:{}", normalized.to_lowercase()),
1464            short,
1465            full: normalized.to_owned(),
1466        })
1467    }
1468
1469    /// Build a local-root identity, qualified by host for remote directories.
1470    pub fn path(path: &Path, remote: Option<&str>) -> Self {
1471        let normalized = path.components().collect::<PathBuf>();
1472        let path_text = normalized.to_string_lossy().into_owned();
1473        let full = remote.map_or_else(|| path_text.clone(), |host| format!("{host}:{path_text}"));
1474        let key = remote.map_or_else(
1475            || format!("path:{path_text}"),
1476            |host| format!("path:{}:{path_text}", host.to_lowercase()),
1477        );
1478        Self {
1479            key,
1480            short: path_leaf(path),
1481            full,
1482        }
1483    }
1484}
1485
1486fn normalize_github_source(source: &str) -> Option<String> {
1487    let source = source.trim();
1488    let path = source
1489        .strip_prefix("https://github.com/")
1490        .or_else(|| source.strip_prefix("http://github.com/"))
1491        .or_else(|| source.strip_prefix("git@github.com:"))
1492        .or_else(|| source.strip_prefix("ssh://git@github.com/"))
1493        .or_else(|| {
1494            (!source.contains("://") && !source.contains('@') && !source.contains(':'))
1495                .then_some(source)
1496        })?
1497        .trim_end_matches(".git");
1498    let mut parts = path.split('/');
1499    let owner = parts.next()?;
1500    let repository = parts.next()?;
1501    (!owner.is_empty() && !repository.is_empty() && parts.next().is_none())
1502        .then(|| format!("{owner}/{repository}"))
1503}
1504
1505/// Last component of a path, falling back to the whole path when it has none.
1506fn path_leaf(path: &Path) -> String {
1507    path.file_name()
1508        .unwrap_or(path.as_os_str())
1509        .to_string_lossy()
1510        .into_owned()
1511}
1512
1513fn created_at_seconds(timestamp: &str) -> Option<i64> {
1514    chrono::DateTime::parse_from_rfc3339(timestamp)
1515        .ok()
1516        .map(|timestamp| timestamp.timestamp())
1517}
1518
1519#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1520#[serde(deny_unknown_fields)]
1521pub struct State {
1522    pub version: u32,
1523    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1524    pub sessions: BTreeMap<String, SessionRecord>,
1525    /// Child sessions keyed by their session id. The relationship lives in
1526    /// controller state so every control surface sees the same session family.
1527    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1528    pub subagents: BTreeMap<String, SubagentRecord>,
1529    /// Recently used source directories, keyed by `local` or SSH host name.
1530    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1531    pub mount_history: BTreeMap<String, Vec<PathBuf>>,
1532    /// Most recently launched container size on each physical target host.
1533    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1534    pub container_sizes: BTreeMap<String, HostContainerSize>,
1535}
1536
1537impl Default for State {
1538    fn default() -> Self {
1539        Self {
1540            version: STATE_VERSION,
1541            sessions: BTreeMap::new(),
1542            subagents: BTreeMap::new(),
1543            mount_history: BTreeMap::new(),
1544            container_sizes: BTreeMap::new(),
1545        }
1546    }
1547}
1548
1549impl State {
1550    /// The session whose project identity names a row.
1551    ///
1552    /// A sub-agent child runs inside its parent's workspace and owns no
1553    /// managed worktree, so its own `project_directory` is the parent's
1554    /// worktree checkout, whose directory is named after the parent session
1555    /// id. Reading the project identity from the parent instead keeps a child
1556    /// under the same project heading and target label as the session it
1557    /// belongs to.
1558    #[must_use]
1559    pub fn project_identity_session<'a>(&'a self, session: &'a SessionRecord) -> &'a SessionRecord {
1560        self.subagents
1561            .get(&session.id)
1562            .and_then(|record| self.sessions.get(&record.parent_session_id))
1563            .unwrap_or(session)
1564    }
1565
1566    /// Whether `id` names a sub-agent rather than a session the user started:
1567    /// a Mjolnir-managed child, or the record a client builds to show a
1568    /// harness-owned child (see [`crate::native_agent::view_id`]).
1569    ///
1570    /// Every list of top-level sessions filters with this, so the lists
1571    /// cannot disagree about what a sub-agent is.
1572    #[must_use]
1573    pub fn is_subagent_session(&self, id: &str) -> bool {
1574        self.subagents.contains_key(id) || crate::native_agent::is_view_id(id)
1575    }
1576
1577    pub fn validate(&self) -> Result<()> {
1578        if self.version != STATE_VERSION {
1579            bail!(
1580                "unsupported Mjolnir state version {}; expected {STATE_VERSION}",
1581                self.version
1582            );
1583        }
1584        for (id, session) in &self.sessions {
1585            session.validate(id)?;
1586        }
1587        for (child_id, subagent) in &self.subagents {
1588            if child_id != &subagent.child_session_id {
1589                bail!("sub-agent key {child_id:?} does not match its child session id");
1590            }
1591            if child_id == &subagent.parent_session_id {
1592                bail!("sub-agent {child_id:?} cannot be its own parent");
1593            }
1594            if !self.sessions.contains_key(child_id) {
1595                bail!("sub-agent {child_id:?} has no child session");
1596            }
1597            if !self.sessions.contains_key(&subagent.parent_session_id) {
1598                bail!(
1599                    "sub-agent {child_id:?} has unknown parent {:?}",
1600                    subagent.parent_session_id
1601                );
1602            }
1603            if self.subagents.contains_key(&subagent.parent_session_id) {
1604                bail!("sub-agent {child_id:?} cannot belong to another sub-agent");
1605            }
1606            if subagent.task_name.trim().is_empty()
1607                || subagent.profile_id.trim().is_empty()
1608                || subagent.request_key.trim().is_empty()
1609            {
1610                bail!("sub-agent {child_id:?} has incomplete relationship metadata");
1611            }
1612        }
1613        for (host, sources) in &self.mount_history {
1614            if host.trim().is_empty() {
1615                bail!("mount history contains an empty host key");
1616            }
1617            if sources.iter().any(|source| !source.is_absolute()) {
1618                bail!("mount history for {host:?} contains a non-absolute source path");
1619            }
1620        }
1621        for (host, size) in &self.container_sizes {
1622            if host.trim().is_empty() {
1623                bail!("container size history contains an empty host key");
1624            }
1625            if size.cpus == 0 || size.memory_bytes == 0 {
1626                bail!("container size history for {host:?} contains a zero value");
1627            }
1628            if size.cpus > i64::MAX as u64 || size.memory_bytes > i64::MAX as u64 {
1629                bail!("container size history for {host:?} exceeds SQLite integer range");
1630            }
1631        }
1632        Ok(())
1633    }
1634
1635    pub fn remember_mount_sources(&mut self, host: &str, mounts: &[AdditionalMount]) {
1636        if mounts.is_empty() {
1637            return;
1638        }
1639        let sources = self.mount_history.entry(host.to_owned()).or_default();
1640        for mount in mounts.iter().rev() {
1641            sources.retain(|source| source != &mount.source);
1642            sources.insert(0, mount.source.clone());
1643        }
1644        sources.truncate(20);
1645    }
1646
1647    pub fn remember_container_size(&mut self, host: &str, size: HostContainerSize) {
1648        self.container_sizes.insert(host.to_owned(), size);
1649    }
1650
1651    pub fn project_directories(&self, host: &str) -> &[PathBuf] {
1652        self.mount_history
1653            .get(&project_history_key(host))
1654            .map(Vec::as_slice)
1655            .unwrap_or_default()
1656    }
1657
1658    pub fn remember_project_directory(&mut self, host: &str, directory: &Path) {
1659        let key = project_history_key(host);
1660        let directories = self.mount_history.entry(key).or_default();
1661        directories.retain(|existing| existing != directory);
1662        directories.insert(0, directory.to_path_buf());
1663        directories.truncate(20);
1664    }
1665
1666    pub fn destroy_stopped_session(&mut self, session_id: &str) -> Result<SessionRecord> {
1667        let session = self
1668            .sessions
1669            .get(session_id)
1670            .with_context(|| format!("unknown session {session_id}"))?;
1671        if session.state.is_active() {
1672            bail!("refusing to destroy active session {session_id}");
1673        }
1674        Ok(self
1675            .sessions
1676            .remove(session_id)
1677            .expect("session checked above"))
1678    }
1679
1680    /// Remove a session record from state regardless of its lifecycle state.
1681    ///
1682    /// Force destruction is the one caller: by the time it runs, every
1683    /// external artifact has been torn down or its loss accepted, so no state
1684    /// is refused here.
1685    pub fn destroy_session_force(&mut self, session_id: &str) -> Result<SessionRecord> {
1686        self.sessions
1687            .get(session_id)
1688            .with_context(|| format!("unknown session {session_id}"))?;
1689        Ok(self
1690            .sessions
1691            .remove(session_id)
1692            .expect("session checked above"))
1693    }
1694
1695    /// Setup may add replacements under new names, but must not rewrite
1696    /// dependencies still owned by active sessions.
1697    pub fn validate_setup_update(&self, before: &Config, after: &Config) -> Result<()> {
1698        for session in self
1699            .sessions
1700            .values()
1701            .filter(|session| session.state.is_active())
1702        {
1703            let protected = if let Some(profile) = before.profiles.get(&session.last_profile) {
1704                let mut comparable = profile.clone();
1705                if let Some(updated) = after.profiles.get(&session.last_profile) {
1706                    comparable.enabled = updated.enabled;
1707                }
1708                // A mismatched harness is already broken; allow repairing it.
1709                profile.kind == session.harness_kind
1710                    && after.profiles.get(&session.last_profile) != Some(&comparable)
1711            } else {
1712                false
1713            };
1714            let bundle_changed = session.project_directory.is_none()
1715                && before
1716                    .bundles
1717                    .get(&session.bundle_id)
1718                    .is_some_and(|bundle| after.bundles.get(&session.bundle_id) != Some(bundle));
1719            // Build cache settings are resolved at provisioning time and kept
1720            // on the session record, so editing them does not disturb a
1721            // running session.
1722            let target_changed =
1723                before
1724                    .targets
1725                    .get(&session.target_template_id)
1726                    .is_some_and(|target| {
1727                        after
1728                            .targets
1729                            .get(&session.target_template_id)
1730                            .map(TargetTemplate::without_launch_only_settings)
1731                            != Some(target.without_launch_only_settings())
1732                    });
1733            if protected || bundle_changed || target_changed {
1734                bail!(
1735                    "Setup would change configuration used by active session {:?}. Keep its profile {:?}, bundle {:?}, and target {:?}; add a separate entry for new settings, or stop the session before editing its configuration.",
1736                    session.id,
1737                    session.last_profile,
1738                    session.bundle_id,
1739                    session.target_template_id
1740                );
1741            }
1742        }
1743        Ok(())
1744    }
1745
1746    /// Strict validation for callers that need all active references intact.
1747    pub fn validate_against_config(&self, config: &Config) -> Result<()> {
1748        self.validate()?;
1749        config.validate()?;
1750        for session in self.sessions.values() {
1751            session.validate_configuration(config)?;
1752        }
1753        Ok(())
1754    }
1755}
1756
1757fn project_history_key(host: &str) -> String {
1758    format!("project:{host}")
1759}
1760
1761/// Generate an opaque, filesystem-safe stable id for a new logical session.
1762pub fn new_session_id() -> Result<String> {
1763    let mut random = [0u8; 16];
1764    getrandom::fill(&mut random)
1765        .map_err(|error| anyhow::anyhow!("generate Mjolnir session id: {error}"))?;
1766    Ok(crate::hex::lower_hex(random))
1767}
1768
1769/// Return the newest clean ACP session title from canonical worker events.
1770pub fn harness_session_title(events: &[SequencedEvent]) -> Option<String> {
1771    events.iter().rev().find_map(|event| {
1772        let WorkerEvent::Adapter { payload, .. } = &event.event else {
1773            return None;
1774        };
1775        let crate::acp::RuntimeEvent::SessionUpdate { update } =
1776            serde_json::from_value(payload.clone()).ok()?
1777        else {
1778            return None;
1779        };
1780        let kind = update
1781            .get("sessionUpdate")
1782            .and_then(serde_json::Value::as_str)?;
1783        let title = match kind {
1784            "session_info_update" | "session_title" => {
1785                update.get("title").and_then(serde_json::Value::as_str)
1786            }
1787            _ => None,
1788        }?;
1789        normalize_session_title(title)
1790    })
1791}
1792
1793pub fn normalize_session_title(title: &str) -> Option<String> {
1794    let normalized = crate::relay::strip_hidden_prompt_context(title)
1795        .split_whitespace()
1796        .collect::<Vec<_>>()
1797        .join(" ");
1798    (!normalized.is_empty()).then_some(normalized)
1799}
1800
1801/// Build the short-lived title shown before the harness supplies its own.
1802///
1803/// The first visible user prompt is immediately useful for identifying a
1804/// session, but it can be arbitrarily large. Keep this fallback bounded; a
1805/// later ACP session-info update remains authoritative and replaces it.
1806pub fn provisional_session_title(prompt: &str) -> Option<String> {
1807    const MAX_TITLE_CHARS: usize = 64;
1808
1809    let normalized = normalize_session_title(prompt)?;
1810    if normalized.chars().count() <= MAX_TITLE_CHARS {
1811        return Some(normalized);
1812    }
1813
1814    let mut truncated = normalized
1815        .chars()
1816        .take(MAX_TITLE_CHARS - 1)
1817        .collect::<String>();
1818    if let Some(boundary) = truncated.rfind(char::is_whitespace) {
1819        truncated.truncate(boundary);
1820    }
1821    truncated.push('…');
1822    Some(truncated)
1823}
1824
1825pub fn short_id(id: &str) -> &str {
1826    id.get(..8).unwrap_or(id)
1827}
1828
1829#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
1830pub struct RecoveryCandidate {
1831    pub session_id: String,
1832    pub target_template_id: String,
1833    pub locator: TargetLocator,
1834    pub ownership: Option<crate::worker_launch::WorkerOwnership>,
1835    /// Instance that created the worker, from its label or tag, else from
1836    /// the ownership marker. `None` means an older build left no stamp.
1837    #[serde(default)]
1838    pub instance_id: Option<String>,
1839    /// State of the session this resource is labelled for, when the
1840    /// controller still tracks that session. A leftover resource the session
1841    /// record no longer names can only be destroyed, never adopted, because
1842    /// the session id is already taken.
1843    #[serde(default, skip_serializing_if = "Option::is_none")]
1844    pub tracked_session: Option<SessionState>,
1845}
1846
1847#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
1848pub struct RecoveryScan {
1849    pub candidates: Vec<RecoveryCandidate>,
1850    pub warnings: Vec<String>,
1851    /// Identity of the instance that ran the scan.
1852    #[serde(default)]
1853    pub instance_id: String,
1854    /// Candidates left out because another or an unknown instance created
1855    /// them and the scan was not widened to all instances.
1856    #[serde(default)]
1857    pub hidden_other_instances: usize,
1858}
1859
1860#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1861#[serde(deny_unknown_fields)]
1862pub struct ResumeRepositorySourceReceipt {
1863    pub session_id: String,
1864    pub bundle_id: String,
1865    pub checkpoint_sha256: String,
1866    pub repositories: Vec<crate::config::ProjectRepository>,
1867}
1868
1869#[cfg(test)]
1870mod tests;