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 api_events_preserve_command_identity_for_failed_completions() {
368        let dir = tempfile::tempdir().unwrap();
369        let path = dir.path().join("events.sqlite");
370        save_session_to(
371            &path,
372            &super::super::tests::session("session-1", "project-1"),
373        )
374        .unwrap();
375        let mut observations = turn();
376        let RelayObservation::CommandCompleted { outcome, .. } = observations.last_mut().unwrap()
377        else {
378            unreachable!()
379        };
380        *outcome = RelayCommandOutcome::Prompt {
381            diagnostic: None,
382            stop_reason: "provider_error".into(),
383            usage: None,
384        };
385        page(&path, observations, false).unwrap();
386        let events = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100)
387            .unwrap()
388            .events;
389        assert!(
390            matches!(&events[1].event, ApiEventData::Error { command_id: Some(id), message } if id == "prompt-1" && message == "provider_error")
391        );
392        assert!(
393            matches!(&events[2].event, ApiEventData::TurnEnded { turn } if turn.accepted_ordinal == Some(1))
394        );
395    }
396
397    #[test]
398    fn api_activity_events_track_changes_without_repeating_or_reviving_forgotten_sessions() {
399        let dir = tempfile::tempdir().unwrap();
400        let path = dir.path().join("events.sqlite");
401        save_session_to(
402            &path,
403            &super::super::tests::session("session-1", "project-1"),
404        )
405        .unwrap();
406        let mut connection = open(&path).unwrap();
407        let idle = ApiActivityState {
408            state: "running".into(),
409            details: Some(ApiActivityDetails {
410                kind: ApiActivityKind::Idle,
411                turn_started_at_ms: None,
412                step_started_at_ms: None,
413                background_started_at_ms: None,
414                idle_since_ms: Some(100),
415                label: None,
416            }),
417            is_idle: true,
418            waiting_for_input: false,
419            capacity_retry: false,
420        };
421        record_api_activities_with(
422            &mut connection,
423            vec![("session-1".into(), idle.clone())],
424            100,
425        )
426        .unwrap();
427        record_api_activities_with(
428            &mut connection,
429            vec![("session-1".into(), idle.clone())],
430            200,
431        )
432        .unwrap();
433        let mut background = idle.clone();
434        background.is_idle = false;
435        let details = background.details.as_mut().unwrap();
436        details.kind = ApiActivityKind::Background;
437        details.idle_since_ms = None;
438        details.background_started_at_ms = Some(250);
439        record_api_activities_with(
440            &mut connection,
441            vec![("session-1".into(), background.clone())],
442            250,
443        )
444        .unwrap();
445        background.waiting_for_input = true;
446        record_api_activities_with(
447            &mut connection,
448            vec![("session-1".into(), background.clone())],
449            300,
450        )
451        .unwrap();
452        let unknown = ApiActivityState {
453            state: "disconnected".into(),
454            details: None,
455            is_idle: false,
456            waiting_for_input: false,
457            capacity_retry: false,
458        };
459        record_api_activities_with(
460            &mut connection,
461            vec![("session-1".into(), unknown.clone())],
462            400,
463        )
464        .unwrap();
465        drop(connection);
466        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(0), 100).unwrap();
467        assert_eq!(page.events.len(), 4);
468        assert_eq!(
469            page.events.last().unwrap().event,
470            ApiEventData::ActivityChanged { activity: unknown }
471        );
472        assert_eq!(
473            page.events[2].event,
474            ApiEventData::ActivityChanged {
475                activity: background
476            }
477        );
478        delete_session_from(&path, "session-1").unwrap();
479        record_api_activities_with(
480            &mut open(&path).unwrap(),
481            vec![("session-1".into(), idle)],
482            500,
483        )
484        .unwrap();
485        let page = load_api_events_from(&path, &ApiEventFilter::default(), Some(4), 100).unwrap();
486        assert!(page.events.is_empty());
487        assert_eq!(page.latest_seq, 4);
488    }
489}