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 presentation_state: None,
153 sprite_mode: None,
154 progression_xp: None,
155 combat_cues: Vec::new(),
156 statuses: Vec::new(),
157 }],
158 resource_nodes: vec![ResourceNodeView {
159 id: "spawn-oak-1".into(),
160 label: "Oak tree".into(),
161 x: 126.0,
162 y: 134.0,
163 z: 0.0,
164 item_template: "oak_log".into(),
165 state: ResourceNodeState::Available,
166 blocking: true,
167 blocking_radius_m: 0.8,
168 harvest_off: false,
169 tile_id: None,
170 yaw: 0.0,
171 pitch: 0.0,
172 roll: 0.0,
173 draw_scale: 1.0,
174 sprite_mode: None,
175 presentation_state: None,
176 growth_progress: None,
177 channel_start_tick: None,
178 channel_end_tick: None,
179 harvest_drop_templates: Vec::new(),
180 }],
181 buildings: Vec::new(),
182 doors: Vec::new(),
183 npcs: Vec::new(),
184 inventory: Vec::new(),
185 blueprints: Vec::new(),
186 building_materials: Vec::new(),
187 world_clock: WorldClock::default(),
188 ground_drops: Vec::new(),
189 placed_containers: Vec::new(),
190 combat: None,
191 interior_map: None,
192 quest_log: Vec::new(),
193 hired_workers: Vec::new(),
194 interactables: Vec::new(),
195 ledger: None,
196 career: None,
197 combat_fx: Vec::new(),
198 property_plots: Vec::new(),
199 terrain_overlays: Vec::new(),
200 };
201 let msg = ServerMessage::Tick(delta);
202 let mut buf = Vec::new();
203 write_server_message(&mut buf, &msg).await.unwrap();
204 let parsed = read_server_message(&mut buf.as_slice()).await.unwrap();
205 assert_eq!(msg, parsed);
206 }
207
208 #[tokio::test]
209 async fn roundtrip_session_auth_hello() {
210 let msg = ClientMessage::Hello(Hello {
211 client_name: "traveler".into(),
212 protocol_version: PROTOCOL_VERSION,
213 auth: AuthCredential::Session {
214 token: "flat_sess_test".into(),
215 },
216 character_id: Some(Uuid::parse_str("019f399e-2b4e-7001-87a2-2e99c005e79d").unwrap()),
217 });
218 let mut buf = Vec::new();
219 write_client_message(&mut buf, &msg).await.unwrap();
220 let parsed = read_client_message(&mut buf.as_slice()).await.unwrap();
221 assert_eq!(msg, parsed);
222 }
223
224 #[tokio::test]
225 async fn rejects_garbage_without_length_desync() {
226 let mut buf = b"GET /health HTTP/1.1\r\n".as_slice();
227 let err = read_client_message(&mut buf).await.unwrap_err();
228 assert!(matches!(err, FrameError::BadMagic { .. }));
229 }
230
231 #[tokio::test]
232 async fn multiplexed_frames_stay_aligned() {
233 use crate::types::{TickDelta, Transform, WorldClock, WorldCoord};
234 use crate::{Intent, ServerMessage};
235
236 let (mut client_io, mut server_io) = tokio::io::duplex(64 * 1024);
237
238 for seq in 1..=200u32 {
239 write_client_message(
240 &mut client_io,
241 &ClientMessage::Intent(Intent::Stop { entity_id: 1, seq }),
242 )
243 .await
244 .unwrap();
245 write_server_message(
246 &mut server_io,
247 &ServerMessage::Tick(TickDelta {
248 tick: seq as u64,
249 entities: vec![EntityState {
250 id: 1,
251 label: "p".into(),
252 transform: Transform {
253 position: WorldCoord::surface(128.0, 128.0),
254 yaw: 0.0,
255 velocity: crate::types::Velocity2D { vx: 0.0, vy: 0.0 },
256 },
257 vitals: None,
258 attributes: None,
259 skills: None,
260 inside_building: None,
261 tile_id: None,
262 paperdoll_ref: None,
263 presentation_state: None,
264 sprite_mode: None,
265 progression_xp: None,
266 combat_cues: Vec::new(),
267 statuses: Vec::new(),
268 }],
269 resource_nodes: vec![],
270 buildings: vec![],
271 doors: vec![],
272 npcs: vec![],
273 inventory: vec![],
274 blueprints: vec![],
275 building_materials: vec![],
276 world_clock: WorldClock::default(),
277 ground_drops: vec![],
278 placed_containers: vec![],
279 combat: None,
280 interior_map: None,
281 quest_log: vec![],
282 hired_workers: Vec::new(),
283 interactables: vec![],
284 ledger: None,
285 career: None,
286 combat_fx: Vec::new(),
287 property_plots: Vec::new(),
288 terrain_overlays: Vec::new(),
289 }),
290 )
291 .await
292 .unwrap();
293 }
294
295 for seq in 1..=200u32 {
296 let intent = read_client_message(&mut server_io).await.unwrap();
297 assert_eq!(
298 intent,
299 ClientMessage::Intent(Intent::Stop { entity_id: 1, seq })
300 );
301 let ServerMessage::Tick(delta) = read_server_message(&mut client_io).await.unwrap()
302 else {
303 panic!("expected tick");
304 };
305 assert_eq!(delta.tick, seq as u64);
306 }
307 }
308}