1use crate::TraceEvent;
8use crate::codec::{self, PoolEntry, StackPoolEntry, WireTypeId};
9use crate::schema::{SchemaEntry, SchemaRegistry};
10use crate::types::{
11 CountingWriter, EncodeState, EventEncoder, InternedStackFrames, InternedString,
12};
13use std::any::TypeId;
14use std::collections::{HashMap, HashSet};
15use std::hash::{BuildHasherDefault, Hasher};
16use std::io::{self, Write};
17use std::sync::Arc;
18
19#[doc(hidden)]
24#[derive(Default)]
25pub struct FxHasher(u64);
26
27impl FxHasher {
28 #[inline]
29 fn hash_word(&mut self, word: u64) {
30 self.0 = (self.0.rotate_left(5) ^ word).wrapping_mul(0x517cc1b727220a95);
31 }
32}
33
34impl Hasher for FxHasher {
35 #[inline]
36 fn write(&mut self, mut bytes: &[u8]) {
37 while bytes.len() >= 8 {
38 self.hash_word(u64::from_ne_bytes(bytes[..8].try_into().unwrap()));
39 bytes = &bytes[8..];
40 }
41 if bytes.len() >= 4 {
42 self.hash_word(u32::from_ne_bytes(bytes[..4].try_into().unwrap()) as u64);
43 bytes = &bytes[4..];
44 }
45 for &b in bytes {
46 self.hash_word(b as u64);
47 }
48 }
49
50 #[inline]
51 fn write_u8(&mut self, i: u8) {
52 self.hash_word(i as u64);
53 }
54
55 #[inline]
56 fn write_u16(&mut self, i: u16) {
57 self.hash_word(i as u64);
58 }
59
60 #[inline]
61 fn write_u32(&mut self, i: u32) {
62 self.hash_word(i as u64);
63 }
64
65 #[inline]
66 fn write_u64(&mut self, i: u64) {
67 self.hash_word(i);
68 }
69
70 #[inline]
71 fn write_usize(&mut self, i: usize) {
72 self.hash_word(i as u64);
73 }
74
75 #[inline]
76 fn write_u128(&mut self, i: u128) {
77 self.hash_word(i as u64);
78 self.hash_word((i >> 64) as u64);
79 }
80
81 #[inline]
82 fn finish(&self) -> u64 {
83 self.0
84 }
85}
86
87#[doc(hidden)]
88pub type FxBuildHasher = BuildHasherDefault<FxHasher>;
89#[doc(hidden)]
90pub type FxHashMap<K, V> = HashMap<K, V, FxBuildHasher>;
91#[doc(hidden)]
92pub type FxHashSet<T> = HashSet<T, FxBuildHasher>;
93
94#[derive(Clone, Debug)]
104pub struct Schema {
105 pub(crate) entry: Arc<SchemaEntry>,
106 name_key: Arc<str>,
109}
110
111impl Schema {
112 pub fn new(name: &str, fields: Vec<crate::schema::FieldDef>) -> Self {
117 let name_key: Arc<str> = Arc::from(name);
118 Self {
119 entry: Arc::new(SchemaEntry {
120 name: name.to_string(),
121 fields,
122 annotations: Vec::new(),
123 }),
124 name_key,
125 }
126 }
127
128 pub fn from_entry(entry: SchemaEntry) -> Self {
130 let name_key: Arc<str> = Arc::from(entry.name.as_str());
131 Self {
132 entry: Arc::new(entry),
133 name_key,
134 }
135 }
136
137 pub fn name(&self) -> &str {
139 &self.entry.name
140 }
141
142 pub fn fields(&self) -> &[crate::schema::FieldDef] {
144 &self.entry.fields
145 }
146}
147
148#[derive(Clone, PartialEq, Eq, Hash)]
151enum SchemaKey {
152 Name(Arc<str>),
153 RustType(TypeId),
154}
155
156const DYNAMIC_SCHEMA_CACHE_LIMIT: usize = 1024;
167
168pub struct Encoder<W: Write = Vec<u8>> {
169 state: EncodeState<W>,
170 registry: SchemaRegistry,
171 string_pool: FxHashMap<String, u32>,
172 next_pool_id: u32,
173 stack_pool: FxHashMap<Box<[u64]>, u32>,
174 next_stack_pool_id: u32,
175 schema_ids: FxHashMap<SchemaKey, WireTypeId>,
176 dynamic_schema_cache: FxHashMap<usize, (Arc<SchemaEntry>, WireTypeId)>,
188 slot_cache: Vec<u32>,
191 registered_ids: [u64; (crate::STATIC_WIRE_ID_LIMIT as usize) / 64],
195}
196
197impl Default for Encoder<Vec<u8>> {
198 fn default() -> Self {
199 Self::new()
200 }
201}
202
203impl Encoder<Vec<u8>> {
204 pub fn new() -> Self {
205 let mut buf = Vec::new();
206 codec::encode_header(&mut buf).expect("Vec::write_all cannot fail");
207 Self {
208 state: EncodeState::new(buf),
209 registry: SchemaRegistry::new(),
210 string_pool: FxHashMap::default(),
211 next_pool_id: 0,
212 stack_pool: FxHashMap::default(),
213 next_stack_pool_id: 0,
214 schema_ids: FxHashMap::default(),
215 dynamic_schema_cache: FxHashMap::default(),
216 slot_cache: Vec::new(),
217 registered_ids: [0; (crate::STATIC_WIRE_ID_LIMIT as usize) / 64],
218 }
219 }
220
221 pub fn finish(self) -> Vec<u8> {
223 self.state.writer.into_inner()
224 }
225}
226
227impl<W: Write> Encoder<W> {
228 pub fn new_to(mut writer: W) -> io::Result<Self> {
231 codec::encode_header(&mut writer)?;
232 Ok(Self {
233 state: EncodeState::new(writer),
234 registry: SchemaRegistry::new(),
235 string_pool: FxHashMap::default(),
236 next_pool_id: 0,
237 stack_pool: FxHashMap::default(),
238 next_stack_pool_id: 0,
239 schema_ids: FxHashMap::default(),
240 dynamic_schema_cache: FxHashMap::default(),
241 slot_cache: Vec::new(),
242 registered_ids: [0; (crate::STATIC_WIRE_ID_LIMIT as usize) / 64],
243 })
244 }
245
246 pub(crate) fn from_decoder(
249 mut registry: SchemaRegistry,
250 string_pool: crate::decoder::StringPool,
251 stack_pool: crate::decoder::StackPool,
252 timestamp_base_ns: u64,
253 writer: W,
254 ) -> Self {
255 let mut pool = FxHashMap::default();
256 let mut next_pool_id: u32 = 0;
257 for (id, value) in string_pool.0.into_iter() {
258 pool.insert(value, id.raw_id());
259 if id.raw_id() >= next_pool_id {
260 next_pool_id = id.raw_id() + 1;
261 }
262 }
263
264 let mut new_stack_pool: FxHashMap<Box<[u64]>, u32> = FxHashMap::default();
265 let mut next_stack_pool_id: u32 = 0;
266 for (id, frames) in stack_pool.0.into_iter() {
267 new_stack_pool.insert(frames.into_boxed_slice(), id.raw_id());
268 if id.raw_id() >= next_stack_pool_id {
269 next_stack_pool_id = id.raw_id() + 1;
270 }
271 }
272
273 let mut schema_ids = FxHashMap::default();
274 for (wire_id, entry) in registry.entries() {
275 schema_ids.insert(SchemaKey::Name(Arc::from(entry.name.as_str())), wire_id);
276 }
277 registry.sync_next_id();
278
279 let mut state = EncodeState::new(writer);
280 state.set_ts_base_unchecked(timestamp_base_ns);
281
282 Self {
283 state,
284 registry,
285 string_pool: pool,
286 next_pool_id,
287 stack_pool: new_stack_pool,
288 next_stack_pool_id,
289 schema_ids,
290 dynamic_schema_cache: FxHashMap::default(),
291 slot_cache: Vec::new(),
292 registered_ids: [0; (crate::STATIC_WIRE_ID_LIMIT as usize) / 64],
293 }
294 }
295
296 pub fn into_inner(self) -> W {
298 self.state.writer.into_inner()
299 }
300
301 pub fn as_inner(&self) -> &W {
303 self.state.writer.inner()
304 }
305
306 pub fn bytes_written(&self) -> u64 {
308 self.state.writer.bytes_written()
309 }
310
311 pub fn reset_to(&mut self, mut new_writer: W) -> io::Result<W> {
314 codec::encode_header(&mut new_writer)?;
315 self.string_pool.clear();
316 self.next_pool_id = 0;
317 self.stack_pool.clear();
318 self.next_stack_pool_id = 0;
319 self.registry.clear();
320 self.schema_ids.clear();
321 self.dynamic_schema_cache.clear();
322 self.slot_cache.fill(0);
323 self.registered_ids.fill(0);
324 let old_state = std::mem::replace(&mut self.state, EncodeState::new(new_writer));
326 Ok(old_state.writer.into_inner())
327 }
328
329 fn ensure_registered(&mut self, schema: &Schema) -> io::Result<WireTypeId> {
335 let identity = Arc::as_ptr(&schema.entry) as usize;
336 if let Some((_, wire_id)) = self.dynamic_schema_cache.get(&identity) {
337 return Ok(*wire_id);
338 }
339 let wire_id = self.ensure_registered_slow(schema)?;
340 if self.dynamic_schema_cache.len() >= DYNAMIC_SCHEMA_CACHE_LIMIT {
341 self.dynamic_schema_cache.clear();
345 }
346 self.dynamic_schema_cache
347 .insert(identity, (Arc::clone(&schema.entry), wire_id));
348 Ok(wire_id)
349 }
350
351 fn ensure_registered_slow(&mut self, schema: &Schema) -> io::Result<WireTypeId> {
354 let key = SchemaKey::Name(Arc::clone(&schema.name_key));
355 if let Some(&wire_id) = self.schema_ids.get(&key) {
356 let Some(existing) = self.registry.get(wire_id) else {
358 return Err(io::Error::other(format!(
359 "corrupted internal state. {wire_id:?} in schema_ids but not in registry."
360 )));
361 };
362 if *existing == *schema.entry {
363 return Ok(wire_id);
364 }
365 return Err(io::Error::new(
366 io::ErrorKind::InvalidInput,
367 format!(
368 "schema already registered with different definition: {}",
369 schema.name()
370 ),
371 ));
372 }
373 let id = self.registry.next_type_id();
374 codec::encode_schema(id, &schema.entry, &mut self.state.writer)?;
375 if !schema.entry.annotations.is_empty() {
376 codec::encode_schema_annotations(
377 id,
378 &schema.entry.annotations,
379 &mut self.state.writer,
380 )?;
381 }
382 self.registry
383 .register(id, (*schema.entry).clone())
384 .expect("schema registration failed");
385 self.schema_ids.insert(key, id);
386 Ok(id)
387 }
388
389 pub fn register_schema(
395 &mut self,
396 name: &str,
397 fields: Vec<crate::schema::FieldDef>,
398 ) -> io::Result<Schema> {
399 let schema = Schema::new(name, fields);
400 self.ensure_registered(&schema)?;
401 Ok(schema)
402 }
403
404 pub fn register_existing(&mut self, schema: &Schema) -> io::Result<WireTypeId> {
409 self.ensure_registered(schema)
410 }
411
412 pub fn write_event(
427 &mut self,
428 schema: &Schema,
429 timestamp_ns: u64,
430 values: &[crate::types::FieldValue],
431 ) -> io::Result<()> {
432 let type_id = self.ensure_registered(schema)?;
433 let expected_fields = schema.entry.fields.len();
434
435 if values.len() != expected_fields {
436 return Err(io::Error::new(
437 io::ErrorKind::InvalidInput,
438 format!(
439 "value count ({}) does not match schema field count ({}) for schema '{}'",
440 values.len(),
441 expected_fields,
442 schema.name(),
443 ),
444 ));
445 }
446
447 let ts_delta = self.state.encode_timestamp_delta(timestamp_ns)?;
448 self.state.writer.write_all(&[codec::TAG_EVENT])?;
449 self.state.writer.write_all(&type_id.0.to_le_bytes())?;
450 codec::encode_u24_le(ts_delta, &mut self.state.writer)?;
451 let mut enc = EventEncoder::new(&mut self.state);
452 for (i, v) in values.iter().enumerate() {
453 enc.write_field_value(v, schema.entry.fields[i].field_type)?;
454 }
455 Ok(())
456 }
457
458 pub fn write<T: TraceEvent>(&mut self, event: &T) -> io::Result<()> {
461 let slot = T::type_slot();
462 let tid = if slot != 0 && slot < crate::STATIC_WIRE_ID_LIMIT {
463 let word = (slot >> 6) as usize;
464 let bit = 1u64 << (slot & 63);
465 if self.registered_ids[word] & bit == 0 {
466 self.register_fast_id::<T>(slot)?;
467 }
468 WireTypeId(slot)
469 } else {
470 let s = slot as usize;
471 let cached = self.slot_cache.get(s).copied().unwrap_or(0);
472 if cached != 0 {
473 WireTypeId((cached - 1) as u16)
474 } else {
475 self.resolve_dynamic_wire_id::<T>(s)?
476 }
477 };
478 let ts_ns = event.timestamp();
479 let ts_delta = self.state.encode_timestamp_delta(ts_ns)?;
480 self.state.writer.write_all(&[codec::TAG_EVENT])?;
481 self.state.writer.write_all(&tid.0.to_le_bytes())?;
482 codec::encode_u24_le(ts_delta, &mut self.state.writer)?;
483 let mut enc = EventEncoder::new(&mut self.state);
484 event.encode_fields(&mut enc)
485 }
486
487 #[cold]
491 fn resolve_dynamic_wire_id<T: TraceEvent>(&mut self, slot: usize) -> io::Result<WireTypeId> {
492 let key = SchemaKey::RustType(typeid::of::<T>());
493 let tid = if let Some(&existing) = self.schema_ids.get(&key) {
494 existing
495 } else {
496 let schema = Schema::from_entry(T::schema_entry());
497 let id = self.ensure_registered(&schema)?;
498 self.schema_ids.insert(key, id);
499 id
500 };
501 if slot != 0 {
502 if self.slot_cache.len() <= slot {
503 self.slot_cache.resize(slot + 1, 0);
504 }
505 self.slot_cache[slot] = (tid.0 as u32) + 1;
506 }
507 Ok(tid)
508 }
509
510 #[cold]
513 fn register_fast_id<T: TraceEvent>(&mut self, id: u16) -> io::Result<()> {
514 let entry = T::schema_entry();
515 let wire = WireTypeId(id);
516 codec::encode_schema(wire, &entry, &mut self.state.writer)?;
517 if !entry.annotations.is_empty() {
518 codec::encode_schema_annotations(wire, &entry.annotations, &mut self.state.writer)?;
519 }
520 self.registry.register(wire, entry).map_err(|e| {
521 io::Error::new(
522 io::ErrorKind::InvalidInput,
523 format!("wire id {id} collision: {e}"),
524 )
525 })?;
526 self.registered_ids[(id >> 6) as usize] |= 1u64 << (id & 63);
528 Ok(())
529 }
530
531 pub fn intern_string(&mut self, s: &str) -> io::Result<InternedString> {
533 if let Some(&id) = self.string_pool.get(s) {
534 return Ok(InternedString(id));
535 }
536 let id = self.next_pool_id;
537 self.next_pool_id += 1;
538 self.string_pool.insert(s.to_string(), id);
539 codec::encode_string_pool(
540 &[PoolEntry {
541 pool_id: id,
542 data: s.as_bytes().to_vec(),
543 }],
544 &mut self.state.writer,
545 )?;
546 Ok(InternedString(id))
547 }
548
549 pub fn write_string_pool(&mut self, entries: &[PoolEntry]) -> io::Result<()> {
550 codec::encode_string_pool(entries, &mut self.state.writer)
551 }
552
553 pub fn intern_stack_frames(&mut self, frames: &[u64]) -> io::Result<InternedStackFrames> {
556 if let Some(&id) = self.stack_pool.get(frames) {
557 return Ok(InternedStackFrames(id));
558 }
559 let id = self.next_stack_pool_id;
560 self.next_stack_pool_id += 1;
561 self.stack_pool.insert(frames.into(), id);
562 codec::encode_stack_pool(
563 &[StackPoolEntry {
564 pool_id: id,
565 frames: frames.to_vec(),
568 }],
569 &mut self.state.writer,
570 )?;
571 Ok(InternedStackFrames(id))
572 }
573
574 pub fn write_stack_pool(&mut self, entries: &[StackPoolEntry]) -> io::Result<()> {
575 codec::encode_stack_pool(entries, &mut self.state.writer)
576 }
577
578 pub fn flush(&mut self) -> io::Result<()> {
580 self.state.writer.flush()
581 }
582
583 pub fn into_raw_encoder(self) -> RawEncoder<W> {
590 RawEncoder {
591 writer: self.state.writer,
592 }
593 }
594}
595
596pub struct RawEncoder<W> {
603 writer: CountingWriter<W>,
604}
605
606impl<W: Write> RawEncoder<W> {
607 pub fn write_raw(&mut self, bytes: &[u8]) -> io::Result<()> {
609 self.writer.write_all(bytes)
610 }
611
612 pub fn bytes_written(&self) -> u64 {
615 self.writer.bytes_written()
616 }
617
618 pub fn flush(&mut self) -> io::Result<()> {
620 self.writer.flush()
621 }
622
623 pub fn into_inner(self) -> W {
625 self.writer.into_inner()
626 }
627}
628
629impl Encoder<Vec<u8>> {
630 pub fn write_infallible<T: TraceEvent>(&mut self, event: &T) {
631 self.write(event).expect("writing to Vec<u8> is infallible")
632 }
633
634 pub fn intern_string_infallible(&mut self, s: &str) -> InternedString {
635 self.intern_string(s)
636 .expect("interning into Vec<u8> is infallible")
637 }
638
639 pub fn intern_stack_frames_infallible(&mut self, frames: &[u64]) -> InternedStackFrames {
640 self.intern_stack_frames(frames)
641 .expect("interning into Vec<u8> is infallible")
642 }
643
644 pub fn reset_to_infallible(&mut self, data: Vec<u8>) -> Vec<u8> {
646 self.reset_to(data)
647 .expect("writing to Vec<u8> is infallible")
648 }
649}
650
651#[cfg(test)]
652mod tests {
653 use super::*;
654 use crate::schema::FieldDef;
655 use crate::types::{FieldType, FieldValue};
656
657 #[test]
658 fn encoder_writes_header() {
659 let enc = Encoder::new();
660 let data = enc.finish();
661 assert_eq!(&data[..5], &[0x54, 0x52, 0x43, 0x00, 1]);
662 }
663
664 #[test]
665 fn encoder_register_and_write_event() {
666 let mut enc = Encoder::new();
667 let schema = enc
668 .register_schema(
669 "Ev",
670 vec![FieldDef {
671 name: "v".into(),
672 field_type: FieldType::Varint,
673 }],
674 )
675 .unwrap();
676 enc.write_event(&schema, 1_000, &[FieldValue::Varint(42)])
677 .unwrap();
678 let data = enc.finish();
679 assert!(data.len() > 5);
680 }
681
682 #[test]
683 fn idempotent_re_registration() {
684 let mut enc = Encoder::new();
685 let fields = vec![FieldDef {
686 name: "v".into(),
687 field_type: FieldType::Varint,
688 }];
689 let _s1 = enc.register_schema("Ev", fields.clone()).unwrap();
690 let _s2 = enc.register_schema("Ev", fields).unwrap();
691 }
693
694 #[test]
695 fn re_registration_different_schema_errors() {
696 let mut enc = Encoder::new();
697 enc.register_schema(
698 "Ev",
699 vec![FieldDef {
700 name: "v".into(),
701 field_type: FieldType::Varint,
702 }],
703 )
704 .unwrap();
705 let result = enc.register_schema(
706 "Ev",
707 vec![FieldDef {
708 name: "different".into(),
709 field_type: FieldType::Bool,
710 }],
711 );
712 assert!(result.is_err());
713 }
714
715 #[test]
716 fn schema_auto_registers_on_write() {
717 use crate::decoder::{DecodedFrame, Decoder};
718
719 let schema = Schema::new(
721 "Lazy",
722 vec![FieldDef {
723 name: "v".into(),
724 field_type: FieldType::Varint,
725 }],
726 );
727
728 let mut enc = Encoder::new();
730 enc.write_event(&schema, 1_000, &[FieldValue::Varint(42)])
731 .unwrap();
732
733 let bytes = enc.finish();
734 let mut dec = Decoder::new(&bytes).unwrap();
735 let frames = dec.decode_all();
736 assert!(matches!(&frames[0], DecodedFrame::Schema(s) if s.name == "Lazy"));
737 if let DecodedFrame::Event { values, .. } = &frames[1] {
738 assert_eq!(*values, vec![FieldValue::Varint(42)]);
739 } else {
740 panic!("expected event");
741 }
742 }
743
744 #[test]
745 fn schema_portable_across_encoders() {
746 use crate::decoder::{DecodedFrame, Decoder};
747
748 let mut enc1 = Encoder::new();
749 let schema = enc1
750 .register_schema(
751 "Shared",
752 vec![FieldDef {
753 name: "v".into(),
754 field_type: FieldType::Varint,
755 }],
756 )
757 .unwrap();
758 enc1.write_event(&schema, 1_000, &[FieldValue::Varint(1)])
759 .unwrap();
760
761 let mut enc2 = Encoder::new();
763 enc2.write_event(&schema, 2_000, &[FieldValue::Varint(2)])
764 .unwrap();
765
766 for (enc, expected_val) in [(enc1, 1u64), (enc2, 2u64)] {
768 let bytes = enc.finish();
769 let mut dec = Decoder::new(&bytes).unwrap();
770 let frames = dec.decode_all();
771 let event = frames
772 .iter()
773 .find(|f| matches!(f, DecodedFrame::Event { .. }))
774 .unwrap();
775 if let DecodedFrame::Event { values, .. } = event {
776 assert_eq!(values[0], FieldValue::Varint(expected_val));
777 }
778 }
779 }
780
781 #[test]
782 fn encoder_intern_string_deduplicates() {
783 let mut enc = Encoder::new();
784 let id1 = enc.intern_string("hello").unwrap();
785 let id2 = enc.intern_string("hello").unwrap();
786 let id3 = enc.intern_string("world").unwrap();
787 assert_eq!(id1, id2);
788 assert_ne!(id1, id3);
789 }
790
791 #[test]
792 fn encoder_intern_stack_frames_deduplicates() {
793 let mut enc = Encoder::new();
794 let stack_a: &[u64] = &[0x1000, 0x2000, 0x3000];
795 let stack_b: &[u64] = &[0x4000, 0x5000];
796 let id1 = enc.intern_stack_frames(stack_a).unwrap();
797 let id2 = enc.intern_stack_frames(stack_a).unwrap();
798 let id3 = enc.intern_stack_frames(stack_b).unwrap();
799 assert_eq!(id1, id2);
800 assert_ne!(id1, id3);
801 }
802
803 #[test]
804 fn stack_pool_round_trip_via_decoder() {
805 use crate::decoder::Decoder;
806 use crate::types::InternedStackFrames;
807
808 let mut enc = Encoder::new();
809 let stack_a: &[u64] = &[0xdead, 0xbeef, 0xcafe];
810 let stack_b: &[u64] = &[0x1, 0x2];
811 let id_a = enc.intern_stack_frames(stack_a).unwrap();
812 let id_b = enc.intern_stack_frames(stack_b).unwrap();
813 let bytes = enc.finish();
814
815 let mut dec = Decoder::new(&bytes).unwrap();
816 let _ = dec.decode_all();
817 assert_eq!(
818 dec.stack_pool().get(InternedStackFrames(id_a.raw_id())),
819 Some(stack_a)
820 );
821 assert_eq!(
822 dec.stack_pool().get(InternedStackFrames(id_b.raw_id())),
823 Some(stack_b)
824 );
825 }
826
827 #[test]
828 fn for_each_event_populates_stack_pool() {
829 use crate::decoder::Decoder;
830 use crate::schema::FieldDef;
831 use crate::types::{FieldType, FieldValue, InternedStackFrames};
832
833 let mut enc = Encoder::new();
834 let schema = enc
835 .register_schema(
836 "CpuSampleEvent",
837 vec![FieldDef {
838 name: "callchain".into(),
839 field_type: FieldType::PooledStackFrames,
840 }],
841 )
842 .unwrap();
843 let stack: &[u64] = &[0x1234, 0x5678, 0x9abc];
844 let id = enc.intern_stack_frames(stack).unwrap();
845 enc.write_event(&schema, 1_000_000, &[FieldValue::PooledStackFrames(id)])
846 .unwrap();
847 let bytes = enc.finish();
848
849 let mut dec = Decoder::new(&bytes).unwrap();
850 let mut event_count = 0;
851 dec.for_each_event(|_ev| {
852 event_count += 1;
853 })
854 .unwrap();
855 assert_eq!(event_count, 1);
856 assert_eq!(
857 dec.stack_pool().get(InternedStackFrames(id.raw_id())),
858 Some(stack),
859 );
860 }
861
862 #[test]
863 fn encoder_intern_empty_stack_frames() {
864 use crate::decoder::Decoder;
865 use crate::types::InternedStackFrames;
866
867 let mut enc = Encoder::new();
868 let id1 = enc.intern_stack_frames(&[]).unwrap();
869 let id2 = enc.intern_stack_frames(&[]).unwrap();
870 assert_eq!(id1, id2);
871 let bytes = enc.finish();
872
873 let mut dec = Decoder::new(&bytes).unwrap();
874 let _ = dec.decode_all();
875 assert_eq!(
876 dec.stack_pool().get(InternedStackFrames(id1.raw_id())),
877 Some(&[][..])
878 );
879 }
880
881 #[test]
882 fn write_stack_pool_multi_entry_round_trip() {
883 use crate::decoder::Decoder;
884 use crate::types::InternedStackFrames;
885
886 let mut enc = Encoder::new();
887 let entries = vec![
888 StackPoolEntry {
889 pool_id: 0,
890 frames: vec![0xaaaa, 0xbbbb, 0xcccc],
891 },
892 StackPoolEntry {
893 pool_id: 1,
894 frames: vec![0x1111],
895 },
896 StackPoolEntry {
897 pool_id: 2,
898 frames: vec![],
899 },
900 ];
901 enc.write_stack_pool(&entries).unwrap();
902 let bytes = enc.finish();
903
904 let mut dec = Decoder::new(&bytes).unwrap();
905 let _ = dec.decode_all();
906 assert_eq!(
907 dec.stack_pool().get(InternedStackFrames(0)),
908 Some(&[0xaaaa, 0xbbbb, 0xcccc][..])
909 );
910 assert_eq!(
911 dec.stack_pool().get(InternedStackFrames(1)),
912 Some(&[0x1111][..])
913 );
914 assert_eq!(dec.stack_pool().get(InternedStackFrames(2)), Some(&[][..]));
915 }
916
917 #[test]
918 fn decoder_into_encoder_deduplicates_interned_stack_frames() {
919 use crate::decoder::Decoder;
920
921 let mut enc = Encoder::new();
922 let id1 = enc.intern_stack_frames(&[0x10, 0x20]).unwrap();
923 let base = enc.finish();
924
925 let mut decoder = Decoder::new(&base).unwrap();
926 while decoder.next_frame_ref().ok().flatten().is_some() {}
927 let mut output = Vec::new();
928 let mut ext = decoder.into_encoder(&mut output);
929 let id2 = ext.intern_stack_frames(&[0x10, 0x20]).unwrap();
930 let id3 = ext.intern_stack_frames(&[0x30]).unwrap();
931 assert_eq!(id1.raw_id(), id2.raw_id());
932 assert_ne!(id2.raw_id(), id3.raw_id());
933 }
934
935 #[test]
936 fn timestamp_round_trip() {
937 use crate::decoder::{DecodedFrame, Decoder};
938
939 let mut enc = Encoder::new();
940 let schema = enc
941 .register_schema(
942 "TS",
943 vec![FieldDef {
944 name: "v".into(),
945 field_type: FieldType::Varint,
946 }],
947 )
948 .unwrap();
949
950 let ts1 = 100_000u64;
951 let ts2 = 50_000u64;
952 let ts3 = 200_000_000u64;
953 let ts4 = 100_000_000u64;
954 enc.write_event(&schema, ts1, &[FieldValue::Varint(1)])
955 .unwrap();
956 enc.write_event(&schema, ts2, &[FieldValue::Varint(2)])
957 .unwrap();
958 enc.write_event(&schema, ts3, &[FieldValue::Varint(3)])
959 .unwrap();
960 enc.write_event(&schema, ts4, &[FieldValue::Varint(4)])
961 .unwrap();
962
963 let bytes = enc.finish();
964 let mut dec = Decoder::new(&bytes).unwrap();
965 let events: Vec<_> = dec
966 .decode_all()
967 .into_iter()
968 .filter_map(|f| match f {
969 DecodedFrame::Event {
970 timestamp_ns,
971 values,
972 ..
973 } => Some((timestamp_ns, values)),
974 _ => None,
975 })
976 .collect();
977
978 assert_eq!(events.len(), 4);
979 assert_eq!(events[0].0, ts1);
980 assert_eq!(events[0].1, vec![FieldValue::Varint(1)]);
981 assert_eq!(events[1].0, ts2);
982 assert_eq!(events[1].1, vec![FieldValue::Varint(2)]);
983 assert_eq!(events[2].0, ts3);
984 assert_eq!(events[2].1, vec![FieldValue::Varint(3)]);
985 assert_eq!(events[3].0, ts4);
986 assert_eq!(events[3].1, vec![FieldValue::Varint(4)]);
987 }
988
989 #[test]
990 fn encoder_new_to_writer() {
991 let mut buf = Vec::new();
992 let enc = Encoder::new_to(&mut buf).unwrap();
993 drop(enc);
994 assert!(buf.len() >= 5);
995 assert_eq!(&buf[..5], &[0x54, 0x52, 0x43, 0x00, 1]);
996 }
997
998 #[test]
999 fn decoder_into_encoder_appends_without_header() {
1000 use crate::decoder::{DecodedFrame, Decoder};
1001
1002 let mut enc = Encoder::new();
1004 let schema = enc
1005 .register_schema(
1006 "Ev",
1007 vec![FieldDef {
1008 name: "v".into(),
1009 field_type: FieldType::Varint,
1010 }],
1011 )
1012 .unwrap();
1013 enc.write_event(&schema, 1_000, &[FieldValue::Varint(1)])
1014 .unwrap();
1015 let base = enc.finish();
1016
1017 let mut decoder = Decoder::new(&base).unwrap();
1019 while decoder.next_frame_ref().ok().flatten().is_some() {}
1020 let mut output = Vec::new();
1021 let mut ext = decoder.into_encoder(&mut output);
1022 ext.write_event(&schema, 2_000, &[FieldValue::Varint(2)])
1024 .unwrap();
1025 drop(ext);
1026
1027 let mut combined = base.clone();
1029 combined.extend_from_slice(&output);
1030 let mut dec = Decoder::new(&combined).unwrap();
1031 let events: Vec<_> = dec
1032 .decode_all()
1033 .into_iter()
1034 .filter_map(|f| match f {
1035 DecodedFrame::Event {
1036 timestamp_ns,
1037 values,
1038 ..
1039 } => Some((timestamp_ns, values)),
1040 _ => None,
1041 })
1042 .collect();
1043 assert_eq!(events.len(), 2);
1044 assert_eq!(events[0].0, 1_000);
1045 assert_eq!(events[1].0, 2_000);
1046 }
1047
1048 #[test]
1049 fn decoder_into_encoder_deduplicates_interned_strings() {
1050 use crate::decoder::{DecodedFrame, Decoder};
1051
1052 let mut enc = Encoder::new();
1054 let id1 = enc.intern_string("hello").unwrap();
1055 let base = enc.finish();
1056
1057 let mut decoder = Decoder::new(&base).unwrap();
1059 while decoder.next_frame_ref().ok().flatten().is_some() {}
1060 let mut output = Vec::new();
1061 let mut ext = decoder.into_encoder(&mut output);
1062 let id2 = ext.intern_string("hello").unwrap();
1064 let id3 = ext.intern_string("world").unwrap();
1065 drop(ext);
1066
1067 assert_eq!(id1, id2, "existing string should reuse pool ID");
1068 assert_ne!(id2, id3);
1069
1070 let mut combined = base.clone();
1072 combined.extend_from_slice(&output);
1073 let mut dec = Decoder::new(&combined).unwrap();
1074 let frames = dec.decode_all();
1075 let pool_frames: Vec<_> = frames
1076 .iter()
1077 .filter(|f| matches!(f, DecodedFrame::StringPool(_)))
1078 .collect();
1079 assert_eq!(pool_frames.len(), 2);
1081 }
1082
1083 struct FastSlot {
1086 ts: u64,
1087 }
1088 impl TraceEvent for FastSlot {
1089 fn type_slot() -> u16 {
1090 5
1091 }
1092 fn event_name() -> &'static str {
1093 "FastSlot"
1094 }
1095 fn field_defs() -> Vec<FieldDef> {
1096 Vec::new()
1097 }
1098 fn timestamp(&self) -> u64 {
1099 self.ts
1100 }
1101 fn encode_fields<W: Write>(&self, _enc: &mut EventEncoder<'_, W>) -> io::Result<()> {
1102 Ok(())
1103 }
1104 }
1105
1106 #[test]
1107 fn fast_slot_registers_at_slot_id() {
1108 use crate::decoder::Decoder;
1109
1110 let mut enc = Encoder::new();
1111 enc.write(&FastSlot { ts: 1_000_000 }).unwrap();
1112 let dynamic = enc
1114 .register_schema(
1115 "Dyn",
1116 vec![FieldDef {
1117 name: "v".into(),
1118 field_type: FieldType::Varint,
1119 }],
1120 )
1121 .unwrap();
1122 enc.write_event(&dynamic, 2_000, &[FieldValue::Varint(1)])
1123 .unwrap();
1124 let bytes = enc.finish();
1125
1126 let mut dec = Decoder::new(&bytes).unwrap();
1127 let _ = dec.decode_all();
1128 assert_eq!(
1130 dec.registry().get(WireTypeId(5)).unwrap().name(),
1131 "FastSlot"
1132 );
1133 assert_eq!(
1135 dec.registry()
1136 .get(WireTypeId(crate::STATIC_WIRE_ID_LIMIT))
1137 .unwrap()
1138 .name(),
1139 "Dyn"
1140 );
1141 }
1142
1143 #[test]
1144 fn register_and_write() {
1145 use crate::decoder::{DecodedFrame, Decoder};
1146
1147 let mut enc = Encoder::new();
1148 let schema = enc
1149 .register_schema(
1150 "MyEvent",
1151 vec![
1152 FieldDef {
1153 name: "count".into(),
1154 field_type: FieldType::Varint,
1155 },
1156 FieldDef {
1157 name: "name".into(),
1158 field_type: FieldType::String,
1159 },
1160 ],
1161 )
1162 .unwrap();
1163
1164 enc.write_event(
1165 &schema,
1166 1_000_000,
1167 &[FieldValue::Varint(42), FieldValue::String("hello".into())],
1168 )
1169 .unwrap();
1170
1171 let bytes = enc.finish();
1172 let mut dec = Decoder::new(&bytes).unwrap();
1173 let frames = dec.decode_all();
1174 let events: Vec<_> = frames
1175 .into_iter()
1176 .filter_map(|f| match f {
1177 DecodedFrame::Event {
1178 timestamp_ns,
1179 values,
1180 ..
1181 } => Some((timestamp_ns, values)),
1182 _ => None,
1183 })
1184 .collect();
1185 assert_eq!(events.len(), 1);
1186 assert_eq!(events[0].0, 1_000_000);
1187 assert_eq!(events[0].1[0], FieldValue::Varint(42));
1188 assert_eq!(events[0].1[1], FieldValue::String("hello".into()));
1189 }
1190
1191 #[test]
1192 fn register_conflict_errors() {
1193 let mut enc = Encoder::new();
1194 enc.register_schema(
1195 "Ev",
1196 vec![FieldDef {
1197 name: "v".into(),
1198 field_type: FieldType::Varint,
1199 }],
1200 )
1201 .unwrap();
1202 let result = enc.register_schema(
1203 "Ev",
1204 vec![FieldDef {
1205 name: "other".into(),
1206 field_type: FieldType::Bool,
1207 }],
1208 );
1209 assert!(result.is_err());
1210 }
1211
1212 #[test]
1213 fn write_wrong_field_count_errors() {
1214 let mut enc = Encoder::new();
1215 let schema = enc
1216 .register_schema(
1217 "Ev",
1218 vec![FieldDef {
1219 name: "v".into(),
1220 field_type: FieldType::Varint,
1221 }],
1222 )
1223 .unwrap();
1224 let result = enc.write_event(&schema, 0, &[FieldValue::Varint(1), FieldValue::Varint(2)]);
1226 assert!(result.is_err());
1227 }
1228
1229 #[test]
1232 fn timestamp_base_advances_per_event() {
1233 use crate::decoder::{DecodedFrame, Decoder};
1234
1235 let mut enc = Encoder::new();
1236 let schema = enc
1237 .register_schema(
1238 "Ev",
1239 vec![FieldDef {
1240 name: "v".into(),
1241 field_type: FieldType::Varint,
1242 }],
1243 )
1244 .unwrap();
1245
1246 let ts1 = 12_000_000u64;
1247 let ts2 = 24_000_000u64;
1248 enc.write_event(&schema, ts1, &[FieldValue::Varint(1)])
1249 .unwrap();
1250 enc.write_event(&schema, ts2, &[FieldValue::Varint(2)])
1251 .unwrap();
1252
1253 let bytes = enc.finish();
1254
1255 let reset_count = bytes.iter().filter(|&&b| b == 0x05).count();
1256 assert_eq!(
1257 reset_count, 0,
1258 "base should advance per event, avoiding unnecessary resets"
1259 );
1260
1261 let mut dec = Decoder::new(&bytes).unwrap();
1262 let events: Vec<_> = dec
1263 .decode_all()
1264 .into_iter()
1265 .filter_map(|f| match f {
1266 DecodedFrame::Event { timestamp_ns, .. } => Some(timestamp_ns),
1267 _ => None,
1268 })
1269 .collect();
1270 assert_eq!(events, vec![ts1, ts2]);
1271 }
1272
1273 #[test]
1274 fn reset_to_preserves_capacity() {
1275 let mut enc = Encoder::new();
1276 for i in 0..100 {
1277 enc.intern_string(&format!("string_{}", i)).unwrap();
1278 }
1279 let cap_before = enc.string_pool.capacity();
1280 let _bytes = enc.reset_to(Vec::new());
1281 let cap_after = enc.string_pool.capacity();
1282 assert_eq!(
1283 cap_before, cap_after,
1284 "string_pool capacity should be preserved after reset_to"
1285 );
1286 }
1287
1288 #[test]
1289 fn reset_to_returns_old_data_and_clears_state() {
1290 use crate::decoder::{DecodedFrame, Decoder};
1291
1292 let mut enc = Encoder::new();
1293 let schema = enc
1294 .register_schema(
1295 "Ev",
1296 vec![FieldDef {
1297 name: "v".into(),
1298 field_type: FieldType::Varint,
1299 }],
1300 )
1301 .unwrap();
1302 enc.write_event(&schema, 1_000, &[FieldValue::Varint(42)])
1303 .unwrap();
1304 let _s = enc.intern_string("hello").unwrap();
1305
1306 let old_bytes_written = enc.bytes_written();
1307 assert!(old_bytes_written > 0);
1308
1309 let old = enc.reset_to_infallible(Vec::new());
1311
1312 let mut dec = Decoder::new(&old).unwrap();
1314 let frames = dec.decode_all();
1315 assert!(frames.iter().any(|f| matches!(f, DecodedFrame::Schema(_))));
1316 assert!(
1317 frames
1318 .iter()
1319 .any(|f| matches!(f, DecodedFrame::Event { .. }))
1320 );
1321 assert!(
1322 frames
1323 .iter()
1324 .any(|f| matches!(f, DecodedFrame::StringPool(_)))
1325 );
1326
1327 assert!(
1329 enc.bytes_written() < old_bytes_written,
1330 "bytes_written should reset (got {} vs old {})",
1331 enc.bytes_written(),
1332 old_bytes_written
1333 );
1334
1335 enc.write_event(&schema, 2_000, &[FieldValue::Varint(99)])
1338 .unwrap();
1339
1340 let _s2 = enc.intern_string("hello").unwrap();
1342
1343 let new_bytes = enc.reset_to_infallible(Vec::new());
1345 let mut dec2 = Decoder::new(&new_bytes).unwrap();
1346 let new_frames = dec2.decode_all();
1347 assert!(
1349 new_frames
1350 .iter()
1351 .any(|f| matches!(f, DecodedFrame::Schema(s) if s.name == "Ev")),
1352 "new trace must contain schema definition"
1353 );
1354 assert!(
1356 new_frames
1357 .iter()
1358 .any(|f| matches!(f, DecodedFrame::StringPool(_))),
1359 "new trace must contain string pool"
1360 );
1361 let event = new_frames
1363 .iter()
1364 .find_map(|f| match f {
1365 DecodedFrame::Event {
1366 timestamp_ns,
1367 values,
1368 ..
1369 } => Some((timestamp_ns, values)),
1370 _ => None,
1371 })
1372 .expect("new trace must contain event");
1373 assert_eq!(*event.0, 2_000);
1374 assert_eq!(event.1[0], FieldValue::Varint(99));
1375 }
1376
1377 #[test]
1378 fn into_raw_encoder_preserves_byte_count() {
1379 let mut enc = Encoder::new();
1380 let schema = enc
1381 .register_schema(
1382 "Ev",
1383 vec![FieldDef {
1384 name: "v".into(),
1385 field_type: FieldType::Varint,
1386 }],
1387 )
1388 .unwrap();
1389 enc.write_event(&schema, 1_000, &[FieldValue::Varint(42)])
1390 .unwrap();
1391
1392 let bytes_before = enc.bytes_written();
1393 assert!(bytes_before > 0);
1394
1395 let raw = enc.into_raw_encoder();
1396 assert_eq!(
1397 raw.bytes_written(),
1398 bytes_before,
1399 "byte count must be preserved across conversion"
1400 );
1401 }
1402
1403 #[test]
1404 fn raw_encoder_write_raw_and_bytes_written() {
1405 let enc = Encoder::new();
1406 let initial = enc.bytes_written();
1407 let mut raw = enc.into_raw_encoder();
1408
1409 let payload = [0xAA; 100];
1410 raw.write_raw(&payload).unwrap();
1411
1412 assert_eq!(
1413 raw.bytes_written(),
1414 initial + payload.len() as u64,
1415 "bytes_written must include raw payload"
1416 );
1417 }
1418
1419 #[test]
1420 fn raw_encoder_into_inner_returns_all_data() {
1421 use crate::decoder::{DecodedFrame, Decoder};
1422
1423 let mut enc = Encoder::new();
1426 let schema = enc
1427 .register_schema(
1428 "Ev",
1429 vec![FieldDef {
1430 name: "v".into(),
1431 field_type: FieldType::Varint,
1432 }],
1433 )
1434 .unwrap();
1435 enc.write_event(&schema, 1_000, &[FieldValue::Varint(1)])
1436 .unwrap();
1437
1438 let raw_batch = {
1440 let mut batch_enc = Encoder::new();
1441 batch_enc
1442 .write_event(&schema, 2_000, &[FieldValue::Varint(2)])
1443 .unwrap();
1444 batch_enc.finish()
1445 };
1446
1447 let mut raw = enc.into_raw_encoder();
1448 raw.write_raw(&raw_batch).unwrap();
1449 let combined = raw.into_inner();
1450
1451 let mut dec = Decoder::new(&combined).unwrap();
1452 let events: Vec<_> = dec
1453 .decode_all()
1454 .into_iter()
1455 .filter_map(|f| match f {
1456 DecodedFrame::Event {
1457 timestamp_ns,
1458 values,
1459 ..
1460 } => Some((timestamp_ns, values)),
1461 _ => None,
1462 })
1463 .collect();
1464
1465 assert_eq!(events.len(), 2);
1466 assert_eq!(events[0].0, 1_000);
1467 assert_eq!(events[0].1, vec![FieldValue::Varint(1)]);
1468 assert_eq!(events[1].0, 2_000);
1469 assert_eq!(events[1].1, vec![FieldValue::Varint(2)]);
1470 }
1471}
1472
1473#[cfg(test)]
1474mod dynamic_schema_cache_tests {
1475 use super::*;
1476 use crate::schema::{FieldDef, SchemaEntry};
1477 use crate::types::{FieldType, FieldValue};
1478
1479 fn schema(name: &str) -> Schema {
1480 Schema::from_entry(SchemaEntry::new(
1481 name,
1482 vec![FieldDef::new("v", FieldType::Varint)],
1483 ))
1484 }
1485
1486 #[test]
1490 fn identity_cache_agrees_with_name_registration() {
1491 let mut enc = Encoder::new();
1492 let a = schema("Ev");
1493 let id1 = enc.ensure_registered(&a).unwrap();
1494 let id2 = enc.ensure_registered(&a).unwrap();
1495 assert_eq!(id1, id2, "same handle must reuse its wire id");
1496
1497 let b = schema("Ev"); let id3 = enc.ensure_registered(&b).unwrap();
1499 assert_eq!(id1, id3, "same name must resolve to the same wire id");
1500 assert_eq!(enc.dynamic_schema_cache.len(), 2);
1501 }
1502
1503 #[test]
1506 fn identity_cache_is_bounded() {
1507 let mut enc = Encoder::new();
1508 for i in 0..(DYNAMIC_SCHEMA_CACHE_LIMIT * 2 + 7) {
1509 let s = schema(&format!("Ev{}", i % 3));
1512 enc.ensure_registered(&s).unwrap();
1513 assert!(
1514 enc.dynamic_schema_cache.len() <= DYNAMIC_SCHEMA_CACHE_LIMIT,
1515 "cache exceeded its bound at iteration {i}"
1516 );
1517 }
1518 }
1519
1520 #[test]
1523 fn events_across_cache_clears_decode() {
1524 let mut enc = Encoder::new();
1525 for i in 0..(DYNAMIC_SCHEMA_CACHE_LIMIT + 3) {
1526 let s = schema("Ev");
1527 enc.write_event(&s, i as u64, &[FieldValue::Varint(i as u64)])
1529 .unwrap();
1530 }
1531 let data = enc.finish();
1532 let mut decoder = crate::decoder::Decoder::new(&data).unwrap();
1533 let mut count = 0u64;
1534 decoder
1535 .for_each_event(|ev| {
1536 assert_eq!(ev.name, "Ev");
1537 count += 1;
1538 })
1539 .unwrap();
1540 assert_eq!(count, (DYNAMIC_SCHEMA_CACHE_LIMIT + 3) as u64);
1541 }
1542}