use std::sync::{Arc, RwLock};
use std::time::Duration;
use flatland_cli::characters::CharacterSummary;
use flatland_cli::game_session;
use flatland_client_lib::{
needs_asset_sync, sync_assets, AssetSyncKind, AssetSyncOptions,
};
use flatland_client_lib::{
connect, load_client_settings, AutoNavigator, ConnectOptions, GameClient, RemoteSession,
};
use flatland_client_ui::{InputAction, UiKeyEvent};
use flatland_gfx_engine::octant_facing_axes;
use flatland_protocol::ChatChannel;
use tokio::sync::mpsc;
use tokio::time::{self, MissedTickBehavior};
use tracing::info;
use flatland_play_loop::{dispatch_action, dispatch_menu_key, ActionCtx};
use crate::play::GfxPlayArgs;
const RECONNECT_INITIAL: Duration = Duration::from_secs(1);
const RECONNECT_MAX: Duration = Duration::from_secs(30);
const RECONNECT_SESSION_COOLDOWN: Duration = Duration::from_millis(500);
#[derive(Debug, Clone, PartialEq)]
pub enum ConnectionUi {
Connected,
Reconnecting {
attempt: u32,
reason: Option<String>,
next_secs: f32,
last_error: Option<String>,
},
Lost {
reason: Option<String>,
last_error: Option<String>,
},
}
impl Default for ConnectionUi {
fn default() -> Self {
Self::Connected
}
}
#[derive(Debug, Clone, Default)]
pub struct GfxPresentation {
pub auto_nav_facing_axes: Option<(f32, f32)>,
pub connection: ConnectionUi,
}
#[derive(Debug, Clone)]
pub struct GfxSnapshot {
pub game: flatland_client_lib::GameState,
pub presentation: GfxPresentation,
}
pub struct GfxNet {
state: Arc<RwLock<Option<GfxSnapshot>>>,
cmd_tx: mpsc::UnboundedSender<NetCmd>,
}
pub enum NetCmd {
Say(ChatChannel, String),
MenuKey(UiKeyEvent),
KeychainSelect(usize),
RouteEditorClick(flatland_client_lib::RouteEditorClick),
Action(InputAction),
Movement {
forward: f32,
strafe: f32,
vertical: f32,
sprint: bool,
},
AutoNavTo {
x: f32,
y: f32,
},
SetCombatTargetSlot {
slot_index: u8,
target_id: flatland_protocol::EntityId,
label: String,
},
CancelAutoNav,
StopMovement,
Esc,
RouteEditorWaypoint {
x: f32,
y: f32,
z: f32,
},
RouteEditorMapClick {
x: f32,
y: f32,
},
RetryConnect,
Shutdown,
}
impl GfxNet {
pub fn spawn_connecting(
args: GfxPlayArgs,
chosen: Option<CharacterSummary>,
) -> anyhow::Result<(Self, std::sync::mpsc::Receiver<anyhow::Result<()>>)> {
let (ready_tx, ready_rx) = std::sync::mpsc::sync_channel::<anyhow::Result<()>>(1);
let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
let state = Arc::new(RwLock::new(None));
std::thread::Builder::new()
.name("flatland-gfx-net".into())
.spawn({
let state = Arc::clone(&state);
move || {
let rt = match tokio::runtime::Runtime::new() {
Ok(rt) => rt,
Err(err) => {
let _ = ready_tx.send(Err(err.into()));
return;
}
};
if let Err(err) =
rt.block_on(net_thread(args, chosen, cmd_rx, state, ready_tx))
{
tracing::error!(error = %err, "gfx network thread exited");
}
}
})?;
Ok((Self { state, cmd_tx }, ready_rx))
}
#[allow(dead_code)]
pub fn spawn(
args: GfxPlayArgs,
chosen: Option<CharacterSummary>,
) -> anyhow::Result<Self> {
let (net, ready_rx) = Self::spawn_connecting(args, chosen)?;
ready_rx
.recv()
.map_err(|_| anyhow::anyhow!("gfx network thread exited before ready"))??;
Ok(net)
}
pub fn send(&self, cmd: NetCmd) {
let _ = self.cmd_tx.send(cmd);
}
pub fn snapshot(&self) -> Option<GfxSnapshot> {
self.state.read().ok().and_then(|g| g.clone())
}
}
impl Drop for GfxNet {
fn drop(&mut self) {
let _ = self.cmd_tx.send(NetCmd::Shutdown);
}
}
fn update_auto_nav_presentation(
presentation: &mut GfxPresentation,
auto_nav_active: bool,
forward: f32,
strafe: f32,
) {
if auto_nav_active {
if forward.abs() > f32::EPSILON || strafe.abs() > f32::EPSILON {
presentation.auto_nav_facing_axes = Some(octant_facing_axes(forward, strafe));
}
} else {
presentation.auto_nav_facing_axes = None;
}
}
fn sync_state(
state: &Arc<RwLock<Option<GfxSnapshot>>>,
client: &GameClient<RemoteSession>,
presentation: &GfxPresentation,
) {
if let Ok(mut guard) = state.write() {
*guard = Some(GfxSnapshot {
game: client.state.clone(),
presentation: presentation.clone(),
});
}
}
fn auto_reconnect_enabled(no_reconnect: bool) -> bool {
if no_reconnect {
return false;
}
match std::env::var("FLATLAND_AUTO_RECONNECT").as_deref() {
Ok("0" | "false" | "FALSE" | "no" | "NO") => false,
_ => true,
}
}
async fn connect_client(opts: &ConnectOptions) -> anyhow::Result<GameClient<RemoteSession>> {
let character_id = opts.character_id;
let session = connect(opts.clone()).await?;
let mut client = GameClient::new(session);
client.state.character_id = character_id;
client.wait_until_ready().await?;
Ok(client)
}
async fn wait_for_server(
opts: &ConnectOptions,
state: &Arc<RwLock<Option<GfxSnapshot>>>,
cmd_rx: &mut mpsc::UnboundedReceiver<NetCmd>,
disconnect_reason: Option<String>,
) -> anyhow::Result<GameClient<RemoteSession>> {
let mut backoff = RECONNECT_INITIAL;
let mut attempt = 0u32;
loop {
attempt += 1;
let next_secs = backoff.as_secs_f32();
if let Ok(mut guard) = state.write() {
if let Some(s) = guard.as_mut() {
s.presentation.connection = ConnectionUi::Reconnecting {
attempt,
reason: disconnect_reason.clone(),
next_secs,
last_error: None,
};
s.game.connected = false;
s.game.push_log(format!(
"Reconnecting to {} (attempt {attempt})…",
opts.server
));
}
}
match connect_client(opts).await {
Ok(client) => return Ok(client),
Err(err) => {
let err_s = err.to_string();
if let Ok(mut guard) = state.write() {
if let Some(s) = guard.as_mut() {
s.presentation.connection = ConnectionUi::Reconnecting {
attempt,
reason: disconnect_reason.clone(),
next_secs,
last_error: Some(err_s.clone()),
};
s.game.push_log(format!(
"Reconnect failed: {err_s} — retry in {next_secs:.1}s"
));
}
}
}
}
let sleep = tokio::time::sleep(backoff);
tokio::pin!(sleep);
let mut force_retry = false;
tokio::select! {
_ = &mut sleep => {}
cmd = cmd_rx.recv() => {
match cmd {
Some(NetCmd::Shutdown) | None => {
anyhow::bail!("quit while reconnecting");
}
Some(NetCmd::RetryConnect) => {
force_retry = true;
}
Some(_) => {}
}
}
}
if force_retry {
backoff = RECONNECT_INITIAL;
} else {
backoff = (backoff * 2).min(RECONNECT_MAX);
}
}
}
async fn net_thread(
args: GfxPlayArgs,
chosen: Option<CharacterSummary>,
mut cmd_rx: mpsc::UnboundedReceiver<NetCmd>,
state: Arc<RwLock<Option<GfxSnapshot>>>,
ready_tx: std::sync::mpsc::SyncSender<anyhow::Result<()>>,
) -> anyhow::Result<()> {
let auto_reconnect = auto_reconnect_enabled(args.no_reconnect);
let asset_sync_msg = if !args.skip_asset_sync {
match sync_assets(AssetSyncOptions {
quiet: true,
..AssetSyncOptions::default()
})
.await
{
Ok(result) => match result.kind {
AssetSyncKind::Remote => {
info!(publish_rev = result.state.publish_rev, "gfx assets synced");
Some(format!("Assets synced (rev {})", result.state.publish_rev))
}
AssetSyncKind::LocalCacheNoRemote => {
info!(
publish_rev = result.state.publish_rev,
"gfx using cached sprites — remote latest.json not published"
);
Some(format!(
"Using cached sprites (rev {}) — remote bundle not published yet",
result.state.publish_rev
))
}
AssetSyncKind::RepoDevNoRemote => {
info!(
publish_rev = result.state.publish_rev,
"gfx using repo sprites — remote latest.json not published"
);
Some(format!(
"Using repo sprites (rev {}) — run flatland-admin content publish for remote sync",
result.state.publish_rev
))
}
},
Err(err) => {
tracing::warn!(error = %err, "gfx asset sync failed — sprites may be missing");
Some(format!("Asset sync failed: {err}"))
}
}
} else {
None
};
let opts = match &chosen {
Some(character) => game_session::prepare_remote_play_as(character, args.server).await?,
None => game_session::prepare_remote_play(args.name.as_deref(), args.server).await?,
};
let mut client = match connect_client(&opts).await {
Ok(c) => c,
Err(err) => {
let _ = ready_tx.send(Err(err));
return Ok(());
}
};
client
.state
.logs
.retain(|line| !line.starts_with("gfx:") && !line.contains("not wired"));
client
.state
.push_log(format!("Gfx client v{} ready.", env!("CARGO_PKG_VERSION")));
if let Some(msg) = asset_sync_msg {
client.state.push_log(msg);
}
if client.state.publish_rev > 0 && needs_asset_sync(client.state.publish_rev) {
client.state.push_log(format!(
"Server publish rev {} — run: flatland3-gfx assets sync",
client.state.publish_rev
));
}
info!(
entity_id = client.entity_id(),
"gfx connected — ? help, Ctrl+C quit"
);
let mut presentation = GfxPresentation {
connection: ConnectionUi::Connected,
..Default::default()
};
sync_state(&state, &client, &presentation);
ready_tx
.send(Ok(()))
.map_err(|_| anyhow::anyhow!("gfx play loop not waiting for ready"))?;
let keys = load_client_settings().keys;
'sessions: loop {
let mut was_moving = false;
let mut intent = time::interval(Duration::from_millis(100));
intent.set_missed_tick_behavior(MissedTickBehavior::Skip);
intent.tick().await;
let mut drain = time::interval(Duration::from_millis(33));
drain.set_missed_tick_behavior(MissedTickBehavior::Skip);
drain.tick().await;
let mut move_forward = 0.0f32;
let mut move_strafe = 0.0f32;
let mut move_vertical = 0.0f32;
let mut move_sprint = false;
let mut auto_nav: Option<AutoNavigator> = None;
let mut block_active = false;
let mut last_move_err: Option<String> = None;
let mut shutdown = false;
presentation = GfxPresentation::default();
'session: loop {
tokio::select! {
cmd = cmd_rx.recv() => {
let Some(cmd) = cmd else {
shutdown = true;
break 'session;
};
match cmd {
NetCmd::Shutdown => {
shutdown = true;
break 'session;
}
NetCmd::Esc => {
if client.state.show_npc_chat
|| client.state.show_npc_verb_menu
|| client.state.show_shop_menu
|| client.state.bank_panel.is_some()
|| client.state.storage_panel.is_some()
|| (client.state.show_quest_offer
&& (client.state.show_npc_chat
|| client.state.npc_verb_target.is_some()))
{
let _ = client.npc_interaction_back().await;
} else {
let _ = client.back_on_esc();
}
}
NetCmd::Say(channel, text) => {
if let Err(err) = client.say(channel, &text).await {
client.state.push_log(format!("Chat failed: {err}"));
}
}
NetCmd::MenuKey(key) => {
let _ = dispatch_menu_key(&mut client, key).await;
}
NetCmd::KeychainSelect(idx) => {
let n = client.state.keychain_entries().len();
if n > 0 {
client.state.keychain_menu_index = idx.min(n - 1);
}
}
NetCmd::RouteEditorClick(click) => {
client.worker_route_editor_ui_click(click);
}
NetCmd::Action(action) => {
let mut ctx = ActionCtx {
block_active: &mut block_active,
auto_nav: &mut auto_nav,
};
dispatch_action(&mut client, &keys, action, &mut ctx).await;
if auto_nav.is_none() {
presentation.auto_nav_facing_axes = None;
}
}
NetCmd::Movement { forward, strafe, vertical, sprint } => {
move_forward = forward;
move_strafe = strafe;
move_vertical = vertical;
move_sprint = sprint;
}
NetCmd::AutoNavTo { x, y } => {
auto_nav = AutoNavigator::plan(&client.state, x, y);
}
NetCmd::SetCombatTargetSlot {
slot_index,
target_id,
label,
} => {
if let Err(err) = client
.set_combat_target_slot(slot_index, target_id, &label)
.await
{
client.state.push_log(format!("Target: {err}"));
}
}
NetCmd::CancelAutoNav => {
auto_nav = None;
presentation.auto_nav_facing_axes = None;
}
NetCmd::StopMovement => {
move_forward = 0.0;
move_strafe = 0.0;
move_vertical = 0.0;
auto_nav = None;
presentation.auto_nav_facing_axes = None;
if was_moving {
if let Err(err) = client.stop().await {
client.state.push_log(format!("Stop failed: {err}"));
}
was_moving = false;
}
}
NetCmd::RouteEditorWaypoint { x, y, z } => {
client.worker_route_editor_add_waypoint(x, y, z);
}
NetCmd::RouteEditorMapClick { x, y } => {
client.worker_route_editor_map_click(x, y);
}
NetCmd::RetryConnect => {
}
}
client.drain_events();
sync_state(&state, &client, &presentation);
if !client.state.connected {
if let Some(reason) = &client.state.disconnect_reason {
info!(%reason, "gfx session ended");
} else {
info!("gfx session ended (server disconnected)");
}
break 'session;
}
}
_ = drain.tick() => {
client.drain_events();
sync_state(&state, &client, &presentation);
if !client.state.connected {
if let Some(reason) = &client.state.disconnect_reason {
info!(%reason, "gfx session ended");
} else {
info!("gfx session ended (server disconnected)");
}
break 'session;
}
}
_ = intent.tick() => {
let (px, py, pz) = client.state.player_position_with_z();
let (forward, strafe, vertical, sprint) = if let Some(ref mut nav) = auto_nav {
match nav.steer(px, py, pz, &client.state) {
Some(v) => v,
None => {
auto_nav = None;
presentation.auto_nav_facing_axes = None;
client.state.push_log("Arrived at map target.");
(0.0, 0.0, 0.0, false)
}
}
} else {
(move_forward, move_strafe, move_vertical, move_sprint)
};
update_auto_nav_presentation(
&mut presentation,
auto_nav.is_some(),
forward,
strafe,
);
let moving = forward.abs() > f32::EPSILON
|| strafe.abs() > f32::EPSILON
|| vertical.abs() > f32::EPSILON;
if moving {
if let Err(err) = client.move_by(forward, strafe, vertical, sprint).await {
let msg = format!("Move failed: {err}");
if last_move_err.as_deref() != Some(&msg) {
client.state.push_log(msg.clone());
last_move_err = Some(msg);
}
} else {
last_move_err = None;
}
was_moving = true;
} else if was_moving {
if let Err(err) = client.stop().await {
client.state.push_log(format!("Stop failed: {err}"));
}
was_moving = false;
}
client.drain_events();
sync_state(&state, &client, &presentation);
if !client.state.connected {
if let Some(reason) = &client.state.disconnect_reason {
info!(%reason, "gfx session ended");
} else {
info!("gfx session ended (server disconnected)");
}
break 'session;
}
}
}
}
let shutdown_requested = shutdown;
let disconnect_reason = client.state.disconnect_reason.clone();
client.disconnect();
if shutdown_requested || !auto_reconnect {
if let Ok(mut guard) = state.write() {
if let Some(s) = guard.as_mut() {
s.game.connected = false;
s.presentation.connection = ConnectionUi::Lost {
reason: disconnect_reason.clone(),
last_error: None,
};
if !shutdown_requested {
s.game.push_log("Disconnected — reconnect disabled.");
}
}
}
if shutdown_requested {
break 'sessions;
}
loop {
match cmd_rx.recv().await {
Some(NetCmd::Shutdown) | None => break 'sessions,
Some(NetCmd::RetryConnect) => {
info!("gfx manual reconnect requested");
break;
}
Some(_) => {}
}
}
} else {
info!("gfx disconnected — reconnecting");
if let Ok(mut guard) = state.write() {
if let Some(s) = guard.as_mut() {
s.game.connected = false;
s.presentation.connection = ConnectionUi::Reconnecting {
attempt: 0,
reason: disconnect_reason.clone(),
next_secs: RECONNECT_INITIAL.as_secs_f32(),
last_error: None,
};
s.game.push_log("Disconnected — reconnecting…");
}
}
}
drop(client);
let reconnected = loop {
tokio::time::sleep(RECONNECT_SESSION_COOLDOWN).await;
match wait_for_server(&opts, &state, &mut cmd_rx, disconnect_reason.clone()).await {
Ok(new_client) => break Some(new_client),
Err(err) => {
let msg = err.to_string();
tracing::warn!(error = %msg, "gfx reconnect aborted");
if msg.contains("quit while reconnecting") {
break None;
}
if let Ok(mut guard) = state.write() {
if let Some(s) = guard.as_mut() {
s.presentation.connection = ConnectionUi::Lost {
reason: disconnect_reason.clone(),
last_error: Some(msg),
};
}
}
let mut retry = false;
loop {
match cmd_rx.recv().await {
Some(NetCmd::Shutdown) | None => break,
Some(NetCmd::RetryConnect) => {
retry = true;
break;
}
Some(_) => {}
}
}
if !retry {
break None;
}
}
}
};
let Some(new_client) = reconnected else {
break 'sessions;
};
client = new_client;
client.close_overlays();
client
.state
.push_log(format!("Reconnected as entity {}.", client.entity_id()));
info!(entity_id = client.entity_id(), "gfx reconnected");
presentation = GfxPresentation {
connection: ConnectionUi::Connected,
..Default::default()
};
sync_state(&state, &client, &presentation);
}
Ok(())
}