1use std::fmt;
2use std::num::NonZeroU16;
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: NonZeroU16,
21 key_resolver: Arc<dyn PartitionKeyResolver>,
22 hasher: Arc<dyn PartitionHasher>,
23 },
24}
25
26impl SqlitePartitioningConfig {
27 pub fn inline(
28 count: NonZeroU16,
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::NonZeroU16;
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 NonZeroU16::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}