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]
604 fn a_lifecycle_save_that_records_a_launch_failure_emits_one_error_event() {
605 let dir = tempfile::tempdir().unwrap();
609 let path = dir.path().join("events.sqlite");
610 let mut record = super::super::tests::session("session-1", "project-1");
611 save_session_to(&path, &record).unwrap();
612
613 record.state = mj_core::state::SessionState::Error;
614 record.last_error = Some("worker bootstrap failed: Connection closed by host".into());
615 save_lifecycle_session_to(&path, &record).unwrap();
616 save_lifecycle_session_to(&path, &record).unwrap();
618
619 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
620 let errors: Vec<_> = page
621 .events
622 .iter()
623 .filter_map(|event| match &event.event {
624 ApiEventData::SessionFault { message, .. } => Some(message.clone()),
625 _ => None,
626 })
627 .collect();
628 assert_eq!(
629 errors,
630 vec!["worker bootstrap failed: Connection closed by host".to_owned()],
631 "a recorded launch failure surfaces as exactly one error event"
632 );
633 }
634
635 #[test]
636 fn api_events_preserve_command_identity_for_failed_completions() {
637 let dir = tempfile::tempdir().unwrap();
638 let path = dir.path().join("events.sqlite");
639 save_session_to(
640 &path,
641 &super::super::tests::session("session-1", "project-1"),
642 )
643 .unwrap();
644 let mut observations = turn();
645 let RelayObservation::CommandCompleted { outcome, .. } = observations.last_mut().unwrap()
646 else {
647 unreachable!()
648 };
649 *outcome = RelayCommandOutcome::Prompt {
650 diagnostic: None,
651 stop_reason: "provider_error".into(),
652 usage: None,
653 };
654 page(&path, observations, false).unwrap();
655 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
656 .unwrap()
657 .events;
658 assert_eq!(events.len(), 2);
659 assert!(matches!(&events[1].event, ApiEventData::TurnEnded { turn }
660 if turn.command_id == "prompt-1" && turn.accepted_ordinal == Some(1)
661 && turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Failed
662 && turn.outcome.stop_reason.as_deref() == Some("provider_error")));
663 }
664
665 #[test]
666 fn api_activity_events_track_changes_without_repeating_or_reviving_forgotten_sessions() {
667 let dir = tempfile::tempdir().unwrap();
668 let path = dir.path().join("events.sqlite");
669 save_session_to(
670 &path,
671 &super::super::tests::session("session-1", "project-1"),
672 )
673 .unwrap();
674 let mut connection = open(&path).unwrap();
675 let idle = ApiActivityState {
676 state: "running".into(),
677 details: Some(ApiActivityDetails {
678 kind: ApiActivityKind::Idle,
679 turn_started_at_ms: None,
680 step_started_at_ms: None,
681 background_started_at_ms: None,
682 idle_since_ms: Some(100),
683 last_activity_at_ms: None,
684 label: None,
685 }),
686 is_idle: true,
687 waiting_for_input: false,
688 capacity_retry: false,
689 };
690 record_api_activities_with(
691 &mut connection,
692 vec![("session-1".into(), idle.clone())],
693 100,
694 )
695 .unwrap();
696 record_api_activities_with(
697 &mut connection,
698 vec![("session-1".into(), idle.clone())],
699 200,
700 )
701 .unwrap();
702 let mut background = idle.clone();
703 background.is_idle = false;
704 let details = background.details.as_mut().unwrap();
705 details.kind = ApiActivityKind::Background;
706 details.idle_since_ms = None;
707 details.background_started_at_ms = Some(250);
708 record_api_activities_with(
709 &mut connection,
710 vec![("session-1".into(), background.clone())],
711 250,
712 )
713 .unwrap();
714 background.waiting_for_input = true;
715 record_api_activities_with(
716 &mut connection,
717 vec![("session-1".into(), background.clone())],
718 300,
719 )
720 .unwrap();
721 let unknown = ApiActivityState {
722 state: "disconnected".into(),
723 details: None,
724 is_idle: false,
725 waiting_for_input: false,
726 capacity_retry: false,
727 };
728 record_api_activities_with(
729 &mut connection,
730 vec![("session-1".into(), unknown.clone())],
731 400,
732 )
733 .unwrap();
734 drop(connection);
735 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
736 assert_eq!(page.events.len(), 4);
737 assert_eq!(
738 page.events.last().unwrap().event,
739 ApiEventData::ActivityChanged { activity: unknown }
740 );
741 assert_eq!(
742 page.events[2].event,
743 ApiEventData::ActivityChanged {
744 activity: background
745 }
746 );
747 delete_session_from(&path, "session-1").unwrap();
748 record_api_activities_with(
749 &mut open(&path).unwrap(),
750 vec![("session-1".into(), idle)],
751 500,
752 )
753 .unwrap();
754 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(4), 100).unwrap();
755 assert!(page.events.is_empty());
756 assert_eq!(page.latest_seq, 4);
757 }
758 #[test]
759 fn checkpoint_cleanup_does_not_end_a_turn_or_hide_control_failure() {
760 use mj_core::event_outcome::TurnResultKind;
761 use mj_core::relay::RelayCommandKind;
762 let dir = tempfile::tempdir().unwrap();
763 let path = dir.path().join("events.sqlite");
764 save_session_to(
765 &path,
766 &super::super::tests::session("session-1", "project-1"),
767 )
768 .unwrap();
769 page(&path, turn().into_iter().take(2).collect(), false).unwrap();
770 for (id, reason) in [
771 ("arbitrary-one", OutcomeReason::ControllerDisconnected),
772 ("arbitrary-two", OutcomeReason::OwnerLostOnRestart),
773 ] {
774 page(
775 &path,
776 vec![RelayObservation::CommandInterrupted {
777 command_id: id.into(),
778 command: RelayCommandKind::BeginCheckpoint,
779 reason: Some(reason),
780 message: "diagnostic wording is not a contract".into(),
781 }],
782 false,
783 )
784 .unwrap();
785 }
786 page(
787 &path,
788 vec![RelayObservation::CommandRejected {
789 command_id: "arbitrary-three".into(),
790 command: RelayCommandKind::CompleteCheckpoint,
791 reason: Some(OutcomeReason::CommandFailed),
792 message: "Cannot save recovery floor".into(),
793 }],
794 false,
795 )
796 .unwrap();
797 let current = load_materialized_session_from(&path, "session-1")
798 .unwrap()
799 .unwrap();
800 assert!(current.active_turn.is_some());
801 assert!(current.last_turn_outcome.is_none());
802 assert!(current.transcript.iter().all(|item| !matches!(&item.body, TranscriptBody::System { text } if text.contains("diagnostic wording"))));
803 assert!(current.transcript.iter().any(|item| matches!(&item.body, TranscriptBody::System { text } if text.contains("Cannot save recovery floor"))));
804 page(&path, vec![turn().pop().unwrap()], false).unwrap();
805 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
806 .unwrap()
807 .events;
808 assert_eq!(events.len(), 5);
809 for event in &events[1..3] {
810 assert!(
811 matches!(&event.event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Cancelled && result.command_kind == "begin_checkpoint")
812 );
813 }
814 assert!(
815 matches!(&events[3].event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Failed)
816 );
817 assert!(
818 matches!(&events[4].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == TurnResultKind::Completed)
819 );
820 }
821
822 #[test]
823 fn terminal_prompt_results_are_semantic_and_not_duplicate_error_events() {
824 use mj_core::event_outcome::TurnResultKind;
825 use mj_core::relay::RelayCommandKind;
826 for (terminal, expected) in [
827 (
828 RelayObservation::CommandCompleted {
829 barrier_command_id: None,
830 command_id: "prompt-1".into(),
831 command: Some(RelayCommandKind::Prompt),
832 outcome: RelayCommandOutcome::Prompt {
833 stop_reason: "Cancelled".into(),
834 usage: None,
835 diagnostic: None,
836 },
837 },
838 TurnResultKind::Cancelled,
839 ),
840 (
841 RelayObservation::CommandCompleted {
842 barrier_command_id: None,
843 command_id: "prompt-1".into(),
844 command: Some(RelayCommandKind::Prompt),
845 outcome: RelayCommandOutcome::Prompt {
846 stop_reason: "QuotaLimit".into(),
847 usage: None,
848 diagnostic: None,
849 },
850 },
851 TurnResultKind::Failed,
852 ),
853 (
854 RelayObservation::CommandRejected {
855 command_id: "prompt-1".into(),
856 command: RelayCommandKind::Prompt,
857 reason: Some(OutcomeReason::AdmissionRejected),
858 message: "unavailable".into(),
859 },
860 TurnResultKind::Rejected,
861 ),
862 (
863 RelayObservation::CommandInterrupted {
864 command_id: "prompt-1".into(),
865 command: RelayCommandKind::Prompt,
866 reason: Some(OutcomeReason::RuntimeStopped),
867 message: "runtime stopped".into(),
868 },
869 TurnResultKind::Interrupted,
870 ),
871 ] {
872 let dir = tempfile::tempdir().unwrap();
873 let path = dir.path().join("events.sqlite");
874 save_session_to(
875 &path,
876 &super::super::tests::session("session-1", "project-1"),
877 )
878 .unwrap();
879 let mut observations = turn();
880 *observations.last_mut().unwrap() = terminal;
881 page(&path, observations, false).unwrap();
882 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
883 .unwrap()
884 .events;
885 assert_eq!(events.len(), 2);
886 assert!(
887 matches!(&events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == expected && turn.command_id == "prompt-1")
888 );
889 }
890 }
891
892 #[test]
893 fn typed_runtime_fault_does_not_rewrite_a_completed_turn() {
894 let dir = tempfile::tempdir().unwrap();
895 let path = dir.path().join("events.sqlite");
896 save_session_to(
897 &path,
898 &super::super::tests::session("session-1", "project-1"),
899 )
900 .unwrap();
901 page(&path, turn(), false).unwrap();
902 page(
903 &path,
904 vec![RelayObservation::SessionFault {
905 reason: OutcomeReason::RuntimeUnavailable,
906 message: "runtime exited".into(),
907 }],
908 false,
909 )
910 .unwrap();
911 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
912 .unwrap()
913 .events;
914 assert_eq!(events.len(), 3);
915 assert!(
916 matches!(&events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Completed)
917 );
918 assert!(matches!(
919 &events[2].event,
920 ApiEventData::SessionFault {
921 reason: OutcomeReason::RuntimeUnavailable,
922 ..
923 }
924 ));
925 }
926
927 #[test]
928 fn checkpoint_recovery_consumes_attempt_identity_once_and_preserves_correlation() {
929 let dir = tempfile::tempdir().unwrap();
930 let path = dir.path().join("events.sqlite");
931 let mut record = super::super::tests::session("session-1", "project-1");
932 record.state = SessionState::Running;
933 save_session_to(&path, &record).unwrap();
934 record.state = SessionState::Checkpointing;
935 let mut connection = open(&path).unwrap();
936 {
937 let tx = connection.transaction().unwrap();
938 begin_checkpoint_operation_with(&tx, &record, "attempt-1").unwrap();
939 assert!(begin_checkpoint_operation_with(&tx, &record, "attempt-2").is_err());
940 tx.execute(
941 "UPDATE checkpoint_operations SET related_command_ids = '[\"barrier-1\"]'",
942 [],
943 )
944 .unwrap();
945 tx.commit().unwrap();
946 }
947 drop(connection);
948 assert_eq!(
949 recover_interrupted_checkpointing_sessions_to(&path, "2026-09-27T12:00:00Z").unwrap(),
950 1
951 );
952 assert_eq!(
953 recover_interrupted_checkpointing_sessions_to(&path, "2026-09-27T12:00:01Z").unwrap(),
954 0
955 );
956 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
957 .unwrap()
958 .events;
959 assert_eq!(events.len(), 1);
960 assert!(
961 matches!(&events[0].event, ApiEventData::CommandEnded { result }
962 if result.command_id == "attempt-1" && result.related_command_ids == ["barrier-1"]
963 && result.owner == CommandOwner::Daemon && result.reason == Some(OutcomeReason::ControllerRestarted))
964 );
965 }
966
967 #[test]
968 fn event_migration_preserves_cursors_and_keeps_unknown_history_explicit() {
969 let dir = tempfile::tempdir().unwrap();
970 let path = dir.path().join("events.sqlite");
971 save_session_to(
972 &path,
973 &super::super::tests::session("session-1", "project-1"),
974 )
975 .unwrap();
976 page(&path, turn(), false).unwrap();
977 let turn = load_materialized_session_from(&path, "session-1")
978 .unwrap()
979 .unwrap()
980 .last_turn_outcome
981 .unwrap();
982 let connection = Connection::open(&path).unwrap();
983 connection
984 .execute(
985 "UPDATE api_events SET body = ?1 WHERE seq = 2",
986 [serde_json::json!({"type":"turn_ended", "data":{"turn":turn}}).to_string()],
987 )
988 .unwrap();
989 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();
990 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();
991 drop(connection);
992 forget_verified_schema(&path);
993 let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
994 assert_eq!(page.latest_seq, 8);
995 assert_eq!(page.next_after_seq, 8);
996 assert_eq!(
997 page.events
998 .iter()
999 .map(|event| event.seq)
1000 .collect::<Vec<_>>(),
1001 [1, 2, 8]
1002 );
1003 assert_eq!(page.events[2].recorded_at_ms, 1234);
1004 assert!(
1005 matches!(&page.events[2].event, ApiEventData::LegacyNotice { original_type, command_id: Some(id), message } if original_type == "error" && id == "arbitrary" && message == "old diagnostic")
1006 );
1007 assert!(
1008 matches!(&page.events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Completed)
1009 );
1010 }
1011 #[test]
1012 fn requested_checkpoint_result_commits_with_metadata_and_only_for_its_owner() {
1013 let dir = tempfile::tempdir().unwrap();
1014 let path = dir.path().join("events.sqlite");
1015 let mut record = super::super::tests::session("session-1", "project-1");
1016 record.state = SessionState::Running;
1017 let previous_checkpoint = record.checkpoint.clone();
1018 save_session_to(&path, &record).unwrap();
1019 let mut connection = open(&path).unwrap();
1020 record.state = SessionState::Checkpointing;
1021 {
1022 let tx = connection.transaction().unwrap();
1023 begin_checkpoint_operation_with(&tx, &record, "owner").unwrap();
1024 tx.commit().unwrap();
1025 }
1026 record.state = SessionState::Running;
1027 record.checkpoint = Some(CheckpointMetadata {
1028 archive_path: dir.path().join("verified.tar"),
1029 sha256: "a".repeat(64),
1030 created_at: "2026-09-27T12:00:00Z".into(),
1031 event_frontier: 0,
1032 });
1033 {
1034 let tx = connection.transaction().unwrap();
1035 assert!(save_requested_checkpoint_with(&tx, &record, "stale-owner").is_err());
1036 finish_failed_checkpoint_with(
1037 &tx,
1038 &record,
1039 "stale-owner",
1040 false,
1041 "late failure".into(),
1042 )
1043 .unwrap();
1044 assert_eq!(
1045 tx.query_row("SELECT count(*) FROM api_events", [], |r| r
1046 .get::<_, usize>(0))
1047 .unwrap(),
1048 0
1049 );
1050 save_requested_checkpoint_with(&tx, &record, "owner").unwrap();
1051 }
1053 assert_eq!(
1054 load_state_from(&path).unwrap().sessions["session-1"].checkpoint,
1055 previous_checkpoint
1056 );
1057 assert!(
1058 load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1059 .unwrap()
1060 .events
1061 .is_empty()
1062 );
1063 {
1064 let tx = connection.transaction().unwrap();
1065 save_requested_checkpoint_with(&tx, &record, "owner").unwrap();
1066 finish_failed_checkpoint_with(
1067 &tx,
1068 &record,
1069 "owner",
1070 false,
1071 "late cleanup failure".into(),
1072 )
1073 .unwrap();
1074 tx.commit().unwrap();
1075 }
1076 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1077 .unwrap()
1078 .events;
1079 assert_eq!(events.len(), 1);
1080 assert!(
1081 matches!(&events[0].event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Succeeded && result.command_id == "owner")
1082 );
1083 assert_eq!(
1084 load_state_from(&path).unwrap().sessions["session-1"].checkpoint,
1085 record.checkpoint
1086 );
1087 }
1088
1089 #[test]
1090 fn checkpoint_failure_and_deferral_are_distinct_and_preserve_previous_archive() {
1091 for deferred in [false, true] {
1092 let dir = tempfile::tempdir().unwrap();
1093 let path = dir.path().join("events.sqlite");
1094 let mut record = super::super::tests::session("session-1", "project-1");
1095 record.state = SessionState::Running;
1096 record.checkpoint = Some(CheckpointMetadata {
1097 archive_path: dir.path().join("previous.tar"),
1098 sha256: "a".repeat(64),
1099 created_at: "2026-09-27T12:00:00Z".into(),
1100 event_frontier: 0,
1101 });
1102 save_session_to(&path, &record).unwrap();
1103 let previous = record.checkpoint.clone();
1104 let previous_warning = record.last_checkpoint_error.clone();
1105 let mut connection = open(&path).unwrap();
1106 let tx = connection.transaction().unwrap();
1107 record.state = SessionState::Checkpointing;
1108 begin_checkpoint_operation_with(&tx, &record, "owner").unwrap();
1109 record.state = SessionState::Running;
1110 if !deferred {
1111 record.last_checkpoint_error = Some("archive verification failed".into());
1112 }
1113 finish_failed_checkpoint_with(&tx, &record, "owner", deferred, "diagnostic".into())
1114 .unwrap();
1115 finish_failed_checkpoint_with(&tx, &record, "owner", deferred, "diagnostic".into())
1116 .unwrap();
1117 tx.commit().unwrap();
1118 let loaded = load_state_from(&path).unwrap();
1119 assert_eq!(loaded.sessions["session-1"].checkpoint, previous);
1120 assert_eq!(
1121 loaded.sessions["session-1"].last_checkpoint_error,
1122 if deferred {
1123 previous_warning
1124 } else {
1125 Some("archive verification failed".into())
1126 }
1127 );
1128 let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1129 .unwrap()
1130 .events;
1131 assert_eq!(events.len(), 1);
1132 assert!(
1133 matches!(&events[0].event, ApiEventData::CommandEnded { result } if result.outcome == if deferred { CommandResultKind::Rejected } else { CommandResultKind::Failed })
1134 );
1135 }
1136 }
1137}