Skip to main content

flatland_protocol/
frame.rs

1//! Length-prefixed postcard frames for TCP transport.
2//!
3//! Wire layout: `[magic:4][len:4 LE][postcard envelope]`
4//! Magic prevents interpreting payload bytes as a length after stream slip.
5
6use std::marker::Unpin;
7
8use crate::codec::Codec;
9use crate::{ClientMessage, Envelope, PostcardCodec, ServerMessage, PROTOCOL_VERSION};
10use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
11
12const MAX_FRAME_BYTES: usize = 1_048_576;
13/// Frame sentinel — must not appear as a misaligned length prefix.
14pub const FRAME_MAGIC: [u8; 4] = *b"FL3\x01";
15
16#[derive(Debug, thiserror::Error)]
17pub enum FrameError {
18    #[error("frame too large: {0} bytes")]
19    TooLarge(usize),
20    #[error("frame magic mismatch: expected FL3\\x01, got {got:?} (TCP frame read was interrupted — file a bug)")]
21    BadMagic { got: [u8; 4] },
22    #[error("io error: {0}")]
23    Io(#[from] std::io::Error),
24    #[error("codec error: {0}")]
25    Codec(#[from] crate::codec::CodecError),
26    #[error("protocol version mismatch: expected {PROTOCOL_VERSION}, got {0}")]
27    VersionMismatch(u16),
28}
29
30pub async fn write_client_message<W>(
31    writer: &mut W,
32    message: &ClientMessage,
33) -> Result<(), FrameError>
34where
35    W: AsyncWrite + Unpin + Send,
36{
37    write_payload(writer, &Envelope::new(message.clone())).await
38}
39
40pub async fn read_client_message<R>(reader: &mut R) -> Result<ClientMessage, FrameError>
41where
42    R: AsyncRead + Unpin + Send,
43{
44    let envelope: Envelope<ClientMessage> = read_payload(reader).await?;
45    Ok(envelope.payload)
46}
47
48pub async fn write_server_message<W>(
49    writer: &mut W,
50    message: &ServerMessage,
51) -> Result<(), FrameError>
52where
53    W: AsyncWrite + Unpin + Send,
54{
55    write_payload(writer, &Envelope::new(message.clone())).await
56}
57
58pub async fn read_server_message<R>(reader: &mut R) -> Result<ServerMessage, FrameError>
59where
60    R: AsyncRead + Unpin + Send,
61{
62    let envelope: Envelope<ServerMessage> = read_payload(reader).await?;
63    Ok(envelope.payload)
64}
65
66async fn write_payload<W, T>(writer: &mut W, envelope: &Envelope<T>) -> Result<(), FrameError>
67where
68    W: AsyncWrite + Unpin + Send,
69    T: serde::Serialize,
70{
71    let codec = PostcardCodec;
72    let bytes = codec.encode(envelope)?;
73    if bytes.len() > MAX_FRAME_BYTES {
74        return Err(FrameError::TooLarge(bytes.len()));
75    }
76    let len = u32::try_from(bytes.len()).map_err(|_| FrameError::TooLarge(bytes.len()))?;
77    writer.write_all(&FRAME_MAGIC).await?;
78    writer.write_u32_le(len).await?;
79    writer.write_all(&bytes).await?;
80    writer.flush().await?;
81    Ok(())
82}
83
84async fn read_payload<R, T>(reader: &mut R) -> Result<Envelope<T>, FrameError>
85where
86    R: AsyncRead + Unpin + Send,
87    T: serde::de::DeserializeOwned,
88{
89    let mut magic = [0_u8; 4];
90    reader.read_exact(&mut magic).await?;
91    if magic != FRAME_MAGIC {
92        return Err(FrameError::BadMagic { got: magic });
93    }
94    let len = reader.read_u32_le().await? as usize;
95    if len > MAX_FRAME_BYTES {
96        return Err(FrameError::TooLarge(len));
97    }
98    let mut buf = vec![0_u8; len];
99    reader.read_exact(&mut buf).await?;
100    let codec = PostcardCodec;
101    let envelope: Envelope<T> = codec.decode(&buf)?;
102    if envelope.protocol_version != PROTOCOL_VERSION {
103        return Err(FrameError::VersionMismatch(envelope.protocol_version));
104    }
105    Ok(envelope)
106}
107
108#[cfg(test)]
109mod tests {
110    use super::*;
111    use crate::types::EntityState;
112    use crate::{AuthCredential, ClientMessage, Hello, PROTOCOL_VERSION};
113    use uuid::Uuid;
114
115    #[tokio::test]
116    async fn roundtrip_client_message() {
117        let msg = ClientMessage::Hello(Hello {
118            client_name: "test".into(),
119            protocol_version: PROTOCOL_VERSION,
120            auth: Default::default(),
121            character_id: None,
122        });
123        let mut buf = Vec::new();
124        write_client_message(&mut buf, &msg).await.unwrap();
125        let parsed = read_client_message(&mut buf.as_slice()).await.unwrap();
126        assert_eq!(msg, parsed);
127    }
128
129    #[tokio::test]
130    async fn roundtrip_full_world_tick() {
131        use crate::types::{
132            PlayerSkills, PlayerVitals, PrimaryAttributes, ResourceNodeState, ResourceNodeView,
133            ServerMessage, TickDelta, Transform, Velocity2D, WorldClock, WorldCoord,
134        };
135
136        let delta = TickDelta {
137            tick: 42,
138            entities: vec![EntityState {
139                id: 1,
140                label: "Traveler".into(),
141                transform: Transform {
142                    position: WorldCoord::surface(128.0, 128.0),
143                    yaw: 0.0,
144                    velocity: Velocity2D { vx: 0.0, vy: 0.0 },
145                },
146                vitals: Some(PlayerVitals::default()),
147                attributes: Some(PrimaryAttributes::default()),
148                skills: Some(PlayerSkills::default()),
149                inside_building: None,
150                tile_id: None,
151                paperdoll_ref: None,
152                draw_scale: 1.0,
153                presentation_state: None,
154                sprite_mode: None,
155                progression_xp: None,
156                combat_cues: Vec::new(),
157                statuses: Vec::new(),
158            }],
159            resource_nodes: vec![ResourceNodeView {
160                id: "spawn-oak-1".into(),
161                label: "Oak tree".into(),
162                x: 126.0,
163                y: 134.0,
164                z: 0.0,
165                item_template: "oak_log".into(),
166                state: ResourceNodeState::Available,
167                blocking: true,
168                blocking_radius_m: 0.8,
169                harvest_off: false,
170                tile_id: None,
171                yaw: 0.0,
172                pitch: 0.0,
173                roll: 0.0,
174                draw_scale: 1.0,
175                sprite_mode: None,
176                presentation_state: None,
177                growth_progress: None,
178                channel_start_tick: None,
179                channel_end_tick: None,
180                harvest_drop_templates: Vec::new(),
181            }],
182            buildings: Vec::new(),
183            doors: Vec::new(),
184            npcs: Vec::new(),
185            inventory: Vec::new(),
186            blueprints: Vec::new(),
187            building_materials: Vec::new(),
188            world_clock: WorldClock::default(),
189            ground_drops: Vec::new(),
190            placed_containers: Vec::new(),
191            combat: None,
192            interior_map: None,
193            quest_log: Vec::new(),
194            hired_workers: Vec::new(),
195            interactables: Vec::new(),
196            ledger: None,
197            career: None,
198            combat_fx: Vec::new(),
199            ground_hazards: Vec::new(),
200            property_plots: Vec::new(),
201            terrain_overlays: Vec::new(),
202        };
203        let msg = ServerMessage::Tick(delta);
204        let mut buf = Vec::new();
205        write_server_message(&mut buf, &msg).await.unwrap();
206        let parsed = read_server_message(&mut buf.as_slice()).await.unwrap();
207        assert_eq!(msg, parsed);
208    }
209
210    #[tokio::test]
211    async fn roundtrip_session_auth_hello() {
212        let msg = ClientMessage::Hello(Hello {
213            client_name: "traveler".into(),
214            protocol_version: PROTOCOL_VERSION,
215            auth: AuthCredential::Session {
216                token: "flat_sess_test".into(),
217            },
218            character_id: Some(Uuid::parse_str("019f399e-2b4e-7001-87a2-2e99c005e79d").unwrap()),
219        });
220        let mut buf = Vec::new();
221        write_client_message(&mut buf, &msg).await.unwrap();
222        let parsed = read_client_message(&mut buf.as_slice()).await.unwrap();
223        assert_eq!(msg, parsed);
224    }
225
226    #[tokio::test]
227    async fn rejects_garbage_without_length_desync() {
228        let mut buf = b"GET /health HTTP/1.1\r\n".as_slice();
229        let err = read_client_message(&mut buf).await.unwrap_err();
230        assert!(matches!(err, FrameError::BadMagic { .. }));
231    }
232
233    #[tokio::test]
234    async fn multiplexed_frames_stay_aligned() {
235        use crate::types::{TickDelta, Transform, WorldClock, WorldCoord};
236        use crate::{Intent, ServerMessage};
237
238        let (mut client_io, mut server_io) = tokio::io::duplex(64 * 1024);
239
240        for seq in 1..=200u32 {
241            write_client_message(
242                &mut client_io,
243                &ClientMessage::Intent(Intent::Stop { entity_id: 1, seq }),
244            )
245            .await
246            .unwrap();
247            write_server_message(
248                &mut server_io,
249                &ServerMessage::Tick(TickDelta {
250                    tick: seq as u64,
251                    entities: vec![EntityState {
252                        id: 1,
253                        label: "p".into(),
254                        transform: Transform {
255                            position: WorldCoord::surface(128.0, 128.0),
256                            yaw: 0.0,
257                            velocity: crate::types::Velocity2D { vx: 0.0, vy: 0.0 },
258                        },
259                        vitals: None,
260                        attributes: None,
261                        skills: None,
262                        inside_building: None,
263                        tile_id: None,
264                        paperdoll_ref: None,
265                        draw_scale: 1.0,
266                        presentation_state: None,
267                        sprite_mode: None,
268                        progression_xp: None,
269                        combat_cues: Vec::new(),
270                        statuses: Vec::new(),
271                    }],
272                    resource_nodes: vec![],
273                    buildings: vec![],
274                    doors: vec![],
275                    npcs: vec![],
276                    inventory: vec![],
277                    blueprints: vec![],
278                    building_materials: vec![],
279                    world_clock: WorldClock::default(),
280                    ground_drops: vec![],
281                    placed_containers: vec![],
282                    combat: None,
283                    interior_map: None,
284                    quest_log: vec![],
285                    hired_workers: Vec::new(),
286                    interactables: vec![],
287                    ledger: None,
288                    career: None,
289                    combat_fx: Vec::new(),
290                    ground_hazards: Vec::new(),
291                    property_plots: Vec::new(),
292                    terrain_overlays: Vec::new(),
293                }),
294            )
295            .await
296            .unwrap();
297        }
298
299        for seq in 1..=200u32 {
300            let intent = read_client_message(&mut server_io).await.unwrap();
301            assert_eq!(
302                intent,
303                ClientMessage::Intent(Intent::Stop { entity_id: 1, seq })
304            );
305            let ServerMessage::Tick(delta) = read_server_message(&mut client_io).await.unwrap()
306            else {
307                panic!("expected tick");
308            };
309            assert_eq!(delta.tick, seq as u64);
310        }
311    }
312}