Skip to main content

runifold_store_postgres/
journal.rs

1//! `PostgreSQL` canonical run-event journal and operational reader.
2
3use runifold_core::{Journal, JournalError, RunEvent, RunId};
4use runifold_ops::{
5    RunEventCursor, RunEventPage, RunEventPageSize, RunEventSource, RunEventSourceError,
6};
7use serde_json::Value;
8
9use crate::PostgresConversationStore;
10
11impl Journal for PostgresConversationStore {
12    fn record(&self, event: &RunEvent) -> Result<(), JournalError> {
13        let table = format!("{}_events", self.table());
14        let event_id = event.meta.event_id.as_uuid();
15        let run_id = event.meta.run_id.as_uuid();
16        let sequence = i64::try_from(event.meta.sequence).map_err(|_| JournalError {
17            message: "event sequence exceeds PostgreSQL BIGINT".into(),
18        })?;
19        let encoded = serde_json::to_value(event).map_err(|_| JournalError {
20            message: "canonical run event could not be encoded".into(),
21        })?;
22        self.blocking()
23            .execute(move |client| {
24                client.execute(
25                    &format!(
26                        "INSERT INTO {table} (event_id, run_id, sequence, event_json) \
27                         VALUES ($1, $2, $3, $4)"
28                    ),
29                    &[&event_id, &run_id, &sequence, &encoded],
30                )
31            })
32            .map_err(|_| JournalError {
33                message: "PostgreSQL journal worker is unavailable".into(),
34            })?
35            .map(|_| ())
36            .map_err(|_| JournalError {
37                message: "PostgreSQL journal write failed".into(),
38            })
39    }
40}
41
42impl RunEventSource for PostgresConversationStore {
43    fn event_page(
44        &self,
45        run_id: RunId,
46        after: Option<RunEventCursor>,
47        limit: RunEventPageSize,
48    ) -> Result<RunEventPage, RunEventSourceError> {
49        let table = format!("{}_events", self.table());
50        let run_id = run_id.as_uuid();
51        let after = after.map_or(-1, |cursor| {
52            i64::try_from(cursor.sequence()).unwrap_or(i64::MAX)
53        });
54        let query_limit = limit
55            .get()
56            .checked_add(1)
57            .and_then(|value| i64::try_from(value).ok())
58            .ok_or_else(|| RunEventSourceError::storage("event page limit overflow"))?;
59        let rows = self
60            .blocking()
61            .execute(move |client| {
62                client.query(
63                    &format!(
64                        "SELECT sequence, event_json FROM {table} \
65                         WHERE run_id = $1 AND sequence > $2 \
66                         ORDER BY sequence ASC LIMIT $3"
67                    ),
68                    &[&run_id, &after, &query_limit],
69                )
70            })
71            .map_err(|_| RunEventSourceError::storage("PostgreSQL journal worker unavailable"))?
72            .map_err(|_| RunEventSourceError::storage("PostgreSQL journal query failed"))?;
73        let mut events = rows
74            .into_iter()
75            .map(|row| {
76                let stored_sequence: i64 = row.get("sequence");
77                let event = decode_event(row.get("event_json"))?;
78                if event.meta.run_id != RunId::from_uuid(run_id)
79                    || i64::try_from(event.meta.sequence).ok() != Some(stored_sequence)
80                {
81                    return Err(RunEventSourceError::corrupt_data(
82                        "event index does not match its canonical envelope",
83                    ));
84                }
85                Ok(event)
86            })
87            .collect::<Result<Vec<_>, _>>()?;
88        let has_more = events.len() > limit.get();
89        if has_more {
90            events.truncate(limit.get());
91        }
92        let next = if has_more {
93            events
94                .last()
95                .map(|event| RunEventCursor::after(event.meta.sequence))
96        } else {
97            None
98        };
99        Ok(RunEventPage { events, next })
100    }
101}
102
103fn decode_event(value: Value) -> Result<RunEvent, RunEventSourceError> {
104    serde_json::from_value(value)
105        .map_err(|_| RunEventSourceError::corrupt_data("persisted run event is invalid"))
106}