flatland_client_lib/
remote.rs1use std::net::SocketAddr;
2use std::sync::atomic::{AtomicBool, Ordering};
3use std::sync::Arc;
4
5use flatland_protocol::{
6 frame::{read_server_message, write_client_message},
7 ClientMessage, Intent, ServerMessage,
8};
9use tokio::io::AsyncWriteExt;
10use tokio::net::TcpStream;
11use tokio::sync::mpsc;
12
13use crate::session::{PlayConnection, SessionEvent};
14
15pub struct RemoteSession {
16 pub session_id: flatland_protocol::SessionId,
17 pub entity_id: flatland_protocol::EntityId,
18 cmd_tx: mpsc::UnboundedSender<ClientMessage>,
19 event_rx: mpsc::UnboundedReceiver<SessionEvent>,
20 closing: Arc<AtomicBool>,
21}
22
23#[async_trait::async_trait]
24impl PlayConnection for RemoteSession {
25 fn session_id(&self) -> flatland_protocol::SessionId {
26 self.session_id
27 }
28
29 fn entity_id(&self) -> flatland_protocol::EntityId {
30 self.entity_id
31 }
32
33 async fn submit_intent(&self, intent: Intent) -> anyhow::Result<()> {
34 if let Intent::Harvest {
35 entity_id,
36 ref node_id,
37 seq,
38 } = intent
39 {
40 crate::harvest_trace!(
41 entity_id,
42 node_id = %node_id,
43 seq,
44 "remote session enqueue harvest intent"
45 );
46 }
47 self.cmd_tx
48 .send(ClientMessage::Intent(intent))
49 .map_err(|_| anyhow::anyhow!("session closed"))?;
50 Ok(())
51 }
52
53 fn try_next_event(&mut self) -> Option<SessionEvent> {
54 self.event_rx.try_recv().ok()
55 }
56
57 async fn next_event(&mut self) -> Option<SessionEvent> {
58 self.event_rx.recv().await
59 }
60
61 fn disconnect(&self) {
62 self.closing.store(true, Ordering::SeqCst);
63 let _ = self.cmd_tx.send(ClientMessage::Disconnect);
64 }
65}
66
67impl Drop for RemoteSession {
68 fn drop(&mut self) {
69 if !self.closing.load(Ordering::SeqCst) {
70 self.disconnect();
71 }
72 }
73}
74
75fn disconnect_reason_from_io(
76 err: &flatland_protocol::FrameError,
77 closing: &AtomicBool,
78) -> Option<String> {
79 if closing.load(Ordering::SeqCst) {
80 return None;
81 }
82 match err {
83 flatland_protocol::FrameError::Io(io)
84 if io.kind() == std::io::ErrorKind::ConnectionReset
85 || io.kind() == std::io::ErrorKind::BrokenPipe
86 || io.kind() == std::io::ErrorKind::UnexpectedEof =>
87 {
88 Some("server closed the connection".into())
89 }
90 flatland_protocol::FrameError::BadMagic { .. } => {
91 Some("protocol desync — restart client and server".into())
92 }
93 flatland_protocol::FrameError::TooLarge(len) => Some(format!(
94 "frame too large ({len} bytes) — stream likely desynced; restart client and server"
95 )),
96 flatland_protocol::FrameError::Codec(inner) => Some(format!(
97 "server message decode failed: {inner} — client/server likely out of sync; \
98 run `make dev-restart` and restart `make dev-server` + `make dev-control-plane`"
99 )),
100 other => Some(other.to_string()),
101 }
102}
103
104pub async fn connect_tcp(
105 addr: SocketAddr,
106 hello: &flatland_protocol::Hello,
107) -> anyhow::Result<RemoteSession> {
108 let stream = TcpStream::connect(addr).await?;
109 let (mut reader, mut writer) = stream.into_split();
110
111 write_client_message(&mut writer, &ClientMessage::Hello(hello.clone())).await?;
112
113 let welcome = match read_server_message(&mut reader).await {
114 Ok(ServerMessage::Welcome(welcome)) => welcome,
115 Ok(ServerMessage::ConnectRejected { reason }) => {
116 anyhow::bail!("{}", crate::connect_reject_message(&reason));
117 }
118 Ok(other) => anyhow::bail!("expected welcome, got {other:?}"),
119 Err(flatland_protocol::FrameError::Codec(err)) => {
120 anyhow::bail!(
121 "codec error: {err}\n\
122 Client and server are out of sync — restart both after `make dev-restart`:\n\
123 1. make dev-stop\n\
124 2. make dev-control-plane (terminal 1)\n\
125 3. make dev-server (terminal 2)\n\
126 4. make dev-play NAME=…"
127 )
128 }
129 Err(err) => return Err(err.into()),
130 };
131
132 let (cmd_tx, mut cmd_rx) = mpsc::unbounded_channel::<ClientMessage>();
133 let (event_tx, event_rx) = mpsc::unbounded_channel::<SessionEvent>();
134 let closing = Arc::new(AtomicBool::new(false));
135
136 let _ = event_tx.send(SessionEvent::Welcome {
137 session_id: welcome.session_id,
138 entity_id: welcome.entity_id,
139 snapshot: welcome.snapshot,
140 });
141
142 let writer_events = event_tx.clone();
143 let closing_writer = Arc::clone(&closing);
144 tokio::spawn(async move {
145 while let Some(msg) = cmd_rx.recv().await {
146 if let Err(err) = write_client_message(&mut writer, &msg).await {
147 if !closing_writer.load(Ordering::SeqCst) {
148 tracing::warn!(error = %err, "client write failed");
149 let _ = writer_events.send(SessionEvent::Disconnected {
150 reason: Some(err.to_string()),
151 });
152 }
153 break;
154 }
155 if matches!(msg, ClientMessage::Disconnect) {
156 break;
157 }
158 }
159 let _ = writer.shutdown().await;
160 });
161
162 let closing_reader = Arc::clone(&closing);
163 tokio::spawn(async move {
164 loop {
165 match read_server_message(&mut reader).await {
166 Ok(ServerMessage::ContentUpdated(snapshot)) => {
167 if event_tx
168 .send(SessionEvent::ContentUpdated { snapshot })
169 .is_err()
170 {
171 break;
172 }
173 }
174 Ok(ServerMessage::Tick(delta)) => {
175 if event_tx.send(SessionEvent::Tick(delta)).is_err() {
176 break;
177 }
178 }
179 Ok(ServerMessage::IntentAck {
180 entity_id,
181 seq,
182 tick,
183 }) => {
184 if event_tx
185 .send(SessionEvent::IntentAck {
186 entity_id,
187 seq,
188 tick,
189 })
190 .is_err()
191 {
192 break;
193 }
194 }
195 Ok(ServerMessage::Chat(msg)) => {
196 if event_tx.send(SessionEvent::Chat(msg)).is_err() {
197 break;
198 }
199 }
200 Ok(ServerMessage::HarvestResult(result)) => {
201 crate::harvest_trace!(
202 node_id = %result.node_id,
203 template = %result.item_template,
204 quantity = result.quantity,
205 "remote session received harvest result from TCP"
206 );
207 if event_tx.send(SessionEvent::HarvestResult(result)).is_err() {
208 break;
209 }
210 }
211 Ok(ServerMessage::CraftResult(result)) => {
212 if event_tx.send(SessionEvent::CraftResult(result)).is_err() {
213 break;
214 }
215 }
216 Ok(ServerMessage::Death(notice)) => {
217 if event_tx.send(SessionEvent::Death(notice)).is_err() {
218 break;
219 }
220 }
221 Ok(ServerMessage::Interaction(notice)) => {
222 if event_tx.send(SessionEvent::Interaction(notice)).is_err() {
223 break;
224 }
225 }
226 Ok(ServerMessage::ShopOpened(catalog)) => {
227 if event_tx.send(SessionEvent::ShopOpened(catalog)).is_err() {
228 break;
229 }
230 }
231 Ok(ServerMessage::BankOpened(panel)) => {
232 if event_tx.send(SessionEvent::BankOpened(panel)).is_err() {
233 break;
234 }
235 }
236 Ok(ServerMessage::StorageOpened(panel)) => {
237 if event_tx.send(SessionEvent::StorageOpened(panel)).is_err() {
238 break;
239 }
240 }
241 Ok(ServerMessage::MarketOpened(panel)) => {
242 if event_tx.send(SessionEvent::MarketOpened(panel)).is_err() {
243 break;
244 }
245 }
246 Ok(ServerMessage::TradeOpened(panel)) => {
247 if event_tx.send(SessionEvent::TradeOpened(panel)).is_err() {
248 break;
249 }
250 }
251 Ok(ServerMessage::TradeClosed { reason }) => {
252 if event_tx
253 .send(SessionEvent::TradeClosed { reason })
254 .is_err()
255 {
256 break;
257 }
258 }
259 Ok(ServerMessage::NpcTalkOpened(opened)) => {
260 if event_tx.send(SessionEvent::NpcTalkOpened(opened)).is_err() {
261 break;
262 }
263 }
264 Ok(ServerMessage::NpcTalkPending(pending)) => {
265 if event_tx
266 .send(SessionEvent::NpcTalkPending(pending))
267 .is_err()
268 {
269 break;
270 }
271 }
272 Ok(ServerMessage::NpcTalkReply(reply)) => {
273 if event_tx.send(SessionEvent::NpcTalkReply(reply)).is_err() {
274 break;
275 }
276 }
277 Ok(ServerMessage::NpcTalkClosed(closed)) => {
278 if event_tx.send(SessionEvent::NpcTalkClosed(closed)).is_err() {
279 break;
280 }
281 }
282 Ok(ServerMessage::NpcTalkError(err)) => {
283 if event_tx.send(SessionEvent::NpcTalkError(err)).is_err() {
284 break;
285 }
286 }
287 Ok(ServerMessage::UseResult(result)) => {
288 if event_tx.send(SessionEvent::UseResult(result)).is_err() {
289 break;
290 }
291 }
292 Ok(ServerMessage::QuestOffer(offer)) => {
293 if event_tx.send(SessionEvent::QuestOffer(offer)).is_err() {
294 break;
295 }
296 }
297 Ok(ServerMessage::QuestAccepted(notice)) => {
298 if event_tx.send(SessionEvent::QuestAccepted(notice)).is_err() {
299 break;
300 }
301 }
302 Ok(ServerMessage::QuestWithdrawn(notice)) => {
303 if event_tx.send(SessionEvent::QuestWithdrawn(notice)).is_err() {
304 break;
305 }
306 }
307 Ok(ServerMessage::QuestStepCompleted(notice)) => {
308 if event_tx
309 .send(SessionEvent::QuestStepCompleted(notice))
310 .is_err()
311 {
312 break;
313 }
314 }
315 Ok(ServerMessage::QuestCompleted(notice)) => {
316 if event_tx.send(SessionEvent::QuestCompleted(notice)).is_err() {
317 break;
318 }
319 }
320 Ok(ServerMessage::Welcome(welcome)) => {
321 if event_tx
324 .send(SessionEvent::Welcome {
325 session_id: welcome.session_id,
326 entity_id: welcome.entity_id,
327 snapshot: welcome.snapshot,
328 })
329 .is_err()
330 {
331 break;
332 }
333 }
334 Ok(ServerMessage::ConnectRejected { reason }) => {
335 let _ = event_tx.send(SessionEvent::Disconnected {
336 reason: Some(reason),
337 });
338 break;
339 }
340 Err(err) => {
341 let reason = disconnect_reason_from_io(&err, &closing_reader);
342 if let Some(ref r) = reason {
343 tracing::warn!(error = %r, "server read failed");
344 }
345 let _ = event_tx.send(SessionEvent::Disconnected { reason });
346 break;
347 }
348 }
349 }
350 });
351
352 Ok(RemoteSession {
353 session_id: welcome.session_id,
354 entity_id: welcome.entity_id,
355 cmd_tx,
356 event_rx,
357 closing,
358 })
359}