1use super::*;
3pub 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 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
103pub 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(¤t, &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 apply_projection_event_to(
330 &path,
331 "session-1",
332 current.applied_event_ordinal,
333 "",
334 ¤t.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 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 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}