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
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                    // Mid-session Welcome = region handoff / worker reconnect.
322                    // Must update entity_id + world bounds or the client freezes on the old shard.
323                    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}