1use 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#[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 #[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
188pub 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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, ¤t_schema_value, &new_schema_value)
727 .await
728 .unwrap();
729
730 manager
732 .update(schema_key, ¤t_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, ¤t_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}