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 fn load_prompt_acceptance(session_id: &str, command_id: &str) -> Result<Option<u64>> {
331 let mut reader = open_reader(&database_path())?;
332 let connection = reader.transaction()?;
333 if let Some(fields) = read_materialized_session_fields(&connection, session_id)? {
334 if let Some(turn) = fields
335 .active_turn
336 .filter(|turn| turn.command_id == command_id)
337 {
338 return turn
339 .accepted_ordinal
340 .map(Some)
341 .context("accepted prompt has no acceptance ordinal");
342 }
343 if let Some(turn) = fields
344 .last_turn_outcome
345 .filter(|turn| turn.command_id == command_id)
346 {
347 return turn
348 .accepted_ordinal
349 .map(Some)
350 .context("completed prompt has no acceptance ordinal");
351 }
352 }
353 let ordinal: Option<u64> = connection.query_row(
354 "SELECT accepted_ordinal FROM materialized_queued_prompts WHERE session_id=?1 AND command_id=?2
355 UNION ALL SELECT json_extract(body, '$.accepted_ordinal') FROM session_turn_usage
356 WHERE session_id=?1 AND command_id=?2 LIMIT 1",
357 params![session_id, command_id], |row| row.get(0),
358 ).optional()?.flatten();
359 if ordinal.is_some() {
360 return Ok(ordinal);
361 }
362 let started: bool = connection.query_row(
363 "SELECT EXISTS(SELECT 1 FROM materialized_transcript_items WHERE session_id=?1 AND stable_id IN (?2, ?3))",
364 params![session_id, format!("user:{command_id}"), format!("{}{command_id}", mj_core::archive::CONTEXT_BOUNDARY_PREFIX)], |row| row.get(0),
365 )?;
366 ensure!(
367 !started,
368 "input was delivered but its acceptance ordinal is unavailable; refusing to send it twice"
369 );
370 Ok(None)
371}
372
373pub(super) fn load_materialized_turn_outcome_from(
374 path: &Path,
375 session_id: &str,
376) -> Result<Option<MaterializedTurnState>> {
377 let connection = open_reader(path)?;
378 read_materialized_turn_state(&connection, session_id)
379}
380
381pub(super) fn read_materialized_turn_state(
383 connection: &Connection,
384 session_id: &str,
385) -> Result<Option<MaterializedTurnState>> {
386 let row = connection.query_row(
387 "SELECT execution_state, running_started_at_ms, active_turn_json, last_turn_outcome_json
388 FROM materialized_sessions WHERE session_id=?1",
389 [session_id],
390 |row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?,
391 row.get::<_, Option<String>>(2)?, row.get::<_, Option<String>>(3)?)),
392 ).optional()?;
393 let Some((execution, started, active, last)) = row else {
394 return Ok(None);
395 };
396 Ok(Some((
397 parse_materialized_execution(&execution, started)?,
398 active
399 .as_deref()
400 .map(serde_json::from_str)
401 .transpose()
402 .with_context(|| format!("parse active turn for session {session_id}"))?,
403 last.as_deref()
404 .map(serde_json::from_str)
405 .transpose()
406 .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
407 )))
408}
409
410pub fn load_materialized_finished_turn_message(session_id: &str) -> Result<Option<String>> {
421 load_materialized_finished_turn_message_from(&database_path(), session_id)
422}
423
424pub(crate) struct ChildAnswerMessages {
427 pub latest: Option<String>,
428 pub finished: Option<String>,
429}
430
431pub(crate) fn load_child_answer_messages(
432 children: &[(String, Option<(u64, u64)>)],
433) -> Result<BTreeMap<String, ChildAnswerMessages>> {
434 let mut reader = open_reader(&database_path())?;
435 let connection = reader.transaction()?;
436 children
437 .iter()
438 .map(|(id, span)| {
439 let latest = last_materialized_agent_message(&connection, id)?.map(|(_, text)| text);
440 let finished = span
441 .map(|(start, end)| {
442 last_materialized_agent_message_within(&connection, id, start, end)
443 })
444 .transpose()?
445 .flatten();
446 Ok((id.clone(), ChildAnswerMessages { latest, finished }))
447 })
448 .collect()
449}
450
451pub(super) fn load_materialized_finished_turn_message_from(
452 path: &Path,
453 session_id: &str,
454) -> Result<Option<String>> {
455 let mut reader = open_reader(path)?;
456 let connection = reader.transaction()?;
457 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
458 return Ok(None);
459 };
460 let Some(turn) = fields.last_turn_outcome else {
461 return Ok(None);
462 };
463 let Some(start_position) = turn.turn_start_position else {
464 return Ok(None);
465 };
466 last_materialized_agent_message_within(
467 &connection,
468 session_id,
469 start_position,
470 turn.completed_ordinal,
471 )
472}
473
474pub fn load_materialized_turn_summary(
482 session_id: &str,
483 turn_start_position: u64,
484 turn_completed_position: u64,
485) -> Result<TurnSummary> {
486 load_materialized_turn_summary_from(
487 &database_path(),
488 session_id,
489 turn_start_position,
490 turn_completed_position,
491 )
492}
493
494pub(super) fn load_materialized_turn_summary_from(
495 path: &Path,
496 session_id: &str,
497 turn_start_position: u64,
498 turn_completed_position: u64,
499) -> Result<TurnSummary> {
500 let mut reader = open_reader(path)?;
501 let connection = reader.transaction()?;
502 let turn_number = connection.query_row(
503 "SELECT COUNT(*)
504 FROM materialized_transcript_items
505 WHERE session_id = ?1
506 AND position <= ?3
507 AND (
508 stable_id GLOB ?2
509 OR json_extract(
510 CASE
511 WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
512 THEN body_json
513 ELSE '{}'
514 END,
515 '$.kind'
516 ) = 'user'
517 )",
518 params![
519 session_id,
520 format!("{}*", mj_core::transcript::HARNESS_TURN_ITEM_PREFIX),
521 turn_start_position
522 ],
523 |row| row.get::<_, u64>(0),
524 )?;
525 let (turn_started_at_ms, last_changed_at_ms) = connection.query_row(
526 "SELECT COALESCE(MIN(created_at_ms), 0), COALESCE(MAX(last_changed_at_ms), 0)
527 FROM materialized_transcript_items
528 WHERE session_id = ?1 AND position >= ?2 AND position <= ?3",
529 params![session_id, turn_start_position, turn_completed_position],
530 |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
531 )?;
532 let final_message = last_materialized_agent_message_within(
533 &connection,
534 session_id,
535 turn_start_position,
536 turn_completed_position,
537 )?;
538 let tool_calls = connection.query_row(
539 "SELECT COUNT(*)
540 FROM materialized_transcript_items
541 WHERE session_id = ?1 AND position >= ?2 AND position <= ?3
542 AND json_extract(body_json, '$.kind') = 'tool'",
543 params![session_id, turn_start_position, turn_completed_position],
544 |row| row.get::<_, u64>(0),
545 )?;
546 Ok(TurnSummary {
547 turn_number,
548 turn_started_at_ms,
549 last_changed_at_ms,
550 final_message,
551 tool_calls,
552 })
553}
554
555pub fn load_materialized_transcript_filtered(
556 session_id: &str,
557 after_seq: u64,
558 limit: usize,
559 roles: Vec<mj_core::transcript::TranscriptRole>,
560 finished_only: bool,
561) -> Result<Option<TranscriptPage>> {
562 load_materialized_transcript_filtered_from(
563 &database_path(),
564 session_id,
565 after_seq,
566 limit,
567 roles,
568 finished_only,
569 )
570}
571
572pub(super) fn load_materialized_transcript_filtered_from(
573 path: &Path,
574 session_id: &str,
575 after_seq: u64,
576 limit: usize,
577 roles: Vec<mj_core::transcript::TranscriptRole>,
578 finished_only: bool,
579) -> Result<Option<TranscriptPage>> {
580 let mut reader = open_reader(path)?;
581 let connection = reader.transaction()?;
582 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
583 return Ok(None);
584 };
585 let role_kinds: Vec<_> = if finished_only {
586 vec!["agent"]
587 } else {
588 roles.iter().map(|role| role.storage_kind()).collect()
589 };
590 let role_kinds_json = serde_json::to_string(&role_kinds)?;
591 let open_agent_seq = if finished_only {
592 connection.query_row(
593 "SELECT MIN(COALESCE(latest_content_event_ordinal, position))
594 FROM materialized_transcript_items
595 WHERE session_id = ?1
596 AND COALESCE(latest_content_event_ordinal, position) > ?2
597 AND json_extract(body_json, '$.kind') = 'agent'
598 AND json_extract(body_json, '$.streaming') = 1",
599 params![session_id, after_seq],
600 |row| row.get::<_, Option<u64>>(0),
601 )?
602 } else {
603 None
604 };
605 let mut statement = connection.prepare(
606 "WITH matches AS (
607 SELECT *, COALESCE(latest_content_event_ordinal, position) AS seq
608 FROM materialized_transcript_items WHERE session_id = ?1
609 AND COALESCE(latest_content_event_ordinal, position) > ?2
610 AND (json_array_length(?4) = 0 OR json_extract(body_json, '$.kind')
611 IN (SELECT value FROM json_each(?4)))
612 AND (?5 = 0 OR (
613 json_extract(body_json, '$.kind') = 'agent'
614 AND json_extract(body_json, '$.streaming') = 0
615 AND (?6 IS NULL OR COALESCE(latest_content_event_ordinal, position) < ?6)
616 ))
617 ), boundary AS (SELECT MAX(seq) AS seq FROM (SELECT seq FROM matches ORDER BY seq LIMIT ?3))
618 SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
619 last_changed_at_ms, body_json
620 FROM matches WHERE seq <= (SELECT seq FROM boundary)
621 ORDER BY seq, stable_id",
622 )?;
623 let rows = statement
624 .query_map(
625 params![
626 session_id,
627 after_seq,
628 limit.clamp(1, 1000) as i64,
629 role_kinds_json,
630 i64::from(finished_only),
631 open_agent_seq,
632 ],
633 |row| {
634 Ok((
635 row.get::<_, String>(0)?,
636 row.get::<_, u64>(1)?,
637 row.get::<_, Option<u64>>(2)?,
638 row.get::<_, i64>(3)?,
639 row.get::<_, i64>(4)?,
640 row.get::<_, String>(5)?,
641 ))
642 },
643 )?
644 .collect::<rusqlite::Result<Vec<_>>>()?;
645 let items = rows
646 .into_iter()
647 .map(
648 |(
649 stable_id,
650 position,
651 latest_content_event_ordinal,
652 created_at_ms,
653 last_changed_at_ms,
654 body_json,
655 )| {
656 Ok(Arc::new(TranscriptItem {
657 stable_id,
658 position,
659 latest_content_event_ordinal,
660 created_at_ms,
661 last_changed_at_ms,
662 body: decode_transcript_body(&body_json, session_id)?,
663 }))
664 },
665 )
666 .collect::<Result<Vec<_>>>()?;
667 let latest_seq = connection.query_row(
668 "SELECT COALESCE(MAX(COALESCE(latest_content_event_ordinal, position)), 0)
669 FROM materialized_transcript_items
670 WHERE session_id = ?1",
671 [session_id],
672 |row| row.get::<_, u64>(0),
673 )?;
674 let last_seq = items.last().map_or(after_seq, |item| item.seq());
675 let more: bool = connection.query_row(
676 "SELECT EXISTS(
677 SELECT 1 FROM materialized_transcript_items
678 WHERE session_id = ?1
679 AND COALESCE(latest_content_event_ordinal, position) > ?2
680 AND (json_array_length(?3) = 0 OR json_extract(body_json, '$.kind')
681 IN (SELECT value FROM json_each(?3)))
682 AND (?4 = 0 OR (
683 json_extract(body_json, '$.kind') = 'agent'
684 AND json_extract(body_json, '$.streaming') = 0
685 AND (?5 IS NULL OR COALESCE(latest_content_event_ordinal, position) < ?5)
686 ))
687 )",
688 params![
689 session_id,
690 last_seq,
691 role_kinds_json,
692 i64::from(finished_only),
693 open_agent_seq
694 ],
695 |row| row.get(0),
696 )?;
697 let next_after_seq = if more {
698 last_seq
699 } else if finished_only {
700 open_agent_seq.map_or_else(
701 || latest_seq.max(after_seq),
702 |open_seq| open_seq.saturating_sub(1).max(after_seq),
703 )
704 } else {
705 latest_seq.max(after_seq)
706 };
707 Ok(Some(TranscriptPage {
708 next_after_seq,
709 items,
710 latest_seq,
711 execution: fields.execution,
712 }))
713}
714
715pub(super) const RETENTION_BATCH_ITEMS: usize = 4_096;
722
723pub(super) const RETENTION_BODY_FLOOR_BYTES: usize = 4 * 1024;
726
727pub fn compact_materialized_transcript_through(
740 session_id: &str,
741 event_frontier: u64,
742) -> Result<TranscriptRetention> {
743 let session_id = session_id.to_owned();
744 submit_database_write("compact_materialized_transcript", move |_| {
745 compact_materialized_transcript_in(&database_path(), &session_id, event_frontier)
746 })
747}
748
749pub(super) fn compact_materialized_transcript_in(
750 path: &Path,
751 session_id: &str,
752 event_frontier: u64,
753) -> Result<TranscriptRetention> {
754 let mut connection = open(path)?;
755 let candidates = {
756 let mut statement = connection.prepare(
757 "SELECT stable_id, body_json
758 FROM materialized_transcript_items
759 WHERE session_id = ?1
760 AND position <= ?2
761 AND length(body_json) > ?3
762 AND json_extract(
763 CASE WHEN json_valid(body_json) THEN body_json ELSE '{}' END,
764 '$.kind'
765 ) = 'tool'
766 ORDER BY position, stable_id
767 LIMIT ?4",
768 )?;
769 statement
770 .query_map(
771 params![
772 session_id,
773 event_frontier,
774 RETENTION_BODY_FLOOR_BYTES as i64,
775 RETENTION_BATCH_ITEMS as i64 + 1
776 ],
777 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
778 )?
779 .collect::<rusqlite::Result<Vec<_>>>()?
780 };
781 let remaining = candidates.len() > RETENTION_BATCH_ITEMS;
782 let mut retention = TranscriptRetention {
783 remaining,
784 ..TranscriptRetention::default()
785 };
786 let transaction =
787 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
788 for (stable_id, body_json) in candidates.into_iter().take(RETENTION_BATCH_ITEMS) {
789 let mut body: TranscriptBody = match serde_json::from_str(&body_json) {
790 Ok(body) => body,
791 Err(error) => {
793 tracing::warn!(%session_id, %stable_id, %error, "skipping unreadable transcript body");
794 continue;
795 }
796 };
797 if !mj_transcript::transcript::compact_tool_call_for_retention(&mut body) {
798 continue;
799 }
800 let compacted = serde_json::to_string(&body)
801 .with_context(|| format!("serialize compacted transcript body {stable_id}"))?;
802 if compacted.len() >= body_json.len() {
803 continue;
804 }
805 transaction.execute(
806 "UPDATE materialized_transcript_items SET body_json = ?3
807 WHERE session_id = ?1 AND stable_id = ?2",
808 params![session_id, stable_id, compacted],
809 )?;
810 retention.items += 1;
811 retention.bytes += body_json.len() - compacted.len();
812 }
813 transaction.commit()?;
814 Ok(retention)
815}
816
817pub const PROJECTION_TAIL_ITEMS: usize = 1_024;
825
826pub fn load_materialized_actor_projection(
829 session_id: &str,
830) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
831 load_materialized_actor_projection_from(&database_path(), session_id)
832}
833
834pub(super) fn load_materialized_actor_projection_from(
835 path: &Path,
836 session_id: &str,
837) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
838 let mut reader = open_reader(path)?;
839 let connection = reader.transaction()?;
840 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
841 return Ok(None);
842 };
843 let latest_turn = last_materialized_turn_start(&connection, session_id)?;
844 let desired: Option<u64> = connection
845 .query_row(
846 "SELECT position FROM materialized_transcript_items WHERE session_id=?1
847 ORDER BY position DESC, stable_id DESC LIMIT 1 OFFSET ?2",
848 params![session_id, PROJECTION_TAIL_ITEMS - 1],
849 |row| row.get(0),
850 )
851 .optional()?;
852 let mutable: Option<u64> = connection.query_row(
855 "SELECT MIN(position) FROM materialized_transcript_items WHERE session_id=?1
856 AND (json_extract(body_json, '$.streaming')=1
857 OR (json_extract(body_json, '$.kind')='tool'
858 AND json_extract(body_json, '$.call.status') IN ('pending','in_progress')))",
859 [session_id],
860 |row| row.get(0),
861 )?;
862 let boundary = desired
863 .unwrap_or(0)
864 .min(latest_turn.unwrap_or(0))
865 .min(mutable.unwrap_or(u64::MAX));
866 let start: u64 = connection.query_row(
867 "SELECT COALESCE(MAX(position),0) FROM materialized_transcript_items
868 WHERE session_id=?1 AND position<=?2
869 AND (json_extract(body_json, '$.kind')='user' OR stable_id LIKE 'harness-turn:%')",
870 params![session_id, boundary],
871 |row| row.get(0),
872 )?;
873 let transcript = read_transcript_range(&connection, session_id, start, None, None)?;
874 let omitted_items = connection.query_row(
875 "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id=?1 AND position<?2",
876 params![session_id, start],
877 |row| row.get(0),
878 )?;
879 let window = ProjectionWindow {
880 omitted_items,
881 provisional_title: first_materialized_user_message(&connection, session_id)?
882 .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
883 latest_turn_start_position: latest_turn,
884 };
885 let materialized = MaterializedSession {
886 session_id: session_id.to_owned(),
887 applied_event_ordinal: fields.applied_event_ordinal,
888 applied_event_digest: fields.applied_event_digest,
889 last_activity_at_ms: fields.last_activity_at_ms,
890 execution: fields.execution,
891 session_title: fields.session_title,
892 configuration: fields.configuration,
893 transcript,
894 queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
895 pending_elicitations: fields.pending_elicitations,
896 active_turn: fields.active_turn,
897 last_turn_outcome: fields.last_turn_outcome,
898 };
899 materialized.validate()?;
900 Ok(Some((materialized, window)))
901}
902
903pub fn load_transcript_history(
904 session_id: &str,
905 before: Option<&mj_core::storage::TranscriptCursor>,
906 limit: usize,
907) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
908 load_transcript_history_from(&database_path(), session_id, before, limit)
909}
910
911pub(super) fn load_transcript_history_from(
912 path: &Path,
913 session_id: &str,
914 before: Option<&mj_core::storage::TranscriptCursor>,
915 limit: usize,
916) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
917 let mut reader = open_reader(path)?;
918 let connection = reader.transaction()?;
919 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
920 return Ok(None);
921 };
922 let limit = limit.clamp(1, 256);
923 let mut items = read_transcript_range(&connection, session_id, 0, before, Some(limit + 1))?;
924 let has_more = items.len() > limit;
925 if has_more {
926 items.remove(0);
927 }
928 let before = has_more.then(|| mj_core::storage::TranscriptCursor::of(&items[0]));
929 Ok(Some(mj_core::storage::TranscriptHistoryPage {
930 items,
931 before,
932 frontier: fields.applied_event_ordinal,
933 }))
934}
935
936fn read_transcript_range(
937 connection: &Connection,
938 session_id: &str,
939 start: u64,
940 before: Option<&mj_core::storage::TranscriptCursor>,
941 limit: Option<usize>,
942) -> Result<Vec<Arc<TranscriptItem>>> {
943 let mut statement = connection.prepare(
944 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
945 last_changed_at_ms, body_json FROM materialized_transcript_items
946 WHERE session_id=?1 AND position>=?2
947 AND (position,stable_id)<(?3,?4)
948 ORDER BY position DESC, stable_id DESC LIMIT ?5",
949 )?;
950 let rows = statement
951 .query_map(
952 params![
953 session_id,
954 start,
955 before.map_or(i64::MAX as u64, |c| c.position),
956 before.map_or("", |c| c.stable_id.as_str()),
957 limit.map_or(-1, |limit| limit as i64)
958 ],
959 |row| {
960 Ok((
961 row.get::<_, String>(0)?,
962 row.get::<_, u64>(1)?,
963 row.get::<_, Option<u64>>(2)?,
964 row.get::<_, i64>(3)?,
965 row.get::<_, i64>(4)?,
966 row.get::<_, String>(5)?,
967 ))
968 },
969 )?
970 .collect::<rusqlite::Result<Vec<_>>>()?;
971 let mut items = rows
972 .into_iter()
973 .map(
974 |(
975 stable_id,
976 position,
977 latest_content_event_ordinal,
978 created_at_ms,
979 last_changed_at_ms,
980 body,
981 )| {
982 Ok(Arc::new(TranscriptItem {
983 stable_id,
984 position,
985 latest_content_event_ordinal,
986 created_at_ms,
987 last_changed_at_ms,
988 body: decode_transcript_body(&body, session_id)?,
989 }))
990 },
991 )
992 .collect::<Result<Vec<_>>>()?;
993 items.reverse();
994 Ok(items)
995}
996
997pub fn load_projection_references(
1000 projection: &MaterializedSession,
1001 events: &[mj_core::relay::RelayEvent],
1002) -> Result<Vec<Arc<TranscriptItem>>> {
1003 load_projection_references_from(&database_path(), projection, events)
1004}
1005
1006pub(super) fn load_projection_references_from(
1007 path: &Path,
1008 projection: &MaterializedSession,
1009 events: &[mj_core::relay::RelayEvent],
1010) -> Result<Vec<Arc<TranscriptItem>>> {
1011 let (mut ids, terminals) = mj_transcript::projection::historical_references(events)?;
1012 let retained = projection
1013 .transcript
1014 .iter()
1015 .map(|item| item.stable_id.as_str())
1016 .collect::<std::collections::HashSet<_>>();
1017 ids.retain(|id| !retained.contains(id.as_str()));
1018 if ids.is_empty() && terminals.is_empty() {
1019 return Ok(Vec::new());
1020 }
1021 let mut reader = open_reader(path)?;
1022 let connection = reader.transaction()?;
1023 let mut parameters = vec![projection.session_id.clone(), serde_json::to_string(&ids)?];
1026 let mut statement = connection.prepare(if terminals.is_empty() {
1027 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1028 last_changed_at_ms, body_json FROM materialized_transcript_items
1029 WHERE session_id=?1 AND stable_id IN (SELECT value FROM json_each(?2))
1030 ORDER BY position,stable_id"
1031 } else {
1032 parameters.push(serde_json::to_string(&retained)?);
1033 parameters.push(serde_json::to_string(&terminals)?);
1034 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1035 last_changed_at_ms, body_json FROM materialized_transcript_items
1036 WHERE session_id=?1 AND stable_id NOT IN (SELECT value FROM json_each(?3))
1037 AND (stable_id IN (SELECT value FROM json_each(?2))
1038 OR EXISTS(SELECT 1 FROM json_each(body_json, '$.terminal_refs')
1039 WHERE value IN (SELECT value FROM json_each(?4))))
1040 ORDER BY position,stable_id"
1041 })?;
1042 let rows = statement
1043 .query_map(rusqlite::params_from_iter(parameters), |row| {
1044 Ok((
1045 row.get::<_, String>(0)?,
1046 row.get::<_, u64>(1)?,
1047 row.get::<_, Option<u64>>(2)?,
1048 row.get::<_, i64>(3)?,
1049 row.get::<_, i64>(4)?,
1050 row.get::<_, String>(5)?,
1051 ))
1052 })?
1053 .collect::<rusqlite::Result<Vec<_>>>()?;
1054 rows.into_iter()
1055 .map(
1056 |(
1057 stable_id,
1058 position,
1059 latest_content_event_ordinal,
1060 created_at_ms,
1061 last_changed_at_ms,
1062 body,
1063 )| {
1064 Ok(Arc::new(TranscriptItem {
1065 stable_id,
1066 position,
1067 latest_content_event_ordinal,
1068 created_at_ms,
1069 last_changed_at_ms,
1070 body: decode_transcript_body(&body, &projection.session_id)?,
1071 }))
1072 },
1073 )
1074 .collect()
1075}
1076
1077pub fn load_materialized_projection_tail(
1087 session_id: &str,
1088 transcript_limit: usize,
1089) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
1090 load_materialized_projection_tail_from(&database_path(), session_id, transcript_limit)
1091}
1092
1093pub(super) fn load_materialized_projection_tail_from(
1094 path: &Path,
1095 session_id: &str,
1096 transcript_limit: usize,
1097) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
1098 let mut reader = open_reader(path)?;
1099 let connection = reader.transaction()?;
1102 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
1103 return Ok(None);
1104 };
1105 let transcript = read_materialized_transcript(&connection, session_id, Some(transcript_limit))?;
1106 let total_items = connection.query_row(
1107 "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id = ?1",
1108 [session_id],
1109 |row| row.get::<_, usize>(0),
1110 )?;
1111 let window = ProjectionWindow {
1112 omitted_items: total_items.saturating_sub(transcript.len()),
1113 provisional_title: first_materialized_user_message(&connection, session_id)?
1114 .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
1115 latest_turn_start_position: last_materialized_turn_start(&connection, session_id)?,
1116 };
1117 let materialized = MaterializedSession {
1118 session_id: session_id.to_owned(),
1119 applied_event_ordinal: fields.applied_event_ordinal,
1120 applied_event_digest: fields.applied_event_digest,
1121 last_activity_at_ms: fields.last_activity_at_ms,
1122 execution: fields.execution,
1123 session_title: fields.session_title,
1124 configuration: fields.configuration,
1125 transcript,
1126 queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
1127 pending_elicitations: fields.pending_elicitations,
1128 active_turn: fields.active_turn,
1129 last_turn_outcome: fields.last_turn_outcome,
1130 };
1131 materialized.validate()?;
1132 Ok(Some((materialized, window)))
1133}
1134
1135pub fn materialized_event_frontier(session_id: &str) -> Result<Option<(u64, String)>> {
1139 materialized_event_frontier_from(&database_path(), session_id)
1140}
1141
1142pub(super) fn materialized_event_frontier_from(
1143 path: &Path,
1144 session_id: &str,
1145) -> Result<Option<(u64, String)>> {
1146 Ok(open_reader(path)?
1147 .query_row(
1148 "SELECT applied_event_ordinal, applied_event_digest
1149 FROM materialized_sessions WHERE session_id = ?1",
1150 [session_id],
1151 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1152 )
1153 .optional()?)
1154}
1155
1156pub fn replace_materialized_queued_prompts(
1160 session_id: &str,
1161 queued_prompts: &[MaterializedQueuedPrompt],
1162) -> Result<()> {
1163 let session_id = session_id.to_owned();
1164 let queued_prompts = queued_prompts.to_vec();
1165 submit_database_write("replace_materialized_queued_prompts", move |_| {
1166 replace_materialized_queued_prompts_in(&database_path(), &session_id, &queued_prompts)
1167 })
1168}
1169
1170pub(super) fn replace_materialized_queued_prompts_in(
1171 path: &Path,
1172 session_id: &str,
1173 queued_prompts: &[MaterializedQueuedPrompt],
1174) -> Result<()> {
1175 let mut connection = open(path)?;
1176 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1177 if !session_exists(&tx, session_id)? {
1178 bail!("unknown session {session_id}");
1179 }
1180 replace_materialized_queue(&tx, session_id, queued_prompts)?;
1181 tx.commit()?;
1182 Ok(())
1183}
1184
1185pub fn load_transcribed_session_activity() -> Result<BTreeMap<String, Option<i64>>> {
1192 load_transcribed_session_activity_from(&database_path())
1193}
1194
1195fn load_transcribed_session_activity_from(path: &Path) -> Result<BTreeMap<String, Option<i64>>> {
1196 let connection = open_reader(path)?;
1197 let mut statement = connection.prepare(
1198 "SELECT session_id, last_activity_at_ms
1199 FROM materialized_sessions s
1200 WHERE EXISTS (
1201 SELECT 1 FROM materialized_transcript_items i
1202 WHERE i.session_id = s.session_id
1203 )",
1204 )?;
1205 let rows = statement.query_map([], |row| {
1206 Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?))
1207 })?;
1208 let mut activity = BTreeMap::new();
1209 for row in rows {
1210 let (session_id, last_activity_at_ms) = row?;
1211 activity.insert(session_id, last_activity_at_ms);
1212 }
1213 Ok(activity)
1214}
1215
1216pub fn load_materialized_queued_prompts() -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>>
1220{
1221 load_materialized_queued_prompts_from(&database_path())
1222}
1223
1224pub(super) fn load_materialized_queued_prompts_from(
1225 path: &Path,
1226) -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>> {
1227 let connection = open_reader(path)?;
1228 let mut statement = connection.prepare(
1229 "SELECT session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1230 FROM materialized_queued_prompts
1231 ORDER BY session_id, ordinal",
1232 )?;
1233 let rows = statement.query_map([], |row| {
1234 Ok((
1235 row.get::<_, String>(0)?,
1236 row.get::<_, String>(1)?,
1237 row.get::<_, String>(2)?,
1238 row.get::<_, String>(3)?,
1239 row.get::<_, i64>(4)?,
1240 row.get::<_, Option<u64>>(5)?,
1241 ))
1242 })?;
1243 let mut queues = BTreeMap::<String, Vec<MaterializedQueuedPrompt>>::new();
1244 for row in rows {
1245 let (session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal) =
1246 row?;
1247 let content = serde_json::from_str(&content_json).with_context(|| {
1248 format!("parse materialized queued prompt for session {session_id}")
1249 })?;
1250 let kind = serde_json::from_str(&kind_json).with_context(|| {
1251 format!("parse materialized queue entry kind for session {session_id}")
1252 })?;
1253 queues
1254 .entry(session_id)
1255 .or_default()
1256 .push(MaterializedQueuedPrompt {
1257 command_id,
1258 kind,
1259 content,
1260 queued_at_ms,
1261 accepted_ordinal,
1262 });
1263 }
1264 Ok(queues)
1265}
1266
1267pub(super) fn load_materialized_session_from(
1268 path: &Path,
1269 session_id: &str,
1270) -> Result<Option<MaterializedSession>> {
1271 let mut reader = open_reader(path)?;
1272 let connection = reader.transaction()?;
1273 load_materialized_session_with(&connection, session_id)
1274}
1275
1276pub(super) fn load_materialized_session_with(
1277 connection: &rusqlite::Transaction<'_>,
1278 session_id: &str,
1279) -> Result<Option<MaterializedSession>> {
1280 let Some(fields) = read_materialized_session_fields(connection, session_id)? else {
1281 return Ok(None);
1282 };
1283 let materialized = MaterializedSession {
1284 session_id: session_id.to_owned(),
1285 applied_event_ordinal: fields.applied_event_ordinal,
1286 applied_event_digest: fields.applied_event_digest,
1287 last_activity_at_ms: fields.last_activity_at_ms,
1288 execution: fields.execution,
1289 session_title: fields.session_title,
1290 configuration: fields.configuration,
1291 transcript: read_materialized_transcript(connection, session_id, None)?,
1292 queued_prompts: read_materialized_queued_prompts(connection, session_id)?,
1293 pending_elicitations: fields.pending_elicitations,
1294 active_turn: fields.active_turn,
1295 last_turn_outcome: fields.last_turn_outcome,
1296 };
1297 materialized.validate()?;
1298 Ok(Some(materialized))
1299}
1300
1301pub(super) struct MaterializedSessionFields {
1303 pub(super) applied_event_ordinal: u64,
1304 pub(super) applied_event_digest: String,
1305 pub(super) last_activity_at_ms: Option<i64>,
1306 pub(super) execution: MaterializedExecutionState,
1307 pub(super) session_title: Option<String>,
1308 pub(super) configuration: mj_core::state::SessionConfiguration,
1309 pub(super) pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
1310 pub(super) active_turn: Option<MaterializedTurn>,
1311 pub(super) last_turn_outcome: Option<MaterializedTurnOutcome>,
1312}
1313
1314pub(super) fn read_materialized_session_fields(
1315 connection: &Connection,
1316 session_id: &str,
1317) -> Result<Option<MaterializedSessionFields>> {
1318 let row = connection
1319 .query_row(
1320 "SELECT applied_event_ordinal, applied_event_digest, last_activity_at_ms,
1321 execution_state, running_started_at_ms, session_title, configuration_json,
1322 pending_elicitations_json, active_turn_json, last_turn_outcome_json
1323 FROM materialized_sessions WHERE session_id = ?1",
1324 [session_id],
1325 |row| {
1326 Ok((
1327 row.get::<_, u64>(0)?,
1328 row.get::<_, String>(1)?,
1329 row.get::<_, Option<i64>>(2)?,
1330 row.get::<_, String>(3)?,
1331 row.get::<_, Option<i64>>(4)?,
1332 row.get::<_, Option<String>>(5)?,
1333 row.get::<_, String>(6)?,
1334 row.get::<_, String>(7)?,
1335 row.get::<_, Option<String>>(8)?,
1336 row.get::<_, Option<String>>(9)?,
1337 ))
1338 },
1339 )
1340 .optional()?;
1341 let Some((
1342 applied_event_ordinal,
1343 applied_event_digest,
1344 last_activity_at_ms,
1345 execution,
1346 running_started_at_ms,
1347 session_title,
1348 configuration_json,
1349 pending_elicitations_json,
1350 active_turn_json,
1351 last_turn_outcome_json,
1352 )) = row
1353 else {
1354 return Ok(None);
1355 };
1356 #[cfg(test)]
1357 super::tests::after_materialized_frontier_read();
1358 Ok(Some(MaterializedSessionFields {
1359 applied_event_ordinal,
1360 applied_event_digest,
1361 last_activity_at_ms,
1362 execution: parse_materialized_execution(&execution, running_started_at_ms)?,
1363 session_title,
1364 configuration: serde_json::from_str(&configuration_json).with_context(|| {
1365 format!("parse materialized configuration for session {session_id}")
1366 })?,
1367 pending_elicitations: serde_json::from_str(&pending_elicitations_json)
1368 .with_context(|| format!("parse pending elicitations for session {session_id}"))?,
1369 active_turn: active_turn_json
1370 .as_deref()
1371 .map(serde_json::from_str)
1372 .transpose()
1373 .with_context(|| format!("parse active turn for session {session_id}"))?,
1374 last_turn_outcome: last_turn_outcome_json
1375 .as_deref()
1376 .map(serde_json::from_str)
1377 .transpose()
1378 .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
1379 }))
1380}
1381
1382pub(super) fn read_materialized_transcript(
1386 connection: &Connection,
1387 session_id: &str,
1388 limit: Option<usize>,
1389) -> Result<Vec<Arc<TranscriptItem>>> {
1390 let mut statement = connection.prepare(match limit {
1391 Some(_) => {
1392 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1393 last_changed_at_ms, body_json
1394 FROM materialized_transcript_items
1395 WHERE session_id = ?1
1396 ORDER BY position DESC, stable_id DESC
1397 LIMIT ?2"
1398 }
1399 None => {
1400 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1401 last_changed_at_ms, body_json
1402 FROM materialized_transcript_items
1403 WHERE session_id = ?1
1404 ORDER BY position, stable_id"
1405 }
1406 })?;
1407 let read = |row: &rusqlite::Row<'_>| {
1408 Ok((
1409 row.get::<_, String>(0)?,
1410 row.get::<_, u64>(1)?,
1411 row.get::<_, Option<u64>>(2)?,
1412 row.get::<_, i64>(3)?,
1413 row.get::<_, i64>(4)?,
1414 row.get::<_, String>(5)?,
1415 ))
1416 };
1417 let rows = match limit {
1418 Some(limit) => statement
1419 .query_map(params![session_id, limit as i64], read)?
1420 .collect::<rusqlite::Result<Vec<_>>>()?,
1421 None => statement
1422 .query_map([session_id], read)?
1423 .collect::<rusqlite::Result<Vec<_>>>()?,
1424 };
1425 let mut transcript = rows
1426 .into_iter()
1427 .map(
1428 |(
1429 stable_id,
1430 position,
1431 latest_content_event_ordinal,
1432 created_at_ms,
1433 last_changed_at_ms,
1434 body_json,
1435 )| {
1436 Ok(Arc::new(TranscriptItem {
1437 stable_id,
1438 position,
1439 latest_content_event_ordinal,
1440 created_at_ms,
1441 last_changed_at_ms,
1442 body: decode_transcript_body(&body_json, session_id)?,
1443 }))
1444 },
1445 )
1446 .collect::<Result<Vec<_>>>()?;
1447 if limit.is_some() {
1448 transcript.reverse();
1451 }
1452 Ok(transcript)
1453}
1454
1455pub(super) fn read_materialized_queued_prompts(
1456 connection: &Connection,
1457 session_id: &str,
1458) -> Result<Vec<MaterializedQueuedPrompt>> {
1459 let mut statement = connection.prepare(
1460 "SELECT command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1461 FROM materialized_queued_prompts
1462 WHERE session_id = ?1
1463 ORDER BY ordinal",
1464 )?;
1465 let rows = statement
1466 .query_map([session_id], |row| {
1467 Ok((
1468 row.get::<_, String>(0)?,
1469 row.get::<_, String>(1)?,
1470 row.get::<_, String>(2)?,
1471 row.get::<_, i64>(3)?,
1472 row.get::<_, Option<u64>>(4)?,
1473 ))
1474 })?
1475 .collect::<rusqlite::Result<Vec<_>>>()?;
1476 rows.into_iter()
1477 .map(
1478 |(command_id, kind_json, content_json, queued_at_ms, accepted_ordinal)| {
1479 Ok(MaterializedQueuedPrompt {
1480 command_id,
1481 kind: serde_json::from_str(&kind_json).with_context(|| {
1482 format!("parse materialized queue entry kind for session {session_id}")
1483 })?,
1484 content: serde_json::from_str(&content_json).with_context(|| {
1485 format!("parse materialized queued prompt for session {session_id}")
1486 })?,
1487 queued_at_ms,
1488 accepted_ordinal,
1489 })
1490 },
1491 )
1492 .collect()
1493}
1494
1495pub fn save_materialized_session(materialized: &MaterializedSession) -> Result<()> {
1499 let materialized = materialized.clone();
1500 submit_database_write("save_materialized_session", move |_| {
1501 save_materialized_session_to(&database_path(), &materialized)
1502 })
1503}
1504
1505pub(super) fn save_materialized_session_to(
1506 path: &Path,
1507 materialized: &MaterializedSession,
1508) -> Result<()> {
1509 materialized.validate()?;
1510 let mut connection = open(path)?;
1511 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1512 if !session_exists(&tx, &materialized.session_id)? {
1513 bail!("unknown session {}", materialized.session_id);
1514 }
1515 write_materialized_session(&tx, materialized)?;
1516 tx.commit()?;
1517 Ok(())
1518}
1519
1520pub struct ProjectionPage<'a> {
1525 pub(super) session_id: &'a str,
1526 pub(super) transaction: Transaction<'a>,
1527 pub(super) applied_ordinal: u64,
1528 pub(super) applied_digest: String,
1529 pub(super) dirty: bool,
1530 pub(super) pending: MaterializedSessionMutation,
1531 pub(super) pending_transcript: BTreeMap<String, PendingTranscriptMutation>,
1532 pub(super) pending_turns: Vec<MaterializedTurnOutcome>,
1533 usage_configuration: mj_core::state::SessionConfiguration,
1536 pub(super) pending_events: Vec<(i64, ApiEventData)>,
1537}
1538
1539pub(super) struct PendingTranscriptMutation {
1540 pub(super) final_mutation: TranscriptMutation,
1541 pub(super) remove_before_upsert: bool,
1542}
1543
1544impl ProjectionPage<'_> {
1545 pub fn apply(
1549 &mut self,
1550 event_ordinal: u64,
1551 previous_event_digest: &str,
1552 event_digest: &str,
1553 mutation: &MaterializedSessionMutation,
1554 ) -> Result<ProjectionApplyOutcome> {
1555 if event_ordinal == 0 {
1556 bail!("relay event ordinal must be positive");
1557 }
1558 let chained = !previous_event_digest.is_empty();
1564 if chained {
1565 validate_relay_event_digest(previous_event_digest, "previous relay event digest")?;
1566 }
1567 validate_relay_event_frontier(event_ordinal, event_digest, "relay event frontier")?;
1568 let session_id = self.session_id;
1569 let applied = self.applied_ordinal;
1570 if event_ordinal < applied {
1571 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1572 }
1573 if event_ordinal == applied {
1574 if event_digest != self.applied_digest {
1575 bail!(
1576 "relay event digest mismatch for session {session_id} at ordinal {event_ordinal}: projection has {}, received {event_digest}",
1577 self.applied_digest
1578 );
1579 }
1580 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1581 }
1582 let expected = applied
1583 .checked_add(1)
1584 .context("materialized event ordinal overflow")?;
1585 if event_ordinal != expected {
1586 bail!(
1587 "relay event gap for session {session_id}: expected ordinal {expected}, received {event_ordinal}"
1588 );
1589 }
1590 if chained && previous_event_digest != self.applied_digest {
1591 bail!(
1592 "relay event chain diverged for session {session_id} before ordinal {event_ordinal}: projection has {}, event follows {previous_event_digest}",
1593 self.applied_digest
1594 );
1595 }
1596
1597 if let Some(event) = &mutation.native_agent {
1598 native_agents::apply_native_agent_event(&self.transaction, session_id, event)?;
1599 }
1600 if let Some(activity_at_ms) = mutation.last_activity_at_ms {
1601 self.pending.last_activity_at_ms = Some(
1602 self.pending
1603 .last_activity_at_ms
1604 .map_or(activity_at_ms, |existing| existing.max(activity_at_ms)),
1605 );
1606 }
1607 if let Some(execution) = mutation.execution {
1608 self.pending.execution = Some(execution);
1609 }
1610 if let Some(title) = &mutation.session_title {
1611 if title.as_ref().is_some_and(|title| title.trim().is_empty()) {
1612 bail!("materialized session title cannot be empty");
1613 }
1614 self.pending.session_title = Some(title.clone());
1615 }
1616 if let Some(configuration) = &mutation.configuration {
1617 self.usage_configuration = configuration.clone();
1618 self.pending.configuration = Some(configuration.clone());
1619 }
1620 for item_mutation in &mutation.transcript {
1621 match item_mutation {
1622 TranscriptMutation::Upsert(item) => {
1623 item.validate(event_ordinal)?;
1624 let stable_id = item.stable_id.clone();
1625 let entry = self.pending_transcript.entry(stable_id).or_insert_with(|| {
1626 PendingTranscriptMutation {
1627 final_mutation: TranscriptMutation::Upsert(item.clone()),
1628 remove_before_upsert: false,
1629 }
1630 });
1631 entry.remove_before_upsert |=
1632 matches!(&entry.final_mutation, TranscriptMutation::Remove { .. });
1633 entry.final_mutation = TranscriptMutation::Upsert(item.clone());
1634 }
1635 TranscriptMutation::Remove { stable_id } => {
1636 if stable_id.trim().is_empty() {
1637 bail!("cannot remove a transcript item with an empty stable id");
1638 }
1639 let removed = TranscriptMutation::Remove {
1640 stable_id: stable_id.clone(),
1641 };
1642 self.pending_transcript
1643 .entry(stable_id.clone())
1644 .and_modify(|entry| entry.final_mutation = removed.clone())
1645 .or_insert(PendingTranscriptMutation {
1646 final_mutation: removed,
1647 remove_before_upsert: false,
1648 });
1649 }
1650 }
1651 }
1652 if let Some(queued_prompts) = &mutation.queued_prompts {
1653 self.pending.queued_prompts = Some(queued_prompts.clone());
1654 }
1655 if let Some(pending_elicitations) = &mutation.pending_elicitations {
1656 self.pending.pending_elicitations = Some(pending_elicitations.clone());
1657 }
1658 self.pending
1659 .config_results
1660 .extend(mutation.config_results.clone());
1661 if let Some(active_turn) = &mutation.active_turn {
1662 if let Some(turn) = active_turn {
1663 super::usage::record_turn_selection(
1664 &self.transaction,
1665 session_id,
1666 &turn.command_id,
1667 &self.usage_configuration,
1668 )?;
1669 }
1670 self.pending.active_turn = Some(active_turn.clone());
1671 }
1672 if mutation.clear_turn_outcome {
1673 self.pending.clear_turn_outcome = true;
1674 self.pending.last_turn_outcome = None;
1675 }
1676 if let Some(last_turn_outcome) = &mutation.last_turn_outcome {
1677 self.pending_turns.push(last_turn_outcome.clone());
1678 self.pending.last_turn_outcome = Some(last_turn_outcome.clone());
1679 }
1680 if let Some(cost) = &mutation.provider_cost {
1681 self.pending.provider_cost = Some(cost.clone());
1682 }
1683 self.pending_events.extend(
1684 mutation
1685 .api_events
1686 .iter()
1687 .cloned()
1688 .map(|event| (mutation.last_activity_at_ms.unwrap_or(0), event)),
1689 );
1690 self.applied_ordinal = event_ordinal;
1691 event_digest.clone_into(&mut self.applied_digest);
1692 self.dirty = true;
1693 Ok(ProjectionApplyOutcome::Applied)
1694 }
1695
1696 pub(super) fn flush(&mut self) -> Result<()> {
1700 if !self.dirty {
1701 return Ok(());
1702 }
1703 let tx = &self.transaction;
1704 let session_id = self.session_id;
1705 if let Some(execution) = self.pending.execution {
1706 let (state, started_at_ms) = materialized_execution_columns(execution);
1707 tx.execute(
1708 "UPDATE materialized_sessions
1709 SET execution_state = ?2, running_started_at_ms = ?3
1710 WHERE session_id = ?1",
1711 params![session_id, state, started_at_ms],
1712 )?;
1713 }
1714 if let Some(title) = &self.pending.session_title {
1715 tx.execute(
1716 "UPDATE materialized_sessions SET session_title = ?2 WHERE session_id = ?1",
1717 params![session_id, title],
1718 )?;
1719 }
1720 if let Some(configuration) = &self.pending.configuration {
1721 tx.execute(
1722 "UPDATE materialized_sessions SET configuration_json = ?2 WHERE session_id = ?1",
1723 params![session_id, serde_json::to_string(configuration)?],
1724 )?;
1725 }
1726 for pending in self.pending_transcript.values() {
1727 match &pending.final_mutation {
1728 TranscriptMutation::Upsert(item) => {
1729 if pending.remove_before_upsert {
1733 tx.execute(
1734 "DELETE FROM materialized_transcript_items
1735 WHERE session_id = ?1 AND stable_id = ?2",
1736 params![session_id, item.stable_id],
1737 )?;
1738 }
1739 upsert_transcript_item(tx, session_id, item)?;
1740 }
1741 TranscriptMutation::Remove { stable_id } => {
1742 tx.execute(
1743 "DELETE FROM materialized_transcript_items
1744 WHERE session_id = ?1 AND stable_id = ?2",
1745 params![session_id, stable_id],
1746 )?;
1747 }
1748 }
1749 }
1750 if let Some(queued_prompts) = &self.pending.queued_prompts {
1751 replace_materialized_queue(tx, session_id, queued_prompts)?;
1752 }
1753 if let Some(pending_elicitations) = &self.pending.pending_elicitations {
1754 tx.execute(
1755 "UPDATE materialized_sessions
1756 SET pending_elicitations_json = ?2 WHERE session_id = ?1",
1757 params![session_id, serde_json::to_string(pending_elicitations)?],
1758 )?;
1759 }
1760 for (recorded_at_ms, event) in &self.pending_events {
1761 events::insert_api_event(tx, session_id, *recorded_at_ms, event)?;
1762 }
1763 for turn in &self.pending_turns {
1764 tx.execute("INSERT OR REPLACE INTO session_turn_usage(session_id, command_id, completed_ordinal, turn_start_position, body) VALUES (?1, ?2, ?3, ?4, ?5)", params![session_id, turn.command_id, turn.completed_ordinal, turn.turn_start_position, serde_json::to_string(turn)?])?;
1765 }
1766 if let Some(cost) = &self.pending.provider_cost {
1767 tx.execute(
1768 "INSERT OR REPLACE INTO session_provider_cost(session_id, body) VALUES (?1, ?2)",
1769 params![session_id, serde_json::to_string(cost)?],
1770 )?;
1771 }
1772 for (command_id, error) in &self.pending.config_results {
1773 tx.execute("INSERT OR REPLACE INTO api_config_results(session_id, command_id, error) VALUES (?1, ?2, ?3)", params![session_id, command_id, error])?;
1774 }
1775 if let Some(active_turn) = &self.pending.active_turn {
1776 tx.execute(
1777 "UPDATE materialized_sessions SET active_turn_json = ?2 WHERE session_id = ?1",
1778 params![
1779 session_id,
1780 active_turn
1781 .as_ref()
1782 .map(serde_json::to_string)
1783 .transpose()?
1784 ],
1785 )?;
1786 }
1787 if self.pending.clear_turn_outcome {
1788 self.transaction.execute("UPDATE materialized_sessions SET last_turn_outcome_json = NULL WHERE session_id = ?1", [session_id])?;
1789 }
1790 if let Some(last_turn_outcome) = &self.pending.last_turn_outcome {
1791 tx.execute(
1792 "UPDATE materialized_sessions
1793 SET last_turn_outcome_json = ?2 WHERE session_id = ?1",
1794 params![session_id, serde_json::to_string(last_turn_outcome)?],
1795 )?;
1796 }
1797 tx.execute(
1798 "UPDATE materialized_sessions
1799 SET last_activity_at_ms = CASE
1800 WHEN ?2 IS NULL THEN last_activity_at_ms
1801 WHEN last_activity_at_ms IS NULL OR last_activity_at_ms < ?2 THEN ?2
1802 ELSE last_activity_at_ms
1803 END,
1804 applied_event_ordinal = ?3,
1805 applied_event_digest = ?4
1806 WHERE session_id = ?1",
1807 params![
1808 session_id,
1809 self.pending.last_activity_at_ms,
1810 self.applied_ordinal,
1811 self.applied_digest,
1812 ],
1813 )?;
1814 Ok(())
1815 }
1816}
1817
1818pub fn apply_projection_page<T>(
1823 session_id: &str,
1824 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T> + Send + 'static,
1825) -> Result<T>
1826where
1827 T: Send + 'static,
1828{
1829 let session_id = session_id.to_owned();
1830 submit_database_write("apply_projection_page", move |connection| {
1831 apply_projection_page_with(connection, &session_id, fill)
1832 })
1833}
1834
1835#[cfg(test)]
1836pub(super) fn apply_projection_page_to<T>(
1837 path: &Path,
1838 session_id: &str,
1839 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1840) -> Result<T> {
1841 let mut connection = open(path)?;
1842 apply_projection_page_with(&mut connection, session_id, fill)
1843}
1844
1845pub(super) fn apply_projection_page_with<T>(
1846 connection: &mut Connection,
1847 session_id: &str,
1848 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1849) -> Result<T> {
1850 let transaction =
1851 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1852 let (applied_ordinal, applied_digest) = transaction
1853 .query_row(
1854 "SELECT applied_event_ordinal, applied_event_digest
1855 FROM materialized_sessions WHERE session_id = ?1",
1856 [session_id],
1857 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1858 )
1859 .optional()?
1860 .with_context(|| format!("unknown session {session_id}"))?;
1861 validate_relay_event_frontier(
1862 applied_ordinal,
1863 &applied_digest,
1864 "persisted relay event frontier",
1865 )?;
1866 let usage_configuration = read_materialized_session_fields(&transaction, session_id)?
1867 .context("projection session disappeared")?
1868 .configuration;
1869 let mut page = ProjectionPage {
1870 session_id,
1871 transaction,
1872 applied_ordinal,
1873 applied_digest,
1874 dirty: false,
1875 pending: MaterializedSessionMutation::default(),
1876 pending_transcript: BTreeMap::new(),
1877 pending_turns: Vec::new(),
1878 usage_configuration,
1879 pending_events: Vec::new(),
1880 };
1881 let filled = fill(&mut page)?;
1884 page.flush()?;
1885 page.transaction.commit()?;
1886 Ok(filled)
1887}
1888
1889pub fn apply_projection_event(
1891 session_id: &str,
1892 event_ordinal: u64,
1893 previous_event_digest: &str,
1894 event_digest: &str,
1895 mutation: &MaterializedSessionMutation,
1896) -> Result<ProjectionApplyOutcome> {
1897 let session_id = session_id.to_owned();
1898 let previous_event_digest = previous_event_digest.to_owned();
1899 let event_digest = event_digest.to_owned();
1900 let mutation = mutation.clone();
1901 submit_database_write("apply_projection_event", move |connection| {
1902 apply_projection_page_with(connection, &session_id, |page| {
1903 page.apply(
1904 event_ordinal,
1905 &previous_event_digest,
1906 &event_digest,
1907 &mutation,
1908 )
1909 })
1910 })
1911}
1912
1913#[cfg(test)]
1914pub(super) fn apply_projection_event_to(
1915 path: &Path,
1916 session_id: &str,
1917 event_ordinal: u64,
1918 previous_event_digest: &str,
1919 event_digest: &str,
1920 mutation: &MaterializedSessionMutation,
1921) -> Result<ProjectionApplyOutcome> {
1922 apply_projection_page_to(path, session_id, |page| {
1923 page.apply(event_ordinal, previous_event_digest, event_digest, mutation)
1924 })
1925}
1926
1927pub fn advance_viewed_through_event_ordinal(session_id: &str, through: u64) -> Result<u64> {
1930 let session_id = session_id.to_owned();
1931 submit_database_write("advance_viewed_through_event_ordinal", move |_| {
1932 advance_viewed_through_event_ordinal_to(&database_path(), &session_id, through)
1933 })
1934}
1935
1936pub(super) fn advance_viewed_through_event_ordinal_to(
1937 path: &Path,
1938 session_id: &str,
1939 through: u64,
1940) -> Result<u64> {
1941 let mut connection = open(path)?;
1942 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1943 let applied = tx
1944 .query_row(
1945 "SELECT applied_event_ordinal FROM materialized_sessions WHERE session_id = ?1",
1946 [session_id],
1947 |row| row.get::<_, u64>(0),
1948 )
1949 .optional()?
1950 .with_context(|| format!("unknown session {session_id}"))?;
1951 if through > applied {
1952 bail!(
1953 "cannot acknowledge event ordinal {through} for session {session_id}; projection is at {applied}"
1954 );
1955 }
1956 tx.execute(
1957 "UPDATE sessions
1958 SET viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?2)
1959 WHERE session_id = ?1",
1960 params![session_id, through],
1961 )?;
1962 let receipt = tx.query_row(
1963 "SELECT viewed_through_event_ordinal FROM sessions WHERE session_id = ?1",
1964 [session_id],
1965 |row| row.get::<_, u64>(0),
1966 )?;
1967 tx.commit()?;
1968 Ok(receipt)
1969}
1970
1971fn decode_transcript_body(body_json: &str, session_id: &str) -> Result<TranscriptBody> {
1977 let mut body: TranscriptBody = serde_json::from_str(body_json)
1978 .with_context(|| format!("parse materialized transcript body for session {session_id}"))?;
1979 match &mut body {
1980 TranscriptBody::Agent { chunks, .. } | TranscriptBody::Thought { chunks, .. } => {
1981 mj_core::transcript::coalesce_content_chunks(chunks);
1982 }
1983 _ => {}
1984 }
1985 Ok(body)
1986}
1987
1988pub(crate) fn load_continuation_evidence(
1991 session_id: &str,
1992 ordinal: u64,
1993 digest: &str,
1994) -> Result<mj_core::continuation::ContinuationEvidence> {
1995 load_continuation_evidence_from(&database_path(), session_id, ordinal, digest)
1996}
1997
1998pub(super) fn load_continuation_evidence_from(
1999 path: &Path,
2000 session_id: &str,
2001 ordinal: u64,
2002 digest: &str,
2003) -> Result<mj_core::continuation::ContinuationEvidence> {
2004 let mut reader = open_reader(path)?;
2005 let connection = reader.transaction()?;
2006 let fields = read_materialized_session_fields(&connection, session_id)?
2007 .context("continuation projection is missing")?;
2008 anyhow::ensure!(
2009 fields.applied_event_ordinal == ordinal && fields.applied_event_digest == digest,
2010 "continuation projection changed"
2011 );
2012 let mut statement = connection.prepare(
2013 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
2014 last_changed_at_ms, body_json
2015 FROM materialized_transcript_items WHERE session_id = ?1
2016 AND (json_extract(body_json, '$.kind') IN ('user', 'agent')
2017 OR stable_id LIKE 'context-cleared:%')
2018 ORDER BY position DESC, stable_id DESC",
2019 )?;
2020 let rows = statement.query_map([session_id], |row| {
2021 Ok((
2022 row.get::<_, String>(0)?,
2023 row.get::<_, u64>(1)?,
2024 row.get::<_, Option<u64>>(2)?,
2025 row.get::<_, i64>(3)?,
2026 row.get::<_, i64>(4)?,
2027 row.get::<_, String>(5)?,
2028 ))
2029 })?;
2030 crate::continuation::evidence_from_items(rows.map(|row| {
2031 let (
2032 stable_id,
2033 position,
2034 latest_content_event_ordinal,
2035 created_at_ms,
2036 last_changed_at_ms,
2037 body_json,
2038 ) = row?;
2039 anyhow::ensure!(
2041 body_json.len() <= 1024 * 1024,
2042 "continuation source message exceeds budget"
2043 );
2044 Ok(Arc::new(TranscriptItem {
2045 stable_id,
2046 position,
2047 latest_content_event_ordinal,
2048 created_at_ms,
2049 last_changed_at_ms,
2050 body: decode_transcript_body(&body_json, session_id)?,
2051 }))
2052 }))
2053}