Skip to main content

mj_controller/database/
events.rs

1//! Durable, ordered events for the native subagent API.
2use super::*;
3#[cfg(test)]
4use mj_core::elicitation::ElicitationRequest;
5
6pub fn load_api_events(
7    filter: &ApiEventFilter,
8    after_seq: Option<u64>,
9    limit: usize,
10) -> Result<ApiEventPage> {
11    load_api_events_from(&database_path(), filter, after_seq, limit)
12}
13
14pub(super) fn load_api_events_from(
15    path: &Path,
16    filter: &ApiEventFilter,
17    after_seq: Option<u64>,
18    limit: usize,
19) -> Result<ApiEventPage> {
20    let mut connection = open_reader(path)?;
21    let tx = connection.transaction()?;
22    // sqlite_sequence survives deletion of the last event; cursors never move back.
23    let latest_seq: u64 = tx.query_row(
24        "SELECT COALESCE((SELECT seq FROM sqlite_sequence WHERE name = 'api_events'), 0)",
25        [],
26        |r| r.get(0),
27    )?;
28    let after_seq = after_seq.unwrap_or(latest_seq);
29    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")?;
30    let events = statement
31        .query_map(
32            params![
33                after_seq,
34                filter.session_id,
35                filter.workspace_id,
36                limit.clamp(1, 1000) as i64
37            ],
38            |r| {
39                Ok((
40                    r.get::<_, u64>(0)?,
41                    r.get::<_, String>(1)?,
42                    r.get::<_, i64>(2)?,
43                    r.get::<_, String>(3)?,
44                ))
45            },
46        )?
47        .map(|r| {
48            let (seq, session_id, recorded_at_ms, body) = r?;
49            Ok(ApiEvent {
50                seq,
51                session_id,
52                recorded_at_ms,
53                event: serde_json::from_str(&body)?,
54            })
55        })
56        .collect::<Result<Vec<_>>>()?;
57    let next_after_seq = events
58        .last()
59        .map_or(latest_seq.max(after_seq), |event| event.seq);
60    Ok(ApiEventPage {
61        events,
62        next_after_seq,
63        latest_seq,
64    })
65}
66
67pub(super) fn insert_api_event(
68    tx: &Transaction<'_>,
69    session_id: &str,
70    recorded_at_ms: i64,
71    event: &ApiEventData,
72) -> Result<()> {
73    tx.execute(
74        "INSERT INTO api_events(session_id, recorded_at_ms, body) VALUES (?1, ?2, ?3)",
75        params![session_id, recorded_at_ms, serde_json::to_string(event)?],
76    )?;
77    Ok(())
78}
79
80/// Called on a blocking task; the shared writer serializes this with relay projection.
81pub fn record_api_activities(
82    activities: Vec<(String, ApiActivityState)>,
83    recorded_at_ms: i64,
84) -> Result<()> {
85    submit_database_write("record_api_activities", move |connection| {
86        record_api_activities_with(connection, activities, recorded_at_ms)
87    })
88}
89
90fn record_api_activities_with(
91    connection: &mut Connection,
92    activities: Vec<(String, ApiActivityState)>,
93    recorded_at_ms: i64,
94) -> Result<()> {
95    let tx = connection.transaction()?;
96    for (session_id, activity) in activities {
97        let exists: bool = tx.query_row(
98            "SELECT EXISTS(SELECT 1 FROM sessions WHERE session_id = ?1)",
99            [&session_id],
100            |r| r.get(0),
101        )?;
102        if !exists {
103            continue;
104        }
105        let body = serde_json::to_string(&activity)?;
106        let previous: Option<String> = tx
107            .query_row(
108                "SELECT body FROM api_session_activity WHERE session_id = ?1",
109                [&session_id],
110                |r| r.get(0),
111            )
112            .optional()?;
113        if previous.as_deref() == Some(body.as_str()) {
114            continue;
115        }
116        insert_api_event(
117            &tx,
118            &session_id,
119            recorded_at_ms,
120            &ApiEventData::ActivityChanged { activity },
121        )?;
122        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])?;
123    }
124    tx.commit()?;
125    Ok(())
126}
127
128pub fn record_api_error(session_id: String, message: String) -> Result<()> {
129    submit_database_write("record_api_error", move |connection| {
130        let tx = connection.transaction()?;
131        let exists: bool = tx.query_row(
132            "SELECT EXISTS(SELECT 1 FROM sessions WHERE session_id = ?1)",
133            [&session_id],
134            |r| r.get(0),
135        )?;
136        if exists {
137            insert_api_event(
138                &tx,
139                &session_id,
140                chrono::Utc::now().timestamp_millis(),
141                &ApiEventData::Error {
142                    message,
143                    command_id: None,
144                },
145            )?;
146        }
147        tx.commit()?;
148        Ok(())
149    })
150}
151
152#[cfg(test)]
153mod tests {
154    use super::*;
155    use mj_core::relay::{
156        RelayCommand, RelayCommandOutcome, RelayEvent, RelayObservation, relay_event_digest,
157    };
158    use mj_transcript::projection::{apply_committed_projection_event, project_relay_event};
159
160    fn page(path: &Path, observations: Vec<RelayObservation>, fail: bool) -> Result<()> {
161        let mut current = load_materialized_session_from(path, "session-1")?.unwrap();
162        apply_projection_page_to(path, "session-1", |page| {
163            for observation in observations {
164                let mut event = RelayEvent {
165                    format: mj_core::relay::RELAY_EVENT_FORMAT_V1,
166                    ordinal: current.applied_event_ordinal + 1,
167                    previous_digest: current.applied_event_digest.clone(),
168                    digest: String::new(),
169                    recorded_at_ms: 100,
170                    command_id: None,
171                    observation,
172                };
173                event.digest = relay_event_digest(&event)?;
174                let mutation = project_relay_event(&current, &event)?.mutation;
175                page.apply(
176                    event.ordinal,
177                    &event.previous_digest,
178                    &event.digest,
179                    &mutation,
180                )?;
181                apply_committed_projection_event(&mut current, &event, mutation)?;
182            }
183            if fail {
184                bail!("injected rollback");
185            }
186            Ok(())
187        })
188    }
189
190    fn turn() -> Vec<RelayObservation> {
191        vec![
192            RelayObservation::CommandQueued {
193                command_id: "prompt-1".into(),
194                command: RelayCommand::Prompt { prompt: vec![] },
195                created_at_ms: 1,
196            },
197            RelayObservation::CommandStarted {
198                command_id: "prompt-1".into(),
199                started_at_ms: 2,
200            },
201            RelayObservation::CommandCompleted {
202                command_id: "prompt-1".into(),
203                outcome: RelayCommandOutcome::Prompt {
204                    diagnostic: None,
205                    stop_reason: "end_turn".into(),
206                    usage: None,
207                },
208            },
209        ]
210    }
211
212    #[test]
213    fn api_events_survive_coalescing_rollback_and_reopen() {
214        let dir = tempfile::tempdir().unwrap();
215        let path = dir.path().join("events.sqlite");
216        save_session_to(
217            &path,
218            &super::super::tests::session("session-1", "project-1"),
219        )
220        .unwrap();
221        assert!(page(&path, turn(), true).is_err());
222        assert!(
223            load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
224                .unwrap()
225                .events
226                .is_empty()
227        );
228        page(&path, turn(), false).unwrap();
229        let first = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 1).unwrap();
230        assert_eq!(first.events.len(), 1);
231        assert_eq!(first.events[0].event.kind(), "turn_started");
232        let rest = load_api_events_from(
233            &path,
234            &ApiEventFilter::default(),
235            Some(first.next_after_seq),
236            100,
237        )
238        .unwrap();
239        assert_eq!(rest.events.len(), 1);
240        assert_eq!(rest.events[0].event.kind(), "turn_ended");
241        let current = load_materialized_session_from(&path, "session-1")
242            .unwrap()
243            .unwrap();
244        assert!(current.active_turn.is_none());
245        // Re-delivery of the committed frontier must not insert a duplicate.
246        apply_projection_event_to(
247            &path,
248            "session-1",
249            current.applied_event_ordinal,
250            "",
251            &current.applied_event_digest,
252            &MaterializedSessionMutation::default(),
253        )
254        .unwrap();
255        assert_eq!(
256            load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
257                .unwrap()
258                .events
259                .len(),
260            2
261        );
262        assert!(
263            load_api_events_from(&path, &ApiEventFilter::default(), None, 100)
264                .unwrap()
265                .events
266                .is_empty()
267        );
268    }
269
270    #[test]
271    fn api_events_filter_and_keep_cursors_after_forgetting_sessions() {
272        let dir = tempfile::tempdir().unwrap();
273        let path = dir.path().join("events.sqlite");
274        save_session_to(
275            &path,
276            &super::super::tests::session("session-1", "project-1"),
277        )
278        .unwrap();
279        page(&path, turn(), false).unwrap();
280        let filter = ApiEventFilter {
281            session_id: Some("other".into()),
282            workspace_id: None,
283        };
284        let empty = load_api_events_from(&path, &filter, Some(0), 100).unwrap();
285        assert!(empty.events.is_empty());
286        assert_eq!(empty.next_after_seq, 2);
287        let filter = ApiEventFilter {
288            session_id: None,
289            workspace_id: Some("default".into()),
290        };
291        assert_eq!(
292            load_api_events_from(&path, &filter, Some(0), 100)
293                .unwrap()
294                .events
295                .len(),
296            2
297        );
298        delete_session_from(&path, "session-1").unwrap();
299        let deleted =
300            load_api_events_from(&path, &ApiEventFilter::default(), Some(2), 100).unwrap();
301        assert!(deleted.events.is_empty());
302        assert_eq!(deleted.latest_seq, 2);
303    }
304
305    #[test]
306    fn api_events_capture_free_text_questions_resolution_and_errors() {
307        let dir = tempfile::tempdir().unwrap();
308        let path = dir.path().join("events.sqlite");
309        let mut record = super::super::tests::session("session-1", "project-1");
310        save_session_to(&path, &record).unwrap();
311        let request = ElicitationRequest::from_acp_params("question-1", serde_json::json!({
312            "mode": "form", "sessionId": "session-1", "message": "Which directory?",
313            "requestedSchema": {"type": "object", "properties": {"directory": {"type": "string"}}, "required": ["directory"]}
314        })).unwrap();
315        let mut observations = turn();
316        observations.splice(
317            2..2,
318            [
319                RelayObservation::ElicitationRequested {
320                    request: request.clone(),
321                },
322                RelayObservation::ElicitationResolved {
323                    elicitation_id: request.id.clone(),
324                    action: "accept".into(),
325                },
326            ],
327        );
328        page(&path, observations, false).unwrap();
329        record.last_error = Some("provisioning failed".into());
330        save_session_to(&path, &record).unwrap();
331        save_session_to(&path, &record).unwrap();
332        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
333        assert_eq!(
334            page.events
335                .iter()
336                .map(|e| e.event.kind())
337                .collect::<Vec<_>>(),
338            [
339                "turn_started",
340                "input_required",
341                "input_resolved",
342                "turn_ended",
343                "error"
344            ]
345        );
346        let ApiEventData::InputRequired {
347            request: actual,
348            turn_id,
349        } = &page.events[1].event
350        else {
351            panic!("question event")
352        };
353        assert_eq!(actual, &request);
354        assert_eq!(*turn_id, Some(1));
355        assert!(
356            matches!(&page.events[2].event, ApiEventData::InputResolved { turn_id: Some(1), action, .. } if action == "accept")
357        );
358        let current = load_materialized_session_from(&path, "session-1")
359            .unwrap()
360            .unwrap();
361        assert!(current.pending_elicitations.is_empty());
362        assert!(current.active_turn.is_none());
363        assert_eq!(current.last_turn_outcome.unwrap().accepted_ordinal, Some(1));
364    }
365
366    #[test]
367    fn a_lifecycle_save_that_records_a_launch_failure_emits_one_error_event() {
368        // The launch-failure path persists through `save_lifecycle_session`
369        // (a plain UPDATE), so prove that path fires the error trigger once so
370        // `mj events` shows the reason exactly once, not zero or twice.
371        let dir = tempfile::tempdir().unwrap();
372        let path = dir.path().join("events.sqlite");
373        let mut record = super::super::tests::session("session-1", "project-1");
374        save_session_to(&path, &record).unwrap();
375
376        record.state = mj_core::state::SessionState::Error;
377        record.last_error = Some("worker bootstrap failed: Connection closed by host".into());
378        save_lifecycle_session_to(&path, &record).unwrap();
379        // An unchanged re-save must not add a second event.
380        save_lifecycle_session_to(&path, &record).unwrap();
381
382        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
383        let errors: Vec<_> = page
384            .events
385            .iter()
386            .filter_map(|event| match &event.event {
387                ApiEventData::Error { message, .. } => Some(message.clone()),
388                _ => None,
389            })
390            .collect();
391        assert_eq!(
392            errors,
393            vec!["worker bootstrap failed: Connection closed by host".to_owned()],
394            "a recorded launch failure surfaces as exactly one error event"
395        );
396    }
397
398    #[test]
399    fn api_events_preserve_command_identity_for_failed_completions() {
400        let dir = tempfile::tempdir().unwrap();
401        let path = dir.path().join("events.sqlite");
402        save_session_to(
403            &path,
404            &super::super::tests::session("session-1", "project-1"),
405        )
406        .unwrap();
407        let mut observations = turn();
408        let RelayObservation::CommandCompleted { outcome, .. } = observations.last_mut().unwrap()
409        else {
410            unreachable!()
411        };
412        *outcome = RelayCommandOutcome::Prompt {
413            diagnostic: None,
414            stop_reason: "provider_error".into(),
415            usage: None,
416        };
417        page(&path, observations, false).unwrap();
418        let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
419            .unwrap()
420            .events;
421        assert!(
422            matches!(&events[1].event, ApiEventData::Error { command_id: Some(id), message } if id == "prompt-1" && message == "provider_error")
423        );
424        assert!(
425            matches!(&events[2].event, ApiEventData::TurnEnded { turn } if turn.accepted_ordinal == Some(1))
426        );
427    }
428
429    #[test]
430    fn api_activity_events_track_changes_without_repeating_or_reviving_forgotten_sessions() {
431        let dir = tempfile::tempdir().unwrap();
432        let path = dir.path().join("events.sqlite");
433        save_session_to(
434            &path,
435            &super::super::tests::session("session-1", "project-1"),
436        )
437        .unwrap();
438        let mut connection = open(&path).unwrap();
439        let idle = ApiActivityState {
440            state: "running".into(),
441            details: Some(ApiActivityDetails {
442                kind: ApiActivityKind::Idle,
443                turn_started_at_ms: None,
444                step_started_at_ms: None,
445                background_started_at_ms: None,
446                idle_since_ms: Some(100),
447                label: None,
448            }),
449            is_idle: true,
450            waiting_for_input: false,
451            capacity_retry: false,
452        };
453        record_api_activities_with(
454            &mut connection,
455            vec![("session-1".into(), idle.clone())],
456            100,
457        )
458        .unwrap();
459        record_api_activities_with(
460            &mut connection,
461            vec![("session-1".into(), idle.clone())],
462            200,
463        )
464        .unwrap();
465        let mut background = idle.clone();
466        background.is_idle = false;
467        let details = background.details.as_mut().unwrap();
468        details.kind = ApiActivityKind::Background;
469        details.idle_since_ms = None;
470        details.background_started_at_ms = Some(250);
471        record_api_activities_with(
472            &mut connection,
473            vec![("session-1".into(), background.clone())],
474            250,
475        )
476        .unwrap();
477        background.waiting_for_input = true;
478        record_api_activities_with(
479            &mut connection,
480            vec![("session-1".into(), background.clone())],
481            300,
482        )
483        .unwrap();
484        let unknown = ApiActivityState {
485            state: "disconnected".into(),
486            details: None,
487            is_idle: false,
488            waiting_for_input: false,
489            capacity_retry: false,
490        };
491        record_api_activities_with(
492            &mut connection,
493            vec![("session-1".into(), unknown.clone())],
494            400,
495        )
496        .unwrap();
497        drop(connection);
498        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
499        assert_eq!(page.events.len(), 4);
500        assert_eq!(
501            page.events.last().unwrap().event,
502            ApiEventData::ActivityChanged { activity: unknown }
503        );
504        assert_eq!(
505            page.events[2].event,
506            ApiEventData::ActivityChanged {
507                activity: background
508            }
509        );
510        delete_session_from(&path, "session-1").unwrap();
511        record_api_activities_with(
512            &mut open(&path).unwrap(),
513            vec![("session-1".into(), idle)],
514            500,
515        )
516        .unwrap();
517        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(4), 100).unwrap();
518        assert!(page.events.is_empty());
519        assert_eq!(page.latest_seq, 4);
520    }
521}