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            property_plots: Vec::new(),
200            terrain_overlays: Vec::new(),
201        };
202        let msg = ServerMessage::Tick(delta);
203        let mut buf = Vec::new();
204        write_server_message(&mut buf, &msg).await.unwrap();
205        let parsed = read_server_message(&mut buf.as_slice()).await.unwrap();
206        assert_eq!(msg, parsed);
207    }
208
209    #[tokio::test]
210    async fn roundtrip_session_auth_hello() {
211        let msg = ClientMessage::Hello(Hello {
212            client_name: "traveler".into(),
213            protocol_version: PROTOCOL_VERSION,
214            auth: AuthCredential::Session {
215                token: "flat_sess_test".into(),
216            },
217            character_id: Some(Uuid::parse_str("019f399e-2b4e-7001-87a2-2e99c005e79d").unwrap()),
218        });
219        let mut buf = Vec::new();
220        write_client_message(&mut buf, &msg).await.unwrap();
221        let parsed = read_client_message(&mut buf.as_slice()).await.unwrap();
222        assert_eq!(msg, parsed);
223    }
224
225    #[tokio::test]
226    async fn rejects_garbage_without_length_desync() {
227        let mut buf = b"GET /health HTTP/1.1\r\n".as_slice();
228        let err = read_client_message(&mut buf).await.unwrap_err();
229        assert!(matches!(err, FrameError::BadMagic { .. }));
230    }
231
232    #[tokio::test]
233    async fn multiplexed_frames_stay_aligned() {
234        use crate::types::{TickDelta, Transform, WorldClock, WorldCoord};
235        use crate::{Intent, ServerMessage};
236
237        let (mut client_io, mut server_io) = tokio::io::duplex(64 * 1024);
238
239        for seq in 1..=200u32 {
240            write_client_message(
241                &mut client_io,
242                &ClientMessage::Intent(Intent::Stop { entity_id: 1, seq }),
243            )
244            .await
245            .unwrap();
246            write_server_message(
247                &mut server_io,
248                &ServerMessage::Tick(TickDelta {
249                    tick: seq as u64,
250                    entities: vec![EntityState {
251                        id: 1,
252                        label: "p".into(),
253                        transform: Transform {
254                            position: WorldCoord::surface(128.0, 128.0),
255                            yaw: 0.0,
256                            velocity: crate::types::Velocity2D { vx: 0.0, vy: 0.0 },
257                        },
258                        vitals: None,
259                        attributes: None,
260                        skills: None,
261                        inside_building: None,
262                        tile_id: None,
263                        paperdoll_ref: None,
264                        draw_scale: 1.0,
265                        presentation_state: None,
266                        sprite_mode: None,
267                        progression_xp: None,
268                        combat_cues: Vec::new(),
269                statuses: Vec::new(),
270                    }],
271                    resource_nodes: vec![],
272                    buildings: vec![],
273                    doors: vec![],
274                    npcs: vec![],
275                    inventory: vec![],
276                    blueprints: vec![],
277            building_materials: vec![],
278                    world_clock: WorldClock::default(),
279                    ground_drops: vec![],
280                    placed_containers: vec![],
281                    combat: None,
282                    interior_map: None,
283                    quest_log: vec![],
284                    hired_workers: Vec::new(),
285                    interactables: vec![],
286                    ledger: None,
287                    career: None,
288                    combat_fx: Vec::new(),
289                    property_plots: Vec::new(),
290                    terrain_overlays: Vec::new(),
291                }),
292            )
293            .await
294            .unwrap();
295        }
296
297        for seq in 1..=200u32 {
298            let intent = read_client_message(&mut server_io).await.unwrap();
299            assert_eq!(
300                intent,
301                ClientMessage::Intent(Intent::Stop { entity_id: 1, seq })
302            );
303            let ServerMessage::Tick(delta) = read_server_message(&mut client_io).await.unwrap()
304            else {
305                panic!("expected tick");
306            };
307            assert_eq!(delta.tick, seq as u64);
308        }
309    }
310}