Skip to main content

mj_core/
storage.rs

1//! Data exchanged with the controller store, independent of its SQLite implementation.
2use crate::elicitation::ElicitationRequest;
3use crate::state::*;
4use crate::usage::ProviderCost;
5use serde::{Deserialize, Serialize};
6use std::collections::BTreeMap;
7use std::sync::Arc;
8/// A deterministic projection integrity violation. Retrying cannot fix it, so
9/// callers must report it separately from transport failures.
10#[derive(Debug)]
11pub struct ProjectionIntegrityError(pub String);
12
13impl std::fmt::Display for ProjectionIntegrityError {
14    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
15        formatter.write_str(&self.0)
16    }
17}
18
19impl std::error::Error for ProjectionIntegrityError {}
20
21/// A store this build cannot safely read and write.
22///
23/// Carried as a typed cause rather than a message so the daemon can tell a
24/// store that moved underneath it from a transport failure. It survives every
25/// `anyhow` hop to the caller, which finds it with `error.chain()`.
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27pub struct StoreSchemaMismatch {
28    pub found: i64,
29    pub supported: i64,
30    pub reason: StoreSchemaMismatchReason,
31}
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq)]
34pub enum StoreSchemaMismatchReason {
35    NeedsMigration,
36    Incompatible { minimum_compatible: i64 },
37    InvalidCompatibilityMetadata,
38    Rollback { previous: i64 },
39}
40
41impl std::fmt::Display for StoreSchemaMismatch {
42    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
43        let Self {
44            found,
45            supported,
46            reason,
47        } = self;
48        match reason {
49            StoreSchemaMismatchReason::Incompatible { minimum_compatible } => write!(
50                formatter,
51                "Mjolnir database schema {found} requires at least build schema {minimum_compatible} for reads and writes; this build supports {supported}; upgrade Mjolnir, run this build with an isolated data directory (--instance NAME or MJ_DATA_DIR), or restore a backup made by the older build"
52            ),
53            StoreSchemaMismatchReason::NeedsMigration => write!(
54                formatter,
55                "Mjolnir database schema {found} is not the supported schema {supported}; start the Mjolnir daemon to migrate it"
56            ),
57            StoreSchemaMismatchReason::InvalidCompatibilityMetadata => write!(
58                formatter,
59                "Mjolnir database schema {found} has missing or invalid compatibility metadata; refusing access from build schema {supported}"
60            ),
61            StoreSchemaMismatchReason::Rollback { previous } => write!(
62                formatter,
63                "Mjolnir database schema rolled back from {previous} to {found} underneath this writer; refusing writes from build schema {supported}"
64            ),
65        }
66    }
67}
68
69impl std::error::Error for StoreSchemaMismatch {}
70
71#[derive(Debug, Clone, Copy, PartialEq, Eq)]
72pub enum HistoryScope {
73    Project,
74    Session,
75    All,
76}
77
78#[derive(Debug, Clone, PartialEq, Eq)]
79pub struct PromptHistoryEntry {
80    pub id: i64,
81    pub session_id: String,
82    pub text: String,
83}
84
85#[derive(Debug, Clone, Copy, PartialEq, Eq)]
86pub enum ProjectionApplyOutcome {
87    Applied,
88    AlreadyApplied,
89}
90
91#[derive(Debug, Clone, PartialEq)]
92pub enum TranscriptMutation {
93    Upsert(TranscriptItem),
94    Remove { stable_id: String },
95}
96
97/// Changes derived from one relay event. `None` leaves a scalar untouched;
98/// the nested option on `session_title` permits explicitly clearing it.
99#[derive(Debug, Clone, PartialEq, Default)]
100pub struct MaterializedSessionMutation {
101    pub native_agent: Option<crate::relay::RelayEvent>,
102    /// Relay receipt time for this event. Persistence and the actor cache both
103    /// take a monotonic maximum so removing detail rows cannot move activity
104    /// backwards.
105    pub last_activity_at_ms: Option<i64>,
106    pub execution: Option<MaterializedExecutionState>,
107    pub session_title: Option<Option<String>>,
108    pub configuration: Option<SessionConfiguration>,
109    pub transcript: Vec<TranscriptMutation>,
110    pub queued_prompts: Option<Vec<MaterializedQueuedPrompt>>,
111    pub pending_elicitations: Option<Vec<crate::elicitation::ElicitationRequest>>,
112    /// The nested option distinguishes "unchanged" from "cleared", which is
113    /// how a completed turn removes the running turn.
114    pub active_turn: Option<Option<MaterializedTurn>>,
115    /// Explicit context resets retire the last turn without deleting its history.
116    pub clear_turn_outcome: bool,
117    pub last_turn_outcome: Option<MaterializedTurnOutcome>,
118    pub config_results: Vec<(String, Option<String>)>,
119    pub provider_cost: Option<crate::usage::ProviderCost>,
120    pub api_events: Vec<ApiEventData>,
121}
122
123/// One viewer's stored state for one session.
124#[derive(Debug, Clone, Default, PartialEq, Eq)]
125pub struct ClientSessionState {
126    pub draft: String,
127    pub through_event_ordinal: u64,
128}
129
130/// A terminal's current composer and the shared value it originally inherited.
131/// The inherited value is retired on detach, never replaced by client-local text.
132#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
133pub struct DetachedSessionDraft {
134    pub text: String,
135    pub inherited_input: Option<String>,
136}
137
138/// What one turn produced, without loading the transcript around it.
139#[derive(Debug, Clone, PartialEq, Eq)]
140pub struct TurnSummary {
141    /// One-based count of turn starts up to and including this one, which is
142    /// what a caller means by "turn 3 of this session".
143    pub turn_number: u64,
144    pub turn_started_at_ms: i64,
145    /// Newest change at or after the turn start, so the caller can measure how
146    /// long the turn took.
147    pub last_changed_at_ms: i64,
148    /// The last nonempty agent message the turn produced, flattened to text.
149    pub final_message: Option<String>,
150    /// How many tool calls the turn made.
151    pub tool_calls: u64,
152}
153
154/// One page of a session's transcript, ordered by the sequence a reader pages
155/// by rather than by creation order.
156#[derive(Debug, Clone)]
157pub struct TranscriptPage {
158    pub items: Vec<Arc<TranscriptItem>>,
159    /// The newest sequence in the whole transcript, so a caller can tell
160    /// whether this page reached the end without asking for another one.
161    pub latest_seq: u64,
162    pub next_after_seq: u64,
163    pub execution: MaterializedExecutionState,
164}
165
166/// Exclusive chronological cursor; update sequences are not creation order.
167#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
168pub struct TranscriptCursor {
169    pub position: u64,
170    pub stable_id: String,
171}
172
173impl TranscriptCursor {
174    pub fn of(item: &TranscriptItem) -> Self {
175        Self {
176            position: item.position,
177            stable_id: item.stable_id.clone(),
178        }
179    }
180}
181
182/// A bounded historical read, independent of the live projection window.
183#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
184pub struct TranscriptHistoryPage {
185    pub items: Vec<Arc<TranscriptItem>>,
186    pub before: Option<TranscriptCursor>,
187    pub frontier: u64,
188}
189
190/// What one retention pass reclaimed.
191#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
192pub struct TranscriptRetention {
193    pub items: usize,
194    pub bytes: usize,
195    /// Rows this pass left for the next one, because of
196    /// `RETENTION_BATCH_ITEMS`.
197    pub remaining: bool,
198}
199
200/// A second-opinion review that was still open when the UI last stopped.
201#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
202pub struct StoredReview {
203    pub workflow: crate::second_opinion::ReviewWorkflow,
204    /// Reviewer lifetime this review belongs to. It is bumped when native
205    /// continuity is lost, so a resumed session starts a new conversation
206    /// rather than pretending to reload one that is gone.
207    pub generation: u64,
208    /// The primary's transcript frontier when the context request went out.
209    pub context_baseline: u64,
210    /// Whether the reviewer's native session is known to be gone.
211    pub native_lost: bool,
212    /// What the controller has read of the reviewer's conversation. The
213    /// reviewer's own journal is the source, but it dies with the target, so
214    /// this copy is what keeps a finished review readable afterwards.
215    pub reviewer_transcript: Vec<std::sync::Arc<crate::state::TranscriptItem>>,
216}
217
218/// How far this session has been reviewed.
219///
220/// `baselines` are Git tree ids by repository root: the working tree as of the
221/// last completed review. They advance only when a review resolves, which is
222/// what makes a cancelled review lossless -- the next one covers both turns.
223#[derive(Debug, Clone, Default, PartialEq, serde::Serialize, serde::Deserialize)]
224pub struct TurnReviewState {
225    pub baselines: std::collections::BTreeMap<std::path::PathBuf, String>,
226    pub reviewed_through_ordinal: u64,
227    /// The last forwarded verdict, which turns the next review into a
228    /// verification pass. Cleared once that pass consumes it.
229    pub prior_review: Option<crate::review::lanes::PriorReviewContext>,
230    /// A review that was running when the daemon stopped. On recovery it is
231    /// cleared without advancing the baseline.
232    pub active: Option<String>,
233    /// A corrective prompt that was submitted while its relay acceptance was
234    /// still ambiguous. It survives a daemon restart so the exact command can
235    /// be retried and reconciled without losing the findings.
236    #[serde(default)]
237    pub pending_forward: Option<crate::review::driver::PendingForward>,
238}
239
240/// What a bounded prompt search found, and whether it stopped early.
241///
242/// The flag is not decoration. Without it a caller cannot tell twenty matches
243/// from the first twenty of many, and will present a partial answer as a whole
244/// one.
245#[derive(Debug, Clone, PartialEq, Eq)]
246pub struct BoundedPromptHistory {
247    pub entries: Vec<PromptHistoryEntry>,
248    pub truncated: bool,
249}
250
251#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
252pub struct UsageCoverage {
253    pub recorded_turns: u64,
254    pub full_turn_reports: u64,
255    pub last_request_reports: u64,
256    pub unspecified_reports: u64,
257    pub missing_reports: u64,
258    /// Started turns for which no completion report has been recorded.
259    #[serde(default)]
260    pub unfinished_turns: u64,
261}
262
263#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
264pub struct UsageCounterTotal {
265    pub tokens: u64,
266    /// Number of full-turn reports supplying this particular counter.
267    pub reported_turns: u64,
268}
269
270#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
271pub struct UsagePage {
272    pub session_id: String,
273    pub turns: Vec<MaterializedTurnOutcome>,
274    pub next_after_seq: u64,
275    pub latest_seq: u64,
276    /// Totals include only reports whose scope is known to be a whole turn.
277    pub totals: BTreeMap<String, UsageCounterTotal>,
278    pub coverage: UsageCoverage,
279    #[serde(default, skip_serializing_if = "Option::is_none")]
280    pub provider_session_cost: Option<ProviderCost>,
281    #[serde(default)]
282    pub turn_selections: BTreeMap<String, UsageSelection>,
283    #[serde(default)]
284    pub by_model: Vec<UsageModelTotal>,
285}
286
287/// Configuration observed when a turn started. None means unknown, not zero.
288#[derive(Debug, Clone, Default, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
289pub struct UsageSelection {
290    pub model: Option<String>,
291    pub effort: Option<String>,
292}
293
294#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
295pub struct UsageModelTotal {
296    #[serde(flatten)]
297    pub selection: UsageSelection,
298    pub totals: BTreeMap<String, UsageCounterTotal>,
299}
300
301#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
302pub struct UsageSummary {
303    pub totals: BTreeMap<String, UsageCounterTotal>,
304    pub coverage: UsageCoverage,
305    pub by_model: Vec<UsageModelTotal>,
306    #[serde(default, skip_serializing_if = "Option::is_none")]
307    pub provider_session_cost: Option<ProviderCost>,
308}
309
310#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
311pub struct UsageTreeSession {
312    pub session_id: String,
313    pub parent_session_id: Option<String>,
314    pub task_name: Option<String>,
315    pub operational_session_present: bool,
316    #[serde(flatten)]
317    pub summary: UsageSummary,
318}
319
320#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
321pub struct UsageTree {
322    pub parent_session_id: String,
323    pub sessions: Vec<UsageTreeSession>,
324    #[serde(flatten)]
325    pub summary: UsageSummary,
326}
327
328#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
329pub struct ApiEvent {
330    pub seq: u64,
331    pub session_id: String,
332    pub recorded_at_ms: i64,
333    #[serde(flatten)]
334    pub event: ApiEventData,
335}
336
337#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
338#[serde(tag = "type", content = "data", rename_all = "snake_case")]
339pub enum ApiEventData {
340    RuntimeResolved {
341        receipt: crate::harness_runtime::RuntimeReceipt,
342    },
343    TurnStarted {
344        turn: MaterializedTurn,
345    },
346    TurnEnded {
347        turn: crate::event_outcome::ApiTurnOutcome,
348    },
349    CommandEnded {
350        #[serde(flatten)]
351        result: crate::event_outcome::CommandResult,
352    },
353    SessionFault {
354        reason: crate::event_outcome::OutcomeReason,
355        message: String,
356        command_id: Option<String>,
357    },
358    LegacyNotice {
359        original_type: String,
360        message: String,
361        command_id: Option<String>,
362    },
363    InputRequired {
364        /// Absent when a completed turn asks for unstructured user input.
365        #[serde(default, skip_serializing_if = "Option::is_none")]
366        request: Option<ElicitationRequest>,
367        turn_id: Option<u64>,
368    },
369    InputResolved {
370        elicitation_id: String,
371        turn_id: Option<u64>,
372        action: String,
373    },
374    ActivityChanged {
375        activity: ApiActivityState,
376    },
377}
378
379impl ApiEventData {
380    pub fn kind(&self) -> &'static str {
381        match self {
382            Self::RuntimeResolved { .. } => "runtime_resolved",
383            Self::TurnStarted { .. } => "turn_started",
384            Self::TurnEnded { .. } => "turn_ended",
385            Self::CommandEnded { .. } => "command_ended",
386            Self::SessionFault { .. } => "session_fault",
387            Self::LegacyNotice { .. } => "legacy_notice",
388            Self::InputRequired { .. } => "input_required",
389            Self::InputResolved { .. } => "input_resolved",
390            Self::ActivityChanged { .. } => "activity_changed",
391        }
392    }
393}
394
395/// These are the same structured facts rendered by the web and terminal UIs.
396#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
397pub struct ApiActivityState {
398    pub state: String,
399    pub details: Option<ApiActivityDetails>,
400    pub is_idle: bool,
401    pub waiting_for_input: bool,
402    pub capacity_retry: bool,
403}
404
405#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
406#[serde(deny_unknown_fields)]
407pub struct ApiActivityDetails {
408    pub kind: ApiActivityKind,
409    #[serde(default, skip_serializing_if = "Option::is_none")]
410    pub turn_started_at_ms: Option<i64>,
411    #[serde(default, skip_serializing_if = "Option::is_none")]
412    pub step_started_at_ms: Option<i64>,
413    #[serde(default, skip_serializing_if = "Option::is_none")]
414    pub background_started_at_ms: Option<i64>,
415    #[serde(default, skip_serializing_if = "Option::is_none")]
416    pub idle_since_ms: Option<i64>,
417    /// When anything last arrived from the harness, while a turn or a tool is
418    /// in flight. A surface subtracts it from the current time to show how
419    /// long a running session has been quiet. Mjolnir never ends a turn for
420    /// this on its own; see `mj_core::activity::silence_note`.
421    #[serde(default, skip_serializing_if = "Option::is_none")]
422    pub last_activity_at_ms: Option<i64>,
423    #[serde(default, skip_serializing_if = "Option::is_none")]
424    pub label: Option<String>,
425}
426
427#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
428#[serde(rename_all = "lowercase")]
429pub enum ApiActivityKind {
430    Turn,
431    Step,
432    Background,
433    Idle,
434    Lifecycle,
435}
436
437#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
438pub struct ApiEventFilter {
439    pub session_id: Option<String>,
440    pub workspace_id: Option<String>,
441}
442
443#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
444pub struct ApiEventPage {
445    pub events: Vec<ApiEvent>,
446    pub next_after_seq: u64,
447    pub latest_seq: u64,
448}