1use super::*;
2
3pub fn load_materialized_session(session_id: &str) -> Result<Option<MaterializedSession>> {
13 load_materialized_session_from(&database_path(), session_id)
14}
15
16pub 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
137pub(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
204pub(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
245pub(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
314pub type MaterializedTurnState = (
316 MaterializedExecutionState,
317 Option<MaterializedTurn>,
318 Option<MaterializedTurnOutcome>,
319);
320
321pub 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
343pub 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
380pub 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
544pub(super) const RETENTION_BATCH_ITEMS: usize = 4_096;
551
552pub(super) const RETENTION_BODY_FLOOR_BYTES: usize = 4 * 1024;
555
556pub 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 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
646pub const PROJECTION_TAIL_ITEMS: usize = 1_024;
654
655pub fn load_materialized_projection_tail(
665 session_id: &str,
666 transcript_limit: usize,
667) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
668 load_materialized_projection_tail_from(&database_path(), session_id, transcript_limit)
669}
670
671pub(super) fn load_materialized_projection_tail_from(
672 path: &Path,
673 session_id: &str,
674 transcript_limit: usize,
675) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
676 let mut reader = open_reader(path)?;
677 let connection = reader.transaction()?;
680 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
681 return Ok(None);
682 };
683 let transcript = read_materialized_transcript(&connection, session_id, Some(transcript_limit))?;
684 let total_items = connection.query_row(
685 "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id = ?1",
686 [session_id],
687 |row| row.get::<_, usize>(0),
688 )?;
689 let window = ProjectionWindow {
690 omitted_items: total_items.saturating_sub(transcript.len()),
691 provisional_title: first_materialized_user_message(&connection, session_id)?
692 .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
693 latest_turn_start_position: last_materialized_turn_start(&connection, session_id)?,
694 };
695 let materialized = MaterializedSession {
696 session_id: session_id.to_owned(),
697 applied_event_ordinal: fields.applied_event_ordinal,
698 applied_event_digest: fields.applied_event_digest,
699 last_activity_at_ms: fields.last_activity_at_ms,
700 execution: fields.execution,
701 session_title: fields.session_title,
702 configuration: fields.configuration,
703 transcript,
704 queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
705 pending_elicitations: fields.pending_elicitations,
706 active_turn: fields.active_turn,
707 last_turn_outcome: fields.last_turn_outcome,
708 };
709 materialized.validate()?;
710 Ok(Some((materialized, window)))
711}
712
713pub fn materialized_event_frontier(session_id: &str) -> Result<Option<(u64, String)>> {
717 materialized_event_frontier_from(&database_path(), session_id)
718}
719
720pub(super) fn materialized_event_frontier_from(
721 path: &Path,
722 session_id: &str,
723) -> Result<Option<(u64, String)>> {
724 Ok(open_reader(path)?
725 .query_row(
726 "SELECT applied_event_ordinal, applied_event_digest
727 FROM materialized_sessions WHERE session_id = ?1",
728 [session_id],
729 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
730 )
731 .optional()?)
732}
733
734pub fn replace_materialized_queued_prompts(
738 session_id: &str,
739 queued_prompts: &[MaterializedQueuedPrompt],
740) -> Result<()> {
741 let session_id = session_id.to_owned();
742 let queued_prompts = queued_prompts.to_vec();
743 submit_database_write("replace_materialized_queued_prompts", move |_| {
744 replace_materialized_queued_prompts_in(&database_path(), &session_id, &queued_prompts)
745 })
746}
747
748pub(super) fn replace_materialized_queued_prompts_in(
749 path: &Path,
750 session_id: &str,
751 queued_prompts: &[MaterializedQueuedPrompt],
752) -> Result<()> {
753 let mut connection = open(path)?;
754 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
755 if !session_exists(&tx, session_id)? {
756 bail!("unknown session {session_id}");
757 }
758 replace_materialized_queue(&tx, session_id, queued_prompts)?;
759 tx.commit()?;
760 Ok(())
761}
762
763pub fn load_transcribed_session_activity() -> Result<BTreeMap<String, Option<i64>>> {
770 load_transcribed_session_activity_from(&database_path())
771}
772
773fn load_transcribed_session_activity_from(path: &Path) -> Result<BTreeMap<String, Option<i64>>> {
774 let connection = open_reader(path)?;
775 let mut statement = connection.prepare(
776 "SELECT session_id, last_activity_at_ms
777 FROM materialized_sessions s
778 WHERE EXISTS (
779 SELECT 1 FROM materialized_transcript_items i
780 WHERE i.session_id = s.session_id
781 )",
782 )?;
783 let rows = statement.query_map([], |row| {
784 Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?))
785 })?;
786 let mut activity = BTreeMap::new();
787 for row in rows {
788 let (session_id, last_activity_at_ms) = row?;
789 activity.insert(session_id, last_activity_at_ms);
790 }
791 Ok(activity)
792}
793
794pub fn load_materialized_queued_prompts() -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>>
798{
799 load_materialized_queued_prompts_from(&database_path())
800}
801
802pub(super) fn load_materialized_queued_prompts_from(
803 path: &Path,
804) -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>> {
805 let connection = open_reader(path)?;
806 let mut statement = connection.prepare(
807 "SELECT session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
808 FROM materialized_queued_prompts
809 ORDER BY session_id, ordinal",
810 )?;
811 let rows = statement.query_map([], |row| {
812 Ok((
813 row.get::<_, String>(0)?,
814 row.get::<_, String>(1)?,
815 row.get::<_, String>(2)?,
816 row.get::<_, String>(3)?,
817 row.get::<_, i64>(4)?,
818 row.get::<_, Option<u64>>(5)?,
819 ))
820 })?;
821 let mut queues = BTreeMap::<String, Vec<MaterializedQueuedPrompt>>::new();
822 for row in rows {
823 let (session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal) =
824 row?;
825 let content = serde_json::from_str(&content_json).with_context(|| {
826 format!("parse materialized queued prompt for session {session_id}")
827 })?;
828 let kind = serde_json::from_str(&kind_json).with_context(|| {
829 format!("parse materialized queue entry kind for session {session_id}")
830 })?;
831 queues
832 .entry(session_id)
833 .or_default()
834 .push(MaterializedQueuedPrompt {
835 command_id,
836 kind,
837 content,
838 queued_at_ms,
839 accepted_ordinal,
840 });
841 }
842 Ok(queues)
843}
844
845pub(super) fn load_materialized_session_from(
846 path: &Path,
847 session_id: &str,
848) -> Result<Option<MaterializedSession>> {
849 let mut reader = open_reader(path)?;
850 let connection = reader.transaction()?;
851 load_materialized_session_with(&connection, session_id)
852}
853
854pub(super) fn load_materialized_session_with(
855 connection: &rusqlite::Transaction<'_>,
856 session_id: &str,
857) -> Result<Option<MaterializedSession>> {
858 let Some(fields) = read_materialized_session_fields(connection, session_id)? else {
859 return Ok(None);
860 };
861 let materialized = MaterializedSession {
862 session_id: session_id.to_owned(),
863 applied_event_ordinal: fields.applied_event_ordinal,
864 applied_event_digest: fields.applied_event_digest,
865 last_activity_at_ms: fields.last_activity_at_ms,
866 execution: fields.execution,
867 session_title: fields.session_title,
868 configuration: fields.configuration,
869 transcript: read_materialized_transcript(connection, session_id, None)?,
870 queued_prompts: read_materialized_queued_prompts(connection, session_id)?,
871 pending_elicitations: fields.pending_elicitations,
872 active_turn: fields.active_turn,
873 last_turn_outcome: fields.last_turn_outcome,
874 };
875 materialized.validate()?;
876 Ok(Some(materialized))
877}
878
879pub(super) struct MaterializedSessionFields {
881 pub(super) applied_event_ordinal: u64,
882 pub(super) applied_event_digest: String,
883 pub(super) last_activity_at_ms: Option<i64>,
884 pub(super) execution: MaterializedExecutionState,
885 pub(super) session_title: Option<String>,
886 pub(super) configuration: BTreeMap<String, serde_json::Value>,
887 pub(super) pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
888 pub(super) active_turn: Option<MaterializedTurn>,
889 pub(super) last_turn_outcome: Option<MaterializedTurnOutcome>,
890}
891
892pub(super) fn read_materialized_session_fields(
893 connection: &Connection,
894 session_id: &str,
895) -> Result<Option<MaterializedSessionFields>> {
896 let row = connection
897 .query_row(
898 "SELECT applied_event_ordinal, applied_event_digest, last_activity_at_ms,
899 execution_state, running_started_at_ms, session_title, configuration_json,
900 pending_elicitations_json, active_turn_json, last_turn_outcome_json
901 FROM materialized_sessions WHERE session_id = ?1",
902 [session_id],
903 |row| {
904 Ok((
905 row.get::<_, u64>(0)?,
906 row.get::<_, String>(1)?,
907 row.get::<_, Option<i64>>(2)?,
908 row.get::<_, String>(3)?,
909 row.get::<_, Option<i64>>(4)?,
910 row.get::<_, Option<String>>(5)?,
911 row.get::<_, String>(6)?,
912 row.get::<_, String>(7)?,
913 row.get::<_, Option<String>>(8)?,
914 row.get::<_, Option<String>>(9)?,
915 ))
916 },
917 )
918 .optional()?;
919 let Some((
920 applied_event_ordinal,
921 applied_event_digest,
922 last_activity_at_ms,
923 execution,
924 running_started_at_ms,
925 session_title,
926 configuration_json,
927 pending_elicitations_json,
928 active_turn_json,
929 last_turn_outcome_json,
930 )) = row
931 else {
932 return Ok(None);
933 };
934 #[cfg(test)]
935 super::tests::after_materialized_frontier_read();
936 Ok(Some(MaterializedSessionFields {
937 applied_event_ordinal,
938 applied_event_digest,
939 last_activity_at_ms,
940 execution: parse_materialized_execution(&execution, running_started_at_ms)?,
941 session_title,
942 configuration: serde_json::from_str(&configuration_json).with_context(|| {
943 format!("parse materialized configuration for session {session_id}")
944 })?,
945 pending_elicitations: serde_json::from_str(&pending_elicitations_json)
946 .with_context(|| format!("parse pending elicitations for session {session_id}"))?,
947 active_turn: active_turn_json
948 .as_deref()
949 .map(serde_json::from_str)
950 .transpose()
951 .with_context(|| format!("parse active turn for session {session_id}"))?,
952 last_turn_outcome: last_turn_outcome_json
953 .as_deref()
954 .map(serde_json::from_str)
955 .transpose()
956 .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
957 }))
958}
959
960pub(super) fn read_materialized_transcript(
964 connection: &Connection,
965 session_id: &str,
966 limit: Option<usize>,
967) -> Result<Vec<Arc<TranscriptItem>>> {
968 let mut statement = connection.prepare(match limit {
969 Some(_) => {
970 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
971 last_changed_at_ms, body_json
972 FROM materialized_transcript_items
973 WHERE session_id = ?1
974 ORDER BY position DESC, stable_id DESC
975 LIMIT ?2"
976 }
977 None => {
978 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
979 last_changed_at_ms, body_json
980 FROM materialized_transcript_items
981 WHERE session_id = ?1
982 ORDER BY position, stable_id"
983 }
984 })?;
985 let read = |row: &rusqlite::Row<'_>| {
986 Ok((
987 row.get::<_, String>(0)?,
988 row.get::<_, u64>(1)?,
989 row.get::<_, Option<u64>>(2)?,
990 row.get::<_, i64>(3)?,
991 row.get::<_, i64>(4)?,
992 row.get::<_, String>(5)?,
993 ))
994 };
995 let rows = match limit {
996 Some(limit) => statement
997 .query_map(params![session_id, limit as i64], read)?
998 .collect::<rusqlite::Result<Vec<_>>>()?,
999 None => statement
1000 .query_map([session_id], read)?
1001 .collect::<rusqlite::Result<Vec<_>>>()?,
1002 };
1003 let mut transcript = rows
1004 .into_iter()
1005 .map(
1006 |(
1007 stable_id,
1008 position,
1009 latest_content_event_ordinal,
1010 created_at_ms,
1011 last_changed_at_ms,
1012 body_json,
1013 )| {
1014 Ok(Arc::new(TranscriptItem {
1015 stable_id,
1016 position,
1017 latest_content_event_ordinal,
1018 created_at_ms,
1019 last_changed_at_ms,
1020 body: decode_transcript_body(&body_json, session_id)?,
1021 }))
1022 },
1023 )
1024 .collect::<Result<Vec<_>>>()?;
1025 if limit.is_some() {
1026 transcript.reverse();
1029 }
1030 Ok(transcript)
1031}
1032
1033pub(super) fn read_materialized_queued_prompts(
1034 connection: &Connection,
1035 session_id: &str,
1036) -> Result<Vec<MaterializedQueuedPrompt>> {
1037 let mut statement = connection.prepare(
1038 "SELECT command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1039 FROM materialized_queued_prompts
1040 WHERE session_id = ?1
1041 ORDER BY ordinal",
1042 )?;
1043 let rows = statement
1044 .query_map([session_id], |row| {
1045 Ok((
1046 row.get::<_, String>(0)?,
1047 row.get::<_, String>(1)?,
1048 row.get::<_, String>(2)?,
1049 row.get::<_, i64>(3)?,
1050 row.get::<_, Option<u64>>(4)?,
1051 ))
1052 })?
1053 .collect::<rusqlite::Result<Vec<_>>>()?;
1054 rows.into_iter()
1055 .map(
1056 |(command_id, kind_json, content_json, queued_at_ms, accepted_ordinal)| {
1057 Ok(MaterializedQueuedPrompt {
1058 command_id,
1059 kind: serde_json::from_str(&kind_json).with_context(|| {
1060 format!("parse materialized queue entry kind for session {session_id}")
1061 })?,
1062 content: serde_json::from_str(&content_json).with_context(|| {
1063 format!("parse materialized queued prompt for session {session_id}")
1064 })?,
1065 queued_at_ms,
1066 accepted_ordinal,
1067 })
1068 },
1069 )
1070 .collect()
1071}
1072
1073pub fn save_materialized_session(materialized: &MaterializedSession) -> Result<()> {
1077 let materialized = materialized.clone();
1078 submit_database_write("save_materialized_session", move |_| {
1079 save_materialized_session_to(&database_path(), &materialized)
1080 })
1081}
1082
1083pub(super) fn save_materialized_session_to(
1084 path: &Path,
1085 materialized: &MaterializedSession,
1086) -> Result<()> {
1087 materialized.validate()?;
1088 let mut connection = open(path)?;
1089 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1090 if !session_exists(&tx, &materialized.session_id)? {
1091 bail!("unknown session {}", materialized.session_id);
1092 }
1093 write_materialized_session(&tx, materialized)?;
1094 tx.commit()?;
1095 Ok(())
1096}
1097
1098pub struct ProjectionPage<'a> {
1103 pub(super) session_id: &'a str,
1104 pub(super) transaction: Transaction<'a>,
1105 pub(super) applied_ordinal: u64,
1106 pub(super) applied_digest: String,
1107 pub(super) dirty: bool,
1108 pub(super) pending: MaterializedSessionMutation,
1109 pub(super) pending_transcript: BTreeMap<String, PendingTranscriptMutation>,
1110 pub(super) pending_turns: Vec<MaterializedTurnOutcome>,
1111 pub(super) pending_events: Vec<(i64, ApiEventData)>,
1112}
1113
1114pub(super) struct PendingTranscriptMutation {
1115 pub(super) final_mutation: TranscriptMutation,
1116 pub(super) remove_before_upsert: bool,
1117}
1118
1119impl ProjectionPage<'_> {
1120 pub fn apply(
1124 &mut self,
1125 event_ordinal: u64,
1126 previous_event_digest: &str,
1127 event_digest: &str,
1128 mutation: &MaterializedSessionMutation,
1129 ) -> Result<ProjectionApplyOutcome> {
1130 if event_ordinal == 0 {
1131 bail!("relay event ordinal must be positive");
1132 }
1133 let chained = !previous_event_digest.is_empty();
1139 if chained {
1140 validate_relay_event_digest(previous_event_digest, "previous relay event digest")?;
1141 }
1142 validate_relay_event_frontier(event_ordinal, event_digest, "relay event frontier")?;
1143 let session_id = self.session_id;
1144 let applied = self.applied_ordinal;
1145 if event_ordinal < applied {
1146 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1147 }
1148 if event_ordinal == applied {
1149 if event_digest != self.applied_digest {
1150 bail!(
1151 "relay event digest mismatch for session {session_id} at ordinal {event_ordinal}: projection has {}, received {event_digest}",
1152 self.applied_digest
1153 );
1154 }
1155 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1156 }
1157 let expected = applied
1158 .checked_add(1)
1159 .context("materialized event ordinal overflow")?;
1160 if event_ordinal != expected {
1161 bail!(
1162 "relay event gap for session {session_id}: expected ordinal {expected}, received {event_ordinal}"
1163 );
1164 }
1165 if chained && previous_event_digest != self.applied_digest {
1166 bail!(
1167 "relay event chain diverged for session {session_id} before ordinal {event_ordinal}: projection has {}, event follows {previous_event_digest}",
1168 self.applied_digest
1169 );
1170 }
1171
1172 if let Some(event) = &mutation.native_agent {
1173 native_agents::apply_native_agent_event(&self.transaction, session_id, event)?;
1174 }
1175 if let Some(activity_at_ms) = mutation.last_activity_at_ms {
1176 self.pending.last_activity_at_ms = Some(
1177 self.pending
1178 .last_activity_at_ms
1179 .map_or(activity_at_ms, |existing| existing.max(activity_at_ms)),
1180 );
1181 }
1182 if let Some(execution) = mutation.execution {
1183 self.pending.execution = Some(execution);
1184 }
1185 if let Some(title) = &mutation.session_title {
1186 if title.as_ref().is_some_and(|title| title.trim().is_empty()) {
1187 bail!("materialized session title cannot be empty");
1188 }
1189 self.pending.session_title = Some(title.clone());
1190 }
1191 if let Some(configuration) = &mutation.configuration {
1192 self.pending.configuration = Some(configuration.clone());
1193 }
1194 for item_mutation in &mutation.transcript {
1195 match item_mutation {
1196 TranscriptMutation::Upsert(item) => {
1197 item.validate(event_ordinal)?;
1198 let stable_id = item.stable_id.clone();
1199 let entry = self.pending_transcript.entry(stable_id).or_insert_with(|| {
1200 PendingTranscriptMutation {
1201 final_mutation: TranscriptMutation::Upsert(item.clone()),
1202 remove_before_upsert: false,
1203 }
1204 });
1205 entry.remove_before_upsert |=
1206 matches!(&entry.final_mutation, TranscriptMutation::Remove { .. });
1207 entry.final_mutation = TranscriptMutation::Upsert(item.clone());
1208 }
1209 TranscriptMutation::Remove { stable_id } => {
1210 if stable_id.trim().is_empty() {
1211 bail!("cannot remove a transcript item with an empty stable id");
1212 }
1213 let removed = TranscriptMutation::Remove {
1214 stable_id: stable_id.clone(),
1215 };
1216 self.pending_transcript
1217 .entry(stable_id.clone())
1218 .and_modify(|entry| entry.final_mutation = removed.clone())
1219 .or_insert(PendingTranscriptMutation {
1220 final_mutation: removed,
1221 remove_before_upsert: false,
1222 });
1223 }
1224 }
1225 }
1226 if let Some(queued_prompts) = &mutation.queued_prompts {
1227 self.pending.queued_prompts = Some(queued_prompts.clone());
1228 }
1229 if let Some(pending_elicitations) = &mutation.pending_elicitations {
1230 self.pending.pending_elicitations = Some(pending_elicitations.clone());
1231 }
1232 self.pending
1233 .config_results
1234 .extend(mutation.config_results.clone());
1235 if let Some(active_turn) = &mutation.active_turn {
1236 self.pending.active_turn = Some(active_turn.clone());
1237 }
1238 if mutation.clear_turn_outcome {
1239 self.pending.clear_turn_outcome = true;
1240 self.pending.last_turn_outcome = None;
1241 }
1242 if let Some(last_turn_outcome) = &mutation.last_turn_outcome {
1243 self.pending_turns.push(last_turn_outcome.clone());
1244 self.pending.last_turn_outcome = Some(last_turn_outcome.clone());
1245 }
1246 if let Some(cost) = &mutation.provider_cost {
1247 self.pending.provider_cost = Some(cost.clone());
1248 }
1249 self.pending_events.extend(
1250 mutation
1251 .api_events
1252 .iter()
1253 .cloned()
1254 .map(|event| (mutation.last_activity_at_ms.unwrap_or(0), event)),
1255 );
1256 self.applied_ordinal = event_ordinal;
1257 event_digest.clone_into(&mut self.applied_digest);
1258 self.dirty = true;
1259 Ok(ProjectionApplyOutcome::Applied)
1260 }
1261
1262 pub(super) fn flush(&mut self) -> Result<()> {
1266 if !self.dirty {
1267 return Ok(());
1268 }
1269 let tx = &self.transaction;
1270 let session_id = self.session_id;
1271 if let Some(execution) = self.pending.execution {
1272 let (state, started_at_ms) = materialized_execution_columns(execution);
1273 tx.execute(
1274 "UPDATE materialized_sessions
1275 SET execution_state = ?2, running_started_at_ms = ?3
1276 WHERE session_id = ?1",
1277 params![session_id, state, started_at_ms],
1278 )?;
1279 }
1280 if let Some(title) = &self.pending.session_title {
1281 tx.execute(
1282 "UPDATE materialized_sessions SET session_title = ?2 WHERE session_id = ?1",
1283 params![session_id, title],
1284 )?;
1285 }
1286 if let Some(configuration) = &self.pending.configuration {
1287 tx.execute(
1288 "UPDATE materialized_sessions SET configuration_json = ?2 WHERE session_id = ?1",
1289 params![session_id, serde_json::to_string(configuration)?],
1290 )?;
1291 }
1292 for pending in self.pending_transcript.values() {
1293 match &pending.final_mutation {
1294 TranscriptMutation::Upsert(item) => {
1295 if pending.remove_before_upsert {
1299 tx.execute(
1300 "DELETE FROM materialized_transcript_items
1301 WHERE session_id = ?1 AND stable_id = ?2",
1302 params![session_id, item.stable_id],
1303 )?;
1304 }
1305 upsert_transcript_item(tx, session_id, item)?;
1306 }
1307 TranscriptMutation::Remove { stable_id } => {
1308 tx.execute(
1309 "DELETE FROM materialized_transcript_items
1310 WHERE session_id = ?1 AND stable_id = ?2",
1311 params![session_id, stable_id],
1312 )?;
1313 }
1314 }
1315 }
1316 if let Some(queued_prompts) = &self.pending.queued_prompts {
1317 replace_materialized_queue(tx, session_id, queued_prompts)?;
1318 }
1319 if let Some(pending_elicitations) = &self.pending.pending_elicitations {
1320 tx.execute(
1321 "UPDATE materialized_sessions
1322 SET pending_elicitations_json = ?2 WHERE session_id = ?1",
1323 params![session_id, serde_json::to_string(pending_elicitations)?],
1324 )?;
1325 }
1326 for (recorded_at_ms, event) in &self.pending_events {
1327 events::insert_api_event(tx, session_id, *recorded_at_ms, event)?;
1328 }
1329 for turn in &self.pending_turns {
1330 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)?])?;
1331 }
1332 if let Some(cost) = &self.pending.provider_cost {
1333 tx.execute(
1334 "INSERT OR REPLACE INTO session_provider_cost(session_id, body) VALUES (?1, ?2)",
1335 params![session_id, serde_json::to_string(cost)?],
1336 )?;
1337 }
1338 for (command_id, error) in &self.pending.config_results {
1339 tx.execute("INSERT OR REPLACE INTO api_config_results(session_id, command_id, error) VALUES (?1, ?2, ?3)", params![session_id, command_id, error])?;
1340 }
1341 if let Some(active_turn) = &self.pending.active_turn {
1342 tx.execute(
1343 "UPDATE materialized_sessions SET active_turn_json = ?2 WHERE session_id = ?1",
1344 params![
1345 session_id,
1346 active_turn
1347 .as_ref()
1348 .map(serde_json::to_string)
1349 .transpose()?
1350 ],
1351 )?;
1352 }
1353 if self.pending.clear_turn_outcome {
1354 self.transaction.execute("UPDATE materialized_sessions SET last_turn_outcome_json = NULL WHERE session_id = ?1", [session_id])?;
1355 }
1356 if let Some(last_turn_outcome) = &self.pending.last_turn_outcome {
1357 tx.execute(
1358 "UPDATE materialized_sessions
1359 SET last_turn_outcome_json = ?2 WHERE session_id = ?1",
1360 params![session_id, serde_json::to_string(last_turn_outcome)?],
1361 )?;
1362 }
1363 tx.execute(
1364 "UPDATE materialized_sessions
1365 SET last_activity_at_ms = CASE
1366 WHEN ?2 IS NULL THEN last_activity_at_ms
1367 WHEN last_activity_at_ms IS NULL OR last_activity_at_ms < ?2 THEN ?2
1368 ELSE last_activity_at_ms
1369 END,
1370 applied_event_ordinal = ?3,
1371 applied_event_digest = ?4
1372 WHERE session_id = ?1",
1373 params![
1374 session_id,
1375 self.pending.last_activity_at_ms,
1376 self.applied_ordinal,
1377 self.applied_digest,
1378 ],
1379 )?;
1380 Ok(())
1381 }
1382}
1383
1384pub fn apply_projection_page<T>(
1389 session_id: &str,
1390 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T> + Send + 'static,
1391) -> Result<T>
1392where
1393 T: Send + 'static,
1394{
1395 let session_id = session_id.to_owned();
1396 submit_database_write("apply_projection_page", move |connection| {
1397 apply_projection_page_with(connection, &session_id, fill)
1398 })
1399}
1400
1401#[cfg(test)]
1402pub(super) fn apply_projection_page_to<T>(
1403 path: &Path,
1404 session_id: &str,
1405 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1406) -> Result<T> {
1407 let mut connection = open(path)?;
1408 apply_projection_page_with(&mut connection, session_id, fill)
1409}
1410
1411pub(super) fn apply_projection_page_with<T>(
1412 connection: &mut Connection,
1413 session_id: &str,
1414 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1415) -> Result<T> {
1416 let transaction =
1417 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1418 let (applied_ordinal, applied_digest) = transaction
1419 .query_row(
1420 "SELECT applied_event_ordinal, applied_event_digest
1421 FROM materialized_sessions WHERE session_id = ?1",
1422 [session_id],
1423 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1424 )
1425 .optional()?
1426 .with_context(|| format!("unknown session {session_id}"))?;
1427 validate_relay_event_frontier(
1428 applied_ordinal,
1429 &applied_digest,
1430 "persisted relay event frontier",
1431 )?;
1432 let mut page = ProjectionPage {
1433 session_id,
1434 transaction,
1435 applied_ordinal,
1436 applied_digest,
1437 dirty: false,
1438 pending: MaterializedSessionMutation::default(),
1439 pending_transcript: BTreeMap::new(),
1440 pending_turns: Vec::new(),
1441 pending_events: Vec::new(),
1442 };
1443 let filled = fill(&mut page)?;
1446 page.flush()?;
1447 page.transaction.commit()?;
1448 Ok(filled)
1449}
1450
1451pub fn apply_projection_event(
1453 session_id: &str,
1454 event_ordinal: u64,
1455 previous_event_digest: &str,
1456 event_digest: &str,
1457 mutation: &MaterializedSessionMutation,
1458) -> Result<ProjectionApplyOutcome> {
1459 let session_id = session_id.to_owned();
1460 let previous_event_digest = previous_event_digest.to_owned();
1461 let event_digest = event_digest.to_owned();
1462 let mutation = mutation.clone();
1463 submit_database_write("apply_projection_event", move |connection| {
1464 apply_projection_page_with(connection, &session_id, |page| {
1465 page.apply(
1466 event_ordinal,
1467 &previous_event_digest,
1468 &event_digest,
1469 &mutation,
1470 )
1471 })
1472 })
1473}
1474
1475#[cfg(test)]
1476pub(super) fn apply_projection_event_to(
1477 path: &Path,
1478 session_id: &str,
1479 event_ordinal: u64,
1480 previous_event_digest: &str,
1481 event_digest: &str,
1482 mutation: &MaterializedSessionMutation,
1483) -> Result<ProjectionApplyOutcome> {
1484 apply_projection_page_to(path, session_id, |page| {
1485 page.apply(event_ordinal, previous_event_digest, event_digest, mutation)
1486 })
1487}
1488
1489pub fn advance_viewed_through_event_ordinal(session_id: &str, through: u64) -> Result<u64> {
1492 let session_id = session_id.to_owned();
1493 submit_database_write("advance_viewed_through_event_ordinal", move |_| {
1494 advance_viewed_through_event_ordinal_to(&database_path(), &session_id, through)
1495 })
1496}
1497
1498pub(super) fn advance_viewed_through_event_ordinal_to(
1499 path: &Path,
1500 session_id: &str,
1501 through: u64,
1502) -> Result<u64> {
1503 let mut connection = open(path)?;
1504 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1505 let applied = tx
1506 .query_row(
1507 "SELECT applied_event_ordinal FROM materialized_sessions WHERE session_id = ?1",
1508 [session_id],
1509 |row| row.get::<_, u64>(0),
1510 )
1511 .optional()?
1512 .with_context(|| format!("unknown session {session_id}"))?;
1513 if through > applied {
1514 bail!(
1515 "cannot acknowledge event ordinal {through} for session {session_id}; projection is at {applied}"
1516 );
1517 }
1518 tx.execute(
1519 "UPDATE sessions
1520 SET viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?2)
1521 WHERE session_id = ?1",
1522 params![session_id, through],
1523 )?;
1524 let receipt = tx.query_row(
1525 "SELECT viewed_through_event_ordinal FROM sessions WHERE session_id = ?1",
1526 [session_id],
1527 |row| row.get::<_, u64>(0),
1528 )?;
1529 tx.commit()?;
1530 Ok(receipt)
1531}
1532
1533fn decode_transcript_body(body_json: &str, session_id: &str) -> Result<TranscriptBody> {
1539 let mut body: TranscriptBody = serde_json::from_str(body_json)
1540 .with_context(|| format!("parse materialized transcript body for session {session_id}"))?;
1541 match &mut body {
1542 TranscriptBody::Agent { chunks, .. } | TranscriptBody::Thought { chunks, .. } => {
1543 mj_core::transcript::coalesce_content_chunks(chunks);
1544 }
1545 _ => {}
1546 }
1547 Ok(body)
1548}
1549
1550pub(crate) fn load_continuation_evidence(
1553 session_id: &str,
1554 ordinal: u64,
1555 digest: &str,
1556) -> Result<mj_core::continuation::ContinuationEvidence> {
1557 load_continuation_evidence_from(&database_path(), session_id, ordinal, digest)
1558}
1559
1560pub(super) fn load_continuation_evidence_from(
1561 path: &Path,
1562 session_id: &str,
1563 ordinal: u64,
1564 digest: &str,
1565) -> Result<mj_core::continuation::ContinuationEvidence> {
1566 let mut reader = open_reader(path)?;
1567 let connection = reader.transaction()?;
1568 let fields = read_materialized_session_fields(&connection, session_id)?
1569 .context("continuation projection is missing")?;
1570 anyhow::ensure!(
1571 fields.applied_event_ordinal == ordinal && fields.applied_event_digest == digest,
1572 "continuation projection changed"
1573 );
1574 let mut statement = connection.prepare(
1575 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1576 last_changed_at_ms, body_json
1577 FROM materialized_transcript_items WHERE session_id = ?1
1578 AND (json_extract(body_json, '$.kind') IN ('user', 'agent')
1579 OR stable_id LIKE 'context-cleared:%')
1580 ORDER BY position DESC, stable_id DESC",
1581 )?;
1582 let rows = statement.query_map([session_id], |row| {
1583 Ok((
1584 row.get::<_, String>(0)?,
1585 row.get::<_, u64>(1)?,
1586 row.get::<_, Option<u64>>(2)?,
1587 row.get::<_, i64>(3)?,
1588 row.get::<_, i64>(4)?,
1589 row.get::<_, String>(5)?,
1590 ))
1591 })?;
1592 crate::continuation::evidence_from_items(rows.map(|row| {
1593 let (
1594 stable_id,
1595 position,
1596 latest_content_event_ordinal,
1597 created_at_ms,
1598 last_changed_at_ms,
1599 body_json,
1600 ) = row?;
1601 anyhow::ensure!(
1603 body_json.len() <= 1024 * 1024,
1604 "continuation source message exceeds budget"
1605 );
1606 Ok(Arc::new(TranscriptItem {
1607 stable_id,
1608 position,
1609 latest_content_event_ordinal,
1610 created_at_ms,
1611 last_changed_at_ms,
1612 body: decode_transcript_body(&body_json, session_id)?,
1613 }))
1614 }))
1615}