runifold_store_postgres/
journal.rs1use 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}