1use fallible_streaming_iterator::FallibleStreamingIterator;
4use rusqlite::hooks::Action;
5use rusqlite::session::ChangesetIter;
6use rusqlite::types::ValueRef;
7
8use coven_foundation::changeset::{ChangeOp, RowChange};
9
10#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub(crate) enum UpdateValue {
12 New,
13 Old,
14}
15
16enum ColumnCell<'value> {
17 Absent,
18 Present(ValueRef<'value>),
19}
20
21#[derive(Debug, thiserror::Error)]
22pub enum ChangesetError {
23 #[error("start changeset iterator: {0}")]
24 Start(#[source] rusqlite::Error),
25 #[error("advance changeset iterator: {0}")]
26 Next(#[source] rusqlite::Error),
27 #[error("read changeset operation: {0}")]
28 Operation(#[source] rusqlite::Error),
29 #[error("read changeset {side:?} value for column {column}: {source}")]
30 Value {
31 side: &'static str,
32 column: usize,
33 #[source]
34 source: rusqlite::Error,
35 },
36}
37
38pub fn walk(changeset_bytes: &[u8]) -> Result<Vec<RowChange>, ChangesetError> {
42 walk_with_update_values(changeset_bytes, UpdateValue::New)
43}
44
45pub fn walk_old(changeset_bytes: &[u8]) -> Result<Vec<RowChange>, ChangesetError> {
46 walk_with_update_values(changeset_bytes, UpdateValue::Old)
47}
48
49fn walk_with_update_values(
50 changeset_bytes: &[u8],
51 update_value: UpdateValue,
52) -> Result<Vec<RowChange>, ChangesetError> {
53 if changeset_bytes.is_empty() {
54 return Ok(Vec::new());
55 }
56
57 let input: &mut dyn std::io::Read = &mut &changeset_bytes[..];
58 let mut iter = ChangesetIter::start_strm(&input).map_err(ChangesetError::Start)?;
59
60 let mut changes = Vec::new();
61 while let Some(item) = iter.next().map_err(ChangesetError::Next)? {
62 let op = item.op().map_err(ChangesetError::Operation)?;
63 let change_op = match op.code() {
64 Action::SQLITE_INSERT => ChangeOp::Insert,
65 Action::SQLITE_UPDATE => ChangeOp::Update,
66 Action::SQLITE_DELETE => ChangeOp::Delete,
67 _ => continue,
68 };
69 let ncol = op.number_of_columns();
70
71 let cells = (0..ncol)
72 .map(|c| extract_col(item, c as usize, change_op, update_value))
73 .collect::<Result<Vec<_>, _>>()?;
74 let columns = cells
75 .iter()
76 .map(|(cell, _)| match cell {
77 ColumnCell::Absent => None,
78 ColumnCell::Present(value) => value_ref_to_string(*value),
79 })
80 .collect();
81 let changed_columns = cells.iter().map(|(_, changed)| *changed).collect();
82 changes.push(RowChange::new(
83 op.table_name().to_string(),
84 change_op,
85 columns,
86 changed_columns,
87 ));
88 }
89
90 Ok(changes)
91}
92
93fn extract_col<'value>(
97 item: &'value rusqlite::session::ChangesetItem,
98 col: usize,
99 op: ChangeOp,
100 update_value: UpdateValue,
101) -> Result<(ColumnCell<'value>, bool), ChangesetError> {
102 match op {
103 ChangeOp::Insert => changeset_value(item, col, UpdateValue::New).map(|cell| (cell, true)),
104 ChangeOp::Delete => changeset_value(item, col, UpdateValue::Old).map(|cell| (cell, true)),
105 ChangeOp::Update => {
106 let new = changeset_value(item, col, UpdateValue::New)?;
107 let old = changeset_value(item, col, UpdateValue::Old)?;
108 let changed = match (&new, &old) {
109 (ColumnCell::Present(new), ColumnCell::Present(old)) => new != old,
110 (ColumnCell::Present(_), ColumnCell::Absent) => true,
111 (ColumnCell::Absent, _) => false,
112 };
113 let cell = match update_value {
114 UpdateValue::New => match new {
115 ColumnCell::Absent => old,
116 present => present,
117 },
118 UpdateValue::Old => match old {
119 ColumnCell::Absent => new,
120 present => present,
121 },
122 };
123 Ok((cell, changed))
124 }
125 }
126}
127
128fn changeset_value<'value>(
129 item: &'value rusqlite::session::ChangesetItem,
130 col: usize,
131 side: UpdateValue,
132) -> Result<ColumnCell<'value>, ChangesetError> {
133 let value = match side {
134 UpdateValue::New => item.new_value(col),
135 UpdateValue::Old => item.old_value(col),
136 };
137 match value {
138 Ok(value) => Ok(ColumnCell::Present(value)),
139 Err(rusqlite::Error::InvalidColumnIndex(_)) => Ok(ColumnCell::Absent),
140 Err(source) => Err(ChangesetError::Value {
141 side: match side {
142 UpdateValue::New => "new",
143 UpdateValue::Old => "old",
144 },
145 column: col,
146 source,
147 }),
148 }
149}
150
151pub fn value_ref_to_string(v: ValueRef<'_>) -> Option<String> {
167 match v {
168 ValueRef::Null => None,
169 ValueRef::Integer(i) => Some(i.to_string()),
170 ValueRef::Real(f) => Some(real_to_sqlite_text(f)),
171 ValueRef::Text(t) | ValueRef::Blob(t) => Some(String::from_utf8_lossy(t).into_owned()),
172 }
173}
174
175fn real_to_sqlite_text(f: f64) -> String {
182 if f.is_nan() {
183 return String::new();
184 }
185 if f.is_infinite() {
186 return if f < 0.0 {
187 "-Inf".to_string()
188 } else {
189 "Inf".to_string()
190 };
191 }
192 let s = f.to_string();
193 if s.contains('.') || s.contains('e') || s.contains('E') {
197 s
198 } else {
199 format!("{s}.0")
200 }
201}
202
203#[cfg(test)]
204mod tests {
205 use super::*;
206
207 #[test]
213 fn real_renders_as_a_faithful_float() {
214 let conn = rusqlite::Connection::open_in_memory().expect("open");
215 for &f in &[0.0_f64, 1.0, 1.5, -2.25, 0.1, 123456.789, 1.0e6, 100.0, 0.5] {
217 let sqlite_text: String = conn
218 .query_row("SELECT CAST(? AS TEXT)", [f], |r| r.get(0))
219 .expect("cast");
220 let ours = real_to_sqlite_text(f);
221 assert_eq!(
222 ours, sqlite_text,
223 "REAL {f} rendered {ours:?}, SQLite renders {sqlite_text:?}",
224 );
225 }
226
227 for &f in &[1.234567890123457_f64, 1.0e-7, 9_999_999_999_999.0, -42.0] {
230 let ours = real_to_sqlite_text(f);
231 assert!(
232 ours.contains('.') || ours.contains('e') || ours.contains('E'),
233 "REAL {f} rendered {ours:?} which doesn't read as a float",
234 );
235 assert_eq!(
236 ours.parse::<f64>().expect("parses back"),
237 f,
238 "REAL {f} rendered {ours:?} which doesn't round-trip",
239 );
240 }
241 }
242}