Skip to main content

eventuary_sqlite/
writer.rs

1use std::fmt;
2use std::num::NonZeroU32;
3use std::sync::Arc;
4
5use eventuary_core::io::Writer;
6use eventuary_core::partition::{
7    PartitionHash, PartitionHasher, PartitionKey, PartitionKeyResolver, PartitionStrategy,
8};
9use eventuary_core::{Error, Event, Result, SerializedEvent};
10
11use crate::database::SqliteConn;
12use crate::event_log::{SqliteEventLogSchema, SqliteEventLogSchemaConfig};
13use crate::relation::SqliteRelationName;
14
15#[derive(Clone, Default)]
16pub enum SqlitePartitioningConfig {
17    #[default]
18    Off,
19    Inline {
20        partition_count: NonZeroU32,
21        key_resolver: Arc<dyn PartitionKeyResolver>,
22        hasher: Arc<dyn PartitionHasher>,
23    },
24}
25
26impl SqlitePartitioningConfig {
27    pub fn inline(
28        count: NonZeroU32,
29        resolver: impl PartitionKeyResolver + 'static,
30        hasher: impl PartitionHasher + 'static,
31    ) -> Self {
32        Self::Inline {
33            partition_count: count,
34            key_resolver: Arc::new(resolver),
35            hasher: Arc::new(hasher),
36        }
37    }
38}
39
40impl fmt::Debug for SqlitePartitioningConfig {
41    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
42        match self {
43            Self::Off => write!(f, "SqlitePartitioningConfig::Off"),
44            Self::Inline {
45                partition_count, ..
46            } => f
47                .debug_struct("SqlitePartitioningConfig::Inline")
48                .field("partition_count", partition_count)
49                .finish(),
50        }
51    }
52}
53
54#[derive(Debug, Clone)]
55pub struct SqliteWriterConfig {
56    pub events_relation: SqliteRelationName,
57    pub partitioning: SqlitePartitioningConfig,
58}
59
60impl Default for SqliteWriterConfig {
61    fn default() -> Self {
62        Self {
63            events_relation: SqliteRelationName::new("events").expect("default events relation"),
64            partitioning: SqlitePartitioningConfig::Off,
65        }
66    }
67}
68
69pub struct SqliteWriter {
70    conn: SqliteConn,
71    insert_sql: Arc<String>,
72    partitioning: SqlitePartitioningConfig,
73}
74
75impl SqliteWriter {
76    pub fn connect(conn: SqliteConn, config: SqliteWriterConfig) -> Result<Self> {
77        Self::prepare_schema(&conn, &config)?;
78        Ok(Self::new_with_config(conn, config))
79    }
80
81    pub fn prepare_schema(conn: &SqliteConn, config: &SqliteWriterConfig) -> Result<()> {
82        let guard = conn.lock().map_err(|e| Error::Store(e.to_string()))?;
83        SqliteEventLogSchema::prepare(
84            &guard,
85            &SqliteEventLogSchemaConfig {
86                events_relation: config.events_relation.clone(),
87            },
88        )
89    }
90
91    pub fn schema_sql(config: &SqliteWriterConfig) -> String {
92        SqliteEventLogSchema::schema_sql(&SqliteEventLogSchemaConfig {
93            events_relation: config.events_relation.clone(),
94        })
95    }
96
97    pub fn new(conn: SqliteConn) -> Self {
98        Self::new_with_config(conn, SqliteWriterConfig::default())
99    }
100
101    pub fn new_with_config(conn: SqliteConn, config: SqliteWriterConfig) -> Self {
102        let insert_sql = format!(
103            "INSERT INTO {events} \
104             (id, organization, namespace, topic, event_key, payload, content_type, metadata, \
105             timestamp, version, parent_id, correlation_id, causation_id, \
106             partition_key, partition_hash, partition_id, partition_count, partition_strategy) \
107             VALUES \
108             (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18)",
109            events = config.events_relation.render(),
110        );
111        Self {
112            conn,
113            insert_sql: Arc::new(insert_sql),
114            partitioning: config.partitioning,
115        }
116    }
117
118    fn partition_data(&self, event: &Event) -> Result<PartitionData> {
119        match &self.partitioning {
120            SqlitePartitioningConfig::Off => Ok(PartitionData::default()),
121            SqlitePartitioningConfig::Inline {
122                partition_count,
123                key_resolver,
124                hasher,
125            } => {
126                let partition_key = key_resolver.partition_key(event)?;
127                let partition_hash = hasher.hash(&partition_key);
128                let partition = hasher.partition_for(&partition_key, *partition_count);
129                let partition_strategy = PartitionStrategy::new(hasher.strategy())?;
130                Ok(PartitionData {
131                    partition_key: Some(partition_key),
132                    partition_hash: Some(partition_hash),
133                    partition_id: Some(partition.id() as i64),
134                    partition_count: Some(partition.count() as i64),
135                    partition_strategy: Some(partition_strategy),
136                })
137            }
138        }
139    }
140}
141
142impl Writer for SqliteWriter {
143    async fn write(&self, event: &Event) -> Result<()> {
144        let conn = Arc::clone(&self.conn);
145        let pd = self.partition_data(event)?;
146        let event = event.clone();
147        let sql = Arc::clone(&self.insert_sql);
148        tokio::task::spawn_blocking(move || {
149            let guard = conn.lock().map_err(|e| Error::Store(e.to_string()))?;
150            insert_event(&guard, &sql, &event, &pd)
151        })
152        .await
153        .map_err(|e| Error::Store(format!("blocking task panicked: {e}")))?
154    }
155
156    async fn write_all(&self, events: &[Event]) -> Result<()> {
157        if events.is_empty() {
158            return Ok(());
159        }
160        let conn = Arc::clone(&self.conn);
161        let partition_data: Result<Vec<PartitionData>> =
162            events.iter().map(|e| self.partition_data(e)).collect();
163        let partition_data = partition_data?;
164        let events = events.to_vec();
165        let sql = Arc::clone(&self.insert_sql);
166        tokio::task::spawn_blocking(move || {
167            let mut guard = conn.lock().map_err(|e| Error::Store(e.to_string()))?;
168            let tx = guard
169                .transaction()
170                .map_err(|e| Error::Store(e.to_string()))?;
171            for (event, pd) in events.iter().zip(partition_data.iter()) {
172                insert_event(&tx, &sql, event, pd)?;
173            }
174            tx.commit().map_err(|e| Error::Store(e.to_string()))?;
175            Ok(())
176        })
177        .await
178        .map_err(|e| Error::Store(format!("blocking task panicked: {e}")))?
179    }
180}
181
182#[derive(Default)]
183struct PartitionData {
184    partition_key: Option<PartitionKey>,
185    partition_hash: Option<PartitionHash>,
186    partition_id: Option<i64>,
187    partition_count: Option<i64>,
188    partition_strategy: Option<PartitionStrategy>,
189}
190
191fn insert_event(
192    conn: &rusqlite::Connection,
193    sql: &str,
194    event: &Event,
195    pd: &PartitionData,
196) -> Result<()> {
197    let serialized = SerializedEvent::from_event(event)?;
198    let content_type = serialized.payload.content_type().to_string();
199    let payload = serde_json::to_string(&serialized.payload)
200        .map_err(|e| Error::Store(format!("encode payload: {e}")))?;
201    let metadata = serde_json::to_string(&serialized.metadata)
202        .map_err(|e| Error::Store(format!("encode metadata: {e}")))?;
203    conn.execute(
204        sql,
205        rusqlite::params![
206            serialized.id.to_string(),
207            serialized.organization,
208            serialized.namespace,
209            serialized.topic,
210            serialized.key,
211            payload,
212            content_type,
213            metadata,
214            serialized.timestamp.to_rfc3339(),
215            serialized.version as i64,
216            serialized.parent_id.map(|id| id.to_string()),
217            serialized.correlation_id,
218            serialized.causation_id,
219            pd.partition_key.as_ref().map(|k| k.as_str()),
220            pd.partition_hash.map(|h| h.to_sql_i64()),
221            pd.partition_id,
222            pd.partition_count,
223            pd.partition_strategy.as_ref().map(|s| s.as_str()),
224        ],
225    )
226    .map_err(|e| Error::Store(e.to_string()))?;
227    Ok(())
228}
229
230#[cfg(test)]
231mod tests {
232    use std::num::NonZeroU32;
233
234    use eventuary_core::io::Writer;
235    use eventuary_core::partition::{EventKeyPartitionKeyResolver, Fnv1a64PartitionHasher};
236    use eventuary_core::{Event, Payload};
237
238    use super::{SqlitePartitioningConfig, SqliteWriter, SqliteWriterConfig};
239    use crate::database::SqliteDatabase;
240
241    type PartitionRow = (
242        Option<String>,
243        Option<i64>,
244        Option<i64>,
245        Option<i64>,
246        Option<String>,
247    );
248
249    fn keyed_event(key: &str) -> Event {
250        Event::builder(
251            "acme",
252            "/orders",
253            "order.created",
254            key,
255            Payload::from_string("{}"),
256        )
257        .unwrap()
258        .build()
259        .unwrap()
260    }
261
262    fn query_partition_row(conn: &rusqlite::Connection, event_id: &str) -> PartitionRow {
263        conn.query_row(
264            "SELECT partition_key, partition_hash, partition_id, partition_count, partition_strategy FROM events WHERE id = ?1",
265            rusqlite::params![event_id],
266            |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?, row.get(4)?)),
267        )
268        .unwrap()
269    }
270
271    #[tokio::test]
272    async fn writer_off_partitioning_leaves_columns_null() {
273        let db = SqliteDatabase::open_in_memory().unwrap();
274        SqliteWriter::prepare_schema(&db.conn(), &SqliteWriterConfig::default()).unwrap();
275        let writer = SqliteWriter::new(db.conn());
276        let event = keyed_event("order-123");
277        writer.write(&event).await.unwrap();
278
279        let conn = db.conn();
280        let guard = conn.lock().unwrap();
281        let (pk, ph, pi, pc, ps) = query_partition_row(&guard, &event.id().as_uuid().to_string());
282
283        assert!(pk.is_none());
284        assert!(ph.is_none());
285        assert!(pi.is_none());
286        assert!(pc.is_none());
287        assert!(ps.is_none());
288    }
289
290    #[tokio::test]
291    async fn writer_inline_partitioning_persists_all_columns() {
292        let db = SqliteDatabase::open_in_memory().unwrap();
293        let config = SqliteWriterConfig {
294            partitioning: SqlitePartitioningConfig::inline(
295                NonZeroU32::new(64).unwrap(),
296                EventKeyPartitionKeyResolver::new(),
297                Fnv1a64PartitionHasher,
298            ),
299            ..SqliteWriterConfig::default()
300        };
301        SqliteWriter::prepare_schema(&db.conn(), &config).unwrap();
302        let writer = SqliteWriter::new_with_config(db.conn(), config);
303        let event = keyed_event("order-123");
304        writer.write(&event).await.unwrap();
305
306        let conn = db.conn();
307        let guard = conn.lock().unwrap();
308        let (pk, ph, pi, pc, ps) = query_partition_row(&guard, &event.id().as_uuid().to_string());
309
310        assert_eq!(pk.as_deref(), Some("order-123"));
311        assert_eq!(ph, Some(0x1b96f9c28b5d5aba_u64 as i64));
312        assert_eq!(pi, Some((0x1b96f9c28b5d5aba_u64 % 64) as i64));
313        assert_eq!(pc, Some(64));
314        assert_eq!(ps.as_deref(), Some("fnv1a64:v1"));
315    }
316}