Skip to main content

flatland_client_lib/
remote.rs

1use 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.send(SessionEvent::TradeClosed { reason }).is_err() {
253                        break;
254                    }
255                }
256                Ok(ServerMessage::NpcTalkOpened(opened)) => {
257                    if event_tx.send(SessionEvent::NpcTalkOpened(opened)).is_err() {
258                        break;
259                    }
260                }
261                Ok(ServerMessage::NpcTalkPending(pending)) => {
262                    if event_tx
263                        .send(SessionEvent::NpcTalkPending(pending))
264                        .is_err()
265                    {
266                        break;
267                    }
268                }
269                Ok(ServerMessage::NpcTalkReply(reply)) => {
270                    if event_tx.send(SessionEvent::NpcTalkReply(reply)).is_err() {
271                        break;
272                    }
273                }
274                Ok(ServerMessage::NpcTalkClosed(closed)) => {
275                    if event_tx.send(SessionEvent::NpcTalkClosed(closed)).is_err() {
276                        break;
277                    }
278                }
279                Ok(ServerMessage::NpcTalkError(err)) => {
280                    if event_tx.send(SessionEvent::NpcTalkError(err)).is_err() {
281                        break;
282                    }
283                }
284                Ok(ServerMessage::UseResult(result)) => {
285                    if event_tx.send(SessionEvent::UseResult(result)).is_err() {
286                        break;
287                    }
288                }
289                Ok(ServerMessage::QuestOffer(offer)) => {
290                    if event_tx.send(SessionEvent::QuestOffer(offer)).is_err() {
291                        break;
292                    }
293                }
294                Ok(ServerMessage::QuestAccepted(notice)) => {
295                    if event_tx.send(SessionEvent::QuestAccepted(notice)).is_err() {
296                        break;
297                    }
298                }
299                Ok(ServerMessage::QuestWithdrawn(notice)) => {
300                    if event_tx.send(SessionEvent::QuestWithdrawn(notice)).is_err() {
301                        break;
302                    }
303                }
304                Ok(ServerMessage::QuestStepCompleted(notice)) => {
305                    if event_tx
306                        .send(SessionEvent::QuestStepCompleted(notice))
307                        .is_err()
308                    {
309                        break;
310                    }
311                }
312                Ok(ServerMessage::QuestCompleted(notice)) => {
313                    if event_tx.send(SessionEvent::QuestCompleted(notice)).is_err() {
314                        break;
315                    }
316                }
317                Ok(ServerMessage::QuestCatalogUpdated(update)) => {
318                    if event_tx
319                        .send(SessionEvent::QuestCatalogUpdated(update))
320                        .is_err()
321                    {
322                        break;
323                    }
324                }
325                Ok(ServerMessage::Welcome(welcome)) => {
326                    // Mid-session Welcome = region handoff / worker reconnect.
327                    // Must update entity_id + world bounds or the client freezes on the old shard.
328                    if event_tx
329                        .send(SessionEvent::Welcome {
330                            session_id: welcome.session_id,
331                            entity_id: welcome.entity_id,
332                            snapshot: welcome.snapshot,
333                        })
334                        .is_err()
335                    {
336                        break;
337                    }
338                }
339                Ok(ServerMessage::ConnectRejected { reason }) => {
340                    let _ = event_tx.send(SessionEvent::Disconnected {
341                        reason: Some(reason),
342                    });
343                    break;
344                }
345                Err(err) => {
346                    let reason = disconnect_reason_from_io(&err, &closing_reader);
347                    if let Some(ref r) = reason {
348                        tracing::warn!(error = %r, "server read failed");
349                    }
350                    let _ = event_tx.send(SessionEvent::Disconnected { reason });
351                    break;
352                }
353            }
354        }
355    });
356
357    Ok(RemoteSession {
358        session_id: welcome.session_id,
359        entity_id: welcome.entity_id,
360        cmd_tx,
361        event_rx,
362        closing,
363    })
364}