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