Skip to main content

mj_controller/database/
events.rs

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