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    read_materialized_turn_state(&connection, session_id)
379}
380
381/// Turn observation does not need configuration, forms, or transcript bodies.
382pub(super) fn read_materialized_turn_state(
383    connection: &Connection,
384    session_id: &str,
385) -> Result<Option<MaterializedTurnState>> {
386    let row = connection.query_row(
387        "SELECT execution_state, running_started_at_ms, active_turn_json, last_turn_outcome_json
388         FROM materialized_sessions WHERE session_id=?1",
389        [session_id],
390        |row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?,
391                  row.get::<_, Option<String>>(2)?, row.get::<_, Option<String>>(3)?)),
392    ).optional()?;
393    let Some((execution, started, active, last)) = row else {
394        return Ok(None);
395    };
396    Ok(Some((
397        parse_materialized_execution(&execution, started)?,
398        active
399            .as_deref()
400            .map(serde_json::from_str)
401            .transpose()
402            .with_context(|| format!("parse active turn for session {session_id}"))?,
403        last.as_deref()
404            .map(serde_json::from_str)
405            .transpose()
406            .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
407    )))
408}
409
410/// The answer a session's last finished turn ended with.
411///
412/// This is what a finished child session reports back to the agent that
413/// delegated to it, so it has to be that turn's own last agent message. The
414/// session-wide last agent message is not the same thing: a harness records
415/// messages of its own outside any turn, and a resume notice arriving after
416/// the child finished would then stand in for the child's report.
417///
418/// `None` means the session has no projection row, or no finished turn whose
419/// span is recorded, and the caller decides what to show instead.
420pub fn load_materialized_finished_turn_message(session_id: &str) -> Result<Option<String>> {
421    load_materialized_finished_turn_message_from(&database_path(), session_id)
422}
423
424/// Read answer text only when returning a child wait. Turn bounds are from
425/// the observation that decided completion, even if another turn starts now.
426pub(crate) struct ChildAnswerMessages {
427    pub latest: Option<String>,
428    pub finished: Option<String>,
429}
430
431pub(crate) fn load_child_answer_messages(
432    children: &[(String, Option<(u64, u64)>)],
433) -> Result<BTreeMap<String, ChildAnswerMessages>> {
434    let mut reader = open_reader(&database_path())?;
435    let connection = reader.transaction()?;
436    children
437        .iter()
438        .map(|(id, span)| {
439            let latest = last_materialized_agent_message(&connection, id)?.map(|(_, text)| text);
440            let finished = span
441                .map(|(start, end)| {
442                    last_materialized_agent_message_within(&connection, id, start, end)
443                })
444                .transpose()?
445                .flatten();
446            Ok((id.clone(), ChildAnswerMessages { latest, finished }))
447        })
448        .collect()
449}
450
451pub(super) fn load_materialized_finished_turn_message_from(
452    path: &Path,
453    session_id: &str,
454) -> Result<Option<String>> {
455    let mut reader = open_reader(path)?;
456    let connection = reader.transaction()?;
457    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
458        return Ok(None);
459    };
460    let Some(turn) = fields.last_turn_outcome else {
461        return Ok(None);
462    };
463    let Some(start_position) = turn.turn_start_position else {
464        return Ok(None);
465    };
466    last_materialized_agent_message_within(
467        &connection,
468        session_id,
469        start_position,
470        turn.completed_ordinal,
471    )
472}
473
474/// Summarize the turn that ran from `turn_start_position` to
475/// `turn_completed_position`.
476///
477/// Both bounds come from the turn's own record: a transcript item's position
478/// and a completed turn's `completed_ordinal` are the same relay ordinal, so
479/// the completion ordinal is the last position the turn can own. Anything the
480/// session records afterwards belongs to no turn, or to the next one.
481pub fn load_materialized_turn_summary(
482    session_id: &str,
483    turn_start_position: u64,
484    turn_completed_position: u64,
485) -> Result<TurnSummary> {
486    load_materialized_turn_summary_from(
487        &database_path(),
488        session_id,
489        turn_start_position,
490        turn_completed_position,
491    )
492}
493
494pub(super) fn load_materialized_turn_summary_from(
495    path: &Path,
496    session_id: &str,
497    turn_start_position: u64,
498    turn_completed_position: u64,
499) -> Result<TurnSummary> {
500    let mut reader = open_reader(path)?;
501    let connection = reader.transaction()?;
502    let turn_number = connection.query_row(
503        "SELECT COUNT(*)
504         FROM materialized_transcript_items
505         WHERE session_id = ?1
506           AND position <= ?3
507           AND (
508               stable_id GLOB ?2
509               OR json_extract(
510                   CASE
511                       WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
512                       THEN body_json
513                       ELSE '{}'
514                   END,
515                   '$.kind'
516               ) = 'user'
517           )",
518        params![
519            session_id,
520            format!("{}*", mj_core::transcript::HARNESS_TURN_ITEM_PREFIX),
521            turn_start_position
522        ],
523        |row| row.get::<_, u64>(0),
524    )?;
525    let (turn_started_at_ms, last_changed_at_ms) = connection.query_row(
526        "SELECT COALESCE(MIN(created_at_ms), 0), COALESCE(MAX(last_changed_at_ms), 0)
527         FROM materialized_transcript_items
528         WHERE session_id = ?1 AND position >= ?2 AND position <= ?3",
529        params![session_id, turn_start_position, turn_completed_position],
530        |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
531    )?;
532    let final_message = last_materialized_agent_message_within(
533        &connection,
534        session_id,
535        turn_start_position,
536        turn_completed_position,
537    )?;
538    let tool_calls = connection.query_row(
539        "SELECT COUNT(*)
540         FROM materialized_transcript_items
541         WHERE session_id = ?1 AND position >= ?2 AND position <= ?3
542           AND json_extract(body_json, '$.kind') = 'tool'",
543        params![session_id, turn_start_position, turn_completed_position],
544        |row| row.get::<_, u64>(0),
545    )?;
546    Ok(TurnSummary {
547        turn_number,
548        turn_started_at_ms,
549        last_changed_at_ms,
550        final_message,
551        tool_calls,
552    })
553}
554
555pub fn load_materialized_transcript_filtered(
556    session_id: &str,
557    after_seq: u64,
558    limit: usize,
559    roles: Vec<mj_core::transcript::TranscriptRole>,
560    finished_only: bool,
561) -> Result<Option<TranscriptPage>> {
562    load_materialized_transcript_filtered_from(
563        &database_path(),
564        session_id,
565        after_seq,
566        limit,
567        roles,
568        finished_only,
569    )
570}
571
572pub(super) fn load_materialized_transcript_filtered_from(
573    path: &Path,
574    session_id: &str,
575    after_seq: u64,
576    limit: usize,
577    roles: Vec<mj_core::transcript::TranscriptRole>,
578    finished_only: bool,
579) -> Result<Option<TranscriptPage>> {
580    let mut reader = open_reader(path)?;
581    let connection = reader.transaction()?;
582    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
583        return Ok(None);
584    };
585    let role_kinds: Vec<_> = if finished_only {
586        vec!["agent"]
587    } else {
588        roles.iter().map(|role| role.storage_kind()).collect()
589    };
590    let role_kinds_json = serde_json::to_string(&role_kinds)?;
591    let open_agent_seq = if finished_only {
592        connection.query_row(
593            "SELECT MIN(COALESCE(latest_content_event_ordinal, position))
594                 FROM materialized_transcript_items
595                 WHERE session_id = ?1
596                   AND COALESCE(latest_content_event_ordinal, position) > ?2
597                   AND json_extract(body_json, '$.kind') = 'agent'
598                   AND json_extract(body_json, '$.streaming') = 1",
599            params![session_id, after_seq],
600            |row| row.get::<_, Option<u64>>(0),
601        )?
602    } else {
603        None
604    };
605    let mut statement = connection.prepare(
606        "WITH matches AS (
607             SELECT *, COALESCE(latest_content_event_ordinal, position) AS seq
608             FROM materialized_transcript_items WHERE session_id = ?1
609             AND COALESCE(latest_content_event_ordinal, position) > ?2
610             AND (json_array_length(?4) = 0 OR json_extract(body_json, '$.kind')
611                  IN (SELECT value FROM json_each(?4)))
612             AND (?5 = 0 OR (
613                 json_extract(body_json, '$.kind') = 'agent'
614                 AND json_extract(body_json, '$.streaming') = 0
615                 AND (?6 IS NULL OR COALESCE(latest_content_event_ordinal, position) < ?6)
616             ))
617         ), boundary AS (SELECT MAX(seq) AS seq FROM (SELECT seq FROM matches ORDER BY seq LIMIT ?3))
618         SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
619                last_changed_at_ms, body_json
620         FROM matches WHERE seq <= (SELECT seq FROM boundary)
621         ORDER BY seq, stable_id",
622    )?;
623    let rows = statement
624        .query_map(
625            params![
626                session_id,
627                after_seq,
628                limit.clamp(1, 1000) as i64,
629                role_kinds_json,
630                i64::from(finished_only),
631                open_agent_seq,
632            ],
633            |row| {
634                Ok((
635                    row.get::<_, String>(0)?,
636                    row.get::<_, u64>(1)?,
637                    row.get::<_, Option<u64>>(2)?,
638                    row.get::<_, i64>(3)?,
639                    row.get::<_, i64>(4)?,
640                    row.get::<_, String>(5)?,
641                ))
642            },
643        )?
644        .collect::<rusqlite::Result<Vec<_>>>()?;
645    let items = rows
646        .into_iter()
647        .map(
648            |(
649                stable_id,
650                position,
651                latest_content_event_ordinal,
652                created_at_ms,
653                last_changed_at_ms,
654                body_json,
655            )| {
656                Ok(Arc::new(TranscriptItem {
657                    stable_id,
658                    position,
659                    latest_content_event_ordinal,
660                    created_at_ms,
661                    last_changed_at_ms,
662                    body: decode_transcript_body(&body_json, session_id)?,
663                }))
664            },
665        )
666        .collect::<Result<Vec<_>>>()?;
667    let latest_seq = connection.query_row(
668        "SELECT COALESCE(MAX(COALESCE(latest_content_event_ordinal, position)), 0)
669         FROM materialized_transcript_items
670         WHERE session_id = ?1",
671        [session_id],
672        |row| row.get::<_, u64>(0),
673    )?;
674    let last_seq = items.last().map_or(after_seq, |item| item.seq());
675    let more: bool = connection.query_row(
676        "SELECT EXISTS(
677             SELECT 1 FROM materialized_transcript_items
678             WHERE session_id = ?1
679               AND COALESCE(latest_content_event_ordinal, position) > ?2
680               AND (json_array_length(?3) = 0 OR json_extract(body_json, '$.kind')
681                    IN (SELECT value FROM json_each(?3)))
682               AND (?4 = 0 OR (
683                   json_extract(body_json, '$.kind') = 'agent'
684                   AND json_extract(body_json, '$.streaming') = 0
685                   AND (?5 IS NULL OR COALESCE(latest_content_event_ordinal, position) < ?5)
686               ))
687         )",
688        params![
689            session_id,
690            last_seq,
691            role_kinds_json,
692            i64::from(finished_only),
693            open_agent_seq
694        ],
695        |row| row.get(0),
696    )?;
697    let next_after_seq = if more {
698        last_seq
699    } else if finished_only {
700        open_agent_seq.map_or_else(
701            || latest_seq.max(after_seq),
702            |open_seq| open_seq.saturating_sub(1).max(after_seq),
703        )
704    } else {
705        latest_seq.max(after_seq)
706    };
707    Ok(Some(TranscriptPage {
708        next_after_seq,
709        items,
710        latest_seq,
711        execution: fields.execution,
712    }))
713}
714
715/// How many transcript rows one retention pass rewrites.
716///
717/// The daemon is the single database writer, so a pass that rewrote every row
718/// of a long session would stall every other write behind it. A capped pass
719/// leaves the rest for the next checkpoint, which is the next time any of it
720/// becomes redundant anyway.
721pub(super) const RETENTION_BATCH_ITEMS: usize = 4_096;
722
723/// Rows below this are already small enough that rewriting them would cost
724/// more than it reclaims.
725pub(super) const RETENTION_BODY_FLOOR_BYTES: usize = 4 * 1024;
726
727/// Drop tool output that a verified checkpoint already holds.
728///
729/// The projection only ever grew: the only deletes were a per-item remove, a
730/// whole-session wipe, and the `sessions` cascade. One measured session reached
731/// 28,066 items and 635 MiB, of which 561 MB was tool-call content.
732///
733/// A checkpoint archive carries the complete transcript up to its event
734/// frontier, and one checkpoint per session is retained, so every item at or
735/// below `event_frontier` is durably recorded elsewhere. What stays here is
736/// what the transcript still shows: which tool ran, on what, with what result,
737/// and each edit's diffstat. See
738/// [`mj_transcript::transcript::compact_tool_call_for_retention`].
739pub fn compact_materialized_transcript_through(
740    session_id: &str,
741    event_frontier: u64,
742) -> Result<TranscriptRetention> {
743    let session_id = session_id.to_owned();
744    submit_database_write("compact_materialized_transcript", move |_| {
745        compact_materialized_transcript_in(&database_path(), &session_id, event_frontier)
746    })
747}
748
749pub(super) fn compact_materialized_transcript_in(
750    path: &Path,
751    session_id: &str,
752    event_frontier: u64,
753) -> Result<TranscriptRetention> {
754    let mut connection = open(path)?;
755    let candidates = {
756        let mut statement = connection.prepare(
757            "SELECT stable_id, body_json
758             FROM materialized_transcript_items
759             WHERE session_id = ?1
760               AND position <= ?2
761               AND length(body_json) > ?3
762               AND json_extract(
763                   CASE WHEN json_valid(body_json) THEN body_json ELSE '{}' END,
764                   '$.kind'
765               ) = 'tool'
766             ORDER BY position, stable_id
767             LIMIT ?4",
768        )?;
769        statement
770            .query_map(
771                params![
772                    session_id,
773                    event_frontier,
774                    RETENTION_BODY_FLOOR_BYTES as i64,
775                    RETENTION_BATCH_ITEMS as i64 + 1
776                ],
777                |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
778            )?
779            .collect::<rusqlite::Result<Vec<_>>>()?
780    };
781    let remaining = candidates.len() > RETENTION_BATCH_ITEMS;
782    let mut retention = TranscriptRetention {
783        remaining,
784        ..TranscriptRetention::default()
785    };
786    let transaction =
787        connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
788    for (stable_id, body_json) in candidates.into_iter().take(RETENTION_BATCH_ITEMS) {
789        let mut body: TranscriptBody = match serde_json::from_str(&body_json) {
790            Ok(body) => body,
791            // A row this cannot read is a row it must not rewrite.
792            Err(error) => {
793                tracing::warn!(%session_id, %stable_id, %error, "skipping unreadable transcript body");
794                continue;
795            }
796        };
797        if !mj_transcript::transcript::compact_tool_call_for_retention(&mut body) {
798            continue;
799        }
800        let compacted = serde_json::to_string(&body)
801            .with_context(|| format!("serialize compacted transcript body {stable_id}"))?;
802        if compacted.len() >= body_json.len() {
803            continue;
804        }
805        transaction.execute(
806            "UPDATE materialized_transcript_items SET body_json = ?3
807             WHERE session_id = ?1 AND stable_id = ?2",
808            params![session_id, stable_id, compacted],
809        )?;
810        retention.items += 1;
811        retention.bytes += body_json.len() - compacted.len();
812    }
813    transaction.commit()?;
814    Ok(retention)
815}
816
817/// How many transcript items a polled projection carries.
818///
819/// Every viewer of a polled projection is bounded already: the conversation
820/// pane keeps `chat::TAIL_SEED_ITEMS` (256) entries, and the browser
821/// transcript keeps 1,000 rendered lines. This is set above both, since an
822/// entry renders to at least one line, so the window is the whole of what any
823/// of them would show.
824pub const PROJECTION_TAIL_ITEMS: usize = 1_024;
825
826/// Load the same safe window the live projector maintains, without first
827/// deserializing all of the conversation. Read metadata and bodies together.
828pub fn load_materialized_actor_projection(
829    session_id: &str,
830) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
831    load_materialized_actor_projection_from(&database_path(), session_id)
832}
833
834pub(super) fn load_materialized_actor_projection_from(
835    path: &Path,
836    session_id: &str,
837) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
838    let mut reader = open_reader(path)?;
839    let connection = reader.transaction()?;
840    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
841        return Ok(None);
842    };
843    let latest_turn = last_materialized_turn_start(&connection, session_id)?;
844    let desired: Option<u64> = connection
845        .query_row(
846            "SELECT position FROM materialized_transcript_items WHERE session_id=?1
847         ORDER BY position DESC, stable_id DESC LIMIT 1 OFFSET ?2",
848            params![session_id, PROJECTION_TAIL_ITEMS - 1],
849            |row| row.get(0),
850        )
851        .optional()?;
852    // Pending tools and streams can outlive their initiating turn. Keep their
853    // turn too; the same rule is applied by ProjectionWindow::trim.
854    let mutable: Option<u64> = connection.query_row(
855        "SELECT MIN(position) FROM materialized_transcript_items WHERE session_id=?1
856         AND (json_extract(body_json, '$.streaming')=1
857              OR (json_extract(body_json, '$.kind')='tool'
858                  AND json_extract(body_json, '$.call.status') IN ('pending','in_progress')))",
859        [session_id],
860        |row| row.get(0),
861    )?;
862    let boundary = desired
863        .unwrap_or(0)
864        .min(latest_turn.unwrap_or(0))
865        .min(mutable.unwrap_or(u64::MAX));
866    let start: u64 = connection.query_row(
867        "SELECT COALESCE(MAX(position),0) FROM materialized_transcript_items
868         WHERE session_id=?1 AND position<=?2
869         AND (json_extract(body_json, '$.kind')='user' OR stable_id LIKE 'harness-turn:%')",
870        params![session_id, boundary],
871        |row| row.get(0),
872    )?;
873    let transcript = read_transcript_range(&connection, session_id, start, None, None)?;
874    let omitted_items = connection.query_row(
875        "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id=?1 AND position<?2",
876        params![session_id, start],
877        |row| row.get(0),
878    )?;
879    let window = ProjectionWindow {
880        omitted_items,
881        provisional_title: first_materialized_user_message(&connection, session_id)?
882            .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
883        latest_turn_start_position: latest_turn,
884    };
885    let materialized = MaterializedSession {
886        session_id: session_id.to_owned(),
887        applied_event_ordinal: fields.applied_event_ordinal,
888        applied_event_digest: fields.applied_event_digest,
889        last_activity_at_ms: fields.last_activity_at_ms,
890        execution: fields.execution,
891        session_title: fields.session_title,
892        configuration: fields.configuration,
893        transcript,
894        queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
895        pending_elicitations: fields.pending_elicitations,
896        active_turn: fields.active_turn,
897        last_turn_outcome: fields.last_turn_outcome,
898    };
899    materialized.validate()?;
900    Ok(Some((materialized, window)))
901}
902
903pub fn load_transcript_history(
904    session_id: &str,
905    before: Option<&mj_core::storage::TranscriptCursor>,
906    limit: usize,
907) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
908    load_transcript_history_from(&database_path(), session_id, before, limit)
909}
910
911pub(super) fn load_transcript_history_from(
912    path: &Path,
913    session_id: &str,
914    before: Option<&mj_core::storage::TranscriptCursor>,
915    limit: usize,
916) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
917    let mut reader = open_reader(path)?;
918    let connection = reader.transaction()?;
919    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
920        return Ok(None);
921    };
922    let limit = limit.clamp(1, 256);
923    let mut items = read_transcript_range(&connection, session_id, 0, before, Some(limit + 1))?;
924    let has_more = items.len() > limit;
925    if has_more {
926        items.remove(0);
927    }
928    let before = has_more.then(|| mj_core::storage::TranscriptCursor::of(&items[0]));
929    Ok(Some(mj_core::storage::TranscriptHistoryPage {
930        items,
931        before,
932        frontier: fields.applied_event_ordinal,
933    }))
934}
935
936fn read_transcript_range(
937    connection: &Connection,
938    session_id: &str,
939    start: u64,
940    before: Option<&mj_core::storage::TranscriptCursor>,
941    limit: Option<usize>,
942) -> Result<Vec<Arc<TranscriptItem>>> {
943    let mut statement = connection.prepare(
944        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
945                last_changed_at_ms, body_json FROM materialized_transcript_items
946         WHERE session_id=?1 AND position>=?2
947           AND (position,stable_id)<(?3,?4)
948         ORDER BY position DESC, stable_id DESC LIMIT ?5",
949    )?;
950    let rows = statement
951        .query_map(
952            params![
953                session_id,
954                start,
955                before.map_or(i64::MAX as u64, |c| c.position),
956                before.map_or("", |c| c.stable_id.as_str()),
957                limit.map_or(-1, |limit| limit as i64)
958            ],
959            |row| {
960                Ok((
961                    row.get::<_, String>(0)?,
962                    row.get::<_, u64>(1)?,
963                    row.get::<_, Option<u64>>(2)?,
964                    row.get::<_, i64>(3)?,
965                    row.get::<_, i64>(4)?,
966                    row.get::<_, String>(5)?,
967                ))
968            },
969        )?
970        .collect::<rusqlite::Result<Vec<_>>>()?;
971    let mut items = rows
972        .into_iter()
973        .map(
974            |(
975                stable_id,
976                position,
977                latest_content_event_ordinal,
978                created_at_ms,
979                last_changed_at_ms,
980                body,
981            )| {
982                Ok(Arc::new(TranscriptItem {
983                    stable_id,
984                    position,
985                    latest_content_event_ordinal,
986                    created_at_ms,
987                    last_changed_at_ms,
988                    body: decode_transcript_body(&body, session_id)?,
989                }))
990            },
991        )
992        .collect::<Result<Vec<_>>>()?;
993    items.reverse();
994    Ok(items)
995}
996
997/// Rehydrate explicitly referenced historical tools before deriving a delta.
998/// A window is a cache, never evidence that a durable tool does not exist.
999pub fn load_projection_references(
1000    projection: &MaterializedSession,
1001    events: &[mj_core::relay::RelayEvent],
1002) -> Result<Vec<Arc<TranscriptItem>>> {
1003    load_projection_references_from(&database_path(), projection, events)
1004}
1005
1006pub(super) fn load_projection_references_from(
1007    path: &Path,
1008    projection: &MaterializedSession,
1009    events: &[mj_core::relay::RelayEvent],
1010) -> Result<Vec<Arc<TranscriptItem>>> {
1011    let (mut ids, terminals) = mj_transcript::projection::historical_references(events)?;
1012    let retained = projection
1013        .transcript
1014        .iter()
1015        .map(|item| item.stable_id.as_str())
1016        .collect::<std::collections::HashSet<_>>();
1017    ids.retain(|id| !retained.contains(id.as_str()));
1018    if ids.is_empty() && terminals.is_empty() {
1019        return Ok(Vec::new());
1020    }
1021    let mut reader = open_reader(path)?;
1022    let connection = reader.transaction()?;
1023    // Ordinary tool updates use the stable-ID index. Only terminal attachment
1024    // events need the JSON reverse-reference scan.
1025    let mut parameters = vec![projection.session_id.clone(), serde_json::to_string(&ids)?];
1026    let mut statement = connection.prepare(if terminals.is_empty() {
1027        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1028                last_changed_at_ms, body_json FROM materialized_transcript_items
1029         WHERE session_id=?1 AND stable_id IN (SELECT value FROM json_each(?2))
1030         ORDER BY position,stable_id"
1031    } else {
1032        parameters.push(serde_json::to_string(&retained)?);
1033        parameters.push(serde_json::to_string(&terminals)?);
1034        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1035                last_changed_at_ms, body_json FROM materialized_transcript_items
1036         WHERE session_id=?1 AND stable_id NOT IN (SELECT value FROM json_each(?3))
1037           AND (stable_id IN (SELECT value FROM json_each(?2))
1038             OR EXISTS(SELECT 1 FROM json_each(body_json, '$.terminal_refs')
1039                       WHERE value IN (SELECT value FROM json_each(?4))))
1040         ORDER BY position,stable_id"
1041    })?;
1042    let rows = statement
1043        .query_map(rusqlite::params_from_iter(parameters), |row| {
1044            Ok((
1045                row.get::<_, String>(0)?,
1046                row.get::<_, u64>(1)?,
1047                row.get::<_, Option<u64>>(2)?,
1048                row.get::<_, i64>(3)?,
1049                row.get::<_, i64>(4)?,
1050                row.get::<_, String>(5)?,
1051            ))
1052        })?
1053        .collect::<rusqlite::Result<Vec<_>>>()?;
1054    rows.into_iter()
1055        .map(
1056            |(
1057                stable_id,
1058                position,
1059                latest_content_event_ordinal,
1060                created_at_ms,
1061                last_changed_at_ms,
1062                body,
1063            )| {
1064                Ok(Arc::new(TranscriptItem {
1065                    stable_id,
1066                    position,
1067                    latest_content_event_ordinal,
1068                    created_at_ms,
1069                    last_changed_at_ms,
1070                    body: decode_transcript_body(&body, &projection.session_id)?,
1071                }))
1072            },
1073        )
1074        .collect()
1075}
1076
1077/// Load a projection carrying only the end of its transcript.
1078///
1079/// The steady-state poll reloads a session's projection every time anything
1080/// about it moves. Loading the whole transcript to do that is work
1081/// proportional to everything that has ever happened in the conversation —
1082/// 635 MiB and 28,066 items on one measured session — for a view that shows
1083/// the last few hundred entries. This reads the window instead, plus the two
1084/// facts that live outside it, each with one indexed query. See
1085/// [`ProjectionWindow`].
1086pub fn load_materialized_projection_tail(
1087    session_id: &str,
1088    transcript_limit: usize,
1089) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
1090    load_materialized_projection_tail_from(&database_path(), session_id, transcript_limit)
1091}
1092
1093pub(super) fn load_materialized_projection_tail_from(
1094    path: &Path,
1095    session_id: &str,
1096    transcript_limit: usize,
1097) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
1098    let mut reader = open_reader(path)?;
1099    // Frontier, mutable transcript bodies, and window metadata must describe
1100    // one WAL snapshot even if the daemon commits between these queries.
1101    let connection = reader.transaction()?;
1102    let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
1103        return Ok(None);
1104    };
1105    let transcript = read_materialized_transcript(&connection, session_id, Some(transcript_limit))?;
1106    let total_items = connection.query_row(
1107        "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id = ?1",
1108        [session_id],
1109        |row| row.get::<_, usize>(0),
1110    )?;
1111    let window = ProjectionWindow {
1112        omitted_items: total_items.saturating_sub(transcript.len()),
1113        provisional_title: first_materialized_user_message(&connection, session_id)?
1114            .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
1115        latest_turn_start_position: last_materialized_turn_start(&connection, session_id)?,
1116    };
1117    let materialized = MaterializedSession {
1118        session_id: session_id.to_owned(),
1119        applied_event_ordinal: fields.applied_event_ordinal,
1120        applied_event_digest: fields.applied_event_digest,
1121        last_activity_at_ms: fields.last_activity_at_ms,
1122        execution: fields.execution,
1123        session_title: fields.session_title,
1124        configuration: fields.configuration,
1125        transcript,
1126        queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
1127        pending_elicitations: fields.pending_elicitations,
1128        active_turn: fields.active_turn,
1129        last_turn_outcome: fields.last_turn_outcome,
1130    };
1131    materialized.validate()?;
1132    Ok(Some((materialized, window)))
1133}
1134
1135/// Read only the projection's event frontier. Deciding whether a stored
1136/// projection already matches an archive costs one row this way, instead of
1137/// deserializing every transcript item to compare two integers.
1138pub fn materialized_event_frontier(session_id: &str) -> Result<Option<(u64, String)>> {
1139    materialized_event_frontier_from(&database_path(), session_id)
1140}
1141
1142pub(super) fn materialized_event_frontier_from(
1143    path: &Path,
1144    session_id: &str,
1145) -> Result<Option<(u64, String)>> {
1146    Ok(open_reader(path)?
1147        .query_row(
1148            "SELECT applied_event_ordinal, applied_event_digest
1149             FROM materialized_sessions WHERE session_id = ?1",
1150            [session_id],
1151            |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1152        )
1153        .optional()?)
1154}
1155
1156/// Replace a session's durable prompt queue without touching its transcript or
1157/// event frontier. Resume uses this when it keeps the stored projection but
1158/// still has to drop the queue the archive carried.
1159pub fn replace_materialized_queued_prompts(
1160    session_id: &str,
1161    queued_prompts: &[MaterializedQueuedPrompt],
1162) -> Result<()> {
1163    let session_id = session_id.to_owned();
1164    let queued_prompts = queued_prompts.to_vec();
1165    submit_database_write("replace_materialized_queued_prompts", move |_| {
1166        replace_materialized_queued_prompts_in(&database_path(), &session_id, &queued_prompts)
1167    })
1168}
1169
1170pub(super) fn replace_materialized_queued_prompts_in(
1171    path: &Path,
1172    session_id: &str,
1173    queued_prompts: &[MaterializedQueuedPrompt],
1174) -> Result<()> {
1175    let mut connection = open(path)?;
1176    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1177    if !session_exists(&tx, session_id)? {
1178        bail!("unknown session {session_id}");
1179    }
1180    replace_materialized_queue(&tx, session_id, queued_prompts)?;
1181    tx.commit()?;
1182    Ok(())
1183}
1184
1185/// The activity watermark of every session whose projection holds at least one
1186/// transcript item, by session id.
1187///
1188/// This is a change token, not a projection: session indexing needs to know
1189/// which live conversations have moved since it last looked, and loading each
1190/// one's transcript to find out would cost the whole corpus every sync.
1191pub fn load_transcribed_session_activity() -> Result<BTreeMap<String, Option<i64>>> {
1192    load_transcribed_session_activity_from(&database_path())
1193}
1194
1195fn load_transcribed_session_activity_from(path: &Path) -> Result<BTreeMap<String, Option<i64>>> {
1196    let connection = open_reader(path)?;
1197    let mut statement = connection.prepare(
1198        "SELECT session_id, last_activity_at_ms
1199         FROM materialized_sessions s
1200         WHERE EXISTS (
1201             SELECT 1 FROM materialized_transcript_items i
1202             WHERE i.session_id = s.session_id
1203         )",
1204    )?;
1205    let rows = statement.query_map([], |row| {
1206        Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?))
1207    })?;
1208    let mut activity = BTreeMap::new();
1209    for row in rows {
1210        let (session_id, last_activity_at_ms) = row?;
1211        activity.insert(session_id, last_activity_at_ms);
1212    }
1213    Ok(activity)
1214}
1215
1216/// Load only the durable prompt queues without deserializing transcript rows.
1217/// Dashboard startup uses this path so work is proportional to queued prompts,
1218/// not to the complete retained conversation history.
1219pub fn load_materialized_queued_prompts() -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>>
1220{
1221    load_materialized_queued_prompts_from(&database_path())
1222}
1223
1224pub(super) fn load_materialized_queued_prompts_from(
1225    path: &Path,
1226) -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>> {
1227    let connection = open_reader(path)?;
1228    let mut statement = connection.prepare(
1229        "SELECT session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1230         FROM materialized_queued_prompts
1231         ORDER BY session_id, ordinal",
1232    )?;
1233    let rows = statement.query_map([], |row| {
1234        Ok((
1235            row.get::<_, String>(0)?,
1236            row.get::<_, String>(1)?,
1237            row.get::<_, String>(2)?,
1238            row.get::<_, String>(3)?,
1239            row.get::<_, i64>(4)?,
1240            row.get::<_, Option<u64>>(5)?,
1241        ))
1242    })?;
1243    let mut queues = BTreeMap::<String, Vec<MaterializedQueuedPrompt>>::new();
1244    for row in rows {
1245        let (session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal) =
1246            row?;
1247        let content = serde_json::from_str(&content_json).with_context(|| {
1248            format!("parse materialized queued prompt for session {session_id}")
1249        })?;
1250        let kind = serde_json::from_str(&kind_json).with_context(|| {
1251            format!("parse materialized queue entry kind for session {session_id}")
1252        })?;
1253        queues
1254            .entry(session_id)
1255            .or_default()
1256            .push(MaterializedQueuedPrompt {
1257                command_id,
1258                kind,
1259                content,
1260                queued_at_ms,
1261                accepted_ordinal,
1262            });
1263    }
1264    Ok(queues)
1265}
1266
1267pub(super) fn load_materialized_session_from(
1268    path: &Path,
1269    session_id: &str,
1270) -> Result<Option<MaterializedSession>> {
1271    let mut reader = open_reader(path)?;
1272    let connection = reader.transaction()?;
1273    load_materialized_session_with(&connection, session_id)
1274}
1275
1276pub(super) fn load_materialized_session_with(
1277    connection: &rusqlite::Transaction<'_>,
1278    session_id: &str,
1279) -> Result<Option<MaterializedSession>> {
1280    let Some(fields) = read_materialized_session_fields(connection, session_id)? else {
1281        return Ok(None);
1282    };
1283    let materialized = MaterializedSession {
1284        session_id: session_id.to_owned(),
1285        applied_event_ordinal: fields.applied_event_ordinal,
1286        applied_event_digest: fields.applied_event_digest,
1287        last_activity_at_ms: fields.last_activity_at_ms,
1288        execution: fields.execution,
1289        session_title: fields.session_title,
1290        configuration: fields.configuration,
1291        transcript: read_materialized_transcript(connection, session_id, None)?,
1292        queued_prompts: read_materialized_queued_prompts(connection, session_id)?,
1293        pending_elicitations: fields.pending_elicitations,
1294        active_turn: fields.active_turn,
1295        last_turn_outcome: fields.last_turn_outcome,
1296    };
1297    materialized.validate()?;
1298    Ok(Some(materialized))
1299}
1300
1301/// Everything a projection holds apart from its transcript and its queue.
1302pub(super) struct MaterializedSessionFields {
1303    pub(super) applied_event_ordinal: u64,
1304    pub(super) applied_event_digest: String,
1305    pub(super) last_activity_at_ms: Option<i64>,
1306    pub(super) execution: MaterializedExecutionState,
1307    pub(super) session_title: Option<String>,
1308    pub(super) configuration: mj_core::state::SessionConfiguration,
1309    pub(super) pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
1310    pub(super) active_turn: Option<MaterializedTurn>,
1311    pub(super) last_turn_outcome: Option<MaterializedTurnOutcome>,
1312}
1313
1314pub(super) fn read_materialized_session_fields(
1315    connection: &Connection,
1316    session_id: &str,
1317) -> Result<Option<MaterializedSessionFields>> {
1318    let row = connection
1319        .query_row(
1320            "SELECT applied_event_ordinal, applied_event_digest, last_activity_at_ms,
1321                    execution_state, running_started_at_ms, session_title, configuration_json,
1322                    pending_elicitations_json, active_turn_json, last_turn_outcome_json
1323             FROM materialized_sessions WHERE session_id = ?1",
1324            [session_id],
1325            |row| {
1326                Ok((
1327                    row.get::<_, u64>(0)?,
1328                    row.get::<_, String>(1)?,
1329                    row.get::<_, Option<i64>>(2)?,
1330                    row.get::<_, String>(3)?,
1331                    row.get::<_, Option<i64>>(4)?,
1332                    row.get::<_, Option<String>>(5)?,
1333                    row.get::<_, String>(6)?,
1334                    row.get::<_, String>(7)?,
1335                    row.get::<_, Option<String>>(8)?,
1336                    row.get::<_, Option<String>>(9)?,
1337                ))
1338            },
1339        )
1340        .optional()?;
1341    let Some((
1342        applied_event_ordinal,
1343        applied_event_digest,
1344        last_activity_at_ms,
1345        execution,
1346        running_started_at_ms,
1347        session_title,
1348        configuration_json,
1349        pending_elicitations_json,
1350        active_turn_json,
1351        last_turn_outcome_json,
1352    )) = row
1353    else {
1354        return Ok(None);
1355    };
1356    #[cfg(test)]
1357    super::tests::after_materialized_frontier_read();
1358    Ok(Some(MaterializedSessionFields {
1359        applied_event_ordinal,
1360        applied_event_digest,
1361        last_activity_at_ms,
1362        execution: parse_materialized_execution(&execution, running_started_at_ms)?,
1363        session_title,
1364        configuration: serde_json::from_str(&configuration_json).with_context(|| {
1365            format!("parse materialized configuration for session {session_id}")
1366        })?,
1367        pending_elicitations: serde_json::from_str(&pending_elicitations_json)
1368            .with_context(|| format!("parse pending elicitations for session {session_id}"))?,
1369        active_turn: active_turn_json
1370            .as_deref()
1371            .map(serde_json::from_str)
1372            .transpose()
1373            .with_context(|| format!("parse active turn for session {session_id}"))?,
1374        last_turn_outcome: last_turn_outcome_json
1375            .as_deref()
1376            .map(serde_json::from_str)
1377            .transpose()
1378            .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
1379    }))
1380}
1381
1382/// Read a session's transcript, oldest first. `limit` reads only that many
1383/// items from the end, walking the `materialized_transcript_position` index
1384/// backwards so the read costs the rows it returns.
1385pub(super) fn read_materialized_transcript(
1386    connection: &Connection,
1387    session_id: &str,
1388    limit: Option<usize>,
1389) -> Result<Vec<Arc<TranscriptItem>>> {
1390    let mut statement = connection.prepare(match limit {
1391        Some(_) => {
1392            "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1393                    last_changed_at_ms, body_json
1394             FROM materialized_transcript_items
1395             WHERE session_id = ?1
1396             ORDER BY position DESC, stable_id DESC
1397             LIMIT ?2"
1398        }
1399        None => {
1400            "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1401                    last_changed_at_ms, body_json
1402             FROM materialized_transcript_items
1403             WHERE session_id = ?1
1404             ORDER BY position, stable_id"
1405        }
1406    })?;
1407    let read = |row: &rusqlite::Row<'_>| {
1408        Ok((
1409            row.get::<_, String>(0)?,
1410            row.get::<_, u64>(1)?,
1411            row.get::<_, Option<u64>>(2)?,
1412            row.get::<_, i64>(3)?,
1413            row.get::<_, i64>(4)?,
1414            row.get::<_, String>(5)?,
1415        ))
1416    };
1417    let rows = match limit {
1418        Some(limit) => statement
1419            .query_map(params![session_id, limit as i64], read)?
1420            .collect::<rusqlite::Result<Vec<_>>>()?,
1421        None => statement
1422            .query_map([session_id], read)?
1423            .collect::<rusqlite::Result<Vec<_>>>()?,
1424    };
1425    let mut transcript = rows
1426        .into_iter()
1427        .map(
1428            |(
1429                stable_id,
1430                position,
1431                latest_content_event_ordinal,
1432                created_at_ms,
1433                last_changed_at_ms,
1434                body_json,
1435            )| {
1436                Ok(Arc::new(TranscriptItem {
1437                    stable_id,
1438                    position,
1439                    latest_content_event_ordinal,
1440                    created_at_ms,
1441                    last_changed_at_ms,
1442                    body: decode_transcript_body(&body_json, session_id)?,
1443                }))
1444            },
1445        )
1446        .collect::<Result<Vec<_>>>()?;
1447    if limit.is_some() {
1448        // The bounded query walks the index backwards to bound what it reads;
1449        // every caller wants the transcript in the order it was written.
1450        transcript.reverse();
1451    }
1452    Ok(transcript)
1453}
1454
1455pub(super) fn read_materialized_queued_prompts(
1456    connection: &Connection,
1457    session_id: &str,
1458) -> Result<Vec<MaterializedQueuedPrompt>> {
1459    let mut statement = connection.prepare(
1460        "SELECT command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1461         FROM materialized_queued_prompts
1462         WHERE session_id = ?1
1463         ORDER BY ordinal",
1464    )?;
1465    let rows = statement
1466        .query_map([session_id], |row| {
1467            Ok((
1468                row.get::<_, String>(0)?,
1469                row.get::<_, String>(1)?,
1470                row.get::<_, String>(2)?,
1471                row.get::<_, i64>(3)?,
1472                row.get::<_, Option<u64>>(4)?,
1473            ))
1474        })?
1475        .collect::<rusqlite::Result<Vec<_>>>()?;
1476    rows.into_iter()
1477        .map(
1478            |(command_id, kind_json, content_json, queued_at_ms, accepted_ordinal)| {
1479                Ok(MaterializedQueuedPrompt {
1480                    command_id,
1481                    kind: serde_json::from_str(&kind_json).with_context(|| {
1482                        format!("parse materialized queue entry kind for session {session_id}")
1483                    })?,
1484                    content: serde_json::from_str(&content_json).with_context(|| {
1485                        format!("parse materialized queued prompt for session {session_id}")
1486                    })?,
1487                    queued_at_ms,
1488                    accepted_ordinal,
1489                })
1490            },
1491        )
1492        .collect()
1493}
1494
1495/// Replace a complete projection, primarily when seeding a restored
1496/// checkpoint. Operational `SessionRecord` metadata and read receipts are not
1497/// modified.
1498pub fn save_materialized_session(materialized: &MaterializedSession) -> Result<()> {
1499    let materialized = materialized.clone();
1500    submit_database_write("save_materialized_session", move |_| {
1501        save_materialized_session_to(&database_path(), &materialized)
1502    })
1503}
1504
1505pub(super) fn save_materialized_session_to(
1506    path: &Path,
1507    materialized: &MaterializedSession,
1508) -> Result<()> {
1509    materialized.validate()?;
1510    let mut connection = open(path)?;
1511    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1512    if !session_exists(&tx, &materialized.session_id)? {
1513        bail!("unknown session {}", materialized.session_id);
1514    }
1515    write_materialized_session(&tx, materialized)?;
1516    tx.commit()?;
1517    Ok(())
1518}
1519
1520/// One relay page being applied inside a single write transaction. The relay
1521/// retains everything past the last acknowledgement, so a page that fails
1522/// part-way rolls back to the previous durable frontier and is simply
1523/// redelivered. Only a committed page may be acknowledged.
1524pub struct ProjectionPage<'a> {
1525    pub(super) session_id: &'a str,
1526    pub(super) transaction: Transaction<'a>,
1527    pub(super) applied_ordinal: u64,
1528    pub(super) applied_digest: String,
1529    pub(super) dirty: bool,
1530    pub(super) pending: MaterializedSessionMutation,
1531    pub(super) pending_transcript: BTreeMap<String, PendingTranscriptMutation>,
1532    pub(super) pending_turns: Vec<MaterializedTurnOutcome>,
1533    /// Configuration changes are coalesced for the session, but selection is
1534    /// captured at each turn start before later events can overwrite it.
1535    usage_configuration: mj_core::state::SessionConfiguration,
1536    pub(super) pending_events: Vec<(i64, ApiEventData)>,
1537}
1538
1539pub(super) struct PendingTranscriptMutation {
1540    pub(super) final_mutation: TranscriptMutation,
1541    pub(super) remove_before_upsert: bool,
1542}
1543
1544impl ProjectionPage<'_> {
1545    /// Apply the projection effects of the next relay event to the open page.
1546    /// The event must continue the chain the page has reached so far, which is
1547    /// the persisted frontier plus every event already applied to this page.
1548    pub fn apply(
1549        &mut self,
1550        event_ordinal: u64,
1551        previous_event_digest: &str,
1552        event_digest: &str,
1553        mutation: &MaterializedSessionMutation,
1554    ) -> Result<ProjectionApplyOutcome> {
1555        if event_ordinal == 0 {
1556            bail!("relay event ordinal must be positive");
1557        }
1558        // A v2 event carries no chain link (empty previous digest). Its
1559        // continuity to the projection frontier is proven by ordinal
1560        // contiguity plus the attach cursor the controller validated against
1561        // the worker, not by an in-record back-reference; divergence is caught
1562        // there, before any event is applied.
1563        let chained = !previous_event_digest.is_empty();
1564        if chained {
1565            validate_relay_event_digest(previous_event_digest, "previous relay event digest")?;
1566        }
1567        validate_relay_event_frontier(event_ordinal, event_digest, "relay event frontier")?;
1568        let session_id = self.session_id;
1569        let applied = self.applied_ordinal;
1570        if event_ordinal < applied {
1571            return Ok(ProjectionApplyOutcome::AlreadyApplied);
1572        }
1573        if event_ordinal == applied {
1574            if event_digest != self.applied_digest {
1575                bail!(
1576                    "relay event digest mismatch for session {session_id} at ordinal {event_ordinal}: projection has {}, received {event_digest}",
1577                    self.applied_digest
1578                );
1579            }
1580            return Ok(ProjectionApplyOutcome::AlreadyApplied);
1581        }
1582        let expected = applied
1583            .checked_add(1)
1584            .context("materialized event ordinal overflow")?;
1585        if event_ordinal != expected {
1586            bail!(
1587                "relay event gap for session {session_id}: expected ordinal {expected}, received {event_ordinal}"
1588            );
1589        }
1590        if chained && previous_event_digest != self.applied_digest {
1591            bail!(
1592                "relay event chain diverged for session {session_id} before ordinal {event_ordinal}: projection has {}, event follows {previous_event_digest}",
1593                self.applied_digest
1594            );
1595        }
1596
1597        if let Some(event) = &mutation.native_agent {
1598            native_agents::apply_native_agent_event(&self.transaction, session_id, event)?;
1599        }
1600        if let Some(activity_at_ms) = mutation.last_activity_at_ms {
1601            self.pending.last_activity_at_ms = Some(
1602                self.pending
1603                    .last_activity_at_ms
1604                    .map_or(activity_at_ms, |existing| existing.max(activity_at_ms)),
1605            );
1606        }
1607        if let Some(execution) = mutation.execution {
1608            self.pending.execution = Some(execution);
1609        }
1610        if let Some(title) = &mutation.session_title {
1611            if title.as_ref().is_some_and(|title| title.trim().is_empty()) {
1612                bail!("materialized session title cannot be empty");
1613            }
1614            self.pending.session_title = Some(title.clone());
1615        }
1616        if let Some(configuration) = &mutation.configuration {
1617            self.usage_configuration = configuration.clone();
1618            self.pending.configuration = Some(configuration.clone());
1619        }
1620        for item_mutation in &mutation.transcript {
1621            match item_mutation {
1622                TranscriptMutation::Upsert(item) => {
1623                    item.validate(event_ordinal)?;
1624                    let stable_id = item.stable_id.clone();
1625                    let entry = self.pending_transcript.entry(stable_id).or_insert_with(|| {
1626                        PendingTranscriptMutation {
1627                            final_mutation: TranscriptMutation::Upsert(item.clone()),
1628                            remove_before_upsert: false,
1629                        }
1630                    });
1631                    entry.remove_before_upsert |=
1632                        matches!(&entry.final_mutation, TranscriptMutation::Remove { .. });
1633                    entry.final_mutation = TranscriptMutation::Upsert(item.clone());
1634                }
1635                TranscriptMutation::Remove { stable_id } => {
1636                    if stable_id.trim().is_empty() {
1637                        bail!("cannot remove a transcript item with an empty stable id");
1638                    }
1639                    let removed = TranscriptMutation::Remove {
1640                        stable_id: stable_id.clone(),
1641                    };
1642                    self.pending_transcript
1643                        .entry(stable_id.clone())
1644                        .and_modify(|entry| entry.final_mutation = removed.clone())
1645                        .or_insert(PendingTranscriptMutation {
1646                            final_mutation: removed,
1647                            remove_before_upsert: false,
1648                        });
1649                }
1650            }
1651        }
1652        if let Some(queued_prompts) = &mutation.queued_prompts {
1653            self.pending.queued_prompts = Some(queued_prompts.clone());
1654        }
1655        if let Some(pending_elicitations) = &mutation.pending_elicitations {
1656            self.pending.pending_elicitations = Some(pending_elicitations.clone());
1657        }
1658        self.pending
1659            .config_results
1660            .extend(mutation.config_results.clone());
1661        if let Some(active_turn) = &mutation.active_turn {
1662            if let Some(turn) = active_turn {
1663                super::usage::record_turn_selection(
1664                    &self.transaction,
1665                    session_id,
1666                    &turn.command_id,
1667                    &self.usage_configuration,
1668                )?;
1669            }
1670            self.pending.active_turn = Some(active_turn.clone());
1671        }
1672        if mutation.clear_turn_outcome {
1673            self.pending.clear_turn_outcome = true;
1674            self.pending.last_turn_outcome = None;
1675        }
1676        if let Some(last_turn_outcome) = &mutation.last_turn_outcome {
1677            self.pending_turns.push(last_turn_outcome.clone());
1678            self.pending.last_turn_outcome = Some(last_turn_outcome.clone());
1679        }
1680        if let Some(cost) = &mutation.provider_cost {
1681            self.pending.provider_cost = Some(cost.clone());
1682        }
1683        self.pending_events.extend(
1684            mutation
1685                .api_events
1686                .iter()
1687                .cloned()
1688                .map(|event| (mutation.last_activity_at_ms.unwrap_or(0), event)),
1689        );
1690        self.applied_ordinal = event_ordinal;
1691        event_digest.clone_into(&mut self.applied_digest);
1692        self.dirty = true;
1693        Ok(ProjectionApplyOutcome::Applied)
1694    }
1695
1696    /// Persist the coalesced final state of this page. Intermediate event
1697    /// frontiers are useful only for chain validation: a page commits or rolls
1698    /// back as a unit, so writing them individually adds no recovery value.
1699    pub(super) fn flush(&mut self) -> Result<()> {
1700        if !self.dirty {
1701            return Ok(());
1702        }
1703        let tx = &self.transaction;
1704        let session_id = self.session_id;
1705        if let Some(execution) = self.pending.execution {
1706            let (state, started_at_ms) = materialized_execution_columns(execution);
1707            tx.execute(
1708                "UPDATE materialized_sessions
1709                 SET execution_state = ?2, running_started_at_ms = ?3
1710                 WHERE session_id = ?1",
1711                params![session_id, state, started_at_ms],
1712            )?;
1713        }
1714        if let Some(title) = &self.pending.session_title {
1715            tx.execute(
1716                "UPDATE materialized_sessions SET session_title = ?2 WHERE session_id = ?1",
1717                params![session_id, title],
1718            )?;
1719        }
1720        if let Some(configuration) = &self.pending.configuration {
1721            tx.execute(
1722                "UPDATE materialized_sessions SET configuration_json = ?2 WHERE session_id = ?1",
1723                params![session_id, serde_json::to_string(configuration)?],
1724            )?;
1725        }
1726        for pending in self.pending_transcript.values() {
1727            match &pending.final_mutation {
1728                TranscriptMutation::Upsert(item) => {
1729                    // A remove followed by an upsert deliberately starts a new
1730                    // item identity. Preserve that boundary even though other
1731                    // repeated updates are coalesced to one write.
1732                    if pending.remove_before_upsert {
1733                        tx.execute(
1734                            "DELETE FROM materialized_transcript_items
1735                             WHERE session_id = ?1 AND stable_id = ?2",
1736                            params![session_id, item.stable_id],
1737                        )?;
1738                    }
1739                    upsert_transcript_item(tx, session_id, item)?;
1740                }
1741                TranscriptMutation::Remove { stable_id } => {
1742                    tx.execute(
1743                        "DELETE FROM materialized_transcript_items
1744                         WHERE session_id = ?1 AND stable_id = ?2",
1745                        params![session_id, stable_id],
1746                    )?;
1747                }
1748            }
1749        }
1750        if let Some(queued_prompts) = &self.pending.queued_prompts {
1751            replace_materialized_queue(tx, session_id, queued_prompts)?;
1752        }
1753        if let Some(pending_elicitations) = &self.pending.pending_elicitations {
1754            tx.execute(
1755                "UPDATE materialized_sessions
1756                 SET pending_elicitations_json = ?2 WHERE session_id = ?1",
1757                params![session_id, serde_json::to_string(pending_elicitations)?],
1758            )?;
1759        }
1760        for (recorded_at_ms, event) in &self.pending_events {
1761            events::insert_api_event(tx, session_id, *recorded_at_ms, event)?;
1762        }
1763        for turn in &self.pending_turns {
1764            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)?])?;
1765        }
1766        if let Some(cost) = &self.pending.provider_cost {
1767            tx.execute(
1768                "INSERT OR REPLACE INTO session_provider_cost(session_id, body) VALUES (?1, ?2)",
1769                params![session_id, serde_json::to_string(cost)?],
1770            )?;
1771        }
1772        for (command_id, error) in &self.pending.config_results {
1773            tx.execute("INSERT OR REPLACE INTO api_config_results(session_id, command_id, error) VALUES (?1, ?2, ?3)", params![session_id, command_id, error])?;
1774        }
1775        if let Some(active_turn) = &self.pending.active_turn {
1776            tx.execute(
1777                "UPDATE materialized_sessions SET active_turn_json = ?2 WHERE session_id = ?1",
1778                params![
1779                    session_id,
1780                    active_turn
1781                        .as_ref()
1782                        .map(serde_json::to_string)
1783                        .transpose()?
1784                ],
1785            )?;
1786        }
1787        if self.pending.clear_turn_outcome {
1788            self.transaction.execute("UPDATE materialized_sessions SET last_turn_outcome_json = NULL WHERE session_id = ?1", [session_id])?;
1789        }
1790        if let Some(last_turn_outcome) = &self.pending.last_turn_outcome {
1791            tx.execute(
1792                "UPDATE materialized_sessions
1793                 SET last_turn_outcome_json = ?2 WHERE session_id = ?1",
1794                params![session_id, serde_json::to_string(last_turn_outcome)?],
1795            )?;
1796        }
1797        tx.execute(
1798            "UPDATE materialized_sessions
1799             SET last_activity_at_ms = CASE
1800                     WHEN ?2 IS NULL THEN last_activity_at_ms
1801                     WHEN last_activity_at_ms IS NULL OR last_activity_at_ms < ?2 THEN ?2
1802                     ELSE last_activity_at_ms
1803                 END,
1804                 applied_event_ordinal = ?3,
1805                 applied_event_digest = ?4
1806             WHERE session_id = ?1",
1807            params![
1808                session_id,
1809                self.pending.last_activity_at_ms,
1810                self.applied_ordinal,
1811                self.applied_digest,
1812            ],
1813        )?;
1814        Ok(())
1815    }
1816}
1817
1818/// Apply one relay page in a single transaction. `fill` feeds the page's
1819/// events through [`ProjectionPage::apply`]; the projection changes and the
1820/// event frontier commit together only when `fill` succeeds, so callers may
1821/// acknowledge the page's last ordinal to the relay after this returns.
1822pub fn apply_projection_page<T>(
1823    session_id: &str,
1824    fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T> + Send + 'static,
1825) -> Result<T>
1826where
1827    T: Send + 'static,
1828{
1829    let session_id = session_id.to_owned();
1830    submit_database_write("apply_projection_page", move |connection| {
1831        apply_projection_page_with(connection, &session_id, fill)
1832    })
1833}
1834
1835#[cfg(test)]
1836pub(super) fn apply_projection_page_to<T>(
1837    path: &Path,
1838    session_id: &str,
1839    fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1840) -> Result<T> {
1841    let mut connection = open(path)?;
1842    apply_projection_page_with(&mut connection, session_id, fill)
1843}
1844
1845pub(super) fn apply_projection_page_with<T>(
1846    connection: &mut Connection,
1847    session_id: &str,
1848    fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1849) -> Result<T> {
1850    let transaction =
1851        connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1852    let (applied_ordinal, applied_digest) = transaction
1853        .query_row(
1854            "SELECT applied_event_ordinal, applied_event_digest
1855             FROM materialized_sessions WHERE session_id = ?1",
1856            [session_id],
1857            |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1858        )
1859        .optional()?
1860        .with_context(|| format!("unknown session {session_id}"))?;
1861    validate_relay_event_frontier(
1862        applied_ordinal,
1863        &applied_digest,
1864        "persisted relay event frontier",
1865    )?;
1866    let usage_configuration = read_materialized_session_fields(&transaction, session_id)?
1867        .context("projection session disappeared")?
1868        .configuration;
1869    let mut page = ProjectionPage {
1870        session_id,
1871        transaction,
1872        applied_ordinal,
1873        applied_digest,
1874        dirty: false,
1875        pending: MaterializedSessionMutation::default(),
1876        pending_transcript: BTreeMap::new(),
1877        pending_turns: Vec::new(),
1878        usage_configuration,
1879        pending_events: Vec::new(),
1880    };
1881    // Dropping the page on failure rolls the whole transaction back, leaving
1882    // the projection at the frontier the relay last saw acknowledged.
1883    let filled = fill(&mut page)?;
1884    page.flush()?;
1885    page.transaction.commit()?;
1886    Ok(filled)
1887}
1888
1889/// Apply exactly one relay event, as a page of one.
1890pub fn apply_projection_event(
1891    session_id: &str,
1892    event_ordinal: u64,
1893    previous_event_digest: &str,
1894    event_digest: &str,
1895    mutation: &MaterializedSessionMutation,
1896) -> Result<ProjectionApplyOutcome> {
1897    let session_id = session_id.to_owned();
1898    let previous_event_digest = previous_event_digest.to_owned();
1899    let event_digest = event_digest.to_owned();
1900    let mutation = mutation.clone();
1901    submit_database_write("apply_projection_event", move |connection| {
1902        apply_projection_page_with(connection, &session_id, |page| {
1903            page.apply(
1904                event_ordinal,
1905                &previous_event_digest,
1906                &event_digest,
1907                &mutation,
1908            )
1909        })
1910    })
1911}
1912
1913#[cfg(test)]
1914pub(super) fn apply_projection_event_to(
1915    path: &Path,
1916    session_id: &str,
1917    event_ordinal: u64,
1918    previous_event_digest: &str,
1919    event_digest: &str,
1920    mutation: &MaterializedSessionMutation,
1921) -> Result<ProjectionApplyOutcome> {
1922    apply_projection_page_to(path, session_id, |page| {
1923        page.apply(event_ordinal, previous_event_digest, event_digest, mutation)
1924    })
1925}
1926
1927/// Advance the persisted detach/read receipt monotonically. A receipt cannot
1928/// acknowledge an event the controller projection has not durably applied.
1929pub fn advance_viewed_through_event_ordinal(session_id: &str, through: u64) -> Result<u64> {
1930    let session_id = session_id.to_owned();
1931    submit_database_write("advance_viewed_through_event_ordinal", move |_| {
1932        advance_viewed_through_event_ordinal_to(&database_path(), &session_id, through)
1933    })
1934}
1935
1936pub(super) fn advance_viewed_through_event_ordinal_to(
1937    path: &Path,
1938    session_id: &str,
1939    through: u64,
1940) -> Result<u64> {
1941    let mut connection = open(path)?;
1942    let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1943    let applied = tx
1944        .query_row(
1945            "SELECT applied_event_ordinal FROM materialized_sessions WHERE session_id = ?1",
1946            [session_id],
1947            |row| row.get::<_, u64>(0),
1948        )
1949        .optional()?
1950        .with_context(|| format!("unknown session {session_id}"))?;
1951    if through > applied {
1952        bail!(
1953            "cannot acknowledge event ordinal {through} for session {session_id}; projection is at {applied}"
1954        );
1955    }
1956    tx.execute(
1957        "UPDATE sessions
1958         SET viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?2)
1959         WHERE session_id = ?1",
1960        params![session_id, through],
1961    )?;
1962    let receipt = tx.query_row(
1963        "SELECT viewed_through_event_ordinal FROM sessions WHERE session_id = ?1",
1964        [session_id],
1965        |row| row.get::<_, u64>(0),
1966    )?;
1967    tx.commit()?;
1968    Ok(receipt)
1969}
1970
1971/// Decode a stored transcript body, merging runs of streamed text chunks.
1972///
1973/// Rows written before streamed text was merged on the way in hold one chunk
1974/// per token, which costs far more memory decoded than the text it carries.
1975/// Collapsing them here shrinks long sessions without rewriting stored JSON.
1976fn decode_transcript_body(body_json: &str, session_id: &str) -> Result<TranscriptBody> {
1977    let mut body: TranscriptBody = serde_json::from_str(body_json)
1978        .with_context(|| format!("parse materialized transcript body for session {session_id}"))?;
1979    match &mut body {
1980        TranscriptBody::Agent { chunks, .. } | TranscriptBody::Thought { chunks, .. } => {
1981            mj_core::transcript::coalesce_content_chunks(chunks);
1982        }
1983        _ => {}
1984    }
1985    Ok(body)
1986}
1987
1988/// Read complete authorization and recent replies in one consistent snapshot.
1989/// Runs on the continuation service's blocking pool, never a render/event loop.
1990pub(crate) fn load_continuation_evidence(
1991    session_id: &str,
1992    ordinal: u64,
1993    digest: &str,
1994) -> Result<mj_core::continuation::ContinuationEvidence> {
1995    load_continuation_evidence_from(&database_path(), session_id, ordinal, digest)
1996}
1997
1998pub(super) fn load_continuation_evidence_from(
1999    path: &Path,
2000    session_id: &str,
2001    ordinal: u64,
2002    digest: &str,
2003) -> Result<mj_core::continuation::ContinuationEvidence> {
2004    let mut reader = open_reader(path)?;
2005    let connection = reader.transaction()?;
2006    let fields = read_materialized_session_fields(&connection, session_id)?
2007        .context("continuation projection is missing")?;
2008    anyhow::ensure!(
2009        fields.applied_event_ordinal == ordinal && fields.applied_event_digest == digest,
2010        "continuation projection changed"
2011    );
2012    let mut statement = connection.prepare(
2013        "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
2014                last_changed_at_ms, body_json
2015         FROM materialized_transcript_items WHERE session_id = ?1
2016         AND (json_extract(body_json, '$.kind') IN ('user', 'agent')
2017              OR stable_id LIKE 'context-cleared:%')
2018         ORDER BY position DESC, stable_id DESC",
2019    )?;
2020    let rows = statement.query_map([session_id], |row| {
2021        Ok((
2022            row.get::<_, String>(0)?,
2023            row.get::<_, u64>(1)?,
2024            row.get::<_, Option<u64>>(2)?,
2025            row.get::<_, i64>(3)?,
2026            row.get::<_, i64>(4)?,
2027            row.get::<_, String>(5)?,
2028        ))
2029    })?;
2030    crate::continuation::evidence_from_items(rows.map(|row| {
2031        let (
2032            stable_id,
2033            position,
2034            latest_content_event_ordinal,
2035            created_at_ms,
2036            last_changed_at_ms,
2037            body_json,
2038        ) = row?;
2039        // Refuse exceptionally large source messages rather than clip consent.
2040        anyhow::ensure!(
2041            body_json.len() <= 1024 * 1024,
2042            "continuation source message exceeds budget"
2043        );
2044        Ok(Arc::new(TranscriptItem {
2045            stable_id,
2046            position,
2047            latest_content_event_ordinal,
2048            created_at_ms,
2049            last_changed_at_ms,
2050            body: decode_transcript_body(&body_json, session_id)?,
2051        }))
2052    }))
2053}