Skip to main content

mj_controller/database/
events.rs

1//! Durable, ordered events for the native subagent API.
2use 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    // sqlite_sequence survives deletion of the last event; cursors never move back.
25    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
82/// Called on a blocking task; the shared writer serializes this with relay projection.
83pub 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
155/// Upgrade bodies in place: sequence identity and the AUTOINCREMENT frontier stay unchanged.
156pub(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
196/// Lifecycle writes own this comparison and event in the same database transaction.
197pub(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
219/// Reserve one requested checkpoint and save its lifecycle state in one writer decision.
220pub 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
291/// Consume identity and publish the terminal fact atomically; repeated finishes do nothing.
292pub(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
325/// Called before submission, so even a lost acknowledgement retains barrier correlation.
326pub 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(&current, &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        // Re-delivery of the committed frontier must not insert a duplicate.
482        apply_projection_event_to(
483            &path,
484            "session-1",
485            current.applied_event_ordinal,
486            "",
487            &current.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    // Hard-won: be5abcca: launch failure reasons were absent from the user-visible event stream.
603    #[test]
604    fn a_lifecycle_save_that_records_a_launch_failure_emits_one_error_event() {
605        // The launch-failure path persists through `save_lifecycle_session`
606        // (a plain UPDATE), so prove that path fires the fault transition once so
607        // `mj events` shows the reason exactly once, not zero or twice.
608        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        // An unchanged re-save must not add a second event.
617        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            // Drop rolls back both metadata and the terminal event.
1052        }
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}