1use std::mem;
5
6use reifydb_value::{
7 Result,
8 value::{datetime::DateTime, diff_type::DiffType},
9};
10use serde::{Deserialize, Serialize};
11use smallvec::SmallVec;
12
13use crate::{
14 common::CommitVersion,
15 interface::{
16 catalog::{flow::OperatorId, object::ObjectId},
17 consolidate::coalesce_diffs,
18 },
19 value::column::columns::Columns,
20};
21
22pub type Diffs = SmallVec<[Diff; 4]>;
23
24pub type StagedBatch = (DiffType, Columns);
25
26#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
27pub enum ChangeOrigin {
28 Object(ObjectId),
29 Flow(OperatorId),
30}
31
32#[derive(Debug, Clone, Serialize, Deserialize)]
33pub enum Diff {
34 Insert {
35 post: Columns,
36 origin: Option<ChangeOrigin>,
37 },
38 Update {
39 pre: Columns,
40 post: Columns,
41 origin: Option<ChangeOrigin>,
42 },
43 Remove {
44 pre: Columns,
45 origin: Option<ChangeOrigin>,
46 },
47}
48
49impl Diff {
50 pub fn insert(post: Columns) -> Self {
51 Self::Insert {
52 post,
53 origin: None,
54 }
55 }
56
57 pub fn update(pre: Columns, post: Columns) -> Self {
58 Self::Update {
59 pre,
60 post,
61 origin: None,
62 }
63 }
64
65 pub fn remove(pre: Columns) -> Self {
66 Self::Remove {
67 pre,
68 origin: None,
69 }
70 }
71
72 pub fn pre(&self) -> Option<&Columns> {
73 match self {
74 Diff::Insert {
75 ..
76 } => None,
77 Diff::Update {
78 pre,
79 ..
80 } => Some(pre),
81 Diff::Remove {
82 pre,
83 ..
84 } => Some(pre),
85 }
86 }
87
88 pub fn post(&self) -> Option<&Columns> {
89 match self {
90 Diff::Insert {
91 post,
92 ..
93 } => Some(post),
94 Diff::Update {
95 post,
96 ..
97 } => Some(post),
98 Diff::Remove {
99 ..
100 } => None,
101 }
102 }
103
104 pub fn columns_mut(&mut self) -> impl Iterator<Item = &mut Columns> {
105 let pair: [Option<&mut Columns>; 2] = match self {
106 Diff::Insert {
107 post,
108 ..
109 } => [Some(post), None],
110 Diff::Update {
111 pre,
112 post,
113 ..
114 } => [Some(pre), Some(post)],
115 Diff::Remove {
116 pre,
117 ..
118 } => [Some(pre), None],
119 };
120 pair.into_iter().flatten()
121 }
122
123 pub fn kind(&self) -> DiffType {
124 match self {
125 Diff::Insert {
126 ..
127 } => DiffType::Insert,
128 Diff::Update {
129 ..
130 } => DiffType::Update,
131 Diff::Remove {
132 ..
133 } => DiffType::Remove,
134 }
135 }
136
137 pub fn row_count(&self) -> usize {
138 match self {
139 Diff::Insert {
140 post,
141 ..
142 } => post.row_count(),
143 Diff::Update {
144 post,
145 ..
146 } => post.row_count(),
147 Diff::Remove {
148 pre,
149 ..
150 } => pre.row_count(),
151 }
152 }
153
154 pub fn origin(&self) -> Option<&ChangeOrigin> {
155 match self {
156 Diff::Insert {
157 origin,
158 ..
159 } => origin.as_ref(),
160 Diff::Update {
161 origin,
162 ..
163 } => origin.as_ref(),
164 Diff::Remove {
165 origin,
166 ..
167 } => origin.as_ref(),
168 }
169 }
170
171 pub fn set_origin(&mut self, new_origin: Option<ChangeOrigin>) {
172 match self {
173 Diff::Insert {
174 origin,
175 ..
176 } => *origin = new_origin,
177 Diff::Update {
178 origin,
179 ..
180 } => *origin = new_origin,
181 Diff::Remove {
182 origin,
183 ..
184 } => *origin = new_origin,
185 }
186 }
187
188 pub fn effective_origin<'a>(&'a self, parent: &'a ChangeOrigin) -> &'a ChangeOrigin {
189 self.origin().unwrap_or(parent)
190 }
191}
192
193#[derive(Debug, Clone, Serialize, Deserialize)]
194pub struct Change {
195 pub origin: ChangeOrigin,
196
197 pub diffs: Diffs,
198
199 pub version: CommitVersion,
200
201 pub changed_at: DateTime,
202}
203
204impl Change {
205 pub fn from_object(
206 object: ObjectId,
207 version: CommitVersion,
208 diffs: impl Into<Diffs>,
209 changed_at: DateTime,
210 ) -> Self {
211 Self {
212 origin: ChangeOrigin::Object(object),
213 diffs: diffs.into(),
214 version,
215 changed_at,
216 }
217 }
218
219 pub fn from_flow(
220 from: OperatorId,
221 version: CommitVersion,
222 diffs: impl Into<Diffs>,
223 changed_at: DateTime,
224 ) -> Self {
225 Self {
226 origin: ChangeOrigin::Flow(from),
227 diffs: diffs.into(),
228 version,
229 changed_at,
230 }
231 }
232
233 pub fn row_count(&self) -> usize {
234 self.diffs.iter().map(Diff::row_count).sum()
235 }
236
237 pub fn merge(changes: Vec<Change>) -> Result<Change> {
238 let mut iter = changes.into_iter();
239 let mut merged = iter.next().expect("Change::merge requires at least one Change");
240 for mut ch in iter {
241 if ch.changed_at > merged.changed_at {
242 merged.changed_at = ch.changed_at;
243 }
244 if ch.origin != merged.origin {
245 for diff in ch.diffs.iter_mut() {
246 if diff.origin().is_none() {
247 diff.set_origin(Some(ch.origin.clone()));
248 }
249 }
250 }
251 merged.diffs.extend(ch.diffs);
252 }
253 merged.coalesce()?;
254 Ok(merged)
255 }
256
257 pub fn coalesce(&mut self) -> Result<()> {
258 if self.diffs.len() <= 1 {
259 return Ok(());
260 }
261 let original = mem::take(&mut self.diffs);
262 self.diffs = SmallVec::from_vec(coalesce_diffs(original.into_vec())?);
263 Ok(())
264 }
265}