1use super::*;
2
3pub fn load_materialized_session(session_id: &str) -> Result<Option<MaterializedSession>> {
13 load_materialized_session_from(&database_path(), session_id)
14}
15
16pub fn load_materialized_session_summary(
20 session_id: &str,
21) -> Result<Option<MaterializedSessionSummary>> {
22 load_materialized_session_summary_from(&database_path(), session_id)
23}
24
25pub(super) fn load_materialized_session_summary_from(
26 path: &Path,
27 session_id: &str,
28) -> Result<Option<MaterializedSessionSummary>> {
29 let mut reader = open_reader(path)?;
30 let connection = reader.transaction()?;
31 let row = connection
32 .query_row(
33 "SELECT applied_event_ordinal, last_activity_at_ms, execution_state,
34 running_started_at_ms, session_title
35 FROM materialized_sessions WHERE session_id = ?1",
36 [session_id],
37 |row| {
38 Ok((
39 row.get::<_, u64>(0)?,
40 row.get::<_, Option<i64>>(1)?,
41 row.get::<_, String>(2)?,
42 row.get::<_, Option<i64>>(3)?,
43 row.get::<_, Option<String>>(4)?,
44 ))
45 },
46 )
47 .optional()?;
48 let Some((
49 applied_event_ordinal,
50 last_activity_at_ms,
51 execution,
52 running_started_at_ms,
53 session_title,
54 )) = row
55 else {
56 return Ok(None);
57 };
58
59 #[cfg(test)]
60 super::tests::after_materialized_frontier_read();
61 let last_user_message = last_materialized_user_message(&connection, session_id)?;
62 let last_agent_message = last_materialized_agent_message(&connection, session_id)?;
63 let last_agent_message_follows_last_user =
64 last_agent_message
65 .as_ref()
66 .is_some_and(|(agent_position, _)| {
67 last_user_message
68 .as_ref()
69 .is_none_or(|(user_position, _)| agent_position > user_position)
70 });
71 let mut ordinal_statement = connection.prepare(
72 "SELECT latest_content_event_ordinal
73 FROM materialized_transcript_items
74 WHERE session_id = ?1
75 AND latest_content_event_ordinal IS NOT NULL
76 AND EXISTS (
77 SELECT 1 FROM json_each(
78 CASE
79 WHEN latest_content_event_ordinal IS NOT NULL
80 AND json_valid(body_json)
81 THEN body_json
82 ELSE '{}'
83 END,
84 '$.chunks'
85 ) AS chunk
86 WHERE json_extract(chunk.value, '$.content.type') IS NOT NULL
87 AND (
88 json_extract(chunk.value, '$.content.type') <> 'text'
89 OR trim(coalesce(json_extract(chunk.value, '$.content.text'), '')) <> ''
90 )
91 )
92 ORDER BY position, stable_id",
93 )?;
94 let agent_message_latest_content_ordinals = ordinal_statement
95 .query_map([session_id], |row| row.get::<_, u64>(0))?
96 .collect::<rusqlite::Result<Vec<_>>>()?;
97 let restart_pattern = format!("{}*", mj_core::transcript::WORK_INTERRUPTED_ITEM_PREFIX);
98 let mut restart_statement = connection.prepare(
99 "SELECT position
100 FROM materialized_transcript_items
101 WHERE session_id = ?1 AND stable_id GLOB ?2
102 ORDER BY position, stable_id",
103 )?;
104 let mut interruption_event_ordinals = restart_statement
105 .query_map((session_id, restart_pattern), |row| row.get::<_, u64>(0))?
106 .collect::<rusqlite::Result<Vec<_>>>()?;
107 let outcome: Option<String> = connection.query_row(
108 "SELECT last_turn_outcome_json FROM materialized_sessions WHERE session_id = ?1",
109 [session_id],
110 |row| row.get(0),
111 )?;
112 if let Some(ordinal) = outcome
113 .map(|json| serde_json::from_str::<MaterializedTurnOutcome>(&json))
114 .transpose()?
115 .as_ref()
116 .and_then(MaterializedTurnOutcome::interruption_ordinal)
117 {
118 interruption_event_ordinals.push(ordinal);
119 }
120 interruption_event_ordinals.sort_unstable();
121 interruption_event_ordinals.dedup();
122
123 Ok(Some(MaterializedSessionSummary {
124 session_id: session_id.to_owned(),
125 applied_event_ordinal,
126 last_activity_at_ms,
127 execution: parse_materialized_execution(&execution, running_started_at_ms)?,
128 session_title,
129 last_agent_message: last_agent_message.map(|(_, message)| message),
130 last_user_message: last_user_message.map(|(_, message)| message),
131 last_agent_message_follows_last_user,
132 agent_message_latest_content_ordinals,
133 interruption_event_ordinals,
134 }))
135}
136
137pub(super) fn first_materialized_user_message(
141 connection: &Connection,
142 session_id: &str,
143) -> Result<Option<(u64, String)>> {
144 materialized_user_message(connection, session_id, true)
145}
146
147pub(super) fn last_materialized_user_message(
148 connection: &Connection,
149 session_id: &str,
150) -> Result<Option<(u64, String)>> {
151 materialized_user_message(connection, session_id, false)
152}
153
154pub(super) fn materialized_user_message(
155 connection: &Connection,
156 session_id: &str,
157 oldest_first: bool,
158) -> Result<Option<(u64, String)>> {
159 let mut statement = connection.prepare(if oldest_first {
160 "SELECT position, body_json
161 FROM materialized_transcript_items
162 WHERE session_id = ?1
163 AND json_extract(
164 CASE
165 WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
166 THEN body_json
167 ELSE '{}'
168 END,
169 '$.kind'
170 ) = 'user'
171 ORDER BY position, stable_id"
172 } else {
173 "SELECT position, body_json
174 FROM materialized_transcript_items
175 WHERE session_id = ?1
176 AND json_extract(
177 CASE
178 WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
179 THEN body_json
180 ELSE '{}'
181 END,
182 '$.kind'
183 ) = 'user'
184 ORDER BY position DESC, stable_id DESC"
185 })?;
186 let rows = statement.query_map([session_id], |row| {
187 Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?))
188 })?;
189 for row in rows {
190 let (position, body_json) = row?;
191 let body: TranscriptBody = serde_json::from_str(&body_json)
192 .with_context(|| format!("parse materialized user message for session {session_id}"))?;
193 let TranscriptBody::User { content } = body else {
194 continue;
195 };
196 let text = mj_core::transcript::materialized_content_text(&content);
197 if !text.trim().is_empty() {
198 return Ok(Some((position, text)));
199 }
200 }
201 Ok(None)
202}
203
204pub(super) fn last_materialized_turn_start(
208 connection: &Connection,
209 session_id: &str,
210) -> Result<Option<u64>> {
211 Ok(connection
212 .query_row(
213 "SELECT position
214 FROM materialized_transcript_items
215 WHERE session_id = ?1
216 AND (
217 stable_id GLOB ?2
218 OR json_extract(
219 CASE
220 WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
221 THEN body_json
222 ELSE '{}'
223 END,
224 '$.kind'
225 ) = 'user'
226 )
227 ORDER BY position DESC, stable_id DESC
228 LIMIT 1",
229 params![
230 session_id,
231 format!("{}*", mj_core::transcript::HARNESS_TURN_ITEM_PREFIX)
232 ],
233 |row| row.get::<_, u64>(0),
234 )
235 .optional()?)
236}
237
238pub(super) fn last_materialized_agent_message(
239 connection: &Connection,
240 session_id: &str,
241) -> Result<Option<(u64, String)>> {
242 last_materialized_agent_message_in(connection, session_id, 0, None)
243}
244
245pub(super) fn last_materialized_agent_message_within(
252 connection: &Connection,
253 session_id: &str,
254 after_position: u64,
255 through_position: u64,
256) -> Result<Option<String>> {
257 Ok(last_materialized_agent_message_in(
258 connection,
259 session_id,
260 after_position,
261 Some(through_position),
262 )?
263 .map(|(_, text)| text))
264}
265
266pub(super) fn last_materialized_agent_message_in(
267 connection: &Connection,
268 session_id: &str,
269 after_position: u64,
270 through_position: Option<u64>,
271) -> Result<Option<(u64, String)>> {
272 let row = connection
273 .query_row(
274 "SELECT position, body_json
275 FROM materialized_transcript_items
276 WHERE session_id = ?1
277 AND position > ?2
278 AND (?3 IS NULL OR position <= ?3)
279 AND latest_content_event_ordinal IS NOT NULL
280 AND EXISTS (
281 SELECT 1 FROM json_each(
282 CASE
283 WHEN latest_content_event_ordinal IS NOT NULL
284 AND json_valid(body_json)
285 THEN body_json
286 ELSE '{}'
287 END,
288 '$.chunks'
289 ) AS chunk
290 WHERE json_extract(chunk.value, '$.content.type') IS NOT NULL
291 AND (
292 json_extract(chunk.value, '$.content.type') <> 'text'
293 OR trim(coalesce(json_extract(chunk.value, '$.content.text'), '')) <> ''
294 )
295 )
296 ORDER BY position DESC, stable_id DESC
297 LIMIT 1",
298 params![session_id, after_position, through_position],
299 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
300 )
301 .optional()?;
302 let Some((position, body_json)) = row else {
303 return Ok(None);
304 };
305 let body: TranscriptBody = serde_json::from_str(&body_json)
306 .with_context(|| format!("parse materialized agent message for session {session_id}"))?;
307 let TranscriptBody::Agent { chunks, .. } = body else {
308 return Ok(None);
309 };
310 let text = mj_core::transcript::materialized_chunks_text(&chunks);
311 Ok((!text.trim().is_empty()).then_some((position, text)))
312}
313
314pub type MaterializedTurnState = (
316 MaterializedExecutionState,
317 Option<MaterializedTurn>,
318 Option<MaterializedTurnOutcome>,
319);
320
321pub fn load_materialized_turn_outcome(session_id: &str) -> Result<Option<MaterializedTurnState>> {
325 load_materialized_turn_outcome_from(&database_path(), session_id)
326}
327
328pub(super) fn load_materialized_turn_outcome_from(
329 path: &Path,
330 session_id: &str,
331) -> Result<Option<MaterializedTurnState>> {
332 let connection = open_reader(path)?;
333 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
334 return Ok(None);
335 };
336 Ok(Some((
337 fields.execution,
338 fields.active_turn,
339 fields.last_turn_outcome,
340 )))
341}
342
343pub fn load_materialized_finished_turn_message(session_id: &str) -> Result<Option<String>> {
354 load_materialized_finished_turn_message_from(&database_path(), session_id)
355}
356
357pub(super) fn load_materialized_finished_turn_message_from(
358 path: &Path,
359 session_id: &str,
360) -> Result<Option<String>> {
361 let mut reader = open_reader(path)?;
362 let connection = reader.transaction()?;
363 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
364 return Ok(None);
365 };
366 let Some(turn) = fields.last_turn_outcome else {
367 return Ok(None);
368 };
369 let Some(start_position) = turn.turn_start_position else {
370 return Ok(None);
371 };
372 last_materialized_agent_message_within(
373 &connection,
374 session_id,
375 start_position,
376 turn.completed_ordinal,
377 )
378}
379
380pub fn load_materialized_turn_summary(
388 session_id: &str,
389 turn_start_position: u64,
390 turn_completed_position: u64,
391) -> Result<TurnSummary> {
392 load_materialized_turn_summary_from(
393 &database_path(),
394 session_id,
395 turn_start_position,
396 turn_completed_position,
397 )
398}
399
400pub(super) fn load_materialized_turn_summary_from(
401 path: &Path,
402 session_id: &str,
403 turn_start_position: u64,
404 turn_completed_position: u64,
405) -> Result<TurnSummary> {
406 let mut reader = open_reader(path)?;
407 let connection = reader.transaction()?;
408 let turn_number = connection.query_row(
409 "SELECT COUNT(*)
410 FROM materialized_transcript_items
411 WHERE session_id = ?1
412 AND position <= ?3
413 AND (
414 stable_id GLOB ?2
415 OR json_extract(
416 CASE
417 WHEN stable_id GLOB 'user:*' OR stable_id GLOB 'user-*'
418 THEN body_json
419 ELSE '{}'
420 END,
421 '$.kind'
422 ) = 'user'
423 )",
424 params![
425 session_id,
426 format!("{}*", mj_core::transcript::HARNESS_TURN_ITEM_PREFIX),
427 turn_start_position
428 ],
429 |row| row.get::<_, u64>(0),
430 )?;
431 let (turn_started_at_ms, last_changed_at_ms) = connection.query_row(
432 "SELECT COALESCE(MIN(created_at_ms), 0), COALESCE(MAX(last_changed_at_ms), 0)
433 FROM materialized_transcript_items
434 WHERE session_id = ?1 AND position >= ?2 AND position <= ?3",
435 params![session_id, turn_start_position, turn_completed_position],
436 |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
437 )?;
438 let final_message = last_materialized_agent_message_within(
439 &connection,
440 session_id,
441 turn_start_position,
442 turn_completed_position,
443 )?;
444 Ok(TurnSummary {
445 turn_number,
446 turn_started_at_ms,
447 last_changed_at_ms,
448 final_message,
449 })
450}
451
452pub fn load_materialized_transcript_filtered(
453 session_id: &str,
454 after_seq: u64,
455 limit: usize,
456 role: Option<mj_core::transcript::TranscriptRole>,
457) -> Result<Option<TranscriptPage>> {
458 load_materialized_transcript_filtered_from(&database_path(), session_id, after_seq, limit, role)
459}
460
461pub(super) fn load_materialized_transcript_filtered_from(
462 path: &Path,
463 session_id: &str,
464 after_seq: u64,
465 limit: usize,
466 role: Option<mj_core::transcript::TranscriptRole>,
467) -> Result<Option<TranscriptPage>> {
468 let mut reader = open_reader(path)?;
469 let connection = reader.transaction()?;
470 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
471 return Ok(None);
472 };
473 let role = role.map(|r| r.storage_kind());
474 let mut statement = connection.prepare(
475 "WITH matches AS (
476 SELECT *, COALESCE(latest_content_event_ordinal, position) AS seq
477 FROM materialized_transcript_items WHERE session_id = ?1
478 AND COALESCE(latest_content_event_ordinal, position) > ?2
479 AND (?4 IS NULL OR json_extract(body_json, '$.kind') = ?4)
480 ), boundary AS (SELECT MAX(seq) AS seq FROM (SELECT seq FROM matches ORDER BY seq LIMIT ?3))
481 SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
482 last_changed_at_ms, body_json
483 FROM matches WHERE seq <= (SELECT seq FROM boundary)
484 ORDER BY seq, stable_id",
485 )?;
486 let rows = statement
487 .query_map(
488 params![session_id, after_seq, limit.clamp(1, 1000) as i64, role],
489 |row| {
490 Ok((
491 row.get::<_, String>(0)?,
492 row.get::<_, u64>(1)?,
493 row.get::<_, Option<u64>>(2)?,
494 row.get::<_, i64>(3)?,
495 row.get::<_, i64>(4)?,
496 row.get::<_, String>(5)?,
497 ))
498 },
499 )?
500 .collect::<rusqlite::Result<Vec<_>>>()?;
501 let items = rows
502 .into_iter()
503 .map(
504 |(
505 stable_id,
506 position,
507 latest_content_event_ordinal,
508 created_at_ms,
509 last_changed_at_ms,
510 body_json,
511 )| {
512 Ok(Arc::new(TranscriptItem {
513 stable_id,
514 position,
515 latest_content_event_ordinal,
516 created_at_ms,
517 last_changed_at_ms,
518 body: decode_transcript_body(&body_json, session_id)?,
519 }))
520 },
521 )
522 .collect::<Result<Vec<_>>>()?;
523 let latest_seq = connection.query_row(
524 "SELECT COALESCE(MAX(COALESCE(latest_content_event_ordinal, position)), 0)
525 FROM materialized_transcript_items
526 WHERE session_id = ?1",
527 [session_id],
528 |row| row.get::<_, u64>(0),
529 )?;
530 let last_seq = items.last().map_or(after_seq, |item| item.seq());
531 let more: bool = connection.query_row("SELECT EXISTS(SELECT 1 FROM materialized_transcript_items WHERE session_id = ?1 AND COALESCE(latest_content_event_ordinal, position) > ?2 AND (?3 IS NULL OR json_extract(body_json, '$.kind') = ?3))", params![session_id, last_seq, role], |r| r.get(0))?;
532 Ok(Some(TranscriptPage {
533 next_after_seq: if more {
534 last_seq
535 } else {
536 latest_seq.max(after_seq)
537 },
538 items,
539 latest_seq,
540 execution: fields.execution,
541 }))
542}
543
544pub(super) const RETENTION_BATCH_ITEMS: usize = 4_096;
551
552pub(super) const RETENTION_BODY_FLOOR_BYTES: usize = 4 * 1024;
555
556pub fn compact_materialized_transcript_through(
569 session_id: &str,
570 event_frontier: u64,
571) -> Result<TranscriptRetention> {
572 let session_id = session_id.to_owned();
573 submit_database_write("compact_materialized_transcript", move |_| {
574 compact_materialized_transcript_in(&database_path(), &session_id, event_frontier)
575 })
576}
577
578pub(super) fn compact_materialized_transcript_in(
579 path: &Path,
580 session_id: &str,
581 event_frontier: u64,
582) -> Result<TranscriptRetention> {
583 let mut connection = open(path)?;
584 let candidates = {
585 let mut statement = connection.prepare(
586 "SELECT stable_id, body_json
587 FROM materialized_transcript_items
588 WHERE session_id = ?1
589 AND position <= ?2
590 AND length(body_json) > ?3
591 AND json_extract(
592 CASE WHEN json_valid(body_json) THEN body_json ELSE '{}' END,
593 '$.kind'
594 ) = 'tool'
595 ORDER BY position, stable_id
596 LIMIT ?4",
597 )?;
598 statement
599 .query_map(
600 params![
601 session_id,
602 event_frontier,
603 RETENTION_BODY_FLOOR_BYTES as i64,
604 RETENTION_BATCH_ITEMS as i64 + 1
605 ],
606 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
607 )?
608 .collect::<rusqlite::Result<Vec<_>>>()?
609 };
610 let remaining = candidates.len() > RETENTION_BATCH_ITEMS;
611 let mut retention = TranscriptRetention {
612 remaining,
613 ..TranscriptRetention::default()
614 };
615 let transaction =
616 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
617 for (stable_id, body_json) in candidates.into_iter().take(RETENTION_BATCH_ITEMS) {
618 let mut body: TranscriptBody = match serde_json::from_str(&body_json) {
619 Ok(body) => body,
620 Err(error) => {
622 tracing::warn!(%session_id, %stable_id, %error, "skipping unreadable transcript body");
623 continue;
624 }
625 };
626 if !mj_transcript::transcript::compact_tool_call_for_retention(&mut body) {
627 continue;
628 }
629 let compacted = serde_json::to_string(&body)
630 .with_context(|| format!("serialize compacted transcript body {stable_id}"))?;
631 if compacted.len() >= body_json.len() {
632 continue;
633 }
634 transaction.execute(
635 "UPDATE materialized_transcript_items SET body_json = ?3
636 WHERE session_id = ?1 AND stable_id = ?2",
637 params![session_id, stable_id, compacted],
638 )?;
639 retention.items += 1;
640 retention.bytes += body_json.len() - compacted.len();
641 }
642 transaction.commit()?;
643 Ok(retention)
644}
645
646pub const PROJECTION_TAIL_ITEMS: usize = 1_024;
654
655pub fn load_materialized_actor_projection(
658 session_id: &str,
659) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
660 load_materialized_actor_projection_from(&database_path(), session_id)
661}
662
663pub(super) fn load_materialized_actor_projection_from(
664 path: &Path,
665 session_id: &str,
666) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
667 let mut reader = open_reader(path)?;
668 let connection = reader.transaction()?;
669 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
670 return Ok(None);
671 };
672 let latest_turn = last_materialized_turn_start(&connection, session_id)?;
673 let desired: Option<u64> = connection
674 .query_row(
675 "SELECT position FROM materialized_transcript_items WHERE session_id=?1
676 ORDER BY position DESC, stable_id DESC LIMIT 1 OFFSET ?2",
677 params![session_id, PROJECTION_TAIL_ITEMS - 1],
678 |row| row.get(0),
679 )
680 .optional()?;
681 let mutable: Option<u64> = connection.query_row(
684 "SELECT MIN(position) FROM materialized_transcript_items WHERE session_id=?1
685 AND (json_extract(body_json, '$.streaming')=1
686 OR (json_extract(body_json, '$.kind')='tool'
687 AND json_extract(body_json, '$.call.status') IN ('pending','in_progress')))",
688 [session_id],
689 |row| row.get(0),
690 )?;
691 let boundary = desired
692 .unwrap_or(0)
693 .min(latest_turn.unwrap_or(0))
694 .min(mutable.unwrap_or(u64::MAX));
695 let start: u64 = connection.query_row(
696 "SELECT COALESCE(MAX(position),0) FROM materialized_transcript_items
697 WHERE session_id=?1 AND position<=?2
698 AND (json_extract(body_json, '$.kind')='user' OR stable_id LIKE 'harness-turn:%')",
699 params![session_id, boundary],
700 |row| row.get(0),
701 )?;
702 let transcript = read_transcript_range(&connection, session_id, start, None, None)?;
703 let omitted_items = connection.query_row(
704 "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id=?1 AND position<?2",
705 params![session_id, start],
706 |row| row.get(0),
707 )?;
708 let window = ProjectionWindow {
709 omitted_items,
710 provisional_title: first_materialized_user_message(&connection, session_id)?
711 .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
712 latest_turn_start_position: latest_turn,
713 };
714 let materialized = MaterializedSession {
715 session_id: session_id.to_owned(),
716 applied_event_ordinal: fields.applied_event_ordinal,
717 applied_event_digest: fields.applied_event_digest,
718 last_activity_at_ms: fields.last_activity_at_ms,
719 execution: fields.execution,
720 session_title: fields.session_title,
721 configuration: fields.configuration,
722 transcript,
723 queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
724 pending_elicitations: fields.pending_elicitations,
725 active_turn: fields.active_turn,
726 last_turn_outcome: fields.last_turn_outcome,
727 };
728 materialized.validate()?;
729 Ok(Some((materialized, window)))
730}
731
732pub fn load_transcript_history(
733 session_id: &str,
734 before: Option<&mj_core::storage::TranscriptCursor>,
735 limit: usize,
736) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
737 load_transcript_history_from(&database_path(), session_id, before, limit)
738}
739
740pub(super) fn load_transcript_history_from(
741 path: &Path,
742 session_id: &str,
743 before: Option<&mj_core::storage::TranscriptCursor>,
744 limit: usize,
745) -> Result<Option<mj_core::storage::TranscriptHistoryPage>> {
746 let mut reader = open_reader(path)?;
747 let connection = reader.transaction()?;
748 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
749 return Ok(None);
750 };
751 let limit = limit.clamp(1, 256);
752 let mut items = read_transcript_range(&connection, session_id, 0, before, Some(limit + 1))?;
753 let has_more = items.len() > limit;
754 if has_more {
755 items.remove(0);
756 }
757 let before = has_more.then(|| mj_core::storage::TranscriptCursor::of(&items[0]));
758 Ok(Some(mj_core::storage::TranscriptHistoryPage {
759 items,
760 before,
761 frontier: fields.applied_event_ordinal,
762 }))
763}
764
765fn read_transcript_range(
766 connection: &Connection,
767 session_id: &str,
768 start: u64,
769 before: Option<&mj_core::storage::TranscriptCursor>,
770 limit: Option<usize>,
771) -> Result<Vec<Arc<TranscriptItem>>> {
772 let mut statement = connection.prepare(
773 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
774 last_changed_at_ms, body_json FROM materialized_transcript_items
775 WHERE session_id=?1 AND position>=?2
776 AND (position,stable_id)<(?3,?4)
777 ORDER BY position DESC, stable_id DESC LIMIT ?5",
778 )?;
779 let rows = statement
780 .query_map(
781 params![
782 session_id,
783 start,
784 before.map_or(i64::MAX as u64, |c| c.position),
785 before.map_or("", |c| c.stable_id.as_str()),
786 limit.map_or(-1, |limit| limit as i64)
787 ],
788 |row| {
789 Ok((
790 row.get::<_, String>(0)?,
791 row.get::<_, u64>(1)?,
792 row.get::<_, Option<u64>>(2)?,
793 row.get::<_, i64>(3)?,
794 row.get::<_, i64>(4)?,
795 row.get::<_, String>(5)?,
796 ))
797 },
798 )?
799 .collect::<rusqlite::Result<Vec<_>>>()?;
800 let mut items = rows
801 .into_iter()
802 .map(
803 |(
804 stable_id,
805 position,
806 latest_content_event_ordinal,
807 created_at_ms,
808 last_changed_at_ms,
809 body,
810 )| {
811 Ok(Arc::new(TranscriptItem {
812 stable_id,
813 position,
814 latest_content_event_ordinal,
815 created_at_ms,
816 last_changed_at_ms,
817 body: decode_transcript_body(&body, session_id)?,
818 }))
819 },
820 )
821 .collect::<Result<Vec<_>>>()?;
822 items.reverse();
823 Ok(items)
824}
825
826pub fn load_projection_references(
829 projection: &MaterializedSession,
830 events: &[mj_core::relay::RelayEvent],
831) -> Result<Vec<Arc<TranscriptItem>>> {
832 load_projection_references_from(&database_path(), projection, events)
833}
834
835pub(super) fn load_projection_references_from(
836 path: &Path,
837 projection: &MaterializedSession,
838 events: &[mj_core::relay::RelayEvent],
839) -> Result<Vec<Arc<TranscriptItem>>> {
840 let (mut ids, terminals) = mj_transcript::projection::historical_references(events)?;
841 let retained = projection
842 .transcript
843 .iter()
844 .map(|item| item.stable_id.as_str())
845 .collect::<std::collections::HashSet<_>>();
846 ids.retain(|id| !retained.contains(id.as_str()));
847 if ids.is_empty() && terminals.is_empty() {
848 return Ok(Vec::new());
849 }
850 let mut reader = open_reader(path)?;
851 let connection = reader.transaction()?;
852 let mut parameters = vec![projection.session_id.clone(), serde_json::to_string(&ids)?];
855 let mut statement = connection.prepare(if terminals.is_empty() {
856 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
857 last_changed_at_ms, body_json FROM materialized_transcript_items
858 WHERE session_id=?1 AND stable_id IN (SELECT value FROM json_each(?2))
859 ORDER BY position,stable_id"
860 } else {
861 parameters.push(serde_json::to_string(&retained)?);
862 parameters.push(serde_json::to_string(&terminals)?);
863 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
864 last_changed_at_ms, body_json FROM materialized_transcript_items
865 WHERE session_id=?1 AND stable_id NOT IN (SELECT value FROM json_each(?3))
866 AND (stable_id IN (SELECT value FROM json_each(?2))
867 OR EXISTS(SELECT 1 FROM json_each(body_json, '$.terminal_refs')
868 WHERE value IN (SELECT value FROM json_each(?4))))
869 ORDER BY position,stable_id"
870 })?;
871 let rows = statement
872 .query_map(rusqlite::params_from_iter(parameters), |row| {
873 Ok((
874 row.get::<_, String>(0)?,
875 row.get::<_, u64>(1)?,
876 row.get::<_, Option<u64>>(2)?,
877 row.get::<_, i64>(3)?,
878 row.get::<_, i64>(4)?,
879 row.get::<_, String>(5)?,
880 ))
881 })?
882 .collect::<rusqlite::Result<Vec<_>>>()?;
883 rows.into_iter()
884 .map(
885 |(
886 stable_id,
887 position,
888 latest_content_event_ordinal,
889 created_at_ms,
890 last_changed_at_ms,
891 body,
892 )| {
893 Ok(Arc::new(TranscriptItem {
894 stable_id,
895 position,
896 latest_content_event_ordinal,
897 created_at_ms,
898 last_changed_at_ms,
899 body: decode_transcript_body(&body, &projection.session_id)?,
900 }))
901 },
902 )
903 .collect()
904}
905
906pub fn load_materialized_projection_tail(
916 session_id: &str,
917 transcript_limit: usize,
918) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
919 load_materialized_projection_tail_from(&database_path(), session_id, transcript_limit)
920}
921
922pub(super) fn load_materialized_projection_tail_from(
923 path: &Path,
924 session_id: &str,
925 transcript_limit: usize,
926) -> Result<Option<(MaterializedSession, ProjectionWindow)>> {
927 let mut reader = open_reader(path)?;
928 let connection = reader.transaction()?;
931 let Some(fields) = read_materialized_session_fields(&connection, session_id)? else {
932 return Ok(None);
933 };
934 let transcript = read_materialized_transcript(&connection, session_id, Some(transcript_limit))?;
935 let total_items = connection.query_row(
936 "SELECT COUNT(*) FROM materialized_transcript_items WHERE session_id = ?1",
937 [session_id],
938 |row| row.get::<_, usize>(0),
939 )?;
940 let window = ProjectionWindow {
941 omitted_items: total_items.saturating_sub(transcript.len()),
942 provisional_title: first_materialized_user_message(&connection, session_id)?
943 .and_then(|(_, text)| mj_core::state::provisional_session_title(&text)),
944 latest_turn_start_position: last_materialized_turn_start(&connection, session_id)?,
945 };
946 let materialized = MaterializedSession {
947 session_id: session_id.to_owned(),
948 applied_event_ordinal: fields.applied_event_ordinal,
949 applied_event_digest: fields.applied_event_digest,
950 last_activity_at_ms: fields.last_activity_at_ms,
951 execution: fields.execution,
952 session_title: fields.session_title,
953 configuration: fields.configuration,
954 transcript,
955 queued_prompts: read_materialized_queued_prompts(&connection, session_id)?,
956 pending_elicitations: fields.pending_elicitations,
957 active_turn: fields.active_turn,
958 last_turn_outcome: fields.last_turn_outcome,
959 };
960 materialized.validate()?;
961 Ok(Some((materialized, window)))
962}
963
964pub fn materialized_event_frontier(session_id: &str) -> Result<Option<(u64, String)>> {
968 materialized_event_frontier_from(&database_path(), session_id)
969}
970
971pub(super) fn materialized_event_frontier_from(
972 path: &Path,
973 session_id: &str,
974) -> Result<Option<(u64, String)>> {
975 Ok(open_reader(path)?
976 .query_row(
977 "SELECT applied_event_ordinal, applied_event_digest
978 FROM materialized_sessions WHERE session_id = ?1",
979 [session_id],
980 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
981 )
982 .optional()?)
983}
984
985pub fn replace_materialized_queued_prompts(
989 session_id: &str,
990 queued_prompts: &[MaterializedQueuedPrompt],
991) -> Result<()> {
992 let session_id = session_id.to_owned();
993 let queued_prompts = queued_prompts.to_vec();
994 submit_database_write("replace_materialized_queued_prompts", move |_| {
995 replace_materialized_queued_prompts_in(&database_path(), &session_id, &queued_prompts)
996 })
997}
998
999pub(super) fn replace_materialized_queued_prompts_in(
1000 path: &Path,
1001 session_id: &str,
1002 queued_prompts: &[MaterializedQueuedPrompt],
1003) -> Result<()> {
1004 let mut connection = open(path)?;
1005 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1006 if !session_exists(&tx, session_id)? {
1007 bail!("unknown session {session_id}");
1008 }
1009 replace_materialized_queue(&tx, session_id, queued_prompts)?;
1010 tx.commit()?;
1011 Ok(())
1012}
1013
1014pub fn load_transcribed_session_activity() -> Result<BTreeMap<String, Option<i64>>> {
1021 load_transcribed_session_activity_from(&database_path())
1022}
1023
1024fn load_transcribed_session_activity_from(path: &Path) -> Result<BTreeMap<String, Option<i64>>> {
1025 let connection = open_reader(path)?;
1026 let mut statement = connection.prepare(
1027 "SELECT session_id, last_activity_at_ms
1028 FROM materialized_sessions s
1029 WHERE EXISTS (
1030 SELECT 1 FROM materialized_transcript_items i
1031 WHERE i.session_id = s.session_id
1032 )",
1033 )?;
1034 let rows = statement.query_map([], |row| {
1035 Ok((row.get::<_, String>(0)?, row.get::<_, Option<i64>>(1)?))
1036 })?;
1037 let mut activity = BTreeMap::new();
1038 for row in rows {
1039 let (session_id, last_activity_at_ms) = row?;
1040 activity.insert(session_id, last_activity_at_ms);
1041 }
1042 Ok(activity)
1043}
1044
1045pub fn load_materialized_queued_prompts() -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>>
1049{
1050 load_materialized_queued_prompts_from(&database_path())
1051}
1052
1053pub(super) fn load_materialized_queued_prompts_from(
1054 path: &Path,
1055) -> Result<BTreeMap<String, Vec<MaterializedQueuedPrompt>>> {
1056 let connection = open_reader(path)?;
1057 let mut statement = connection.prepare(
1058 "SELECT session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1059 FROM materialized_queued_prompts
1060 ORDER BY session_id, ordinal",
1061 )?;
1062 let rows = statement.query_map([], |row| {
1063 Ok((
1064 row.get::<_, String>(0)?,
1065 row.get::<_, String>(1)?,
1066 row.get::<_, String>(2)?,
1067 row.get::<_, String>(3)?,
1068 row.get::<_, i64>(4)?,
1069 row.get::<_, Option<u64>>(5)?,
1070 ))
1071 })?;
1072 let mut queues = BTreeMap::<String, Vec<MaterializedQueuedPrompt>>::new();
1073 for row in rows {
1074 let (session_id, command_id, kind_json, content_json, queued_at_ms, accepted_ordinal) =
1075 row?;
1076 let content = serde_json::from_str(&content_json).with_context(|| {
1077 format!("parse materialized queued prompt for session {session_id}")
1078 })?;
1079 let kind = serde_json::from_str(&kind_json).with_context(|| {
1080 format!("parse materialized queue entry kind for session {session_id}")
1081 })?;
1082 queues
1083 .entry(session_id)
1084 .or_default()
1085 .push(MaterializedQueuedPrompt {
1086 command_id,
1087 kind,
1088 content,
1089 queued_at_ms,
1090 accepted_ordinal,
1091 });
1092 }
1093 Ok(queues)
1094}
1095
1096pub(super) fn load_materialized_session_from(
1097 path: &Path,
1098 session_id: &str,
1099) -> Result<Option<MaterializedSession>> {
1100 let mut reader = open_reader(path)?;
1101 let connection = reader.transaction()?;
1102 load_materialized_session_with(&connection, session_id)
1103}
1104
1105pub(super) fn load_materialized_session_with(
1106 connection: &rusqlite::Transaction<'_>,
1107 session_id: &str,
1108) -> Result<Option<MaterializedSession>> {
1109 let Some(fields) = read_materialized_session_fields(connection, session_id)? else {
1110 return Ok(None);
1111 };
1112 let materialized = MaterializedSession {
1113 session_id: session_id.to_owned(),
1114 applied_event_ordinal: fields.applied_event_ordinal,
1115 applied_event_digest: fields.applied_event_digest,
1116 last_activity_at_ms: fields.last_activity_at_ms,
1117 execution: fields.execution,
1118 session_title: fields.session_title,
1119 configuration: fields.configuration,
1120 transcript: read_materialized_transcript(connection, session_id, None)?,
1121 queued_prompts: read_materialized_queued_prompts(connection, session_id)?,
1122 pending_elicitations: fields.pending_elicitations,
1123 active_turn: fields.active_turn,
1124 last_turn_outcome: fields.last_turn_outcome,
1125 };
1126 materialized.validate()?;
1127 Ok(Some(materialized))
1128}
1129
1130pub(super) struct MaterializedSessionFields {
1132 pub(super) applied_event_ordinal: u64,
1133 pub(super) applied_event_digest: String,
1134 pub(super) last_activity_at_ms: Option<i64>,
1135 pub(super) execution: MaterializedExecutionState,
1136 pub(super) session_title: Option<String>,
1137 pub(super) configuration: BTreeMap<String, serde_json::Value>,
1138 pub(super) pending_elicitations: Vec<mj_core::elicitation::ElicitationRequest>,
1139 pub(super) active_turn: Option<MaterializedTurn>,
1140 pub(super) last_turn_outcome: Option<MaterializedTurnOutcome>,
1141}
1142
1143pub(super) fn read_materialized_session_fields(
1144 connection: &Connection,
1145 session_id: &str,
1146) -> Result<Option<MaterializedSessionFields>> {
1147 let row = connection
1148 .query_row(
1149 "SELECT applied_event_ordinal, applied_event_digest, last_activity_at_ms,
1150 execution_state, running_started_at_ms, session_title, configuration_json,
1151 pending_elicitations_json, active_turn_json, last_turn_outcome_json
1152 FROM materialized_sessions WHERE session_id = ?1",
1153 [session_id],
1154 |row| {
1155 Ok((
1156 row.get::<_, u64>(0)?,
1157 row.get::<_, String>(1)?,
1158 row.get::<_, Option<i64>>(2)?,
1159 row.get::<_, String>(3)?,
1160 row.get::<_, Option<i64>>(4)?,
1161 row.get::<_, Option<String>>(5)?,
1162 row.get::<_, String>(6)?,
1163 row.get::<_, String>(7)?,
1164 row.get::<_, Option<String>>(8)?,
1165 row.get::<_, Option<String>>(9)?,
1166 ))
1167 },
1168 )
1169 .optional()?;
1170 let Some((
1171 applied_event_ordinal,
1172 applied_event_digest,
1173 last_activity_at_ms,
1174 execution,
1175 running_started_at_ms,
1176 session_title,
1177 configuration_json,
1178 pending_elicitations_json,
1179 active_turn_json,
1180 last_turn_outcome_json,
1181 )) = row
1182 else {
1183 return Ok(None);
1184 };
1185 #[cfg(test)]
1186 super::tests::after_materialized_frontier_read();
1187 Ok(Some(MaterializedSessionFields {
1188 applied_event_ordinal,
1189 applied_event_digest,
1190 last_activity_at_ms,
1191 execution: parse_materialized_execution(&execution, running_started_at_ms)?,
1192 session_title,
1193 configuration: serde_json::from_str(&configuration_json).with_context(|| {
1194 format!("parse materialized configuration for session {session_id}")
1195 })?,
1196 pending_elicitations: serde_json::from_str(&pending_elicitations_json)
1197 .with_context(|| format!("parse pending elicitations for session {session_id}"))?,
1198 active_turn: active_turn_json
1199 .as_deref()
1200 .map(serde_json::from_str)
1201 .transpose()
1202 .with_context(|| format!("parse active turn for session {session_id}"))?,
1203 last_turn_outcome: last_turn_outcome_json
1204 .as_deref()
1205 .map(serde_json::from_str)
1206 .transpose()
1207 .with_context(|| format!("parse last turn outcome for session {session_id}"))?,
1208 }))
1209}
1210
1211pub(super) fn read_materialized_transcript(
1215 connection: &Connection,
1216 session_id: &str,
1217 limit: Option<usize>,
1218) -> Result<Vec<Arc<TranscriptItem>>> {
1219 let mut statement = connection.prepare(match limit {
1220 Some(_) => {
1221 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1222 last_changed_at_ms, body_json
1223 FROM materialized_transcript_items
1224 WHERE session_id = ?1
1225 ORDER BY position DESC, stable_id DESC
1226 LIMIT ?2"
1227 }
1228 None => {
1229 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1230 last_changed_at_ms, body_json
1231 FROM materialized_transcript_items
1232 WHERE session_id = ?1
1233 ORDER BY position, stable_id"
1234 }
1235 })?;
1236 let read = |row: &rusqlite::Row<'_>| {
1237 Ok((
1238 row.get::<_, String>(0)?,
1239 row.get::<_, u64>(1)?,
1240 row.get::<_, Option<u64>>(2)?,
1241 row.get::<_, i64>(3)?,
1242 row.get::<_, i64>(4)?,
1243 row.get::<_, String>(5)?,
1244 ))
1245 };
1246 let rows = match limit {
1247 Some(limit) => statement
1248 .query_map(params![session_id, limit as i64], read)?
1249 .collect::<rusqlite::Result<Vec<_>>>()?,
1250 None => statement
1251 .query_map([session_id], read)?
1252 .collect::<rusqlite::Result<Vec<_>>>()?,
1253 };
1254 let mut transcript = rows
1255 .into_iter()
1256 .map(
1257 |(
1258 stable_id,
1259 position,
1260 latest_content_event_ordinal,
1261 created_at_ms,
1262 last_changed_at_ms,
1263 body_json,
1264 )| {
1265 Ok(Arc::new(TranscriptItem {
1266 stable_id,
1267 position,
1268 latest_content_event_ordinal,
1269 created_at_ms,
1270 last_changed_at_ms,
1271 body: decode_transcript_body(&body_json, session_id)?,
1272 }))
1273 },
1274 )
1275 .collect::<Result<Vec<_>>>()?;
1276 if limit.is_some() {
1277 transcript.reverse();
1280 }
1281 Ok(transcript)
1282}
1283
1284pub(super) fn read_materialized_queued_prompts(
1285 connection: &Connection,
1286 session_id: &str,
1287) -> Result<Vec<MaterializedQueuedPrompt>> {
1288 let mut statement = connection.prepare(
1289 "SELECT command_id, kind_json, content_json, queued_at_ms, accepted_ordinal
1290 FROM materialized_queued_prompts
1291 WHERE session_id = ?1
1292 ORDER BY ordinal",
1293 )?;
1294 let rows = statement
1295 .query_map([session_id], |row| {
1296 Ok((
1297 row.get::<_, String>(0)?,
1298 row.get::<_, String>(1)?,
1299 row.get::<_, String>(2)?,
1300 row.get::<_, i64>(3)?,
1301 row.get::<_, Option<u64>>(4)?,
1302 ))
1303 })?
1304 .collect::<rusqlite::Result<Vec<_>>>()?;
1305 rows.into_iter()
1306 .map(
1307 |(command_id, kind_json, content_json, queued_at_ms, accepted_ordinal)| {
1308 Ok(MaterializedQueuedPrompt {
1309 command_id,
1310 kind: serde_json::from_str(&kind_json).with_context(|| {
1311 format!("parse materialized queue entry kind for session {session_id}")
1312 })?,
1313 content: serde_json::from_str(&content_json).with_context(|| {
1314 format!("parse materialized queued prompt for session {session_id}")
1315 })?,
1316 queued_at_ms,
1317 accepted_ordinal,
1318 })
1319 },
1320 )
1321 .collect()
1322}
1323
1324pub fn save_materialized_session(materialized: &MaterializedSession) -> Result<()> {
1328 let materialized = materialized.clone();
1329 submit_database_write("save_materialized_session", move |_| {
1330 save_materialized_session_to(&database_path(), &materialized)
1331 })
1332}
1333
1334pub(super) fn save_materialized_session_to(
1335 path: &Path,
1336 materialized: &MaterializedSession,
1337) -> Result<()> {
1338 materialized.validate()?;
1339 let mut connection = open(path)?;
1340 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1341 if !session_exists(&tx, &materialized.session_id)? {
1342 bail!("unknown session {}", materialized.session_id);
1343 }
1344 write_materialized_session(&tx, materialized)?;
1345 tx.commit()?;
1346 Ok(())
1347}
1348
1349pub struct ProjectionPage<'a> {
1354 pub(super) session_id: &'a str,
1355 pub(super) transaction: Transaction<'a>,
1356 pub(super) applied_ordinal: u64,
1357 pub(super) applied_digest: String,
1358 pub(super) dirty: bool,
1359 pub(super) pending: MaterializedSessionMutation,
1360 pub(super) pending_transcript: BTreeMap<String, PendingTranscriptMutation>,
1361 pub(super) pending_turns: Vec<MaterializedTurnOutcome>,
1362 pub(super) pending_events: Vec<(i64, ApiEventData)>,
1363}
1364
1365pub(super) struct PendingTranscriptMutation {
1366 pub(super) final_mutation: TranscriptMutation,
1367 pub(super) remove_before_upsert: bool,
1368}
1369
1370impl ProjectionPage<'_> {
1371 pub fn apply(
1375 &mut self,
1376 event_ordinal: u64,
1377 previous_event_digest: &str,
1378 event_digest: &str,
1379 mutation: &MaterializedSessionMutation,
1380 ) -> Result<ProjectionApplyOutcome> {
1381 if event_ordinal == 0 {
1382 bail!("relay event ordinal must be positive");
1383 }
1384 let chained = !previous_event_digest.is_empty();
1390 if chained {
1391 validate_relay_event_digest(previous_event_digest, "previous relay event digest")?;
1392 }
1393 validate_relay_event_frontier(event_ordinal, event_digest, "relay event frontier")?;
1394 let session_id = self.session_id;
1395 let applied = self.applied_ordinal;
1396 if event_ordinal < applied {
1397 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1398 }
1399 if event_ordinal == applied {
1400 if event_digest != self.applied_digest {
1401 bail!(
1402 "relay event digest mismatch for session {session_id} at ordinal {event_ordinal}: projection has {}, received {event_digest}",
1403 self.applied_digest
1404 );
1405 }
1406 return Ok(ProjectionApplyOutcome::AlreadyApplied);
1407 }
1408 let expected = applied
1409 .checked_add(1)
1410 .context("materialized event ordinal overflow")?;
1411 if event_ordinal != expected {
1412 bail!(
1413 "relay event gap for session {session_id}: expected ordinal {expected}, received {event_ordinal}"
1414 );
1415 }
1416 if chained && previous_event_digest != self.applied_digest {
1417 bail!(
1418 "relay event chain diverged for session {session_id} before ordinal {event_ordinal}: projection has {}, event follows {previous_event_digest}",
1419 self.applied_digest
1420 );
1421 }
1422
1423 if let Some(event) = &mutation.native_agent {
1424 native_agents::apply_native_agent_event(&self.transaction, session_id, event)?;
1425 }
1426 if let Some(activity_at_ms) = mutation.last_activity_at_ms {
1427 self.pending.last_activity_at_ms = Some(
1428 self.pending
1429 .last_activity_at_ms
1430 .map_or(activity_at_ms, |existing| existing.max(activity_at_ms)),
1431 );
1432 }
1433 if let Some(execution) = mutation.execution {
1434 self.pending.execution = Some(execution);
1435 }
1436 if let Some(title) = &mutation.session_title {
1437 if title.as_ref().is_some_and(|title| title.trim().is_empty()) {
1438 bail!("materialized session title cannot be empty");
1439 }
1440 self.pending.session_title = Some(title.clone());
1441 }
1442 if let Some(configuration) = &mutation.configuration {
1443 self.pending.configuration = Some(configuration.clone());
1444 }
1445 for item_mutation in &mutation.transcript {
1446 match item_mutation {
1447 TranscriptMutation::Upsert(item) => {
1448 item.validate(event_ordinal)?;
1449 let stable_id = item.stable_id.clone();
1450 let entry = self.pending_transcript.entry(stable_id).or_insert_with(|| {
1451 PendingTranscriptMutation {
1452 final_mutation: TranscriptMutation::Upsert(item.clone()),
1453 remove_before_upsert: false,
1454 }
1455 });
1456 entry.remove_before_upsert |=
1457 matches!(&entry.final_mutation, TranscriptMutation::Remove { .. });
1458 entry.final_mutation = TranscriptMutation::Upsert(item.clone());
1459 }
1460 TranscriptMutation::Remove { stable_id } => {
1461 if stable_id.trim().is_empty() {
1462 bail!("cannot remove a transcript item with an empty stable id");
1463 }
1464 let removed = TranscriptMutation::Remove {
1465 stable_id: stable_id.clone(),
1466 };
1467 self.pending_transcript
1468 .entry(stable_id.clone())
1469 .and_modify(|entry| entry.final_mutation = removed.clone())
1470 .or_insert(PendingTranscriptMutation {
1471 final_mutation: removed,
1472 remove_before_upsert: false,
1473 });
1474 }
1475 }
1476 }
1477 if let Some(queued_prompts) = &mutation.queued_prompts {
1478 self.pending.queued_prompts = Some(queued_prompts.clone());
1479 }
1480 if let Some(pending_elicitations) = &mutation.pending_elicitations {
1481 self.pending.pending_elicitations = Some(pending_elicitations.clone());
1482 }
1483 self.pending
1484 .config_results
1485 .extend(mutation.config_results.clone());
1486 if let Some(active_turn) = &mutation.active_turn {
1487 self.pending.active_turn = Some(active_turn.clone());
1488 }
1489 if mutation.clear_turn_outcome {
1490 self.pending.clear_turn_outcome = true;
1491 self.pending.last_turn_outcome = None;
1492 }
1493 if let Some(last_turn_outcome) = &mutation.last_turn_outcome {
1494 self.pending_turns.push(last_turn_outcome.clone());
1495 self.pending.last_turn_outcome = Some(last_turn_outcome.clone());
1496 }
1497 if let Some(cost) = &mutation.provider_cost {
1498 self.pending.provider_cost = Some(cost.clone());
1499 }
1500 self.pending_events.extend(
1501 mutation
1502 .api_events
1503 .iter()
1504 .cloned()
1505 .map(|event| (mutation.last_activity_at_ms.unwrap_or(0), event)),
1506 );
1507 self.applied_ordinal = event_ordinal;
1508 event_digest.clone_into(&mut self.applied_digest);
1509 self.dirty = true;
1510 Ok(ProjectionApplyOutcome::Applied)
1511 }
1512
1513 pub(super) fn flush(&mut self) -> Result<()> {
1517 if !self.dirty {
1518 return Ok(());
1519 }
1520 let tx = &self.transaction;
1521 let session_id = self.session_id;
1522 if let Some(execution) = self.pending.execution {
1523 let (state, started_at_ms) = materialized_execution_columns(execution);
1524 tx.execute(
1525 "UPDATE materialized_sessions
1526 SET execution_state = ?2, running_started_at_ms = ?3
1527 WHERE session_id = ?1",
1528 params![session_id, state, started_at_ms],
1529 )?;
1530 }
1531 if let Some(title) = &self.pending.session_title {
1532 tx.execute(
1533 "UPDATE materialized_sessions SET session_title = ?2 WHERE session_id = ?1",
1534 params![session_id, title],
1535 )?;
1536 }
1537 if let Some(configuration) = &self.pending.configuration {
1538 tx.execute(
1539 "UPDATE materialized_sessions SET configuration_json = ?2 WHERE session_id = ?1",
1540 params![session_id, serde_json::to_string(configuration)?],
1541 )?;
1542 }
1543 for pending in self.pending_transcript.values() {
1544 match &pending.final_mutation {
1545 TranscriptMutation::Upsert(item) => {
1546 if pending.remove_before_upsert {
1550 tx.execute(
1551 "DELETE FROM materialized_transcript_items
1552 WHERE session_id = ?1 AND stable_id = ?2",
1553 params![session_id, item.stable_id],
1554 )?;
1555 }
1556 upsert_transcript_item(tx, session_id, item)?;
1557 }
1558 TranscriptMutation::Remove { stable_id } => {
1559 tx.execute(
1560 "DELETE FROM materialized_transcript_items
1561 WHERE session_id = ?1 AND stable_id = ?2",
1562 params![session_id, stable_id],
1563 )?;
1564 }
1565 }
1566 }
1567 if let Some(queued_prompts) = &self.pending.queued_prompts {
1568 replace_materialized_queue(tx, session_id, queued_prompts)?;
1569 }
1570 if let Some(pending_elicitations) = &self.pending.pending_elicitations {
1571 tx.execute(
1572 "UPDATE materialized_sessions
1573 SET pending_elicitations_json = ?2 WHERE session_id = ?1",
1574 params![session_id, serde_json::to_string(pending_elicitations)?],
1575 )?;
1576 }
1577 for (recorded_at_ms, event) in &self.pending_events {
1578 events::insert_api_event(tx, session_id, *recorded_at_ms, event)?;
1579 }
1580 for turn in &self.pending_turns {
1581 tx.execute("INSERT OR REPLACE INTO session_turn_usage(session_id, command_id, completed_ordinal, turn_start_position, body) VALUES (?1, ?2, ?3, ?4, ?5)", params![session_id, turn.command_id, turn.completed_ordinal, turn.turn_start_position, serde_json::to_string(turn)?])?;
1582 }
1583 if let Some(cost) = &self.pending.provider_cost {
1584 tx.execute(
1585 "INSERT OR REPLACE INTO session_provider_cost(session_id, body) VALUES (?1, ?2)",
1586 params![session_id, serde_json::to_string(cost)?],
1587 )?;
1588 }
1589 for (command_id, error) in &self.pending.config_results {
1590 tx.execute("INSERT OR REPLACE INTO api_config_results(session_id, command_id, error) VALUES (?1, ?2, ?3)", params![session_id, command_id, error])?;
1591 }
1592 if let Some(active_turn) = &self.pending.active_turn {
1593 tx.execute(
1594 "UPDATE materialized_sessions SET active_turn_json = ?2 WHERE session_id = ?1",
1595 params![
1596 session_id,
1597 active_turn
1598 .as_ref()
1599 .map(serde_json::to_string)
1600 .transpose()?
1601 ],
1602 )?;
1603 }
1604 if self.pending.clear_turn_outcome {
1605 self.transaction.execute("UPDATE materialized_sessions SET last_turn_outcome_json = NULL WHERE session_id = ?1", [session_id])?;
1606 }
1607 if let Some(last_turn_outcome) = &self.pending.last_turn_outcome {
1608 tx.execute(
1609 "UPDATE materialized_sessions
1610 SET last_turn_outcome_json = ?2 WHERE session_id = ?1",
1611 params![session_id, serde_json::to_string(last_turn_outcome)?],
1612 )?;
1613 }
1614 tx.execute(
1615 "UPDATE materialized_sessions
1616 SET last_activity_at_ms = CASE
1617 WHEN ?2 IS NULL THEN last_activity_at_ms
1618 WHEN last_activity_at_ms IS NULL OR last_activity_at_ms < ?2 THEN ?2
1619 ELSE last_activity_at_ms
1620 END,
1621 applied_event_ordinal = ?3,
1622 applied_event_digest = ?4
1623 WHERE session_id = ?1",
1624 params![
1625 session_id,
1626 self.pending.last_activity_at_ms,
1627 self.applied_ordinal,
1628 self.applied_digest,
1629 ],
1630 )?;
1631 Ok(())
1632 }
1633}
1634
1635pub fn apply_projection_page<T>(
1640 session_id: &str,
1641 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T> + Send + 'static,
1642) -> Result<T>
1643where
1644 T: Send + 'static,
1645{
1646 let session_id = session_id.to_owned();
1647 submit_database_write("apply_projection_page", move |connection| {
1648 apply_projection_page_with(connection, &session_id, fill)
1649 })
1650}
1651
1652#[cfg(test)]
1653pub(super) fn apply_projection_page_to<T>(
1654 path: &Path,
1655 session_id: &str,
1656 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1657) -> Result<T> {
1658 let mut connection = open(path)?;
1659 apply_projection_page_with(&mut connection, session_id, fill)
1660}
1661
1662pub(super) fn apply_projection_page_with<T>(
1663 connection: &mut Connection,
1664 session_id: &str,
1665 fill: impl FnOnce(&mut ProjectionPage<'_>) -> Result<T>,
1666) -> Result<T> {
1667 let transaction =
1668 connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1669 let (applied_ordinal, applied_digest) = transaction
1670 .query_row(
1671 "SELECT applied_event_ordinal, applied_event_digest
1672 FROM materialized_sessions WHERE session_id = ?1",
1673 [session_id],
1674 |row| Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?)),
1675 )
1676 .optional()?
1677 .with_context(|| format!("unknown session {session_id}"))?;
1678 validate_relay_event_frontier(
1679 applied_ordinal,
1680 &applied_digest,
1681 "persisted relay event frontier",
1682 )?;
1683 let mut page = ProjectionPage {
1684 session_id,
1685 transaction,
1686 applied_ordinal,
1687 applied_digest,
1688 dirty: false,
1689 pending: MaterializedSessionMutation::default(),
1690 pending_transcript: BTreeMap::new(),
1691 pending_turns: Vec::new(),
1692 pending_events: Vec::new(),
1693 };
1694 let filled = fill(&mut page)?;
1697 page.flush()?;
1698 page.transaction.commit()?;
1699 Ok(filled)
1700}
1701
1702pub fn apply_projection_event(
1704 session_id: &str,
1705 event_ordinal: u64,
1706 previous_event_digest: &str,
1707 event_digest: &str,
1708 mutation: &MaterializedSessionMutation,
1709) -> Result<ProjectionApplyOutcome> {
1710 let session_id = session_id.to_owned();
1711 let previous_event_digest = previous_event_digest.to_owned();
1712 let event_digest = event_digest.to_owned();
1713 let mutation = mutation.clone();
1714 submit_database_write("apply_projection_event", move |connection| {
1715 apply_projection_page_with(connection, &session_id, |page| {
1716 page.apply(
1717 event_ordinal,
1718 &previous_event_digest,
1719 &event_digest,
1720 &mutation,
1721 )
1722 })
1723 })
1724}
1725
1726#[cfg(test)]
1727pub(super) fn apply_projection_event_to(
1728 path: &Path,
1729 session_id: &str,
1730 event_ordinal: u64,
1731 previous_event_digest: &str,
1732 event_digest: &str,
1733 mutation: &MaterializedSessionMutation,
1734) -> Result<ProjectionApplyOutcome> {
1735 apply_projection_page_to(path, session_id, |page| {
1736 page.apply(event_ordinal, previous_event_digest, event_digest, mutation)
1737 })
1738}
1739
1740pub fn advance_viewed_through_event_ordinal(session_id: &str, through: u64) -> Result<u64> {
1743 let session_id = session_id.to_owned();
1744 submit_database_write("advance_viewed_through_event_ordinal", move |_| {
1745 advance_viewed_through_event_ordinal_to(&database_path(), &session_id, through)
1746 })
1747}
1748
1749pub(super) fn advance_viewed_through_event_ordinal_to(
1750 path: &Path,
1751 session_id: &str,
1752 through: u64,
1753) -> Result<u64> {
1754 let mut connection = open(path)?;
1755 let tx = connection.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
1756 let applied = tx
1757 .query_row(
1758 "SELECT applied_event_ordinal FROM materialized_sessions WHERE session_id = ?1",
1759 [session_id],
1760 |row| row.get::<_, u64>(0),
1761 )
1762 .optional()?
1763 .with_context(|| format!("unknown session {session_id}"))?;
1764 if through > applied {
1765 bail!(
1766 "cannot acknowledge event ordinal {through} for session {session_id}; projection is at {applied}"
1767 );
1768 }
1769 tx.execute(
1770 "UPDATE sessions
1771 SET viewed_through_event_ordinal = max(viewed_through_event_ordinal, ?2)
1772 WHERE session_id = ?1",
1773 params![session_id, through],
1774 )?;
1775 let receipt = tx.query_row(
1776 "SELECT viewed_through_event_ordinal FROM sessions WHERE session_id = ?1",
1777 [session_id],
1778 |row| row.get::<_, u64>(0),
1779 )?;
1780 tx.commit()?;
1781 Ok(receipt)
1782}
1783
1784fn decode_transcript_body(body_json: &str, session_id: &str) -> Result<TranscriptBody> {
1790 let mut body: TranscriptBody = serde_json::from_str(body_json)
1791 .with_context(|| format!("parse materialized transcript body for session {session_id}"))?;
1792 match &mut body {
1793 TranscriptBody::Agent { chunks, .. } | TranscriptBody::Thought { chunks, .. } => {
1794 mj_core::transcript::coalesce_content_chunks(chunks);
1795 }
1796 _ => {}
1797 }
1798 Ok(body)
1799}
1800
1801pub(crate) fn load_continuation_evidence(
1804 session_id: &str,
1805 ordinal: u64,
1806 digest: &str,
1807) -> Result<mj_core::continuation::ContinuationEvidence> {
1808 load_continuation_evidence_from(&database_path(), session_id, ordinal, digest)
1809}
1810
1811pub(super) fn load_continuation_evidence_from(
1812 path: &Path,
1813 session_id: &str,
1814 ordinal: u64,
1815 digest: &str,
1816) -> Result<mj_core::continuation::ContinuationEvidence> {
1817 let mut reader = open_reader(path)?;
1818 let connection = reader.transaction()?;
1819 let fields = read_materialized_session_fields(&connection, session_id)?
1820 .context("continuation projection is missing")?;
1821 anyhow::ensure!(
1822 fields.applied_event_ordinal == ordinal && fields.applied_event_digest == digest,
1823 "continuation projection changed"
1824 );
1825 let mut statement = connection.prepare(
1826 "SELECT stable_id, position, latest_content_event_ordinal, created_at_ms,
1827 last_changed_at_ms, body_json
1828 FROM materialized_transcript_items WHERE session_id = ?1
1829 AND (json_extract(body_json, '$.kind') IN ('user', 'agent')
1830 OR stable_id LIKE 'context-cleared:%')
1831 ORDER BY position DESC, stable_id DESC",
1832 )?;
1833 let rows = statement.query_map([session_id], |row| {
1834 Ok((
1835 row.get::<_, String>(0)?,
1836 row.get::<_, u64>(1)?,
1837 row.get::<_, Option<u64>>(2)?,
1838 row.get::<_, i64>(3)?,
1839 row.get::<_, i64>(4)?,
1840 row.get::<_, String>(5)?,
1841 ))
1842 })?;
1843 crate::continuation::evidence_from_items(rows.map(|row| {
1844 let (
1845 stable_id,
1846 position,
1847 latest_content_event_ordinal,
1848 created_at_ms,
1849 last_changed_at_ms,
1850 body_json,
1851 ) = row?;
1852 anyhow::ensure!(
1854 body_json.len() <= 1024 * 1024,
1855 "continuation source message exceeds budget"
1856 );
1857 Ok(Arc::new(TranscriptItem {
1858 stable_id,
1859 position,
1860 latest_content_event_ordinal,
1861 created_at_ms,
1862 last_changed_at_ms,
1863 body: decode_transcript_body(&body_json, session_id)?,
1864 }))
1865 }))
1866}