1use std::collections::HashMap;
4use std::path::{Path, PathBuf};
5use std::sync::atomic::{AtomicU64, Ordering};
6use std::sync::{Arc, Weak};
7
8use async_trait::async_trait;
9use everruns_core::error::{AgentLoopError, Result as CoreResult};
10use everruns_core::events::{Event, EventData, EventRequest, OutputMessageCompletedData};
11use everruns_core::message::{ContentPart, Message};
12use everruns_core::message_filter::{MessageFilter, MessageQuery};
13use everruns_core::message_retriever::{MessageHistory, MessageRetriever};
14use everruns_core::tools::ToolResultImage;
15use everruns_core::traits::EventEmitter;
16use everruns_core::typed_id::{EventId, MessageId, SessionId};
17use serde::{Deserialize, Serialize};
18use tokio::io::AsyncWriteExt;
19use tokio::sync::{Mutex, RwLock};
20
21pub const DEFAULT_EVENT_READ_LIMIT: usize = 256;
23pub const MAX_EVENT_PAGE_SIZE: usize = 1024;
25pub const MAX_EVENT_HISTORY_REPLAY: usize = 100_000;
27pub const MAX_EVENT_HISTORY_PAGE_SIZE: usize = 256;
29
30#[derive(Clone, Copy, Debug, PartialEq, Eq)]
32pub enum EventDurability {
33 Volatile,
35 CrashDurable,
37}
38
39#[derive(Clone, Debug, thiserror::Error, PartialEq, Eq)]
41#[non_exhaustive]
42pub enum EventLogError {
43 #[error("invalid event read: {detail}")]
45 InvalidRead { detail: String },
46 #[error("event cursor belongs to another session: {detail}")]
48 CrossSessionCursor { detail: String },
49 #[error("incompatible event cursor: {detail}")]
51 IncompatibleCursor { detail: String },
52 #[error("expired event cursor: {detail}")]
54 ExpiredCursor { detail: String },
55 #[error("invalid event append: {detail}")]
57 InvalidAppend { detail: String },
58 #[error("event log corruption: {detail}")]
60 Corruption { detail: String },
61 #[error("event log backend failure: {detail}")]
63 Backend { detail: String },
64}
65
66impl From<std::io::Error> for EventLogError {
67 fn from(error: std::io::Error) -> Self {
68 Self::Backend {
69 detail: error.to_string(),
70 }
71 }
72}
73
74#[derive(Clone, Copy, Debug, PartialEq, Eq)]
76pub struct EventReadLimit(u16);
77
78impl EventReadLimit {
79 pub fn new(limit: usize) -> Result<Self, EventLogError> {
81 if limit == 0 || limit > MAX_EVENT_PAGE_SIZE {
82 return Err(EventLogError::InvalidRead {
83 detail: format!("limit must be between 1 and {MAX_EVENT_PAGE_SIZE}, got {limit}"),
84 });
85 }
86 Ok(Self(limit as u16))
87 }
88
89 pub fn get(self) -> usize {
91 self.0 as usize
92 }
93}
94
95impl Default for EventReadLimit {
96 fn default() -> Self {
97 Self(DEFAULT_EVENT_READ_LIMIT as u16)
98 }
99}
100
101#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
103pub struct EventCursor {
104 session_id: SessionId,
105 after_sequence: i32,
106 snapshot_high_watermark: Option<i32>,
107}
108
109impl EventCursor {
110 pub fn session_id(&self) -> SessionId {
112 self.session_id
113 }
114
115 pub fn after_sequence(&self) -> i32 {
117 self.after_sequence
118 }
119
120 pub fn snapshot_high_watermark(&self) -> Option<i32> {
122 self.snapshot_high_watermark
123 }
124
125 pub fn after(session_id: SessionId, after_sequence: i32) -> Result<Self, EventLogError> {
130 if after_sequence < 0 {
131 return Err(EventLogError::InvalidRead {
132 detail: "poll cursor sequence cannot be negative".into(),
133 });
134 }
135 Ok(Self {
136 session_id,
137 after_sequence,
138 snapshot_high_watermark: None,
139 })
140 }
141}
142
143#[derive(Clone, Debug)]
145pub struct EventReadRequest {
146 session_id: SessionId,
147 cursor: Option<EventCursor>,
148 limit: EventReadLimit,
149}
150
151impl EventReadRequest {
152 pub fn new(session_id: SessionId, limit: EventReadLimit) -> Self {
154 Self {
155 session_id,
156 cursor: None,
157 limit,
158 }
159 }
160
161 pub fn from_cursor(cursor: EventCursor, limit: EventReadLimit) -> Self {
163 Self {
164 session_id: cursor.session_id,
165 cursor: Some(cursor),
166 limit,
167 }
168 }
169
170 pub fn with_cursor(mut self, cursor: EventCursor) -> Self {
174 self.cursor = Some(cursor);
175 self
176 }
177
178 pub fn session_id(&self) -> SessionId {
180 self.session_id
181 }
182
183 pub fn limit(&self) -> EventReadLimit {
185 self.limit
186 }
187}
188
189#[derive(Clone, Debug)]
191pub struct EventPage {
192 pub events: Vec<Event>,
194 pub next_cursor: Option<EventCursor>,
196 snapshot_high_watermark: i32,
197}
198
199#[derive(Clone, Copy, Debug, PartialEq, Eq)]
201pub struct EventHistoryReadLimit(u16);
202
203impl EventHistoryReadLimit {
204 pub fn new(limit: usize) -> Result<Self, EventLogError> {
207 if limit == 0 || limit > MAX_EVENT_HISTORY_PAGE_SIZE {
208 return Err(EventLogError::InvalidRead {
209 detail: format!(
210 "history limit must be between 1 and {MAX_EVENT_HISTORY_PAGE_SIZE}, got {limit}"
211 ),
212 });
213 }
214 Ok(Self(limit as u16))
215 }
216
217 pub fn get(self) -> usize {
219 self.0 as usize
220 }
221}
222
223impl Default for EventHistoryReadLimit {
224 fn default() -> Self {
225 Self(MAX_EVENT_HISTORY_PAGE_SIZE as u16)
226 }
227}
228
229#[derive(Clone, Debug)]
231pub struct EventHistoryReadRequest {
232 session_id: SessionId,
233 cursor: Option<EventCursor>,
234 limit: EventHistoryReadLimit,
235}
236
237impl EventHistoryReadRequest {
238 pub fn new(session_id: SessionId, limit: EventHistoryReadLimit) -> Self {
240 Self {
241 session_id,
242 cursor: None,
243 limit,
244 }
245 }
246
247 pub fn with_cursor(mut self, cursor: EventCursor) -> Self {
249 self.cursor = Some(cursor);
250 self
251 }
252
253 pub fn session_id(&self) -> SessionId {
255 self.session_id
256 }
257
258 pub fn limit(&self) -> EventHistoryReadLimit {
260 self.limit
261 }
262}
263
264#[derive(Clone, Debug)]
266pub struct EventHistoryPage {
267 pub messages: Vec<Message>,
269 pub next_cursor: Option<EventCursor>,
271 snapshot_high_watermark: i32,
272}
273
274impl EventHistoryPage {
275 pub fn snapshot_high_watermark(&self) -> i32 {
277 self.snapshot_high_watermark
278 }
279}
280
281impl EventPage {
282 pub fn snapshot_high_watermark(&self) -> i32 {
284 self.snapshot_high_watermark
285 }
286}
287
288#[async_trait]
290pub trait EventReader: Send + Sync {
291 async fn read_page(&self, request: EventReadRequest) -> Result<EventPage, EventLogError>;
294}
295
296#[async_trait]
302pub trait EventLog: EventReader {
303 async fn append(&self, request: EventRequest) -> Result<Event, EventLogError>;
305
306 fn durability(&self) -> EventDurability;
308}
309
310pub trait EventSink: Send + Sync {
312 fn try_send(&self, event: Event) -> Result<(), EventSinkError>;
314}
315
316#[derive(Clone, Copy, Debug, thiserror::Error, PartialEq, Eq)]
318pub enum EventSinkError {
319 #[error("event sink is full")]
321 Full {
322 dropped: u64,
324 },
325 #[error("event sink is closed")]
327 Closed,
328}
329
330#[derive(Clone, Copy, Debug, Default)]
332pub struct NoopEventSink;
333
334impl EventSink for NoopEventSink {
335 fn try_send(&self, _event: Event) -> Result<(), EventSinkError> {
336 Ok(())
337 }
338}
339
340#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
342pub struct EventDeliveryStats {
343 pub full: u64,
345 pub closed: u64,
347}
348
349#[derive(Clone)]
356pub struct HostEventEmitter {
357 log: Arc<dyn EventLog>,
358 sink: Arc<dyn EventSink>,
359 session_locks: Arc<Mutex<HashMap<SessionId, Weak<Mutex<()>>>>>,
360 full: Arc<AtomicU64>,
361 closed: Arc<AtomicU64>,
362}
363
364impl HostEventEmitter {
365 pub fn new(log: Arc<dyn EventLog>, sink: Arc<dyn EventSink>) -> Self {
367 Self {
368 log,
369 sink,
370 session_locks: Arc::new(Mutex::new(HashMap::new())),
371 full: Arc::new(AtomicU64::new(0)),
372 closed: Arc::new(AtomicU64::new(0)),
373 }
374 }
375
376 pub fn event_log(&self) -> Arc<dyn EventLog> {
378 self.log.clone()
379 }
380
381 pub fn delivery_stats(&self) -> EventDeliveryStats {
383 EventDeliveryStats {
384 full: self.full.load(Ordering::Relaxed),
385 closed: self.closed.load(Ordering::Relaxed),
386 }
387 }
388
389 async fn session_lock(&self, session_id: SessionId) -> Arc<Mutex<()>> {
390 let mut locks = self.session_locks.lock().await;
391 locks.retain(|_, lock| lock.strong_count() > 0);
392 if let Some(lock) = locks.get(&session_id).and_then(Weak::upgrade) {
393 return lock;
394 }
395 let lock = Arc::new(Mutex::new(()));
396 locks.insert(session_id, Arc::downgrade(&lock));
397 lock
398 }
399
400 fn notify(&self, event: Event) {
401 match self.sink.try_send(event) {
402 Ok(()) => {}
403 Err(EventSinkError::Full { dropped }) => {
404 self.full.fetch_add(dropped, Ordering::Relaxed);
405 tracing::debug!("live event sink full; canonical append remains committed");
406 }
407 Err(EventSinkError::Closed) => {
408 self.closed.fetch_add(1, Ordering::Relaxed);
409 tracing::debug!("live event sink closed; canonical append remains committed");
410 }
411 }
412 }
413}
414
415#[async_trait]
416impl EventEmitter for HostEventEmitter {
417 async fn emit(&self, request: EventRequest) -> CoreResult<Event> {
418 let session_id = request.session_id;
419 let lock = self.session_lock(session_id).await;
420 let _guard = lock.lock().await;
421 let event = if request.is_ephemeral() {
422 ephemeral_event(request)
423 } else {
424 self.log
425 .append(request)
426 .await
427 .map_err(|error| AgentLoopError::store(error.to_string()))?
428 };
429 self.notify(event.clone());
430 Ok(event)
431 }
432}
433
434fn ephemeral_event(request: EventRequest) -> Event {
435 Event {
436 id: EventId::new(),
437 event_type: request.event_type,
438 ts: request.ts,
439 session_id: request.session_id,
440 context: request.context,
441 data: request.data,
442 metadata: request.metadata,
443 tags: request.tags,
444 sequence: None,
445 }
446}
447
448#[derive(Default)]
449struct EventIndex {
450 by_session: HashMap<SessionId, Vec<Event>>,
451 by_id: HashMap<EventId, Vec<u8>>,
452 by_sequence: HashMap<(SessionId, i32), EventId>,
453}
454
455impl EventIndex {
456 fn next_sequence(&self, session_id: SessionId) -> Result<i32, EventLogError> {
457 self.by_session
458 .get(&session_id)
459 .and_then(|events| events.last())
460 .and_then(|event| event.sequence)
461 .unwrap_or(0)
462 .checked_add(1)
463 .ok_or_else(|| EventLogError::InvalidAppend {
464 detail: "session sequence exhausted".into(),
465 })
466 }
467
468 fn insert_existing(&mut self, event: Event) -> Result<bool, EventLogError> {
469 let sequence = event.sequence.ok_or_else(|| EventLogError::Corruption {
470 detail: format!("durable event {} has no sequence", event.id),
471 })?;
472 if sequence <= 0 {
473 return Err(EventLogError::Corruption {
474 detail: format!("event {} has non-positive sequence {sequence}", event.id),
475 });
476 }
477 let canonical = serde_json::to_vec(&event).map_err(|error| EventLogError::Corruption {
478 detail: error.to_string(),
479 })?;
480 if let Some(existing) = self.by_id.get(&event.id) {
481 if existing == &canonical {
482 return Ok(false);
483 }
484 return Err(EventLogError::Corruption {
485 detail: format!("event id {} has conflicting canonical envelopes", event.id),
486 });
487 }
488 if let Some(existing_id) = self.by_sequence.get(&(event.session_id, sequence)) {
489 return Err(EventLogError::Corruption {
490 detail: format!(
491 "session {} sequence {sequence} conflicts between {} and {}",
492 event.session_id, existing_id, event.id
493 ),
494 });
495 }
496 if let Some(previous) = self
497 .by_session
498 .get(&event.session_id)
499 .and_then(|events| events.last())
500 .and_then(|event| event.sequence)
501 && sequence <= previous
502 {
503 return Err(EventLogError::Corruption {
504 detail: format!(
505 "session {} sequence {sequence} follows {previous}",
506 event.session_id
507 ),
508 });
509 }
510 self.by_id.insert(event.id, canonical);
511 self.by_sequence
512 .insert((event.session_id, sequence), event.id);
513 self.by_session
514 .entry(event.session_id)
515 .or_default()
516 .push(event);
517 Ok(true)
518 }
519
520 fn page(&self, request: EventReadRequest) -> Result<EventPage, EventLogError> {
521 let events = self
522 .by_session
523 .get(&request.session_id)
524 .map(Vec::as_slice)
525 .unwrap_or_default();
526 let current_high = events.last().and_then(|event| event.sequence);
527 let (after, snapshot) = match request.cursor {
528 Some(cursor) => {
529 if cursor.session_id != request.session_id {
530 return Err(EventLogError::CrossSessionCursor {
531 detail: "cursor belongs to another session".into(),
532 });
533 }
534 let snapshot = cursor
535 .snapshot_high_watermark
536 .unwrap_or(current_high.unwrap_or(0));
537 if cursor.after_sequence > snapshot {
538 return Err(EventLogError::IncompatibleCursor {
539 detail: "cursor position exceeds its snapshot".into(),
540 });
541 }
542 if snapshot > current_high.unwrap_or(0) {
543 return Err(EventLogError::ExpiredCursor {
544 detail: "cursor snapshot is not available in this log".into(),
545 });
546 }
547 (cursor.after_sequence, snapshot)
548 }
549 None => (0, current_high.unwrap_or(0)),
550 };
551 if snapshot == 0 {
552 return Ok(EventPage {
553 events: Vec::new(),
554 next_cursor: None,
555 snapshot_high_watermark: 0,
556 });
557 }
558 let mut selected = events
559 .iter()
560 .filter(|event| {
561 event
562 .sequence
563 .is_some_and(|sequence| sequence > after && sequence <= snapshot)
564 })
565 .take(request.limit.get() + 1)
566 .cloned()
567 .collect::<Vec<_>>();
568 let has_more = selected.len() > request.limit.get();
569 if has_more {
570 selected.pop();
571 }
572 let next_cursor = has_more.then(|| EventCursor {
573 session_id: request.session_id,
574 after_sequence: selected
575 .last()
576 .and_then(|event| event.sequence)
577 .expect("a page with more events returned at least one event"),
578 snapshot_high_watermark: Some(snapshot),
579 });
580 Ok(EventPage {
581 events: selected,
582 next_cursor,
583 snapshot_high_watermark: snapshot,
584 })
585 }
586}
587
588#[derive(Default)]
590pub struct InMemoryEventLog {
591 index: RwLock<EventIndex>,
592}
593
594impl InMemoryEventLog {
595 pub fn new() -> Self {
597 Self::default()
598 }
599}
600
601#[async_trait]
602impl EventReader for InMemoryEventLog {
603 async fn read_page(&self, request: EventReadRequest) -> Result<EventPage, EventLogError> {
604 self.index.read().await.page(request)
605 }
606}
607
608#[async_trait]
609impl EventLog for InMemoryEventLog {
610 async fn append(&self, request: EventRequest) -> Result<Event, EventLogError> {
611 if request.is_ephemeral() {
612 return Err(EventLogError::InvalidAppend {
613 detail: format!(
614 "ephemeral event {} must be routed sink-only",
615 request.event_type
616 ),
617 });
618 }
619 let mut index = self.index.write().await;
620 let sequence = index.next_sequence(request.session_id)?;
621 let event = request.into_event(EventId::new(), sequence);
622 index.insert_existing(event.clone())?;
623 Ok(event)
624 }
625
626 fn durability(&self) -> EventDurability {
627 EventDurability::Volatile
628 }
629}
630
631struct JsonlState {
632 file: tokio::fs::File,
633 committed_len: u64,
634 index: EventIndex,
635}
636
637pub struct JsonlEventLog {
644 path: PathBuf,
645 state: Mutex<JsonlState>,
646}
647
648impl JsonlEventLog {
649 pub async fn open(path: impl AsRef<Path>) -> Result<Self, EventLogError> {
651 let path = path.as_ref().to_path_buf();
652 if let Some(parent) = path.parent()
653 && !parent.as_os_str().is_empty()
654 {
655 tokio::fs::create_dir_all(parent).await?;
656 }
657 let bytes = match tokio::fs::read(&path).await {
658 Ok(bytes) => bytes,
659 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Vec::new(),
660 Err(error) => return Err(error.into()),
661 };
662 let committed_len = bytes
663 .iter()
664 .rposition(|byte| *byte == b'\n')
665 .map_or(0, |position| position + 1);
666 let mut index = EventIndex::default();
667 for (line_index, line) in bytes[..committed_len]
668 .split(|byte| *byte == b'\n')
669 .filter(|line| !line.is_empty())
670 .enumerate()
671 {
672 let event: Event =
673 serde_json::from_slice(line).map_err(|error| EventLogError::Corruption {
674 detail: format!("line {}: {error}", line_index + 1),
675 })?;
676 index.insert_existing(event)?;
677 }
678 let mut options = tokio::fs::OpenOptions::new();
679 options.create(true).read(true).append(true);
680 #[cfg(unix)]
681 {
682 options.mode(0o600);
685 }
686 let file = options.open(&path).await?;
687 if bytes.len() != committed_len {
688 file.set_len(committed_len as u64).await?;
689 file.sync_data().await?;
690 }
691 Ok(Self {
692 path,
693 state: Mutex::new(JsonlState {
694 file,
695 committed_len: committed_len as u64,
696 index,
697 }),
698 })
699 }
700
701 pub fn path(&self) -> &Path {
703 &self.path
704 }
705}
706
707#[async_trait]
708impl EventReader for JsonlEventLog {
709 async fn read_page(&self, request: EventReadRequest) -> Result<EventPage, EventLogError> {
710 self.state.lock().await.index.page(request)
711 }
712}
713
714#[async_trait]
715impl EventLog for JsonlEventLog {
716 async fn append(&self, request: EventRequest) -> Result<Event, EventLogError> {
717 if request.is_ephemeral() {
718 return Err(EventLogError::InvalidAppend {
719 detail: format!(
720 "ephemeral event {} must be routed sink-only",
721 request.event_type
722 ),
723 });
724 }
725 let mut state = self.state.lock().await;
726 let sequence = state.index.next_sequence(request.session_id)?;
727 let event = request.into_event(EventId::new(), sequence);
728 let mut encoded =
729 serde_json::to_vec(&event).map_err(|error| EventLogError::InvalidAppend {
730 detail: error.to_string(),
731 })?;
732 encoded.push(b'\n');
733 let previous_len = state.committed_len;
734 let write_result = async {
735 state.file.write_all(&encoded).await?;
736 state.file.flush().await?;
737 state.file.sync_data().await
738 }
739 .await;
740 if let Err(error) = write_result {
741 let _ = state.file.set_len(previous_len).await;
742 let _ = state.file.sync_data().await;
743 return Err(EventLogError::Backend {
744 detail: error.to_string(),
745 });
746 }
747 state.committed_len += encoded.len() as u64;
748 state.index.insert_existing(event.clone())?;
749 Ok(event)
750 }
751
752 fn durability(&self) -> EventDurability {
753 EventDurability::CrashDurable
754 }
755}
756
757#[derive(Clone)]
758struct ProjectedMessage {
759 event_id: EventId,
760 event_type: String,
761 sequence: i32,
762 tool_name: Option<String>,
763 message: Message,
764}
765
766#[derive(Clone)]
768pub struct EventHistory {
769 reader: Arc<dyn EventReader>,
770}
771
772impl EventHistory {
773 pub fn new(reader: Arc<dyn EventReader>) -> Self {
775 Self { reader }
776 }
777
778 pub fn event_reader(&self) -> Arc<dyn EventReader> {
780 self.reader.clone()
781 }
782
783 pub async fn has_history(&self, session_id: SessionId) -> Result<bool, EventLogError> {
789 let page = self
790 .reader
791 .read_page(EventReadRequest::new(
792 session_id,
793 EventReadLimit::new(1).expect("one is a valid event read limit"),
794 ))
795 .await?;
796 Ok(!page.events.is_empty())
797 }
798
799 pub(crate) async fn contains_event_type(
800 &self,
801 session_id: SessionId,
802 event_type: &str,
803 ) -> Result<bool, EventLogError> {
804 let limit = EventReadLimit::default();
805 let mut request = EventReadRequest::new(session_id, limit);
806 let mut seen = 0usize;
807 loop {
808 let page = self.reader.read_page(request).await?;
809 seen = seen.saturating_add(page.events.len());
810 if page
811 .events
812 .iter()
813 .any(|event| event.event_type == event_type)
814 {
815 return Ok(true);
816 }
817 if seen > MAX_EVENT_HISTORY_REPLAY {
818 return Err(EventLogError::InvalidRead {
819 detail: format!(
820 "event-type replay exceeds the {MAX_EVENT_HISTORY_REPLAY}-event bound"
821 ),
822 });
823 }
824 let Some(cursor) = page.next_cursor else {
825 return Ok(false);
826 };
827 request = EventReadRequest::from_cursor(cursor, limit);
828 }
829 }
830
831 pub async fn read_page(
838 &self,
839 request: EventHistoryReadRequest,
840 ) -> Result<EventHistoryPage, EventLogError> {
841 let message_limit = request.limit.get();
842 let first_raw_limit = EventReadLimit::new(message_limit.min(MAX_EVENT_PAGE_SIZE))?;
843 let mut raw_request = EventReadRequest::new(request.session_id, first_raw_limit);
844 if let Some(cursor) = request.cursor {
845 raw_request = raw_request.with_cursor(cursor);
846 }
847 let mut messages = Vec::with_capacity(message_limit);
848 let mut examined = 0usize;
849 loop {
850 let page = self.reader.read_page(raw_request).await?;
851 let snapshot_high_watermark = page.snapshot_high_watermark();
852 examined = examined.saturating_add(page.events.len());
853 if examined > MAX_EVENT_HISTORY_REPLAY {
854 return Err(EventLogError::InvalidRead {
855 detail: format!(
856 "history page examined more than {MAX_EVENT_HISTORY_REPLAY} events"
857 ),
858 });
859 }
860 for event in page.events {
861 if let Some(message) = message_from_event(&event) {
862 messages.push(message);
863 }
864 }
865 if messages.len() >= message_limit {
866 let boundary_cursor = page.next_cursor.clone();
872 let mut next_cursor = None;
873 if let Some(boundary_cursor) = boundary_cursor {
874 let mut probe_request = EventReadRequest::from_cursor(
875 boundary_cursor.clone(),
876 EventReadLimit::default(),
877 );
878 loop {
879 let probe = self.reader.read_page(probe_request).await?;
880 examined = examined.saturating_add(probe.events.len());
881 if examined > MAX_EVENT_HISTORY_REPLAY {
882 return Err(EventLogError::InvalidRead {
883 detail: format!(
884 "history page examined more than {MAX_EVENT_HISTORY_REPLAY} events"
885 ),
886 });
887 }
888 if probe.events.iter().any(event_projects_message) {
889 next_cursor = Some(boundary_cursor);
890 break;
891 }
892 let Some(cursor) = probe.next_cursor else {
893 break;
894 };
895 probe_request =
896 EventReadRequest::from_cursor(cursor, EventReadLimit::default());
897 }
898 }
899 return Ok(EventHistoryPage {
900 messages,
901 next_cursor,
902 snapshot_high_watermark,
903 });
904 }
905 if page.next_cursor.is_none() {
906 return Ok(EventHistoryPage {
907 messages,
908 next_cursor: None,
909 snapshot_high_watermark,
910 });
911 }
912 let remaining = message_limit - messages.len();
913 let raw_limit = EventReadLimit::new(remaining.min(MAX_EVENT_PAGE_SIZE))?;
914 raw_request = EventReadRequest::from_cursor(
915 page.next_cursor.expect("checked continuation above"),
916 raw_limit,
917 );
918 }
919 }
920
921 async fn project(&self, session_id: SessionId) -> Result<Vec<ProjectedMessage>, EventLogError> {
922 let limit = EventReadLimit::default();
923 let mut request = EventReadRequest::new(session_id, limit);
924 let mut projected = Vec::new();
925 let mut examined = 0usize;
926 loop {
927 let page = self.reader.read_page(request).await?;
928 examined = examined.saturating_add(page.events.len());
929 if examined > MAX_EVENT_HISTORY_REPLAY {
930 return Err(EventLogError::InvalidRead {
931 detail: format!(
932 "history replay examined more than {MAX_EVENT_HISTORY_REPLAY} events"
933 ),
934 });
935 }
936 for event in page.events {
937 if let Some(message) = message_from_event(&event) {
938 projected.push(ProjectedMessage {
939 event_id: event.id,
940 event_type: event.event_type,
941 sequence: event.sequence.expect("reader returns durable events"),
942 tool_name: match &event.data {
943 EventData::ToolCompleted(data) => Some(data.tool_name.clone()),
944 _ => None,
945 },
946 message,
947 });
948 }
949 }
950 let Some(cursor) = page.next_cursor else {
951 break;
952 };
953 request = EventReadRequest::from_cursor(cursor, limit);
954 }
955 Ok(projected)
956 }
957
958 async fn filtered(&self, query: &MessageQuery) -> Result<Vec<ProjectedMessage>, EventLogError> {
959 let mut projected = self.project(query.session_id).await?;
960 if let Some(after) = query.after_sequence {
961 projected.retain(|item| i64::from(item.sequence) > after);
962 }
963 for filter in &query.filters {
964 match filter {
965 MessageFilter::TimeRange { from, to } => projected.retain(|item| {
966 from.is_none_or(|from| item.message.created_at >= from)
967 && to.is_none_or(|to| item.message.created_at <= to)
968 }),
969 MessageFilter::EventTypes(types) => {
970 projected.retain(|item| types.contains(&item.event_type))
971 }
972 MessageFilter::ToolName(name) => {
973 projected.retain(|item| item.tool_name.as_ref() == Some(name))
974 }
975 MessageFilter::Search(search) => {
976 let search = search.to_lowercase();
977 projected.retain(|item| {
978 item.message
979 .text()
980 .is_some_and(|text| text.to_lowercase().contains(&search))
981 });
982 }
983 MessageFilter::ExcludeIds(ids) => {
984 projected.retain(|item| !ids.contains(&item.event_id))
985 }
986 MessageFilter::IncludeIds(ids) => {
987 projected.retain(|item| ids.contains(&item.event_id))
988 }
989 MessageFilter::Custom(predicate) => {
990 projected.retain(|item| predicate(&item.message))
991 }
992 }
993 }
994 Ok(projected)
995 }
996}
997
998fn event_projects_message(event: &Event) -> bool {
999 matches!(
1000 &event.data,
1001 EventData::InputMessage(_)
1002 | EventData::OutputMessageCompleted(_)
1003 | EventData::ToolCompleted(_)
1004 )
1005}
1006
1007#[async_trait]
1008impl MessageRetriever for EventHistory {
1009 async fn get(
1010 &self,
1011 session_id: SessionId,
1012 message_id: MessageId,
1013 ) -> CoreResult<Option<Message>> {
1014 Ok(self
1015 .project(session_id)
1016 .await
1017 .map_err(core_event_error)?
1018 .into_iter()
1019 .find(|item| item.message.id == message_id)
1020 .map(|item| item.message))
1021 }
1022
1023 async fn load(&self, session_id: SessionId) -> CoreResult<Vec<Message>> {
1024 Ok(self
1025 .project(session_id)
1026 .await
1027 .map_err(core_event_error)?
1028 .into_iter()
1029 .map(|item| item.message)
1030 .collect())
1031 }
1032
1033 async fn load_filtered(&self, query: MessageQuery) -> CoreResult<Vec<Message>> {
1034 let mut messages = self
1035 .filtered(&query)
1036 .await
1037 .map_err(core_event_error)?
1038 .into_iter()
1039 .map(|item| item.message)
1040 .collect::<Vec<_>>();
1041 let count_before_limit = messages.len();
1042 query.apply_window_bounds(&mut messages);
1043 query.prepend_excluded_notice(&mut messages, count_before_limit);
1044 query.apply_injections(&mut messages);
1045 Ok(messages)
1046 }
1047
1048 async fn load_filtered_history(&self, query: MessageQuery) -> CoreResult<MessageHistory> {
1049 let source_sequence = self
1050 .project(query.session_id)
1051 .await
1052 .map_err(core_event_error)?
1053 .last()
1054 .map(|item| i64::from(item.sequence));
1055 Ok(MessageHistory {
1056 messages: self.load_filtered(query).await?,
1057 source_sequence,
1058 })
1059 }
1060
1061 async fn load_page(
1062 &self,
1063 session_id: SessionId,
1064 offset: usize,
1065 limit: usize,
1066 ) -> CoreResult<Vec<Message>> {
1067 if limit == 0 {
1068 return Ok(Vec::new());
1069 }
1070 let mut cursor = None;
1071 let mut skipped = 0usize;
1072 let mut messages = Vec::with_capacity(limit.min(MAX_EVENT_HISTORY_PAGE_SIZE));
1073 loop {
1074 let requested = if skipped < offset {
1075 (offset - skipped).min(MAX_EVENT_HISTORY_PAGE_SIZE)
1076 } else {
1077 (limit - messages.len()).min(MAX_EVENT_HISTORY_PAGE_SIZE)
1078 };
1079 let mut request = EventHistoryReadRequest::new(
1080 session_id,
1081 EventHistoryReadLimit::new(requested).map_err(core_event_error)?,
1082 );
1083 if let Some(previous) = cursor {
1084 request = request.with_cursor(previous);
1085 }
1086 let page = self.read_page(request).await.map_err(core_event_error)?;
1087 if skipped < offset {
1088 skipped = skipped.saturating_add(page.messages.len());
1089 } else {
1090 messages.extend(page.messages);
1091 }
1092 cursor = page.next_cursor;
1093 if messages.len() >= limit || cursor.is_none() {
1094 messages.truncate(limit);
1095 return Ok(messages);
1096 }
1097 }
1098 }
1099
1100 async fn count(&self, session_id: SessionId) -> CoreResult<usize> {
1101 Ok(self
1102 .project(session_id)
1103 .await
1104 .map_err(core_event_error)?
1105 .len())
1106 }
1107}
1108
1109fn core_event_error(error: EventLogError) -> AgentLoopError {
1110 AgentLoopError::store(error.to_string())
1111}
1112
1113fn message_from_event(event: &Event) -> Option<Message> {
1114 match &event.data {
1115 EventData::InputMessage(data) => Some(data.message.clone()),
1116 EventData::OutputMessageCompleted(OutputMessageCompletedData { message, .. }) => {
1117 Some(message.clone())
1118 }
1119 EventData::ToolCompleted(data) => {
1120 let mut message = tool_completed_to_message(data.clone());
1121 message.id = MessageId::from_uuid(event.id.uuid());
1122 message.created_at = event.ts;
1123 Some(message)
1124 }
1125 _ => None,
1128 }
1129}
1130
1131fn tool_completed_to_message(data: everruns_core::events::ToolCompletedData) -> Message {
1132 let mut images = Vec::<ToolResultImage>::new();
1133 let result = data.result.map(|parts| {
1134 for part in &parts {
1135 if let ContentPart::Image(image) = part
1136 && let (Some(base64), Some(media_type)) = (&image.base64, &image.media_type)
1137 {
1138 images.push(ToolResultImage {
1139 base64: base64.clone(),
1140 media_type: media_type.clone(),
1141 });
1142 }
1143 }
1144 let text_parts = parts
1145 .iter()
1146 .filter(|part| matches!(part, ContentPart::Text(_)))
1147 .collect::<Vec<_>>();
1148 if text_parts.len() == 1
1149 && let ContentPart::Text(text) = text_parts[0]
1150 {
1151 parse_structured_tool_result_text(&text.text)
1152 } else if text_parts.is_empty() {
1153 serde_json::Value::Null
1154 } else {
1155 serde_json::to_value(text_parts).unwrap_or_default()
1156 }
1157 });
1158 let mut message = if images.is_empty() {
1159 Message::tool_result(&data.tool_call_id, result, data.error)
1160 } else {
1161 Message::tool_result_with_images(&data.tool_call_id, result, images)
1162 };
1163 let mut metadata = std::collections::HashMap::new();
1164 metadata.insert("tool_name".into(), serde_json::json!(data.tool_name));
1165 if let Some(value) = data.tool_call_fingerprint {
1166 metadata.insert("tool_call_fingerprint".into(), serde_json::json!(value));
1167 }
1168 if let Some(value) = data.tool_result_fingerprint {
1169 metadata.insert("tool_result_fingerprint".into(), serde_json::json!(value));
1170 }
1171 message.metadata = Some(metadata);
1172 message
1173}
1174
1175fn parse_structured_tool_result_text(text: &str) -> serde_json::Value {
1176 let trimmed = text.trim_start();
1177 if !trimmed.starts_with('{') && !trimmed.starts_with('[') {
1178 return serde_json::Value::String(text.to_string());
1179 }
1180 match serde_json::from_str(trimmed) {
1181 Ok(value @ (serde_json::Value::Object(_) | serde_json::Value::Array(_))) => value,
1182 _ => serde_json::Value::String(text.to_string()),
1183 }
1184}
1185
1186#[cfg(test)]
1187mod tests {
1188 use super::*;
1189 use everruns_core::events::{
1190 EventContext, InputMessageData, OutputMessageDeltaData, SessionStartedData,
1191 };
1192 use everruns_core::typed_id::{HarnessId, TurnId};
1193
1194 struct LifecycleHeavyReader;
1195
1196 #[async_trait]
1197 impl EventReader for LifecycleHeavyReader {
1198 async fn read_page(&self, request: EventReadRequest) -> Result<EventPage, EventLogError> {
1199 let after = request
1200 .cursor
1201 .as_ref()
1202 .map_or(0, EventCursor::after_sequence);
1203 let high_watermark = (MAX_EVENT_HISTORY_REPLAY + 1) as i32;
1204 let end = after
1205 .saturating_add(request.limit.get() as i32)
1206 .min(high_watermark);
1207 let events = ((after + 1)..=end)
1208 .map(|sequence| {
1209 EventRequest::new(
1210 request.session_id,
1211 EventContext::empty(),
1212 OutputMessageDeltaData {
1213 turn_id: TurnId::new(),
1214 message_id: MessageId::new(),
1215 delta: String::new(),
1216 accumulated: String::new(),
1217 phase: None,
1218 },
1219 )
1220 .into_event(EventId::new(), sequence)
1221 })
1222 .collect();
1223 let next_cursor = (end < high_watermark).then_some(EventCursor {
1224 session_id: request.session_id,
1225 after_sequence: end,
1226 snapshot_high_watermark: Some(high_watermark),
1227 });
1228 Ok(EventPage {
1229 events,
1230 next_cursor,
1231 snapshot_high_watermark: high_watermark,
1232 })
1233 }
1234 }
1235
1236 #[tokio::test]
1237 async fn full_projection_caps_examined_lifecycle_envelopes() {
1238 let history = EventHistory::new(Arc::new(LifecycleHeavyReader));
1239 let error = match history.project(SessionId::new()).await {
1240 Ok(_) => panic!("lifecycle-heavy replay must be bounded"),
1241 Err(error) => error,
1242 };
1243 assert!(matches!(error, EventLogError::InvalidRead { .. }));
1244 assert!(error.to_string().contains("examined more than"));
1245 }
1246
1247 #[tokio::test]
1248 async fn exact_message_boundary_ignores_trailing_lifecycle_events() {
1249 let session_id = SessionId::new();
1250 let log = Arc::new(InMemoryEventLog::new());
1251 log.append(EventRequest::new(
1252 session_id,
1253 EventContext::empty(),
1254 InputMessageData::new(Message::user("hello")),
1255 ))
1256 .await
1257 .expect("append input message");
1258 log.append(EventRequest::new(
1259 session_id,
1260 EventContext::empty(),
1261 OutputMessageCompletedData::new(Message::assistant("hi")),
1262 ))
1263 .await
1264 .expect("append output message");
1265 log.append(EventRequest::new(
1266 session_id,
1267 EventContext::empty(),
1268 SessionStartedData {
1269 harness_id: HarnessId::new(),
1270 agent_id: None,
1271 model_id: None,
1272 },
1273 ))
1274 .await
1275 .expect("append lifecycle event");
1276
1277 let page = EventHistory::new(log)
1278 .read_page(EventHistoryReadRequest::new(
1279 session_id,
1280 EventHistoryReadLimit::new(2).expect("valid history limit"),
1281 ))
1282 .await
1283 .expect("read history");
1284
1285 assert_eq!(page.messages.len(), 2);
1286 assert!(page.next_cursor.is_none());
1287 }
1288
1289 #[test]
1290 fn tool_completion_projection_preserves_structured_result_and_fingerprints() {
1291 let event = EventRequest::new(
1292 SessionId::new(),
1293 EventContext::empty(),
1294 everruns_core::events::ToolCompletedData::success(
1295 "call_read".into(),
1296 "read_file".into(),
1297 vec![ContentPart::text(
1298 serde_json::json!({
1299 "path": "/workspace/src/lib.rs",
1300 "content": "1|fn main() {}"
1301 })
1302 .to_string(),
1303 )],
1304 Some(1),
1305 )
1306 .with_fingerprints("sha256:call".into(), "sha256:result".into()),
1307 )
1308 .into_event(EventId::new(), 1);
1309 let message = message_from_event(&event).expect("tool result message");
1310 let result = message
1311 .tool_result_content()
1312 .and_then(|content| content.result.as_ref())
1313 .expect("projected result");
1314 assert_eq!(result["path"], "/workspace/src/lib.rs");
1315 let metadata = message.metadata.expect("tool metadata");
1316 assert_eq!(metadata["tool_name"], "read_file");
1317 assert_eq!(metadata["tool_call_fingerprint"], "sha256:call");
1318 assert_eq!(metadata["tool_result_fingerprint"], "sha256:result");
1319 }
1320
1321 #[test]
1322 fn tool_completion_projection_keeps_scalar_json_as_text() {
1323 let event = EventRequest::new(
1324 SessionId::new(),
1325 EventContext::empty(),
1326 everruns_core::events::ToolCompletedData::success(
1327 "call_scalar".into(),
1328 "custom_tool".into(),
1329 vec![ContentPart::text("123")],
1330 Some(1),
1331 ),
1332 )
1333 .into_event(EventId::new(), 1);
1334 let message = message_from_event(&event).expect("tool result message");
1335 let result = message
1336 .tool_result_content()
1337 .and_then(|content| content.result.as_ref())
1338 .expect("projected result");
1339 assert_eq!(result, &serde_json::Value::String("123".into()));
1340 }
1341}