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