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(other) => anyhow::bail!("expected welcome, got {other:?}"),
116        Err(flatland_protocol::FrameError::Codec(err)) => {
117            anyhow::bail!(
118                "codec error: {err}\n\
119                 Client and server are out of sync — restart both after `make dev-restart`:\n\
120                   1. make dev-stop\n\
121                   2. make dev-control-plane   (terminal 1)\n\
122                   3. make dev-server          (terminal 2)\n\
123                   4. make dev-play NAME=…"
124            )
125        }
126        Err(err) => return Err(err.into()),
127    };
128
129    let (cmd_tx, mut cmd_rx) = mpsc::unbounded_channel::<ClientMessage>();
130    let (event_tx, event_rx) = mpsc::unbounded_channel::<SessionEvent>();
131    let closing = Arc::new(AtomicBool::new(false));
132
133    let _ = event_tx.send(SessionEvent::Welcome {
134        session_id: welcome.session_id,
135        entity_id: welcome.entity_id,
136        snapshot: welcome.snapshot,
137    });
138
139    let writer_events = event_tx.clone();
140    let closing_writer = Arc::clone(&closing);
141    tokio::spawn(async move {
142        while let Some(msg) = cmd_rx.recv().await {
143            if let Err(err) = write_client_message(&mut writer, &msg).await {
144                if !closing_writer.load(Ordering::SeqCst) {
145                    tracing::warn!(error = %err, "client write failed");
146                    let _ = writer_events.send(SessionEvent::Disconnected {
147                        reason: Some(err.to_string()),
148                    });
149                }
150                break;
151            }
152            if matches!(msg, ClientMessage::Disconnect) {
153                break;
154            }
155        }
156        let _ = writer.shutdown().await;
157    });
158
159    let closing_reader = Arc::clone(&closing);
160    tokio::spawn(async move {
161        loop {
162            match read_server_message(&mut reader).await {
163                Ok(ServerMessage::ContentUpdated(snapshot)) => {
164                    if event_tx
165                        .send(SessionEvent::ContentUpdated { snapshot })
166                        .is_err()
167                    {
168                        break;
169                    }
170                }
171                Ok(ServerMessage::Tick(delta)) => {
172                    if event_tx.send(SessionEvent::Tick(delta)).is_err() {
173                        break;
174                    }
175                }
176                Ok(ServerMessage::IntentAck {
177                    entity_id,
178                    seq,
179                    tick,
180                }) => {
181                    if event_tx
182                        .send(SessionEvent::IntentAck {
183                            entity_id,
184                            seq,
185                            tick,
186                        })
187                        .is_err()
188                    {
189                        break;
190                    }
191                }
192                Ok(ServerMessage::Chat(msg)) => {
193                    if event_tx.send(SessionEvent::Chat(msg)).is_err() {
194                        break;
195                    }
196                }
197                Ok(ServerMessage::HarvestResult(result)) => {
198                    crate::harvest_trace!(
199                        node_id = %result.node_id,
200                        template = %result.item_template,
201                        quantity = result.quantity,
202                        "remote session received harvest result from TCP"
203                    );
204                    if event_tx.send(SessionEvent::HarvestResult(result)).is_err() {
205                        break;
206                    }
207                }
208                Ok(ServerMessage::CraftResult(result)) => {
209                    if event_tx.send(SessionEvent::CraftResult(result)).is_err() {
210                        break;
211                    }
212                }
213                Ok(ServerMessage::Death(notice)) => {
214                    if event_tx.send(SessionEvent::Death(notice)).is_err() {
215                        break;
216                    }
217                }
218                Ok(ServerMessage::Interaction(notice)) => {
219                    if event_tx.send(SessionEvent::Interaction(notice)).is_err() {
220                        break;
221                    }
222                }
223                Ok(ServerMessage::ShopOpened(catalog)) => {
224                    if event_tx.send(SessionEvent::ShopOpened(catalog)).is_err() {
225                        break;
226                    }
227                }
228                Ok(ServerMessage::NpcTalkOpened(opened)) => {
229                    if event_tx.send(SessionEvent::NpcTalkOpened(opened)).is_err() {
230                        break;
231                    }
232                }
233                Ok(ServerMessage::NpcTalkPending(pending)) => {
234                    if event_tx
235                        .send(SessionEvent::NpcTalkPending(pending))
236                        .is_err()
237                    {
238                        break;
239                    }
240                }
241                Ok(ServerMessage::NpcTalkReply(reply)) => {
242                    if event_tx.send(SessionEvent::NpcTalkReply(reply)).is_err() {
243                        break;
244                    }
245                }
246                Ok(ServerMessage::NpcTalkClosed(closed)) => {
247                    if event_tx.send(SessionEvent::NpcTalkClosed(closed)).is_err() {
248                        break;
249                    }
250                }
251                Ok(ServerMessage::NpcTalkError(err)) => {
252                    if event_tx.send(SessionEvent::NpcTalkError(err)).is_err() {
253                        break;
254                    }
255                }
256                Ok(ServerMessage::UseResult(result)) => {
257                    if event_tx.send(SessionEvent::UseResult(result)).is_err() {
258                        break;
259                    }
260                }
261                Ok(ServerMessage::QuestOffer(offer)) => {
262                    if event_tx.send(SessionEvent::QuestOffer(offer)).is_err() {
263                        break;
264                    }
265                }
266                Ok(ServerMessage::QuestAccepted(notice)) => {
267                    if event_tx.send(SessionEvent::QuestAccepted(notice)).is_err() {
268                        break;
269                    }
270                }
271                Ok(ServerMessage::QuestWithdrawn(notice)) => {
272                    if event_tx.send(SessionEvent::QuestWithdrawn(notice)).is_err() {
273                        break;
274                    }
275                }
276                Ok(ServerMessage::QuestStepCompleted(notice)) => {
277                    if event_tx
278                        .send(SessionEvent::QuestStepCompleted(notice))
279                        .is_err()
280                    {
281                        break;
282                    }
283                }
284                Ok(ServerMessage::QuestCompleted(notice)) => {
285                    if event_tx.send(SessionEvent::QuestCompleted(notice)).is_err() {
286                        break;
287                    }
288                }
289                Ok(ServerMessage::Welcome(_)) => {}
290                Err(err) => {
291                    let reason = disconnect_reason_from_io(&err, &closing_reader);
292                    if let Some(ref r) = reason {
293                        tracing::warn!(error = %r, "server read failed");
294                    }
295                    let _ = event_tx.send(SessionEvent::Disconnected { reason });
296                    break;
297                }
298            }
299        }
300    });
301
302    Ok(RemoteSession {
303        session_id: welcome.session_id,
304        entity_id: welcome.entity_id,
305        cmd_tx,
306        event_rx,
307        closing,
308    })
309}