feat: publish HoloLake model-native living system source

This commit is contained in:
冰朔 2026-08-03 10:04:41 +08:00
commit c395dd3a99
2467 changed files with 615073 additions and 0 deletions

View file

@ -0,0 +1,11 @@
[package]
name = "guanghu-broadcast-tower"
version = "0.1.0"
edition = "2021"
license = "AGPL-3.0-or-later"
description = "Hosted bootstrap executor for the unique HLDP Guanghu broadcast tower"
[dependencies]
guanghu-hldp-runtime = { path = "../hldp-runtime" }
serde = { version = "1", features = ["derive"] }
serde_json = "1"

View file

@ -0,0 +1,343 @@
use std::{
fs,
io::{self, Read, Write},
net::{SocketAddr, TcpListener, TcpStream},
path::Path,
sync::atomic::{AtomicU64, Ordering},
time::{Duration, SystemTime, UNIX_EPOCH},
};
use guanghu_hldp_runtime::{validate_world_seed, WorldManifest};
use serde::Serialize;
static EPOCH_SEQUENCE: AtomicU64 = AtomicU64::new(0);
const USAGE: &str =
"usage: guanghu-broadcast-tower serve <world-root> <loopback-address> <epoch-path>";
#[derive(Clone, Debug, Serialize)]
pub struct TowerSnapshot {
pub status: &'static str,
pub tower_id: String,
pub world_id: String,
pub world_version: String,
pub phase: String,
pub domain_count: usize,
pub code_channel_id: String,
pub authority_language: &'static str,
pub runtime: &'static str,
pub linux_exited: bool,
}
impl TowerSnapshot {
pub fn from_manifest(manifest: &WorldManifest) -> Self {
Self {
status: "ok",
tower_id: manifest.broadcast_tower.id.clone(),
world_id: manifest.world_id.clone(),
world_version: manifest.version.clone(),
phase: manifest.phase.clone(),
domain_count: manifest.domains.len(),
code_channel_id: manifest.code_channel.id.clone(),
authority_language: "HLDP",
runtime: "HOSTED_BOOTSTRAP",
linux_exited: false,
}
}
}
pub fn run_command(arguments: Vec<String>, connection_limit: Option<usize>) -> Result<(), String> {
let mut arguments = arguments.into_iter();
let command = arguments.next().ok_or_else(|| USAGE.to_owned())?;
let world_root = arguments.next().ok_or_else(|| USAGE.to_owned())?;
let address = arguments.next().ok_or_else(|| USAGE.to_owned())?;
let epoch_path = arguments.next().ok_or_else(|| USAGE.to_owned())?;
if command != "serve" || arguments.next().is_some() {
return Err(USAGE.to_owned());
}
let address = validate_loopback_address(&address)?;
let manifest =
validate_world_seed(Path::new(&world_root)).map_err(|error| error.to_string())?;
let listener = TcpListener::bind(address).map_err(|error| error.to_string())?;
let snapshot = TowerSnapshot::from_manifest(&manifest);
serve(listener, snapshot, Path::new(&epoch_path), connection_limit)
.map_err(|error| error.to_string())
}
pub fn validate_loopback_address(address: &str) -> Result<SocketAddr, String> {
let parsed: SocketAddr = address
.parse()
.map_err(|error| format!("invalid broadcast address {address}: {error}"))?;
if !parsed.ip().is_loopback() {
return Err("hosted broadcast tower must bind to loopback".to_owned());
}
Ok(parsed)
}
pub fn write_epoch(path: &Path, snapshot: &TowerSnapshot, address: SocketAddr) -> io::Result<()> {
let parent = path
.parent()
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "epoch path has no parent"))?;
fs::create_dir_all(parent)?;
let started_at_unix = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_err(io::Error::other)?
.as_secs();
let temporary = parent.join(format!(
".epoch.{}.{}.tmp",
std::process::id(),
EPOCH_SEQUENCE.fetch_add(1, Ordering::Relaxed)
));
let contents = format!(
concat!(
"schema: guanghu.broadcast-epoch/v1\n",
"tower_id: {}\n",
"world_id: {}\n",
"world_version: {}\n",
"phase: {}\n",
"status: RUNNING_HOSTED\n",
"authority_language: HLDP\n",
"bind: {}\n",
"pid: {}\n",
"started_at_unix: {}\n",
"linux_dependency: true\n",
"native_claim: false\n"
),
snapshot.tower_id,
snapshot.world_id,
snapshot.world_version,
snapshot.phase,
address,
std::process::id(),
started_at_unix
);
fs::write(&temporary, contents)?;
fs::rename(temporary, path)
}
pub fn serve(
listener: TcpListener,
snapshot: TowerSnapshot,
epoch_path: &Path,
connection_limit: Option<usize>,
) -> io::Result<()> {
let address = listener.local_addr()?;
if !address.ip().is_loopback() {
return Err(io::Error::new(
io::ErrorKind::PermissionDenied,
"hosted broadcast tower may only listen on loopback",
));
}
write_epoch(epoch_path, &snapshot, address)?;
for connection in listener
.incoming()
.take(connection_limit.unwrap_or(usize::MAX))
{
handle_connection(connection?, &snapshot)?;
}
Ok(())
}
fn handle_connection(mut stream: TcpStream, snapshot: &TowerSnapshot) -> io::Result<()> {
stream.set_read_timeout(Some(Duration::from_secs(5)))?;
let mut request = [0_u8; 8192];
let length = stream.read(&mut request)?;
let request = String::from_utf8_lossy(&request[..length]);
let request_line = request.lines().next().unwrap_or_default();
let (status, content_type, body) = match request_line {
"GET /healthz HTTP/1.1" | "GET /healthz HTTP/1.0" => (
"200 OK",
"application/json",
serde_json::to_string(snapshot).map_err(io::Error::other)?,
),
"GET /v1/world HTTP/1.1" | "GET /v1/world HTTP/1.0" => (
"200 OK",
"application/json",
serde_json::to_string_pretty(snapshot).map_err(io::Error::other)?,
),
line if line.starts_with("GET ") => (
"404 Not Found",
"application/json",
"{\"error\":\"route_not_found\"}".to_owned(),
),
_ => (
"405 Method Not Allowed",
"application/json",
"{\"error\":\"method_not_allowed\"}".to_owned(),
),
};
let response = format!(
"HTTP/1.1 {status}\r\nContent-Type: {content_type}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
);
stream.write_all(response.as_bytes())?;
stream.flush()
}
#[cfg(test)]
mod tests {
use std::{
io::{Read, Write},
net::{TcpListener, TcpStream},
path::PathBuf,
thread,
};
use guanghu_hldp_runtime::load_world_manifest;
use super::{run_command, serve, validate_loopback_address, write_epoch, TowerSnapshot, USAGE};
fn snapshot() -> TowerSnapshot {
TowerSnapshot {
status: "ok",
tower_id: "BT-GH-ROOT-0001".to_owned(),
world_id: "GLW-ROOT-0001".to_owned(),
world_version: "0.1.0-stage1".to_owned(),
phase: "HOSTED_BOOTSTRAP_PROTOTYPE".to_owned(),
domain_count: 5,
code_channel_id: "HLP-MOD-CODE-CHANNEL".to_owned(),
authority_language: "HLDP",
runtime: "HOSTED_BOOTSTRAP",
linux_exited: false,
}
}
fn epoch_path() -> PathBuf {
std::env::temp_dir().join(format!(
"guanghu-broadcast-epoch-{}-{}.hldp",
std::process::id(),
thread::current().name().unwrap_or("test")
))
}
fn request(request: &[u8]) -> String {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind test tower");
let address = listener.local_addr().expect("test address");
let epoch = epoch_path();
let epoch_for_server = epoch.clone();
let handle = thread::spawn(move || {
serve(listener, snapshot(), &epoch_for_server, Some(1)).expect("serve one request");
});
let mut stream = TcpStream::connect(address).expect("connect to test tower");
stream.write_all(request).expect("write request");
let mut response = String::new();
stream.read_to_string(&mut response).expect("read response");
handle.join().expect("tower thread");
std::fs::remove_file(epoch).expect("remove test epoch");
response
}
#[test]
fn rejects_non_loopback_hosted_bindings() {
assert!(validate_loopback_address("not-an-address").is_err());
assert!(validate_loopback_address("0.0.0.0:8077").is_err());
assert!(validate_loopback_address("127.0.0.1:8077").is_ok());
}
#[test]
fn command_parser_rejects_bad_shapes_and_missing_worlds() {
assert_eq!(run_command(vec![], Some(0)), Err(USAGE.to_owned()));
assert_eq!(
run_command(
vec![
"wrong".to_owned(),
"world".to_owned(),
"127.0.0.1:0".to_owned(),
"/tmp/epoch".to_owned(),
],
Some(0),
),
Err(USAGE.to_owned())
);
let error = run_command(
vec![
"serve".to_owned(),
"/definitely/missing".to_owned(),
"127.0.0.1:0".to_owned(),
"/tmp/missing-epoch".to_owned(),
],
Some(0),
)
.expect_err("missing world must fail");
assert!(error.contains("WORLD-MANIFEST.hldp"));
}
#[test]
fn snapshot_is_derived_from_the_registered_world() {
let path =
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../world-seed/WORLD-MANIFEST.hldp");
let manifest = load_world_manifest(&path).expect("world manifest");
let derived = TowerSnapshot::from_manifest(&manifest);
assert_eq!(derived.tower_id, manifest.broadcast_tower.id);
assert_eq!(derived.world_id, manifest.world_id);
assert_eq!(derived.world_version, manifest.version);
assert_eq!(derived.phase, manifest.phase);
assert_eq!(derived.domain_count, 5);
assert_eq!(derived.code_channel_id, manifest.code_channel.id);
assert!(!derived.linux_exited);
}
#[test]
fn epoch_and_listener_paths_fail_closed() {
let address = validate_loopback_address("127.0.0.1:8077").expect("loopback");
assert!(write_epoch(PathBuf::new().as_path(), &snapshot(), address).is_err());
let public = TcpListener::bind("0.0.0.0:0").expect("bind wildcard test listener");
let error = serve(public, snapshot(), &epoch_path(), Some(0))
.expect_err("wildcard listener must fail closed");
assert_eq!(error.kind(), std::io::ErrorKind::PermissionDenied);
}
#[test]
fn serves_hldp_world_health_and_writes_epoch() {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind test tower");
let address = listener.local_addr().expect("test address");
let epoch = epoch_path();
let epoch_for_server = epoch.clone();
let handle = thread::spawn(move || {
serve(listener, snapshot(), &epoch_for_server, Some(1)).expect("serve one request");
});
let mut stream = TcpStream::connect(address).expect("connect to test tower");
stream
.write_all(b"GET /healthz HTTP/1.1\r\nHost: localhost\r\n\r\n")
.expect("write request");
let mut response = String::new();
stream.read_to_string(&mut response).expect("read response");
handle.join().expect("tower thread");
assert!(response.starts_with("HTTP/1.1 200 OK"));
assert!(response.contains("\"tower_id\":\"BT-GH-ROOT-0001\""));
assert!(response.contains("\"domain_count\":5"));
assert!(response.contains("\"linux_exited\":false"));
let epoch_contents = std::fs::read_to_string(&epoch).expect("read epoch");
assert!(epoch_contents.contains("status: RUNNING_HOSTED"));
assert!(epoch_contents.contains("authority_language: HLDP"));
assert!(epoch_contents.contains("native_claim: false"));
std::fs::remove_file(epoch).expect("remove test epoch");
}
#[test]
fn unknown_routes_fail_closed() {
let response = request(b"GET /guess HTTP/1.1\r\nHost: localhost\r\n\r\n");
assert!(response.starts_with("HTTP/1.1 404 Not Found"));
assert!(response.contains("route_not_found"));
}
#[test]
fn serves_pretty_world_and_rejects_non_get_methods() {
let world = request(b"GET /v1/world HTTP/1.0\r\nHost: localhost\r\n\r\n");
assert!(world.starts_with("HTTP/1.1 200 OK"));
assert!(world.contains("\n \"world_id\": \"GLW-ROOT-0001\""));
let rejected = request(b"POST /healthz HTTP/1.1\r\nHost: localhost\r\n\r\n");
assert!(rejected.starts_with("HTTP/1.1 405 Method Not Allowed"));
assert!(rejected.contains("method_not_allowed"));
}
}

View file

@ -0,0 +1,74 @@
use std::{env, process::ExitCode};
fn main() -> ExitCode {
exit_code(guanghu_broadcast_tower::run_command(
env::args().skip(1).collect(),
None,
))
}
fn exit_code(result: Result<(), String>) -> ExitCode {
match result {
Ok(()) => ExitCode::SUCCESS,
Err(error) => {
eprintln!("GUANGHU_BROADCAST_TOWER_ERROR: {error}");
ExitCode::FAILURE
}
}
}
#[cfg(test)]
mod tests {
use std::{
io::{Read, Write},
net::{TcpListener, TcpStream},
process::ExitCode,
thread,
time::Duration,
};
use super::exit_code;
#[test]
fn exit_code_is_binary() {
assert_eq!(exit_code(Ok(())), ExitCode::SUCCESS);
assert_eq!(exit_code(Err("failed".to_owned())), ExitCode::FAILURE);
}
#[test]
fn bounded_service_completes_after_one_connection() {
let reservation = TcpListener::bind("127.0.0.1:0").expect("reserve port");
let address = reservation.local_addr().expect("reserved address");
drop(reservation);
let world = format!("{}/../../world-seed", env!("CARGO_MANIFEST_DIR"));
let epoch = std::env::temp_dir().join(format!("tower-main-epoch-{}", std::process::id()));
let epoch_for_server = epoch.clone();
let address_for_server = address.to_string();
let handle = thread::spawn(move || {
guanghu_broadcast_tower::run_command(
vec![
"serve".to_owned(),
world,
address_for_server,
epoch_for_server.to_string_lossy().into_owned(),
],
Some(1),
)
});
let mut stream = loop {
match TcpStream::connect(address) {
Ok(stream) => break stream,
Err(_) => thread::sleep(Duration::from_millis(10)),
}
};
stream
.write_all(b"GET /healthz HTTP/1.0\r\n\r\n")
.expect("request");
let mut response = String::new();
stream.read_to_string(&mut response).expect("response");
assert!(response.starts_with("HTTP/1.1 200 OK"));
assert_eq!(handle.join().expect("service thread"), Ok(()));
std::fs::remove_file(epoch).expect("remove epoch");
}
}

View file

@ -0,0 +1,183 @@
use std::{
io::{Read, Write},
net::{TcpListener, TcpStream},
path::{Path, PathBuf},
thread,
time::Duration,
};
use guanghu_broadcast_tower::{
run_command, serve, validate_loopback_address, write_epoch, TowerSnapshot,
};
use guanghu_hldp_runtime::load_world_manifest;
fn world_root() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../world-seed")
}
fn snapshot() -> TowerSnapshot {
let manifest =
load_world_manifest(&world_root().join("WORLD-MANIFEST.hldp")).expect("world manifest");
TowerSnapshot::from_manifest(&manifest)
}
fn epoch_path(label: &str) -> PathBuf {
std::env::temp_dir().join(format!(
"broadcast-library-{label}-{}-{}.hldp",
std::process::id(),
thread::current().name().unwrap_or("test")
))
}
fn request(label: &str, request: &[u8]) -> String {
let listener = TcpListener::bind("127.0.0.1:0").expect("bind test tower");
let address = listener.local_addr().expect("test address");
let epoch = epoch_path(label);
let epoch_for_server = epoch.clone();
let handle = thread::spawn(move || {
serve(listener, snapshot(), &epoch_for_server, Some(1)).expect("serve request");
});
let mut stream = TcpStream::connect(address).expect("connect");
stream.write_all(request).expect("request");
let mut response = String::new();
stream.read_to_string(&mut response).expect("response");
handle.join().expect("tower thread");
std::fs::remove_file(epoch).expect("remove epoch");
response
}
#[test]
fn complete_library_surface_is_binary_and_loopback_only() {
assert!(validate_loopback_address("not-an-address").is_err());
assert!(validate_loopback_address("0.0.0.0:8077").is_err());
let address = validate_loopback_address("127.0.0.1:8077").expect("loopback");
assert!(write_epoch(Path::new(""), &snapshot(), address).is_err());
let public = TcpListener::bind("0.0.0.0:0").expect("wildcard listener");
assert!(serve(public, snapshot(), &epoch_path("public"), Some(0)).is_err());
assert!(run_command(vec![], Some(0)).is_err());
for incomplete in [
vec!["serve".to_owned()],
vec!["serve".to_owned(), "world".to_owned()],
vec![
"serve".to_owned(),
"world".to_owned(),
"127.0.0.1:0".to_owned(),
],
] {
assert!(run_command(incomplete, Some(0)).is_err());
}
assert!(run_command(
vec![
"wrong".to_owned(),
"world".to_owned(),
"127.0.0.1:0".to_owned(),
"/tmp/epoch".to_owned(),
],
Some(0),
)
.is_err());
let occupied = TcpListener::bind("127.0.0.1:0").expect("occupied loopback");
assert!(run_command(
vec![
"serve".to_owned(),
world_root().to_string_lossy().into_owned(),
occupied.local_addr().expect("occupied address").to_string(),
epoch_path("occupied").to_string_lossy().into_owned(),
],
Some(0),
)
.is_err());
drop(occupied);
assert!(run_command(
vec![
"serve".to_owned(),
world_root().to_string_lossy().into_owned(),
"127.0.0.1:0".to_owned(),
String::new(),
],
Some(0),
)
.is_err());
assert!(run_command(
vec![
"serve".to_owned(),
"/definitely/missing".to_owned(),
"127.0.0.1:0".to_owned(),
"/tmp/epoch".to_owned(),
],
Some(0),
)
.is_err());
let reservation = TcpListener::bind("127.0.0.1:0").expect("reserve loopback");
let address = reservation.local_addr().expect("reserved address");
drop(reservation);
let epoch = epoch_path("command");
let epoch_for_server = epoch.clone();
let world = world_root().to_string_lossy().into_owned();
let handle = thread::spawn(move || {
run_command(
vec![
"serve".to_owned(),
world,
address.to_string(),
epoch_for_server.to_string_lossy().into_owned(),
],
Some(1),
)
});
let mut stream = loop {
match TcpStream::connect(address) {
Ok(stream) => break stream,
Err(_) => thread::sleep(Duration::from_millis(10)),
}
};
stream
.write_all(b"GET /healthz HTTP/1.0\r\n\r\n")
.expect("command request");
let mut response = String::new();
stream
.read_to_string(&mut response)
.expect("command response");
assert_eq!(handle.join().expect("command thread"), Ok(()));
assert!(response.contains("200 OK"));
std::fs::remove_file(epoch).expect("remove command epoch");
}
#[test]
fn every_registered_http_decision_is_exercised() {
for (label, raw, status, body) in [
(
"health",
b"GET /healthz HTTP/1.1\r\n\r\n".as_slice(),
"200 OK",
"\"tower_id\":\"BT-GH-ROOT-0001\"",
),
(
"world",
b"GET /v1/world HTTP/1.0\r\n\r\n".as_slice(),
"200 OK",
"\"world_id\": \"GLW-ROOT-0001\"",
),
(
"missing",
b"GET /missing HTTP/1.1\r\n\r\n".as_slice(),
"404 Not Found",
"route_not_found",
),
(
"method",
b"POST /healthz HTTP/1.1\r\n\r\n".as_slice(),
"405 Method Not Allowed",
"method_not_allowed",
),
] {
let response = request(label, raw);
assert!(response.contains(status));
assert!(response.contains(body));
}
}

View file

@ -0,0 +1,98 @@
use std::{
fs,
io::{Read, Write},
net::{TcpListener, TcpStream},
path::{Path, PathBuf},
process::{Child, Command, Output, Stdio},
thread,
time::{Duration, Instant},
};
fn world_root() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../world-seed")
}
fn binary() -> &'static str {
env!("CARGO_BIN_EXE_guanghu-broadcast-tower")
}
fn run(arguments: &[&str]) -> Output {
Command::new(binary())
.args(arguments)
.output()
.expect("broadcast tower should start")
}
fn wait_for_epoch(child: &mut Child, epoch: &Path) {
let deadline = Instant::now() + Duration::from_secs(5);
while Instant::now() < deadline {
if epoch.is_file() {
return;
}
if let Some(status) = child.try_wait().expect("inspect tower process") {
panic!("broadcast tower exited before epoch with {status}");
}
thread::sleep(Duration::from_millis(20));
}
panic!("broadcast tower did not write its epoch");
}
#[test]
fn command_rejects_missing_arguments_and_public_bindings() {
let usage = run(&[]);
assert!(!usage.status.success());
assert!(String::from_utf8_lossy(&usage.stderr).contains("usage:"));
let root = world_root();
let public = run(&[
"serve",
root.to_str().expect("UTF-8 world root"),
"0.0.0.0:8077",
"/tmp/guanghu-public-epoch.hldp",
]);
assert!(!public.status.success());
assert!(String::from_utf8_lossy(&public.stderr).contains("loopback"));
}
#[test]
fn command_serves_the_registered_world_on_loopback() {
let reservation = TcpListener::bind("127.0.0.1:0").expect("reserve loopback port");
let address = reservation.local_addr().expect("reserved address");
drop(reservation);
let epoch =
std::env::temp_dir().join(format!("guanghu-command-epoch-{}.hldp", std::process::id()));
let root = world_root();
let mut child = Command::new(binary())
.args([
"serve",
root.to_str().expect("UTF-8 world root"),
&address.to_string(),
epoch.to_str().expect("UTF-8 epoch path"),
])
.stdout(Stdio::null())
.stderr(Stdio::piped())
.spawn()
.expect("start broadcast tower");
wait_for_epoch(&mut child, &epoch);
let mut stream = TcpStream::connect(address).expect("connect to broadcast tower");
stream
.write_all(b"GET /v1/world HTTP/1.1\r\nHost: localhost\r\n\r\n")
.expect("write world request");
let mut response = String::new();
stream
.read_to_string(&mut response)
.expect("read world response");
child.kill().expect("stop test broadcast tower");
child.wait().expect("reap test broadcast tower");
let epoch_contents = fs::read_to_string(&epoch).expect("read command epoch");
fs::remove_file(&epoch).expect("remove command epoch");
assert!(response.starts_with("HTTP/1.1 200 OK"));
assert!(response.contains("\"world_id\": \"GLW-ROOT-0001\""));
assert!(response.contains("\"domain_count\": 5"));
assert!(response.contains("\"linux_exited\": false"));
assert!(epoch_contents.contains("tower_id: BT-GH-ROOT-0001"));
}