Skip to main content

pylon_replication/
lib.rs

1//! Entity replication for Pylon realtime shards.
2//!
3//! A simulation keeps its replicated state in a [`Replicated`] store:
4//! entities with stable u64 ids, a 3D position, and components (bytes in
5//! the game's own encoding, keyed by a u8 id). Every change bumps the
6//! store's change counter, so the server can tell what changed since it
7//! last wrote to a given subscriber without diffing snapshots.
8//!
9//! The server sends each subscriber replication frames ([`frame`]): spawn
10//! (full state) for entities that came into view, update (only changed
11//! axes and components) for entities in view, and despawn for entities
12//! that left view. A client applies them to a [`ReplicaTable`].
13//!
14//! A WebAssembly shard module records its store's changes
15//! ([`Replicated::record_changes`]) and hands them to the host each tick;
16//! the host applies them to its copy ([`Replicated::apply_changes`]).
17//!
18//! This crate has no dependencies, so guests build it for
19//! wasm32-unknown-unknown.
20
21use std::collections::BTreeMap;
22
23pub mod datagram;
24pub mod frame;
25pub mod varint;
26
27pub use frame::{DecodeError, FrameSummary, ReplicaEntity, ReplicaTable};
28
29/// A stable entity id, chosen by the simulation.
30pub type EntityId = u64;
31
32/// A component id: the game decides what each one holds.
33pub type ComponentId = u8;
34
35/// One replicated entity.
36#[derive(Debug, Clone, PartialEq)]
37pub struct Entity {
38    pub pos: [f32; 3],
39    /// Change counter value of the last position change.
40    pub pos_seq: u64,
41    /// Components. A removed component stays as `bytes: None` so readers
42    /// that last saw it can be told it is gone.
43    pub components: BTreeMap<ComponentId, Component>,
44    /// Change counter value of the last change of any kind.
45    pub seq: u64,
46    /// Change counter value of this entity's spawn. An id despawned and
47    /// spawned again gets a new value, so readers can tell a new entity
48    /// from the old one.
49    pub spawn_seq: u64,
50}
51
52impl Entity {
53    /// True when anything changed after change counter value `since`.
54    pub fn changed_since(&self, since: u64) -> bool {
55        self.seq > since
56    }
57
58    /// A component's bytes, if present.
59    pub fn component(&self, id: ComponentId) -> Option<&[u8]> {
60        self.components.get(&id).and_then(|c| c.bytes.as_deref())
61    }
62}
63
64#[derive(Debug, Clone, PartialEq)]
65pub struct Component {
66    pub bytes: Option<Vec<u8>>,
67    pub seq: u64,
68}
69
70// Heap an entity and its component entries take beyond the component
71// bytes, for limits. Measured on 64-bit targets (see the heap test in
72// tests/cost.rs): an entry in the store's map, the first leaf of an
73// entity's component map, one entry in that map, and the allocation of a
74// component's bytes.
75/// An entity's entry in the store's map.
76const ENTITY_COST: usize = 160;
77/// The first node of an entity's component map, allocated with its first
78/// component.
79const COMPONENT_MAP_COST: usize = 400;
80/// A component entry (present or removed), with its bytes' allocation.
81const COMPONENT_COST: usize = 64;
82
83/// The replicated state of a simulation.
84#[derive(Debug)]
85pub struct Replicated {
86    entities: BTreeMap<EntityId, Entity>,
87    seq: u64,
88    log: Option<Vec<u8>>,
89    /// Approximate memory held: a fixed cost per entity and per component
90    /// entry (removed ones too) plus component bytes.
91    cost: usize,
92    /// Unique per store made in this process, so a WebAssembly module can
93    /// tell when the game swapped in a new store.
94    store_id: u64,
95}
96
97impl Default for Replicated {
98    fn default() -> Self {
99        static NEXT: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
100        Self {
101            entities: BTreeMap::new(),
102            seq: 0,
103            log: None,
104            cost: 0,
105            store_id: NEXT.fetch_add(1, std::sync::atomic::Ordering::Relaxed),
106        }
107    }
108}
109
110/// A clone is a different store: it gets its own id and no change log, so
111/// a WebAssembly module that swaps one in (a round reset from a template, a
112/// rollback) sends the host a full dump of it.
113impl Clone for Replicated {
114    fn clone(&self) -> Self {
115        Self {
116            entities: self.entities.clone(),
117            seq: self.seq,
118            cost: self.cost,
119            ..Self::default()
120        }
121    }
122}
123
124/// Stores are equal when they hold the same entities; the id and the log
125/// are bookkeeping.
126impl PartialEq for Replicated {
127    fn eq(&self, other: &Self) -> bool {
128        self.entities == other.entities && self.seq == other.seq
129    }
130}
131
132mod op {
133    pub const SPAWN: u8 = 1;
134    pub const DESPAWN: u8 = 2;
135    pub const POS: u8 = 3;
136    pub const SET: u8 = 4;
137    pub const REMOVE: u8 = 5;
138}
139
140/// Why a change log did not apply.
141#[derive(Debug, Clone, PartialEq, Eq)]
142pub struct ChangeError(pub String);
143
144impl std::fmt::Display for ChangeError {
145    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
146        f.write_str(&self.0)
147    }
148}
149
150impl std::error::Error for ChangeError {}
151
152impl Replicated {
153    pub fn new() -> Self {
154        Self::default()
155    }
156
157    /// The change counter. Every change raises it by one.
158    pub fn seq(&self) -> u64 {
159        self.seq
160    }
161
162    pub fn len(&self) -> usize {
163        self.entities.len()
164    }
165
166    pub fn is_empty(&self) -> bool {
167        self.entities.is_empty()
168    }
169
170    /// Approximate memory the store holds, for limits: a fixed cost per
171    /// entity and per component entry (removed ones too) plus the bytes.
172    pub fn cost(&self) -> usize {
173        self.cost
174    }
175
176    /// This store's process-unique id.
177    pub fn store_id(&self) -> u64 {
178        self.store_id
179    }
180
181    /// Remove every entity but keep the change counter running, so an id
182    /// spawned again afterwards gets a spawn counter no reader has seen and
183    /// reads as a new entity.
184    pub fn clear(&mut self) {
185        self.entities.clear();
186        self.cost = 0;
187        self.bump();
188    }
189
190    pub fn get(&self, id: EntityId) -> Option<&Entity> {
191        self.entities.get(&id)
192    }
193
194    pub fn contains(&self, id: EntityId) -> bool {
195        self.entities.contains_key(&id)
196    }
197
198    /// Entities in id order.
199    pub fn iter(&self) -> impl Iterator<Item = (EntityId, &Entity)> {
200        self.entities.iter().map(|(id, e)| (*id, e))
201    }
202
203    fn bump(&mut self) -> u64 {
204        self.seq += 1;
205        self.seq
206    }
207
208    /// Add an entity. Returns false (and changes nothing) when the id is
209    /// taken.
210    pub fn spawn(&mut self, id: EntityId, pos: [f32; 3]) -> bool {
211        if self.entities.contains_key(&id) {
212            return false;
213        }
214        self.cost += ENTITY_COST;
215        let seq = self.bump();
216        self.entities.insert(
217            id,
218            Entity {
219                pos,
220                pos_seq: seq,
221                components: BTreeMap::new(),
222                seq,
223                spawn_seq: seq,
224            },
225        );
226        if let Some(log) = &mut self.log {
227            log.push(op::SPAWN);
228            varint::write_u64(log, id);
229            write_pos(log, pos);
230        }
231        true
232    }
233
234    /// Remove an entity. Returns false when there was none.
235    pub fn despawn(&mut self, id: EntityId) -> bool {
236        let Some(e) = self.entities.remove(&id) else {
237            return false;
238        };
239        self.cost -= entity_cost(&e);
240        self.bump();
241        if let Some(log) = &mut self.log {
242            log.push(op::DESPAWN);
243            varint::write_u64(log, id);
244        }
245        true
246    }
247
248    /// Move an entity. A position equal to the current one is not a change.
249    pub fn set_pos(&mut self, id: EntityId, pos: [f32; 3]) -> bool {
250        let Some(current) = self.entities.get(&id).map(|e| e.pos) else {
251            return false;
252        };
253        if current.map(f32::to_bits) == pos.map(f32::to_bits) {
254            return true;
255        }
256        let seq = self.bump();
257        let e = self.entities.get_mut(&id).expect("checked above");
258        e.pos = pos;
259        e.pos_seq = seq;
260        e.seq = seq;
261        if let Some(log) = &mut self.log {
262            log.push(op::POS);
263            varint::write_u64(log, id);
264            write_pos(log, pos);
265        }
266        true
267    }
268
269    /// Set a component. Bytes equal to the current ones are not a change.
270    pub fn set_component(&mut self, id: EntityId, component: ComponentId, bytes: &[u8]) -> bool {
271        let Some(e) = self.entities.get(&id) else {
272            return false;
273        };
274        if e.component(component) == Some(bytes) {
275            return true;
276        }
277        let seq = self.bump();
278        let e = self.entities.get_mut(&id).expect("checked above");
279        if e.components.is_empty() {
280            self.cost += COMPONENT_MAP_COST;
281        }
282        let old = e.components.insert(
283            component,
284            Component {
285                bytes: Some(bytes.to_vec()),
286                seq,
287            },
288        );
289        e.seq = seq;
290        self.cost += bytes.len();
291        match old {
292            Some(c) => self.cost -= c.bytes.map_or(0, |b| b.len()),
293            None => self.cost += COMPONENT_COST,
294        }
295        if let Some(log) = &mut self.log {
296            log.push(op::SET);
297            varint::write_u64(log, id);
298            log.push(component);
299            varint::write_u64(log, bytes.len() as u64);
300            log.extend_from_slice(bytes);
301        }
302        true
303    }
304
305    /// Remove a component. Returns false when the entity or component is
306    /// missing.
307    pub fn remove_component(&mut self, id: EntityId, component: ComponentId) -> bool {
308        let present = self
309            .entities
310            .get(&id)
311            .is_some_and(|e| e.component(component).is_some());
312        if !present {
313            return false;
314        }
315        let seq = self.bump();
316        let e = self.entities.get_mut(&id).expect("checked above");
317        let old = e
318            .components
319            .insert(component, Component { bytes: None, seq });
320        e.seq = seq;
321        // The entry stays as a tombstone: only its bytes are freed.
322        self.cost -= old.and_then(|c| c.bytes).map_or(0, |b| b.len());
323        if let Some(log) = &mut self.log {
324            log.push(op::REMOVE);
325            varint::write_u64(log, id);
326            log.push(component);
327        }
328        true
329    }
330
331    /// Start or stop keeping a log of changes for [`Replicated::take_changes`].
332    /// A WebAssembly shard module turns this on.
333    pub fn record_changes(&mut self, on: bool) {
334        self.log = on.then(Vec::new);
335    }
336
337    /// The changes since the last call, as bytes for
338    /// [`Replicated::apply_changes`]. Empty when recording is off.
339    pub fn take_changes(&mut self) -> Vec<u8> {
340        match &mut self.log {
341            Some(log) => std::mem::take(log),
342            None => Vec::new(),
343        }
344    }
345
346    /// Every entity and component as a change log that rebuilds this
347    /// store from empty (spawns, then component sets).
348    pub fn full_changes(&self) -> Vec<u8> {
349        let mut out = Vec::new();
350        for (id, e) in &self.entities {
351            out.push(op::SPAWN);
352            varint::write_u64(&mut out, *id);
353            write_pos(&mut out, e.pos);
354            for (cid, c) in &e.components {
355                if let Some(bytes) = &c.bytes {
356                    out.push(op::SET);
357                    varint::write_u64(&mut out, *id);
358                    out.push(*cid);
359                    varint::write_u64(&mut out, bytes.len() as u64);
360                    out.extend_from_slice(bytes);
361                }
362            }
363        }
364        out
365    }
366
367    /// Apply changes another store recorded. Stops at the first malformed
368    /// or inconsistent change; the changes before it stay applied.
369    pub fn apply_changes(&mut self, bytes: &[u8]) -> Result<(), ChangeError> {
370        self.apply_changes_limited(bytes, usize::MAX, usize::MAX)
371    }
372
373    /// `apply_changes` for a log from an untrusted source: stops with an
374    /// error at the first change that would take the store past
375    /// `max_entities` or `max_cost` (see [`Replicated::cost`]). That change
376    /// and the rest do not apply.
377    pub fn apply_changes_limited(
378        &mut self,
379        mut bytes: &[u8],
380        max_entities: usize,
381        max_cost: usize,
382    ) -> Result<(), ChangeError> {
383        let err = |m: &str| ChangeError(m.to_string());
384        let over = |entities: usize, cost: usize| {
385            ChangeError(format!(
386                "the change takes the store past its limit ({entities} entities, {cost} bytes; at most {max_entities} and {max_cost})"
387            ))
388        };
389        while let Some((&tag, rest)) = bytes.split_first() {
390            bytes = rest;
391            let id = varint::read_u64(&mut bytes).ok_or_else(|| err("truncated entity id"))?;
392            match tag {
393                op::SPAWN => {
394                    let pos = read_pos(&mut bytes).ok_or_else(|| err("truncated position"))?;
395                    let (entities, cost) = (self.entities.len() + 1, self.cost + ENTITY_COST);
396                    if !self.entities.contains_key(&id)
397                        && (entities > max_entities || cost > max_cost)
398                    {
399                        return Err(over(entities, cost));
400                    }
401                    if !self.spawn(id, pos) {
402                        return Err(ChangeError(format!("spawn of existing entity {id}")));
403                    }
404                }
405                op::DESPAWN => {
406                    self.despawn(id);
407                }
408                op::POS => {
409                    let pos = read_pos(&mut bytes).ok_or_else(|| err("truncated position"))?;
410                    if !self.set_pos(id, pos) {
411                        return Err(ChangeError(format!("move of missing entity {id}")));
412                    }
413                }
414                op::SET => {
415                    let (&cid, rest) = bytes
416                        .split_first()
417                        .ok_or_else(|| err("truncated component"))?;
418                    bytes = rest;
419                    let len =
420                        varint::read_u64(&mut bytes).ok_or_else(|| err("truncated length"))?;
421                    let len = usize::try_from(len).map_err(|_| err("component too large"))?;
422                    if bytes.len() < len {
423                        return Err(err("component runs past the end"));
424                    }
425                    let (value, rest) = bytes.split_at(len);
426                    bytes = rest;
427                    let Some(e) = self.entities.get(&id) else {
428                        return Err(ChangeError(format!("component on missing entity {id}")));
429                    };
430                    let added = match e.components.get(&cid) {
431                        Some(c) => len.saturating_sub(c.bytes.as_ref().map_or(0, Vec::len)),
432                        None if e.components.is_empty() => {
433                            COMPONENT_MAP_COST + COMPONENT_COST + len
434                        }
435                        None => COMPONENT_COST + len,
436                    };
437                    let cost = self.cost.saturating_add(added);
438                    if cost > max_cost {
439                        return Err(over(self.entities.len(), cost));
440                    }
441                    self.set_component(id, cid, value);
442                }
443                op::REMOVE => {
444                    let (&cid, rest) = bytes
445                        .split_first()
446                        .ok_or_else(|| err("truncated component"))?;
447                    bytes = rest;
448                    self.remove_component(id, cid);
449                }
450                other => return Err(ChangeError(format!("unknown change tag {other}"))),
451            }
452        }
453        Ok(())
454    }
455}
456
457/// What an entity adds to [`Replicated::cost`].
458fn entity_cost(e: &Entity) -> usize {
459    let map = if e.components.is_empty() {
460        0
461    } else {
462        COMPONENT_MAP_COST
463    };
464    ENTITY_COST
465        + map
466        + e.components
467            .values()
468            .map(|c| COMPONENT_COST + c.bytes.as_ref().map_or(0, Vec::len))
469            .sum::<usize>()
470}
471
472fn write_pos(out: &mut Vec<u8>, pos: [f32; 3]) {
473    for v in pos {
474        out.extend_from_slice(&v.to_le_bytes());
475    }
476}
477
478fn read_pos(bytes: &mut &[u8]) -> Option<[f32; 3]> {
479    if bytes.len() < 12 {
480        return None;
481    }
482    let (head, rest) = bytes.split_at(12);
483    *bytes = rest;
484    let f = |i: usize| f32::from_le_bytes(head[i..i + 4].try_into().unwrap());
485    Some([f(0), f(4), f(8)])
486}
487
488#[cfg(test)]
489mod tests {
490    use super::*;
491
492    #[test]
493    fn changes_bump_the_counter_and_no_ops_do_not() {
494        let mut r = Replicated::new();
495        assert!(r.spawn(1, [0.0, 0.0, 0.0]));
496        assert!(!r.spawn(1, [5.0, 0.0, 0.0]));
497        let s = r.seq();
498        r.set_pos(1, [0.0, 0.0, 0.0]);
499        r.set_component(1, 3, b"x");
500        let after = r.seq();
501        assert_eq!(after, s + 1);
502        r.set_component(1, 3, b"x");
503        assert_eq!(r.seq(), after);
504        assert!(r.get(1).unwrap().changed_since(s));
505        assert!(!r.get(1).unwrap().changed_since(after));
506        assert!(r.remove_component(1, 3));
507        assert_eq!(r.get(1).unwrap().component(3), None);
508        assert!(!r.remove_component(1, 3));
509    }
510
511    #[test]
512    fn a_recorded_log_rebuilds_the_same_entities() {
513        let mut guest = Replicated::new();
514        guest.record_changes(true);
515        guest.spawn(7, [1.0, 2.0, 3.0]);
516        guest.spawn(9, [0.0, 0.0, 0.0]);
517        guest.set_pos(7, [1.5, 2.0, 3.0]);
518        guest.set_component(7, 1, b"hp=10");
519        guest.set_component(9, 2, &[]);
520        guest.remove_component(9, 2);
521        guest.despawn(9);
522        let log = guest.take_changes();
523        assert!(guest.take_changes().is_empty());
524
525        let mut host = Replicated::new();
526        host.apply_changes(&log).unwrap();
527        let summary = |r: &Replicated| {
528            r.iter()
529                .map(|(id, e)| (id, e.pos, e.component(1).map(<[u8]>::to_vec)))
530                .collect::<Vec<_>>()
531        };
532        assert_eq!(summary(&host), summary(&guest));
533    }
534
535    #[test]
536    fn a_full_dump_rebuilds_the_store() {
537        let mut a = Replicated::new();
538        a.spawn(1, [1.0, 2.0, 3.0]);
539        a.spawn(4, [0.0; 3]);
540        a.set_component(1, 2, b"x");
541        a.set_component(4, 2, b"y");
542        a.remove_component(4, 2);
543        let mut b = Replicated::new();
544        b.apply_changes(&a.full_changes()).unwrap();
545        let view = |r: &Replicated| {
546            r.iter()
547                .map(|(id, e)| (id, e.pos, e.component(2).map(<[u8]>::to_vec)))
548                .collect::<Vec<_>>()
549        };
550        assert_eq!(view(&a), view(&b));
551    }
552
553    #[test]
554    fn cost_counts_entities_entries_tombstones_and_bytes() {
555        let mut r = Replicated::new();
556        r.spawn(1, [0.0; 3]);
557        assert_eq!(r.cost(), ENTITY_COST);
558        let base = ENTITY_COST + COMPONENT_MAP_COST;
559        r.set_component(1, 1, &[0; 10]);
560        r.set_component(1, 2, &[0; 5]);
561        assert_eq!(r.cost(), base + 2 * COMPONENT_COST + 15);
562        r.set_component(1, 1, &[0; 3]);
563        assert_eq!(r.cost(), base + 2 * COMPONENT_COST + 8);
564        // A removed component is a tombstone: its entry still costs.
565        r.remove_component(1, 2);
566        assert_eq!(r.cost(), base + 2 * COMPONENT_COST + 3);
567        // Empty components cost too.
568        r.set_component(1, 9, &[]);
569        assert_eq!(r.cost(), base + 3 * COMPONENT_COST + 3);
570        r.despawn(1);
571        assert_eq!(r.cost(), 0);
572    }
573
574    #[test]
575    fn a_limited_apply_stops_at_the_change_that_passes_the_limit() {
576        let mut source = Replicated::new();
577        source.record_changes(true);
578        for id in 0..1000 {
579            source.spawn(id, [0.0; 3]);
580        }
581        let log = source.take_changes();
582        let mut host = Replicated::new();
583        assert!(host.apply_changes_limited(&log, 10, usize::MAX).is_err());
584        // It stopped at the limit, not after the whole log.
585        assert_eq!(host.len(), 10);
586        let mut host = Replicated::new();
587        assert!(host
588            .apply_changes_limited(&log, usize::MAX, 20 * ENTITY_COST)
589            .is_err());
590        assert_eq!(host.len(), 20);
591        assert!(host.cost() <= 20 * ENTITY_COST);
592    }
593
594    #[test]
595    fn a_limited_apply_refuses_one_large_component_before_it_applies() {
596        let mut source = Replicated::new();
597        source.record_changes(true);
598        source.spawn(1, [0.0; 3]);
599        source.set_component(1, 1, &[7; 100]);
600        source.set_component(1, 1, &vec![7; 1 << 20]);
601        let log = source.take_changes();
602        let limit = ENTITY_COST + COMPONENT_MAP_COST + COMPONENT_COST + 1000;
603        let mut host = Replicated::new();
604        assert!(host.apply_changes_limited(&log, usize::MAX, limit).is_err());
605        // The megabyte never landed; the store holds the 100 bytes.
606        assert_eq!(
607            host.get(1).unwrap().component(1).map(<[u8]>::len),
608            Some(100)
609        );
610        assert!(host.cost() <= limit);
611    }
612
613    #[test]
614    fn a_clone_is_a_new_store_without_a_log() {
615        let mut a = Replicated::new();
616        a.record_changes(true);
617        a.spawn(1, [0.0; 3]);
618        let b = a.clone();
619        assert_ne!(a.store_id(), b.store_id());
620        let mut b = b;
621        assert!(b.take_changes().is_empty());
622        assert_eq!(a, b);
623        assert_eq!(a.cost(), b.cost());
624    }
625
626    #[test]
627    fn clear_keeps_the_counter_so_a_respawned_id_is_new() {
628        let mut r = Replicated::new();
629        r.spawn(7, [0.0; 3]);
630        let first = r.get(7).unwrap().spawn_seq;
631        r.clear();
632        r.spawn(7, [0.0; 3]);
633        assert!(r.get(7).unwrap().spawn_seq > first);
634        assert_ne!(Replicated::new().store_id(), Replicated::new().store_id());
635    }
636
637    #[test]
638    fn bad_logs_are_refused_without_panicking() {
639        let mut r = Replicated::new();
640        assert!(r.apply_changes(&[op::SPAWN]).is_err());
641        assert!(r.apply_changes(&[op::SPAWN, 1, 0, 0]).is_err());
642        assert!(r
643            .apply_changes(&[op::POS, 5, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0])
644            .is_err());
645        assert!(r.apply_changes(&[99, 1]).is_err());
646        let mut long = vec![op::SPAWN, 1];
647        long.extend_from_slice(&[0; 12]);
648        long.extend_from_slice(&[op::SET, 1, 4, 0xff, 0xff, 0xff, 0xff, 0x0f]);
649        assert!(r.apply_changes(&long).is_err());
650    }
651}