scroll.server/crates/gateway/src/message_manager.rs
WiseDev 2ea283cdda push messages to the client during a battle
the client never receives individual commands in a battle - SectorManager
has receiveSectorState, receiveCompressedSectorState and a heartbeat, and
nothing else. so everything the opponent does has to reach the player
inside server state, which means the server needs to be able to speak
first. it could not: the rpc only ever answered.

the gateway now runs a ticker for the length of a battle, calling a new
battle_tick on the service five times a second and writing whatever it
returns straight to the socket. the simulation stays in the service and
the socket stays in the gateway.

the first rider is emotes. the player's SendBattleEventMessage reaches
the service instead of being logged and dropped, and the bot answers with
a taunt of its own; it also sends one unprompted every twelve to thirty
seconds, drawn from the rows of taunts.csv that TauntMenu marks as
usable. the reply carries the opponent account so it renders on their
side of the arena.
2026-08-23 13:03:07 +03:00

340 lines
13 KiB
Rust

use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use tokio::task::JoinHandle;
use logic::{
message_type, CancelMatchmakeDoneMessage, ClientCapabilitiesMessage, GoHomeMessage,
KeepAliveServerMessage, LoginFailedMessage, LoginMessage, LoginOkMessage, ServerErrorMessage,
};
use service_rpc::{
AccountRef, DeviceInfo, HomeRequestKind, LoginOutcome, Session as AuthSession, WireMessage,
};
use titan::{Incoming, LogicLong, MessagingSender, OutboundMessage};
use crate::backend::Backends;
use crate::config::GatewayConfig;
#[derive(Debug, thiserror::Error)]
pub enum RoutingError {
#[error("transport: {0}")]
Transport(#[from] titan::MessagingError),
#[error("backend: {0}")]
Backend(#[from] service_rpc::RpcError),
#[error("client sent message type {0} before authenticating")]
Unauthenticated(u16),
}
type RoutingResult<T> = Result<T, RoutingError>;
pub const BATTLE_TICK_INTERVAL: Duration = Duration::from_millis(200);
pub struct MessageManager {
peer: SocketAddr,
config: Arc<GatewayConfig>,
backends: Arc<Backends>,
sender: MessagingSender,
account: Option<AccountRef>,
battle_ticker: Option<JoinHandle<()>>,
}
impl Drop for MessageManager {
fn drop(&mut self) {
self.stop_battle_ticker();
}
}
impl MessageManager {
pub fn new(
peer: SocketAddr,
config: Arc<GatewayConfig>,
backends: Arc<Backends>,
sender: MessagingSender,
) -> Self {
Self {
peer,
config,
backends,
sender,
account: None,
battle_ticker: None,
}
}
pub fn account(&self) -> Option<AccountRef> {
self.account
}
pub async fn receive_message(&mut self, incoming: Incoming) -> RoutingResult<()> {
let message_type = incoming.message_type();
if self.account.is_none() && message_type != message_type::LOGIN {
return Err(RoutingError::Unauthenticated(message_type));
}
let Some(message) = incoming.message.as_ref() else {
return Ok(());
};
if !message.direction().accepts_from_client() {
tracing::warn!(
peer = %self.peer,
message = message.name(),
"client sent a server to client message"
);
return Ok(());
}
match message_type {
message_type::LOGIN => {
let login = incoming.downcast::<LoginMessage>().expect("login message");
self.handle_login(login).await
}
message_type::KEEP_ALIVE => {
self.sender
.send_message(&KeepAliveServerMessage::default())
.await?;
Ok(())
}
message_type::CLIENT_CAPABILITIES => {
let capabilities = incoming
.downcast::<ClientCapabilitiesMessage>()
.expect("client capabilities message");
if let Some(account) = self.account {
self.backends
.game
.client_capabilities(account, capabilities.ping_ms)
.await?;
}
Ok(())
}
message_type::GO_HOME => {
let _ = incoming.downcast::<GoHomeMessage>();
self.stop_battle_ticker();
self.push_home(HomeRequestKind::GoHome).await
}
message_type::SEND_BATTLE_EVENT => {
let Some(account) = self.account else {
return Ok(());
};
if let Err(error) = self
.backends
.game
.battle_event(account, incoming.payload)
.await
{
tracing::warn!(peer = %self.peer, %error, "could not deliver the battle event");
}
Ok(())
}
message_type::CANCEL_MATCHMAKE => {
tracing::info!(peer = %self.peer, "client cancelled matchmaking");
self.sender
.send_message(&CancelMatchmakeDoneMessage::default())
.await?;
Ok(())
}
message_type::HOME_LOGIC_STOPPED => {
let Some(account) = self.account else {
return Ok(());
};
match self
.backends
.game
.home_logic_stopped(account, incoming.payload)
.await
{
Ok(replies) => {
self.push_wire_messages(replies).await?;
self.start_battle_ticker();
}
Err(error) => {
tracing::warn!(peer = %self.peer, %error, "could not match the player");
self.sender
.send_message(&CancelMatchmakeDoneMessage::default())
.await?;
}
}
Ok(())
}
message_type::START_MISSION => {
let Some(account) = self.account else {
return Ok(());
};
match self
.backends
.game
.start_mission(account, incoming.payload)
.await
{
Ok(replies) => {
self.push_wire_messages(replies).await?;
self.start_battle_ticker();
}
Err(error) => {
tracing::warn!(peer = %self.peer, %error, "could not start the mission");
self.sender
.send_message(&ServerErrorMessage::new(
"this server could not build the battle",
))
.await?;
}
}
Ok(())
}
message_type::END_CLIENT_TURN => {
if let Some(account) = self.account {
let replies = self
.backends
.game
.end_client_turn(account, incoming.payload)
.await?;
self.push_wire_messages(replies).await?;
}
Ok(())
}
_ => {
tracing::debug!(
peer = %self.peer,
message = message.name(),
"no gateway route for this message"
);
Ok(())
}
}
}
pub async fn on_disconnect(&self) {
if let Some(account) = self.account {
if let Err(error) = self.backends.game.disconnect(account).await {
tracing::debug!(peer = %self.peer, %error, "disconnect notification failed");
}
}
}
async fn handle_login(&mut self, login: &LoginMessage) -> RoutingResult<()> {
let account = AccountRef::new(login.account_id.high, login.account_id.low);
let device = DeviceInfo {
udid: login.udid.clone(),
open_udid: login.open_udid.clone(),
device: login.device.clone(),
os_version: login.os_version.clone(),
android: login.android,
preferred_language: login.preferred_device_language.clone(),
client_major_version: login.client_major_version,
client_minor_version: login.client_minor_version,
client_build: login.client_build,
resource_sha: login.resource_sha.clone(),
};
tracing::info!(
peer = %self.peer,
%account,
version = format!(
"{}.{}.{}",
login.client_major_version, login.client_minor_version, login.client_build
),
device = login.device.as_deref().unwrap_or("unknown"),
"login attempt"
);
match self
.backends
.auth
.login(account, login.pass_token.clone(), device)
.await?
{
LoginOutcome::Rejected {
error_code,
message,
} => {
tracing::info!(
peer = %self.peer,
error_code,
reason = message.as_deref().unwrap_or("unspecified"),
"login rejected"
);
self.sender
.send_message(&LoginFailedMessage {
error_code,
fingerprint_json: None,
redirect_domain: None,
content_url: None,
remaining_seconds: 0,
})
.await?;
Ok(())
}
LoginOutcome::Accepted(session) => {
self.account = Some(session.account);
self.sender.send_message(&self.login_ok(&session)).await?;
self.push_home(HomeRequestKind::Login).await
}
}
}
fn login_ok(&self, session: &AuthSession) -> LoginOkMessage {
let account = LogicLong::new(session.account.high, session.account.low);
LoginOkMessage {
account_id: account,
home_id: account,
pass_token: Some(session.pass_token.clone()),
facebook_id: None,
gamecenter_id: None,
server_major_version: self.config.server_major_version,
server_minor_version: self.config.server_minor_version,
server_build: self.config.server_build,
content_version: self.config.content_version,
server_environment: Some(self.config.server_environment.clone()),
session_count: session.session_count,
play_time_seconds: session.play_time_seconds,
days_since_started_playing: session.days_since_started_playing,
facebook_app_id: self.config.facebook_app_id.clone(),
server_time: Some(session.server_time.clone()),
account_created_date: Some(session.account_created_date.clone()),
startup_cooldown_seconds: self.config.startup_cooldown_seconds,
}
}
async fn push_home(&self, kind: HomeRequestKind) -> RoutingResult<()> {
let Some(account) = self.account else {
return Ok(());
};
let messages = self.backends.game.load_home(account, kind).await?;
self.push_wire_messages(messages).await
}
fn start_battle_ticker(&mut self) {
let Some(account) = self.account else {
return;
};
self.stop_battle_ticker();
let backends = Arc::clone(&self.backends);
let sender = self.sender.clone();
let peer = self.peer;
self.battle_ticker = Some(tokio::spawn(async move {
let mut ticker = tokio::time::interval(BATTLE_TICK_INTERVAL);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
ticker.tick().await;
if sender.is_closed() {
return;
}
match backends.game.battle_tick(account).await {
Ok(messages) => {
for message in messages {
let outbound = OutboundMessage::raw(
message.message_type,
message.message_version,
message.payload,
);
if sender.send(outbound).await.is_err() {
return;
}
}
}
Err(error) => {
tracing::warn!(%peer, %error, "the battle tick failed, stopping the ticker");
return;
}
}
}
}));
}
fn stop_battle_ticker(&mut self) {
if let Some(handle) = self.battle_ticker.take() {
handle.abort();
}
}
async fn push_wire_messages(&self, messages: Vec<WireMessage>) -> RoutingResult<()> {
for message in messages {
self.sender
.send(OutboundMessage::raw(
message.message_type,
message.message_version,
message.payload,
))
.await?;
}
Ok(())
}
}