feat(bridge): add TCP transport for Godot client connection
Godot has no Unix socket API, so TCP localhost is required for client-server IPC. TcpBridge implements SimBridge with the same framing protocol as LocalBridge. Includes accept/connect methods and three integration tests over TCP. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,114 @@
|
||||
// TcpBridge - TCP localhost IPC implementation
|
||||
// Implements D-020 subprocess/IPC architecture
|
||||
// Deterministic client-server communication via TCP sockets
|
||||
// Used for Godot client which lacks Unix socket support
|
||||
|
||||
use super::{BridgeError, ObserverSnapshot, PlayerInput, SimBridge};
|
||||
use crate::bridge::framing::{read_framed, write_framed};
|
||||
use std::io::{BufReader, BufWriter};
|
||||
use std::net::{SocketAddr, TcpListener, TcpStream};
|
||||
use std::sync::Mutex;
|
||||
|
||||
/// TcpBridge: TCP transport for client-server IPC
|
||||
/// Same semantics as LocalBridge but over TCP localhost
|
||||
pub struct TcpBridge {
|
||||
reader: Mutex<BufReader<TcpStream>>,
|
||||
writer: Mutex<BufWriter<TcpStream>>,
|
||||
local_addr: SocketAddr,
|
||||
}
|
||||
|
||||
impl TcpBridge {
|
||||
/// Server-side: bind TCP listener and accept one connection.
|
||||
/// Binds to the specified address (e.g., "127.0.0.1:0" for OS-assigned port).
|
||||
/// Returns the bridge with the actual bound address available via local_addr().
|
||||
pub fn accept(addr: &str) -> Result<Self, BridgeError> {
|
||||
tracing::info!("TcpBridge binding to {}", addr);
|
||||
|
||||
let listener = TcpListener::bind(addr)
|
||||
.map_err(|e| BridgeError::Transport(format!("failed to bind TCP socket: {}", e)))?;
|
||||
|
||||
let local_addr = listener
|
||||
.local_addr()
|
||||
.map_err(|e| BridgeError::Transport(format!("failed to get local address: {}", e)))?;
|
||||
|
||||
tracing::info!("TcpBridge listening on {}", local_addr);
|
||||
|
||||
// Accept one connection
|
||||
let (stream, peer_addr) = listener
|
||||
.accept()
|
||||
.map_err(|e| BridgeError::Transport(format!("failed to accept connection: {}", e)))?;
|
||||
|
||||
tracing::info!(
|
||||
"TcpBridge accepted connection from {} on {}",
|
||||
peer_addr,
|
||||
local_addr
|
||||
);
|
||||
|
||||
// Clone stream for reader and writer
|
||||
let reader_stream = stream.try_clone().map_err(|e| {
|
||||
BridgeError::Transport(format!("failed to clone stream for reader: {}", e))
|
||||
})?;
|
||||
|
||||
Ok(Self {
|
||||
reader: Mutex::new(BufReader::new(reader_stream)),
|
||||
writer: Mutex::new(BufWriter::new(stream)),
|
||||
local_addr,
|
||||
})
|
||||
}
|
||||
|
||||
/// Client-side: connect to TCP address (for tests).
|
||||
pub fn connect(addr: &str) -> Result<Self, BridgeError> {
|
||||
tracing::info!("TcpBridge connecting to {}", addr);
|
||||
|
||||
let stream = TcpStream::connect(addr).map_err(|e| {
|
||||
BridgeError::Transport(format!("failed to connect to TCP socket: {}", e))
|
||||
})?;
|
||||
|
||||
let local_addr = stream
|
||||
.local_addr()
|
||||
.map_err(|e| BridgeError::Transport(format!("failed to get local address: {}", e)))?;
|
||||
|
||||
tracing::trace!("TcpBridge connected to {}", addr);
|
||||
|
||||
// Clone stream for reader and writer
|
||||
let reader_stream = stream.try_clone().map_err(|e| {
|
||||
BridgeError::Transport(format!("failed to clone stream for reader: {}", e))
|
||||
})?;
|
||||
|
||||
Ok(Self {
|
||||
reader: Mutex::new(BufReader::new(reader_stream)),
|
||||
writer: Mutex::new(BufWriter::new(stream)),
|
||||
local_addr,
|
||||
})
|
||||
}
|
||||
|
||||
/// Get the local address (useful for OS-assigned port discovery in tests).
|
||||
pub fn local_addr(&self) -> SocketAddr {
|
||||
self.local_addr
|
||||
}
|
||||
}
|
||||
|
||||
impl SimBridge for TcpBridge {
|
||||
fn send_snapshot(&self, snapshot: &ObserverSnapshot) -> Result<(), BridgeError> {
|
||||
let payload = rmp_serde::to_vec_named(snapshot)?;
|
||||
|
||||
let mut writer = self.writer.lock().expect("writer mutex poisoned");
|
||||
write_framed(writer.get_mut(), &payload)?;
|
||||
|
||||
tracing::trace!("sent snapshot: tick={}", snapshot.tick);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn receive_inputs(&self) -> Result<Vec<PlayerInput>, BridgeError> {
|
||||
let mut reader = self.reader.lock().expect("reader mutex poisoned");
|
||||
|
||||
match read_framed(reader.get_mut())? {
|
||||
Some(payload) => {
|
||||
let inputs: Vec<PlayerInput> = rmp_serde::from_slice(&payload)?;
|
||||
tracing::trace!("received {} inputs", inputs.len());
|
||||
Ok(inputs)
|
||||
}
|
||||
None => Err(BridgeError::Transport("client disconnected (EOF)".into())),
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user