1use 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;
13pub 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}