1use 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 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
80pub 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(¤t, &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 apply_projection_event_to(
247 &path,
248 "session-1",
249 current.applied_event_ordinal,
250 "",
251 ¤t.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 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 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}