1use std::cmp::Ordering;
22use std::collections::HashMap;
23use std::fmt::{Display, Formatter};
24use std::sync::Arc;
25
26use _serde::ViewMetadataEnum;
27use chrono::{DateTime, Utc};
28use serde::{Deserialize, Serialize};
29use serde_repr::{Deserialize_repr, Serialize_repr};
30use uuid::Uuid;
31
32pub use super::view_metadata_builder::ViewMetadataBuilder;
33use super::view_version::{ViewVersionId, ViewVersionRef};
34use super::{SchemaId, SchemaRef};
35use crate::error::{Result, invalid_data, timestamp_ms_to_utc};
36
37pub type ViewMetadataRef = Arc<ViewMetadata>;
39
40pub(crate) static INITIAL_VIEW_VERSION_ID: i32 = 1;
42
43pub const VIEW_PROPERTY_REPLACE_DROP_DIALECT_ALLOWED: &str = "replace.drop-dialect.allowed";
45pub const VIEW_PROPERTY_REPLACE_DROP_DIALECT_ALLOWED_DEFAULT: bool = false;
47pub const VIEW_PROPERTY_VERSION_HISTORY_SIZE: &str = "version.history.num-entries";
49pub const VIEW_PROPERTY_VERSION_HISTORY_SIZE_DEFAULT: usize = 10;
51
52#[derive(Debug, PartialEq, Deserialize, Eq, Clone)]
53#[serde(try_from = "ViewMetadataEnum", into = "ViewMetadataEnum")]
54pub struct ViewMetadata {
59 pub(crate) format_version: ViewFormatVersion,
61 pub(crate) view_uuid: Uuid,
63 pub(crate) location: String,
65 pub(crate) current_version_id: ViewVersionId,
67 pub(crate) versions: HashMap<ViewVersionId, ViewVersionRef>,
69 pub(crate) version_log: Vec<ViewVersionLog>,
72 pub(crate) schemas: HashMap<SchemaId, SchemaRef>,
74 pub(crate) properties: HashMap<String, String>,
78}
79
80impl ViewMetadata {
81 #[must_use]
83 pub fn into_builder(self) -> ViewMetadataBuilder {
84 ViewMetadataBuilder::new_from_metadata(self)
85 }
86
87 #[inline]
89 pub fn format_version(&self) -> ViewFormatVersion {
90 self.format_version
91 }
92
93 #[inline]
95 pub fn uuid(&self) -> Uuid {
96 self.view_uuid
97 }
98
99 #[inline]
101 pub fn location(&self) -> &str {
102 self.location.as_str()
103 }
104
105 #[inline]
107 pub fn current_version_id(&self) -> ViewVersionId {
108 self.current_version_id
109 }
110
111 #[inline]
113 pub fn versions(&self) -> impl ExactSizeIterator<Item = &ViewVersionRef> {
114 self.versions.values()
115 }
116
117 #[inline]
119 pub fn version_by_id(&self, version_id: ViewVersionId) -> Option<&ViewVersionRef> {
120 self.versions.get(&version_id)
121 }
122
123 #[inline]
125 pub fn current_version(&self) -> &ViewVersionRef {
126 self.versions
127 .get(&self.current_version_id)
128 .expect("Current version id set, but not found in view versions")
129 }
130
131 #[inline]
133 pub fn schemas_iter(&self) -> impl ExactSizeIterator<Item = &SchemaRef> {
134 self.schemas.values()
135 }
136
137 #[inline]
139 pub fn schema_by_id(&self, schema_id: SchemaId) -> Option<&SchemaRef> {
140 self.schemas.get(&schema_id)
141 }
142
143 #[inline]
145 pub fn current_schema(&self) -> &SchemaRef {
146 let schema_id = self.current_version().schema_id();
147 self.schema_by_id(schema_id)
148 .expect("Current schema id set, but not found in view metadata")
149 }
150
151 #[inline]
153 pub fn properties(&self) -> &HashMap<String, String> {
154 &self.properties
155 }
156
157 #[inline]
159 pub fn history(&self) -> &[ViewVersionLog] {
160 &self.version_log
161 }
162
163 pub(super) fn validate(&self) -> Result<()> {
165 self.validate_current_version_id()?;
166 self.validate_current_schema_id()?;
167 Ok(())
168 }
169
170 fn validate_current_version_id(&self) -> Result<()> {
171 if !self.versions.contains_key(&self.current_version_id) {
172 return Err(invalid_data!(
173 "No version exists with the current version id {}.",
174 self.current_version_id
175 ));
176 }
177 Ok(())
178 }
179
180 fn validate_current_schema_id(&self) -> Result<()> {
181 let schema_id = self.current_version().schema_id();
182 if !self.schemas.contains_key(&schema_id) {
183 return Err(invalid_data!(
184 "No schema exists with the schema id {schema_id}."
185 ));
186 }
187 Ok(())
188 }
189}
190
191#[derive(Debug, Serialize, Deserialize, PartialEq, Eq, Clone)]
192#[serde(rename_all = "kebab-case")]
193pub struct ViewVersionLog {
195 version_id: ViewVersionId,
197 timestamp_ms: i64,
199}
200
201impl ViewVersionLog {
202 #[inline]
203 pub fn new(version_id: ViewVersionId, timestamp: i64) -> Self {
205 Self {
206 version_id,
207 timestamp_ms: timestamp,
208 }
209 }
210
211 #[inline]
213 pub fn version_id(&self) -> ViewVersionId {
214 self.version_id
215 }
216
217 #[inline]
219 pub fn timestamp_ms(&self) -> i64 {
220 self.timestamp_ms
221 }
222
223 pub fn timestamp(&self) -> Result<DateTime<Utc>> {
225 timestamp_ms_to_utc(self.timestamp_ms)
226 }
227
228 pub(crate) fn set_timestamp_ms(&mut self, timestamp_ms: i64) -> &mut Self {
230 self.timestamp_ms = timestamp_ms;
231 self
232 }
233}
234
235pub(super) mod _serde {
236 use std::{collections::HashMap, sync::Arc};
241
242 use serde::{Deserialize, Serialize};
243 use uuid::Uuid;
244
245 use super::{ViewFormatVersion, ViewVersionId, ViewVersionLog};
246 use crate::Error;
247 use crate::spec::schema::_serde::SchemaV2;
248 use crate::spec::table_metadata::_serde::VersionNumber;
249 use crate::spec::view_version::_serde::ViewVersionV1;
250 use crate::spec::{ViewMetadata, ViewVersion};
251
252 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
253 #[serde(untagged)]
254 pub(super) enum ViewMetadataEnum {
255 V1(ViewMetadataV1),
256 }
257
258 #[derive(Debug, Serialize, Deserialize, PartialEq, Eq)]
259 #[serde(rename_all = "kebab-case")]
260 pub(super) struct ViewMetadataV1 {
262 pub format_version: VersionNumber<1>,
263 pub(super) view_uuid: Uuid,
264 pub(super) location: String,
265 pub(super) current_version_id: ViewVersionId,
266 pub(super) versions: Vec<ViewVersionV1>,
267 pub(super) version_log: Vec<ViewVersionLog>,
268 pub(super) schemas: Vec<SchemaV2>,
269 pub(super) properties: Option<HashMap<String, String>>,
270 }
271
272 impl Serialize for ViewMetadata {
273 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
274 where S: serde::Serializer {
275 let metadata_enum: ViewMetadataEnum =
277 self.clone().try_into().map_err(serde::ser::Error::custom)?;
278
279 metadata_enum.serialize(serializer)
280 }
281 }
282
283 impl TryFrom<ViewMetadataEnum> for ViewMetadata {
284 type Error = Error;
285 fn try_from(value: ViewMetadataEnum) -> Result<Self, Error> {
286 match value {
287 ViewMetadataEnum::V1(value) => value.try_into(),
288 }
289 }
290 }
291
292 impl TryFrom<ViewMetadata> for ViewMetadataEnum {
293 type Error = Error;
294 fn try_from(value: ViewMetadata) -> Result<Self, Error> {
295 Ok(match value.format_version {
296 ViewFormatVersion::V1 => ViewMetadataEnum::V1(value.into()),
297 })
298 }
299 }
300
301 impl TryFrom<ViewMetadataV1> for ViewMetadata {
302 type Error = Error;
303 fn try_from(value: ViewMetadataV1) -> Result<Self, Error> {
304 let schemas = HashMap::from_iter(
305 value
306 .schemas
307 .into_iter()
308 .map(|schema| Ok((schema.schema_id, Arc::new(schema.try_into()?))))
309 .collect::<Result<Vec<_>, Error>>()?,
310 );
311 let versions = HashMap::from_iter(
312 value
313 .versions
314 .into_iter()
315 .map(|x| Ok((x.version_id, Arc::new(ViewVersion::from(x)))))
316 .collect::<Result<Vec<_>, Error>>()?,
317 );
318
319 let view_metadata = ViewMetadata {
320 format_version: ViewFormatVersion::V1,
321 view_uuid: value.view_uuid,
322 location: value.location,
323 schemas,
324 properties: value.properties.unwrap_or_default(),
325 current_version_id: value.current_version_id,
326 versions,
327 version_log: value.version_log,
328 };
329 view_metadata.validate()?;
330 Ok(view_metadata)
331 }
332 }
333
334 impl From<ViewMetadata> for ViewMetadataV1 {
335 fn from(v: ViewMetadata) -> Self {
336 let schemas = v
337 .schemas
338 .into_values()
339 .map(|x| {
340 Arc::try_unwrap(x)
341 .unwrap_or_else(|schema| schema.as_ref().clone())
342 .into()
343 })
344 .collect();
345 let versions = v
346 .versions
347 .into_values()
348 .map(|x| {
349 Arc::try_unwrap(x)
350 .unwrap_or_else(|version| version.as_ref().clone())
351 .into()
352 })
353 .collect();
354 ViewMetadataV1 {
355 format_version: VersionNumber::<1>,
356 view_uuid: v.view_uuid,
357 location: v.location,
358 schemas,
359 properties: Some(v.properties),
360 current_version_id: v.current_version_id,
361 versions,
362 version_log: v.version_log,
363 }
364 }
365 }
366}
367
368#[derive(Debug, Serialize_repr, Deserialize_repr, PartialEq, Eq, Clone, Copy)]
369#[repr(u8)]
370pub enum ViewFormatVersion {
372 V1 = 1u8,
374}
375
376impl PartialOrd for ViewFormatVersion {
377 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
378 Some(self.cmp(other))
379 }
380}
381
382impl Ord for ViewFormatVersion {
383 fn cmp(&self, other: &Self) -> Ordering {
384 (*self as u8).cmp(&(*other as u8))
385 }
386}
387
388impl Display for ViewFormatVersion {
389 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
390 match self {
391 ViewFormatVersion::V1 => write!(f, "v1"),
392 }
393 }
394}
395
396#[cfg(test)]
397pub(crate) mod tests {
398 use std::collections::HashMap;
399 use std::fs;
400 use std::sync::Arc;
401
402 use anyhow::Result;
403 use pretty_assertions::assert_eq;
404 use uuid::Uuid;
405
406 use super::{ViewFormatVersion, ViewMetadataBuilder, ViewVersionLog};
407 use crate::spec::{
408 INITIAL_VIEW_VERSION_ID, NestedField, PrimitiveType, Schema, SqlViewRepresentation, Type,
409 ViewMetadata, ViewRepresentations, ViewVersion,
410 };
411 use crate::{NamespaceIdent, ViewCreation};
412
413 fn check_view_metadata_serde(json: &str, expected_type: ViewMetadata) {
414 let desered_type: ViewMetadata = serde_json::from_str(json).unwrap();
415 assert_eq!(desered_type, expected_type);
416
417 let sered_json = serde_json::to_string(&expected_type).unwrap();
418 let parsed_json_value = serde_json::from_str::<ViewMetadata>(&sered_json).unwrap();
419
420 assert_eq!(parsed_json_value, desered_type);
421 }
422
423 pub(crate) fn get_test_view_metadata(file_name: &str) -> ViewMetadata {
424 let path = format!("testdata/view_metadata/{file_name}");
425 let metadata: String = fs::read_to_string(path).unwrap();
426
427 serde_json::from_str(&metadata).unwrap()
428 }
429
430 #[test]
431 fn test_view_data_v1() {
432 let data = r#"
433 {
434 "view-uuid": "fa6506c3-7681-40c8-86dc-e36561f83385",
435 "format-version" : 1,
436 "location" : "s3://bucket/warehouse/default.db/event_agg",
437 "current-version-id" : 1,
438 "properties" : {
439 "comment" : "Daily event counts"
440 },
441 "versions" : [ {
442 "version-id" : 1,
443 "timestamp-ms" : 1573518431292,
444 "schema-id" : 1,
445 "default-catalog" : "prod",
446 "default-namespace" : [ "default" ],
447 "summary" : {
448 "engine-name" : "Spark",
449 "engineVersion" : "3.3.2"
450 },
451 "representations" : [ {
452 "type" : "sql",
453 "sql" : "SELECT\n COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2",
454 "dialect" : "spark"
455 } ]
456 } ],
457 "schemas": [ {
458 "schema-id": 1,
459 "type" : "struct",
460 "fields" : [ {
461 "id" : 1,
462 "name" : "event_count",
463 "required" : false,
464 "type" : "int",
465 "doc" : "Count of events"
466 } ]
467 } ],
468 "version-log" : [ {
469 "timestamp-ms" : 1573518431292,
470 "version-id" : 1
471 } ]
472 }
473 "#;
474
475 let schema = Schema::builder()
476 .with_schema_id(1)
477 .with_fields(vec![Arc::new(
478 NestedField::optional(1, "event_count", Type::Primitive(PrimitiveType::Int))
479 .with_doc("Count of events"),
480 )])
481 .build()
482 .unwrap();
483 let version = ViewVersion::builder()
484 .with_version_id(1)
485 .with_timestamp_ms(1573518431292)
486 .with_schema_id(1)
487 .with_default_catalog("prod".to_string().into())
488 .with_default_namespace(NamespaceIdent::from_vec(vec!["default".to_string()]).unwrap())
489 .with_summary(HashMap::from_iter(vec![
490 ("engineVersion".to_string(), "3.3.2".to_string()),
491 ("engine-name".to_string(), "Spark".to_string()),
492 ]))
493 .with_representations(ViewRepresentations(vec![
494 SqlViewRepresentation {
495 sql: "SELECT\n COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2"
496 .to_string(),
497 dialect: "spark".to_string(),
498 }
499 .into(),
500 ]))
501 .build();
502
503 let expected = ViewMetadata {
504 format_version: ViewFormatVersion::V1,
505 view_uuid: Uuid::parse_str("fa6506c3-7681-40c8-86dc-e36561f83385").unwrap(),
506 location: "s3://bucket/warehouse/default.db/event_agg".to_string(),
507 current_version_id: 1,
508 versions: HashMap::from_iter(vec![(1, Arc::new(version))]),
509 version_log: vec![ViewVersionLog {
510 timestamp_ms: 1573518431292,
511 version_id: 1,
512 }],
513 schemas: HashMap::from_iter(vec![(1, Arc::new(schema))]),
514 properties: HashMap::from_iter(vec![(
515 "comment".to_string(),
516 "Daily event counts".to_string(),
517 )]),
518 };
519
520 check_view_metadata_serde(data, expected);
521 }
522
523 #[test]
524 fn test_invalid_view_uuid() -> Result<()> {
525 let data = r#"
526 {
527 "format-version" : 1,
528 "view-uuid": "xxxx"
529 }
530 "#;
531 assert!(serde_json::from_str::<ViewMetadata>(data).is_err());
532 Ok(())
533 }
534
535 #[test]
536 fn test_view_builder_from_view_creation() {
537 let representations = ViewRepresentations(vec![
538 SqlViewRepresentation {
539 sql: "SELECT\n COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2"
540 .to_string(),
541 dialect: "spark".to_string(),
542 }
543 .into(),
544 ]);
545 let creation = ViewCreation::builder()
546 .location("s3://bucket/warehouse/default.db/event_agg".to_string())
547 .name("view".to_string())
548 .schema(Schema::builder().build().unwrap())
549 .default_namespace(NamespaceIdent::from_vec(vec!["default".to_string()]).unwrap())
550 .representations(representations)
551 .build();
552
553 let metadata = ViewMetadataBuilder::from_view_creation(creation)
554 .unwrap()
555 .build()
556 .unwrap()
557 .metadata;
558
559 assert_eq!(
560 metadata.location(),
561 "s3://bucket/warehouse/default.db/event_agg"
562 );
563 assert_eq!(metadata.current_version_id(), INITIAL_VIEW_VERSION_ID);
564 assert_eq!(metadata.versions().count(), 1);
565 assert_eq!(metadata.schemas_iter().count(), 1);
566 assert_eq!(metadata.properties().len(), 0);
567 }
568
569 #[test]
570 fn test_view_metadata_v1_file_valid() {
571 let metadata =
572 fs::read_to_string("testdata/view_metadata/ViewMetadataV1Valid.json").unwrap();
573
574 let schema = Schema::builder()
575 .with_schema_id(1)
576 .with_fields(vec![
577 Arc::new(
578 NestedField::optional(1, "event_count", Type::Primitive(PrimitiveType::Int))
579 .with_doc("Count of events"),
580 ),
581 Arc::new(NestedField::optional(
582 2,
583 "event_date",
584 Type::Primitive(PrimitiveType::Date),
585 )),
586 ])
587 .build()
588 .unwrap();
589
590 let version = ViewVersion::builder()
591 .with_version_id(1)
592 .with_timestamp_ms(1573518431292)
593 .with_schema_id(1)
594 .with_default_catalog("prod".to_string().into())
595 .with_default_namespace(NamespaceIdent::from_vec(vec!["default".to_string()]).unwrap())
596 .with_summary(HashMap::from_iter(vec![
597 ("engineVersion".to_string(), "3.3.2".to_string()),
598 ("engine-name".to_string(), "Spark".to_string()),
599 ]))
600 .with_representations(ViewRepresentations(vec![
601 SqlViewRepresentation {
602 sql: "SELECT\n COUNT(1), CAST(event_ts AS DATE)\nFROM events\nGROUP BY 2"
603 .to_string(),
604 dialect: "spark".to_string(),
605 }
606 .into(),
607 ]))
608 .build();
609
610 let expected = ViewMetadata {
611 format_version: ViewFormatVersion::V1,
612 view_uuid: Uuid::parse_str("fa6506c3-7681-40c8-86dc-e36561f83385").unwrap(),
613 location: "s3://bucket/warehouse/default.db/event_agg".to_string(),
614 current_version_id: 1,
615 versions: HashMap::from_iter(vec![(1, Arc::new(version))]),
616 version_log: vec![ViewVersionLog {
617 timestamp_ms: 1573518431292,
618 version_id: 1,
619 }],
620 schemas: HashMap::from_iter(vec![(1, Arc::new(schema))]),
621 properties: HashMap::from_iter(vec![(
622 "comment".to_string(),
623 "Daily event counts".to_string(),
624 )]),
625 };
626
627 check_view_metadata_serde(&metadata, expected);
628 }
629
630 #[test]
631 fn test_view_builder_assign_uuid() {
632 let metadata = get_test_view_metadata("ViewMetadataV1Valid.json");
633 let metadata_builder = metadata.into_builder();
634 let uuid = Uuid::new_v4();
635 let metadata = metadata_builder.assign_uuid(uuid).build().unwrap().metadata;
636 assert_eq!(metadata.uuid(), uuid);
637 }
638
639 #[test]
640 fn test_view_metadata_v1_unsupported_version() {
641 let metadata =
642 fs::read_to_string("testdata/view_metadata/ViewMetadataUnsupportedVersion.json")
643 .unwrap();
644
645 let desered: Result<ViewMetadata, serde_json::Error> = serde_json::from_str(&metadata);
646
647 assert_eq!(
648 desered.unwrap_err().to_string(),
649 "data did not match any variant of untagged enum ViewMetadataEnum"
650 )
651 }
652
653 #[test]
654 fn test_view_metadata_v1_version_not_found() {
655 let metadata =
656 fs::read_to_string("testdata/view_metadata/ViewMetadataV1CurrentVersionNotFound.json")
657 .unwrap();
658
659 let desered: Result<ViewMetadata, serde_json::Error> = serde_json::from_str(&metadata);
660
661 assert_eq!(
662 desered.unwrap_err().to_string(),
663 "DataInvalid => No version exists with the current version id 2."
664 )
665 }
666
667 #[test]
668 fn test_view_metadata_v1_schema_not_found() {
669 let metadata =
670 fs::read_to_string("testdata/view_metadata/ViewMetadataV1SchemaNotFound.json").unwrap();
671
672 let desered: Result<ViewMetadata, serde_json::Error> = serde_json::from_str(&metadata);
673
674 assert_eq!(
675 desered.unwrap_err().to_string(),
676 "DataInvalid => No schema exists with the schema id 2."
677 )
678 }
679
680 #[test]
681 fn test_view_metadata_v1_missing_schema_for_version() {
682 let metadata =
683 fs::read_to_string("testdata/view_metadata/ViewMetadataV1MissingSchema.json").unwrap();
684
685 let desered: Result<ViewMetadata, serde_json::Error> = serde_json::from_str(&metadata);
686
687 assert_eq!(
688 desered.unwrap_err().to_string(),
689 "data did not match any variant of untagged enum ViewMetadataEnum"
690 )
691 }
692
693 #[test]
694 fn test_view_metadata_v1_missing_current_version() {
695 let metadata =
696 fs::read_to_string("testdata/view_metadata/ViewMetadataV1MissingCurrentVersion.json")
697 .unwrap();
698
699 let desered: Result<ViewMetadata, serde_json::Error> = serde_json::from_str(&metadata);
700
701 assert_eq!(
702 desered.unwrap_err().to_string(),
703 "data did not match any variant of untagged enum ViewMetadataEnum"
704 )
705 }
706}