Skip to main content

common_meta/key/
schema_name.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{BTreeMap, HashMap};
16use std::fmt::Display;
17
18use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
19use common_time::DatabaseTimeToLive;
20use futures::stream::BoxStream;
21use humantime_serde::re::humantime;
22use serde::{Deserialize, Serialize};
23use snafu::{OptionExt, ResultExt, ensure};
24use store_api::mito_engine_options::{
25    TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_TRIGGER_FILE_NUM, normalize_twcs_trigger_options,
26};
27
28use crate::ensure_values;
29use crate::error::{
30    self, ConflictingSchemaOptionsSnafu, Error, InvalidMetadataSnafu, ParseOptionSnafu, Result,
31};
32use crate::key::txn_helper::TxnOpGetResponseSet;
33use crate::key::{
34    DeserializedValueWithBytes, MetadataKey, SCHEMA_NAME_KEY_PATTERN, SCHEMA_NAME_KEY_PREFIX,
35};
36use crate::kv_backend::KvBackendRef;
37use crate::kv_backend::txn::Txn;
38use crate::range_stream::{DEFAULT_PAGE_SIZE, PaginationStream};
39use crate::rpc::KeyValue;
40use crate::rpc::store::RangeRequest;
41
42const OPT_KEY_TTL: &str = "ttl";
43
44/// The schema name key, indices all schema names belong to the {catalog_name}
45///
46/// The layout:  `__schema_name/{catalog_name}/{schema_name}`.
47#[derive(Debug, Clone, Copy, PartialEq)]
48pub struct SchemaNameKey<'a> {
49    pub catalog: &'a str,
50    pub schema: &'a str,
51}
52
53impl Default for SchemaNameKey<'_> {
54    fn default() -> Self {
55        Self {
56            catalog: DEFAULT_CATALOG_NAME,
57            schema: DEFAULT_SCHEMA_NAME,
58        }
59    }
60}
61
62#[derive(Debug, Default, Clone, PartialEq, Serialize, Deserialize)]
63pub struct SchemaNameValue {
64    #[serde(default)]
65    pub ttl: Option<DatabaseTimeToLive>,
66    #[serde(default)]
67    pub extra_options: BTreeMap<String, String>,
68    /// Identifies the create-database procedure that wrote this value.
69    #[serde(default, skip_serializing_if = "Option::is_none")]
70    pub create_procedure_id: Option<String>,
71}
72
73impl Display for SchemaNameValue {
74    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
75        if let Some(ttl) = self.ttl.map(|i| i.to_string()) {
76            writeln!(f, "'ttl'='{}'", ttl)?;
77        }
78        for (k, v) in self.extra_options.iter() {
79            writeln!(f, "'{k}'='{v}'")?;
80        }
81
82        Ok(())
83    }
84}
85
86impl TryFrom<&HashMap<String, String>> for SchemaNameValue {
87    type Error = Error;
88
89    fn try_from(value: &HashMap<String, String>) -> std::result::Result<Self, Self::Error> {
90        let mut value = value.clone();
91        normalize_twcs_trigger_options(&mut value).map_err(|conflict| {
92            ConflictingSchemaOptionsSnafu {
93                first_key: TWCS_TRIGGER_FILE_NUM,
94                first_value: conflict.legacy_value,
95                second_key: TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM,
96                second_value: conflict.canonical_value,
97            }
98            .build()
99        })?;
100
101        let ttl = value
102            .get(OPT_KEY_TTL)
103            .map(|ttl_str| {
104                ttl_str.parse::<humantime::Duration>().map_err(|_| {
105                    ParseOptionSnafu {
106                        key: OPT_KEY_TTL,
107                        value: ttl_str.clone(),
108                    }
109                    .build()
110                })
111            })
112            .transpose()?
113            .map(|ttl| ttl.into());
114        let extra_options = value
115            .iter()
116            .filter_map(|(k, v)| {
117                if k == OPT_KEY_TTL {
118                    None
119                } else {
120                    Some((k.clone(), v.clone()))
121                }
122            })
123            .collect();
124
125        Ok(Self {
126            ttl,
127            extra_options,
128            ..Default::default()
129        })
130    }
131}
132
133impl From<SchemaNameValue> for HashMap<String, String> {
134    fn from(value: SchemaNameValue) -> Self {
135        let mut opts = HashMap::new();
136        if let Some(ttl) = value.ttl.map(|ttl| ttl.to_string()) {
137            opts.insert(OPT_KEY_TTL.to_string(), ttl);
138        }
139        opts.extend(
140            value
141                .extra_options
142                .iter()
143                .map(|(k, v)| (k.clone(), v.clone())),
144        );
145        opts
146    }
147}
148
149impl<'a> SchemaNameKey<'a> {
150    pub fn new(catalog: &'a str, schema: &'a str) -> Self {
151        Self { catalog, schema }
152    }
153
154    pub fn range_start_key(catalog: &str) -> String {
155        format!("{}/{}/", SCHEMA_NAME_KEY_PREFIX, catalog)
156    }
157}
158
159impl Display for SchemaNameKey<'_> {
160    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
161        write!(
162            f,
163            "{}/{}/{}",
164            SCHEMA_NAME_KEY_PREFIX, self.catalog, self.schema
165        )
166    }
167}
168
169impl<'a> MetadataKey<'a, SchemaNameKey<'a>> for SchemaNameKey<'_> {
170    fn to_bytes(&self) -> Vec<u8> {
171        self.to_string().into_bytes()
172    }
173
174    fn from_bytes(bytes: &'a [u8]) -> Result<SchemaNameKey<'a>> {
175        let key = std::str::from_utf8(bytes).map_err(|e| {
176            InvalidMetadataSnafu {
177                err_msg: format!(
178                    "SchemaNameKey '{}' is not a valid UTF8 string: {e}",
179                    String::from_utf8_lossy(bytes)
180                ),
181            }
182            .build()
183        })?;
184        SchemaNameKey::try_from(key)
185    }
186}
187
188/// Decodes `KeyValue` to {schema}
189pub fn schema_decoder(kv: KeyValue) -> Result<String> {
190    let str = std::str::from_utf8(&kv.key).context(error::ConvertRawKeySnafu)?;
191    let schema_name = SchemaNameKey::try_from(str)?;
192
193    Ok(schema_name.schema.to_string())
194}
195
196impl<'a> TryFrom<&'a str> for SchemaNameKey<'a> {
197    type Error = Error;
198
199    fn try_from(s: &'a str) -> Result<Self> {
200        let captures = SCHEMA_NAME_KEY_PATTERN
201            .captures(s)
202            .context(InvalidMetadataSnafu {
203                err_msg: format!("Illegal SchemaNameKey format: '{s}'"),
204            })?;
205
206        // Safety: pass the regex check above
207        Ok(Self {
208            catalog: captures.get(1).unwrap().as_str(),
209            schema: captures.get(2).unwrap().as_str(),
210        })
211    }
212}
213
214#[derive(Clone)]
215pub struct SchemaManager {
216    kv_backend: KvBackendRef,
217}
218
219pub type SchemaNameDecodeResult = Result<Option<DeserializedValueWithBytes<SchemaNameValue>>>;
220
221impl SchemaManager {
222    pub fn new(kv_backend: KvBackendRef) -> Self {
223        Self { kv_backend }
224    }
225
226    /// Creates `SchemaNameKey`.
227    pub async fn create(
228        &self,
229        schema: SchemaNameKey<'_>,
230        value: Option<SchemaNameValue>,
231        if_not_exists: bool,
232    ) -> Result<()> {
233        let _timer = crate::metrics::METRIC_META_CREATE_SCHEMA.start_timer();
234
235        let raw_key = schema.to_bytes();
236        let raw_value = value.unwrap_or_default().try_as_raw_value()?;
237        if self
238            .kv_backend
239            .put_conditionally(raw_key, raw_value, if_not_exists)
240            .await?
241        {
242            crate::metrics::METRIC_META_CREATE_SCHEMA_COUNTER.inc();
243        }
244
245        Ok(())
246    }
247
248    pub async fn exists(&self, schema: SchemaNameKey<'_>) -> Result<bool> {
249        let raw_key = schema.to_bytes();
250
251        self.kv_backend.exists(&raw_key).await
252    }
253
254    pub async fn get(
255        &self,
256        schema: SchemaNameKey<'_>,
257    ) -> Result<Option<DeserializedValueWithBytes<SchemaNameValue>>> {
258        let raw_key = schema.to_bytes();
259        self.kv_backend
260            .get(&raw_key)
261            .await?
262            .map(|x| DeserializedValueWithBytes::from_inner_slice(&x.value))
263            .transpose()
264    }
265
266    /// Deletes a [SchemaNameKey].
267    pub async fn delete(&self, schema: SchemaNameKey<'_>) -> Result<()> {
268        let raw_key = schema.to_bytes();
269        self.kv_backend.delete(&raw_key, false).await?;
270
271        Ok(())
272    }
273
274    pub(crate) fn build_update_txn(
275        &self,
276        schema: SchemaNameKey<'_>,
277        current_schema_value: &DeserializedValueWithBytes<SchemaNameValue>,
278        new_schema_value: &SchemaNameValue,
279    ) -> Result<(
280        Txn,
281        impl FnOnce(&mut TxnOpGetResponseSet) -> SchemaNameDecodeResult,
282    )> {
283        let raw_key = schema.to_bytes();
284        let raw_value = current_schema_value.get_raw_bytes();
285        let new_raw_value: Vec<u8> = new_schema_value.try_as_raw_value()?;
286
287        let txn = Txn::compare_and_put(raw_key.clone(), raw_value, new_raw_value);
288
289        Ok((
290            txn,
291            TxnOpGetResponseSet::decode_with(TxnOpGetResponseSet::filter(raw_key)),
292        ))
293    }
294
295    /// Updates a [SchemaNameKey].
296    pub async fn update(
297        &self,
298        schema: SchemaNameKey<'_>,
299        current_schema_value: &DeserializedValueWithBytes<SchemaNameValue>,
300        new_schema_value: &SchemaNameValue,
301    ) -> Result<()> {
302        let (txn, on_failure) =
303            self.build_update_txn(schema, current_schema_value, new_schema_value)?;
304        let mut r = self.kv_backend.txn(txn).await?;
305
306        if !r.succeeded {
307            let mut set = TxnOpGetResponseSet::from(&mut r.responses);
308            let remote_schema_value = on_failure(&mut set)?
309                .context(error::UnexpectedSnafu {
310                    err_msg:
311                        "Reads the empty schema name value in comparing operation of updating schema name value",
312                })?
313                .into_inner();
314
315            let op_name = "the updating schema name value";
316            ensure_values!(&remote_schema_value, new_schema_value, op_name);
317        }
318
319        Ok(())
320    }
321
322    /// Returns a schema stream, it lists all schemas belong to the target `catalog`.
323    pub fn schema_names(&self, catalog: &str) -> BoxStream<'static, Result<String>> {
324        let start_key = SchemaNameKey::range_start_key(catalog);
325        let req = RangeRequest::new().with_prefix(start_key.as_bytes());
326
327        let stream = PaginationStream::new(
328            self.kv_backend.clone(),
329            req,
330            DEFAULT_PAGE_SIZE,
331            schema_decoder,
332        )
333        .into_stream();
334
335        Box::pin(stream)
336    }
337
338    /// Returns schema names and values belonging to the target `catalog`.
339    /// Legacy `null` values are returned as [`SchemaNameValue::default()`].
340    pub fn schemas(&self, catalog: &str) -> BoxStream<'static, Result<(String, SchemaNameValue)>> {
341        let start_key = SchemaNameKey::range_start_key(catalog);
342        let req = RangeRequest::new().with_prefix(start_key.as_bytes());
343
344        let stream = PaginationStream::new(self.kv_backend.clone(), req, DEFAULT_PAGE_SIZE, |kv| {
345            let value = SchemaNameValue::try_from_raw_value(&kv.value)?.unwrap_or_default();
346            Ok((schema_decoder(kv)?, value))
347        })
348        .into_stream();
349
350        Box::pin(stream)
351    }
352}
353
354#[derive(Debug, Clone, Hash, Eq, PartialEq, Deserialize, Serialize)]
355pub struct SchemaName {
356    pub catalog_name: String,
357    pub schema_name: String,
358}
359
360impl<'a> From<&'a SchemaName> for SchemaNameKey<'a> {
361    fn from(value: &'a SchemaName) -> Self {
362        Self {
363            catalog: &value.catalog_name,
364            schema: &value.schema_name,
365        }
366    }
367}
368
369#[cfg(test)]
370mod tests {
371    use std::sync::Arc;
372    use std::time::Duration;
373
374    use common_error::ext::ErrorExt;
375    use common_error::status_code::StatusCode;
376    use futures::TryStreamExt;
377    use store_api::mito_engine_options::{
378        TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_TRIGGER_FILE_NUM,
379    };
380
381    use super::*;
382    use crate::kv_backend::KvBackend;
383    use crate::kv_backend::memory::MemoryKvBackend;
384    use crate::rpc::store::PutRequest;
385
386    #[tokio::test]
387    async fn test_schemas() {
388        let manager = SchemaManager::new(Arc::new(MemoryKvBackend::default()));
389        assert!(
390            manager
391                .schemas("catalog")
392                .try_collect::<Vec<_>>()
393                .await
394                .unwrap()
395                .is_empty()
396        );
397
398        let mut expected = BTreeMap::new();
399        for i in 0..=DEFAULT_PAGE_SIZE {
400            let name = format!("schema_{i}");
401            let value = if i == 0 {
402                SchemaNameValue::default()
403            } else {
404                SchemaNameValue {
405                    ttl: Some(Duration::from_secs(i as u64).into()),
406                    extra_options: BTreeMap::from([("foo".to_string(), i.to_string())]),
407                    create_procedure_id: Some(format!("procedure_{i}")),
408                }
409            };
410            manager
411                .create(
412                    SchemaNameKey::new("catalog", &name),
413                    Some(value.clone()),
414                    false,
415                )
416                .await
417                .unwrap();
418            expected.insert(name, value);
419        }
420        manager
421            .create(SchemaNameKey::new("catalog_other", "schema"), None, false)
422            .await
423            .unwrap();
424
425        let schemas = manager
426            .schemas("catalog")
427            .try_collect::<Vec<_>>()
428            .await
429            .unwrap();
430        assert_eq!(schemas, expected.into_iter().collect::<Vec<_>>());
431        assert_eq!(
432            manager
433                .schema_names("catalog")
434                .try_collect::<Vec<_>>()
435                .await
436                .unwrap(),
437            schemas
438                .into_iter()
439                .map(|(name, _)| name)
440                .collect::<Vec<_>>()
441        );
442    }
443
444    #[tokio::test]
445    async fn test_schemas_legacy_and_invalid_values() {
446        let kv_backend = Arc::new(MemoryKvBackend::default());
447        let manager = SchemaManager::new(kv_backend.clone());
448        let key = SchemaNameKey::new("catalog", "schema").to_bytes();
449
450        for (raw, expected) in [
451            (b"null".as_slice(), SchemaNameValue::default()),
452            (
453                br#"{"ttl":"10s"}"#.as_slice(),
454                SchemaNameValue {
455                    ttl: Some(Duration::from_secs(10).into()),
456                    ..Default::default()
457                },
458            ),
459        ] {
460            kv_backend
461                .put(PutRequest::new().with_key(key.clone()).with_value(raw))
462                .await
463                .unwrap();
464            assert_eq!(
465                manager
466                    .schemas("catalog")
467                    .try_collect::<Vec<_>>()
468                    .await
469                    .unwrap(),
470                vec![("schema".to_string(), expected)]
471            );
472        }
473
474        kv_backend
475            .put(PutRequest::new().with_key(key).with_value(b"invalid"))
476            .await
477            .unwrap();
478        assert!(
479            manager
480                .schemas("catalog")
481                .try_collect::<Vec<_>>()
482                .await
483                .is_err()
484        );
485        assert_eq!(
486            manager
487                .schema_names("catalog")
488                .try_collect::<Vec<_>>()
489                .await
490                .unwrap(),
491            vec!["schema".to_string()]
492        );
493    }
494
495    #[test]
496    fn test_display_schema_value() {
497        let schema_value = SchemaNameValue {
498            ttl: None,
499            ..Default::default()
500        };
501        assert_eq!("", schema_value.to_string());
502
503        let schema_value = SchemaNameValue {
504            ttl: Some(Duration::from_secs(9).into()),
505            ..Default::default()
506        };
507        assert_eq!("'ttl'='9s'\n", schema_value.to_string());
508
509        let schema_value = SchemaNameValue {
510            ttl: Some(Duration::from_secs(0).into()),
511            ..Default::default()
512        };
513        assert_eq!("'ttl'='forever'\n", schema_value.to_string());
514    }
515
516    #[test]
517    fn test_serialization() {
518        let key = SchemaNameKey::new("my-catalog", "my-schema");
519        assert_eq!(key.to_string(), "__schema_name/my-catalog/my-schema");
520
521        let parsed = SchemaNameKey::from_bytes(b"__schema_name/my-catalog/my-schema").unwrap();
522
523        assert_eq!(key, parsed);
524
525        let value = SchemaNameValue {
526            ttl: Some(Duration::from_secs(10).into()),
527            ..Default::default()
528        };
529        let mut opts: HashMap<String, String> = HashMap::new();
530        opts.insert("ttl".to_string(), "10s".to_string());
531        let from_value = SchemaNameValue::try_from(&opts).unwrap();
532        assert_eq!(value, from_value);
533
534        let parsed = SchemaNameValue::try_from_raw_value(
535            serde_json::json!({"ttl": "10s"}).to_string().as_bytes(),
536        )
537        .unwrap();
538        assert_eq!(Some(value), parsed);
539
540        let forever = SchemaNameValue {
541            ttl: Some(Default::default()),
542            ..Default::default()
543        };
544        let parsed = SchemaNameValue::try_from_raw_value(
545            serde_json::json!({"ttl": "forever"}).to_string().as_bytes(),
546        )
547        .unwrap();
548        assert_eq!(Some(forever), parsed);
549
550        let instant_err = SchemaNameValue::try_from_raw_value(
551            serde_json::json!({"ttl": "instant"}).to_string().as_bytes(),
552        );
553        assert!(instant_err.is_err());
554
555        let none = SchemaNameValue::try_from_raw_value("null".as_bytes()).unwrap();
556        assert!(none.is_none());
557
558        let err_empty = SchemaNameValue::try_from_raw_value("".as_bytes());
559        assert!(err_empty.is_err());
560    }
561
562    #[test]
563    fn test_schema_value_normalizes_twcs_trigger_aliases() {
564        for options in [
565            HashMap::from([(
566                TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM.to_string(),
567                "4".to_string(),
568            )]),
569            HashMap::from([
570                (TWCS_TRIGGER_FILE_NUM.to_string(), "4".to_string()),
571                (
572                    TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM.to_string(),
573                    "4".to_string(),
574                ),
575            ]),
576        ] {
577            let value = SchemaNameValue::try_from(&options).unwrap();
578            assert_eq!(
579                BTreeMap::from([(TWCS_TRIGGER_FILE_NUM.to_string(), "4".to_string())]),
580                value.extra_options
581            );
582        }
583    }
584
585    #[test]
586    fn test_schema_value_rejects_conflicting_twcs_trigger_aliases() {
587        let options = HashMap::from([
588            (TWCS_TRIGGER_FILE_NUM.to_string(), "4".to_string()),
589            (
590                TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM.to_string(),
591                "8".to_string(),
592            ),
593        ]);
594
595        let error = SchemaNameValue::try_from(&options).unwrap_err();
596
597        assert_eq!(StatusCode::InvalidArguments, error.status_code());
598        assert_eq!(
599            "Conflicting schema options: compaction.twcs.trigger_file_num=4 and compaction.twcs.active_window.trigger_file_num=8",
600            error.to_string()
601        );
602    }
603
604    #[test]
605    fn test_extra_options_compatibility() {
606        // Test with extra_options only
607        let mut opts: HashMap<String, String> = HashMap::new();
608        opts.insert("foo".to_string(), "bar".to_string());
609        opts.insert("baz".to_string(), "qux".to_string());
610        let value = SchemaNameValue::try_from(&opts).unwrap();
611        assert_eq!(value.ttl, None);
612        assert_eq!(value.extra_options.get("foo"), Some(&"bar".to_string()));
613        assert_eq!(value.extra_options.get("baz"), Some(&"qux".to_string()));
614
615        // Test round-trip conversion
616        let opts_back: HashMap<String, String> = value.clone().into();
617        assert_eq!(opts_back.get("foo"), Some(&"bar".to_string()));
618        assert_eq!(opts_back.get("baz"), Some(&"qux".to_string()));
619        assert!(!opts_back.contains_key("ttl"));
620
621        // Test with both ttl and extra_options
622        let mut opts: HashMap<String, String> = HashMap::new();
623        opts.insert("ttl".to_string(), "5m".to_string());
624        opts.insert("opt1".to_string(), "val1".to_string());
625        let value = SchemaNameValue::try_from(&opts).unwrap();
626        assert_eq!(value.ttl, Some(Duration::from_secs(300).into()));
627        assert_eq!(value.extra_options.get("opt1"), Some(&"val1".to_string()));
628
629        // Test serialization/deserialization compatibility
630        let json = serde_json::to_string(&value).unwrap();
631        let deserialized: SchemaNameValue = serde_json::from_str(&json).unwrap();
632        assert_eq!(value, deserialized);
633
634        // Test display includes extra_options
635        let mut value = SchemaNameValue::default();
636        value
637            .extra_options
638            .insert("foo".to_string(), "bar".to_string());
639        let display = value.to_string();
640        assert!(display.contains("'foo'='bar'"));
641    }
642
643    #[test]
644    fn test_backward_compatibility_with_old_format() {
645        // Simulate old format: only ttl, no extra_options
646        let json = r#"{"ttl":"10s"}"#;
647        let parsed = SchemaNameValue::try_from_raw_value(json.as_bytes()).unwrap();
648        assert_eq!(
649            parsed,
650            Some(SchemaNameValue {
651                ttl: Some(Duration::from_secs(10).into()),
652                extra_options: BTreeMap::new(),
653                create_procedure_id: None,
654            })
655        );
656
657        // Simulate old format: null value
658        let json = r#"null"#;
659        let parsed = SchemaNameValue::try_from_raw_value(json.as_bytes()).unwrap();
660        assert!(parsed.is_none());
661    }
662
663    #[test]
664    fn test_forward_compatibility_with_new_options() {
665        // Simulate new format: ttl + extra_options
666        let json = r#"{"ttl":"15s","extra_options":{"foo":"bar","baz":"qux"}}"#;
667        let parsed = SchemaNameValue::try_from_raw_value(json.as_bytes()).unwrap();
668        let mut expected_options = BTreeMap::new();
669        expected_options.insert("foo".to_string(), "bar".to_string());
670        expected_options.insert("baz".to_string(), "qux".to_string());
671        assert_eq!(
672            parsed,
673            Some(SchemaNameValue {
674                ttl: Some(Duration::from_secs(15).into()),
675                extra_options: expected_options,
676                create_procedure_id: None,
677            })
678        );
679    }
680
681    #[test]
682    fn test_create_procedure_id_serialization() {
683        let value = SchemaNameValue {
684            create_procedure_id: Some("4ee0ba94-11f0-4d4d-9468-5ebf732e3ab2".to_string()),
685            ..Default::default()
686        };
687        let raw = value.try_as_raw_value().unwrap();
688        assert_eq!(
689            SchemaNameValue::try_from_raw_value(&raw).unwrap(),
690            Some(value)
691        );
692
693        let raw = SchemaNameValue::default().try_as_raw_value().unwrap();
694        assert!(
695            !String::from_utf8(raw)
696                .unwrap()
697                .contains("create_procedure_id")
698        );
699    }
700
701    #[tokio::test]
702    async fn test_key_exist() {
703        let manager = SchemaManager::new(Arc::new(MemoryKvBackend::default()));
704        let schema_key = SchemaNameKey::new("my-catalog", "my-schema");
705        manager.create(schema_key, None, false).await.unwrap();
706
707        assert!(manager.exists(schema_key).await.unwrap());
708
709        let wrong_schema_key = SchemaNameKey::new("my-catalog", "my-wrong");
710
711        assert!(!manager.exists(wrong_schema_key).await.unwrap());
712    }
713
714    #[tokio::test]
715    async fn test_update_schema_value() {
716        let manager = SchemaManager::new(Arc::new(MemoryKvBackend::default()));
717        let schema_key = SchemaNameKey::new("my-catalog", "my-schema");
718        manager.create(schema_key, None, false).await.unwrap();
719
720        let current_schema_value = manager.get(schema_key).await.unwrap().unwrap();
721        let new_schema_value = SchemaNameValue {
722            ttl: Some(Duration::from_secs(10).into()),
723            ..Default::default()
724        };
725        manager
726            .update(schema_key, &current_schema_value, &new_schema_value)
727            .await
728            .unwrap();
729
730        // Update with the same value, should be ok
731        manager
732            .update(schema_key, &current_schema_value, &new_schema_value)
733            .await
734            .unwrap();
735
736        let new_schema_value = SchemaNameValue {
737            ttl: Some(Duration::from_secs(40).into()),
738            ..Default::default()
739        };
740        let incorrect_schema_value = SchemaNameValue {
741            ttl: Some(Duration::from_secs(20).into()),
742            ..Default::default()
743        }
744        .try_as_raw_value()
745        .unwrap();
746        let incorrect_schema_value =
747            DeserializedValueWithBytes::from_inner_slice(&incorrect_schema_value).unwrap();
748
749        manager
750            .update(schema_key, &incorrect_schema_value, &new_schema_value)
751            .await
752            .unwrap_err();
753
754        let current_schema_value = manager.get(schema_key).await.unwrap().unwrap();
755        let new_schema_value = SchemaNameValue {
756            ttl: None,
757            ..Default::default()
758        };
759        manager
760            .update(schema_key, &current_schema_value, &new_schema_value)
761            .await
762            .unwrap();
763
764        let current_schema_value = manager.get(schema_key).await.unwrap().unwrap();
765        assert_eq!(new_schema_value, *current_schema_value);
766    }
767}