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