use crate::io::ErrorKind; use crate::log::warn; use crate::records::memory::connections::get_client_type_from_memory; use crate::records::memory::enums::ClientType; use crate::rpc::command_maps::RPC_REPLY; use crate::rpc::server::connection_memory_manager::remove_stream_from_memory; use crate::rpc::server::flood_protection::check_request_frequency_with_client_type; use crate::rpc::server::structs::IncomingCommand; use crate::sled::Db; use crate::sleep; use crate::Arc; use crate::Duration; use crate::Mutex; use crate::TcpStream; async fn read_next_command_byte( stream_locked: &Arc>, ) -> Result, String> { // Poll for a command byte with a nonblocking read. This avoids holding // the stream lock across an awaited read while still removing the need // for peek-based readiness checks. let timeout_duration = Duration::from_millis(100); loop { let stream = stream_locked.lock().await; let mut buffer = [0; 1]; match stream.try_read(&mut buffer) { Ok(1) => return Ok(Some(buffer[0])), Ok(0) => { drop(stream); remove_stream_from_memory(stream_locked).await; return Ok(None); } Ok(_) => {} Err(err) if err.kind() == ErrorKind::WouldBlock => { drop(stream); sleep(timeout_duration).await; } Err(err) => { warn!("Dropped stream: {err:?}"); drop(stream); remove_stream_from_memory(stream_locked).await; return Ok(None); } } } } async fn peer_ip(stream_locked: &Arc>) -> String { // Resolve the peer address once per incoming command so scoring and // command handlers can use the same caller identity. let stream = stream_locked.lock().await; stream .peer_addr() .map(|addr| addr.ip().to_string()) .unwrap_or_else(|_| "unknown".into()) } pub async fn next_incoming_command( stream_locked: Arc>, db: &Db, connections_key: &str, wallet_key: &str, ) -> Result, String> { // A disconnected socket returns None so the caller can end the RPC // loop without treating a clean disconnect as a command failure. let Some(command) = read_next_command_byte(&stream_locked).await? else { return Ok(None); }; let ip = peer_ip(&stream_locked).await; // Connection memory is the source of truth for whether this stream // belongs to a node/miner or a wallet-backed RPC client. let client_type = get_client_type_from_memory(connections_key) .await .unwrap_or(ClientType::Miner); // Replies belong to an existing request path, so only new inbound // commands are counted against flood protection. if command != RPC_REPLY { check_request_frequency_with_client_type(db, ip.clone(), client_type, wallet_key).await; } Ok(Some(IncomingCommand { command, ip, client_type, })) }