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        let RelayObservation::CommandQueued {
553            command: RelayCommand::Prompt { prompt },
554            ..
555        } = &mut observations[0]
556        else {
557            unreachable!("turn fixture starts with a prompt")
558        };
559        prompt.push(agent_client_protocol::schema::v1::ContentBlock::from(
560            "Inspect the chosen directory",
561        ));
562        observations.splice(
563            2..2,
564            [
565                RelayObservation::ElicitationRequested {
566                    request: request.clone(),
567                },
568                RelayObservation::ElicitationResolved {
569                    elicitation_id: request.id.clone(),
570                    action: "accept".into(),
571                    reply: Some(request.reply_text(
572                        &mj_core::elicitation::ElicitationResponse::Accept {
573                            content: std::collections::BTreeMap::from([(
574                                "directory".into(),
575                                mj_core::elicitation::ElicitationValue::String("src".into()),
576                            )]),
577                        },
578                    )),
579                },
580            ],
581        );
582        page(&path, observations, false).unwrap();
583        record.last_error = Some("provisioning failed".into());
584        save_session_to(&path, &record).unwrap();
585        save_session_to(&path, &record).unwrap();
586        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
587        assert_eq!(
588            page.events
589                .iter()
590                .map(|e| e.event.kind())
591                .collect::<Vec<_>>(),
592            [
593                "turn_started",
594                "input_required",
595                "input_resolved",
596                "turn_ended",
597                "session_fault"
598            ]
599        );
600        let ApiEventData::InputRequired {
601            request: actual,
602            turn_id,
603        } = &page.events[1].event
604        else {
605            panic!("question event")
606        };
607        assert_eq!(actual.as_ref(), Some(&request));
608        assert_eq!(*turn_id, Some(1));
609        assert!(
610            matches!(&page.events[2].event, ApiEventData::InputResolved { turn_id: Some(1), action, .. } if action == "accept")
611        );
612        let current = load_materialized_session_from(&path, "session-1")
613            .unwrap()
614            .unwrap();
615        assert!(current.pending_elicitations.is_empty());
616        assert!(current.active_turn.is_none());
617        let reply = current
618            .transcript
619            .iter()
620            .find(|item| {
621                item.stable_id
622                    .starts_with(mj_core::transcript::ELICITATION_REPLY_ITEM_PREFIX)
623            })
624            .unwrap();
625        assert_eq!(
626            mj_transcript::transcript::transcript_item_text(reply),
627            "Which directory?\n\ndirectory: src"
628        );
629        assert!(!reply.is_turn_start());
630        let turn_start = current
631            .transcript
632            .iter()
633            .find(|item| item.is_turn_start())
634            .unwrap()
635            .position;
636        assert_eq!(
637            mj_core::state::latest_completed_turn_ordinal(&current),
638            Some(turn_start)
639        );
640        let connection = Connection::open(&path).unwrap();
641        assert_eq!(
642            super::super::materialized::last_materialized_turn_start(&connection, "session-1")
643                .unwrap(),
644            Some(turn_start)
645        );
646        assert_eq!(
647            super::super::materialized::last_materialized_user_message(&connection, "session-1")
648                .unwrap()
649                .unwrap()
650                .0,
651            turn_start
652        );
653        assert_eq!(current.last_turn_outcome.unwrap().accepted_ordinal, Some(1));
654    }
655
656    // Hard-won: be5abcca: launch failure reasons were absent from the user-visible event stream.
657    #[test]
658    fn a_lifecycle_save_that_records_a_launch_failure_emits_one_error_event() {
659        // The launch-failure path persists through `save_lifecycle_session`
660        // (a plain UPDATE), so prove that path fires the fault transition once so
661        // `mj events` shows the reason exactly once, not zero or twice.
662        let dir = tempfile::tempdir().unwrap();
663        let path = dir.path().join("events.sqlite");
664        let mut record = super::super::tests::session("session-1", "project-1");
665        save_session_to(&path, &record).unwrap();
666
667        record.state = mj_core::state::SessionState::Error;
668        record.last_error = Some("worker bootstrap failed: Connection closed by host".into());
669        save_lifecycle_session_to(&path, &record).unwrap();
670        // An unchanged re-save must not add a second event.
671        save_lifecycle_session_to(&path, &record).unwrap();
672
673        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
674        let errors: Vec<_> = page
675            .events
676            .iter()
677            .filter_map(|event| match &event.event {
678                ApiEventData::SessionFault { message, .. } => Some(message.clone()),
679                _ => None,
680            })
681            .collect();
682        assert_eq!(
683            errors,
684            vec!["worker bootstrap failed: Connection closed by host".to_owned()],
685            "a recorded launch failure surfaces as exactly one error event"
686        );
687    }
688
689    #[test]
690    fn api_events_preserve_command_identity_for_failed_completions() {
691        let dir = tempfile::tempdir().unwrap();
692        let path = dir.path().join("events.sqlite");
693        save_session_to(
694            &path,
695            &super::super::tests::session("session-1", "project-1"),
696        )
697        .unwrap();
698        let mut observations = turn();
699        let RelayObservation::CommandCompleted { outcome, .. } = observations.last_mut().unwrap()
700        else {
701            unreachable!()
702        };
703        *outcome = RelayCommandOutcome::Prompt {
704            diagnostic: None,
705            stop_reason: "provider_error".into(),
706            usage: None,
707        };
708        page(&path, observations, false).unwrap();
709        let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
710            .unwrap()
711            .events;
712        assert_eq!(events.len(), 2);
713        assert!(matches!(&events[1].event, ApiEventData::TurnEnded { turn }
714            if turn.command_id == "prompt-1" && turn.accepted_ordinal == Some(1)
715                && turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Failed
716                && turn.outcome.stop_reason.as_deref() == Some("provider_error")));
717    }
718
719    #[test]
720    fn api_activity_events_track_changes_without_repeating_or_reviving_forgotten_sessions() {
721        let dir = tempfile::tempdir().unwrap();
722        let path = dir.path().join("events.sqlite");
723        save_session_to(
724            &path,
725            &super::super::tests::session("session-1", "project-1"),
726        )
727        .unwrap();
728        let mut connection = open(&path).unwrap();
729        let idle = ApiActivityState {
730            state: "running".into(),
731            details: Some(ApiActivityDetails {
732                kind: ApiActivityKind::Idle,
733                turn_started_at_ms: None,
734                step_started_at_ms: None,
735                background_started_at_ms: None,
736                idle_since_ms: Some(100),
737                last_activity_at_ms: None,
738                label: None,
739            }),
740            is_idle: true,
741            waiting_for_input: false,
742            capacity_retry: false,
743        };
744        record_api_activities_with(
745            &mut connection,
746            vec![("session-1".into(), idle.clone())],
747            100,
748        )
749        .unwrap();
750        record_api_activities_with(
751            &mut connection,
752            vec![("session-1".into(), idle.clone())],
753            200,
754        )
755        .unwrap();
756        let mut background = idle.clone();
757        background.is_idle = false;
758        let details = background.details.as_mut().unwrap();
759        details.kind = ApiActivityKind::Background;
760        details.idle_since_ms = None;
761        details.background_started_at_ms = Some(250);
762        record_api_activities_with(
763            &mut connection,
764            vec![("session-1".into(), background.clone())],
765            250,
766        )
767        .unwrap();
768        background.waiting_for_input = true;
769        record_api_activities_with(
770            &mut connection,
771            vec![("session-1".into(), background.clone())],
772            300,
773        )
774        .unwrap();
775        let unknown = ApiActivityState {
776            state: "disconnected".into(),
777            details: None,
778            is_idle: false,
779            waiting_for_input: false,
780            capacity_retry: false,
781        };
782        record_api_activities_with(
783            &mut connection,
784            vec![("session-1".into(), unknown.clone())],
785            400,
786        )
787        .unwrap();
788        drop(connection);
789        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
790        assert_eq!(page.events.len(), 4);
791        assert_eq!(
792            page.events.last().unwrap().event,
793            ApiEventData::ActivityChanged { activity: unknown }
794        );
795        assert_eq!(
796            page.events[2].event,
797            ApiEventData::ActivityChanged {
798                activity: background
799            }
800        );
801        delete_session_from(&path, "session-1").unwrap();
802        record_api_activities_with(
803            &mut open(&path).unwrap(),
804            vec![("session-1".into(), idle)],
805            500,
806        )
807        .unwrap();
808        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(4), 100).unwrap();
809        assert!(page.events.is_empty());
810        assert_eq!(page.latest_seq, 4);
811    }
812    #[test]
813    fn checkpoint_cleanup_does_not_end_a_turn_or_hide_control_failure() {
814        use mj_core::event_outcome::TurnResultKind;
815        use mj_core::relay::RelayCommandKind;
816        let dir = tempfile::tempdir().unwrap();
817        let path = dir.path().join("events.sqlite");
818        save_session_to(
819            &path,
820            &super::super::tests::session("session-1", "project-1"),
821        )
822        .unwrap();
823        page(&path, turn().into_iter().take(2).collect(), false).unwrap();
824        for (id, reason) in [
825            ("arbitrary-one", OutcomeReason::ControllerDisconnected),
826            ("arbitrary-two", OutcomeReason::OwnerLostOnRestart),
827        ] {
828            page(
829                &path,
830                vec![RelayObservation::CommandInterrupted {
831                    command_id: id.into(),
832                    command: RelayCommandKind::BeginCheckpoint,
833                    reason: Some(reason),
834                    message: "diagnostic wording is not a contract".into(),
835                }],
836                false,
837            )
838            .unwrap();
839        }
840        page(
841            &path,
842            vec![RelayObservation::CommandRejected {
843                command_id: "arbitrary-three".into(),
844                command: RelayCommandKind::CompleteCheckpoint,
845                reason: Some(OutcomeReason::CommandFailed),
846                message: "Cannot save recovery floor".into(),
847            }],
848            false,
849        )
850        .unwrap();
851        let current = load_materialized_session_from(&path, "session-1")
852            .unwrap()
853            .unwrap();
854        assert!(current.active_turn.is_some());
855        assert!(current.last_turn_outcome.is_none());
856        assert!(current.transcript.iter().all(|item| !matches!(&item.body, TranscriptBody::System { text } if text.contains("diagnostic wording"))));
857        assert!(current.transcript.iter().any(|item| matches!(&item.body, TranscriptBody::System { text } if text.contains("Cannot save recovery floor"))));
858        page(&path, vec![turn().pop().unwrap()], false).unwrap();
859        let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
860            .unwrap()
861            .events;
862        assert_eq!(events.len(), 5);
863        for event in &events[1..3] {
864            assert!(
865                matches!(&event.event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Cancelled && result.command_kind == "begin_checkpoint")
866            );
867        }
868        assert!(
869            matches!(&events[3].event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Failed)
870        );
871        assert!(
872            matches!(&events[4].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == TurnResultKind::Completed)
873        );
874    }
875
876    #[test]
877    fn terminal_prompt_results_are_semantic_and_not_duplicate_error_events() {
878        use mj_core::event_outcome::TurnResultKind;
879        use mj_core::relay::RelayCommandKind;
880        for (terminal, expected) in [
881            (
882                RelayObservation::CommandCompleted {
883                    barrier_command_id: None,
884                    command_id: "prompt-1".into(),
885                    command: Some(RelayCommandKind::Prompt),
886                    outcome: RelayCommandOutcome::Prompt {
887                        stop_reason: "Cancelled".into(),
888                        usage: None,
889                        diagnostic: None,
890                    },
891                },
892                TurnResultKind::Cancelled,
893            ),
894            (
895                RelayObservation::CommandCompleted {
896                    barrier_command_id: None,
897                    command_id: "prompt-1".into(),
898                    command: Some(RelayCommandKind::Prompt),
899                    outcome: RelayCommandOutcome::Prompt {
900                        stop_reason: "QuotaLimit".into(),
901                        usage: None,
902                        diagnostic: None,
903                    },
904                },
905                TurnResultKind::Failed,
906            ),
907            (
908                RelayObservation::CommandRejected {
909                    command_id: "prompt-1".into(),
910                    command: RelayCommandKind::Prompt,
911                    reason: Some(OutcomeReason::AdmissionRejected),
912                    message: "unavailable".into(),
913                },
914                TurnResultKind::Rejected,
915            ),
916            (
917                RelayObservation::CommandInterrupted {
918                    command_id: "prompt-1".into(),
919                    command: RelayCommandKind::Prompt,
920                    reason: Some(OutcomeReason::RuntimeStopped),
921                    message: "runtime stopped".into(),
922                },
923                TurnResultKind::Interrupted,
924            ),
925        ] {
926            let dir = tempfile::tempdir().unwrap();
927            let path = dir.path().join("events.sqlite");
928            save_session_to(
929                &path,
930                &super::super::tests::session("session-1", "project-1"),
931            )
932            .unwrap();
933            let mut observations = turn();
934            *observations.last_mut().unwrap() = terminal;
935            page(&path, observations, false).unwrap();
936            let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
937                .unwrap()
938                .events;
939            assert_eq!(events.len(), 2);
940            assert!(
941                matches!(&events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == expected && turn.command_id == "prompt-1")
942            );
943        }
944    }
945
946    #[test]
947    fn typed_runtime_fault_does_not_rewrite_a_completed_turn() {
948        let dir = tempfile::tempdir().unwrap();
949        let path = dir.path().join("events.sqlite");
950        save_session_to(
951            &path,
952            &super::super::tests::session("session-1", "project-1"),
953        )
954        .unwrap();
955        page(&path, turn(), false).unwrap();
956        page(
957            &path,
958            vec![RelayObservation::SessionFault {
959                reason: OutcomeReason::RuntimeUnavailable,
960                message: "runtime exited".into(),
961            }],
962            false,
963        )
964        .unwrap();
965        let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
966            .unwrap()
967            .events;
968        assert_eq!(events.len(), 3);
969        assert!(
970            matches!(&events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Completed)
971        );
972        assert!(matches!(
973            &events[2].event,
974            ApiEventData::SessionFault {
975                reason: OutcomeReason::RuntimeUnavailable,
976                ..
977            }
978        ));
979    }
980
981    #[test]
982    fn checkpoint_recovery_consumes_attempt_identity_once_and_preserves_correlation() {
983        let dir = tempfile::tempdir().unwrap();
984        let path = dir.path().join("events.sqlite");
985        let mut record = super::super::tests::session("session-1", "project-1");
986        record.state = SessionState::Running;
987        save_session_to(&path, &record).unwrap();
988        record.state = SessionState::Checkpointing;
989        let mut connection = open(&path).unwrap();
990        {
991            let tx = connection.transaction().unwrap();
992            begin_checkpoint_operation_with(&tx, &record, "attempt-1").unwrap();
993            assert!(begin_checkpoint_operation_with(&tx, &record, "attempt-2").is_err());
994            tx.execute(
995                "UPDATE checkpoint_operations SET related_command_ids = '[\"barrier-1\"]'",
996                [],
997            )
998            .unwrap();
999            tx.commit().unwrap();
1000        }
1001        drop(connection);
1002        assert_eq!(
1003            recover_interrupted_checkpointing_sessions_to(&path, "2026-09-27T12:00:00Z").unwrap(),
1004            1
1005        );
1006        assert_eq!(
1007            recover_interrupted_checkpointing_sessions_to(&path, "2026-09-27T12:00:01Z").unwrap(),
1008            0
1009        );
1010        let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1011            .unwrap()
1012            .events;
1013        assert_eq!(events.len(), 1);
1014        assert!(
1015            matches!(&events[0].event, ApiEventData::CommandEnded { result }
1016            if result.command_id == "attempt-1" && result.related_command_ids == ["barrier-1"]
1017            && result.owner == CommandOwner::Daemon && result.reason == Some(OutcomeReason::ControllerRestarted))
1018        );
1019    }
1020
1021    #[test]
1022    fn event_migration_preserves_cursors_and_keeps_unknown_history_explicit() {
1023        let dir = tempfile::tempdir().unwrap();
1024        let path = dir.path().join("events.sqlite");
1025        save_session_to(
1026            &path,
1027            &super::super::tests::session("session-1", "project-1"),
1028        )
1029        .unwrap();
1030        page(&path, turn(), false).unwrap();
1031        let turn = load_materialized_session_from(&path, "session-1")
1032            .unwrap()
1033            .unwrap()
1034            .last_turn_outcome
1035            .unwrap();
1036        let connection = Connection::open(&path).unwrap();
1037        connection
1038            .execute(
1039                "UPDATE api_events SET body = ?1 WHERE seq = 2",
1040                [serde_json::json!({"type":"turn_ended", "data":{"turn":turn}}).to_string()],
1041            )
1042            .unwrap();
1043        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();
1044        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; ALTER TABLE mailbox_outbox DROP COLUMN failure; 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();
1045        drop(connection);
1046        forget_verified_schema(&path);
1047        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
1048        assert_eq!(page.latest_seq, 8);
1049        assert_eq!(page.next_after_seq, 8);
1050        assert_eq!(
1051            page.events
1052                .iter()
1053                .map(|event| event.seq)
1054                .collect::<Vec<_>>(),
1055            [1, 2, 8]
1056        );
1057        assert_eq!(page.events[2].recorded_at_ms, 1234);
1058        assert!(
1059            matches!(&page.events[2].event, ApiEventData::LegacyNotice { original_type, command_id: Some(id), message } if original_type == "error" && id == "arbitrary" && message == "old diagnostic")
1060        );
1061        assert!(
1062            matches!(&page.events[1].event, ApiEventData::TurnEnded { turn } if turn.outcome.kind == mj_core::event_outcome::TurnResultKind::Completed)
1063        );
1064    }
1065    #[test]
1066    fn requested_checkpoint_result_commits_with_metadata_and_only_for_its_owner() {
1067        let dir = tempfile::tempdir().unwrap();
1068        let path = dir.path().join("events.sqlite");
1069        let mut record = super::super::tests::session("session-1", "project-1");
1070        record.state = SessionState::Running;
1071        let previous_checkpoint = record.checkpoint.clone();
1072        save_session_to(&path, &record).unwrap();
1073        let mut connection = open(&path).unwrap();
1074        record.state = SessionState::Checkpointing;
1075        {
1076            let tx = connection.transaction().unwrap();
1077            begin_checkpoint_operation_with(&tx, &record, "owner").unwrap();
1078            tx.commit().unwrap();
1079        }
1080        record.state = SessionState::Running;
1081        record.checkpoint = Some(CheckpointMetadata {
1082            archive_path: dir.path().join("verified.tar"),
1083            sha256: "a".repeat(64),
1084            created_at: "2026-09-27T12:00:00Z".into(),
1085            event_frontier: 0,
1086        });
1087        {
1088            let tx = connection.transaction().unwrap();
1089            assert!(save_requested_checkpoint_with(&tx, &record, "stale-owner").is_err());
1090            finish_failed_checkpoint_with(
1091                &tx,
1092                &record,
1093                "stale-owner",
1094                false,
1095                "late failure".into(),
1096            )
1097            .unwrap();
1098            assert_eq!(
1099                tx.query_row("SELECT count(*) FROM api_events", [], |r| r
1100                    .get::<_, usize>(0))
1101                    .unwrap(),
1102                0
1103            );
1104            save_requested_checkpoint_with(&tx, &record, "owner").unwrap();
1105            // Drop rolls back both metadata and the terminal event.
1106        }
1107        assert_eq!(
1108            load_state_from(&path).unwrap().sessions["session-1"].checkpoint,
1109            previous_checkpoint
1110        );
1111        assert!(
1112            load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1113                .unwrap()
1114                .events
1115                .is_empty()
1116        );
1117        {
1118            let tx = connection.transaction().unwrap();
1119            save_requested_checkpoint_with(&tx, &record, "owner").unwrap();
1120            finish_failed_checkpoint_with(
1121                &tx,
1122                &record,
1123                "owner",
1124                false,
1125                "late cleanup failure".into(),
1126            )
1127            .unwrap();
1128            tx.commit().unwrap();
1129        }
1130        let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1131            .unwrap()
1132            .events;
1133        assert_eq!(events.len(), 1);
1134        assert!(
1135            matches!(&events[0].event, ApiEventData::CommandEnded { result } if result.outcome == CommandResultKind::Succeeded && result.command_id == "owner")
1136        );
1137        assert_eq!(
1138            load_state_from(&path).unwrap().sessions["session-1"].checkpoint,
1139            record.checkpoint
1140        );
1141    }
1142
1143    #[test]
1144    fn checkpoint_failure_and_deferral_are_distinct_and_preserve_previous_archive() {
1145        for deferred in [false, true] {
1146            let dir = tempfile::tempdir().unwrap();
1147            let path = dir.path().join("events.sqlite");
1148            let mut record = super::super::tests::session("session-1", "project-1");
1149            record.state = SessionState::Running;
1150            record.checkpoint = Some(CheckpointMetadata {
1151                archive_path: dir.path().join("previous.tar"),
1152                sha256: "a".repeat(64),
1153                created_at: "2026-09-27T12:00:00Z".into(),
1154                event_frontier: 0,
1155            });
1156            save_session_to(&path, &record).unwrap();
1157            let previous = record.checkpoint.clone();
1158            let previous_warning = record.last_checkpoint_error.clone();
1159            let mut connection = open(&path).unwrap();
1160            let tx = connection.transaction().unwrap();
1161            record.state = SessionState::Checkpointing;
1162            begin_checkpoint_operation_with(&tx, &record, "owner").unwrap();
1163            record.state = SessionState::Running;
1164            if !deferred {
1165                record.last_checkpoint_error = Some("archive verification failed".into());
1166            }
1167            finish_failed_checkpoint_with(&tx, &record, "owner", deferred, "diagnostic".into())
1168                .unwrap();
1169            finish_failed_checkpoint_with(&tx, &record, "owner", deferred, "diagnostic".into())
1170                .unwrap();
1171            tx.commit().unwrap();
1172            let loaded = load_state_from(&path).unwrap();
1173            assert_eq!(loaded.sessions["session-1"].checkpoint, previous);
1174            assert_eq!(
1175                loaded.sessions["session-1"].last_checkpoint_error,
1176                if deferred {
1177                    previous_warning
1178                } else {
1179                    Some("archive verification failed".into())
1180                }
1181            );
1182            let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
1183                .unwrap()
1184                .events;
1185            assert_eq!(events.len(), 1);
1186            assert!(
1187                matches!(&events[0].event, ApiEventData::CommandEnded { result } if result.outcome == if deferred { CommandResultKind::Rejected } else { CommandResultKind::Failed })
1188            );
1189        }
1190    }
1191}