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 role: Option<mj_core::transcript::TranscriptRole>,
560) -> Result<Option<TranscriptPage>> {
561 load_materialized_transcript_filtered_from(&database_path(), session_id, after_seq, limit, role)
562}
563
564pub(super) fn load_materialized_transcript_filtered_from(
565 path: &Path,
566 session_id: &str,
567 after_seq: u64,
568 limit: usize,
569 role: Option<mj_core::transcript::TranscriptRole>,
570) -> Result<Option<TranscriptPage>> {
571 let mut reader = open_reader(path)?;
572 let connection = reader.transaction()?;
573 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
574 return Ok(None);
575 };
576 let role = role.map(|r| r.storage_kind());
577 let mut statement = connection.prepare(
578 "WITH matches AS (
579 SELECT *, COALESCE(latest_content_event_ordinal, position) AS seq
580 FROM materialized_transcript_items WHERE session_id = ?1
581 AND COALESCE(latest_content_event_ordinal, position) > ?2
582 AND (?4 IS NULL OR json_extract(body_json, '$.kind') = ?4)
583 ), boundary AS (SELECT MAX(seq) AS seq FROM (SELECT seq FROM matches ORDER BY seq LIMIT ?3))
584 SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
585 last_changed_at_ms, body_json
586 FROM matches WHERE seq <= (SELECT seq FROM boundary)
587 ORDER BY seq, stable_id",
588 )?;
589 let rows = statement
590 .query_map(
591 params![session_id, after_seq, limit.clamp(1, 1000) as i64, role],
592 |row| {
593 Ok((
594 row.get::<_, String>(0)?,
595 row.get::<_, u64>(1)?,
596 row.get::<_, Option<u64>>(2)?,
597 row.get::<_, i64>(3)?,
598 row.get::<_, i64>(4)?,
599 row.get::<_, String>(5)?,
600 ))
601 },
602 )?
603 .collect::<rusqlite::Result<Vec<_>>>()?;
604 let items = rows
605 .into_iter()
606 .map(
607 |(
608 stable_id,
609 position,
610 latest_content_event_ordinal,
611 created_at_ms,
612 last_changed_at_ms,
613 body_json,
614 )| {
615 Ok(Arc::new(TranscriptItem {
616 stable_id,
617 position,
618 latest_content_event_ordinal,
619 created_at_ms,
620 last_changed_at_ms,
621 body: decode_transcript_body(&body_json, session_id)?,
622 }))
623 },
624 )
625 .collect::<Result<Vec<_>>>()?;
626 let latest_seq = connection.query_row(
627 "SELECT COALESCE(MAX(COALESCE(latest_content_event_ordinal, position)), 0)
628 FROM materialized_transcript_items
629 WHERE session_id = ?1",
630 [session_id],
631 |row| row.get::<_, u64>(0),
632 )?;
633 let last_seq = items.last().map_or(after_seq, |item| item.seq());
634 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))?;
635 Ok(Some(TranscriptPage {
636 next_after_seq: if more {
637 last_seq
638 } else {
639 latest_seq.max(after_seq)
640 },
641 items,
642 latest_seq,
643 execution: fields.execution,
644 }))
645}
646
647pub(super) const RETENTION_BATCH_ITEMS: usize = 4_096;
654
655pub(super) const RETENTION_BODY_FLOOR_BYTES: usize = 4 * 1024;
658
659pub fn compact_materialized_transcript_through(
672 session_id: &str,
673 event_frontier: u64,
674) -> Result<TranscriptRetention> {
675 let session_id = session_id.to_owned();
676 submit_database_write("compact_materialized_transcript", move |_| {
677 compact_materialized_transcript_in(&database_path(), &session_id, event_frontier)
678 })
679}
680
681pub(super) fn compact_materialized_transcript_in(
682 path: &Path,
683 session_id: &str,
684 event_frontier: u64,
685) -> Result<TranscriptRetention> {
686 let mut connection = open(path)?;
687 let candidates = {
688 let mut statement = connection.prepare(
689 "SELECT stable_id, body_json
690 FROM materialized_transcript_items
691 WHERE session_id = ?1
692 AND position <= ?2
693 AND length(body_json) > ?3
694 AND json_extract(
695 CASE WHEN json_valid(body_json) THEN body_json ELSE '{}' END,
696 '$.kind'
697 ) = 'tool'
698 ORDER BY position, stable_id
699 LIMIT ?4",
700 )?;
701 statement
702 .query_map(
703 params![
704 session_id,
705 event_frontier,
706 RETENTION_BODY_FLOOR_BYTES as i64,
707 RETENTION_BATCH_ITEMS as i64 + 1
708 ],
709 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
710 )?
711 .collect::<rusqlite::Result<Vec<_>>>()?
712 };
713 let remaining = candidates.len() > RETENTION_BATCH_ITEMS;
714 let mut retention = TranscriptRetention {
715 remaining,
716 ..TranscriptRetention::default()
717 };
718 let transaction =
719 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
720 for (stable_id, body_json) in candidates.into_iter().take(RETENTION_BATCH_ITEMS) {
721 let mut body: TranscriptBody = match serde_json::from_str(&body_json) {
722 Ok(body) => body,
723 Err(error) => {
725 tracing::warn!(%session_id, %stable_id, %error, "skipping unreadable transcript body");
726 continue;
727 }
728 };
729 if !mj_transcript::transcript::compact_tool_call_for_retention(&mut body) {
730 continue;
731 }
732 let compacted = serde_json::to_string(&body)
733 .with_context(|| format!("serialize compacted transcript body {stable_id}"))?;
734 if compacted.len() >= body_json.len() {
735 continue;
736 }
737 transaction.execute(
738 "UPDATE materialized_transcript_items SET body_json = ?3
739 WHERE session_id = ?1 AND stable_id = ?2",
740 params![session_id, stable_id, compacted],
741 )?;
742 retention.items += 1;
743 retention.bytes += body_json.len() - compacted.len();
744 }
745 transaction.commit()?;
746 Ok(retention)
747}
748
749pub const PROJECTION_TAIL_ITEMS: usize = 1_024;
757
758pub fn load_materialized_actor_projection(
761 session_id: &str,
762) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
763 load_materialized_actor_projection_from(&database_path(), session_id)
764}
765
766pub(super) fn load_materialized_actor_projection_from(
767 path: &Path,
768 session_id: &str,
769) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
770 let mut reader = open_reader(path)?;
771 let connection = reader.transaction()?;
772 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
773 return Ok(None);
774 };
775 let latest_turn = last_materialized_turn_start(&connection, session_id)?;
776 let desired: Option<u64> = connection
777 .query_row(
778 "SELECT position FROM materialized_transcript_items WHERE session_id=?1
779 ORDER BY position DESC, stable_id DESC LIMIT 1 OFFSET ?2",
780 params![session_id, PROJECTION_TAIL_ITEMS - 1],
781 |row| row.get(0),
782 )
783 .optional()?;
784 let mutable: Option<u64> = connection.query_row(
787 "SELECT MIN(position) FROM materialized_transcript_items WHERE session_id=?1
788 AND (json_extract(body_json, '$.streaming')=1
789 OR (json_extract(body_json, '$.kind')='tool'
790 AND json_extract(body_json, '$.call.status') IN ('pending','in_progress')))",
791 [session_id],
792 |row| row.get(0),
793 )?;
794 let boundary = desired
795 .unwrap_or(0)
796 .min(latest_turn.unwrap_or(0))
797 .min(mutable.unwrap_or(u64::MAX));
798 let start: u64 = connection.query_row(
799 "SELECT COALESCE(MAX(position),0) FROM materialized_transcript_items
800 WHERE session_id=?1 AND position<=?2
801 AND (json_extract(body_json, '$.kind')='user' OR stable_id LIKE 'harness-turn:%')",
802 params![session_id, boundary],
803 |row| row.get(0),
804 )?;
805 let transcript = read_transcript_range(&connection, session_id, start, None, None)?;
806 let omitted_items = connection.query_row(
807 "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id=?1 AND position<?2",
808 params![session_id, start],
809 |row| row.get(0),
810 )?;
811 let window = ProjectionWindow {
812 omitted_items,
813 provisional_title: first_materialized_user_message(&connection, session_id)?
814 .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
815 latest_turn_start_position: latest_turn,
816 };
817 let materialized = MaterializedSession {
818 session_id: session_id.to_owned(),
819 applied_event_ordinal: fields.applied_event_ordinal,
820 applied_event_digest: fields.applied_event_digest,
821 last_activity_at_ms: fields.last_activity_at_ms,
822 execution: fields.execution,
823 session_title: fields.session_title,
824 configuration: fields.configuration,
825 transcript,
826 queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
827 pending_elicitations: fields.pending_elicitations,
828 active_turn: fields.active_turn,
829 last_turn_outcome: fields.last_turn_outcome,
830 };
831 materialized.validate()?;
832 Ok(Some((materialized, window)))
833}
834
835pub fn load_transcript_history(
836 session_id: &str,
837 before: Option<&mj_core::storage::TranscriptCursor>,
838 limit: usize,
839) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
840 load_transcript_history_from(&database_path(), session_id, before, limit)
841}
842
843pub(super) fn load_transcript_history_from(
844 path: &Path,
845 session_id: &str,
846 before: Option<&mj_core::storage::TranscriptCursor>,
847 limit: usize,
848) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
849 let mut reader = open_reader(path)?;
850 let connection = reader.transaction()?;
851 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
852 return Ok(None);
853 };
854 let limit = limit.clamp(1, 256);
855 let mut items = read_transcript_range(&connection, session_id, 0, before, Some(limit + 1))?;
856 let has_more = items.len() > limit;
857 if has_more {
858 items.remove(0);
859 }
860 let before = has_more.then(|| mj_core::storage::TranscriptCursor::of(&items[0]));
861 Ok(Some(mj_core::storage::TranscriptHistoryPage {
862 items,
863 before,
864 frontier: fields.applied_event_ordinal,
865 }))
866}
867
868fn read_transcript_range(
869 connection: &Connection,
870 session_id: &str,
871 start: u64,
872 before: Option<&mj_core::storage::TranscriptCursor>,
873 limit: Option<usize>,
874) -> Result<Vec<Arc<TranscriptItem>>> {
875 let mut statement = connection.prepare(
876 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
877 last_changed_at_ms, body_json FROM materialized_transcript_items
878 WHERE session_id=?1 AND position>=?2
879 AND (position,stable_id)<(?3,?4)
880 ORDER BY position DESC, stable_id DESC LIMIT ?5",
881 )?;
882 let rows = statement
883 .query_map(
884 params![
885 session_id,
886 start,
887 before.map_or(i64::MAX as u64, |c| c.position),
888 before.map_or("", |c| c.stable_id.as_str()),
889 limit.map_or(-1, |limit| limit as i64)
890 ],
891 |row| {
892 Ok((
893 row.get::<_, String>(0)?,
894 row.get::<_, u64>(1)?,
895 row.get::<_, Option<u64>>(2)?,
896 row.get::<_, i64>(3)?,
897 row.get::<_, i64>(4)?,
898 row.get::<_, String>(5)?,
899 ))
900 },
901 )?
902 .collect::<rusqlite::Result<Vec<_>>>()?;
903 let mut items = rows
904 .into_iter()
905 .map(
906 |(
907 stable_id,
908 position,
909 latest_content_event_ordinal,
910 created_at_ms,
911 last_changed_at_ms,
912 body,
913 )| {
914 Ok(Arc::new(TranscriptItem {
915 stable_id,
916 position,
917 latest_content_event_ordinal,
918 created_at_ms,
919 last_changed_at_ms,
920 body: decode_transcript_body(&body, session_id)?,
921 }))
922 },
923 )
924 .collect::<Result<Vec<_>>>()?;
925 items.reverse();
926 Ok(items)
927}
928
929pub fn load_projection_references(
932 projection: &MaterializedSession,
933 events: &[mj_core::relay::RelayEvent],
934) -> Result<Vec<Arc<TranscriptItem>>> {
935 load_projection_references_from(&database_path(), projection, events)
936}
937
938pub(super) fn load_projection_references_from(
939 path: &Path,
940 projection: &MaterializedSession,
941 events: &[mj_core::relay::RelayEvent],
942) -> Result<Vec<Arc<TranscriptItem>>> {
943 let (mut ids, terminals) = mj_transcript::projection::historical_references(events)?;
944 let retained = projection
945 .transcript
946 .iter()
947 .map(|item| item.stable_id.as_str())
948 .collect::<std::collections::HashSet<_>>();
949 ids.retain(|id| !retained.contains(id.as_str()));
950 if ids.is_empty() && terminals.is_empty() {
951 return Ok(Vec::new());
952 }
953 let mut reader = open_reader(path)?;
954 let connection = reader.transaction()?;
955 let mut parameters = vec![projection.session_id.clone(), serde_json::to_string(&ids)?];
958 let mut statement = connection.prepare(if terminals.is_empty() {
959 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
960 last_changed_at_ms, body_json FROM materialized_transcript_items
961 WHERE session_id=?1 AND stable_id IN (SELECT value FROM json_each(?2))
962 ORDER BY position,stable_id"
963 } else {
964 parameters.push(serde_json::to_string(&retained)?);
965 parameters.push(serde_json::to_string(&terminals)?);
966 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
967 last_changed_at_ms, body_json FROM materialized_transcript_items
968 WHERE session_id=?1 AND stable_id NOT IN (SELECT value FROM json_each(?3))
969 AND (stable_id IN (SELECT value FROM json_each(?2))
970 OR EXISTS(SELECT 1 FROM json_each(body_json, '$.terminal_refs')
971 WHERE value IN (SELECT value FROM json_each(?4))))
972 ORDER BY position,stable_id"
973 })?;
974 let rows = statement
975 .query_map(rusqlite::params_from_iter(parameters), |row| {
976 Ok((
977 row.get::<_, String>(0)?,
978 row.get::<_, u64>(1)?,
979 row.get::<_, Option<u64>>(2)?,
980 row.get::<_, i64>(3)?,
981 row.get::<_, i64>(4)?,
982 row.get::<_, String>(5)?,
983 ))
984 })?
985 .collect::<rusqlite::Result<Vec<_>>>()?;
986 rows.into_iter()
987 .map(
988 |(
989 stable_id,
990 position,
991 latest_content_event_ordinal,
992 created_at_ms,
993 last_changed_at_ms,
994 body,
995 )| {
996 Ok(Arc::new(TranscriptItem {
997 stable_id,
998 position,
999 latest_content_event_ordinal,
1000 created_at_ms,
1001 last_changed_at_ms,
1002 body: decode_transcript_body(&body, &projection.session_id)?,
1003 }))
1004 },
1005 )
1006 .collect()
1007}
1008
1009pub fn load_materialized_projection_tail(
1019 session_id: &str,
1020 transcript_limit: usize,
1021) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
1022 load_materialized_projection_tail_from(&database_path(), session_id, transcript_limit)
1023}
1024
1025pub(super) fn load_materialized_projection_tail_from(
1026 path: &Path,
1027 session_id: &str,
1028 transcript_limit: usize,
1029) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
1030 let mut reader = open_reader(path)?;
1031 let connection = reader.transaction()?;
1034 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
1035 return Ok(None);
1036 };
1037 let transcript = read_materialized_transcript(&connection, session_id, Some(transcript_limit))?;
1038 let total_items = connection.query_row(
1039 "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id = ?1",
1040 [session_id],
1041 |row| row.get::<_, usize>(0),
1042 )?;
1043 let window = ProjectionWindow {
1044 omitted_items: total_items.saturating_sub(transcript.len()),
1045 provisional_title: first_materialized_user_message(&connection, session_id)?
1046 .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
1047 latest_turn_start_position: last_materialized_turn_start(&connection, session_id)?,
1048 };
1049 let materialized = MaterializedSession {
1050 session_id: session_id.to_owned(),
1051 applied_event_ordinal: fields.applied_event_ordinal,
1052 applied_event_digest: fields.applied_event_digest,
1053 last_activity_at_ms: fields.last_activity_at_ms,
1054 execution: fields.execution,
1055 session_title: fields.session_title,
1056 configuration: fields.configuration,
1057 transcript,
1058 queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
1059 pending_elicitations: fields.pending_elicitations,
1060 active_turn: fields.active_turn,
1061 last_turn_outcome: fields.last_turn_outcome,
1062 };
1063 materialized.validate()?;
1064 Ok(Some((materialized, window)))
1065}
1066
1067pub fn materialized_event_frontier(session_id: &str) -> Result<Option<(u64, String)>> {
1071 materialized_event_frontier_from(&database_path(), session_id)
1072}
1073
1074pub(super) fn materialized_event_frontier_from(
1075 path: &Path,
1076 session_id: &str,
1077) -> Result<Option<(u64, String)>> {
1078 Ok(open_reader(path)?
1079 .query_row(
1080 "SELECT applied_event_ordinal, applied_event_digest
1081 FROM materialized_sessions WHERE session_id = ?1",
1082 [session_id],
1083 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1084 )
1085 .optional()?)
1086}
1087
1088pub fn replace_materialized_queued_prompts(
1092 session_id: &str,
1093 queued_prompts: &[MaterializedQueuedPrompt],
1094) -> Result<()> {
1095 let session_id = session_id.to_owned();
1096 let queued_prompts = queued_prompts.to_vec();
1097 submit_database_write("replace_materialized_queued_prompts", move |_| {
1098 replace_materialized_queued_prompts_in(&database_path(), &session_id, &queued_prompts)
1099 })
1100}
1101
1102pub(super) fn replace_materialized_queued_prompts_in(
1103 path: &Path,
1104 session_id: &str,
1105 queued_prompts: &[MaterializedQueuedPrompt],
1106) -> Result<()> {
1107 let mut connection = open(path)?;
1108 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1109 if !session_exists(&tx, session_id)? {
1110 bail!("unknown session {session_id}");
1111 }
1112 replace_materialized_queue(&tx, session_id, queued_prompts)?;
1113 tx.commit()?;
1114 Ok(())
1115}
1116
1117pub fn load_transcribed_session_activity() -> Result<BTreeMap<String, Option<i64>>> {
1124 load_transcribed_session_activity_from(&database_path())
1125}
1126
1127fn load_transcribed_session_activity_from(path: &Path) -> Result<BTreeMap<String, Option<i64>>> {
1128 let connection = open_reader(path)?;
1129 let mut statement = connection.prepare(
1130 "SELECT session_id, last_activity_at_ms
1131 FROM materialized_sessions s
1132 WHERE EXISTS (
1133 SELECT 1 FROM materialized_transcript_items i
1134 WHERE i.session_id = s.session_id
1135 )",
1136 )?;
1137 let rows = statement.query_map([], |row| {
1138 Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?))
1139 })?;
1140 let mut activity = BTreeMap::new();
1141 for row in rows {
1142 let (session_id, last_activity_at_ms) = row?;
1143 activity.insert(session_id, last_activity_at_ms);
1144 }
1145 Ok(activity)
1146}
1147
1148pub fn load_materialized_queued_prompts() -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>>
1152{
1153 load_materialized_queued_prompts_from(&database_path())
1154}
1155
1156pub(super) fn load_materialized_queued_prompts_from(
1157 path: &Path,
1158) -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>> {
1159 let connection = open_reader(path)?;
1160 let mut statement = connection.prepare(
1161 "SELECT session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1162 FROM materialized_queued_prompts
1163 ORDER BY session_id, ordinal",
1164 )?;
1165 let rows = statement.query_map([], |row| {
1166 Ok((
1167 row.get::<_, String>(0)?,
1168 row.get::<_, String>(1)?,
1169 row.get::<_, String>(2)?,
1170 row.get::<_, String>(3)?,
1171 row.get::<_, i64>(4)?,
1172 row.get::<_, Option<u64>>(5)?,
1173 ))
1174 })?;
1175 let mut queues = BTreeMap::<String, Vec<MaterializedQueuedPrompt>>::new();
1176 for row in rows {
1177 let (session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal) =
1178 row?;
1179 let content = serde_json::from_str(&content_json).with_context(|| {
1180 format!("parse materialized queued prompt for session {session_id}")
1181 })?;
1182 let kind = serde_json::from_str(&kind_json).with_context(|| {
1183 format!("parse materialized queue entry kind for session {session_id}")
1184 })?;
1185 queues
1186 .entry(session_id)
1187 .or_default()
1188 .push(MaterializedQueuedPrompt {
1189 command_id,
1190 kind,
1191 content,
1192 queued_at_ms,
1193 accepted_ordinal,
1194 });
1195 }
1196 Ok(queues)
1197}
1198
1199pub(super) fn load_materialized_session_from(
1200 path: &Path,
1201 session_id: &str,
1202) -> Result<Option<MaterializedSession>> {
1203 let mut reader = open_reader(path)?;
1204 let connection = reader.transaction()?;
1205 load_materialized_session_with(&connection, session_id)
1206}
1207
1208pub(super) fn load_materialized_session_with(
1209 connection: &rusqlite::Transaction<'_>,
1210 session_id: &str,
1211) -> Result<Option<MaterializedSession>> {
1212 let Some(fields) = read_materialized_session_fields(connection, session_id)? else {
1213 return Ok(None);
1214 };
1215 let materialized = MaterializedSession {
1216 session_id: session_id.to_owned(),
1217 applied_event_ordinal: fields.applied_event_ordinal,
1218 applied_event_digest: fields.applied_event_digest,
1219 last_activity_at_ms: fields.last_activity_at_ms,
1220 execution: fields.execution,
1221 session_title: fields.session_title,
1222 configuration: fields.configuration,
1223 transcript: read_materialized_transcript(connection, session_id, None)?,
1224 queued_prompts: read_materialized_queued_prompts(connection, session_id)?,
1225 pending_elicitations: fields.pending_elicitations,
1226 active_turn: fields.active_turn,
1227 last_turn_outcome: fields.last_turn_outcome,
1228 };
1229 materialized.validate()?;
1230 Ok(Some(materialized))
1231}
1232
1233pub(super) struct MaterializedSessionFields {
1235 pub(super) applied_event_ordinal: u64,
1236 pub(super) applied_event_digest: String,
1237 pub(super) last_activity_at_ms: Option<i64>,
1238 pub(super) execution: MaterializedExecutionState,
1239 pub(super) session_title: Option<String>,
1240 pub(super) configuration: mj_core::state::SessionConfiguration,
1241 pub(super) pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
1242 pub(super) active_turn: Option<MaterializedTurn>,
1243 pub(super) last_turn_outcome: Option<MaterializedTurnOutcome>,
1244}
1245
1246pub(super) fn read_materialized_session_fields(
1247 connection: &Connection,
1248 session_id: &str,
1249) -> Result<Option<MaterializedSessionFields>> {
1250 let row = connection
1251 .query_row(
1252 "SELECT applied_event_ordinal, applied_event_digest, last_activity_at_ms,
1253 execution_state, running_started_at_ms, session_title, configuration_json,
1254 pending_elicitations_json, active_turn_json, last_turn_outcome_json
1255 FROM materialized_sessions WHERE session_id = ?1",
1256 [session_id],
1257 |row| {
1258 Ok((
1259 row.get::<_, u64>(0)?,
1260 row.get::<_, String>(1)?,
1261 row.get::<_, Option<i64>>(2)?,
1262 row.get::<_, String>(3)?,
1263 row.get::<_, Option<i64>>(4)?,
1264 row.get::<_, Option<String>>(5)?,
1265 row.get::<_, String>(6)?,
1266 row.get::<_, String>(7)?,
1267 row.get::<_, Option<String>>(8)?,
1268 row.get::<_, Option<String>>(9)?,
1269 ))
1270 },
1271 )
1272 .optional()?;
1273 let Some((
1274 applied_event_ordinal,
1275 applied_event_digest,
1276 last_activity_at_ms,
1277 execution,
1278 running_started_at_ms,
1279 session_title,
1280 configuration_json,
1281 pending_elicitations_json,
1282 active_turn_json,
1283 last_turn_outcome_json,
1284 )) = row
1285 else {
1286 return Ok(None);
1287 };
1288 #[cfg(test)]
1289 super::tests::after_materialized_frontier_read();
1290 Ok(Some(MaterializedSessionFields {
1291 applied_event_ordinal,
1292 applied_event_digest,
1293 last_activity_at_ms,
1294 execution: parse_materialized_execution(&execution, running_started_at_ms)?,
1295 session_title,
1296 configuration: serde_json::from_str(&configuration_json).with_context(|| {
1297 format!("parse materialized configuration for session {session_id}")
1298 })?,
1299 pending_elicitations: serde_json::from_str(&pending_elicitations_json)
1300 .with_context(|| format!("parse pending elicitations for session {session_id}"))?,
1301 active_turn: active_turn_json
1302 .as_deref()
1303 .map(serde_json::from_str)
1304 .transpose()
1305 .with_context(|| format!("parse active turn for session {session_id}"))?,
1306 last_turn_outcome: last_turn_outcome_json
1307 .as_deref()
1308 .map(serde_json::from_str)
1309 .transpose()
1310 .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
1311 }))
1312}
1313
1314pub(super) fn read_materialized_transcript(
1318 connection: &Connection,
1319 session_id: &str,
1320 limit: Option<usize>,
1321) -> Result<Vec<Arc<TranscriptItem>>> {
1322 let mut statement = connection.prepare(match limit {
1323 Some(_) => {
1324 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1325 last_changed_at_ms, body_json
1326 FROM materialized_transcript_items
1327 WHERE session_id = ?1
1328 ORDER BY position DESC, stable_id DESC
1329 LIMIT ?2"
1330 }
1331 None => {
1332 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1333 last_changed_at_ms, body_json
1334 FROM materialized_transcript_items
1335 WHERE session_id = ?1
1336 ORDER BY position, stable_id"
1337 }
1338 })?;
1339 let read = |row: &rusqlite::Row<'_>| {
1340 Ok((
1341 row.get::<_, String>(0)?,
1342 row.get::<_, u64>(1)?,
1343 row.get::<_, Option<u64>>(2)?,
1344 row.get::<_, i64>(3)?,
1345 row.get::<_, i64>(4)?,
1346 row.get::<_, String>(5)?,
1347 ))
1348 };
1349 let rows = match limit {
1350 Some(limit) => statement
1351 .query_map(params![session_id, limit as i64], read)?
1352 .collect::<rusqlite::Result<Vec<_>>>()?,
1353 None => statement
1354 .query_map([session_id], read)?
1355 .collect::<rusqlite::Result<Vec<_>>>()?,
1356 };
1357 let mut transcript = rows
1358 .into_iter()
1359 .map(
1360 |(
1361 stable_id,
1362 position,
1363 latest_content_event_ordinal,
1364 created_at_ms,
1365 last_changed_at_ms,
1366 body_json,
1367 )| {
1368 Ok(Arc::new(TranscriptItem {
1369 stable_id,
1370 position,
1371 latest_content_event_ordinal,
1372 created_at_ms,
1373 last_changed_at_ms,
1374 body: decode_transcript_body(&body_json, session_id)?,
1375 }))
1376 },
1377 )
1378 .collect::<Result<Vec<_>>>()?;
1379 if limit.is_some() {
1380 transcript.reverse();
1383 }
1384 Ok(transcript)
1385}
1386
1387pub(super) fn read_materialized_queued_prompts(
1388 connection: &Connection,
1389 session_id: &str,
1390) -> Result<Vec<MaterializedQueuedPrompt>> {
1391 let mut statement = connection.prepare(
1392 "SELECT command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1393 FROM materialized_queued_prompts
1394 WHERE session_id = ?1
1395 ORDER BY ordinal",
1396 )?;
1397 let rows = statement
1398 .query_map([session_id], |row| {
1399 Ok((
1400 row.get::<_, String>(0)?,
1401 row.get::<_, String>(1)?,
1402 row.get::<_, String>(2)?,
1403 row.get::<_, i64>(3)?,
1404 row.get::<_, Option<u64>>(4)?,
1405 ))
1406 })?
1407 .collect::<rusqlite::Result<Vec<_>>>()?;
1408 rows.into_iter()
1409 .map(
1410 |(command_id, kind_json, content_json, queued_at_ms, accepted_ordinal)| {
1411 Ok(MaterializedQueuedPrompt {
1412 command_id,
1413 kind: serde_json::from_str(&kind_json).with_context(|| {
1414 format!("parse materialized queue entry kind for session {session_id}")
1415 })?,
1416 content: serde_json::from_str(&content_json).with_context(|| {
1417 format!("parse materialized queued prompt for session {session_id}")
1418 })?,
1419 queued_at_ms,
1420 accepted_ordinal,
1421 })
1422 },
1423 )
1424 .collect()
1425}
1426
1427pub fn save_materialized_session(materialized: &MaterializedSession) -> Result<()> {
1431 let materialized = materialized.clone();
1432 submit_database_write("save_materialized_session", move |_| {
1433 save_materialized_session_to(&database_path(), &materialized)
1434 })
1435}
1436
1437pub(super) fn save_materialized_session_to(
1438 path: &Path,
1439 materialized: &MaterializedSession,
1440) -> Result<()> {
1441 materialized.validate()?;
1442 let mut connection = open(path)?;
1443 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1444 if !session_exists(&tx, &materialized.session_id)? {
1445 bail!("unknown session {}", materialized.session_id);
1446 }
1447 write_materialized_session(&tx, materialized)?;
1448 tx.commit()?;
1449 Ok(())
1450}
1451
1452pub struct ProjectionPage<'a> {
1457 pub(super) session_id: &'a str,
1458 pub(super) transaction: Transaction<'a>,
1459 pub(super) applied_ordinal: u64,
1460 pub(super) applied_digest: String,
1461 pub(super) dirty: bool,
1462 pub(super) pending: MaterializedSessionMutation,
1463 pub(super) pending_transcript: BTreeMap<String, PendingTranscriptMutation>,
1464 pub(super) pending_turns: Vec<MaterializedTurnOutcome>,
1465 usage_configuration: mj_core::state::SessionConfiguration,
1468 pub(super) pending_events: Vec<(i64, ApiEventData)>,
1469}
1470
1471pub(super) struct PendingTranscriptMutation {
1472 pub(super) final_mutation: TranscriptMutation,
1473 pub(super) remove_before_upsert: bool,
1474}
1475
1476impl ProjectionPage<'_> {
1477 pub fn apply(
1481 &mut self,
1482 event_ordinal: u64,
1483 previous_event_digest: &str,
1484 event_digest: &str,
1485 mutation: &MaterializedSessionMutation,
1486 ) -> Result<ProjectionApplyOutcome> {
1487 if event_ordinal == 0 {
1488 bail!("relay event ordinal must be positive");
1489 }
1490 let chained = !previous_event_digest.is_empty();
1496 if chained {
1497 validate_relay_event_digest(previous_event_digest, "previous relay event digest")?;
1498 }
1499 validate_relay_event_frontier(event_ordinal, event_digest, "relay event frontier")?;
1500 let session_id = self.session_id;
1501 let applied = self.applied_ordinal;
1502 if event_ordinal < applied {
1503 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1504 }
1505 if event_ordinal == applied {
1506 if event_digest != self.applied_digest {
1507 bail!(
1508 "relay event digest mismatch for session {session_id} at ordinal {event_ordinal}: projection has {}, received {event_digest}",
1509 self.applied_digest
1510 );
1511 }
1512 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1513 }
1514 let expected = applied
1515 .checked_add(1)
1516 .context("materialized event ordinal overflow")?;
1517 if event_ordinal != expected {
1518 bail!(
1519 "relay event gap for session {session_id}: expected ordinal {expected}, received {event_ordinal}"
1520 );
1521 }
1522 if chained && previous_event_digest != self.applied_digest {
1523 bail!(
1524 "relay event chain diverged for session {session_id} before ordinal {event_ordinal}: projection has {}, event follows {previous_event_digest}",
1525 self.applied_digest
1526 );
1527 }
1528
1529 if let Some(event) = &mutation.native_agent {
1530 native_agents::apply_native_agent_event(&self.transaction, session_id, event)?;
1531 }
1532 if let Some(activity_at_ms) = mutation.last_activity_at_ms {
1533 self.pending.last_activity_at_ms = Some(
1534 self.pending
1535 .last_activity_at_ms
1536 .map_or(activity_at_ms, |existing| existing.max(activity_at_ms)),
1537 );
1538 }
1539 if let Some(execution) = mutation.execution {
1540 self.pending.execution = Some(execution);
1541 }
1542 if let Some(title) = &mutation.session_title {
1543 if title.as_ref().is_some_and(|title| title.trim().is_empty()) {
1544 bail!("materialized session title cannot be empty");
1545 }
1546 self.pending.session_title = Some(title.clone());
1547 }
1548 if let Some(configuration) = &mutation.configuration {
1549 self.usage_configuration = configuration.clone();
1550 self.pending.configuration = Some(configuration.clone());
1551 }
1552 for item_mutation in &mutation.transcript {
1553 match item_mutation {
1554 TranscriptMutation::Upsert(item) => {
1555 item.validate(event_ordinal)?;
1556 let stable_id = item.stable_id.clone();
1557 let entry = self.pending_transcript.entry(stable_id).or_insert_with(|| {
1558 PendingTranscriptMutation {
1559 final_mutation: TranscriptMutation::Upsert(item.clone()),
1560 remove_before_upsert: false,
1561 }
1562 });
1563 entry.remove_before_upsert |=
1564 matches!(&entry.final_mutation, TranscriptMutation::Remove { .. });
1565 entry.final_mutation = TranscriptMutation::Upsert(item.clone());
1566 }
1567 TranscriptMutation::Remove { stable_id } => {
1568 if stable_id.trim().is_empty() {
1569 bail!("cannot remove a transcript item with an empty stable id");
1570 }
1571 let removed = TranscriptMutation::Remove {
1572 stable_id: stable_id.clone(),
1573 };
1574 self.pending_transcript
1575 .entry(stable_id.clone())
1576 .and_modify(|entry| entry.final_mutation = removed.clone())
1577 .or_insert(PendingTranscriptMutation {
1578 final_mutation: removed,
1579 remove_before_upsert: false,
1580 });
1581 }
1582 }
1583 }
1584 if let Some(queued_prompts) = &mutation.queued_prompts {
1585 self.pending.queued_prompts = Some(queued_prompts.clone());
1586 }
1587 if let Some(pending_elicitations) = &mutation.pending_elicitations {
1588 self.pending.pending_elicitations = Some(pending_elicitations.clone());
1589 }
1590 self.pending
1591 .config_results
1592 .extend(mutation.config_results.clone());
1593 if let Some(active_turn) = &mutation.active_turn {
1594 if let Some(turn) = active_turn {
1595 super::usage::record_turn_selection(
1596 &self.transaction,
1597 session_id,
1598 &turn.command_id,
1599 &self.usage_configuration,
1600 )?;
1601 }
1602 self.pending.active_turn = Some(active_turn.clone());
1603 }
1604 if mutation.clear_turn_outcome {
1605 self.pending.clear_turn_outcome = true;
1606 self.pending.last_turn_outcome = None;
1607 }
1608 if let Some(last_turn_outcome) = &mutation.last_turn_outcome {
1609 self.pending_turns.push(last_turn_outcome.clone());
1610 self.pending.last_turn_outcome = Some(last_turn_outcome.clone());
1611 }
1612 if let Some(cost) = &mutation.provider_cost {
1613 self.pending.provider_cost = Some(cost.clone());
1614 }
1615 self.pending_events.extend(
1616 mutation
1617 .api_events
1618 .iter()
1619 .cloned()
1620 .map(|event| (mutation.last_activity_at_ms.unwrap_or(0), event)),
1621 );
1622 self.applied_ordinal = event_ordinal;
1623 event_digest.clone_into(&mut self.applied_digest);
1624 self.dirty = true;
1625 Ok(ProjectionApplyOutcome::Applied)
1626 }
1627
1628 pub(super) fn flush(&mut self) -> Result<()> {
1632 if !self.dirty {
1633 return Ok(());
1634 }
1635 let tx = &self.transaction;
1636 let session_id = self.session_id;
1637 if let Some(execution) = self.pending.execution {
1638 let (state, started_at_ms) = materialized_execution_columns(execution);
1639 tx.execute(
1640 "UPDATE materialized_sessions
1641 SET execution_state = ?2, running_started_at_ms = ?3
1642 WHERE session_id = ?1",
1643 params![session_id, state, started_at_ms],
1644 )?;
1645 }
1646 if let Some(title) = &self.pending.session_title {
1647 tx.execute(
1648 "UPDATE materialized_sessions SET session_title = ?2 WHERE session_id = ?1",
1649 params![session_id, title],
1650 )?;
1651 }
1652 if let Some(configuration) = &self.pending.configuration {
1653 tx.execute(
1654 "UPDATE materialized_sessions SET configuration_json = ?2 WHERE session_id = ?1",
1655 params![session_id, serde_json::to_string(configuration)?],
1656 )?;
1657 }
1658 for pending in self.pending_transcript.values() {
1659 match &pending.final_mutation {
1660 TranscriptMutation::Upsert(item) => {
1661 if pending.remove_before_upsert {
1665 tx.execute(
1666 "DELETE FROM materialized_transcript_items
1667 WHERE session_id = ?1 AND stable_id = ?2",
1668 params![session_id, item.stable_id],
1669 )?;
1670 }
1671 upsert_transcript_item(tx, session_id, item)?;
1672 }
1673 TranscriptMutation::Remove { stable_id } => {
1674 tx.execute(
1675 "DELETE FROM materialized_transcript_items
1676 WHERE session_id = ?1 AND stable_id = ?2",
1677 params![session_id, stable_id],
1678 )?;
1679 }
1680 }
1681 }
1682 if let Some(queued_prompts) = &self.pending.queued_prompts {
1683 replace_materialized_queue(tx, session_id, queued_prompts)?;
1684 }
1685 if let Some(pending_elicitations) = &self.pending.pending_elicitations {
1686 tx.execute(
1687 "UPDATE materialized_sessions
1688 SET pending_elicitations_json = ?2 WHERE session_id = ?1",
1689 params![session_id, serde_json::to_string(pending_elicitations)?],
1690 )?;
1691 }
1692 for (recorded_at_ms, event) in &self.pending_events {
1693 events::insert_api_event(tx, session_id, *recorded_at_ms, event)?;
1694 }
1695 for turn in &self.pending_turns {
1696 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)?])?;
1697 }
1698 if let Some(cost) = &self.pending.provider_cost {
1699 tx.execute(
1700 "INSERT OR REPLACE INTO session_provider_cost(session_id, body) VALUES (?1, ?2)",
1701 params![session_id, serde_json::to_string(cost)?],
1702 )?;
1703 }
1704 for (command_id, error) in &self.pending.config_results {
1705 tx.execute("INSERT OR REPLACE INTO api_config_results(session_id, command_id, error) VALUES (?1, ?2, ?3)", params![session_id, command_id, error])?;
1706 }
1707 if let Some(active_turn) = &self.pending.active_turn {
1708 tx.execute(
1709 "UPDATE materialized_sessions SET active_turn_json = ?2 WHERE session_id = ?1",
1710 params![
1711 session_id,
1712 active_turn
1713 .as_ref()
1714 .map(serde_json::to_string)
1715 .transpose()?
1716 ],
1717 )?;
1718 }
1719 if self.pending.clear_turn_outcome {
1720 self.transaction.execute("UPDATE materialized_sessions SET last_turn_outcome_json = NULL WHERE session_id = ?1", [session_id])?;
1721 }
1722 if let Some(last_turn_outcome) = &self.pending.last_turn_outcome {
1723 tx.execute(
1724 "UPDATE materialized_sessions
1725 SET last_turn_outcome_json = ?2 WHERE session_id = ?1",
1726 params![session_id, serde_json::to_string(last_turn_outcome)?],
1727 )?;
1728 }
1729 tx.execute(
1730 "UPDATE materialized_sessions
1731 SET last_activity_at_ms = CASE
1732 WHEN ?2 IS NULL THEN last_activity_at_ms
1733 WHEN last_activity_at_ms IS NULL OR last_activity_at_ms < ?2 THEN ?2
1734 ELSE last_activity_at_ms
1735 END,
1736 applied_event_ordinal = ?3,
1737 applied_event_digest = ?4
1738 WHERE session_id = ?1",
1739 params![
1740 session_id,
1741 self.pending.last_activity_at_ms,
1742 self.applied_ordinal,
1743 self.applied_digest,
1744 ],
1745 )?;
1746 Ok(())
1747 }
1748}
1749
1750pub fn apply_projection_page<T>(
1755 session_id: &str,
1756 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T> + Send + 'static,
1757) -> Result<T>
1758where
1759 T: Send + 'static,
1760{
1761 let session_id = session_id.to_owned();
1762 submit_database_write("apply_projection_page", move |connection| {
1763 apply_projection_page_with(connection, &session_id, fill)
1764 })
1765}
1766
1767#[cfg(test)]
1768pub(super) fn apply_projection_page_to<T>(
1769 path: &Path,
1770 session_id: &str,
1771 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1772) -> Result<T> {
1773 let mut connection = open(path)?;
1774 apply_projection_page_with(&mut connection, session_id, fill)
1775}
1776
1777pub(super) fn apply_projection_page_with<T>(
1778 connection: &mut Connection,
1779 session_id: &str,
1780 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1781) -> Result<T> {
1782 let transaction =
1783 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1784 let (applied_ordinal, applied_digest) = transaction
1785 .query_row(
1786 "SELECT applied_event_ordinal, applied_event_digest
1787 FROM materialized_sessions WHERE session_id = ?1",
1788 [session_id],
1789 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1790 )
1791 .optional()?
1792 .with_context(|| format!("unknown session {session_id}"))?;
1793 validate_relay_event_frontier(
1794 applied_ordinal,
1795 &applied_digest,
1796 "persisted relay event frontier",
1797 )?;
1798 let usage_configuration = read_materialized_session_fields(&transaction, session_id)?
1799 .context("projection session disappeared")?
1800 .configuration;
1801 let mut page = ProjectionPage {
1802 session_id,
1803 transaction,
1804 applied_ordinal,
1805 applied_digest,
1806 dirty: false,
1807 pending: MaterializedSessionMutation::default(),
1808 pending_transcript: BTreeMap::new(),
1809 pending_turns: Vec::new(),
1810 usage_configuration,
1811 pending_events: Vec::new(),
1812 };
1813 let filled = fill(&mut page)?;
1816 page.flush()?;
1817 page.transaction.commit()?;
1818 Ok(filled)
1819}
1820
1821pub fn apply_projection_event(
1823 session_id: &str,
1824 event_ordinal: u64,
1825 previous_event_digest: &str,
1826 event_digest: &str,
1827 mutation: &MaterializedSessionMutation,
1828) -> Result<ProjectionApplyOutcome> {
1829 let session_id = session_id.to_owned();
1830 let previous_event_digest = previous_event_digest.to_owned();
1831 let event_digest = event_digest.to_owned();
1832 let mutation = mutation.clone();
1833 submit_database_write("apply_projection_event", move |connection| {
1834 apply_projection_page_with(connection, &session_id, |page| {
1835 page.apply(
1836 event_ordinal,
1837 &previous_event_digest,
1838 &event_digest,
1839 &mutation,
1840 )
1841 })
1842 })
1843}
1844
1845#[cfg(test)]
1846pub(super) fn apply_projection_event_to(
1847 path: &Path,
1848 session_id: &str,
1849 event_ordinal: u64,
1850 previous_event_digest: &str,
1851 event_digest: &str,
1852 mutation: &MaterializedSessionMutation,
1853) -> Result<ProjectionApplyOutcome> {
1854 apply_projection_page_to(path, session_id, |page| {
1855 page.apply(event_ordinal, previous_event_digest, event_digest, mutation)
1856 })
1857}
1858
1859pub fn advance_viewed_through_event_ordinal(session_id: &str, through: u64) -> Result<u64> {
1862 let session_id = session_id.to_owned();
1863 submit_database_write("advance_viewed_through_event_ordinal", move |_| {
1864 advance_viewed_through_event_ordinal_to(&database_path(), &session_id, through)
1865 })
1866}
1867
1868pub(super) fn advance_viewed_through_event_ordinal_to(
1869 path: &Path,
1870 session_id: &str,
1871 through: u64,
1872) -> Result<u64> {
1873 let mut connection = open(path)?;
1874 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1875 let applied = tx
1876 .query_row(
1877 "SELECT applied_event_ordinal FROM materialized_sessions WHERE session_id = ?1",
1878 [session_id],
1879 |row| row.get::<_, u64>(0),
1880 )
1881 .optional()?
1882 .with_context(|| format!("unknown session {session_id}"))?;
1883 if through > applied {
1884 bail!(
1885 "cannot acknowledge event ordinal {through} for session {session_id}; projection is at {applied}"
1886 );
1887 }
1888 tx.execute(
1889 "UPDATE sessions
1890 SET viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?2)
1891 WHERE session_id = ?1",
1892 params![session_id, through],
1893 )?;
1894 let receipt = tx.query_row(
1895 "SELECT viewed_through_event_ordinal FROM sessions WHERE session_id = ?1",
1896 [session_id],
1897 |row| row.get::<_, u64>(0),
1898 )?;
1899 tx.commit()?;
1900 Ok(receipt)
1901}
1902
1903fn decode_transcript_body(body_json: &str, session_id: &str) -> Result<TranscriptBody> {
1909 let mut body: TranscriptBody = serde_json::from_str(body_json)
1910 .with_context(|| format!("parse materialized transcript body for session {session_id}"))?;
1911 match &mut body {
1912 TranscriptBody::Agent { chunks, .. } | TranscriptBody::Thought { chunks, .. } => {
1913 mj_core::transcript::coalesce_content_chunks(chunks);
1914 }
1915 _ => {}
1916 }
1917 Ok(body)
1918}
1919
1920pub(crate) fn load_continuation_evidence(
1923 session_id: &str,
1924 ordinal: u64,
1925 digest: &str,
1926) -> Result<mj_core::continuation::ContinuationEvidence> {
1927 load_continuation_evidence_from(&database_path(), session_id, ordinal, digest)
1928}
1929
1930pub(super) fn load_continuation_evidence_from(
1931 path: &Path,
1932 session_id: &str,
1933 ordinal: u64,
1934 digest: &str,
1935) -> Result<mj_core::continuation::ContinuationEvidence> {
1936 let mut reader = open_reader(path)?;
1937 let connection = reader.transaction()?;
1938 let fields = read_materialized_session_fields(&connection, session_id)?
1939 .context("continuation projection is missing")?;
1940 anyhow::ensure!(
1941 fields.applied_event_ordinal == ordinal && fields.applied_event_digest == digest,
1942 "continuation projection changed"
1943 );
1944 let mut statement = connection.prepare(
1945 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1946 last_changed_at_ms, body_json
1947 FROM materialized_transcript_items WHERE session_id = ?1
1948 AND (json_extract(body_json, '$.kind') IN ('user', 'agent')
1949 OR stable_id LIKE 'context-cleared:%')
1950 ORDER BY position DESC, stable_id DESC",
1951 )?;
1952 let rows = statement.query_map([session_id], |row| {
1953 Ok((
1954 row.get::<_, String>(0)?,
1955 row.get::<_, u64>(1)?,
1956 row.get::<_, Option<u64>>(2)?,
1957 row.get::<_, i64>(3)?,
1958 row.get::<_, i64>(4)?,
1959 row.get::<_, String>(5)?,
1960 ))
1961 })?;
1962 crate::continuation::evidence_from_items(rows.map(|row| {
1963 let (
1964 stable_id,
1965 position,
1966 latest_content_event_ordinal,
1967 created_at_ms,
1968 last_changed_at_ms,
1969 body_json,
1970 ) = row?;
1971 anyhow::ensure!(
1973 body_json.len() <= 1024 * 1024,
1974 "continuation source message exceeds budget"
1975 );
1976 Ok(Arc::new(TranscriptItem {
1977 stable_id,
1978 position,
1979 latest_content_event_ordinal,
1980 created_at_ms,
1981 last_changed_at_ms,
1982 body: decode_transcript_body(&body_json, session_id)?,
1983 }))
1984 }))
1985}