Skip to main content

reifydb_core/interface/
change.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}