Skip to main content

everruns_host/
events.rs

1//! Canonical host event persistence, bounded replay, and message projection.
2
3use 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
21/// Default number of canonical envelopes returned by a host event read.
22pub const DEFAULT_EVENT_READ_LIMIT: usize = 256;
23/// Hard upper bound for one host event read.
24pub const MAX_EVENT_PAGE_SIZE: usize = 1024;
25/// Maximum envelopes one [`EventHistory`] projection may examine in one call.
26pub const MAX_EVENT_HISTORY_REPLAY: usize = 100_000;
27/// Maximum projected messages returned by one bounded history page.
28pub const MAX_EVENT_HISTORY_PAGE_SIZE: usize = 256;
29
30/// Persistence guarantee provided by an [`EventLog`] implementation.
31#[derive(Clone, Copy, Debug, PartialEq, Eq)]
32pub enum EventDurability {
33    /// Accepted writes survive for the lifetime of the log instance only.
34    Volatile,
35    /// Accepted writes survive process crashes according to the backend contract.
36    CrashDurable,
37}
38
39/// Errors from canonical event append or bounded replay.
40#[derive(Clone, Debug, thiserror::Error, PartialEq, Eq)]
41#[non_exhaustive]
42pub enum EventLogError {
43    /// A cursor, snapshot, session binding, or limit was invalid.
44    #[error("invalid event read: {detail}")]
45    InvalidRead { detail: String },
46    /// A cursor was presented for a different session.
47    #[error("event cursor belongs to another session: {detail}")]
48    CrossSessionCursor { detail: String },
49    /// A cursor's position and stable snapshot are internally inconsistent.
50    #[error("incompatible event cursor: {detail}")]
51    IncompatibleCursor { detail: String },
52    /// The cursor references a snapshot no longer available from this reader.
53    #[error("expired event cursor: {detail}")]
54    ExpiredCursor { detail: String },
55    /// The request cannot be appended to a durable log.
56    #[error("invalid event append: {detail}")]
57    InvalidAppend { detail: String },
58    /// Persisted canonical envelopes conflict or are malformed.
59    #[error("event log corruption: {detail}")]
60    Corruption { detail: String },
61    /// The storage backend failed.
62    #[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/// Validated bound for one [`EventReader::read_page`] call.
75#[derive(Clone, Copy, Debug, PartialEq, Eq)]
76pub struct EventReadLimit(u16);
77
78impl EventReadLimit {
79    /// Validate a non-zero read limit no larger than [`MAX_EVENT_PAGE_SIZE`].
80    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    /// Return the validated limit as a `usize`.
90    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/// Snapshot-bound continuation cursor for canonical event replay.
102#[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    /// Session this continuation is bound to.
111    pub fn session_id(&self) -> SessionId {
112        self.session_id
113    }
114
115    /// Last sequence returned before this continuation.
116    pub fn after_sequence(&self) -> i32 {
117        self.after_sequence
118    }
119
120    /// Stable high-watermark captured by the first page.
121    pub fn snapshot_high_watermark(&self) -> Option<i32> {
122        self.snapshot_high_watermark
123    }
124
125    /// Start a new snapshot after a previously observed high-watermark.
126    ///
127    /// This is the polling form: the next read captures the then-current
128    /// high-watermark instead of remaining pinned to an earlier snapshot.
129    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/// One bounded canonical event read.
144#[derive(Clone, Debug)]
145pub struct EventReadRequest {
146    session_id: SessionId,
147    cursor: Option<EventCursor>,
148    limit: EventReadLimit,
149}
150
151impl EventReadRequest {
152    /// Start a stable snapshot read for `session_id`.
153    pub fn new(session_id: SessionId, limit: EventReadLimit) -> Self {
154        Self {
155            session_id,
156            cursor: None,
157            limit,
158        }
159    }
160
161    /// Continue the stable snapshot represented by `cursor`.
162    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    /// Attach a cursor while retaining the explicitly selected session.
171    ///
172    /// Readers reject a cursor originating from another session.
173    pub fn with_cursor(mut self, cursor: EventCursor) -> Self {
174        self.cursor = Some(cursor);
175        self
176    }
177
178    /// Session selected for this read.
179    pub fn session_id(&self) -> SessionId {
180        self.session_id
181    }
182
183    /// Validated page bound.
184    pub fn limit(&self) -> EventReadLimit {
185        self.limit
186    }
187}
188
189/// A page of full canonical event envelopes in persisted sequence order.
190#[derive(Clone, Debug)]
191pub struct EventPage {
192    /// Full canonical envelopes, never projection records.
193    pub events: Vec<Event>,
194    /// Continuation within the first page's stable snapshot.
195    pub next_cursor: Option<EventCursor>,
196    snapshot_high_watermark: i32,
197}
198
199/// Validated projected-message bound for [`EventHistory::read_page`].
200#[derive(Clone, Copy, Debug, PartialEq, Eq)]
201pub struct EventHistoryReadLimit(u16);
202
203impl EventHistoryReadLimit {
204    /// Validate a non-zero message limit no larger than
205    /// [`MAX_EVENT_HISTORY_PAGE_SIZE`].
206    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    /// Validated projected-message count.
218    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/// One bounded projected-history read.
230#[derive(Clone, Debug)]
231pub struct EventHistoryReadRequest {
232    session_id: SessionId,
233    cursor: Option<EventCursor>,
234    limit: EventHistoryReadLimit,
235}
236
237impl EventHistoryReadRequest {
238    /// Start a projected-message snapshot for `session_id`.
239    pub fn new(session_id: SessionId, limit: EventHistoryReadLimit) -> Self {
240        Self {
241            session_id,
242            cursor: None,
243            limit,
244        }
245    }
246
247    /// Continue or poll using a session-bound event cursor.
248    pub fn with_cursor(mut self, cursor: EventCursor) -> Self {
249        self.cursor = Some(cursor);
250        self
251    }
252
253    /// Selected session.
254    pub fn session_id(&self) -> SessionId {
255        self.session_id
256    }
257
258    /// Validated message-page bound.
259    pub fn limit(&self) -> EventHistoryReadLimit {
260        self.limit
261    }
262}
263
264/// Bounded message projection and its stable event snapshot continuation.
265#[derive(Clone, Debug)]
266pub struct EventHistoryPage {
267    /// Canonical messages in persisted event sequence order.
268    pub messages: Vec<Message>,
269    /// Continuation after the last examined canonical envelope.
270    pub next_cursor: Option<EventCursor>,
271    snapshot_high_watermark: i32,
272}
273
274impl EventHistoryPage {
275    /// Stable canonical-event high-watermark used for this page.
276    pub fn snapshot_high_watermark(&self) -> i32 {
277        self.snapshot_high_watermark
278    }
279}
280
281impl EventPage {
282    /// High-watermark captured by the first read (`0` for an empty session).
283    pub fn snapshot_high_watermark(&self) -> i32 {
284        self.snapshot_high_watermark
285    }
286}
287
288/// Bounded, ordered reader for canonical session events.
289#[async_trait]
290pub trait EventReader: Send + Sync {
291    /// Read one page. The first page fixes a snapshot high-watermark; continuations
292    /// cannot observe concurrent appends and therefore neither skip nor duplicate.
293    async fn read_page(&self, request: EventReadRequest) -> Result<EventPage, EventLogError>;
294}
295
296/// A coherent canonical event log.
297///
298/// Implementations explicitly pair their reader and writer and guarantee
299/// read-your-accepted-writes, unique monotonic per-session sequences, and the
300/// declared durability class. There is deliberately no blanket implementation.
301#[async_trait]
302pub trait EventLog: EventReader {
303    /// Commit a durable request, assigning its event id and persisted sequence.
304    async fn append(&self, request: EventRequest) -> Result<Event, EventLogError>;
305
306    /// Persistence guarantee used for acknowledged appends.
307    fn durability(&self) -> EventDurability;
308}
309
310/// Non-blocking destination for already-finalized canonical envelopes.
311pub trait EventSink: Send + Sync {
312    /// Attempt immediate observational delivery without delaying execution.
313    fn try_send(&self, event: Event) -> Result<(), EventSinkError>;
314}
315
316/// Observational delivery failures. They never change append or turn success.
317#[derive(Clone, Copy, Debug, thiserror::Error, PartialEq, Eq)]
318pub enum EventSinkError {
319    /// The bounded sink cannot accept another envelope now.
320    #[error("event sink is full")]
321    Full {
322        /// Number of envelopes the sink reports dropping.
323        dropped: u64,
324    },
325    /// The sink has no remaining receiver.
326    #[error("event sink is closed")]
327    Closed,
328}
329
330/// Sink used when a host has no live observers.
331#[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/// Counts observational failures without making them execution failures.
341#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
342pub struct EventDeliveryStats {
343    /// Envelopes rejected because the sink was full.
344    pub full: u64,
345    /// Envelopes rejected because the sink was closed.
346    pub closed: u64,
347}
348
349/// Canonical event router shared by in-process, worker, and facade hosts.
350///
351/// Durable requests commit to the log before notification. Ephemeral requests
352/// use the protocol's [`EventRequest::is_ephemeral`] classification, receive
353/// `sequence = None`, and are sink-only. Per-session serialization preserves
354/// append acceptance order through the post-commit sink call.
355#[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    /// Route requests through a coherent log and optional live sink.
366    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    /// Canonical log used by this router.
377    pub fn event_log(&self) -> Arc<dyn EventLog> {
378        self.log.clone()
379    }
380
381    /// Current observational failure counters.
382    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/// Volatile coherent event log used by default in-process hosts.
589#[derive(Default)]
590pub struct InMemoryEventLog {
591    index: RwLock<EventIndex>,
592}
593
594impl InMemoryEventLog {
595    /// Create an empty volatile log.
596    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
637/// Crash-durable JSONL log containing complete canonical [`Event`] envelopes.
638///
639/// An append is acknowledged only after newline-delimited bytes are flushed and
640/// `sync_data` succeeds. A torn final record is an uncommitted tail and is
641/// truncated on open; any malformed complete line is corruption. The in-memory
642/// index is disposable and rebuilt solely from the canonical log.
643pub struct JsonlEventLog {
644    path: PathBuf,
645    state: Mutex<JsonlState>,
646}
647
648impl JsonlEventLog {
649    /// Open or create a canonical event JSONL file and rebuild its index.
650    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            // THREAT[TM-FS-014]: canonical envelopes can contain prompts,
683            // tool results, and metadata; newly created logs are owner-only.
684            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    /// Backing canonical event-log path.
702    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/// Read-only message/history projection rebuilt from canonical events.
767#[derive(Clone)]
768pub struct EventHistory {
769    reader: Arc<dyn EventReader>,
770}
771
772impl EventHistory {
773    /// Project messages from a bounded canonical event reader.
774    pub fn new(reader: Arc<dyn EventReader>) -> Self {
775        Self { reader }
776    }
777
778    /// Underlying canonical reader.
779    pub fn event_reader(&self) -> Arc<dyn EventReader> {
780        self.reader.clone()
781    }
782
783    /// Whether the canonical log contains any event for this session id.
784    ///
785    /// This reports history presence only. It is deliberately not a session
786    /// identity/catalog check: a valid session can have no events, and retained
787    /// events can outlive the session record that originally produced them.
788    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    /// Read one bounded projected-message page without collecting the session.
832    ///
833    /// Raw event reads stay bounded by the remaining message slots, lifecycle
834    /// envelopes are skipped, the first event snapshot is retained across all
835    /// internal reads, and no more than [`MAX_EVENT_HISTORY_REPLAY`] envelopes
836    /// are examined.
837    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                // A raw continuation can point only to trailing lifecycle
867                // envelopes. Probe within the same bounded snapshot so the
868                // projected cursor truthfully means another message exists and
869                // callers never need an empty terminal page after an exact
870                // message boundary.
871                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        // `output.message.replaced` is live/persisted narration only. The later
1126        // completed envelope supplies the one canonical assistant message.
1127        _ => 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}