Files
settled-reach/server/src/bridge/local.rs
T
jpmschweitzerandClaude Opus 4.6 6ed8d11502 feat(simulation): Sprint 19 — 7 server systems
Protocol handshake (#555): HandshakeMessage as first IPC frame,
HandshakeState resource, forward-compatible input handling.

State serialization (#96): serialize_npc_to_frozen/deserialize with
full D-024 axis coverage (10 new optional fields on NpcSaveState).

Scope tags (#98): ScopeTagKind enum, ScopePinned marker, automatic
assignment from KnowledgeGraph and RelationshipGraph.

Timestamp eviction (#97): LastInteractionTick, SimSpacePressure,
BinaryHeap LRU eviction respecting ScopePinned entities.

Save/load (#553): save_to_file/load_from_file via MessagePack,
SaveGame/LoadGame IPC commands, SaveLoadResultWire on snapshot.

Test infrastructure (#200): Layer 3 integration test entry point,
three-layer architecture documented per D-030.

Information boundary tests (#272): 4 negative tests proving no
passive KG leakage, LOS fog holds, tier boundary holds, save
isolation per NPC.

1063 tests passing, 0 failures.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-02-25 12:13:03 +01:00

153 lines
5.3 KiB
Rust

// LocalBridge - Unix socket IPC implementation
// Implements D-020 subprocess/IPC architecture
// Deterministic client-server communication via Unix domain sockets
use super::{BridgeError, ObserverSnapshot, PlayerInput, SimBridge};
use crate::bridge::framing::{read_framed, write_framed};
use std::fs;
use std::io::{BufReader, BufWriter};
use std::os::unix::net::{UnixListener, UnixStream};
use std::path::{Path, PathBuf};
use std::sync::Mutex;
/// LocalBridge: Unix socket transport for client-server IPC
pub struct LocalBridge {
reader: Mutex<BufReader<UnixStream>>,
writer: Mutex<BufWriter<UnixStream>>,
socket_path: PathBuf,
}
impl LocalBridge {
/// Server-side: create a Unix socket listener and accept one connection.
/// Removes any stale socket file before binding.
pub fn accept(path: &Path) -> Result<Self, BridgeError> {
// Remove stale socket if it exists
if path.exists() {
fs::remove_file(path).map_err(|e| {
BridgeError::Transport(format!("failed to remove stale socket: {}", e))
})?;
}
tracing::info!("LocalBridge listening on {:?}", path);
let listener = UnixListener::bind(path)
.map_err(|e| BridgeError::Transport(format!("failed to bind Unix socket: {}", e)))?;
// Accept one connection
let (stream, _addr) = listener
.accept()
.map_err(|e| BridgeError::Transport(format!("failed to accept connection: {}", e)))?;
tracing::info!("LocalBridge accepted connection on {:?}", path);
// 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)),
socket_path: path.to_path_buf(),
})
}
/// Client-side: connect to an existing Unix socket.
pub fn connect(path: &Path) -> Result<Self, BridgeError> {
tracing::info!("LocalBridge connecting to {:?}", path);
let stream = UnixStream::connect(path).map_err(|e| {
BridgeError::Transport(format!("failed to connect to Unix socket: {}", e))
})?;
tracing::trace!("LocalBridge connected to {:?}", path);
// 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)),
socket_path: path.to_path_buf(),
})
}
pub fn socket_path(&self) -> &Path {
&self.socket_path
}
}
impl SimBridge for LocalBridge {
fn send_handshake(&self) -> Result<(), BridgeError> {
use super::types::{HandshakeMessage, PROTOCOL_VERSION};
let msg = HandshakeMessage {
protocol_version: PROTOCOL_VERSION,
};
let payload = rmp_serde::to_vec_named(&msg)?;
let mut writer = self
.writer
.lock()
.map_err(|e| BridgeError::MutexPoisoned(format!("writer: {}", e)))?;
write_framed(writer.get_mut(), &payload)?;
tracing::info!("sent handshake: protocol_version={}", PROTOCOL_VERSION);
Ok(())
}
fn send_snapshot(&self, snapshot: &ObserverSnapshot) -> Result<(), BridgeError> {
let payload = rmp_serde::to_vec_named(snapshot)?;
let mut writer = self
.writer
.lock()
.map_err(|e| BridgeError::MutexPoisoned(format!("writer: {}", e)))?;
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()
.map_err(|e| BridgeError::MutexPoisoned(format!("reader: {}", e)))?;
match read_framed(reader.get_mut())? {
Some(payload) => match rmp_serde::from_slice::<Vec<PlayerInput>>(&payload) {
Ok(inputs) => {
tracing::trace!("received {} inputs", inputs.len());
Ok(inputs)
}
Err(e) => {
let dump_len = payload.len().min(256);
tracing::error!(
"deserialization failed: {}. Raw bytes ({} of {} total): {:02x?}",
e,
dump_len,
payload.len(),
&payload[..dump_len]
);
Err(BridgeError::DeserializationWithDump(format!(
"{} (payload {} bytes)",
e,
payload.len()
)))
}
},
None => Err(BridgeError::Disconnected),
}
}
}
impl Drop for LocalBridge {
fn drop(&mut self) {
// Best-effort socket cleanup
if self.socket_path.exists() {
let _ = fs::remove_file(&self.socket_path);
tracing::trace!("removed socket file {:?}", self.socket_path);
}
}
}