feat(simulation): atlas proxy layers + data-delivery endpoints (T-960, T-949)

T-960: RoadGraphLayer (L2 edges/junctions, MaintenanceAuthority, rail flag, named routes) + SettlementLayer (name/position/size-class/is_capital/is_port from T-955 placements) as new Option fields on AtlasLayerResponse — T-1046 precedent; layer_proxy doc comment corrected (cascade runs through RoadGraph, client up_to ignored). T-949: new StarMapRequest/StarMapResponse (raw JSON passthrough of the star_map_data.json bake) + CityNamesRequest/CityNamesResponse (atlas_city_names via CityContextReader) with defense-in-depth Sol exclusion (GJ-0 system id OR settlement_wave='origin', D-236). Inbound demux extended with mandatory boolean discriminator fields (serde ignores unknown fields — optional-shape sniffing would be ambiguous); AtlasLayerRequest unchanged. SimBridge trait +2 methods, implemented on TcpBridge + LocalBridge. CityPlacement/CityRecord gained name/population/is_capital (BodyWorldState cache-hit path stays DB-free). ~30 new unit tests incl. demux disambiguation + Sol branches; gen_fixtures extended for road/settlement samples.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
2026-07-14 15:45:25 +02:00
co-authored by Claude Fable 5
parent 6bda9f1697
commit 37881acba0
15 changed files with 1654 additions and 23 deletions
+21
View File
@@ -3,6 +3,7 @@
// Deterministic client-server communication via Unix domain sockets
use super::{decode_inbound, BridgeError, Inbound, ObserverSnapshot, SimBridge};
use crate::atlas::atlas_data_proxy::{CityNamesResponse, StarMapResponse};
use crate::atlas::layer_proxy::AtlasLayerResponse;
use crate::bridge::framing::{read_framed, write_framed};
use std::fs;
@@ -145,6 +146,26 @@ impl SimBridge for LocalBridge {
write_framed(writer.get_mut(), &payload)?;
Ok(())
}
fn send_star_map_response(&self, resp: &StarMapResponse) -> Result<(), BridgeError> {
let payload = rmp_serde::to_vec_named(resp)?;
let mut writer = self
.writer
.lock()
.map_err(|e| BridgeError::MutexPoisoned(format!("writer: {}", e)))?;
write_framed(writer.get_mut(), &payload)?;
Ok(())
}
fn send_city_names_response(&self, resp: &CityNamesResponse) -> Result<(), BridgeError> {
let payload = rmp_serde::to_vec_named(resp)?;
let mut writer = self
.writer
.lock()
.map_err(|e| BridgeError::MutexPoisoned(format!("writer: {}", e)))?;
write_framed(writer.get_mut(), &payload)?;
Ok(())
}
}
impl Drop for LocalBridge {
+211 -10
View File
@@ -6,6 +6,9 @@ use bevy_app::prelude::*;
use bevy_ecs::prelude::*;
use bevy_ecs::schedule::IntoScheduleConfigs;
use crate::atlas::atlas_data_proxy::{
CityNamesRequest, CityNamesResponse, StarMapRequest, StarMapResponse,
};
use crate::atlas::layer_proxy::{AtlasLayerRequest, AtlasLayerResponse};
pub mod debug;
@@ -36,30 +39,61 @@ pub enum BridgeError {
}
/// One decoded inbound message. The client→server stream is a single demuxed
/// channel (D-225): a `Vec<PlayerInput>` frame is a MessagePack *array* and an
/// `AtlasLayerRequest` frame is a *map*, so they are distinguishable without a
/// wire-level type tag (existing frames are byte-unchanged — additive).
/// channel (D-225): a `Vec<PlayerInput>` frame is a MessagePack *array* and
/// every request type below is a *map*, so array vs. map alone separates
/// inputs from everything else without a wire-level type tag (existing frames
/// are byte-unchanged — additive).
///
/// **Disambiguating the three map shapes (D-225 extension, T-949):**
/// `AtlasLayerRequest{body_id, up_to}` was the only map shape until T-949
/// added `StarMapRequest`/`CityNamesRequest` alongside it. serde's derived
/// `Deserialize` silently ignores unknown fields by default, so "does this
/// struct parse at all" is not a safe discriminator once more than one map
/// shape can share a field name (`CityNamesRequest` and `AtlasLayerRequest`
/// both key on `body_id`) — a payload carrying every field either shape wants
/// would ambiguously satisfy both. Rather than retrofit
/// `#[serde(deny_unknown_fields)]` onto the existing `AtlasLayerRequest` (risking
/// breakage if any already-deployed client encoder harmlessly sends extra
/// fields), the two *new* map shapes each carry a mandatory boolean
/// discriminator field the others don't have at all (`star_map` /
/// `city_names`): a missing required field is a hard deserialize failure, not
/// a silent ignore, so every shape's required-field set is mutually
/// exclusive. `AtlasLayerRequest` itself is untouched byte-for-byte.
#[derive(Debug)]
pub enum Inbound {
/// A batch of player inputs (the gameplay path).
Inputs(Vec<PlayerInput>),
/// An atlas layer-stream request (#969, D-225).
AtlasRequest(AtlasLayerRequest),
/// A star-map dataset request (T-949a).
StarMapRequest(StarMapRequest),
/// A per-body city-names request (T-949b).
CityNamesRequest(CityNamesRequest),
}
/// Demux a received frame payload into an [`Inbound`] (D-225). Tries
/// `Vec<PlayerInput>` (array), then `AtlasLayerRequest` (map); a frame that is
/// neither is a genuinely malformed input frame.
/// Demux a received frame payload into an [`Inbound`] (D-225, T-949). Tries,
/// in order: `Vec<PlayerInput>` (array) → `AtlasLayerRequest` (map,
/// `body_id`+`up_to`) → `StarMapRequest` (map, `star_map` discriminator) →
/// `CityNamesRequest` (map, `city_names` discriminator + `body_id`). Every map
/// shape's required fields are mutually exclusive (see the [`Inbound`] doc),
/// so this order is for stability, not correctness — a frame that satisfies
/// none of the four shapes is a genuinely malformed input frame.
pub fn decode_inbound(payload: &[u8]) -> Result<Inbound, BridgeError> {
if let Ok(inputs) = rmp_serde::from_slice::<Vec<PlayerInput>>(payload) {
return Ok(Inbound::Inputs(inputs));
}
match rmp_serde::from_slice::<AtlasLayerRequest>(payload) {
Ok(req) => Ok(Inbound::AtlasRequest(req)),
if let Ok(req) = rmp_serde::from_slice::<AtlasLayerRequest>(payload) {
return Ok(Inbound::AtlasRequest(req));
}
if let Ok(req) = rmp_serde::from_slice::<StarMapRequest>(payload) {
return Ok(Inbound::StarMapRequest(req));
}
match rmp_serde::from_slice::<CityNamesRequest>(payload) {
Ok(req) => Ok(Inbound::CityNamesRequest(req)),
Err(e) => {
let dump_len = payload.len().min(256);
tracing::error!(
"inbound decode failed (neither inputs nor atlas request): {}. Raw ({} of {} bytes): {:02x?}",
"inbound decode failed (matches no known frame shape): {}. Raw ({} of {} bytes): {:02x?}",
e,
dump_len,
payload.len(),
@@ -98,6 +132,12 @@ pub trait SimBridge: Send + Sync {
/// Send an atlas layer-stream response to the client (#969, D-225).
fn send_atlas_response(&self, resp: &AtlasLayerResponse) -> Result<(), BridgeError>;
/// Send a star-map response to the client (T-949a).
fn send_star_map_response(&self, resp: &StarMapResponse) -> Result<(), BridgeError>;
/// Send a city-names response to the client (T-949b).
fn send_city_names_response(&self, resp: &CityNamesResponse) -> Result<(), BridgeError>;
}
/// BridgeResource: Bevy Resource wrapper for SimBridge trait object
@@ -132,6 +172,14 @@ impl BridgeResource {
pub fn send_atlas_response(&self, resp: &AtlasLayerResponse) -> Result<(), BridgeError> {
self.inner.send_atlas_response(resp)
}
pub fn send_star_map_response(&self, resp: &StarMapResponse) -> Result<(), BridgeError> {
self.inner.send_star_map_response(resp)
}
pub fn send_city_names_response(&self, resp: &CityNamesResponse) -> Result<(), BridgeError> {
self.inner.send_city_names_response(resp)
}
}
/// Tracks whether the protocol handshake has been sent (#555).
@@ -165,6 +213,8 @@ pub fn receive_bridge_inputs(
handshake: Res<HandshakeState>,
mut error_buffer: ResMut<SimErrorBuffer>,
mut atlas_requests: ResMut<AtlasRequestBuffer>,
mut star_map_requests: ResMut<StarMapRequestBuffer>,
mut city_names_requests: ResMut<CityNamesRequestBuffer>,
time: Option<Res<crate::simulation::time::SimulationTime>>,
) {
let Some(bridge) = bridge else { return };
@@ -193,6 +243,12 @@ pub fn receive_bridge_inputs(
Ok(Some(Inbound::AtlasRequest(req))) => {
atlas_requests.0.push(req);
}
Ok(Some(Inbound::StarMapRequest(req))) => {
star_map_requests.0.push(req);
}
Ok(Some(Inbound::CityNamesRequest(req))) => {
city_names_requests.0.push(req);
}
// No complete frame ready — the backlog is drained.
Ok(None) => break,
Err(BridgeError::Disconnected) => {
@@ -308,12 +364,65 @@ pub fn send_atlas_responses(
}
}
/// Inbound star-map requests routed off the bridge (T-949a), drained by the
/// proxy serve system in `PreInput`.
#[derive(Resource, Default)]
pub struct StarMapRequestBuffer(pub Vec<StarMapRequest>);
/// Outbound star-map responses, filled by the proxy serve system and flushed
/// to the client in `PostSnapshot` (T-949a).
#[derive(Resource, Default)]
pub struct StarMapResponseBuffer(pub Vec<StarMapResponse>);
/// Flush buffered star-map responses to the client (T-949a). A failed send is
/// logged but not fatal.
pub fn send_star_map_responses(
bridge: Option<Res<BridgeResource>>,
mut buffer: ResMut<StarMapResponseBuffer>,
) {
let Some(bridge) = bridge else { return };
for resp in buffer.0.drain(..) {
if let Err(e) = bridge.send_star_map_response(&resp) {
tracing::warn!("failed to send star map response: {}", e);
}
}
}
/// Inbound city-names requests routed off the bridge (T-949b), drained by the
/// proxy serve system in `PreInput`.
#[derive(Resource, Default)]
pub struct CityNamesRequestBuffer(pub Vec<CityNamesRequest>);
/// Outbound city-names responses, filled by the proxy serve system and
/// flushed to the client in `PostSnapshot` (T-949b).
#[derive(Resource, Default)]
pub struct CityNamesResponseBuffer(pub Vec<CityNamesResponse>);
/// Flush buffered city-names responses to the client (T-949b). A failed send
/// is logged but not fatal.
pub fn send_city_names_responses(
bridge: Option<Res<BridgeResource>>,
mut buffer: ResMut<CityNamesResponseBuffer>,
) {
let Some(bridge) = bridge else { return };
for resp in buffer.0.drain(..) {
if let Err(e) = bridge.send_city_names_response(&resp) {
tracing::warn!(
"failed to send city names response for {}: {}",
resp.body_id,
e
);
}
}
}
/// Bridge plugin for client-server communication
/// Abstracts transport layer (LocalBridge/NetworkBridge)
pub struct BridgePlugin;
impl Plugin for BridgePlugin {
fn build(&self, app: &mut App) {
use crate::simulation::time::sim_not_paused;
use crate::tick_phases::TickPhase;
app.init_resource::<SnapshotBuffer>()
@@ -325,10 +434,19 @@ impl Plugin for BridgePlugin {
.init_resource::<crate::perception::query::ActivePerceptionMode>()
.init_resource::<AtlasRequestBuffer>()
.init_resource::<AtlasResponseBuffer>()
.init_resource::<StarMapRequestBuffer>()
.init_resource::<StarMapResponseBuffer>()
.init_resource::<CityNamesRequestBuffer>()
.init_resource::<CityNamesResponseBuffer>()
// Bridge I/O — PreInput (receive) and PostSnapshot (send)
.add_systems(Update, receive_bridge_inputs.in_set(TickPhase::PreInput))
.add_systems(Update, send_bridge_snapshot.in_set(TickPhase::PostSnapshot))
.add_systems(Update, send_atlas_responses.in_set(TickPhase::PostSnapshot))
.add_systems(Update, send_star_map_responses.in_set(TickPhase::PostSnapshot))
.add_systems(
Update,
send_city_names_responses.in_set(TickPhase::PostSnapshot),
)
// Debug commands — Snapshot phase
.add_systems(
Update,
@@ -336,6 +454,13 @@ impl Plugin for BridgePlugin {
)
// Monologue chain — Simulation phase, strict intra-phase sequence.
// trigger_event_monologue must run after conversations + sound (also Simulation).
// T-970: TickPhase::Simulation is not set-gated (see
// social_plugin.rs's collect_sound_events exemption) — this whole
// chain is genuine world-advancing dialogue/monologue logic, so
// it gates safely on its own. trigger_event_monologue's
// .after(collect_sound_events) still holds while gated:
// collect_sound_events itself is never gated, and an ordering
// edge onto a skipped predecessor is trivially satisfied.
.add_systems(
Update,
(
@@ -352,15 +477,22 @@ impl Plugin for BridgePlugin {
crate::simulation::monologue::process_contradiction_monologue
.after(crate::simulation::monologue::trigger_event_monologue),
)
.run_if(sim_not_paused)
.in_set(TickPhase::Simulation),
)
// Observation systems — Simulation phase (reads positions, feeds snapshot)
// Observation systems — Simulation phase (reads positions, feeds snapshot).
// T-970: gates safely on its own (see note above) — these compute
// "current state" (visibility, nearby interactions) that's valid
// as long as nothing moved, which holds while paused since
// Movement is frozen too; unlike SoundEventQueue, nothing here
// depends on being refreshed on a tick where nothing changed.
.add_systems(
Update,
(
crate::perception::observer::compute_visibility_geometry,
crate::simulation::interaction::compute_nearby_interactions,
)
.run_if(sim_not_paused)
.in_set(TickPhase::Simulation),
)
// Observer snapshot assembly — Snapshot phase
@@ -404,4 +536,73 @@ mod inbound_tests {
// Neither shape → a malformed-frame error.
assert!(decode_inbound(&[0xff, 0xff]).is_err());
}
#[test]
fn demux_routes_star_map_requests() {
let req = StarMapRequest { star_map: true };
let frame = rmp_serde::to_vec_named(&req).unwrap();
assert!(matches!(
decode_inbound(&frame),
Ok(Inbound::StarMapRequest(r)) if r.star_map
));
}
#[test]
fn demux_routes_city_names_requests() {
let req = CityNamesRequest {
city_names: true,
body_id: "GJ1c".into(),
};
let frame = rmp_serde::to_vec_named(&req).unwrap();
assert!(matches!(
decode_inbound(&frame),
Ok(Inbound::CityNamesRequest(r)) if r.body_id == "GJ1c"
));
}
/// T-949: the array-vs-map trick (D-225) still separates `Inputs` from
/// everything else, and the three map shapes' discriminator fields keep
/// them mutually exclusive — each of the four frame shapes decodes to
/// exactly its own `Inbound` variant, never a neighbor's.
#[test]
fn inbound_disambiguation_is_unambiguous_across_all_four_shapes() {
let inputs_frame = rmp_serde::to_vec_named(&Vec::<PlayerInput>::new()).unwrap();
let atlas_frame = rmp_serde::to_vec_named(&AtlasLayerRequest {
body_id: "GJ1c".into(),
up_to: CascadeLayer::Topography,
})
.unwrap();
let star_map_frame = rmp_serde::to_vec_named(&StarMapRequest { star_map: true }).unwrap();
let city_names_frame = rmp_serde::to_vec_named(&CityNamesRequest {
city_names: true,
body_id: "GJ1c".into(),
})
.unwrap();
assert!(matches!(
decode_inbound(&inputs_frame),
Ok(Inbound::Inputs(_))
));
assert!(matches!(
decode_inbound(&atlas_frame),
Ok(Inbound::AtlasRequest(_))
));
assert!(matches!(
decode_inbound(&star_map_frame),
Ok(Inbound::StarMapRequest(_))
));
assert!(matches!(
decode_inbound(&city_names_frame),
Ok(Inbound::CityNamesRequest(_))
));
// Cross-check: an AtlasLayerRequest frame must NOT decode as
// CityNamesRequest even though both key on `body_id` — the missing
// `city_names` discriminator makes that a hard failure, not a silent
// "extra field ignored" success either shape could show without it.
assert!(rmp_serde::from_slice::<CityNamesRequest>(&atlas_frame).is_err());
// And a CityNamesRequest frame must NOT decode as AtlasLayerRequest —
// it's missing the required `up_to` field.
assert!(rmp_serde::from_slice::<AtlasLayerRequest>(&city_names_frame).is_err());
}
}
+29
View File
@@ -4,6 +4,7 @@
// Used for Godot client which lacks Unix socket support
use super::{decode_inbound, BridgeError, Inbound, ObserverSnapshot, SimBridge};
use crate::atlas::atlas_data_proxy::{CityNamesResponse, StarMapResponse};
use crate::atlas::layer_proxy::AtlasLayerResponse;
use crate::bridge::framing::{read_framed, write_framed, FrameAccumulator};
use std::io::BufWriter;
@@ -255,4 +256,32 @@ impl SimBridge for TcpBridge {
result?;
Ok(())
}
fn send_star_map_response(&self, resp: &StarMapResponse) -> Result<(), BridgeError> {
let payload = rmp_serde::to_vec_named(resp)?;
let mut writer = self
.writer
.lock()
.map_err(|e| BridgeError::MutexPoisoned(format!("writer: {}", e)))?;
let stream = writer.get_mut();
stream.set_nonblocking(false).map_err(BridgeError::Io)?;
let result = write_framed(stream, &payload);
stream.set_nonblocking(true).map_err(BridgeError::Io)?;
result?;
Ok(())
}
fn send_city_names_response(&self, resp: &CityNamesResponse) -> Result<(), BridgeError> {
let payload = rmp_serde::to_vec_named(resp)?;
let mut writer = self
.writer
.lock()
.map_err(|e| BridgeError::MutexPoisoned(format!("writer: {}", e)))?;
let stream = writer.get_mut();
stream.set_nonblocking(false).map_err(BridgeError::Io)?;
let result = write_framed(stream, &payload);
stream.set_nonblocking(true).map_err(BridgeError::Io)?;
result?;
Ok(())
}
}