1use super::*;
3use mj_core::event_outcome::{CommandOwner, CommandResult, CommandResultKind, OutcomeReason};
4
5#[cfg(test)]
6use mj_core::elicitation::ElicitationRequest;
7
8pub fn load_api_events(
9 filter: &ApiEventFilter,
10 after_seq: Option<u64>,
11 limit: usize,
12) -> Result<ApiEventPage> {
13 load_api_events_from(&database_path(), filter, after_seq, limit)
14}
15
16pub(super) fn load_api_events_from(
17 path: &Path,
18 filter: &ApiEventFilter,
19 after_seq: Option<u64>,
20 limit: usize,
21) -> Result<ApiEventPage> {
22 let mut connection = open_reader(path)?;
23 let tx = connection.transaction()?;
24 let latest_seq: u64 = tx.query_row(
26 "SELECT COALESCE((SELECT seq FROM sqlite_sequence WHERE name = 'api_events'), 0)",
27 [],
28 |r| r.get(0),
29 )?;
30 let after_seq = after_seq.unwrap_or(latest_seq);
31 let mut statement = tx.prepare("SELECT e.seq, e.session_id, e.recorded_at_ms, e.body FROM api_events e JOIN session_contexts w ON w.session_id = e.session_id WHERE e.seq > ?1 AND (?2 IS NULL OR e.session_id = ?2) AND (?3 IS NULL OR w.workspace_id = ?3) ORDER BY e.seq LIMIT ?4")?;
32 let events = statement
33 .query_map(
34 params![
35 after_seq,
36 filter.session_id,
37 filter.workspace_id,
38 limit.clamp(1, 1000) as i64
39 ],
40 |r| {
41 Ok((
42 r.get::<_, u64>(0)?,
43 r.get::<_, String>(1)?,
44 r.get::<_, i64>(2)?,
45 r.get::<_, String>(3)?,
46 ))
47 },
48 )?
49 .map(|r| {
50 let (seq, session_id, recorded_at_ms, body) = r?;
51 Ok(ApiEvent {
52 seq,
53 session_id,
54 recorded_at_ms,
55 event: serde_json::from_str(&body)?,
56 })
57 })
58 .collect::<Result<Vec<_>>>()?;
59 let next_after_seq = events
60 .last()
61 .map_or(latest_seq.max(after_seq), |event| event.seq);
62 Ok(ApiEventPage {
63 events,
64 next_after_seq,
65 latest_seq,
66 })
67}
68
69pub(super) fn insert_api_event(
70 tx: &Transaction<'_>,
71 session_id: &str,
72 recorded_at_ms: i64,
73 event: &ApiEventData,
74) -> Result<()> {
75 tx.execute(
76 "INSERT INTO api_events(session_id, recorded_at_ms, body) VALUES (?1, ?2, ?3)",
77 params![session_id, recorded_at_ms, serde_json::to_string(event)?],
78 )?;
79 Ok(())
80}
81
82pub fn record_api_activities(
84 activities: Vec<(String, ApiActivityState)>,
85 recorded_at_ms: i64,
86) -> Result<()> {
87 submit_database_write("record_api_activities", move |connection| {
88 record_api_activities_with(connection, activities, recorded_at_ms)
89 })
90}
91
92pub(super) fn record_api_activities_with(
93 connection: &mut Connection,
94 activities: Vec<(String, ApiActivityState)>,
95 recorded_at_ms: i64,
96) -> Result<()> {
97 let tx = connection.transaction()?;
98 for (session_id, activity) in activities {
99 let exists: bool = tx.query_row(
100 "SELECT EXISTS(SELECT 1 FROM sessions WHERE session_id = ?1)",
101 [&session_id],
102 |r| r.get(0),
103 )?;
104 if !exists {
105 continue;
106 }
107 let body = serde_json::to_string(&activity)?;
108 let previous: Option<String> = tx
109 .query_row(
110 "SELECT body FROM api_session_activity WHERE session_id = ?1",
111 [&session_id],
112 |r| r.get(0),
113 )
114 .optional()?;
115 if previous.as_deref() == Some(body.as_str()) {
116 continue;
117 }
118 insert_api_event(
119 &tx,
120 &session_id,
121 recorded_at_ms,
122 &ApiEventData::ActivityChanged { activity },
123 )?;
124 tx.execute("INSERT INTO api_session_activity(session_id, body) VALUES (?1, ?2) ON CONFLICT(session_id) DO UPDATE SET body = excluded.body", params![session_id, body])?;
125 }
126 tx.commit()?;
127 Ok(())
128}
129
130pub fn record_startup_fault(session_id: String, message: String) -> Result<()> {
131 submit_database_write("record_startup_fault", move |connection| {
132 let tx = connection.transaction()?;
133 let exists: bool = tx.query_row(
134 "SELECT EXISTS(SELECT 1 FROM sessions WHERE session_id = ?1)",
135 [&session_id],
136 |r| r.get(0),
137 )?;
138 if exists {
139 insert_api_event(
140 &tx,
141 &session_id,
142 chrono::Utc::now().timestamp_millis(),
143 &ApiEventData::SessionFault {
144 reason: OutcomeReason::StartupFailed,
145 message,
146 command_id: None,
147 },
148 )?;
149 }
150 tx.commit()?;
151 Ok(())
152 })
153}
154
155pub(super) fn migrate_event_outcomes(tx: &Transaction<'_>) -> Result<()> {
157 let mut query = tx.prepare("SELECT seq, body FROM api_events WHERE json_extract(body, '$.type') IN ('error', 'turn_ended')")?;
158 let rows = query
159 .query_map([], |row| {
160 Ok((row.get::<_, u64>(0)?, row.get::<_, String>(1)?))
161 })?
162 .collect::<rusqlite::Result<Vec<_>>>()?;
163 for (seq, body) in rows {
164 let mut value: serde_json::Value = serde_json::from_str(&body)?;
165 if value["type"] == "error" {
166 value["type"] = "legacy_notice".into();
167 value["data"]["original_type"] = "error".into();
168 } else {
169 let turn: MaterializedTurnOutcome =
170 serde_json::from_value(value["data"]["turn"].clone())?;
171 value["data"]["turn"] =
172 serde_json::to_value(mj_core::event_outcome::ApiTurnOutcome::from(&turn))?;
173 }
174 tx.execute(
175 "UPDATE api_events SET body = ?2 WHERE seq = ?1",
176 params![seq, serde_json::to_string(&value)?],
177 )?;
178 }
179 Ok(())
180}
181
182pub(super) fn previous_session_error(
183 tx: &Transaction<'_>,
184 session_id: &str,
185) -> Result<Option<String>> {
186 Ok(tx
187 .query_row(
188 "SELECT last_error FROM sessions WHERE session_id = ?1",
189 [session_id],
190 |row| row.get::<_, Option<String>>(0),
191 )
192 .optional()?
193 .flatten())
194}
195
196pub(super) fn record_session_fault_transition(
198 tx: &Transaction<'_>,
199 session: &SessionRecord,
200 previous: Option<String>,
201) -> Result<()> {
202 if let Some(message) = &session.last_error
203 && previous.as_ref() != Some(message)
204 {
205 insert_api_event(
206 tx,
207 &session.id,
208 Utc::now().timestamp_millis(),
209 &ApiEventData::SessionFault {
210 reason: OutcomeReason::LifecycleFailed,
211 message: message.clone(),
212 command_id: None,
213 },
214 )?;
215 }
216 Ok(())
217}
218
219pub fn begin_checkpoint_operation(session: &SessionRecord, command_id: &str) -> Result<()> {
221 let session = session.clone();
222 let command_id = command_id.to_owned();
223 submit_database_write("begin_checkpoint_operation", move |connection| {
224 let tx = connection.transaction()?;
225 begin_checkpoint_operation_with(&tx, &session, &command_id)?;
226 tx.commit()?;
227 Ok(())
228 })
229}
230
231pub(super) fn begin_checkpoint_operation_with(
232 tx: &Transaction<'_>,
233 session: &SessionRecord,
234 command_id: &str,
235) -> Result<()> {
236 super::sessions::validate_session_record(session)?;
237 let state: String = tx.query_row(
238 "SELECT state FROM sessions WHERE session_id = ?1",
239 [&session.id],
240 |row| row.get(0),
241 )?;
242 ensure!(
243 !matches!(state.as_str(), "checkpointing" | "closing" | "destroying"),
244 "session {} is already in a lifecycle operation",
245 session.id
246 );
247 tx.execute(
248 "INSERT INTO checkpoint_operations(session_id, command_id) VALUES (?1, ?2)",
249 params![session.id, command_id],
250 )?;
251 super::state_io::update_lifecycle_fields(tx, session)?;
252 Ok(())
253}
254
255pub fn save_requested_checkpoint(session: &SessionRecord, command_id: &str) -> Result<()> {
256 let session = session.clone();
257 let command_id = command_id.to_owned();
258 submit_database_write("save_requested_checkpoint", move |connection| {
259 let tx = connection.transaction()?;
260 save_requested_checkpoint_with(&tx, &session, &command_id)?;
261 tx.commit()?;
262 Ok(())
263 })
264}
265
266fn save_requested_checkpoint_with(
267 tx: &Transaction<'_>,
268 session: &SessionRecord,
269 command_id: &str,
270) -> Result<()> {
271 super::sessions::validate_session_record(session)?;
272 let owns: bool = tx.query_row("SELECT EXISTS(SELECT 1 FROM checkpoint_operations WHERE session_id = ?1 AND command_id = ?2)", params![session.id, command_id], |row| row.get(0))?;
273 ensure!(
274 owns,
275 "checkpoint operation no longer owns session {}",
276 session.id
277 );
278 ensure!(
279 session.checkpoint.is_some(),
280 "successful checkpoint has no archive metadata"
281 );
282 super::state_io::update_lifecycle_fields(tx, session)?;
283 tx.execute(
284 "UPDATE sessions SET native_session_id = ?2 WHERE session_id = ?1",
285 params![session.id, session.native_session_id],
286 )?;
287 super::state_io::replace_checkpoint(tx, session)?;
288 finish_checkpoint_operation(tx, &session.id, CommandResultKind::Succeeded, None, None)
289}
290
291pub(super) fn finish_checkpoint_operation(
293 tx: &Transaction<'_>,
294 session_id: &str,
295 outcome: CommandResultKind,
296 reason: Option<OutcomeReason>,
297 message: Option<String>,
298) -> Result<()> {
299 let operation = tx.query_row("SELECT command_id, related_command_ids FROM checkpoint_operations WHERE session_id = ?1", [session_id], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))).optional()?;
300 if let Some((command_id, related)) = operation {
301 insert_api_event(
302 tx,
303 session_id,
304 Utc::now().timestamp_millis(),
305 &ApiEventData::CommandEnded {
306 result: CommandResult {
307 owner: CommandOwner::Daemon,
308 command_id,
309 command_kind: "checkpoint".into(),
310 outcome,
311 reason,
312 message,
313 related_command_ids: serde_json::from_str(&related)?,
314 },
315 },
316 )?;
317 tx.execute(
318 "DELETE FROM checkpoint_operations WHERE session_id = ?1",
319 [session_id],
320 )?;
321 }
322 Ok(())
323}
324
325pub fn correlate_checkpoint_barrier(
327 session_id: &str,
328 operation_id: &str,
329 command_id: &str,
330) -> Result<()> {
331 let session_id = session_id.to_owned();
332 let command_id = command_id.to_owned();
333 let operation_id = operation_id.to_owned();
334 submit_database_write("correlate_checkpoint_barrier", move |connection| {
335 connection.execute("UPDATE checkpoint_operations SET related_command_ids = json_insert(related_command_ids, '$[#]', ?2) WHERE session_id = ?1 AND command_id = ?3", params![session_id, command_id, operation_id])?;
336 Ok(())
337 })
338}
339
340pub fn finish_failed_checkpoint(
341 session: &SessionRecord,
342 command_id: &str,
343 deferred: bool,
344 message: String,
345) -> Result<()> {
346 let session = session.clone();
347 let command_id = command_id.to_owned();
348 submit_database_write("finish_failed_checkpoint", move |connection| {
349 let tx = connection.transaction()?;
350 finish_failed_checkpoint_with(&tx, &session, &command_id, deferred, message)?;
351 tx.commit()?;
352 Ok(())
353 })
354}
355
356fn finish_failed_checkpoint_with(
357 tx: &Transaction<'_>,
358 session: &SessionRecord,
359 command_id: &str,
360 deferred: bool,
361 message: String,
362) -> Result<()> {
363 super::sessions::validate_session_record(session)?;
364 let owns: bool = tx.query_row("SELECT EXISTS(SELECT 1 FROM checkpoint_operations WHERE session_id = ?1 AND command_id = ?2)", params![session.id, command_id], |row| row.get(0))?;
365 if !owns {
366 return Ok(());
367 }
368 super::state_io::update_lifecycle_fields(tx, session)?;
369 finish_checkpoint_operation(
370 tx,
371 &session.id,
372 if deferred {
373 CommandResultKind::Rejected
374 } else {
375 CommandResultKind::Failed
376 },
377 Some(if deferred {
378 OutcomeReason::CheckpointDeferred
379 } else {
380 OutcomeReason::CheckpointFailed
381 }),
382 Some(message),
383 )
384}
385
386#[cfg(test)]
387mod tests {
388 use super::*;
389 use mj_core::relay::{
390 RelayCommand, RelayCommandOutcome, RelayEvent, RelayObservation, relay_event_digest,
391 };
392 use mj_transcript::projection::{apply_committed_projection_event, project_relay_event};
393
394 fn page(path: &Path, observations: Vec<RelayObservation>, fail: bool) -> Result<()> {
395 let mut current = load_materialized_session_from(path, "session-1")?.unwrap();
396 apply_projection_page_to(path, "session-1", |page| {
397 for observation in observations {
398 let mut event = RelayEvent {
399 format: mj_core::relay::RELAY_EVENT_FORMAT_V1,
400 ordinal: current.applied_event_ordinal + 1,
401 previous_digest: current.applied_event_digest.clone(),
402 digest: String::new(),
403 recorded_at_ms: 100,
404 command_id: None,
405 observation,
406 };
407 event.digest = relay_event_digest(&event)?;
408 let mutation = project_relay_event(¤t, &event)?.mutation;
409 page.apply(
410 event.ordinal,
411 &event.previous_digest,
412 &event.digest,
413 &mutation,
414 )?;
415 apply_committed_projection_event(&mut current, &event, mutation)?;
416 }
417 if fail {
418 bail!("injected rollback");
419 }
420 Ok(())
421 })
422 }
423
424 fn turn() -> Vec<RelayObservation> {
425 vec![
426 RelayObservation::CommandQueued {
427 command_id: "prompt-1".into(),
428 command: RelayCommand::Prompt { prompt: vec![] },
429 created_at_ms: 1,
430 },
431 RelayObservation::CommandStarted {
432 command_id: "prompt-1".into(),
433 started_at_ms: 2,
434 },
435 RelayObservation::CommandCompleted {
436 barrier_command_id: None,
437 command: None,
438 command_id: "prompt-1".into(),
439 outcome: RelayCommandOutcome::Prompt {
440 diagnostic: None,
441 stop_reason: "end_turn".into(),
442 usage: None,
443 },
444 },
445 ]
446 }
447
448 #[test]
449 fn api_events_survive_coalescing_rollback_and_reopen() {
450 let dir = tempfile::tempdir().unwrap();
451 let path = dir.path().join("events.sqlite");
452 save_session_to(
453 &path,
454 &super::super::tests::session("session-1", "project-1"),
455 )
456 .unwrap();
457 assert!(page(&path, turn(), true).is_err());
458 assert!(
459 load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
460 .unwrap()
461 .events
462 .is_empty()
463 );
464 page(&path, turn(), false).unwrap();
465 let first = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 1).unwrap();
466 assert_eq!(first.events.len(), 1);
467 assert_eq!(first.events[0].event.kind(), "turn_started");
468 let rest = load_api_events_from(
469 &path,
470 &ApiEventFilter::default(),
471 Some(first.next_after_seq),
472 100,
473 )
474 .unwrap();
475 assert_eq!(rest.events.len(), 1);
476 assert_eq!(rest.events[0].event.kind(), "turn_ended");
477 let current = load_materialized_session_from(&path, "session-1")
478 .unwrap()
479 .unwrap();
480 assert!(current.active_turn.is_none());
481 apply_projection_event_to(
483 &path,
484 "session-1",
485 current.applied_event_ordinal,
486 "",
487 ¤t.applied_event_digest,
488 &MaterializedSessionMutation::default(),
489 )
490 .unwrap();
491 assert_eq!(
492 load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
493 .unwrap()
494 .events
495 .len(),
496 2
497 );
498 assert!(
499 load_api_events_from(&path, &ApiEventFilter::default(), None, 100)
500 .unwrap()
501 .events
502 .is_empty()
503 );
504 }
505
506 #[test]
507 fn api_events_filter_and_keep_cursors_after_forgetting_sessions() {
508 let dir = tempfile::tempdir().unwrap();
509 let path = dir.path().join("events.sqlite");
510 save_session_to(
511 &path,
512 &super::super::tests::session("session-1", "project-1"),
513 )
514 .unwrap();
515 page(&path, turn(), false).unwrap();
516 let filter = ApiEventFilter {
517 session_id: Some("other".into()),
518 workspace_id: None,
519 };
520 let empty = load_api_events_from(&path, &filter, Some(0), 100).unwrap();
521 assert!(empty.events.is_empty());
522 assert_eq!(empty.next_after_seq, 2);
523 let filter = ApiEventFilter {
524 session_id: None,
525 workspace_id: Some("default".into()),
526 };
527 assert_eq!(
528 load_api_events_from(&path, &filter, Some(0), 100)
529 .unwrap()
530 .events
531 .len(),
532 2
533 );
534 delete_session_from(&path, "session-1").unwrap();
535 let deleted =
536 load_api_events_from(&path, &ApiEventFilter::default(), Some(2), 100).unwrap();
537 assert!(deleted.events.is_empty());
538 assert_eq!(deleted.latest_seq, 2);
539 }
540
541 #[test]
542 fn api_events_capture_free_text_questions_resolution_and_errors() {
543 let dir = tempfile::tempdir().unwrap();
544 let path = dir.path().join("events.sqlite");
545 let mut record = super::super::tests::session("session-1", "project-1");
546 save_session_to(&path, &record).unwrap();
547 let request = ElicitationRequest::from_acp_params("question-1", serde_json::json!({
548 "mode": "form", "sessionId": "session-1", "message": "Which directory?",
549 "requestedSchema": {"type": "object", "properties": {"directory": {"type": "string"}}, "required": ["directory"]}
550 })).unwrap();
551 let mut observations = turn();
552 observations.splice(
553 2..2,
554 [
555 RelayObservation::ElicitationRequested {
556 request: request.clone(),
557 },
558 RelayObservation::ElicitationResolved {
559 elicitation_id: request.id.clone(),
560 action: "accept".into(),
561 },
562 ],
563 );
564 page(&path, observations, false).unwrap();
565 record.last_error = Some("provisioning failed".into());
566 save_session_to(&path, &record).unwrap();
567 save_session_to(&path, &record).unwrap();
568 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
569 assert_eq!(
570 page.events
571 .iter()
572 .map(|e| e.event.kind())
573 .collect::<Vec<_>>(),
574 [
575 "turn_started",
576 "input_required",
577 "input_resolved",
578 "turn_ended",
579 "session_fault"
580 ]
581 );
582 let ApiEventData::InputRequired {
583 request: actual,
584 turn_id,
585 } = &page.events[1].event
586 else {
587 panic!("question event")
588 };
589 assert_eq!(actual.as_ref(), Some(&request));
590 assert_eq!(*turn_id, Some(1));
591 assert!(
592 matches!(&page.events[2].event, ApiEventData::InputResolved { turn_id: Some(1), action, .. } if action == "accept")
593 );
594 let current = load_materialized_session_from(&path, "session-1")
595 .unwrap()
596 .unwrap();
597 assert!(current.pending_elicitations.is_empty());
598 assert!(current.active_turn.is_none());
599 assert_eq!(current.last_turn_outcome.unwrap().accepted_ordinal, Some(1));
600 }
601
602 #[test]
603 fn a_lifecycle_save_that_records_a_launch_failure_emits_one_error_event() {
604 let dir = tempfile::tempdir().unwrap();
608 let path = dir.path().join("events.sqlite");
609 let mut record = super::super::tests::session("session-1", "project-1");
610 save_session_to(&path, &record).unwrap();
611
612 record.state = mj_core::state::SessionState::Error;
613 record.last_error = Some("worker bootstrap failed: Connection closed by host".into());
614 save_lifecycle_session_to(&path, &record).unwrap();
615 save_lifecycle_session_to(&path, &record).unwrap();
617
618 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
619 let errors: Vec<_> = page
620 .events
621 .iter()
622 .filter_map(|event| match &event.event {
623 ApiEventData::SessionFault { message, .. } => Some(message.clone()),
624 _ => None,
625 })
626 .collect();
627 assert_eq!(
628 errors,
629 vec!["worker bootstrap failed: Connection closed by host".to_owned()],
630 "a recorded launch failure surfaces as exactly one error event"
631 );
632 }
633
634 #[test]
635 fn api_events_preserve_command_identity_for_failed_completions() {
636 let dir = tempfile::tempdir().unwrap();
637 let path = dir.path().join("events.sqlite");
638 save_session_to(
639 &path,
640 &super::super::tests::session("session-1", "project-1"),
641 )
642 .unwrap();
643 let mut observations = turn();
644 let RelayObservation::CommandCompleted { outcome, .. } = observations.last_mut().unwrap()
645 else {
646 unreachable!()
647 };
648 *outcome = RelayCommandOutcome::Prompt {
649 diagnostic: None,
650 stop_reason: "provider_error".into(),
651 usage: None,
652 };
653 page(&path, observations, false).unwrap();
654 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
655 .unwrap()
656 .events;
657 assert_eq!(events.len(), 2);
658 assert!(matches!(&events[1].event, ApiEventData::TurnEnded { turn }
659 if turn.command_id == "prompt-1" && turn.accepted_ordinal == Some(1)
660 && turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Failed
661 && turn.outcome.stop_reason.as_deref() == Some("provider_error")));
662 }
663
664 #[test]
665 fn api_activity_events_track_changes_without_repeating_or_reviving_forgotten_sessions() {
666 let dir = tempfile::tempdir().unwrap();
667 let path = dir.path().join("events.sqlite");
668 save_session_to(
669 &path,
670 &super::super::tests::session("session-1", "project-1"),
671 )
672 .unwrap();
673 let mut connection = open(&path).unwrap();
674 let idle = ApiActivityState {
675 state: "running".into(),
676 details: Some(ApiActivityDetails {
677 kind: ApiActivityKind::Idle,
678 turn_started_at_ms: None,
679 step_started_at_ms: None,
680 background_started_at_ms: None,
681 idle_since_ms: Some(100),
682 last_activity_at_ms: None,
683 label: None,
684 }),
685 is_idle: true,
686 waiting_for_input: false,
687 capacity_retry: false,
688 };
689 record_api_activities_with(
690 &mut connection,
691 vec![("session-1".into(), idle.clone())],
692 100,
693 )
694 .unwrap();
695 record_api_activities_with(
696 &mut connection,
697 vec![("session-1".into(), idle.clone())],
698 200,
699 )
700 .unwrap();
701 let mut background = idle.clone();
702 background.is_idle = false;
703 let details = background.details.as_mut().unwrap();
704 details.kind = ApiActivityKind::Background;
705 details.idle_since_ms = None;
706 details.background_started_at_ms = Some(250);
707 record_api_activities_with(
708 &mut connection,
709 vec![("session-1".into(), background.clone())],
710 250,
711 )
712 .unwrap();
713 background.waiting_for_input = true;
714 record_api_activities_with(
715 &mut connection,
716 vec![("session-1".into(), background.clone())],
717 300,
718 )
719 .unwrap();
720 let unknown = ApiActivityState {
721 state: "disconnected".into(),
722 details: None,
723 is_idle: false,
724 waiting_for_input: false,
725 capacity_retry: false,
726 };
727 record_api_activities_with(
728 &mut connection,
729 vec![("session-1".into(), unknown.clone())],
730 400,
731 )
732 .unwrap();
733 drop(connection);
734 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
735 assert_eq!(page.events.len(), 4);
736 assert_eq!(
737 page.events.last().unwrap().event,
738 ApiEventData::ActivityChanged { activity: unknown }
739 );
740 assert_eq!(
741 page.events[2].event,
742 ApiEventData::ActivityChanged {
743 activity: background
744 }
745 );
746 delete_session_from(&path, "session-1").unwrap();
747 record_api_activities_with(
748 &mut open(&path).unwrap(),
749 vec![("session-1".into(), idle)],
750 500,
751 )
752 .unwrap();
753 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(4), 100).unwrap();
754 assert!(page.events.is_empty());
755 assert_eq!(page.latest_seq, 4);
756 }
757 #[test]
758 fn checkpoint_cleanup_does_not_end_a_turn_or_hide_control_failure() {
759 use mj_core::event_outcome::TurnResultKind;
760 use mj_core::relay::RelayCommandKind;
761 let dir = tempfile::tempdir().unwrap();
762 let path = dir.path().join("events.sqlite");
763 save_session_to(
764 &path,
765 &super::super::tests::session("session-1", "project-1"),
766 )
767 .unwrap();
768 page(&path, turn().into_iter().take(2).collect(), false).unwrap();
769 for (id, reason) in [
770 ("arbitrary-one", OutcomeReason::ControllerDisconnected),
771 ("arbitrary-two", OutcomeReason::OwnerLostOnRestart),
772 ] {
773 page(
774 &path,
775 vec![RelayObservation::CommandInterrupted {
776 command_id: id.into(),
777 command: RelayCommandKind::BeginCheckpoint,
778 reason: Some(reason),
779 message: "diagnostic wording is not a contract".into(),
780 }],
781 false,
782 )
783 .unwrap();
784 }
785 page(
786 &path,
787 vec![RelayObservation::CommandRejected {
788 command_id: "arbitrary-three".into(),
789 command: RelayCommandKind::CompleteCheckpoint,
790 reason: Some(OutcomeReason::CommandFailed),
791 message: "Cannot save recovery floor".into(),
792 }],
793 false,
794 )
795 .unwrap();
796 let current = load_materialized_session_from(&path, "session-1")
797 .unwrap()
798 .unwrap();
799 assert!(current.active_turn.is_some());
800 assert!(current.last_turn_outcome.is_none());
801 assert!(current.transcript.iter().all(|item| !matches!(&item.body, TranscriptBody::System { text } if text.contains("diagnostic wording"))));
802 assert!(current.transcript.iter().any(|item| matches!(&item.body, TranscriptBody::System { text } if text.contains("Cannot save recovery floor"))));
803 page(&path, vec![turn().pop().unwrap()], false).unwrap();
804 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
805 .unwrap()
806 .events;
807 assert_eq!(events.len(), 5);
808 for event in &events[1..3] {
809 assert!(
810 matches!(&event.event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Cancelled && result.command_kind == "begin_checkpoint")
811 );
812 }
813 assert!(
814 matches!(&events[3].event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Failed)
815 );
816 assert!(
817 matches!(&events[4].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == TurnResultKind::Completed)
818 );
819 }
820
821 #[test]
822 fn terminal_prompt_results_are_semantic_and_not_duplicate_error_events() {
823 use mj_core::event_outcome::TurnResultKind;
824 use mj_core::relay::RelayCommandKind;
825 for (terminal, expected) in [
826 (
827 RelayObservation::CommandCompleted {
828 barrier_command_id: None,
829 command_id: "prompt-1".into(),
830 command: Some(RelayCommandKind::Prompt),
831 outcome: RelayCommandOutcome::Prompt {
832 stop_reason: "Cancelled".into(),
833 usage: None,
834 diagnostic: None,
835 },
836 },
837 TurnResultKind::Cancelled,
838 ),
839 (
840 RelayObservation::CommandCompleted {
841 barrier_command_id: None,
842 command_id: "prompt-1".into(),
843 command: Some(RelayCommandKind::Prompt),
844 outcome: RelayCommandOutcome::Prompt {
845 stop_reason: "QuotaLimit".into(),
846 usage: None,
847 diagnostic: None,
848 },
849 },
850 TurnResultKind::Failed,
851 ),
852 (
853 RelayObservation::CommandRejected {
854 command_id: "prompt-1".into(),
855 command: RelayCommandKind::Prompt,
856 reason: Some(OutcomeReason::AdmissionRejected),
857 message: "unavailable".into(),
858 },
859 TurnResultKind::Rejected,
860 ),
861 (
862 RelayObservation::CommandInterrupted {
863 command_id: "prompt-1".into(),
864 command: RelayCommandKind::Prompt,
865 reason: Some(OutcomeReason::RuntimeStopped),
866 message: "runtime stopped".into(),
867 },
868 TurnResultKind::Interrupted,
869 ),
870 ] {
871 let dir = tempfile::tempdir().unwrap();
872 let path = dir.path().join("events.sqlite");
873 save_session_to(
874 &path,
875 &super::super::tests::session("session-1", "project-1"),
876 )
877 .unwrap();
878 let mut observations = turn();
879 *observations.last_mut().unwrap() = terminal;
880 page(&path, observations, false).unwrap();
881 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
882 .unwrap()
883 .events;
884 assert_eq!(events.len(), 2);
885 assert!(
886 matches!(&events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == expected && turn.command_id == "prompt-1")
887 );
888 }
889 }
890
891 #[test]
892 fn typed_runtime_fault_does_not_rewrite_a_completed_turn() {
893 let dir = tempfile::tempdir().unwrap();
894 let path = dir.path().join("events.sqlite");
895 save_session_to(
896 &path,
897 &super::super::tests::session("session-1", "project-1"),
898 )
899 .unwrap();
900 page(&path, turn(), false).unwrap();
901 page(
902 &path,
903 vec![RelayObservation::SessionFault {
904 reason: OutcomeReason::RuntimeUnavailable,
905 message: "runtime exited".into(),
906 }],
907 false,
908 )
909 .unwrap();
910 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
911 .unwrap()
912 .events;
913 assert_eq!(events.len(), 3);
914 assert!(
915 matches!(&events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Completed)
916 );
917 assert!(matches!(
918 &events[2].event,
919 ApiEventData::SessionFault {
920 reason: OutcomeReason::RuntimeUnavailable,
921 ..
922 }
923 ));
924 }
925
926 #[test]
927 fn checkpoint_recovery_consumes_attempt_identity_once_and_preserves_correlation() {
928 let dir = tempfile::tempdir().unwrap();
929 let path = dir.path().join("events.sqlite");
930 let mut record = super::super::tests::session("session-1", "project-1");
931 record.state = SessionState::Running;
932 save_session_to(&path, &record).unwrap();
933 record.state = SessionState::Checkpointing;
934 let mut connection = open(&path).unwrap();
935 {
936 let tx = connection.transaction().unwrap();
937 begin_checkpoint_operation_with(&tx, &record, "attempt-1").unwrap();
938 assert!(begin_checkpoint_operation_with(&tx, &record, "attempt-2").is_err());
939 tx.execute(
940 "UPDATE checkpoint_operations SET related_command_ids = '[\"barrier-1\"]'",
941 [],
942 )
943 .unwrap();
944 tx.commit().unwrap();
945 }
946 drop(connection);
947 assert_eq!(
948 recover_interrupted_checkpointing_sessions_to(&path, "2026-09-27T12:00:00Z").unwrap(),
949 1
950 );
951 assert_eq!(
952 recover_interrupted_checkpointing_sessions_to(&path, "2026-09-27T12:00:01Z").unwrap(),
953 0
954 );
955 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
956 .unwrap()
957 .events;
958 assert_eq!(events.len(), 1);
959 assert!(
960 matches!(&events[0].event, ApiEventData::CommandEnded { result }
961 if result.command_id == "attempt-1" && result.related_command_ids == ["barrier-1"]
962 && result.owner == CommandOwner::Daemon && result.reason == Some(OutcomeReason::ControllerRestarted))
963 );
964 }
965
966 #[test]
967 fn event_migration_preserves_cursors_and_keeps_unknown_history_explicit() {
968 let dir = tempfile::tempdir().unwrap();
969 let path = dir.path().join("events.sqlite");
970 save_session_to(
971 &path,
972 &super::super::tests::session("session-1", "project-1"),
973 )
974 .unwrap();
975 page(&path, turn(), false).unwrap();
976 let turn = load_materialized_session_from(&path, "session-1")
977 .unwrap()
978 .unwrap()
979 .last_turn_outcome
980 .unwrap();
981 let connection = Connection::open(&path).unwrap();
982 connection
983 .execute(
984 "UPDATE api_events SET body = ?1 WHERE seq = 2",
985 [serde_json::json!({"type":"turn_ended", "data":{"turn":turn}}).to_string()],
986 )
987 .unwrap();
988 connection.execute("INSERT INTO api_events(seq, session_id, recorded_at_ms, body) VALUES (8, 'session-1', 1234, ?1)", [r#"{"type":"error","data":{"command_id":"arbitrary","message":"old diagnostic"}}"#]).unwrap();
989 connection.execute_batch("DROP TABLE checkpoint_operations; DELETE FROM schema_migrations WHERE version >= 57; UPDATE schema_compatibility SET minimum_compatible_version = 56; DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 56;").unwrap();
990 drop(connection);
991 forget_verified_schema(&path);
992 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
993 assert_eq!(page.latest_seq, 8);
994 assert_eq!(page.next_after_seq, 8);
995 assert_eq!(
996 page.events
997 .iter()
998 .map(|event| event.seq)
999 .collect::<Vec<_>>(),
1000 [1, 2, 8]
1001 );
1002 assert_eq!(page.events[2].recorded_at_ms, 1234);
1003 assert!(
1004 matches!(&page.events[2].event, ApiEventData::LegacyNotice { original_type, command_id: Some(id), message } if original_type == "error" && id == "arbitrary" && message == "old diagnostic")
1005 );
1006 assert!(
1007 matches!(&page.events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Completed)
1008 );
1009 }
1010 #[test]
1011 fn requested_checkpoint_result_commits_with_metadata_and_only_for_its_owner() {
1012 let dir = tempfile::tempdir().unwrap();
1013 let path = dir.path().join("events.sqlite");
1014 let mut record = super::super::tests::session("session-1", "project-1");
1015 record.state = SessionState::Running;
1016 let previous_checkpoint = record.checkpoint.clone();
1017 save_session_to(&path, &record).unwrap();
1018 let mut connection = open(&path).unwrap();
1019 record.state = SessionState::Checkpointing;
1020 {
1021 let tx = connection.transaction().unwrap();
1022 begin_checkpoint_operation_with(&tx, &record, "owner").unwrap();
1023 tx.commit().unwrap();
1024 }
1025 record.state = SessionState::Running;
1026 record.checkpoint = Some(CheckpointMetadata {
1027 archive_path: dir.path().join("verified.tar"),
1028 sha256: "a".repeat(64),
1029 created_at: "2026-09-27T12:00:00Z".into(),
1030 event_frontier: 0,
1031 });
1032 {
1033 let tx = connection.transaction().unwrap();
1034 assert!(save_requested_checkpoint_with(&tx, &record, "stale-owner").is_err());
1035 finish_failed_checkpoint_with(
1036 &tx,
1037 &record,
1038 "stale-owner",
1039 false,
1040 "late failure".into(),
1041 )
1042 .unwrap();
1043 assert_eq!(
1044 tx.query_row("SELECT count(*) FROM api_events", [], |r| r
1045 .get::<_, usize>(0))
1046 .unwrap(),
1047 0
1048 );
1049 save_requested_checkpoint_with(&tx, &record, "owner").unwrap();
1050 }
1052 assert_eq!(
1053 load_state_from(&path).unwrap().sessions["session-1"].checkpoint,
1054 previous_checkpoint
1055 );
1056 assert!(
1057 load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1058 .unwrap()
1059 .events
1060 .is_empty()
1061 );
1062 {
1063 let tx = connection.transaction().unwrap();
1064 save_requested_checkpoint_with(&tx, &record, "owner").unwrap();
1065 finish_failed_checkpoint_with(
1066 &tx,
1067 &record,
1068 "owner",
1069 false,
1070 "late cleanup failure".into(),
1071 )
1072 .unwrap();
1073 tx.commit().unwrap();
1074 }
1075 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1076 .unwrap()
1077 .events;
1078 assert_eq!(events.len(), 1);
1079 assert!(
1080 matches!(&events[0].event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Succeeded && result.command_id == "owner")
1081 );
1082 assert_eq!(
1083 load_state_from(&path).unwrap().sessions["session-1"].checkpoint,
1084 record.checkpoint
1085 );
1086 }
1087
1088 #[test]
1089 fn checkpoint_failure_and_deferral_are_distinct_and_preserve_previous_archive() {
1090 for deferred in [false, true] {
1091 let dir = tempfile::tempdir().unwrap();
1092 let path = dir.path().join("events.sqlite");
1093 let mut record = super::super::tests::session("session-1", "project-1");
1094 record.state = SessionState::Running;
1095 record.checkpoint = Some(CheckpointMetadata {
1096 archive_path: dir.path().join("previous.tar"),
1097 sha256: "a".repeat(64),
1098 created_at: "2026-09-27T12:00:00Z".into(),
1099 event_frontier: 0,
1100 });
1101 save_session_to(&path, &record).unwrap();
1102 let previous = record.checkpoint.clone();
1103 let previous_warning = record.last_checkpoint_error.clone();
1104 let mut connection = open(&path).unwrap();
1105 let tx = connection.transaction().unwrap();
1106 record.state = SessionState::Checkpointing;
1107 begin_checkpoint_operation_with(&tx, &record, "owner").unwrap();
1108 record.state = SessionState::Running;
1109 if !deferred {
1110 record.last_checkpoint_error = Some("archive verification failed".into());
1111 }
1112 finish_failed_checkpoint_with(&tx, &record, "owner", deferred, "diagnostic".into())
1113 .unwrap();
1114 finish_failed_checkpoint_with(&tx, &record, "owner", deferred, "diagnostic".into())
1115 .unwrap();
1116 tx.commit().unwrap();
1117 let loaded = load_state_from(&path).unwrap();
1118 assert_eq!(loaded.sessions["session-1"].checkpoint, previous);
1119 assert_eq!(
1120 loaded.sessions["session-1"].last_checkpoint_error,
1121 if deferred {
1122 previous_warning
1123 } else {
1124 Some("archive verification failed".into())
1125 }
1126 );
1127 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1128 .unwrap()
1129 .events;
1130 assert_eq!(events.len(), 1);
1131 assert!(
1132 matches!(&events[0].event, ApiEventData::CommandEnded { result } if result.outcome == if deferred { CommandResultKind::Rejected } else { CommandResultKind::Failed })
1133 );
1134 }
1135 }
1136}