Skip to main content

fory_core/resolver/
meta_resolver.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use crate::buffer::{Reader, Writer};
19use crate::config::Config;
20use crate::error::Error;
21use crate::meta::TypeMeta;
22use crate::resolver::type_resolver::NO_USER_TYPE_ID;
23use crate::resolver::{TypeInfo, TypeResolver};
24use std::collections::HashMap;
25use std::rc::Rc;
26
27/// Streaming meta writer that writes TypeMeta inline during serialization.
28/// Uses the streaming protocol:
29/// - (index << 1) | 0 for new type definition (followed by TypeMeta bytes)
30/// - (index << 1) | 1 for reference to previously written type
31#[derive(Default)]
32pub struct MetaWriterResolver {
33    // Provider and target indexes share one Rc<TypeInfo>; pointer identity keeps
34    // their streaming metadata references in one sequence without probing maps.
35    type_info_index_map: HashMap<*const TypeInfo, usize>,
36    type_index_index_map: Vec<usize>,
37    next_index: usize,
38}
39
40const MIN_REMOTE_TYPE_META_VERSIONS: u64 = 8192;
41const MAX_REMOTE_TYPE_META_KEYS: usize = 8192;
42const NO_WRITTEN_TYPE_INDEX: usize = usize::MAX;
43
44#[allow(dead_code)]
45impl MetaWriterResolver {
46    /// Write type meta inline using streaming protocol.
47    /// Returns the index assigned to this type.
48    #[inline(always)]
49    pub fn write_type_meta(
50        &mut self,
51        writer: &mut Writer,
52        provider_type_id: std::any::TypeId,
53        type_resolver: &TypeResolver,
54    ) -> Result<(), Error> {
55        let type_info = type_resolver.get_provider_type_info(&provider_type_id)?;
56        self.write_resolved_type_meta(writer, &type_info)
57    }
58
59    #[inline(always)]
60    pub(crate) fn write_resolved_type_meta(
61        &mut self,
62        writer: &mut Writer,
63        type_info: &Rc<TypeInfo>,
64    ) -> Result<(), Error> {
65        let identity = Rc::as_ptr(type_info);
66        match self.type_info_index_map.get(&identity) {
67            Some(&index) => {
68                // Reference to previously written type: (index << 1) | 1, LSB=1
69                writer.write_var_u32(((index as u32) << 1) | 1);
70            }
71            None => {
72                // New type: index << 1, LSB=0, followed by TypeMeta bytes inline
73                let index = self.next_index;
74                self.next_index += 1;
75                writer.write_var_u32((index as u32) << 1);
76                self.type_info_index_map.insert(identity, index);
77                let type_def = type_info.get_type_def();
78                writer.write_bytes(&type_def);
79            }
80        }
81        Ok(())
82    }
83
84    /// Write type meta by generated struct type index, avoiding Rust TypeId hash lookup.
85    #[inline(always)]
86    pub fn write_type_meta_fast(
87        &mut self,
88        writer: &mut Writer,
89        type_id: std::any::TypeId,
90        type_index: u32,
91        type_resolver: &TypeResolver,
92    ) -> Result<(), Error> {
93        let type_index = type_index as usize;
94        if let Some(&index) = self.type_index_index_map.get(type_index) {
95            if index != NO_WRITTEN_TYPE_INDEX {
96                writer.write_var_u32(((index as u32) << 1) | 1);
97                return Ok(());
98            }
99        }
100
101        let index = self.next_index;
102        self.next_index += 1;
103        writer.write_var_u32((index as u32) << 1);
104        if type_index >= self.type_index_index_map.len() {
105            self.type_index_index_map
106                .resize(type_index + 1, NO_WRITTEN_TYPE_INDEX);
107        }
108        self.type_index_index_map[type_index] = index;
109        let type_meta = type_resolver.get_type_meta_by_index_ref(&type_id, type_index as u32)?;
110        writer.write_bytes(type_meta.get_bytes());
111        Ok(())
112    }
113
114    #[inline(always)]
115    pub fn reset(&mut self) {
116        self.type_info_index_map.clear();
117        self.type_index_index_map.clear();
118        self.next_index = 0;
119    }
120}
121
122/// Streaming meta reader that reads TypeMeta inline during deserialization.
123/// Uses the streaming protocol:
124/// - (index << 1) | 0 for new type definition (followed by TypeMeta bytes)
125/// - (index << 1) | 1 for reference to previously read type
126#[derive(Default)]
127pub struct MetaReaderResolver {
128    pub reading_type_infos: Vec<Rc<TypeInfo>>,
129    parsed_type_infos: HashMap<i64, Rc<TypeInfo>>,
130    remote_schema_versions_by_type: HashMap<String, usize>,
131    total_accepted_schema_versions: u64,
132    cached_meta_header: i64,
133    cached_type_info: Option<Rc<TypeInfo>>,
134}
135
136impl MetaReaderResolver {
137    #[inline(always)]
138    pub fn get(&self, index: usize) -> Option<&Rc<TypeInfo>> {
139        self.reading_type_infos.get(index)
140    }
141
142    /// Read type meta inline using streaming protocol.
143    /// Returns the TypeInfo for this type.
144    #[inline(always)]
145    pub fn read_type_meta(
146        &mut self,
147        reader: &mut Reader,
148        type_resolver: &TypeResolver,
149        config: &Config,
150    ) -> Result<Rc<TypeInfo>, Error> {
151        let index_marker = reader.read_var_u32()?;
152        let is_ref = (index_marker & 1) == 1;
153        let index = (index_marker >> 1) as usize;
154
155        if is_ref {
156            // Reference to previously read type
157            self.reading_type_infos.get(index).cloned().ok_or_else(|| {
158                Error::type_error(format!("TypeInfo not found for type index: {}", index))
159            })
160        } else {
161            // New type - read TypeMeta inline
162            let meta_header = reader.read_i64()?;
163            if let Some(type_info) = self
164                .cached_type_info
165                .as_ref()
166                .filter(|_| self.cached_meta_header == meta_header)
167            {
168                // Header-cache hits intentionally skip without rehashing. Entries reach this cache
169                // only after a successful TypeMeta parse and 52-bit metadata-hash validation. Do
170                // not add body/hash/schema-limit/exact-local checks here; the miss path owns them
171                // before publish.
172                self.reading_type_infos.push(type_info.clone());
173                TypeMeta::skip_bytes_for_validated_header(reader, meta_header)?;
174                return Ok(type_info.clone());
175            }
176            if let Some(type_info) = self.parsed_type_infos.get(&meta_header) {
177                // Header-cache hits intentionally skip without rehashing. Entries reach this cache
178                // only after a successful TypeMeta parse and 52-bit metadata-hash validation. Do
179                // not add body/hash/schema-limit/exact-local checks here; the miss path owns them
180                // before publish.
181                self.cached_meta_header = meta_header;
182                self.cached_type_info = Some(type_info.clone());
183                self.reading_type_infos.push(type_info.clone());
184                TypeMeta::skip_bytes_for_validated_header(reader, meta_header)?;
185                Ok(type_info.clone())
186            } else {
187                let type_def_start = reader.get_cursor() - std::mem::size_of::<i64>();
188                self.read_remote_type_meta(
189                    reader,
190                    type_resolver,
191                    config,
192                    meta_header,
193                    type_def_start,
194                )
195            }
196        }
197    }
198
199    #[cold]
200    #[inline(never)]
201    fn read_remote_type_meta(
202        &mut self,
203        reader: &mut Reader,
204        type_resolver: &TypeResolver,
205        config: &Config,
206        meta_header: i64,
207        type_def_start: usize,
208    ) -> Result<Rc<TypeInfo>, Error> {
209        let type_meta = Rc::new(TypeMeta::from_bytes_with_header(
210            reader,
211            type_resolver,
212            meta_header,
213            config.max_type_fields(),
214            config.max_type_meta_bytes(),
215        )?);
216        let remote_type_def = reader.sub_slice(type_def_start, reader.get_cursor())?;
217
218        let namespace = type_meta.get_namespace();
219        let type_name = type_meta.get_type_name();
220        let register_by_name = !namespace.original.is_empty() || !type_name.original.is_empty();
221        let mut remote_schema_key = None;
222        let type_info = if register_by_name {
223            if let Some(local_type_info) =
224                type_resolver.get_type_info_by_name(&namespace.original, &type_name.original)
225            {
226                if local_type_info.get_type_meta_ref().get_bytes() == remote_type_def {
227                    local_type_info
228                } else {
229                    remote_schema_key =
230                        Some(self.check_remote_type_meta_limit(&type_meta, config)?);
231                    Rc::new(TypeInfo::from_remote_meta(
232                        type_meta.clone(),
233                        Some(local_type_info.get_harness()),
234                        Some(local_type_info.get_type_id() as u32),
235                        Some(local_type_info.get_user_type_id()),
236                    ))
237                }
238            } else {
239                remote_schema_key = Some(self.check_remote_type_meta_limit(&type_meta, config)?);
240                Rc::new(TypeInfo::from_remote_meta(
241                    type_meta.clone(),
242                    None,
243                    None,
244                    None,
245                ))
246            }
247        } else {
248            let type_id = type_meta.get_type_id();
249            let user_type_id = type_meta.get_user_type_id();
250            let local_type_info = if user_type_id != NO_USER_TYPE_ID {
251                type_resolver.get_user_type_info_by_id(user_type_id)
252            } else {
253                type_resolver.get_type_info_by_id(type_id)
254            };
255            if let Some(local_type_info) = local_type_info {
256                if local_type_info.get_type_meta_ref().get_bytes() == remote_type_def {
257                    local_type_info
258                } else {
259                    remote_schema_key =
260                        Some(self.check_remote_type_meta_limit(&type_meta, config)?);
261                    Rc::new(TypeInfo::from_remote_meta(
262                        type_meta.clone(),
263                        Some(local_type_info.get_harness()),
264                        Some(local_type_info.get_type_id() as u32),
265                        Some(local_type_info.get_user_type_id()),
266                    ))
267                }
268            } else {
269                remote_schema_key = Some(self.check_remote_type_meta_limit(&type_meta, config)?);
270                Rc::new(TypeInfo::from_remote_meta(
271                    type_meta.clone(),
272                    None,
273                    None,
274                    None,
275                ))
276            }
277        };
278
279        self.parsed_type_infos
280            .insert(meta_header, type_info.clone());
281        self.cached_meta_header = meta_header;
282        self.cached_type_info = Some(type_info.clone());
283        self.reading_type_infos.push(type_info.clone());
284        if let Some(remote_schema_key) = remote_schema_key {
285            self.record_remote_type_meta(remote_schema_key);
286        }
287        Ok(type_info)
288    }
289
290    #[cold]
291    #[inline(never)]
292    fn check_remote_type_meta_limit(
293        &self,
294        type_meta: &TypeMeta,
295        config: &Config,
296    ) -> Result<String, Error> {
297        let namespace = type_meta.get_namespace();
298        let type_name = type_meta.get_type_name();
299        let key = if !namespace.original.is_empty() || !type_name.original.is_empty() {
300            format!("n{}\0{}", namespace.original, type_name.original)
301        } else {
302            format!("i{}", type_meta.get_user_type_id())
303        };
304
305        let versions_for_type = self
306            .remote_schema_versions_by_type
307            .get(&key)
308            .copied()
309            .unwrap_or(0);
310        // Reaching the key cap must not disable schema evolution for keys that were already
311        // accepted.
312        if versions_for_type == 0
313            && self.remote_schema_versions_by_type.len() >= MAX_REMOTE_TYPE_META_KEYS
314        {
315            return Err(Error::invalid_data(
316                "remote logical TypeMeta key limit exceeded. The data may be malicious",
317            ));
318        }
319        if versions_for_type >= config.max_schema_versions_per_type() {
320            return Err(Error::invalid_data(format!(
321                "remote schema version limit exceeded for one type. The data may be malicious. If the data is not malicious, please increase max_schema_versions_per_type={}",
322                config.max_schema_versions_per_type()
323            )));
324        }
325
326        let accepted_type_count = (self.remote_schema_versions_by_type.len()
327            + if versions_for_type == 0 { 1 } else { 0 }) as u64;
328        let max_average = config.max_average_schema_versions_per_type() as u64;
329        let reached_average_limit = max_average == 0
330            || self.total_accepted_schema_versions / max_average >= accepted_type_count;
331        if self.total_accepted_schema_versions == u64::MAX
332            || (self.total_accepted_schema_versions >= MIN_REMOTE_TYPE_META_VERSIONS
333                && reached_average_limit)
334        {
335            return Err(Error::invalid_data(format!(
336                "remote schema version limit exceeded globally. The data may be malicious. If the data is not malicious, please increase max_average_schema_versions_per_type={}",
337                config.max_average_schema_versions_per_type()
338            )));
339        }
340
341        Ok(key)
342    }
343
344    fn record_remote_type_meta(&mut self, key: String) {
345        let versions_for_type = self
346            .remote_schema_versions_by_type
347            .get(&key)
348            .copied()
349            .unwrap_or(0);
350        self.remote_schema_versions_by_type
351            .insert(key, versions_for_type + 1);
352        // The cold miss check rejects u64::MAX before its caller publishes the TypeInfo and reaches
353        // this mutation.
354        self.total_accepted_schema_versions += 1;
355    }
356
357    #[inline(always)]
358    pub fn reset(&mut self) {
359        self.reading_type_infos.clear();
360    }
361}
362
363#[cfg(test)]
364mod tests {
365    use super::*;
366    use crate::config::Config;
367    use crate::context::{ReadContext, WriteContext};
368    use crate::meta::{
369        FieldInfo, FieldType, MetaString, NAMESPACE_ENCODER, NAMESPACE_ENCODINGS,
370        TYPE_NAME_ENCODER, TYPE_NAME_ENCODINGS,
371    };
372    use crate::serializer::Serializer;
373    use crate::TypeId;
374
375    struct LocalExt;
376
377    impl Serializer for LocalExt {
378        type Target = Self;
379
380        fn write_data(_value: &Self, _context: &mut WriteContext) -> Result<(), Error> {
381            Ok(())
382        }
383
384        fn read_data(_context: &mut ReadContext) -> Result<Self, Error> {
385            Ok(LocalExt)
386        }
387    }
388
389    fn read_type_def(
390        resolver: &mut MetaReaderResolver,
391        config: &Config,
392        type_def: &[u8],
393    ) -> Result<Rc<TypeInfo>, Error> {
394        let type_resolver = TypeResolver::default();
395        read_type_def_with_type_resolver(resolver, config, &type_resolver, type_def)
396    }
397
398    fn read_type_def_with_type_resolver(
399        resolver: &mut MetaReaderResolver,
400        config: &Config,
401        type_resolver: &TypeResolver,
402        type_def: &[u8],
403    ) -> Result<Rc<TypeInfo>, Error> {
404        let mut bytes = vec![];
405        let mut writer = Writer::from_buffer(&mut bytes);
406        writer.write_var_u32(0);
407        writer.write_bytes(type_def);
408        let mut reader = Reader::new(&bytes);
409        resolver.read_type_meta(&mut reader, type_resolver, config)
410    }
411
412    fn remote_struct_meta(user_type_id: u32, field_name: &str) -> TypeMeta {
413        TypeMeta::new(
414            TypeId::STRUCT as u32,
415            user_type_id,
416            MetaString::get_empty().clone(),
417            MetaString::get_empty().clone(),
418            false,
419            vec![FieldInfo::new(
420                field_name,
421                FieldType::new(crate::type_id::INT32, false, vec![]),
422            )],
423        )
424        .unwrap()
425    }
426
427    fn fill_remote_schema_keys(resolver: &mut MetaReaderResolver, count: usize, versions: usize) {
428        assert!(count <= MAX_REMOTE_TYPE_META_KEYS);
429        for user_type_id in 0..count {
430            resolver
431                .remote_schema_versions_by_type
432                .insert(format!("i{user_type_id}"), versions);
433        }
434        resolver.total_accepted_schema_versions = count as u64 * versions as u64;
435    }
436
437    #[test]
438    fn logical_type_key_cap() {
439        let config = Config::default();
440        let mut resolver = MetaReaderResolver::default();
441        fill_remote_schema_keys(&mut resolver, MAX_REMOTE_TYPE_META_KEYS - 1, 1);
442
443        let last = remote_struct_meta((MAX_REMOTE_TYPE_META_KEYS - 1) as u32, "a");
444        read_type_def(&mut resolver, &config, last.get_bytes()).unwrap();
445        assert_eq!(
446            resolver.remote_schema_versions_by_type.len(),
447            MAX_REMOTE_TYPE_META_KEYS
448        );
449        assert_eq!(
450            resolver.total_accepted_schema_versions,
451            MAX_REMOTE_TYPE_META_KEYS as u64
452        );
453
454        let parsed_count = resolver.parsed_type_infos.len();
455        let reading_count = resolver.reading_type_infos.len();
456        let cached_header = resolver.cached_meta_header;
457        let cached_type_info = resolver.cached_type_info.as_ref().map(Rc::as_ptr);
458        let rejected = remote_struct_meta(MAX_REMOTE_TYPE_META_KEYS as u32, "a");
459        let err = read_type_def(&mut resolver, &config, rejected.get_bytes())
460            .unwrap_err()
461            .to_string();
462
463        assert!(err.contains("logical TypeMeta key limit"));
464        assert_eq!(
465            resolver.remote_schema_versions_by_type.len(),
466            MAX_REMOTE_TYPE_META_KEYS
467        );
468        assert_eq!(
469            resolver.total_accepted_schema_versions,
470            MAX_REMOTE_TYPE_META_KEYS as u64
471        );
472        assert_eq!(resolver.parsed_type_infos.len(), parsed_count);
473        assert_eq!(resolver.reading_type_infos.len(), reading_count);
474        assert_eq!(resolver.cached_meta_header, cached_header);
475        assert_eq!(
476            resolver.cached_type_info.as_ref().map(Rc::as_ptr),
477            cached_type_info
478        );
479    }
480
481    #[test]
482    fn existing_key_keeps_limits() {
483        let mut per_type_resolver = MetaReaderResolver::default();
484        fill_remote_schema_keys(&mut per_type_resolver, MAX_REMOTE_TYPE_META_KEYS, 1);
485        let per_type_config = Config {
486            max_schema_versions_per_type: 1,
487            ..Default::default()
488        };
489        let changed = remote_struct_meta(0, "b");
490        let err = read_type_def(
491            &mut per_type_resolver,
492            &per_type_config,
493            changed.get_bytes(),
494        )
495        .unwrap_err()
496        .to_string();
497        assert!(err.contains("max_schema_versions_per_type"));
498
499        let mut average_resolver = MetaReaderResolver::default();
500        fill_remote_schema_keys(&mut average_resolver, MAX_REMOTE_TYPE_META_KEYS, 3);
501        *average_resolver
502            .remote_schema_versions_by_type
503            .get_mut("i0")
504            .unwrap() = 2;
505        average_resolver.total_accepted_schema_versions -= 1;
506        let average_config = Config {
507            max_schema_versions_per_type: 10,
508            max_average_schema_versions_per_type: 3,
509            ..Default::default()
510        };
511
512        let accepted = remote_struct_meta(0, "b");
513        read_type_def(&mut average_resolver, &average_config, accepted.get_bytes()).unwrap();
514        assert_eq!(average_resolver.total_accepted_schema_versions, 24_576);
515
516        let rejected = remote_struct_meta(0, "c");
517        let err = read_type_def(&mut average_resolver, &average_config, rejected.get_bytes())
518            .unwrap_err()
519            .to_string();
520        assert!(err.contains("max_average_schema_versions_per_type"));
521        assert_eq!(average_resolver.total_accepted_schema_versions, 24_576);
522    }
523
524    #[test]
525    fn schema_total_does_not_wrap() {
526        let config = Config {
527            max_schema_versions_per_type: u32::MAX,
528            max_average_schema_versions_per_type: u32::MAX,
529            ..Default::default()
530        };
531        let mut resolver = MetaReaderResolver::default();
532        fill_remote_schema_keys(&mut resolver, 1, 1);
533        resolver.total_accepted_schema_versions = u64::MAX;
534        let meta = remote_struct_meta(0, "b");
535
536        let err = read_type_def(&mut resolver, &config, meta.get_bytes())
537            .unwrap_err()
538            .to_string();
539
540        assert!(err.contains("remote schema version limit exceeded globally"));
541        assert_eq!(resolver.total_accepted_schema_versions, u64::MAX);
542        assert_eq!(resolver.remote_schema_versions_by_type.get("i0"), Some(&1));
543        assert!(resolver.parsed_type_infos.is_empty());
544        assert!(resolver.cached_type_info.is_none());
545        assert!(resolver.reading_type_infos.is_empty());
546    }
547
548    #[test]
549    fn checked_cache_bypasses_key_cap() {
550        let config = Config::default();
551        let mut resolver = MetaReaderResolver::default();
552        fill_remote_schema_keys(&mut resolver, MAX_REMOTE_TYPE_META_KEYS - 1, 1);
553        let meta = remote_struct_meta((MAX_REMOTE_TYPE_META_KEYS - 1) as u32, "a");
554        let first = read_type_def(&mut resolver, &config, meta.get_bytes()).unwrap();
555
556        resolver.reset();
557        resolver.cached_type_info = None;
558        let strict_config = Config {
559            max_schema_versions_per_type: 1,
560            max_average_schema_versions_per_type: 1,
561            ..Default::default()
562        };
563        let cached = read_type_def(&mut resolver, &strict_config, meta.get_bytes()).unwrap();
564
565        assert!(Rc::ptr_eq(&first, &cached));
566        assert_eq!(resolver.reading_type_infos.len(), 1);
567        assert_eq!(
568            resolver.remote_schema_versions_by_type.len(),
569            MAX_REMOTE_TYPE_META_KEYS
570        );
571        assert_eq!(
572            resolver.total_accepted_schema_versions,
573            MAX_REMOTE_TYPE_META_KEYS as u64
574        );
575    }
576
577    #[test]
578    fn exact_local_bypasses_key_cap() {
579        let mut type_resolver = TypeResolver::default();
580        type_resolver
581            .register_serializer_by_name::<LocalExt>("example.SharedExt")
582            .unwrap();
583        let type_resolver = type_resolver.build_final_type_resolver().unwrap();
584        let local_info = type_resolver
585            .get_type_info_by_name("example", "SharedExt")
586            .unwrap();
587        let exact = local_info.get_type_meta_ref().get_bytes().to_vec();
588
589        let mut resolver = MetaReaderResolver::default();
590        fill_remote_schema_keys(&mut resolver, MAX_REMOTE_TYPE_META_KEYS, 1);
591        let strict_config = Config {
592            max_schema_versions_per_type: 1,
593            max_average_schema_versions_per_type: 1,
594            ..Default::default()
595        };
596        let resolved =
597            read_type_def_with_type_resolver(&mut resolver, &strict_config, &type_resolver, &exact)
598                .unwrap();
599
600        assert!(Rc::ptr_eq(&local_info, &resolved));
601        assert_eq!(
602            resolver.remote_schema_versions_by_type.len(),
603            MAX_REMOTE_TYPE_META_KEYS
604        );
605        assert_eq!(
606            resolver.total_accepted_schema_versions,
607            MAX_REMOTE_TYPE_META_KEYS as u64
608        );
609    }
610
611    #[test]
612    fn type_meta_field_limit_rejects_large_struct() {
613        let meta = TypeMeta::new(
614            TypeId::STRUCT as u32,
615            9001,
616            MetaString::get_empty().clone(),
617            MetaString::get_empty().clone(),
618            false,
619            vec![
620                FieldInfo::new("a", FieldType::new(crate::type_id::INT32, false, vec![])),
621                FieldInfo::new("b", FieldType::new(crate::type_id::INT32, false, vec![])),
622            ],
623        )
624        .unwrap();
625        let config = Config {
626            max_type_fields: 1,
627            ..Default::default()
628        };
629        let err = read_type_def(
630            &mut MetaReaderResolver::default(),
631            &config,
632            meta.get_bytes(),
633        )
634        .unwrap_err()
635        .to_string();
636        assert!(err.contains("max_type_fields"));
637    }
638
639    #[test]
640    fn type_meta_body_limit_rejects_large_metadata() {
641        let meta = TypeMeta::new(
642            TypeId::STRUCT as u32,
643            9001,
644            MetaString::get_empty().clone(),
645            MetaString::get_empty().clone(),
646            false,
647            vec![FieldInfo::new(
648                "a",
649                FieldType::new(crate::type_id::INT32, false, vec![]),
650            )],
651        )
652        .unwrap();
653        let config = Config {
654            max_type_meta_bytes: 1,
655            ..Default::default()
656        };
657        let err = read_type_def(
658            &mut MetaReaderResolver::default(),
659            &config,
660            meta.get_bytes(),
661        )
662        .unwrap_err()
663        .to_string();
664        assert!(err.contains("max_type_meta_bytes"));
665    }
666
667    #[test]
668    fn schema_limit_tracks_unknown_struct_types_separately() {
669        fn type_def(user_type_id: u32, field_name: &str) -> Vec<u8> {
670            TypeMeta::new(
671                TypeId::STRUCT as u32,
672                user_type_id,
673                MetaString::get_empty().clone(),
674                MetaString::get_empty().clone(),
675                false,
676                vec![FieldInfo::new(
677                    field_name,
678                    FieldType::new(crate::type_id::INT32, false, vec![]),
679                )],
680            )
681            .unwrap()
682            .get_bytes()
683            .to_vec()
684        }
685
686        let config = Config {
687            max_schema_versions_per_type: 1,
688            ..Default::default()
689        };
690
691        let mut resolver = MetaReaderResolver::default();
692        read_type_def(&mut resolver, &config, &type_def(9001, "a")).unwrap();
693        read_type_def(&mut resolver, &config, &type_def(9002, "a")).unwrap();
694
695        let err = read_type_def(&mut resolver, &config, &type_def(9001, "b"))
696            .unwrap_err()
697            .to_string();
698        assert!(err.contains("max_schema_versions_per_type"));
699    }
700
701    #[test]
702    fn schema_limit_rejects_extra_versions_for_type() {
703        let meta = TypeMeta::new(
704            TypeId::STRUCT as u32,
705            9001,
706            MetaString::get_empty().clone(),
707            MetaString::get_empty().clone(),
708            false,
709            vec![FieldInfo::new(
710                "a",
711                FieldType::new(crate::type_id::INT32, false, vec![]),
712            )],
713        )
714        .unwrap();
715        let type_def = meta.get_bytes().to_vec();
716
717        let config = Config {
718            max_schema_versions_per_type: 1,
719            ..Default::default()
720        };
721        let mut resolver = MetaReaderResolver::default();
722        let mut bytes = vec![];
723        let mut writer = Writer::from_buffer(&mut bytes);
724        writer.write_var_u32(0);
725        writer.write_bytes(&type_def);
726        let mut reader = Reader::new(&bytes);
727        resolver
728            .read_type_meta(&mut reader, &TypeResolver::default(), &config)
729            .unwrap();
730
731        let changed = TypeMeta::new(
732            TypeId::STRUCT as u32,
733            9001,
734            MetaString::get_empty().clone(),
735            MetaString::get_empty().clone(),
736            false,
737            vec![FieldInfo::new(
738                "b",
739                FieldType::new(crate::type_id::INT32, false, vec![]),
740            )],
741        )
742        .unwrap();
743        let mut bytes = vec![];
744        let mut writer = Writer::from_buffer(&mut bytes);
745        writer.write_var_u32(0);
746        writer.write_bytes(changed.get_bytes());
747        let mut reader = Reader::new(&bytes);
748        let err = resolver
749            .read_type_meta(&mut reader, &TypeResolver::default(), &config)
750            .unwrap_err()
751            .to_string();
752        assert!(err.contains("max_schema_versions_per_type"));
753    }
754
755    #[test]
756    fn schema_limit_check_is_not_recorded() {
757        let config = Config {
758            max_schema_versions_per_type: 1,
759            ..Default::default()
760        };
761        let mut resolver = MetaReaderResolver::default();
762        let checked = TypeMeta::new(
763            TypeId::STRUCT as u32,
764            9001,
765            MetaString::get_empty().clone(),
766            MetaString::get_empty().clone(),
767            false,
768            vec![FieldInfo::new(
769                "a",
770                FieldType::new(crate::type_id::INT32, false, vec![]),
771            )],
772        )
773        .unwrap();
774        let accepted = TypeMeta::new(
775            TypeId::STRUCT as u32,
776            9001,
777            MetaString::get_empty().clone(),
778            MetaString::get_empty().clone(),
779            false,
780            vec![FieldInfo::new(
781                "b",
782                FieldType::new(crate::type_id::INT32, false, vec![]),
783            )],
784        )
785        .unwrap();
786
787        resolver
788            .check_remote_type_meta_limit(&checked, &config)
789            .unwrap();
790
791        let mut bytes = vec![];
792        let mut writer = Writer::from_buffer(&mut bytes);
793        writer.write_var_u32(0);
794        writer.write_bytes(accepted.get_bytes());
795        let mut reader = Reader::new(&bytes);
796        resolver
797            .read_type_meta(&mut reader, &TypeResolver::default(), &config)
798            .unwrap();
799    }
800
801    #[test]
802    fn non_struct_type_meta_uses_limit() {
803        let config = Config {
804            max_schema_versions_per_type: 1,
805            ..Default::default()
806        };
807        let mut resolver = MetaReaderResolver::default();
808        let namespace = NAMESPACE_ENCODER
809            .encode_with_encodings("example", NAMESPACE_ENCODINGS)
810            .unwrap();
811        let type_name = TYPE_NAME_ENCODER
812            .encode_with_encodings("RemoteEnum", TYPE_NAME_ENCODINGS)
813            .unwrap();
814        let first = TypeMeta::new(
815            TypeId::NAMED_ENUM as u32,
816            NO_USER_TYPE_ID,
817            namespace.clone(),
818            type_name.clone(),
819            true,
820            vec![],
821        )
822        .unwrap();
823        let second = TypeMeta::new(
824            TypeId::NAMED_EXT as u32,
825            NO_USER_TYPE_ID,
826            namespace,
827            type_name,
828            true,
829            vec![],
830        )
831        .unwrap();
832
833        let key = resolver
834            .check_remote_type_meta_limit(&first, &config)
835            .unwrap();
836        resolver.record_remote_type_meta(key);
837
838        let err = resolver
839            .check_remote_type_meta_limit(&second, &config)
840            .unwrap_err()
841            .to_string();
842        assert!(err.contains("max_schema_versions_per_type"));
843    }
844
845    #[test]
846    fn exact_local_non_struct_type_meta_bypasses_limit() {
847        let config = Config {
848            max_schema_versions_per_type: 1,
849            ..Default::default()
850        };
851        let mut type_resolver = TypeResolver::default();
852        type_resolver
853            .register_serializer_by_name::<LocalExt>("example.SharedExt")
854            .unwrap();
855        let type_resolver = type_resolver.build_final_type_resolver().unwrap();
856        let local_info = type_resolver
857            .get_type_info_by_name("example", "SharedExt")
858            .unwrap();
859        let exact = local_info.get_type_meta_ref().get_bytes().to_vec();
860
861        let mut resolver = MetaReaderResolver::default();
862        read_type_def_with_type_resolver(&mut resolver, &config, &type_resolver, &exact).unwrap();
863
864        let namespace = NAMESPACE_ENCODER
865            .encode_with_encodings("example", NAMESPACE_ENCODINGS)
866            .unwrap();
867        let type_name = TYPE_NAME_ENCODER
868            .encode_with_encodings("SharedExt", TYPE_NAME_ENCODINGS)
869            .unwrap();
870        let second = TypeMeta::new(
871            TypeId::NAMED_ENUM as u32,
872            NO_USER_TYPE_ID,
873            namespace,
874            type_name,
875            true,
876            vec![],
877        )
878        .unwrap();
879        resolver
880            .check_remote_type_meta_limit(&second, &config)
881            .unwrap();
882    }
883}