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(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}