Skip to main content

mj_controller/database/
materialized.rs

1use super::*;
2
3/// Load a session's whole projection, transcript and all.
4///
5/// Crate-private on purpose. The cost of this call is everything that has ever
6/// happened in the conversation, and the callers that made that a visible
7/// problem — the runtime poll and the resume reply — were both outside this
8/// crate. What they wanted was [`load_materialized_projection_tail`]; what
9/// they reached for was this, because it was public and its name did not say
10/// otherwise. The remaining caller owns a live projection and genuinely needs
11/// all of it.
12pub fn load_materialized_session(session_id: &str) -> Result<Option<MaterializedSession>> {
13    load_materialized_session_from(&database_path(), session_id)
14}
15
16/// Load only the projection fields needed by dashboard session summaries.
17/// Transcript bodies for tools, plans, thoughts, and old messages stay in
18/// SQLite, which keeps dashboard startup independent of transcript size.
19pub fn load_materialized_session_summary(
20    session_id: &str,
21) -> Result<Option<MaterializedSessionSummary>> {
22    load_materialized_session_summary_from(&database_path(), session_id)
23}
24
25pub(super) fn load_materialized_session_summary_from(
26    path: &Path,
27    session_id: &str,
28) -> Result<Option<MaterializedSessionSummary>> {
29    let mut reader = open_reader(path)?;
30    let connection = reader.transaction()?;
31    let row = connection
32        .query_row(
33            "SELECT applied_event_ordinal, last_activity_at_ms, execution_state,
34                    running_started_at_ms, session_title
35             FROM materialized_sessions WHERE session_id = ?1",
36            [session_id],
37            |row| {
38                Ok((
39                    row.get::<_, u64>(0)?,
40                    row.get::<_, Option<i64>>(1)?,
41                    row.get::<_, String>(2)?,
42                    row.get::<_, Option<i64>>(3)?,
43                    row.get::<_, Option<String>>(4)?,
44                ))
45            },
46        )
47        .optional()?;
48    let Some((
49        applied_event_ordinal,
50        last_activity_at_ms,
51        execution,
52        running_started_at_ms,
53        session_title,
54    )) = row
55    else {
56        return Ok(None);
57    };
58
59    #[cfg(test)]
60    super::tests::after_materialized_frontier_read();
61    let last_user_message = last_materialized_user_message(&connection, session_id)?;
62    let last_agent_message = last_materialized_agent_message(&connection, session_id)?;
63    let last_agent_message_follows_last_user =
64        last_agent_message
65            .as_ref()
66            .is_some_and(|(agent_position, _)| {
67                last_user_message
68                    .as_ref()
69                    .is_none_or(|(user_position, _)| agent_position > user_position)
70            });
71    let mut ordinal_statement = connection.prepare(
72        "SELECT latest_content_event_ordinal
73         FROM materialized_transcript_items
74         WHERE session_id = ?1
75           AND latest_content_event_ordinal IS NOT NULL
76           AND EXISTS (
77               SELECT 1 FROM json_each(
78                   CASE
79                       WHEN latest_content_event_ordinal IS NOT NULL
80                           AND json_valid(body_json)
81                       THEN body_json
82                       ELSE '{}'
83                   END,
84                   '$.chunks'
85               ) AS chunk
86               WHERE json_extract(chunk.value, '$.content.type') IS NOT NULL
87                 AND (
88                     json_extract(chunk.value, '$.content.type') <> 'text'
89                     OR trim(coalesce(json_extract(chunk.value, '$.content.text'), '')) <> ''
90                 )
91           )
92         ORDER BY position, stable_id",
93    )?;
94    let agent_message_latest_content_ordinals = ordinal_statement
95        .query_map([session_id], |row| row.get::<_, u64>(0))?
96        .collect::<rusqlite::Result<Vec<_>>>()?;
97    let restart_pattern = format!("{}*", mj_core::transcript::WORK_INTERRUPTED_ITEM_PREFIX);
98    let mut restart_statement = connection.prepare(
99        "SELECT position
100         FROM materialized_transcript_items
101         WHERE session_id = ?1 AND stable_id GLOB ?2
102         ORDER BY position, stable_id",
103    )?;
104    let mut interruption_event_ordinals = restart_statement
105        .query_map((session_id, restart_pattern), |row| row.get::<_, u64>(0))?
106        .collect::<rusqlite::Result<Vec<_>>>()?;
107    let outcome: Option<String> = connection.query_row(
108        "SELECT last_turn_outcome_json FROM materialized_sessions WHERE session_id = ?1",
109        [session_id],
110        |row| row.get(0),
111    )?;
112    if let Some(ordinal) = outcome
113        .map(|json| serde_json::from_str::<MaterializedTurnOutcome>(&json))
114        .transpose()?
115        .as_ref()
116        .and_then(MaterializedTurnOutcome::interruption_ordinal)
117    {
118        interruption_event_ordinals.push(ordinal);
119    }
120    interruption_event_ordinals.sort_unstable();
121    interruption_event_ordinals.dedup();
122
123    Ok(Some(MaterializedSessionSummary {
124        session_id: session_id.to_owned(),
125        applied_event_ordinal,
126        last_activity_at_ms,
127        execution: parse_materialized_execution(&execution, running_started_at_ms)?,
128        session_title,
129        last_agent_message: last_agent_message.map(|(_, message)| message),
130        last_user_message: last_user_message.map(|(_, message)| message),
131        last_agent_message_follows_last_user,
132        agent_message_latest_content_ordinals,
133        interruption_event_ordinals,
134    }))
135}
136
137/// The oldest visible user message, which is where a session's provisional
138/// title comes from. It sits at the head of the transcript, so a projection
139/// loaded as a tail cannot find it by scanning; this reads it directly.
140pub(super) fn first_materialized_user_message(
141    connection: &Connection,
142    session_id: &str,
143) -> Result<Option<(u64, String)>> {
144    materialized_user_message(connection, session_id, true)
145}
146
147pub(super) fn last_materialized_user_message(
148    connection: &Connection,
149    session_id: &str,
150) -> Result<Option<(u64, String)>> {
151    materialized_user_message(connection, session_id, false)
152}
153
154pub(super) fn materialized_user_message(
155    connection: &Connection,
156    session_id: &str,
157    oldest_first: bool,
158) -> Result<Option<(u64, String)>> {
159    let mut statement = connection.prepare(if oldest_first {
160        "SELECT position, body_json
161         FROM materialized_transcript_items
162         WHERE session_id = ?1
163           AND json_extract(
164               CASE
165                   WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
166                   THEN body_json
167                   ELSE '{}'
168               END,
169               '$.kind'
170           ) = 'user'
171         ORDER BY position, stable_id"
172    } else {
173        "SELECT position, body_json
174         FROM materialized_transcript_items
175         WHERE session_id = ?1
176           AND json_extract(
177               CASE
178                   WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
179                   THEN body_json
180                   ELSE '{}'
181               END,
182               '$.kind'
183           ) = 'user'
184         ORDER BY position DESC, stable_id DESC"
185    })?;
186    let rows = statement.query_map([session_id], |row| {
187        Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?))
188    })?;
189    for row in rows {
190        let (position, body_json) = row?;
191        let body: TranscriptBody = serde_json::from_str(&body_json)
192            .with_context(|| format!("parse materialized user message for session {session_id}"))?;
193        let TranscriptBody::User { content } = body else {
194            continue;
195        };
196        let text = mj_core::transcript::materialized_content_text(&content);
197        if !text.trim().is_empty() {
198            return Ok(Some((position, text)));
199        }
200    }
201    Ok(None)
202}
203
204/// Where the newest turn began: a user message, or the marker for a turn the
205/// harness started on its own. This is the recovery boundary, so it reads a
206/// position only and never has to decode a transcript body.
207pub(super) fn last_materialized_turn_start(
208    connection: &Connection,
209    session_id: &str,
210) -> Result<Option<u64>> {
211    Ok(connection
212        .query_row(
213            "SELECT position
214             FROM materialized_transcript_items
215             WHERE session_id = ?1
216               AND (
217                   stable_id GLOB ?2
218                   OR json_extract(
219                       CASE
220                           WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
221                           THEN body_json
222                           ELSE '{}'
223                       END,
224                       '$.kind'
225                   ) = 'user'
226               )
227             ORDER BY position DESC, stable_id DESC
228             LIMIT 1",
229            params![
230                session_id,
231                format!("{}*", mj_core::transcript::HARNESS_TURN_ITEM_PREFIX)
232            ],
233            |row| row.get::<_, u64>(0),
234        )
235        .optional()?)
236}
237
238pub(super) fn last_materialized_agent_message(
239    connection: &Connection,
240    session_id: &str,
241) -> Result<Option<(u64, String)>> {
242    last_materialized_agent_message_in(connection, session_id, 0, None)
243}
244
245/// The newest nonempty agent message inside one turn, flattened to text.
246///
247/// The turn is a span, not a starting point. A harness can put an agent
248/// message in the transcript after the turn it answered has ended — a resume
249/// notice is one — and reading everything after the turn's start would return
250/// that notice as the turn's answer.
251pub(super) fn last_materialized_agent_message_within(
252    connection: &Connection,
253    session_id: &str,
254    after_position: u64,
255    through_position: u64,
256) -> Result<Option<String>> {
257    Ok(last_materialized_agent_message_in(
258        connection,
259        session_id,
260        after_position,
261        Some(through_position),
262    )?
263    .map(|(_, text)| text))
264}
265
266pub(super) fn last_materialized_agent_message_in(
267    connection: &Connection,
268    session_id: &str,
269    after_position: u64,
270    through_position: Option<u64>,
271) -> Result<Option<(u64, String)>> {
272    let row = connection
273        .query_row(
274            "SELECT position, body_json
275             FROM materialized_transcript_items
276             WHERE session_id = ?1
277               AND position > ?2
278               AND (?3 IS NULL OR position <= ?3)
279               AND latest_content_event_ordinal IS NOT NULL
280               AND EXISTS (
281                   SELECT 1 FROM json_each(
282                       CASE
283                           WHEN latest_content_event_ordinal IS NOT NULL
284                               AND json_valid(body_json)
285                           THEN body_json
286                           ELSE '{}'
287                       END,
288                       '$.chunks'
289                   ) AS chunk
290                   WHERE json_extract(chunk.value, '$.content.type') IS NOT NULL
291                     AND (
292                         json_extract(chunk.value, '$.content.type') <> 'text'
293                         OR trim(coalesce(json_extract(chunk.value, '$.content.text'), '')) <> ''
294                     )
295               )
296             ORDER BY position DESC, stable_id DESC
297             LIMIT 1",
298            params![session_id, after_position, through_position],
299            |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
300        )
301        .optional()?;
302    let Some((position, body_json)) = row else {
303        return Ok(None);
304    };
305    let body: TranscriptBody = serde_json::from_str(&body_json)
306        .with_context(|| format!("parse materialized agent message for session {session_id}"))?;
307    let TranscriptBody::Agent { chunks, .. } = body else {
308        return Ok(None);
309    };
310    let text = mj_core::transcript::materialized_chunks_text(&chunks);
311    Ok((!text.trim().is_empty()).then_some((position, text)))
312}
313
314/// Execution state, the running turn, and the last finished turn's outcome.
315pub type MaterializedTurnState = (
316    MaterializedExecutionState,
317    Option<MaterializedTurn>,
318    Option<MaterializedTurnOutcome>,
319);
320
321/// Where a session stands turn by turn: what is running now, and how the last
322/// finished prompt ended. Returns `None` when the session has no projection
323/// row. The API's wait loop reads this for sessions whose actor is gone.
324pub fn load_materialized_turn_outcome(session_id: &str) -> Result<Option<MaterializedTurnState>> {
325    load_materialized_turn_outcome_from(&database_path(), session_id)
326}
327
328/// Reconcile a replayed input after the worker has collected its command ledger.
329/// Queued, active and historical turns are read from one projection transaction.
330pub fn load_prompt_acceptance(session_id: &str, command_id: &str) -> Result<Option<u64>> {
331    let mut reader = open_reader(&database_path())?;
332    let connection = reader.transaction()?;
333    if let Some(fields) = read_materialized_session_fields(&connection, session_id)? {
334        if let Some(turn) = fields
335            .active_turn
336            .filter(|turn| turn.command_id == command_id)
337        {
338            return turn
339                .accepted_ordinal
340                .map(Some)
341                .context("accepted prompt has no acceptance ordinal");
342        }
343        if let Some(turn) = fields
344            .last_turn_outcome
345            .filter(|turn| turn.command_id == command_id)
346        {
347            return turn
348                .accepted_ordinal
349                .map(Some)
350                .context("completed prompt has no acceptance ordinal");
351        }
352    }
353    let ordinal: Option<u64> = connection.query_row(
354        "SELECT accepted_ordinal FROM materialized_queued_prompts WHERE session_id=?1 AND command_id=?2
355         UNION ALL SELECT json_extract(body, '$.accepted_ordinal') FROM session_turn_usage
356         WHERE session_id=?1 AND command_id=?2 LIMIT 1",
357        params![session_id, command_id], |row| row.get(0),
358    ).optional()?.flatten();
359    if ordinal.is_some() {
360        return Ok(ordinal);
361    }
362    let started: bool = connection.query_row(
363        "SELECT EXISTS(SELECT 1 FROM materialized_transcript_items WHERE session_id=?1 AND stable_id IN (?2, ?3))",
364        params![session_id, format!("user:{command_id}"), format!("{}{command_id}", mj_core::archive::CONTEXT_BOUNDARY_PREFIX)], |row| row.get(0),
365    )?;
366    ensure!(
367        !started,
368        "input was delivered but its acceptance ordinal is unavailable; refusing to send it twice"
369    );
370    Ok(None)
371}
372
373pub(super) fn load_materialized_turn_outcome_from(
374    path: &Path,
375    session_id: &str,
376) -> Result<Option<MaterializedTurnState>> {
377    let connection = open_reader(path)?;
378    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
379        return Ok(None);
380    };
381    Ok(Some((
382        fields.execution,
383        fields.active_turn,
384        fields.last_turn_outcome,
385    )))
386}
387
388/// The answer a session's last finished turn ended with.
389///
390/// This is what a finished child session reports back to the agent that
391/// delegated to it, so it has to be that turn's own last agent message. The
392/// session-wide last agent message is not the same thing: a harness records
393/// messages of its own outside any turn, and a resume notice arriving after
394/// the child finished would then stand in for the child's report.
395///
396/// `None` means the session has no projection row, or no finished turn whose
397/// span is recorded, and the caller decides what to show instead.
398pub fn load_materialized_finished_turn_message(session_id: &str) -> Result<Option<String>> {
399    load_materialized_finished_turn_message_from(&database_path(), session_id)
400}
401
402pub(super) fn load_materialized_finished_turn_message_from(
403    path: &Path,
404    session_id: &str,
405) -> Result<Option<String>> {
406    let mut reader = open_reader(path)?;
407    let connection = reader.transaction()?;
408    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
409        return Ok(None);
410    };
411    let Some(turn) = fields.last_turn_outcome else {
412        return Ok(None);
413    };
414    let Some(start_position) = turn.turn_start_position else {
415        return Ok(None);
416    };
417    last_materialized_agent_message_within(
418        &connection,
419        session_id,
420        start_position,
421        turn.completed_ordinal,
422    )
423}
424
425/// Summarize the turn that ran from `turn_start_position` to
426/// `turn_completed_position`.
427///
428/// Both bounds come from the turn's own record: a transcript item's position
429/// and a completed turn's `completed_ordinal` are the same relay ordinal, so
430/// the completion ordinal is the last position the turn can own. Anything the
431/// session records afterwards belongs to no turn, or to the next one.
432pub fn load_materialized_turn_summary(
433    session_id: &str,
434    turn_start_position: u64,
435    turn_completed_position: u64,
436) -> Result<TurnSummary> {
437    load_materialized_turn_summary_from(
438        &database_path(),
439        session_id,
440        turn_start_position,
441        turn_completed_position,
442    )
443}
444
445pub(super) fn load_materialized_turn_summary_from(
446    path: &Path,
447    session_id: &str,
448    turn_start_position: u64,
449    turn_completed_position: u64,
450) -> Result<TurnSummary> {
451    let mut reader = open_reader(path)?;
452    let connection = reader.transaction()?;
453    let turn_number = connection.query_row(
454        "SELECT COUNT(*)
455         FROM materialized_transcript_items
456         WHERE session_id = ?1
457           AND position <= ?3
458           AND (
459               stable_id GLOB ?2
460               OR json_extract(
461                   CASE
462                       WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
463                       THEN body_json
464                       ELSE '{}'
465                   END,
466                   '$.kind'
467               ) = 'user'
468           )",
469        params![
470            session_id,
471            format!("{}*", mj_core::transcript::HARNESS_TURN_ITEM_PREFIX),
472            turn_start_position
473        ],
474        |row| row.get::<_, u64>(0),
475    )?;
476    let (turn_started_at_ms, last_changed_at_ms) = connection.query_row(
477        "SELECT COALESCE(MIN(created_at_ms), 0), COALESCE(MAX(last_changed_at_ms), 0)
478         FROM materialized_transcript_items
479         WHERE session_id = ?1 AND position >= ?2 AND position <= ?3",
480        params![session_id, turn_start_position, turn_completed_position],
481        |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
482    )?;
483    let final_message = last_materialized_agent_message_within(
484        &connection,
485        session_id,
486        turn_start_position,
487        turn_completed_position,
488    )?;
489    let tool_calls = connection.query_row(
490        "SELECT COUNT(*)
491         FROM materialized_transcript_items
492         WHERE session_id = ?1 AND position >= ?2 AND position <= ?3
493           AND json_extract(body_json, '$.kind') = 'tool'",
494        params![session_id, turn_start_position, turn_completed_position],
495        |row| row.get::<_, u64>(0),
496    )?;
497    Ok(TurnSummary {
498        turn_number,
499        turn_started_at_ms,
500        last_changed_at_ms,
501        final_message,
502        tool_calls,
503    })
504}
505
506pub fn load_materialized_transcript_filtered(
507    session_id: &str,
508    after_seq: u64,
509    limit: usize,
510    role: Option<mj_core::transcript::TranscriptRole>,
511) -> Result<Option<TranscriptPage>> {
512    load_materialized_transcript_filtered_from(&database_path(), session_id, after_seq, limit, role)
513}
514
515pub(super) fn load_materialized_transcript_filtered_from(
516    path: &Path,
517    session_id: &str,
518    after_seq: u64,
519    limit: usize,
520    role: Option<mj_core::transcript::TranscriptRole>,
521) -> Result<Option<TranscriptPage>> {
522    let mut reader = open_reader(path)?;
523    let connection = reader.transaction()?;
524    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
525        return Ok(None);
526    };
527    let role = role.map(|r| r.storage_kind());
528    let mut statement = connection.prepare(
529        "WITH matches AS (
530             SELECT *, COALESCE(latest_content_event_ordinal, position) AS seq
531             FROM materialized_transcript_items WHERE session_id = ?1
532             AND COALESCE(latest_content_event_ordinal, position) > ?2
533             AND (?4 IS NULL OR json_extract(body_json, '$.kind') = ?4)
534         ), boundary AS (SELECT MAX(seq) AS seq FROM (SELECT seq FROM matches ORDER BY seq LIMIT ?3))
535         SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
536                last_changed_at_ms, body_json
537         FROM matches WHERE seq <= (SELECT seq FROM boundary)
538         ORDER BY seq, stable_id",
539    )?;
540    let rows = statement
541        .query_map(
542            params![session_id, after_seq, limit.clamp(1, 1000) as i64, role],
543            |row| {
544                Ok((
545                    row.get::<_, String>(0)?,
546                    row.get::<_, u64>(1)?,
547                    row.get::<_, Option<u64>>(2)?,
548                    row.get::<_, i64>(3)?,
549                    row.get::<_, i64>(4)?,
550                    row.get::<_, String>(5)?,
551                ))
552            },
553        )?
554        .collect::<rusqlite::Result<Vec<_>>>()?;
555    let items = rows
556        .into_iter()
557        .map(
558            |(
559                stable_id,
560                position,
561                latest_content_event_ordinal,
562                created_at_ms,
563                last_changed_at_ms,
564                body_json,
565            )| {
566                Ok(Arc::new(TranscriptItem {
567                    stable_id,
568                    position,
569                    latest_content_event_ordinal,
570                    created_at_ms,
571                    last_changed_at_ms,
572                    body: decode_transcript_body(&body_json, session_id)?,
573                }))
574            },
575        )
576        .collect::<Result<Vec<_>>>()?;
577    let latest_seq = connection.query_row(
578        "SELECT COALESCE(MAX(COALESCE(latest_content_event_ordinal, position)), 0)
579         FROM materialized_transcript_items
580         WHERE session_id = ?1",
581        [session_id],
582        |row| row.get::<_, u64>(0),
583    )?;
584    let last_seq = items.last().map_or(after_seq, |item| item.seq());
585    let more: bool = connection.query_row("SELECT EXISTS(SELECT 1 FROM materialized_transcript_items WHERE session_id = ?1 AND COALESCE(latest_content_event_ordinal, position) > ?2 AND (?3 IS NULL OR json_extract(body_json, '$.kind') = ?3))", params![session_id, last_seq, role], |r| r.get(0))?;
586    Ok(Some(TranscriptPage {
587        next_after_seq: if more {
588            last_seq
589        } else {
590            latest_seq.max(after_seq)
591        },
592        items,
593        latest_seq,
594        execution: fields.execution,
595    }))
596}
597
598/// How many transcript rows one retention pass rewrites.
599///
600/// The daemon is the single database writer, so a pass that rewrote every row
601/// of a long session would stall every other write behind it. A capped pass
602/// leaves the rest for the next checkpoint, which is the next time any of it
603/// becomes redundant anyway.
604pub(super) const RETENTION_BATCH_ITEMS: usize = 4_096;
605
606/// Rows below this are already small enough that rewriting them would cost
607/// more than it reclaims.
608pub(super) const RETENTION_BODY_FLOOR_BYTES: usize = 4 * 1024;
609
610/// Drop tool output that a verified checkpoint already holds.
611///
612/// The projection only ever grew: the only deletes were a per-item remove, a
613/// whole-session wipe, and the `sessions` cascade. One measured session reached
614/// 28,066 items and 635 MiB, of which 561 MB was tool-call content.
615///
616/// A checkpoint archive carries the complete transcript up to its event
617/// frontier, and one checkpoint per session is retained, so every item at or
618/// below `event_frontier` is durably recorded elsewhere. What stays here is
619/// what the transcript still shows: which tool ran, on what, with what result,
620/// and each edit's diffstat. See
621/// [`mj_transcript::transcript::compact_tool_call_for_retention`].
622pub fn compact_materialized_transcript_through(
623    session_id: &str,
624    event_frontier: u64,
625) -> Result<TranscriptRetention> {
626    let session_id = session_id.to_owned();
627    submit_database_write("compact_materialized_transcript", move |_| {
628        compact_materialized_transcript_in(&database_path(), &session_id, event_frontier)
629    })
630}
631
632pub(super) fn compact_materialized_transcript_in(
633    path: &Path,
634    session_id: &str,
635    event_frontier: u64,
636) -> Result<TranscriptRetention> {
637    let mut connection = open(path)?;
638    let candidates = {
639        let mut statement = connection.prepare(
640            "SELECT stable_id, body_json
641             FROM materialized_transcript_items
642             WHERE session_id = ?1
643               AND position <= ?2
644               AND length(body_json) > ?3
645               AND json_extract(
646                   CASE WHEN json_valid(body_json) THEN body_json ELSE '{}' END,
647                   '$.kind'
648               ) = 'tool'
649             ORDER BY position, stable_id
650             LIMIT ?4",
651        )?;
652        statement
653            .query_map(
654                params![
655                    session_id,
656                    event_frontier,
657                    RETENTION_BODY_FLOOR_BYTES as i64,
658                    RETENTION_BATCH_ITEMS as i64 + 1
659                ],
660                |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
661            )?
662            .collect::<rusqlite::Result<Vec<_>>>()?
663    };
664    let remaining = candidates.len() > RETENTION_BATCH_ITEMS;
665    let mut retention = TranscriptRetention {
666        remaining,
667        ..TranscriptRetention::default()
668    };
669    let transaction =
670        connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
671    for (stable_id, body_json) in candidates.into_iter().take(RETENTION_BATCH_ITEMS) {
672        let mut body: TranscriptBody = match serde_json::from_str(&body_json) {
673            Ok(body) => body,
674            // A row this cannot read is a row it must not rewrite.
675            Err(error) => {
676                tracing::warn!(%session_id, %stable_id, %error, "skipping unreadable transcript body");
677                continue;
678            }
679        };
680        if !mj_transcript::transcript::compact_tool_call_for_retention(&mut body) {
681            continue;
682        }
683        let compacted = serde_json::to_string(&body)
684            .with_context(|| format!("serialize compacted transcript body {stable_id}"))?;
685        if compacted.len() >= body_json.len() {
686            continue;
687        }
688        transaction.execute(
689            "UPDATE materialized_transcript_items SET body_json = ?3
690             WHERE session_id = ?1 AND stable_id = ?2",
691            params![session_id, stable_id, compacted],
692        )?;
693        retention.items += 1;
694        retention.bytes += body_json.len() - compacted.len();
695    }
696    transaction.commit()?;
697    Ok(retention)
698}
699
700/// How many transcript items a polled projection carries.
701///
702/// Every viewer of a polled projection is bounded already: the conversation
703/// pane keeps `chat::TAIL_SEED_ITEMS` (256) entries, and the browser
704/// transcript keeps 1,000 rendered lines. This is set above both, since an
705/// entry renders to at least one line, so the window is the whole of what any
706/// of them would show.
707pub const PROJECTION_TAIL_ITEMS: usize = 1_024;
708
709/// Load the same safe window the live projector maintains, without first
710/// deserializing all of the conversation. Read metadata and bodies together.
711pub fn load_materialized_actor_projection(
712    session_id: &str,
713) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
714    load_materialized_actor_projection_from(&database_path(), session_id)
715}
716
717pub(super) fn load_materialized_actor_projection_from(
718    path: &Path,
719    session_id: &str,
720) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
721    let mut reader = open_reader(path)?;
722    let connection = reader.transaction()?;
723    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
724        return Ok(None);
725    };
726    let latest_turn = last_materialized_turn_start(&connection, session_id)?;
727    let desired: Option<u64> = connection
728        .query_row(
729            "SELECT position FROM materialized_transcript_items WHERE session_id=?1
730         ORDER BY position DESC, stable_id DESC LIMIT 1 OFFSET ?2",
731            params![session_id, PROJECTION_TAIL_ITEMS - 1],
732            |row| row.get(0),
733        )
734        .optional()?;
735    // Pending tools and streams can outlive their initiating turn. Keep their
736    // turn too; the same rule is applied by ProjectionWindow::trim.
737    let mutable: Option<u64> = connection.query_row(
738        "SELECT MIN(position) FROM materialized_transcript_items WHERE session_id=?1
739         AND (json_extract(body_json, '$.streaming')=1
740              OR (json_extract(body_json, '$.kind')='tool'
741                  AND json_extract(body_json, '$.call.status') IN ('pending','in_progress')))",
742        [session_id],
743        |row| row.get(0),
744    )?;
745    let boundary = desired
746        .unwrap_or(0)
747        .min(latest_turn.unwrap_or(0))
748        .min(mutable.unwrap_or(u64::MAX));
749    let start: u64 = connection.query_row(
750        "SELECT COALESCE(MAX(position),0) FROM materialized_transcript_items
751         WHERE session_id=?1 AND position<=?2
752         AND (json_extract(body_json, '$.kind')='user' OR stable_id LIKE 'harness-turn:%')",
753        params![session_id, boundary],
754        |row| row.get(0),
755    )?;
756    let transcript = read_transcript_range(&connection, session_id, start, None, None)?;
757    let omitted_items = connection.query_row(
758        "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id=?1 AND position<?2",
759        params![session_id, start],
760        |row| row.get(0),
761    )?;
762    let window = ProjectionWindow {
763        omitted_items,
764        provisional_title: first_materialized_user_message(&connection, session_id)?
765            .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
766        latest_turn_start_position: latest_turn,
767    };
768    let materialized = MaterializedSession {
769        session_id: session_id.to_owned(),
770        applied_event_ordinal: fields.applied_event_ordinal,
771        applied_event_digest: fields.applied_event_digest,
772        last_activity_at_ms: fields.last_activity_at_ms,
773        execution: fields.execution,
774        session_title: fields.session_title,
775        configuration: fields.configuration,
776        transcript,
777        queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
778        pending_elicitations: fields.pending_elicitations,
779        active_turn: fields.active_turn,
780        last_turn_outcome: fields.last_turn_outcome,
781    };
782    materialized.validate()?;
783    Ok(Some((materialized, window)))
784}
785
786pub fn load_transcript_history(
787    session_id: &str,
788    before: Option<&mj_core::storage::TranscriptCursor>,
789    limit: usize,
790) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
791    load_transcript_history_from(&database_path(), session_id, before, limit)
792}
793
794pub(super) fn load_transcript_history_from(
795    path: &Path,
796    session_id: &str,
797    before: Option<&mj_core::storage::TranscriptCursor>,
798    limit: usize,
799) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
800    let mut reader = open_reader(path)?;
801    let connection = reader.transaction()?;
802    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
803        return Ok(None);
804    };
805    let limit = limit.clamp(1, 256);
806    let mut items = read_transcript_range(&connection, session_id, 0, before, Some(limit + 1))?;
807    let has_more = items.len() > limit;
808    if has_more {
809        items.remove(0);
810    }
811    let before = has_more.then(|| mj_core::storage::TranscriptCursor::of(&items[0]));
812    Ok(Some(mj_core::storage::TranscriptHistoryPage {
813        items,
814        before,
815        frontier: fields.applied_event_ordinal,
816    }))
817}
818
819fn read_transcript_range(
820    connection: &Connection,
821    session_id: &str,
822    start: u64,
823    before: Option<&mj_core::storage::TranscriptCursor>,
824    limit: Option<usize>,
825) -> Result<Vec<Arc<TranscriptItem>>> {
826    let mut statement = connection.prepare(
827        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
828                last_changed_at_ms, body_json FROM materialized_transcript_items
829         WHERE session_id=?1 AND position>=?2
830           AND (position,stable_id)<(?3,?4)
831         ORDER BY position DESC, stable_id DESC LIMIT ?5",
832    )?;
833    let rows = statement
834        .query_map(
835            params![
836                session_id,
837                start,
838                before.map_or(i64::MAX as u64, |c| c.position),
839                before.map_or("", |c| c.stable_id.as_str()),
840                limit.map_or(-1, |limit| limit as i64)
841            ],
842            |row| {
843                Ok((
844                    row.get::<_, String>(0)?,
845                    row.get::<_, u64>(1)?,
846                    row.get::<_, Option<u64>>(2)?,
847                    row.get::<_, i64>(3)?,
848                    row.get::<_, i64>(4)?,
849                    row.get::<_, String>(5)?,
850                ))
851            },
852        )?
853        .collect::<rusqlite::Result<Vec<_>>>()?;
854    let mut items = rows
855        .into_iter()
856        .map(
857            |(
858                stable_id,
859                position,
860                latest_content_event_ordinal,
861                created_at_ms,
862                last_changed_at_ms,
863                body,
864            )| {
865                Ok(Arc::new(TranscriptItem {
866                    stable_id,
867                    position,
868                    latest_content_event_ordinal,
869                    created_at_ms,
870                    last_changed_at_ms,
871                    body: decode_transcript_body(&body, session_id)?,
872                }))
873            },
874        )
875        .collect::<Result<Vec<_>>>()?;
876    items.reverse();
877    Ok(items)
878}
879
880/// Rehydrate explicitly referenced historical tools before deriving a delta.
881/// A window is a cache, never evidence that a durable tool does not exist.
882pub fn load_projection_references(
883    projection: &MaterializedSession,
884    events: &[mj_core::relay::RelayEvent],
885) -> Result<Vec<Arc<TranscriptItem>>> {
886    load_projection_references_from(&database_path(), projection, events)
887}
888
889pub(super) fn load_projection_references_from(
890    path: &Path,
891    projection: &MaterializedSession,
892    events: &[mj_core::relay::RelayEvent],
893) -> Result<Vec<Arc<TranscriptItem>>> {
894    let (mut ids, terminals) = mj_transcript::projection::historical_references(events)?;
895    let retained = projection
896        .transcript
897        .iter()
898        .map(|item| item.stable_id.as_str())
899        .collect::<std::collections::HashSet<_>>();
900    ids.retain(|id| !retained.contains(id.as_str()));
901    if ids.is_empty() && terminals.is_empty() {
902        return Ok(Vec::new());
903    }
904    let mut reader = open_reader(path)?;
905    let connection = reader.transaction()?;
906    // Ordinary tool updates use the stable-ID index. Only terminal attachment
907    // events need the JSON reverse-reference scan.
908    let mut parameters = vec![projection.session_id.clone(), serde_json::to_string(&ids)?];
909    let mut statement = connection.prepare(if terminals.is_empty() {
910        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
911                last_changed_at_ms, body_json FROM materialized_transcript_items
912         WHERE session_id=?1 AND stable_id IN (SELECT value FROM json_each(?2))
913         ORDER BY position,stable_id"
914    } else {
915        parameters.push(serde_json::to_string(&retained)?);
916        parameters.push(serde_json::to_string(&terminals)?);
917        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
918                last_changed_at_ms, body_json FROM materialized_transcript_items
919         WHERE session_id=?1 AND stable_id NOT IN (SELECT value FROM json_each(?3))
920           AND (stable_id IN (SELECT value FROM json_each(?2))
921             OR EXISTS(SELECT 1 FROM json_each(body_json, '$.terminal_refs')
922                       WHERE value IN (SELECT value FROM json_each(?4))))
923         ORDER BY position,stable_id"
924    })?;
925    let rows = statement
926        .query_map(rusqlite::params_from_iter(parameters), |row| {
927            Ok((
928                row.get::<_, String>(0)?,
929                row.get::<_, u64>(1)?,
930                row.get::<_, Option<u64>>(2)?,
931                row.get::<_, i64>(3)?,
932                row.get::<_, i64>(4)?,
933                row.get::<_, String>(5)?,
934            ))
935        })?
936        .collect::<rusqlite::Result<Vec<_>>>()?;
937    rows.into_iter()
938        .map(
939            |(
940                stable_id,
941                position,
942                latest_content_event_ordinal,
943                created_at_ms,
944                last_changed_at_ms,
945                body,
946            )| {
947                Ok(Arc::new(TranscriptItem {
948                    stable_id,
949                    position,
950                    latest_content_event_ordinal,
951                    created_at_ms,
952                    last_changed_at_ms,
953                    body: decode_transcript_body(&body, &projection.session_id)?,
954                }))
955            },
956        )
957        .collect()
958}
959
960/// Load a projection carrying only the end of its transcript.
961///
962/// The steady-state poll reloads a session's projection every time anything
963/// about it moves. Loading the whole transcript to do that is work
964/// proportional to everything that has ever happened in the conversation —
965/// 635 MiB and 28,066 items on one measured session — for a view that shows
966/// the last few hundred entries. This reads the window instead, plus the two
967/// facts that live outside it, each with one indexed query. See
968/// [`ProjectionWindow`].
969pub fn load_materialized_projection_tail(
970    session_id: &str,
971    transcript_limit: usize,
972) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
973    load_materialized_projection_tail_from(&database_path(), session_id, transcript_limit)
974}
975
976pub(super) fn load_materialized_projection_tail_from(
977    path: &Path,
978    session_id: &str,
979    transcript_limit: usize,
980) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
981    let mut reader = open_reader(path)?;
982    // Frontier, mutable transcript bodies, and window metadata must describe
983    // one WAL snapshot even if the daemon commits between these queries.
984    let connection = reader.transaction()?;
985    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
986        return Ok(None);
987    };
988    let transcript = read_materialized_transcript(&connection, session_id, Some(transcript_limit))?;
989    let total_items = connection.query_row(
990        "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id = ?1",
991        [session_id],
992        |row| row.get::<_, usize>(0),
993    )?;
994    let window = ProjectionWindow {
995        omitted_items: total_items.saturating_sub(transcript.len()),
996        provisional_title: first_materialized_user_message(&connection, session_id)?
997            .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
998        latest_turn_start_position: last_materialized_turn_start(&connection, session_id)?,
999    };
1000    let materialized = MaterializedSession {
1001        session_id: session_id.to_owned(),
1002        applied_event_ordinal: fields.applied_event_ordinal,
1003        applied_event_digest: fields.applied_event_digest,
1004        last_activity_at_ms: fields.last_activity_at_ms,
1005        execution: fields.execution,
1006        session_title: fields.session_title,
1007        configuration: fields.configuration,
1008        transcript,
1009        queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
1010        pending_elicitations: fields.pending_elicitations,
1011        active_turn: fields.active_turn,
1012        last_turn_outcome: fields.last_turn_outcome,
1013    };
1014    materialized.validate()?;
1015    Ok(Some((materialized, window)))
1016}
1017
1018/// Read only the projection's event frontier. Deciding whether a stored
1019/// projection already matches an archive costs one row this way, instead of
1020/// deserializing every transcript item to compare two integers.
1021pub fn materialized_event_frontier(session_id: &str) -> Result<Option<(u64, String)>> {
1022    materialized_event_frontier_from(&database_path(), session_id)
1023}
1024
1025pub(super) fn materialized_event_frontier_from(
1026    path: &Path,
1027    session_id: &str,
1028) -> Result<Option<(u64, String)>> {
1029    Ok(open_reader(path)?
1030        .query_row(
1031            "SELECT applied_event_ordinal, applied_event_digest
1032             FROM materialized_sessions WHERE session_id = ?1",
1033            [session_id],
1034            |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1035        )
1036        .optional()?)
1037}
1038
1039/// Replace a session's durable prompt queue without touching its transcript or
1040/// event frontier. Resume uses this when it keeps the stored projection but
1041/// still has to drop the queue the archive carried.
1042pub fn replace_materialized_queued_prompts(
1043    session_id: &str,
1044    queued_prompts: &[MaterializedQueuedPrompt],
1045) -> Result<()> {
1046    let session_id = session_id.to_owned();
1047    let queued_prompts = queued_prompts.to_vec();
1048    submit_database_write("replace_materialized_queued_prompts", move |_| {
1049        replace_materialized_queued_prompts_in(&database_path(), &session_id, &queued_prompts)
1050    })
1051}
1052
1053pub(super) fn replace_materialized_queued_prompts_in(
1054    path: &Path,
1055    session_id: &str,
1056    queued_prompts: &[MaterializedQueuedPrompt],
1057) -> Result<()> {
1058    let mut connection = open(path)?;
1059    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1060    if !session_exists(&tx, session_id)? {
1061        bail!("unknown session {session_id}");
1062    }
1063    replace_materialized_queue(&tx, session_id, queued_prompts)?;
1064    tx.commit()?;
1065    Ok(())
1066}
1067
1068/// The activity watermark of every session whose projection holds at least one
1069/// transcript item, by session id.
1070///
1071/// This is a change token, not a projection: session indexing needs to know
1072/// which live conversations have moved since it last looked, and loading each
1073/// one's transcript to find out would cost the whole corpus every sync.
1074pub fn load_transcribed_session_activity() -> Result<BTreeMap<String, Option<i64>>> {
1075    load_transcribed_session_activity_from(&database_path())
1076}
1077
1078fn load_transcribed_session_activity_from(path: &Path) -> Result<BTreeMap<String, Option<i64>>> {
1079    let connection = open_reader(path)?;
1080    let mut statement = connection.prepare(
1081        "SELECT session_id, last_activity_at_ms
1082         FROM materialized_sessions s
1083         WHERE EXISTS (
1084             SELECT 1 FROM materialized_transcript_items i
1085             WHERE i.session_id = s.session_id
1086         )",
1087    )?;
1088    let rows = statement.query_map([], |row| {
1089        Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?))
1090    })?;
1091    let mut activity = BTreeMap::new();
1092    for row in rows {
1093        let (session_id, last_activity_at_ms) = row?;
1094        activity.insert(session_id, last_activity_at_ms);
1095    }
1096    Ok(activity)
1097}
1098
1099/// Load only the durable prompt queues without deserializing transcript rows.
1100/// Dashboard startup uses this path so work is proportional to queued prompts,
1101/// not to the complete retained conversation history.
1102pub fn load_materialized_queued_prompts() -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>>
1103{
1104    load_materialized_queued_prompts_from(&database_path())
1105}
1106
1107pub(super) fn load_materialized_queued_prompts_from(
1108    path: &Path,
1109) -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>> {
1110    let connection = open_reader(path)?;
1111    let mut statement = connection.prepare(
1112        "SELECT session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1113         FROM materialized_queued_prompts
1114         ORDER BY session_id, ordinal",
1115    )?;
1116    let rows = statement.query_map([], |row| {
1117        Ok((
1118            row.get::<_, String>(0)?,
1119            row.get::<_, String>(1)?,
1120            row.get::<_, String>(2)?,
1121            row.get::<_, String>(3)?,
1122            row.get::<_, i64>(4)?,
1123            row.get::<_, Option<u64>>(5)?,
1124        ))
1125    })?;
1126    let mut queues = BTreeMap::<String, Vec<MaterializedQueuedPrompt>>::new();
1127    for row in rows {
1128        let (session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal) =
1129            row?;
1130        let content = serde_json::from_str(&content_json).with_context(|| {
1131            format!("parse materialized queued prompt for session {session_id}")
1132        })?;
1133        let kind = serde_json::from_str(&kind_json).with_context(|| {
1134            format!("parse materialized queue entry kind for session {session_id}")
1135        })?;
1136        queues
1137            .entry(session_id)
1138            .or_default()
1139            .push(MaterializedQueuedPrompt {
1140                command_id,
1141                kind,
1142                content,
1143                queued_at_ms,
1144                accepted_ordinal,
1145            });
1146    }
1147    Ok(queues)
1148}
1149
1150pub(super) fn load_materialized_session_from(
1151    path: &Path,
1152    session_id: &str,
1153) -> Result<Option<MaterializedSession>> {
1154    let mut reader = open_reader(path)?;
1155    let connection = reader.transaction()?;
1156    load_materialized_session_with(&connection, session_id)
1157}
1158
1159pub(super) fn load_materialized_session_with(
1160    connection: &rusqlite::Transaction<'_>,
1161    session_id: &str,
1162) -> Result<Option<MaterializedSession>> {
1163    let Some(fields) = read_materialized_session_fields(connection, session_id)? else {
1164        return Ok(None);
1165    };
1166    let materialized = MaterializedSession {
1167        session_id: session_id.to_owned(),
1168        applied_event_ordinal: fields.applied_event_ordinal,
1169        applied_event_digest: fields.applied_event_digest,
1170        last_activity_at_ms: fields.last_activity_at_ms,
1171        execution: fields.execution,
1172        session_title: fields.session_title,
1173        configuration: fields.configuration,
1174        transcript: read_materialized_transcript(connection, session_id, None)?,
1175        queued_prompts: read_materialized_queued_prompts(connection, session_id)?,
1176        pending_elicitations: fields.pending_elicitations,
1177        active_turn: fields.active_turn,
1178        last_turn_outcome: fields.last_turn_outcome,
1179    };
1180    materialized.validate()?;
1181    Ok(Some(materialized))
1182}
1183
1184/// Everything a projection holds apart from its transcript and its queue.
1185pub(super) struct MaterializedSessionFields {
1186    pub(super) applied_event_ordinal: u64,
1187    pub(super) applied_event_digest: String,
1188    pub(super) last_activity_at_ms: Option<i64>,
1189    pub(super) execution: MaterializedExecutionState,
1190    pub(super) session_title: Option<String>,
1191    pub(super) configuration: mj_core::state::SessionConfiguration,
1192    pub(super) pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
1193    pub(super) active_turn: Option<MaterializedTurn>,
1194    pub(super) last_turn_outcome: Option<MaterializedTurnOutcome>,
1195}
1196
1197pub(super) fn read_materialized_session_fields(
1198    connection: &Connection,
1199    session_id: &str,
1200) -> Result<Option<MaterializedSessionFields>> {
1201    let row = connection
1202        .query_row(
1203            "SELECT applied_event_ordinal, applied_event_digest, last_activity_at_ms,
1204                    execution_state, running_started_at_ms, session_title, configuration_json,
1205                    pending_elicitations_json, active_turn_json, last_turn_outcome_json
1206             FROM materialized_sessions WHERE session_id = ?1",
1207            [session_id],
1208            |row| {
1209                Ok((
1210                    row.get::<_, u64>(0)?,
1211                    row.get::<_, String>(1)?,
1212                    row.get::<_, Option<i64>>(2)?,
1213                    row.get::<_, String>(3)?,
1214                    row.get::<_, Option<i64>>(4)?,
1215                    row.get::<_, Option<String>>(5)?,
1216                    row.get::<_, String>(6)?,
1217                    row.get::<_, String>(7)?,
1218                    row.get::<_, Option<String>>(8)?,
1219                    row.get::<_, Option<String>>(9)?,
1220                ))
1221            },
1222        )
1223        .optional()?;
1224    let Some((
1225        applied_event_ordinal,
1226        applied_event_digest,
1227        last_activity_at_ms,
1228        execution,
1229        running_started_at_ms,
1230        session_title,
1231        configuration_json,
1232        pending_elicitations_json,
1233        active_turn_json,
1234        last_turn_outcome_json,
1235    )) = row
1236    else {
1237        return Ok(None);
1238    };
1239    #[cfg(test)]
1240    super::tests::after_materialized_frontier_read();
1241    Ok(Some(MaterializedSessionFields {
1242        applied_event_ordinal,
1243        applied_event_digest,
1244        last_activity_at_ms,
1245        execution: parse_materialized_execution(&execution, running_started_at_ms)?,
1246        session_title,
1247        configuration: serde_json::from_str(&configuration_json).with_context(|| {
1248            format!("parse materialized configuration for session {session_id}")
1249        })?,
1250        pending_elicitations: serde_json::from_str(&pending_elicitations_json)
1251            .with_context(|| format!("parse pending elicitations for session {session_id}"))?,
1252        active_turn: active_turn_json
1253            .as_deref()
1254            .map(serde_json::from_str)
1255            .transpose()
1256            .with_context(|| format!("parse active turn for session {session_id}"))?,
1257        last_turn_outcome: last_turn_outcome_json
1258            .as_deref()
1259            .map(serde_json::from_str)
1260            .transpose()
1261            .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
1262    }))
1263}
1264
1265/// Read a session's transcript, oldest first. `limit` reads only that many
1266/// items from the end, walking the `materialized_transcript_position` index
1267/// backwards so the read costs the rows it returns.
1268pub(super) fn read_materialized_transcript(
1269    connection: &Connection,
1270    session_id: &str,
1271    limit: Option<usize>,
1272) -> Result<Vec<Arc<TranscriptItem>>> {
1273    let mut statement = connection.prepare(match limit {
1274        Some(_) => {
1275            "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1276                    last_changed_at_ms, body_json
1277             FROM materialized_transcript_items
1278             WHERE session_id = ?1
1279             ORDER BY position DESC, stable_id DESC
1280             LIMIT ?2"
1281        }
1282        None => {
1283            "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1284                    last_changed_at_ms, body_json
1285             FROM materialized_transcript_items
1286             WHERE session_id = ?1
1287             ORDER BY position, stable_id"
1288        }
1289    })?;
1290    let read = |row: &rusqlite::Row<'_>| {
1291        Ok((
1292            row.get::<_, String>(0)?,
1293            row.get::<_, u64>(1)?,
1294            row.get::<_, Option<u64>>(2)?,
1295            row.get::<_, i64>(3)?,
1296            row.get::<_, i64>(4)?,
1297            row.get::<_, String>(5)?,
1298        ))
1299    };
1300    let rows = match limit {
1301        Some(limit) => statement
1302            .query_map(params![session_id, limit as i64], read)?
1303            .collect::<rusqlite::Result<Vec<_>>>()?,
1304        None => statement
1305            .query_map([session_id], read)?
1306            .collect::<rusqlite::Result<Vec<_>>>()?,
1307    };
1308    let mut transcript = rows
1309        .into_iter()
1310        .map(
1311            |(
1312                stable_id,
1313                position,
1314                latest_content_event_ordinal,
1315                created_at_ms,
1316                last_changed_at_ms,
1317                body_json,
1318            )| {
1319                Ok(Arc::new(TranscriptItem {
1320                    stable_id,
1321                    position,
1322                    latest_content_event_ordinal,
1323                    created_at_ms,
1324                    last_changed_at_ms,
1325                    body: decode_transcript_body(&body_json, session_id)?,
1326                }))
1327            },
1328        )
1329        .collect::<Result<Vec<_>>>()?;
1330    if limit.is_some() {
1331        // The bounded query walks the index backwards to bound what it reads;
1332        // every caller wants the transcript in the order it was written.
1333        transcript.reverse();
1334    }
1335    Ok(transcript)
1336}
1337
1338pub(super) fn read_materialized_queued_prompts(
1339    connection: &Connection,
1340    session_id: &str,
1341) -> Result<Vec<MaterializedQueuedPrompt>> {
1342    let mut statement = connection.prepare(
1343        "SELECT command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1344         FROM materialized_queued_prompts
1345         WHERE session_id = ?1
1346         ORDER BY ordinal",
1347    )?;
1348    let rows = statement
1349        .query_map([session_id], |row| {
1350            Ok((
1351                row.get::<_, String>(0)?,
1352                row.get::<_, String>(1)?,
1353                row.get::<_, String>(2)?,
1354                row.get::<_, i64>(3)?,
1355                row.get::<_, Option<u64>>(4)?,
1356            ))
1357        })?
1358        .collect::<rusqlite::Result<Vec<_>>>()?;
1359    rows.into_iter()
1360        .map(
1361            |(command_id, kind_json, content_json, queued_at_ms, accepted_ordinal)| {
1362                Ok(MaterializedQueuedPrompt {
1363                    command_id,
1364                    kind: serde_json::from_str(&kind_json).with_context(|| {
1365                        format!("parse materialized queue entry kind for session {session_id}")
1366                    })?,
1367                    content: serde_json::from_str(&content_json).with_context(|| {
1368                        format!("parse materialized queued prompt for session {session_id}")
1369                    })?,
1370                    queued_at_ms,
1371                    accepted_ordinal,
1372                })
1373            },
1374        )
1375        .collect()
1376}
1377
1378/// Replace a complete projection, primarily when seeding a restored
1379/// checkpoint. Operational `SessionRecord` metadata and read receipts are not
1380/// modified.
1381pub fn save_materialized_session(materialized: &MaterializedSession) -> Result<()> {
1382    let materialized = materialized.clone();
1383    submit_database_write("save_materialized_session", move |_| {
1384        save_materialized_session_to(&database_path(), &materialized)
1385    })
1386}
1387
1388pub(super) fn save_materialized_session_to(
1389    path: &Path,
1390    materialized: &MaterializedSession,
1391) -> Result<()> {
1392    materialized.validate()?;
1393    let mut connection = open(path)?;
1394    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1395    if !session_exists(&tx, &materialized.session_id)? {
1396        bail!("unknown session {}", materialized.session_id);
1397    }
1398    write_materialized_session(&tx, materialized)?;
1399    tx.commit()?;
1400    Ok(())
1401}
1402
1403/// One relay page being applied inside a single write transaction. The relay
1404/// retains everything past the last acknowledgement, so a page that fails
1405/// part-way rolls back to the previous durable frontier and is simply
1406/// redelivered. Only a committed page may be acknowledged.
1407pub struct ProjectionPage<'a> {
1408    pub(super) session_id: &'a str,
1409    pub(super) transaction: Transaction<'a>,
1410    pub(super) applied_ordinal: u64,
1411    pub(super) applied_digest: String,
1412    pub(super) dirty: bool,
1413    pub(super) pending: MaterializedSessionMutation,
1414    pub(super) pending_transcript: BTreeMap<String, PendingTranscriptMutation>,
1415    pub(super) pending_turns: Vec<MaterializedTurnOutcome>,
1416    /// Configuration changes are coalesced for the session, but selection is
1417    /// captured at each turn start before later events can overwrite it.
1418    usage_configuration: mj_core::state::SessionConfiguration,
1419    pub(super) pending_events: Vec<(i64, ApiEventData)>,
1420}
1421
1422pub(super) struct PendingTranscriptMutation {
1423    pub(super) final_mutation: TranscriptMutation,
1424    pub(super) remove_before_upsert: bool,
1425}
1426
1427impl ProjectionPage<'_> {
1428    /// Apply the projection effects of the next relay event to the open page.
1429    /// The event must continue the chain the page has reached so far, which is
1430    /// the persisted frontier plus every event already applied to this page.
1431    pub fn apply(
1432        &mut self,
1433        event_ordinal: u64,
1434        previous_event_digest: &str,
1435        event_digest: &str,
1436        mutation: &MaterializedSessionMutation,
1437    ) -> Result<ProjectionApplyOutcome> {
1438        if event_ordinal == 0 {
1439            bail!("relay event ordinal must be positive");
1440        }
1441        // A v2 event carries no chain link (empty previous digest). Its
1442        // continuity to the projection frontier is proven by ordinal
1443        // contiguity plus the attach cursor the controller validated against
1444        // the worker, not by an in-record back-reference; divergence is caught
1445        // there, before any event is applied.
1446        let chained = !previous_event_digest.is_empty();
1447        if chained {
1448            validate_relay_event_digest(previous_event_digest, "previous relay event digest")?;
1449        }
1450        validate_relay_event_frontier(event_ordinal, event_digest, "relay event frontier")?;
1451        let session_id = self.session_id;
1452        let applied = self.applied_ordinal;
1453        if event_ordinal < applied {
1454            return Ok(ProjectionApplyOutcome::AlreadyApplied);
1455        }
1456        if event_ordinal == applied {
1457            if event_digest != self.applied_digest {
1458                bail!(
1459                    "relay event digest mismatch for session {session_id} at ordinal {event_ordinal}: projection has {}, received {event_digest}",
1460                    self.applied_digest
1461                );
1462            }
1463            return Ok(ProjectionApplyOutcome::AlreadyApplied);
1464        }
1465        let expected = applied
1466            .checked_add(1)
1467            .context("materialized event ordinal overflow")?;
1468        if event_ordinal != expected {
1469            bail!(
1470                "relay event gap for session {session_id}: expected ordinal {expected}, received {event_ordinal}"
1471            );
1472        }
1473        if chained && previous_event_digest != self.applied_digest {
1474            bail!(
1475                "relay event chain diverged for session {session_id} before ordinal {event_ordinal}: projection has {}, event follows {previous_event_digest}",
1476                self.applied_digest
1477            );
1478        }
1479
1480        if let Some(event) = &mutation.native_agent {
1481            native_agents::apply_native_agent_event(&self.transaction, session_id, event)?;
1482        }
1483        if let Some(activity_at_ms) = mutation.last_activity_at_ms {
1484            self.pending.last_activity_at_ms = Some(
1485                self.pending
1486                    .last_activity_at_ms
1487                    .map_or(activity_at_ms, |existing| existing.max(activity_at_ms)),
1488            );
1489        }
1490        if let Some(execution) = mutation.execution {
1491            self.pending.execution = Some(execution);
1492        }
1493        if let Some(title) = &mutation.session_title {
1494            if title.as_ref().is_some_and(|title| title.trim().is_empty()) {
1495                bail!("materialized session title cannot be empty");
1496            }
1497            self.pending.session_title = Some(title.clone());
1498        }
1499        if let Some(configuration) = &mutation.configuration {
1500            self.usage_configuration = configuration.clone();
1501            self.pending.configuration = Some(configuration.clone());
1502        }
1503        for item_mutation in &mutation.transcript {
1504            match item_mutation {
1505                TranscriptMutation::Upsert(item) => {
1506                    item.validate(event_ordinal)?;
1507                    let stable_id = item.stable_id.clone();
1508                    let entry = self.pending_transcript.entry(stable_id).or_insert_with(|| {
1509                        PendingTranscriptMutation {
1510                            final_mutation: TranscriptMutation::Upsert(item.clone()),
1511                            remove_before_upsert: false,
1512                        }
1513                    });
1514                    entry.remove_before_upsert |=
1515                        matches!(&entry.final_mutation, TranscriptMutation::Remove { .. });
1516                    entry.final_mutation = TranscriptMutation::Upsert(item.clone());
1517                }
1518                TranscriptMutation::Remove { stable_id } => {
1519                    if stable_id.trim().is_empty() {
1520                        bail!("cannot remove a transcript item with an empty stable id");
1521                    }
1522                    let removed = TranscriptMutation::Remove {
1523                        stable_id: stable_id.clone(),
1524                    };
1525                    self.pending_transcript
1526                        .entry(stable_id.clone())
1527                        .and_modify(|entry| entry.final_mutation = removed.clone())
1528                        .or_insert(PendingTranscriptMutation {
1529                            final_mutation: removed,
1530                            remove_before_upsert: false,
1531                        });
1532                }
1533            }
1534        }
1535        if let Some(queued_prompts) = &mutation.queued_prompts {
1536            self.pending.queued_prompts = Some(queued_prompts.clone());
1537        }
1538        if let Some(pending_elicitations) = &mutation.pending_elicitations {
1539            self.pending.pending_elicitations = Some(pending_elicitations.clone());
1540        }
1541        self.pending
1542            .config_results
1543            .extend(mutation.config_results.clone());
1544        if let Some(active_turn) = &mutation.active_turn {
1545            if let Some(turn) = active_turn {
1546                super::usage::record_turn_selection(
1547                    &self.transaction,
1548                    session_id,
1549                    &turn.command_id,
1550                    &self.usage_configuration,
1551                )?;
1552            }
1553            self.pending.active_turn = Some(active_turn.clone());
1554        }
1555        if mutation.clear_turn_outcome {
1556            self.pending.clear_turn_outcome = true;
1557            self.pending.last_turn_outcome = None;
1558        }
1559        if let Some(last_turn_outcome) = &mutation.last_turn_outcome {
1560            self.pending_turns.push(last_turn_outcome.clone());
1561            self.pending.last_turn_outcome = Some(last_turn_outcome.clone());
1562        }
1563        if let Some(cost) = &mutation.provider_cost {
1564            self.pending.provider_cost = Some(cost.clone());
1565        }
1566        self.pending_events.extend(
1567            mutation
1568                .api_events
1569                .iter()
1570                .cloned()
1571                .map(|event| (mutation.last_activity_at_ms.unwrap_or(0), event)),
1572        );
1573        self.applied_ordinal = event_ordinal;
1574        event_digest.clone_into(&mut self.applied_digest);
1575        self.dirty = true;
1576        Ok(ProjectionApplyOutcome::Applied)
1577    }
1578
1579    /// Persist the coalesced final state of this page. Intermediate event
1580    /// frontiers are useful only for chain validation: a page commits or rolls
1581    /// back as a unit, so writing them individually adds no recovery value.
1582    pub(super) fn flush(&mut self) -> Result<()> {
1583        if !self.dirty {
1584            return Ok(());
1585        }
1586        let tx = &self.transaction;
1587        let session_id = self.session_id;
1588        if let Some(execution) = self.pending.execution {
1589            let (state, started_at_ms) = materialized_execution_columns(execution);
1590            tx.execute(
1591                "UPDATE materialized_sessions
1592                 SET execution_state = ?2, running_started_at_ms = ?3
1593                 WHERE session_id = ?1",
1594                params![session_id, state, started_at_ms],
1595            )?;
1596        }
1597        if let Some(title) = &self.pending.session_title {
1598            tx.execute(
1599                "UPDATE materialized_sessions SET session_title = ?2 WHERE session_id = ?1",
1600                params![session_id, title],
1601            )?;
1602        }
1603        if let Some(configuration) = &self.pending.configuration {
1604            tx.execute(
1605                "UPDATE materialized_sessions SET configuration_json = ?2 WHERE session_id = ?1",
1606                params![session_id, serde_json::to_string(configuration)?],
1607            )?;
1608        }
1609        for pending in self.pending_transcript.values() {
1610            match &pending.final_mutation {
1611                TranscriptMutation::Upsert(item) => {
1612                    // A remove followed by an upsert deliberately starts a new
1613                    // item identity. Preserve that boundary even though other
1614                    // repeated updates are coalesced to one write.
1615                    if pending.remove_before_upsert {
1616                        tx.execute(
1617                            "DELETE FROM materialized_transcript_items
1618                             WHERE session_id = ?1 AND stable_id = ?2",
1619                            params![session_id, item.stable_id],
1620                        )?;
1621                    }
1622                    upsert_transcript_item(tx, session_id, item)?;
1623                }
1624                TranscriptMutation::Remove { stable_id } => {
1625                    tx.execute(
1626                        "DELETE FROM materialized_transcript_items
1627                         WHERE session_id = ?1 AND stable_id = ?2",
1628                        params![session_id, stable_id],
1629                    )?;
1630                }
1631            }
1632        }
1633        if let Some(queued_prompts) = &self.pending.queued_prompts {
1634            replace_materialized_queue(tx, session_id, queued_prompts)?;
1635        }
1636        if let Some(pending_elicitations) = &self.pending.pending_elicitations {
1637            tx.execute(
1638                "UPDATE materialized_sessions
1639                 SET pending_elicitations_json = ?2 WHERE session_id = ?1",
1640                params![session_id, serde_json::to_string(pending_elicitations)?],
1641            )?;
1642        }
1643        for (recorded_at_ms, event) in &self.pending_events {
1644            events::insert_api_event(tx, session_id, *recorded_at_ms, event)?;
1645        }
1646        for turn in &self.pending_turns {
1647            tx.execute("INSERT OR REPLACE INTO session_turn_usage(session_id, command_id, completed_ordinal, turn_start_position, body) VALUES (?1, ?2, ?3, ?4, ?5)", params![session_id, turn.command_id, turn.completed_ordinal, turn.turn_start_position, serde_json::to_string(turn)?])?;
1648        }
1649        if let Some(cost) = &self.pending.provider_cost {
1650            tx.execute(
1651                "INSERT OR REPLACE INTO session_provider_cost(session_id, body) VALUES (?1, ?2)",
1652                params![session_id, serde_json::to_string(cost)?],
1653            )?;
1654        }
1655        for (command_id, error) in &self.pending.config_results {
1656            tx.execute("INSERT OR REPLACE INTO api_config_results(session_id, command_id, error) VALUES (?1, ?2, ?3)", params![session_id, command_id, error])?;
1657        }
1658        if let Some(active_turn) = &self.pending.active_turn {
1659            tx.execute(
1660                "UPDATE materialized_sessions SET active_turn_json = ?2 WHERE session_id = ?1",
1661                params![
1662                    session_id,
1663                    active_turn
1664                        .as_ref()
1665                        .map(serde_json::to_string)
1666                        .transpose()?
1667                ],
1668            )?;
1669        }
1670        if self.pending.clear_turn_outcome {
1671            self.transaction.execute("UPDATE materialized_sessions SET last_turn_outcome_json = NULL WHERE session_id = ?1", [session_id])?;
1672        }
1673        if let Some(last_turn_outcome) = &self.pending.last_turn_outcome {
1674            tx.execute(
1675                "UPDATE materialized_sessions
1676                 SET last_turn_outcome_json = ?2 WHERE session_id = ?1",
1677                params![session_id, serde_json::to_string(last_turn_outcome)?],
1678            )?;
1679        }
1680        tx.execute(
1681            "UPDATE materialized_sessions
1682             SET last_activity_at_ms = CASE
1683                     WHEN ?2 IS NULL THEN last_activity_at_ms
1684                     WHEN last_activity_at_ms IS NULL OR last_activity_at_ms < ?2 THEN ?2
1685                     ELSE last_activity_at_ms
1686                 END,
1687                 applied_event_ordinal = ?3,
1688                 applied_event_digest = ?4
1689             WHERE session_id = ?1",
1690            params![
1691                session_id,
1692                self.pending.last_activity_at_ms,
1693                self.applied_ordinal,
1694                self.applied_digest,
1695            ],
1696        )?;
1697        Ok(())
1698    }
1699}
1700
1701/// Apply one relay page in a single transaction. `fill` feeds the page's
1702/// events through [`ProjectionPage::apply`]; the projection changes and the
1703/// event frontier commit together only when `fill` succeeds, so callers may
1704/// acknowledge the page's last ordinal to the relay after this returns.
1705pub fn apply_projection_page<T>(
1706    session_id: &str,
1707    fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T> + Send + 'static,
1708) -> Result<T>
1709where
1710    T: Send + 'static,
1711{
1712    let session_id = session_id.to_owned();
1713    submit_database_write("apply_projection_page", move |connection| {
1714        apply_projection_page_with(connection, &session_id, fill)
1715    })
1716}
1717
1718#[cfg(test)]
1719pub(super) fn apply_projection_page_to<T>(
1720    path: &Path,
1721    session_id: &str,
1722    fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1723) -> Result<T> {
1724    let mut connection = open(path)?;
1725    apply_projection_page_with(&mut connection, session_id, fill)
1726}
1727
1728pub(super) fn apply_projection_page_with<T>(
1729    connection: &mut Connection,
1730    session_id: &str,
1731    fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1732) -> Result<T> {
1733    let transaction =
1734        connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1735    let (applied_ordinal, applied_digest) = transaction
1736        .query_row(
1737            "SELECT applied_event_ordinal, applied_event_digest
1738             FROM materialized_sessions WHERE session_id = ?1",
1739            [session_id],
1740            |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1741        )
1742        .optional()?
1743        .with_context(|| format!("unknown session {session_id}"))?;
1744    validate_relay_event_frontier(
1745        applied_ordinal,
1746        &applied_digest,
1747        "persisted relay event frontier",
1748    )?;
1749    let usage_configuration = read_materialized_session_fields(&transaction, session_id)?
1750        .context("projection session disappeared")?
1751        .configuration;
1752    let mut page = ProjectionPage {
1753        session_id,
1754        transaction,
1755        applied_ordinal,
1756        applied_digest,
1757        dirty: false,
1758        pending: MaterializedSessionMutation::default(),
1759        pending_transcript: BTreeMap::new(),
1760        pending_turns: Vec::new(),
1761        usage_configuration,
1762        pending_events: Vec::new(),
1763    };
1764    // Dropping the page on failure rolls the whole transaction back, leaving
1765    // the projection at the frontier the relay last saw acknowledged.
1766    let filled = fill(&mut page)?;
1767    page.flush()?;
1768    page.transaction.commit()?;
1769    Ok(filled)
1770}
1771
1772/// Apply exactly one relay event, as a page of one.
1773pub fn apply_projection_event(
1774    session_id: &str,
1775    event_ordinal: u64,
1776    previous_event_digest: &str,
1777    event_digest: &str,
1778    mutation: &MaterializedSessionMutation,
1779) -> Result<ProjectionApplyOutcome> {
1780    let session_id = session_id.to_owned();
1781    let previous_event_digest = previous_event_digest.to_owned();
1782    let event_digest = event_digest.to_owned();
1783    let mutation = mutation.clone();
1784    submit_database_write("apply_projection_event", move |connection| {
1785        apply_projection_page_with(connection, &session_id, |page| {
1786            page.apply(
1787                event_ordinal,
1788                &previous_event_digest,
1789                &event_digest,
1790                &mutation,
1791            )
1792        })
1793    })
1794}
1795
1796#[cfg(test)]
1797pub(super) fn apply_projection_event_to(
1798    path: &Path,
1799    session_id: &str,
1800    event_ordinal: u64,
1801    previous_event_digest: &str,
1802    event_digest: &str,
1803    mutation: &MaterializedSessionMutation,
1804) -> Result<ProjectionApplyOutcome> {
1805    apply_projection_page_to(path, session_id, |page| {
1806        page.apply(event_ordinal, previous_event_digest, event_digest, mutation)
1807    })
1808}
1809
1810/// Advance the persisted detach/read receipt monotonically. A receipt cannot
1811/// acknowledge an event the controller projection has not durably applied.
1812pub fn advance_viewed_through_event_ordinal(session_id: &str, through: u64) -> Result<u64> {
1813    let session_id = session_id.to_owned();
1814    submit_database_write("advance_viewed_through_event_ordinal", move |_| {
1815        advance_viewed_through_event_ordinal_to(&database_path(), &session_id, through)
1816    })
1817}
1818
1819pub(super) fn advance_viewed_through_event_ordinal_to(
1820    path: &Path,
1821    session_id: &str,
1822    through: u64,
1823) -> Result<u64> {
1824    let mut connection = open(path)?;
1825    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1826    let applied = tx
1827        .query_row(
1828            "SELECT applied_event_ordinal FROM materialized_sessions WHERE session_id = ?1",
1829            [session_id],
1830            |row| row.get::<_, u64>(0),
1831        )
1832        .optional()?
1833        .with_context(|| format!("unknown session {session_id}"))?;
1834    if through > applied {
1835        bail!(
1836            "cannot acknowledge event ordinal {through} for session {session_id}; projection is at {applied}"
1837        );
1838    }
1839    tx.execute(
1840        "UPDATE sessions
1841         SET viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?2)
1842         WHERE session_id = ?1",
1843        params![session_id, through],
1844    )?;
1845    let receipt = tx.query_row(
1846        "SELECT viewed_through_event_ordinal FROM sessions WHERE session_id = ?1",
1847        [session_id],
1848        |row| row.get::<_, u64>(0),
1849    )?;
1850    tx.commit()?;
1851    Ok(receipt)
1852}
1853
1854/// Decode a stored transcript body, merging runs of streamed text chunks.
1855///
1856/// Rows written before streamed text was merged on the way in hold one chunk
1857/// per token, which costs far more memory decoded than the text it carries.
1858/// Collapsing them here shrinks long sessions without rewriting stored JSON.
1859fn decode_transcript_body(body_json: &str, session_id: &str) -> Result<TranscriptBody> {
1860    let mut body: TranscriptBody = serde_json::from_str(body_json)
1861        .with_context(|| format!("parse materialized transcript body for session {session_id}"))?;
1862    match &mut body {
1863        TranscriptBody::Agent { chunks, .. } | TranscriptBody::Thought { chunks, .. } => {
1864            mj_core::transcript::coalesce_content_chunks(chunks);
1865        }
1866        _ => {}
1867    }
1868    Ok(body)
1869}
1870
1871/// Read complete authorization and recent replies in one consistent snapshot.
1872/// Runs on the continuation service's blocking pool, never a render/event loop.
1873pub(crate) fn load_continuation_evidence(
1874    session_id: &str,
1875    ordinal: u64,
1876    digest: &str,
1877) -> Result<mj_core::continuation::ContinuationEvidence> {
1878    load_continuation_evidence_from(&database_path(), session_id, ordinal, digest)
1879}
1880
1881pub(super) fn load_continuation_evidence_from(
1882    path: &Path,
1883    session_id: &str,
1884    ordinal: u64,
1885    digest: &str,
1886) -> Result<mj_core::continuation::ContinuationEvidence> {
1887    let mut reader = open_reader(path)?;
1888    let connection = reader.transaction()?;
1889    let fields = read_materialized_session_fields(&connection, session_id)?
1890        .context("continuation projection is missing")?;
1891    anyhow::ensure!(
1892        fields.applied_event_ordinal == ordinal && fields.applied_event_digest == digest,
1893        "continuation projection changed"
1894    );
1895    let mut statement = connection.prepare(
1896        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1897                last_changed_at_ms, body_json
1898         FROM materialized_transcript_items WHERE session_id = ?1
1899         AND (json_extract(body_json, '$.kind') IN ('user', 'agent')
1900              OR stable_id LIKE 'context-cleared:%')
1901         ORDER BY position DESC, stable_id DESC",
1902    )?;
1903    let rows = statement.query_map([session_id], |row| {
1904        Ok((
1905            row.get::<_, String>(0)?,
1906            row.get::<_, u64>(1)?,
1907            row.get::<_, Option<u64>>(2)?,
1908            row.get::<_, i64>(3)?,
1909            row.get::<_, i64>(4)?,
1910            row.get::<_, String>(5)?,
1911        ))
1912    })?;
1913    crate::continuation::evidence_from_items(rows.map(|row| {
1914        let (
1915            stable_id,
1916            position,
1917            latest_content_event_ordinal,
1918            created_at_ms,
1919            last_changed_at_ms,
1920            body_json,
1921        ) = row?;
1922        // Refuse exceptionally large source messages rather than clip consent.
1923        anyhow::ensure!(
1924            body_json.len() <= 1024 * 1024,
1925            "continuation source message exceeds budget"
1926        );
1927        Ok(Arc::new(TranscriptItem {
1928            stable_id,
1929            position,
1930            latest_content_event_ordinal,
1931            created_at_ms,
1932            last_changed_at_ms,
1933            body: decode_transcript_body(&body_json, session_id)?,
1934        }))
1935    }))
1936}