use crate::sim::board::{MAX_PLAYERS, PlayerAction, PlayerId};
use crate::sim::direction::Direction;
use std::collections::BTreeMap;
pub const DEFAULT_DELAY: u32 = 3;
pub const PAUSE_LEAD: u32 = 10;
pub const MAX_COMMIT_LEAD: u32 = 30;
pub const HASH_INTERVAL: u32 = 30;
pub fn encode_action(action: PlayerAction) -> [u8; 3] {
match action {
PlayerAction::None => [0, 0, 0],
PlayerAction::Place { x, y, dir } => [x, y, 1 | (dir.id() << 2)],
PlayerAction::Remove { x, y } => [x, y, 2],
}
}
pub fn decode_action(bytes: [u8; 3]) -> PlayerAction {
let (x, y) = (bytes[0], bytes[1]);
match bytes[2] & 0b11 {
1 => PlayerAction::Place {
x,
y,
dir: Direction::from_id((bytes[2] >> 2) & 0b11),
},
2 => PlayerAction::Remove { x, y },
_ => PlayerAction::None,
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct InputMsg {
pub player: PlayerId,
pub frame: u32,
pub action: PlayerAction,
}
pub const INPUT_BYTES: usize = 8;
impl InputMsg {
pub fn encode(self) -> [u8; INPUT_BYTES] {
let f = self.frame.to_le_bytes();
let a = encode_action(self.action);
[self.player, f[0], f[1], f[2], f[3], a[0], a[1], a[2]]
}
pub fn decode(bytes: [u8; INPUT_BYTES]) -> InputMsg {
InputMsg {
player: bytes[0],
frame: u32::from_le_bytes([bytes[1], bytes[2], bytes[3], bytes[4]]),
action: decode_action([bytes[5], bytes[6], bytes[7]]),
}
}
}
pub struct Lockstep {
local: Option<PlayerId>,
players: Vec<PlayerId>,
delay: u32,
frame: u32,
next_commit: u32,
pending: BTreeMap<u32, [Option<PlayerAction>; MAX_PLAYERS]>,
history: Vec<InputMsg>,
pause_at: Option<u32>,
lifted: Option<u32>,
}
fn resend_span(delay: u32) -> u32 {
delay + MAX_COMMIT_LEAD
}
impl Lockstep {
pub fn new(local: PlayerId, players: Vec<PlayerId>, delay: u32) -> Lockstep {
assert!(players.contains(&local));
Lockstep::seated(Some(local), players, delay)
}
pub fn observer(players: Vec<PlayerId>, delay: u32) -> Lockstep {
Lockstep::seated(None, players, delay)
}
fn seated(local: Option<PlayerId>, players: Vec<PlayerId>, delay: u32) -> Lockstep {
assert!(
players
.iter()
.all(|&p| crate::sim::board::seat(p).is_some()),
"player id out of range"
);
let mut session = Lockstep {
local,
players,
delay,
frame: 0,
next_commit: delay,
pending: BTreeMap::new(),
history: Vec::new(),
pause_at: None,
lifted: None,
};
for frame in 0..delay {
let slot = session.slot(frame);
for entry in slot.iter_mut() {
*entry = Some(PlayerAction::None);
}
}
session
}
fn slot(&mut self, frame: u32) -> &mut [Option<PlayerAction>; MAX_PLAYERS] {
let players = &self.players;
self.pending.entry(frame).or_insert_with(|| {
let mut slot = [None; MAX_PLAYERS];
for p in 0..MAX_PLAYERS as u8 {
if !players.contains(&p) {
slot[p as usize] = Some(PlayerAction::None); }
}
slot
})
}
pub fn request_pause(&mut self) -> u32 {
if self.watching() {
return self.pause_at.unwrap_or(u32::MAX);
}
let mut frame = self.next_commit + PAUSE_LEAD;
if let Some(lifted) = self.lifted {
frame = frame.max(lifted + 1);
}
self.receive_pause(frame);
self.pause_at.unwrap_or(frame)
}
pub fn receive_pause(&mut self, frame: u32) {
if frame < self.frame || self.lifted.is_some_and(|lifted| frame <= lifted) {
return;
}
self.pause_at = Some(match self.pause_at {
Some(existing) => existing.min(frame),
None => frame,
});
}
pub fn resume(&mut self) -> u32 {
self.receive_resume(self.pause_at.unwrap_or(0))
}
pub fn receive_resume(&mut self, frame: u32) -> u32 {
let lifted = [self.lifted, self.pause_at, Some(frame)]
.into_iter()
.flatten()
.max()
.unwrap_or(0);
self.lifted = Some(lifted);
self.pause_at = None;
lifted
}
pub fn pause_frame(&self) -> Option<u32> {
self.pause_at
}
pub fn lifted_pause(&self) -> Option<u32> {
self.lifted
}
pub fn paused(&self) -> bool {
self.pause_at.is_some()
}
pub fn frozen(&self) -> bool {
self.pause_at.is_some_and(|at| self.frame >= at)
}
pub fn commit_local(&mut self, action: PlayerAction) -> Option<InputMsg> {
if self.watching() {
return None; }
if self.pause_at.is_some_and(|at| self.next_commit >= at) {
return None;
}
if self.next_commit >= self.frame + self.delay + MAX_COMMIT_LEAD {
return None;
}
let frame = self.next_commit;
self.next_commit += 1;
let local = self.local?;
let action = decode_action(encode_action(action));
debug_assert!(
usize::from(local) < MAX_PLAYERS,
"committing for seat {local}, which is not at the table"
);
self.slot(frame)[local as usize] = Some(action);
let msg = InputMsg {
player: local,
frame,
action,
};
self.history.push(msg);
self.trim_history();
Some(msg)
}
fn trim_history(&mut self) {
let oldest_wanted = self.frame.saturating_sub(resend_span(self.delay));
let keep = self
.history
.iter()
.position(|msg| msg.frame >= oldest_wanted)
.unwrap_or(self.history.len());
if keep > 0 {
self.history.drain(..keep);
}
}
pub fn recent_commits(&self) -> &[InputMsg] {
&self.history
}
pub fn awaiting(&self) -> Vec<PlayerId> {
let Some(slot) = self.pending.get(&self.frame) else {
return self.players.clone();
};
self.players
.iter()
.copied()
.filter(|player| slot[*player as usize].is_none())
.collect()
}
pub fn abandon(&mut self, player: PlayerId, frame: u32) {
debug_assert!(
usize::from(player) < MAX_PLAYERS,
"no such seat to abandon: {player}"
);
self.players.retain(|seated| *seated != player);
for (&at, slot) in self.pending.iter_mut() {
if at >= frame || slot[player as usize].is_none() {
slot[player as usize] = Some(PlayerAction::None);
}
}
}
pub fn receive(&mut self, msg: InputMsg) {
let horizon = self.frame + 2 * resend_span(self.delay);
if msg.frame < self.frame || msg.frame > horizon || !self.players.contains(&msg.player) {
return;
}
debug_assert!(usize::from(msg.player) < MAX_PLAYERS);
let slot = self.slot(msg.frame);
if slot[msg.player as usize].is_none() {
slot[msg.player as usize] = Some(msg.action);
}
}
pub fn advance(&mut self) -> Option<[PlayerAction; MAX_PLAYERS]> {
let slot = self.slot(self.frame);
if slot.iter().any(|a| a.is_none()) {
return None;
}
let actions = std::array::from_fn(|i| slot[i].unwrap_or(PlayerAction::None));
self.pending.remove(&self.frame);
self.frame += 1;
Some(actions)
}
pub fn frame(&self) -> u32 {
self.frame
}
pub fn watching(&self) -> bool {
self.local.is_none()
}
pub fn seat(&self) -> Option<PlayerId> {
self.local
}
pub fn player_count(&self) -> usize {
self.players.len()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sim::board::Board;
use crate::sim::{CrabKind, Handedness, Spawner, TileKind};
#[test]
fn commit_lead_is_bounded_and_recovers() {
let mut session = Lockstep::new(0, vec![0, 1], DEFAULT_DELAY);
let mut accepted = 0;
for _ in 0..100 {
if session.commit_local(PlayerAction::None).is_some() {
accepted += 1;
}
}
assert_eq!(
accepted,
(MAX_COMMIT_LEAD) as usize,
"commits stop at the lead cap while the sim is stalled"
);
for frame in DEFAULT_DELAY..DEFAULT_DELAY + 5 {
session.receive(InputMsg {
player: 1,
frame,
action: PlayerAction::None,
});
}
let mut advanced = 0;
while session.advance().is_some() {
advanced += 1;
}
assert_eq!(advanced, DEFAULT_DELAY as usize + 5);
assert!(session.commit_local(PlayerAction::None).is_some());
}
#[test]
fn a_pause_halts_both_peers_on_one_frame() {
let players = vec![0u8, 1];
let mut a = Lockstep::new(0, players.clone(), DEFAULT_DELAY);
let mut b = Lockstep::new(1, players, DEFAULT_DELAY);
let mut board_a = test_board(3);
let mut board_b = test_board(3);
let mut pause_at = None;
let run = |a: &mut Lockstep, b: &mut Lockstep, ba: &mut Board, bb: &mut Board| {
let msgs: Vec<_> = [
a.commit_local(PlayerAction::None),
b.commit_local(PlayerAction::None),
]
.into_iter()
.flatten()
.collect();
for msg in msgs {
if msg.player == 0 {
b.receive(msg);
} else {
a.receive(msg);
}
}
while let Some(actions) = a.advance() {
ba.tick(&actions);
}
while let Some(actions) = b.advance() {
bb.tick(&actions);
}
};
for step in 0..80 {
if step == 20 {
let frame = a.request_pause();
b.receive_pause(frame);
pause_at = Some(frame);
}
run(&mut a, &mut b, &mut board_a, &mut board_b);
}
let pause_at = pause_at.expect("a pause was called");
assert!(a.frozen() && b.frozen(), "both peers came to rest");
assert_eq!(a.frame(), pause_at, "stopped on the agreed frame");
assert_eq!(b.frame(), pause_at, "and so did the peer");
assert_eq!(board_a.state_hash(), board_b.state_hash());
b.resume();
a.resume();
for _ in 0..40 {
run(&mut a, &mut b, &mut board_a, &mut board_b);
}
assert!(a.frame() > pause_at + 20, "the match ran on");
assert_eq!(a.frame(), b.frame());
assert_eq!(
board_a.state_hash(),
board_b.state_hash(),
"pausing desynced the peers"
);
}
#[test]
fn a_watcher_follows_without_being_waited_for() {
let players = vec![0u8, 1];
let mut a = Lockstep::new(0, players.clone(), DEFAULT_DELAY);
let mut b = Lockstep::new(1, players.clone(), DEFAULT_DELAY);
let mut watcher = Lockstep::observer(players, DEFAULT_DELAY);
assert!(watcher.watching());
assert_eq!(watcher.seat(), None);
assert_eq!(a.seat(), Some(0));
assert!(
watcher
.commit_local(PlayerAction::Remove { x: 1, y: 1 })
.is_none()
);
assert!(watcher.recent_commits().is_empty());
let place = PlayerAction::Place {
x: 2,
y: 3,
dir: Direction::Up,
};
for frame in 0..20u32 {
let from_a = a.commit_local(if frame == 4 {
place
} else {
PlayerAction::None
});
let from_b = b.commit_local(PlayerAction::None);
for msg in [from_a, from_b].into_iter().flatten() {
a.receive(msg);
b.receive(msg);
watcher.receive(msg);
}
assert!(a.advance().is_some(), "player a stalled at {frame}");
assert!(b.advance().is_some(), "player b stalled at {frame}");
}
let mut seen = Vec::new();
while let Some(actions) = watcher.advance() {
seen.push(actions);
}
assert_eq!(
seen.len(),
20 + DEFAULT_DELAY as usize,
"the watcher saw every frame"
);
assert!(
seen.iter().any(|frame| frame[0] == place),
"and the placement among them"
);
watcher.request_pause();
assert!(watcher.pause_frame().is_none(), "a watcher cannot pause");
}
#[test]
fn simultaneous_pauses_settle_on_the_earlier_frame() {
let mut session = Lockstep::new(0, vec![0, 1], DEFAULT_DELAY);
let mine = session.request_pause();
session.receive_pause(mine + 4);
assert_eq!(session.pause_frame(), Some(mine), "later proposal ignored");
session.receive_pause(mine - 4);
assert_eq!(session.pause_frame(), Some(mine - 4), "earlier one wins");
session.receive_pause(mine - 4);
assert_eq!(session.pause_frame(), Some(mine - 4));
session.resume();
assert!(!session.paused() && !session.frozen());
}
#[test]
fn three_player_lockstep_stays_bit_identical() {
let players = vec![0u8, 1, 2];
let mut sessions: Vec<Lockstep> = (0..3u8)
.map(|p| Lockstep::new(p, players.clone(), DEFAULT_DELAY))
.collect();
let mut boards: Vec<Board> = (0..3).map(|_| test_board(4)).collect();
for step in 0u32..300 {
let mut outgoing = Vec::new();
for (i, session) in sessions.iter_mut().enumerate() {
let action = if step % (20 + i as u32 * 7) == 3 {
PlayerAction::Place {
x: (step % 12) as u8,
y: (i as u8 * 2 + 1) % 9,
dir: Direction::Down,
}
} else {
PlayerAction::None
};
outgoing.extend(session.commit_local(action));
}
for msg in outgoing {
for (i, session) in sessions.iter_mut().enumerate() {
if i as u8 != msg.player {
session.receive(msg);
}
}
}
for (session, board) in sessions.iter_mut().zip(&mut boards) {
while let Some(actions) = session.advance() {
board.tick(&actions);
}
}
}
assert!(sessions[0].frame() > 250, "made progress");
assert_eq!(boards[0].state_hash(), boards[1].state_hash());
assert_eq!(boards[1].state_hash(), boards[2].state_hash());
}
fn roundtrip(action: PlayerAction) {
let decoded = decode_action(encode_action(action));
assert_eq!(format!("{action:?}"), format!("{decoded:?}"));
}
#[test]
fn packed_input_round_trips() {
roundtrip(PlayerAction::None);
for dir in [
Direction::Up,
Direction::Right,
Direction::Down,
Direction::Left,
] {
roundtrip(PlayerAction::Place { x: 11, y: 8, dir });
}
roundtrip(PlayerAction::Remove { x: 0, y: 15 });
for x in 16..20u8 {
roundtrip(PlayerAction::Place {
x,
y: 11,
dir: Direction::Right,
});
roundtrip(PlayerAction::Remove { x, y: 12 });
}
let msg = InputMsg {
player: 3,
frame: 123_456,
action: PlayerAction::Place {
x: 7,
y: 2,
dir: Direction::Left,
},
};
assert_eq!(InputMsg::decode(msg.encode()), msg);
}
fn test_board(seed: u64) -> Board {
let mut board = Board::new(12, 9, seed);
board.set_tile(11, 4, TileKind::Castle(0));
board.set_tile(0, 4, TileKind::Castle(1));
board.set_tile(
5,
0,
TileKind::Spawner(Spawner {
dir: Direction::Down,
period: 25,
}),
);
board.spawn_crab(2, 2, Direction::Right, Handedness::Left, CrabKind::Common);
board.spawn_gull(9, 7, Direction::Left);
board
}
#[test]
fn lockstep_peers_stay_bit_identical() {
let players = vec![0u8, 1u8];
let mut a = Lockstep::new(0, players.clone(), DEFAULT_DELAY);
let mut b = Lockstep::new(1, players, DEFAULT_DELAY);
let mut board_a = test_board(9);
let mut board_b = test_board(9);
let (mut queue_ab, mut queue_ba): (Vec<InputMsg>, Vec<InputMsg>) = (vec![], vec![]);
for step in 0u32..600 {
let act_a = if step % 50 == 7 {
PlayerAction::Place {
x: (step % 12) as u8,
y: 3,
dir: Direction::Down,
}
} else {
PlayerAction::None
};
let act_b = if step % 70 == 11 {
PlayerAction::Place {
x: 4,
y: (step % 9) as u8,
dir: Direction::Left,
}
} else {
PlayerAction::None
};
queue_ab.extend(a.commit_local(act_a));
queue_ba.extend(b.commit_local(act_b));
if !(step % 90 < 8) {
while let Some(msg) = queue_ab.pop() {
b.receive(msg);
}
while let Some(msg) = queue_ba.pop() {
a.receive(msg);
}
}
while let Some(actions) = a.advance() {
board_a.tick(&actions);
}
while let Some(actions) = b.advance() {
board_b.tick(&actions);
}
}
while let Some(msg) = queue_ab.pop() {
b.receive(msg);
}
while let Some(msg) = queue_ba.pop() {
a.receive(msg);
}
while let Some(actions) = a.advance() {
board_a.tick(&actions);
}
while let Some(actions) = b.advance() {
board_b.tick(&actions);
}
let min_frame = a.frame().min(b.frame());
assert!(min_frame > 500, "sessions made progress ({min_frame})");
assert_eq!(a.frame(), b.frame());
assert_eq!(
board_a.state_hash(),
board_b.state_hash(),
"lockstep peers diverged"
);
}
}
#[cfg(test)]
mod abandon_tests {
use super::*;
#[test]
fn abandoning_a_player_unsticks_the_frame_they_were_holding() {
let mut session = Lockstep::new(0, vec![0, 1], 0);
session.commit_local(PlayerAction::None);
assert_eq!(session.awaiting(), vec![1], "held up by the one who left");
assert!(session.advance().is_none(), "and going nowhere");
session.abandon(1, 0);
assert!(session.awaiting().is_empty(), "nobody left to wait for");
assert!(session.advance().is_some(), "the round moves again");
}
#[test]
fn an_abandoned_seat_stays_abandoned() {
let mut session = Lockstep::new(0, vec![0, 1], 0);
session.abandon(1, 0);
for frame in 0..30 {
session.commit_local(PlayerAction::None);
assert!(
session.advance().is_some(),
"stalled again at frame {frame}"
);
}
}
#[test]
fn an_abandoned_seat_is_handed_over_empty() {
let mut session = Lockstep::new(0, vec![0, 1], 0);
session.abandon(1, 0);
session.commit_local(PlayerAction::None);
let actions = session.advance().expect("moves");
assert_eq!(actions[1], PlayerAction::None);
}
#[test]
fn the_others_are_still_waited_for() {
let mut session = Lockstep::new(0, vec![0, 1, 2], 0);
session.commit_local(PlayerAction::None);
session.abandon(1, 0);
assert_eq!(session.awaiting(), vec![2], "still owed seat two");
assert!(session.advance().is_none());
session.receive(InputMsg {
player: 2,
frame: 0,
action: PlayerAction::None,
});
assert!(session.advance().is_some());
}
}
#[cfg(test)]
mod loss_tests {
use super::*;
fn run(steps: usize, mut deliver: impl FnMut(usize, PlayerId) -> bool) -> (Lockstep, Lockstep) {
let players = vec![0u8, 1];
let mut a = Lockstep::new(0, players.clone(), DEFAULT_DELAY);
let mut b = Lockstep::new(1, players, DEFAULT_DELAY);
for step in 0..steps {
a.commit_local(PlayerAction::None);
b.commit_local(PlayerAction::None);
let from_a: Vec<InputMsg> = a.recent_commits().to_vec();
let from_b: Vec<InputMsg> = b.recent_commits().to_vec();
if deliver(step, 0) {
for msg in from_a {
b.receive(msg);
}
}
if deliver(step, 1) {
for msg in from_b {
a.receive(msg);
}
}
while a.advance().is_some() {}
while b.advance().is_some() {}
}
(a, b)
}
#[test]
fn a_one_way_loss_burst_does_not_deadlock_the_session() {
let blackout = 20..90;
let (a, b) = run(300, |step, from| from != 0 || !blackout.contains(&step));
assert_eq!(a.frame(), b.frame(), "back in step");
assert!(a.frame() > 250, "and the round ran on ({})", a.frame());
}
#[test]
fn the_resend_tail_stays_within_two_leads() {
let (a, _) = run(200, |_, _| true);
let span = resend_span(DEFAULT_DELAY) as usize;
assert!(
a.recent_commits().len() <= 2 * span,
"{}",
a.recent_commits().len()
);
}
}
#[cfg(test)]
mod pause_echo_tests {
use super::*;
#[test]
fn a_stale_pause_echo_after_the_resume_is_ignored() {
let mut a = Lockstep::new(0, vec![0, 1], DEFAULT_DELAY);
let at = a.request_pause();
while a.frame() < at {
a.commit_local(PlayerAction::None);
a.receive(InputMsg {
player: 1,
frame: a.frame(),
action: PlayerAction::None,
});
assert!(a.advance().is_some());
}
assert!(a.frozen());
a.receive_resume(at);
assert!(!a.paused());
a.receive_pause(at);
assert!(!a.paused(), "an echo of a lifted pause is not a pause");
a.receive_pause(at.saturating_sub(1));
assert!(!a.paused());
let again = a.request_pause();
assert!(again > at);
assert!(a.paused());
}
#[test]
fn a_resume_names_the_pause_it_lifts() {
let mut a = Lockstep::new(0, vec![0, 1], DEFAULT_DELAY);
a.receive_resume(40);
assert_eq!(a.lifted_pause(), Some(40));
a.receive_pause(40);
assert!(!a.paused());
a.receive_pause(41);
assert!(a.paused(), "a later frame is a new pause");
assert_eq!(a.resume(), 41);
}
#[test]
fn a_new_pause_is_always_past_the_lifted_one() {
let mut a = Lockstep::new(0, vec![0, 1], DEFAULT_DELAY);
a.receive_resume(1000);
let at = a.request_pause();
assert!(at > 1000, "{at}");
assert_eq!(a.pause_frame(), Some(at));
}
}
#[cfg(test)]
mod abandon_frame_tests {
use super::*;
#[test]
fn peers_holding_different_inputs_agree_after_the_abandonment() {
let mut host = Lockstep::new(0, vec![0, 1, 2], 0);
let mut peer = Lockstep::new(2, vec![0, 1, 2], 0);
let place = PlayerAction::Place {
x: 1,
y: 1,
dir: Direction::Up,
};
for frame in [1, 2] {
host.receive(InputMsg {
player: 1,
frame,
action: place,
});
}
peer.receive(InputMsg {
player: 1,
frame: 2,
action: place,
});
let at = host.frame();
host.abandon(1, at);
peer.abandon(1, at);
for frame in 0..3 {
for session in [&mut host, &mut peer] {
session.commit_local(PlayerAction::None);
let other = if session.seat() == Some(0) { 2 } else { 0 };
session.receive(InputMsg {
player: other,
frame,
action: PlayerAction::None,
});
}
let from_host = host.advance().expect("host moves");
let from_peer = peer.advance().expect("peer moves");
assert_eq!(from_host[1], PlayerAction::None, "frame {frame}");
assert_eq!(from_host, from_peer, "frame {frame}");
}
}
#[test]
fn a_departed_seat_is_no_longer_heard() {
let mut session = Lockstep::new(0, vec![0, 1], 0);
session.abandon(1, 0);
session.receive(InputMsg {
player: 1,
frame: 5,
action: PlayerAction::Remove { x: 0, y: 0 },
});
for _ in 0..6 {
session.commit_local(PlayerAction::None);
let actions = session.advance().expect("moves");
assert_eq!(actions[1], PlayerAction::None);
}
}
#[test]
fn inputs_beyond_the_horizon_are_dropped() {
let mut session = Lockstep::new(0, vec![0, 1], DEFAULT_DELAY);
let before = session.pending.len();
session.receive(InputMsg {
player: 1,
frame: u32::MAX,
action: PlayerAction::None,
});
session.receive(InputMsg {
player: 1,
frame: 1000,
action: PlayerAction::None,
});
assert_eq!(session.pending.len(), before);
session.receive(InputMsg {
player: 1,
frame: 2 * resend_span(DEFAULT_DELAY),
action: PlayerAction::None,
});
assert_eq!(
session.pending.len(),
before + 1,
"the horizon itself is in"
);
}
}