Skip to main content

pylon_replication/
datagram.rs

1//! The replication datagram: entity updates that may be lost, duplicated,
2//! or reordered (a WebTransport datagram, or a relay of one).
3//!
4//! A subscription on an unreliable transport gets spawns, despawns, and
5//! full frames as ordinary replication frames ([`crate::frame`]) on a
6//! reliable, ordered stream, and its updates in datagrams. Every datagram
7//! stands alone: an update carries the entity's **absolute** quantized
8//! position, not a difference, so a lost datagram only delays what it held.
9//! The client acks the frame numbers it receives; the server keeps sending
10//! whatever changed since the last acked state until an ack covers it.
11//!
12//! ```text
13//! u8      version (2)
14//! varint  frame number (per subscription, from 1)
15//! varint  tick
16//! varint  input ack: the highest client_seq the shard has processed
17//! varint  stream tick: the tick of the last stream frame sent by this tick
18//! f32 LE  precision
19//! varint  parts: datagrams sent for this tick
20//! varint  update count, then per entity (ids ascend, delta-coded):
21//!           id, u16 LE spawn tick, u8 mask (1 x, 2 y, 4 z, 8 components),
22//!           the axes in the mask (zigzag, absolute quantized),
23//!           components when bit 8 is set (as in a frame)
24//! ```
25//!
26//! **Spawn tick.** An update names the tick of the stream frame that
27//! spawned the entity (its low 16 bits), and the client applies it only to
28//! the entity spawned then. A late datagram for an entity that has since
29//! been despawned and spawned again under the same id, or one built before
30//! a full frame, changes nothing. Ticks travel in every frame header, so
31//! the two sides agree even when the server dropped frames the client never
32//! got, which counting spawns could not survive.
33//!
34//! **Acks.** The client acks a datagram with its number and the tick of
35//! the last stream frame it had applied (`ReplicaTable::stream_tick`). An
36//! ack covers an entity only when that tick is at or after its spawn.
37//!
38//! **Order.** The client keeps, per entity, the number of the last datagram
39//! it applied, and skips an update from an older one.
40//!
41//! **Whole ticks.** Every tick without a full frame sends at least one
42//! datagram, and each names how many the tick has and the tick of the last
43//! stream frame sent by then. A tick is whole on the client when it has
44//! all of its datagrams and that stream frame. Only then does the client's
45//! table hold the server's state as of the tick, so only then does the
46//! input ack describe the table (prediction depends on that): the client
47//! applies a tick's datagrams, and takes its ack, when the tick is whole or
48//! a later tick is.
49
50use crate::frame::{mask, ComponentChange, DecodeError, ReplicaTable};
51use crate::varint;
52use crate::EntityId;
53
54pub const VERSION: u8 = 2;
55
56/// Bytes before the first update, at most: version, four varints, the
57/// precision, the parts, and the update count.
58pub const MAX_HEADER_LEN: usize = 1 + 10 + 10 + 10 + 10 + 4 + PARTS_LEN + 3;
59
60/// Bytes kept for the parts count, which is known only when the tick's
61/// last datagram is built (up to 2^21 parts).
62const PARTS_LEN: usize = 3;
63
64/// The part of a spawn tick an update carries.
65pub fn spawn_tag(spawn_tick: u64) -> u16 {
66    spawn_tick as u16
67}
68
69/// Encode one update body (everything after the id): the spawn tick's tag,
70/// mask, the axes given, and the component changes.
71pub fn encode_entry_into(
72    out: &mut Vec<u8>,
73    spawn_tick: u64,
74    axes: [Option<i64>; 3],
75    components: &[ComponentChange<'_>],
76) {
77    out.extend_from_slice(&spawn_tag(spawn_tick).to_le_bytes());
78    let mut m = 0u8;
79    for (i, bit) in [mask::X, mask::Y, mask::Z].into_iter().enumerate() {
80        if axes[i].is_some() {
81            m |= bit;
82        }
83    }
84    if !components.is_empty() {
85        m |= mask::COMPONENTS;
86    }
87    out.push(m);
88    for v in axes.into_iter().flatten() {
89        varint::write_i64(out, v);
90    }
91    if !components.is_empty() {
92        varint::write_u64(out, components.len() as u64);
93        for (id, bytes) in components {
94            out.push(*id);
95            match bytes {
96                Some(b) => {
97                    varint::write_u64(out, b.len() as u64 + 1);
98                    out.extend_from_slice(b);
99                }
100                None => varint::write_u64(out, 0),
101            }
102        }
103    }
104}
105
106/// Builds one datagram. Add updates in ascending id order.
107pub struct DatagramBuilder {
108    header: Vec<u8>,
109    body: Vec<u8>,
110    count: u64,
111    last: Option<EntityId>,
112}
113
114impl DatagramBuilder {
115    pub fn new(frame: u64, tick: u64, ack: u64, stream_tick: u64, precision: f32) -> Self {
116        let mut header = Vec::with_capacity(MAX_HEADER_LEN);
117        header.push(VERSION);
118        varint::write_u64(&mut header, frame);
119        varint::write_u64(&mut header, tick);
120        varint::write_u64(&mut header, ack);
121        varint::write_u64(&mut header, stream_tick);
122        header.extend_from_slice(&precision.to_le_bytes());
123        Self {
124            header,
125            body: Vec::new(),
126            count: 0,
127            last: None,
128        }
129    }
130
131    /// The datagram's length if an update with `entry` (from
132    /// [`encode_entry_into`]) for `id` were added.
133    pub fn len_with(&self, id: EntityId, entry: &[u8]) -> usize {
134        let id_len = match self.last {
135            Some(prev) => varint::len_u64(id.wrapping_sub(prev)),
136            None => varint::len_u64(id),
137        };
138        self.header.len()
139            + PARTS_LEN
140            + varint::len_u64(self.count + 1)
141            + self.body.len()
142            + id_len
143            + entry.len()
144    }
145
146    pub fn is_empty(&self) -> bool {
147        self.count == 0
148    }
149
150    /// Add an update. Ids must ascend across calls.
151    pub fn update(&mut self, id: EntityId, entry: &[u8]) {
152        match self.last {
153            Some(prev) => {
154                debug_assert!(id > prev, "ids must ascend");
155                varint::write_u64(&mut self.body, id.wrapping_sub(prev));
156            }
157            None => varint::write_u64(&mut self.body, id),
158        }
159        self.last = Some(id);
160        self.body.extend_from_slice(entry);
161        self.count += 1;
162    }
163
164    /// The datagram, as one of `parts` for its tick.
165    pub fn finish(mut self, parts: u64) -> Vec<u8> {
166        debug_assert!(parts < 1 << 21, "parts must fit PARTS_LEN");
167        varint::write_u64(&mut self.header, parts);
168        varint::write_u64(&mut self.header, self.count);
169        self.header.extend_from_slice(&self.body);
170        self.header
171    }
172}
173
174/// What one datagram did.
175#[derive(Debug, Clone, Default, PartialEq, Eq)]
176pub struct DatagramSummary {
177    pub frame: u64,
178    pub tick: u64,
179    pub ack: u64,
180    /// The tick of the last stream frame the server had sent by `tick`.
181    pub stream_tick: u64,
182    /// Datagrams the server sent for `tick`.
183    pub parts: u64,
184    /// Entities it updated.
185    pub updated: Vec<EntityId>,
186    /// Updates it skipped: an unknown entity, another spawn, an older
187    /// datagram than the entity's last, or a precision the table is not in.
188    pub skipped: usize,
189}
190
191/// Longest update list a datagram may declare.
192const MAX_COUNT: u64 = 1 << 16;
193
194impl ReplicaTable {
195    /// Apply one datagram. It changes nothing it cannot apply exactly (see
196    /// [`DatagramSummary::skipped`]). An error means the bytes are not a
197    /// datagram; the table then may hold part of it.
198    pub fn apply_datagram(&mut self, datagram: &[u8]) -> Result<DatagramSummary, DecodeError> {
199        let err = |m: &str| DecodeError(m.to_string());
200        let mut b = datagram;
201        let (&version, rest) = b.split_first().ok_or_else(|| err("datagram is empty"))?;
202        b = rest;
203        if version != VERSION {
204            return Err(DecodeError(format!("datagram version {version}")));
205        }
206        let frame = varint::read_u64(&mut b).ok_or_else(|| err("bad frame number"))?;
207        let tick = varint::read_u64(&mut b).ok_or_else(|| err("bad tick"))?;
208        let ack = varint::read_u64(&mut b).ok_or_else(|| err("bad ack"))?;
209        let stream_tick = varint::read_u64(&mut b).ok_or_else(|| err("bad stream tick"))?;
210        if b.len() < 4 {
211            return Err(err("datagram ends early"));
212        }
213        let precision = f32::from_le_bytes(b[..4].try_into().unwrap());
214        b = &b[4..];
215        let parts = varint::read_u64(&mut b).ok_or_else(|| err("bad parts"))?;
216        let n = varint::read_u64(&mut b).ok_or_else(|| err("bad count"))?;
217        if n > MAX_COUNT {
218            return Err(err("count too large"));
219        }
220        let mut summary = DatagramSummary {
221            frame,
222            tick,
223            ack,
224            stream_tick,
225            parts,
226            ..DatagramSummary::default()
227        };
228        // Positions in another precision mean nothing to this table: the
229        // full frame of the new precision is still on its way.
230        let usable = precision == self.precision;
231        let mut last: Option<EntityId> = None;
232        for _ in 0..n {
233            let v = varint::read_u64(&mut b).ok_or_else(|| err("bad id"))?;
234            let id = match last {
235                Some(prev) => prev
236                    .checked_add(v)
237                    .filter(|_| v > 0)
238                    .ok_or_else(|| err("ids do not ascend"))?,
239                None => v,
240            };
241            last = Some(id);
242            if b.len() < 3 {
243                return Err(err("datagram ends early"));
244            }
245            let tag = u16::from_le_bytes([b[0], b[1]]);
246            let m = b[2];
247            b = &b[3..];
248            let mut axes = [None; 3];
249            for (i, bit) in [mask::X, mask::Y, mask::Z].into_iter().enumerate() {
250                if m & bit != 0 {
251                    axes[i] = Some(varint::read_i64(&mut b).ok_or_else(|| err("bad position"))?);
252                }
253            }
254            let mut changes: Vec<(u8, Option<&[u8]>)> = Vec::new();
255            if m & mask::COMPONENTS != 0 {
256                let count = varint::read_u64(&mut b).ok_or_else(|| err("bad component count"))?;
257                if count > 256 {
258                    return Err(err("more than 256 components"));
259                }
260                for _ in 0..count {
261                    let (&c, rest) = b.split_first().ok_or_else(|| err("datagram ends early"))?;
262                    b = rest;
263                    let len =
264                        varint::read_u64(&mut b).ok_or_else(|| err("bad component length"))?;
265                    if len == 0 {
266                        changes.push((c, None));
267                        continue;
268                    }
269                    let len = usize::try_from(len - 1).map_err(|_| err("component too large"))?;
270                    if b.len() < len {
271                        return Err(err("component runs past the end"));
272                    }
273                    let (value, rest) = b.split_at(len);
274                    b = rest;
275                    changes.push((c, Some(value)));
276                }
277            }
278            let current = usable
279                && self.spawn_ticks.get(&id).map(|&t| spawn_tag(t)) == Some(tag)
280                && self.datagram_frames.get(&id).is_none_or(|&f| frame > f);
281            let Some(entity) = self.entities.get_mut(&id).filter(|_| current) else {
282                summary.skipped += 1;
283                continue;
284            };
285            for (i, v) in axes.into_iter().enumerate() {
286                if let Some(v) = v {
287                    entity.q[i] = v;
288                }
289            }
290            for (c, bytes) in changes {
291                match bytes {
292                    Some(v) => {
293                        entity.components.insert(c, v.to_vec());
294                    }
295                    None => {
296                        entity.components.remove(&c);
297                    }
298                }
299            }
300            self.datagram_frames.insert(id, frame);
301            summary.updated.push(id);
302        }
303        if !b.is_empty() {
304            return Err(err("trailing bytes"));
305        }
306        Ok(summary)
307    }
308}
309
310#[cfg(test)]
311mod tests {
312    use super::*;
313    use crate::frame::FrameBuilder;
314
315    /// A table holding entities 1 and 2, spawned at tick 1, precision 0.5.
316    fn table() -> ReplicaTable {
317        let mut t = ReplicaTable::new();
318        let mut f = FrameBuilder::new(true, 0.5);
319        f.spawn(1, [0, 0, 0], std::iter::empty());
320        f.spawn(2, [10, 10, 0], [(7u8, &b"a"[..])].into_iter());
321        t.apply_stream(&f.finish(), 1).unwrap();
322        t
323    }
324
325    /// Updates are (id, spawn tick, axes).
326    fn datagram(
327        frame: u64,
328        precision: f32,
329        updates: &[(EntityId, u64, [Option<i64>; 3])],
330    ) -> Vec<u8> {
331        let mut d = DatagramBuilder::new(frame, 40 + frame, 3, 1, precision);
332        for (id, spawn_tick, axes) in updates {
333            let mut e = Vec::new();
334            encode_entry_into(&mut e, *spawn_tick, *axes, &[]);
335            d.update(*id, &e);
336        }
337        d.finish(1)
338    }
339
340    #[test]
341    fn an_update_sets_absolute_axes_and_components() {
342        let mut t = table();
343        assert_eq!(t.spawn_ticks.get(&2), Some(&1));
344        assert_eq!(t.stream_tick, 1);
345        let mut d = DatagramBuilder::new(1, 41, 9, 1, 0.5);
346        let mut e = Vec::new();
347        encode_entry_into(
348            &mut e,
349            1,
350            [Some(12), None, Some(-3)],
351            &[(7, None), (8, Some(b"xy"))],
352        );
353        let predicted = d.len_with(2, &e);
354        d.update(2, &e);
355        let bytes = d.finish(2);
356        assert!(bytes.len() <= predicted && bytes.len() + 2 >= predicted);
357        let s = t.apply_datagram(&bytes).unwrap();
358        assert_eq!(
359            (s.frame, s.tick, s.ack, s.stream_tick, s.parts),
360            (1, 41, 9, 1, 2)
361        );
362        assert_eq!(s.updated, vec![2]);
363        let e = &t.entities[&2];
364        assert_eq!(e.q, [12, 10, -3]);
365        assert_eq!(e.components.get(&7), None);
366        assert_eq!(e.components.get(&8).map(Vec::as_slice), Some(&b"xy"[..]));
367    }
368
369    #[test]
370    fn an_older_datagram_for_an_entity_changes_nothing() {
371        let mut t = table();
372        t.apply_datagram(&datagram(5, 0.5, &[(1, 1, [Some(50), None, None])]))
373            .unwrap();
374        let s = t
375            .apply_datagram(&datagram(
376                4,
377                0.5,
378                &[
379                    (1, 1, [Some(40), None, None]),
380                    (2, 1, [Some(1), None, None]),
381                ],
382            ))
383            .unwrap();
384        assert_eq!(s.updated, vec![2], "entity 2 had no newer datagram");
385        assert_eq!(s.skipped, 1);
386        assert_eq!(t.entities[&1].q[0], 50);
387    }
388
389    #[test]
390    fn another_spawn_unknown_entity_and_other_precision_are_skipped() {
391        let mut t = table();
392        // Entity 1 despawns and spawns again at tick 2.
393        let mut f = FrameBuilder::new(false, 0.5);
394        f.despawn(1);
395        f.spawn(1, [3, 3, 0], std::iter::empty());
396        t.apply_stream(&f.finish(), 2).unwrap();
397        assert_eq!(t.spawn_ticks.get(&1), Some(&2));
398        let s = t
399            .apply_datagram(&datagram(
400                1,
401                0.5,
402                &[
403                    (1, 1, [Some(99), None, None]),
404                    (9, 1, [Some(1), None, None]),
405                ],
406            ))
407            .unwrap();
408        assert!(s.updated.is_empty());
409        assert_eq!(s.skipped, 2);
410        assert_eq!(t.entities[&1].q[0], 3);
411        let s = t
412            .apply_datagram(&datagram(2, 0.25, &[(1, 2, [Some(99), None, None])]))
413            .unwrap();
414        assert_eq!(s.skipped, 1);
415        assert_eq!(t.entities[&1].q[0], 3);
416    }
417
418    #[test]
419    fn a_respawn_starts_the_entity_over_for_datagram_order() {
420        let mut t = table();
421        t.apply_datagram(&datagram(9, 0.5, &[(1, 1, [Some(9), None, None])]))
422            .unwrap();
423        let mut f = FrameBuilder::new(false, 0.5);
424        f.despawn(1);
425        f.spawn(1, [0, 0, 0], std::iter::empty());
426        t.apply_stream(&f.finish(), 2).unwrap();
427        // A datagram numbered below 9 but for the new spawn applies.
428        let s = t
429            .apply_datagram(&datagram(5, 0.5, &[(1, 2, [Some(4), None, None])]))
430            .unwrap();
431        assert_eq!(s.updated, vec![1]);
432    }
433
434    #[test]
435    fn a_full_frame_starts_over_whatever_the_client_missed() {
436        let mut t = table();
437        // The server dropped frames the client never got, then sent a full
438        // frame at tick 30: every entity counts as spawned then.
439        let mut f = FrameBuilder::new(true, 0.5);
440        f.spawn(1, [5, 5, 0], std::iter::empty());
441        t.apply_stream(&f.finish(), 30).unwrap();
442        assert_eq!(t.stream_tick, 30);
443        assert_eq!(t.spawn_ticks.get(&1), Some(&30));
444        assert_eq!(t.spawn_ticks.get(&2), None);
445        // A datagram built before it names the old spawn: skipped.
446        let s = t
447            .apply_datagram(&datagram(7, 0.5, &[(1, 1, [Some(99), None, None])]))
448            .unwrap();
449        assert_eq!(s.skipped, 1);
450        let s = t
451            .apply_datagram(&datagram(8, 0.5, &[(1, 30, [Some(6), None, None])]))
452            .unwrap();
453        assert_eq!(s.updated, vec![1]);
454        assert_eq!(t.entities[&1].q[0], 6);
455    }
456
457    #[test]
458    fn spawn_ticks_wrap_at_16_bits() {
459        let mut t = ReplicaTable::new();
460        let mut f = FrameBuilder::new(true, 0.5);
461        f.spawn(1, [0, 0, 0], std::iter::empty());
462        t.apply_stream(&f.finish(), 70_000).unwrap();
463        let s = t
464            .apply_datagram(&datagram(1, 0.5, &[(1, 70_000, [Some(1), None, None])]))
465            .unwrap();
466        assert_eq!(s.updated, vec![1]);
467    }
468
469    #[test]
470    fn malformed_datagrams_are_errors_not_panics() {
471        let mut t = table();
472        let good = datagram(1, 0.5, &[(1, 1, [Some(1), Some(2), None])]);
473        for cut in 0..good.len() {
474            let _ = t.clone().apply_datagram(&good[..cut]);
475        }
476        assert!(
477            t.apply_datagram(&[1, 0, 0]).is_err(),
478            "a frame, not a datagram"
479        );
480        let mut extra = good.clone();
481        extra.push(0);
482        assert!(t.apply_datagram(&extra).is_err());
483    }
484}