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    #[test]
603    fn a_lifecycle_save_that_records_a_launch_failure_emits_one_error_event() {
604        // The launch-failure path persists through `save_lifecycle_session`
605        // (a plain UPDATE), so prove that path fires the fault transition once so
606        // `mj events` shows the reason exactly once, not zero or twice.
607        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        // An unchanged re-save must not add a second event.
616        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            // Drop rolls back both metadata and the terminal event.
1051        }
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}