1use std::collections::HashMap;
22use std::sync::Arc;
23
24use _serde::ViewVersionV1;
25use chrono::{DateTime, Utc};
26use serde::{Deserialize, Serialize};
27use typed_builder::TypedBuilder;
28
29use super::INITIAL_VIEW_VERSION_ID;
30use super::view_metadata::ViewVersionLog;
31use crate::catalog::NamespaceIdent;
32use crate::error::{Result, invalid_data, timestamp_ms_to_utc};
33use crate::spec::{SchemaId, SchemaRef, ViewMetadata};
34
35pub type ViewVersionRef = Arc<ViewVersion>;
37
38pub type ViewVersionId = i32;
40
41#[derive(Debug, PartialEq, Eq, Clone, Serialize, Deserialize, TypedBuilder)]
42#[serde(from = "ViewVersionV1", into = "ViewVersionV1")]
43#[builder(field_defaults(setter(prefix = "with_")))]
44pub struct ViewVersion {
46 #[builder(default = INITIAL_VIEW_VERSION_ID)]
48 version_id: ViewVersionId,
49 schema_id: SchemaId,
51 timestamp_ms: i64,
53 #[builder(default = HashMap::new())]
55 summary: HashMap<String, String>,
56 representations: ViewRepresentations,
58 #[builder(default = None)]
60 default_catalog: Option<String>,
61 default_namespace: NamespaceIdent,
63}
64
65impl ViewVersion {
66 #[inline]
68 pub fn version_id(&self) -> ViewVersionId {
69 self.version_id
70 }
71
72 #[inline]
74 pub fn schema_id(&self) -> SchemaId {
75 self.schema_id
76 }
77
78 #[inline]
80 pub fn timestamp(&self) -> Result<DateTime<Utc>> {
81 timestamp_ms_to_utc(self.timestamp_ms)
82 }
83
84 #[inline]
86 pub fn timestamp_ms(&self) -> i64 {
87 self.timestamp_ms
88 }
89
90 #[inline]
92 pub fn summary(&self) -> &HashMap<String, String> {
93 &self.summary
94 }
95
96 #[inline]
98 pub fn representations(&self) -> &ViewRepresentations {
99 &self.representations
100 }
101
102 #[inline]
104 pub fn default_catalog(&self) -> Option<&String> {
105 self.default_catalog.as_ref()
106 }
107
108 #[inline]
110 pub fn default_namespace(&self) -> &NamespaceIdent {
111 &self.default_namespace
112 }
113
114 pub fn schema(&self, view_metadata: &ViewMetadata) -> Result<SchemaRef> {
116 view_metadata
117 .schema_by_id(self.schema_id())
118 .ok_or_else(|| invalid_data!("Schema with id {} not found", self.schema_id()))
119 .cloned()
120 }
121
122 pub(crate) fn log(&self) -> ViewVersionLog {
124 ViewVersionLog::new(self.version_id, self.timestamp_ms)
125 }
126
127 pub(crate) fn behaves_identical_to(&self, other: &Self) -> bool {
137 self.summary == other.summary
138 && self.representations == other.representations
139 && self.default_catalog == other.default_catalog
140 && self.default_namespace == other.default_namespace
141 && self.schema_id == other.schema_id
142 }
143
144 pub fn with_version_id(self, version_id: i32) -> Self {
146 Self { version_id, ..self }
147 }
148
149 pub fn with_schema_id(self, schema_id: SchemaId) -> Self {
151 Self { schema_id, ..self }
152 }
153}
154
155#[derive(Debug, PartialEq, Eq, Clone, Serialize, Deserialize)]
157pub struct ViewRepresentations(pub(crate) Vec<ViewRepresentation>);
158
159impl ViewRepresentations {
160 #[inline]
161 pub fn len(&self) -> usize {
163 self.0.len()
164 }
165
166 #[inline]
167 pub fn is_empty(&self) -> bool {
169 self.0.is_empty()
170 }
171
172 pub fn iter(&self) -> impl ExactSizeIterator<Item = &'_ ViewRepresentation> {
174 self.0.iter()
175 }
176}
177
178impl IntoIterator for ViewRepresentations {
180 type Item = ViewRepresentation;
181 type IntoIter = std::vec::IntoIter<Self::Item>;
182
183 fn into_iter(self) -> Self::IntoIter {
184 self.0.into_iter()
185 }
186}
187
188#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
189#[serde(tag = "type")]
190pub enum ViewRepresentation {
194 #[serde(rename = "sql")]
195 Sql(SqlViewRepresentation),
197}
198
199#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
200#[serde(rename_all = "kebab-case")]
201pub struct SqlViewRepresentation {
204 #[serde(rename = "sql")]
205 pub sql: String,
207 #[serde(rename = "dialect")]
208 pub dialect: String,
210}
211
212pub(super) mod _serde {
213 use serde::{Deserialize, Serialize};
218
219 use super::{ViewRepresentation, ViewRepresentations, ViewVersion};
220 use crate::catalog::NamespaceIdent;
221
222 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
223 #[serde(rename_all = "kebab-case")]
224 pub(crate) struct ViewVersionV1 {
226 pub version_id: i32,
227 pub schema_id: i32,
228 pub timestamp_ms: i64,
229 pub summary: std::collections::HashMap<String, String>,
230 pub representations: Vec<ViewRepresentation>,
231 #[serde(skip_serializing_if = "Option::is_none")]
232 pub default_catalog: Option<String>,
233 pub default_namespace: NamespaceIdent,
234 }
235
236 impl From<ViewVersionV1> for ViewVersion {
237 fn from(v1: ViewVersionV1) -> Self {
238 ViewVersion {
239 version_id: v1.version_id,
240 schema_id: v1.schema_id,
241 timestamp_ms: v1.timestamp_ms,
242 summary: v1.summary,
243 representations: ViewRepresentations(v1.representations),
244 default_catalog: v1.default_catalog,
245 default_namespace: v1.default_namespace,
246 }
247 }
248 }
249
250 impl From<ViewVersion> for ViewVersionV1 {
251 fn from(v1: ViewVersion) -> Self {
252 ViewVersionV1 {
253 version_id: v1.version_id,
254 schema_id: v1.schema_id,
255 timestamp_ms: v1.timestamp_ms,
256 summary: v1.summary,
257 representations: v1.representations.0,
258 default_catalog: v1.default_catalog,
259 default_namespace: v1.default_namespace,
260 }
261 }
262 }
263}
264
265impl From<SqlViewRepresentation> for ViewRepresentation {
266 fn from(sql: SqlViewRepresentation) -> Self {
267 ViewRepresentation::Sql(sql)
268 }
269}
270
271#[cfg(test)]
272mod tests {
273 use chrono::{TimeZone, Utc};
274
275 use crate::NamespaceIdent;
276 use crate::spec::ViewRepresentations;
277 use crate::spec::view_version::_serde::ViewVersionV1;
278 use crate::spec::view_version::ViewVersion;
279
280 #[test]
281 fn view_version() {
282 let record = serde_json::json!(
283 {
284 "version-id" : 1,
285 "timestamp-ms" : 1573518431292i64,
286 "schema-id" : 1,
287 "default-catalog" : "prod",
288 "default-namespace" : [ "default" ],
289 "summary" : {
290 "engine-name" : "Spark",
291 "engineVersion" : "3.3.2"
292 },
293 "representations" : [ {
294 "type" : "sql",
295 "sql" : "SELECT\n COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2",
296 "dialect" : "spark"
297 } ]
298 }
299 );
300
301 let result: ViewVersion = serde_json::from_value::<ViewVersionV1>(record.clone())
302 .unwrap()
303 .into();
304
305 assert_eq!(serde_json::to_value(result.clone()).unwrap(), record);
307
308 assert_eq!(result.version_id(), 1);
309 assert_eq!(
310 result.timestamp().unwrap(),
311 Utc.timestamp_millis_opt(1573518431292).unwrap()
312 );
313 assert_eq!(result.schema_id(), 1);
314 assert_eq!(result.default_catalog, Some("prod".to_string()));
315 assert_eq!(result.summary(), &{
316 let mut map = std::collections::HashMap::new();
317 map.insert("engine-name".to_string(), "Spark".to_string());
318 map.insert("engineVersion".to_string(), "3.3.2".to_string());
319 map
320 });
321 assert_eq!(
322 result.representations().to_owned(),
323 ViewRepresentations(vec![super::ViewRepresentation::Sql(
324 super::SqlViewRepresentation {
325 sql: "SELECT\n COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2"
326 .to_string(),
327 dialect: "spark".to_string(),
328 },
329 )])
330 );
331 assert_eq!(result.default_namespace.inner(), vec![
332 "default".to_string()
333 ]);
334 }
335
336 #[test]
337 fn test_behaves_identical_to() {
338 let view_version = ViewVersion::builder()
339 .with_version_id(1)
340 .with_schema_id(1)
341 .with_timestamp_ms(1573518431292)
342 .with_summary({
343 let mut map = std::collections::HashMap::new();
344 map.insert("engine-name".to_string(), "Spark".to_string());
345 map.insert("engineVersion".to_string(), "3.3.2".to_string());
346 map
347 })
348 .with_representations(ViewRepresentations(vec![super::ViewRepresentation::Sql(
349 super::SqlViewRepresentation {
350 sql: "SELECT\n COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2"
351 .to_string(),
352 dialect: "spark".to_string(),
353 },
354 )]))
355 .with_default_catalog(Some("prod".to_string()))
356 .with_default_namespace(NamespaceIdent::new("default".to_string()))
357 .build();
358
359 let mut identical_view_version = view_version.clone();
360 identical_view_version.version_id = 2;
361 identical_view_version.timestamp_ms = 1573518431293;
362
363 let different_view_version = ViewVersion::builder()
364 .with_version_id(view_version.version_id())
365 .with_schema_id(view_version.schema_id())
366 .with_timestamp_ms(view_version.timestamp_ms())
367 .with_summary(view_version.summary().clone())
368 .with_representations(ViewRepresentations(vec![super::ViewRepresentation::Sql(
369 super::SqlViewRepresentation {
370 sql: "SELECT * from events".to_string(),
371 dialect: "spark".to_string(),
372 },
373 )]))
374 .with_default_catalog(view_version.default_catalog().cloned())
375 .with_default_namespace(view_version.default_namespace().clone())
376 .build();
377
378 assert!(view_version.behaves_identical_to(&identical_view_version));
379 assert!(!view_version.behaves_identical_to(&different_view_version));
380 }
381}