1use std::collections::BTreeMap;
22
23pub mod datagram;
24pub mod frame;
25pub mod varint;
26
27pub use frame::{DecodeError, FrameSummary, ReplicaEntity, ReplicaTable};
28
29pub type EntityId = u64;
31
32pub type ComponentId = u8;
34
35#[derive(Debug, Clone, PartialEq)]
37pub struct Entity {
38 pub pos: [f32; 3],
39 pub pos_seq: u64,
41 pub components: BTreeMap<ComponentId, Component>,
44 pub seq: u64,
46 pub spawn_seq: u64,
50}
51
52impl Entity {
53 pub fn changed_since(&self, since: u64) -> bool {
55 self.seq > since
56 }
57
58 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
70const ENTITY_COST: usize = 160;
77const COMPONENT_MAP_COST: usize = 400;
80const COMPONENT_COST: usize = 64;
82
83#[derive(Debug)]
85pub struct Replicated {
86 entities: BTreeMap<EntityId, Entity>,
87 seq: u64,
88 log: Option<Vec<u8>>,
89 cost: usize,
92 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
110impl 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
124impl 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#[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 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 pub fn cost(&self) -> usize {
173 self.cost
174 }
175
176 pub fn store_id(&self) -> u64 {
178 self.store_id
179 }
180
181 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 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 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 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 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 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 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 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 pub fn record_changes(&mut self, on: bool) {
334 self.log = on.then(Vec::new);
335 }
336
337 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 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 pub fn apply_changes(&mut self, bytes: &[u8]) -> Result<(), ChangeError> {
370 self.apply_changes_limited(bytes, usize::MAX, usize::MAX)
371 }
372
373 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
457fn 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 r.remove_component(1, 2);
566 assert_eq!(r.cost(), base + 2 * COMPONENT_COST + 3);
567 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 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 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}