#![cfg_attr(
dylint_lib = "running_process_env_literal",
deny(running_process_env_direct)
)]
use std::env;
use std::io::Write;
use std::process::ExitCode;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::{Duration, Instant};
use prost::Message;
use running_process_platform_internal::platform::ipc::endpoint_is_filesystem_backed;
use running_process_platform_internal::platform::process::{
install_shutdown_request_handler, ShutdownRequest,
};
use running_process::broker::broker_http_discovery;
use running_process::broker::broker_http_port::BrokerHttpPort;
use running_process::broker::broker_http_server::BrokerHttpServer;
use running_process::broker::http_endpoint_registry::HttpEndpointRegistry;
use running_process::broker::lifecycle::names_v2::{broker_v2_runtime_dir, v2_program_pipe};
use running_process::broker::lifecycle::privilege::refuse_privileged_run;
use running_process::broker::lifecycle::sid::user_sid_hash;
use running_process::broker::protocol::{
hello_reply, read_frame_with_cap, validate_frame_envelope, write_frame, ErrorCode, Frame,
FrameKind, FrameValidationError, FramingError, Hello, HelloReply, Negotiated, PayloadEncoding,
Refused, CONTROL_PAYLOAD_PROTOCOL, ENVELOPE_VERSION, MAX_HELLO_BYTES, PROTOCOL_VERSION,
};
use running_process::broker::protocol_v2::ServiceDefinitionLoader;
use running_process::broker::server::deadline_stream::{hello_read_deadline, DeadlineStream};
use running_process::broker::server::service_def_loader::ServiceDefinitionError;
use running_process::broker::server::singleton_bind;
use running_process::client::{IpcListener, IpcStream, ListenerNonblockingMode};
macro_rules! say {
($($arg:tt)*) => {{
use std::io::Write as _;
let _ = writeln!(std::io::stdout(), $($arg)*);
}};
}
macro_rules! say_err {
($($arg:tt)*) => {{
use std::io::Write as _;
let _ = writeln!(std::io::stderr(), $($arg)*);
}};
}
const DEFAULT_PROGRAM: &str = "broker-v2-scaffold";
const SCAFFOLD_PIPE_IDX: u32 = 0;
const MAX_INFLIGHT_HANDLERS: usize = 256;
const ACCEPT_POLL_INTERVAL: Duration = Duration::from_millis(20);
const HANDLER_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
#[cfg(coverage)]
unsafe extern "C" {
fn __llvm_profile_write_file() -> std::ffi::c_int;
}
fn flush_coverage_profile() -> Result<(), String> {
#[cfg(coverage)]
{
let result = unsafe { __llvm_profile_write_file() };
if result != 0 {
return Err(format!(
"__llvm_profile_write_file returned nonzero status {result}"
));
}
}
Ok(())
}
fn max_inflight_handlers() -> usize {
running_process::env_vars::BROKER_MAX_INFLIGHT_HANDLERS.count_or(MAX_INFLIGHT_HANDLERS)
}
struct InflightGuard(Arc<AtomicUsize>);
impl Drop for InflightGuard {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::SeqCst);
}
}
#[derive(Debug, Clone)]
struct CliOptions {
no_bind: bool,
once: bool,
program: String,
http_port: Option<BrokerHttpPort>,
}
fn parse_http_port(value: &str) -> Result<BrokerHttpPort, String> {
if value.eq_ignore_ascii_case("dynamic") {
return Ok(BrokerHttpPort::Dynamic);
}
match value.parse::<u16>() {
Ok(0) => Ok(BrokerHttpPort::Dynamic),
Ok(preferred) => Ok(BrokerHttpPort::StaticOrFallback { preferred }),
Err(_) => Err(format!(
"--http-port expects a port number or `dynamic`, got {value:?}"
)),
}
}
fn parse_cli(args: &[String]) -> Result<CliOptions, String> {
let mut opts = CliOptions {
no_bind: false,
once: false,
program: DEFAULT_PROGRAM.to_owned(),
http_port: None,
};
let mut i = 1; while i < args.len() {
match args[i].as_str() {
"--no-bind" => opts.no_bind = true,
"--once" => opts.once = true,
"--program" => {
i += 1;
if i >= args.len() {
return Err("--program requires a value".to_owned());
}
opts.program = args[i].clone();
}
"--http-port" => {
i += 1;
if i >= args.len() {
return Err("--http-port requires a value".to_owned());
}
opts.http_port = Some(parse_http_port(&args[i])?);
}
"--help" | "-h" => {
return Err(format!(
"running-process-broker-v2 {} — usage:\n \
[--program <name>] (default: {DEFAULT_PROGRAM})\n \
[--once] (accept one connection then exit)\n \
[--no-bind] (exit 0 immediately; for integration test)\n \
[--http-port <n|dynamic>] (serve the aggregation page; off by default)",
env!("CARGO_PKG_VERSION")
));
}
unknown => return Err(format!("unknown argument: {unknown}")),
}
i += 1;
}
Ok(opts)
}
fn main() -> ExitCode {
let args: Vec<String> = env::args().collect();
let opts = match parse_cli(&args) {
Ok(o) => o,
Err(msg) => {
say_err!("{msg}");
return ExitCode::from(2);
}
};
say!(
"running-process-broker-v2 {} (slice 1 of running-process#532)",
env!("CARGO_PKG_VERSION")
);
if opts.no_bind {
say!("running-process-broker-v2 --no-bind: skipping listener bind");
return ExitCode::SUCCESS;
}
if let Err(err) = refuse_privileged_run() {
say_err!(
"running-process-broker-v2: refusing privileged startup: {err}. \
Run as an unprivileged user, or set \
RUNNING_PROCESS_BROKER_ALLOW_PRIVILEGED=1 for isolated test environments only."
);
return ExitCode::from(77); }
let sid = match user_sid_hash() {
Ok(s) => s,
Err(err) => {
say_err!("running-process-broker-v2: user_sid_hash failed: {err}");
return ExitCode::from(1);
}
};
let pipe_name = match v2_program_pipe(&opts.program, &sid, SCAFFOLD_PIPE_IDX) {
Ok(n) => n,
Err(err) => {
say_err!("running-process-broker-v2: v2_program_pipe failed: {err}");
return ExitCode::from(1);
}
};
let socket_path = match singleton_bind::resolve_socket_path(&pipe_name) {
Ok(p) => p,
Err(err) => {
say_err!("running-process-broker-v2: resolve_socket_path failed: {err}");
return ExitCode::from(1);
}
};
let listener = match singleton_bind::bind_singleton(&socket_path) {
Ok(l) => l,
Err(singleton_bind::BindSingletonError::AlreadyBound(_)) => {
say_err!(
"running-process-broker-v2: another broker is already \
bound at {socket_path} (program={}). Refusing to \
start to avoid double-bind. Stop the other broker \
first, or pass `--program <other-name>` to bind a \
distinct namespace.",
opts.program,
);
return ExitCode::from(75); }
Err(err) => {
say_err!("running-process-broker-v2: bind failed at {socket_path}: {err:?}");
return ExitCode::from(1);
}
};
say!(
"running-process-broker-v2 bound at {socket_path} (program={}, mode={})",
opts.program,
if opts.once { "once" } else { "loop" }
);
if let Err(err) = std::io::stdout().flush() {
say_err!("running-process-broker-v2: stdout flush failed: {err}");
}
let loader = Arc::new(ServiceDefinitionLoader::default_root());
let inflight = Arc::new(AtomicUsize::new(0));
let http = opts.http_port.and_then(|config| {
match start_http_surface(config, &opts.program, &broker_v2_runtime_dir()) {
Ok(started) => Some(started),
Err(err) => {
say_err!("running-process-broker-v2: HTTP surface disabled: {err}");
None
}
}
});
let exit_code = if opts.once {
accept_one(
&listener,
Arc::clone(&loader),
http.as_ref().map(|h| Arc::clone(&h.registry)),
)
} else {
let shutdown = match install_shutdown_request_handler() {
Ok(shutdown) => shutdown,
Err(err) => {
say_err!(
"running-process-broker-v2: install shutdown request handler failed: {err}"
);
return ExitCode::from(1);
}
};
accept_loop(
&listener,
Arc::clone(&loader),
Arc::clone(&inflight),
&shutdown,
http.as_ref().map(|h| Arc::clone(&h.registry)),
)
};
if endpoint_is_filesystem_backed() {
let _ = std::fs::remove_file(&socket_path);
}
if let Some(started) = &http {
if let Err(err) =
broker_http_discovery::unpublish_http_port(&broker_v2_runtime_dir(), &started.program)
{
say_err!("running-process-broker-v2: could not unpublish HTTP endpoint: {err}");
}
}
if let Err(err) = flush_coverage_profile() {
say_err!("running-process-broker-v2: coverage profile flush failed: {err}");
return ExitCode::from(1);
}
exit_code
}
struct HttpSurface {
program: String,
registry: Arc<HttpEndpointRegistry>,
}
fn start_http_surface(
config: BrokerHttpPort,
program: &str,
runtime_dir: &std::path::Path,
) -> Result<HttpSurface, String> {
let registry = Arc::new(HttpEndpointRegistry::new());
let server =
BrokerHttpServer::bind(config, Arc::clone(®istry)).map_err(|e| e.to_string())?;
let local = server.local_addr();
let path =
broker_http_discovery::publish_http_port(runtime_dir, program, local.ip(), local.port())
.map_err(|e| format!("publishing {local} to {}: {e}", runtime_dir.display()))?;
say!(
"running-process-broker-v2 http at http://{local} (published to {})",
path.display()
);
if let Err(err) = std::io::stdout().flush() {
say_err!("running-process-broker-v2: stdout flush failed: {err}");
}
thread::Builder::new()
.name("rpb-v2-http".to_string())
.spawn(move || loop {
if let Err(err) = server.serve_once() {
say_err!("running-process-broker-v2: http accept failed: {err}");
}
})
.map_err(|e| format!("spawning the http thread: {e}"))?;
Ok(HttpSurface {
program: program.to_owned(),
registry,
})
}
fn accept_loop(
listener: &IpcListener,
loader: Arc<ServiceDefinitionLoader>,
inflight: Arc<AtomicUsize>,
shutdown: &ShutdownRequest,
http: Option<Arc<HttpEndpointRegistry>>,
) -> ExitCode {
if let Err(err) = listener.set_nonblocking(ListenerNonblockingMode::Accept) {
say_err!("running-process-broker-v2: set listener nonblocking failed: {err}");
return ExitCode::from(1);
}
let max_inflight = max_inflight_handlers();
let mut handlers = Vec::new();
loop {
reap_finished_handlers(&mut handlers);
match poll_accept_until_shutdown(shutdown, || listener.accept()) {
Ok(Some(stream)) => {
let n = inflight.fetch_add(1, Ordering::SeqCst);
if n >= max_inflight {
inflight.fetch_sub(1, Ordering::SeqCst);
say_err!(
"running-process-broker-v2: at MAX_INFLIGHT_HANDLERS ({max_inflight}); dropping connection",
);
drop(stream);
continue;
}
let loader = Arc::clone(&loader);
let inflight_handler = Arc::clone(&inflight);
let http_handler = http.clone();
let spawn_result = thread::Builder::new()
.name("rpb-v2-handler".to_string())
.spawn(move || {
let _inflight_guard = InflightGuard(inflight_handler);
let mut s = stream;
let result = handle_hello_with_deadline(&mut s, &loader);
match result {
Ok(svc) => {
if let Some(reg) = &http_handler {
reg.track(svc.clone());
}
say!("running-process-broker-v2 Hello service={svc:?} negotiated",)
}
Err(err) => {
say_err!("running-process-broker-v2 Hello handler failed: {err}")
}
}
});
match spawn_result {
Ok(handler) => handlers.push(handler),
Err(err) => {
say_err!(
"running-process-broker-v2: thread spawn failed: {err}; \
dropping connection"
);
inflight.fetch_sub(1, Ordering::SeqCst);
}
}
}
Ok(None) => {
say!("running-process-broker-v2: shutdown requested; draining handlers");
drain_handlers(&mut handlers, HANDLER_DRAIN_TIMEOUT);
return ExitCode::SUCCESS;
}
Err(err) => {
say_err!("running-process-broker-v2: accept failed: {err}");
return ExitCode::from(1);
}
}
}
}
fn poll_accept_until_shutdown<T>(
shutdown: &ShutdownRequest,
mut accept: impl FnMut() -> std::io::Result<T>,
) -> std::io::Result<Option<T>> {
while !shutdown.requested() {
match accept() {
Ok(value) => return Ok(Some(value)),
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(ACCEPT_POLL_INTERVAL);
}
Err(err) if err.kind() == std::io::ErrorKind::Interrupted => {}
Err(err) => return Err(err),
}
}
Ok(None)
}
fn reap_finished_handlers(handlers: &mut Vec<thread::JoinHandle<()>>) {
let mut index = 0;
while index < handlers.len() {
if handlers[index].is_finished() {
let handler = handlers.swap_remove(index);
if handler.join().is_err() {
say_err!("running-process-broker-v2: handler thread panicked");
}
} else {
index += 1;
}
}
}
fn drain_handlers(handlers: &mut Vec<thread::JoinHandle<()>>, timeout: Duration) {
let deadline = Instant::now() + timeout;
while !handlers.is_empty() && Instant::now() < deadline {
reap_finished_handlers(handlers);
if !handlers.is_empty() {
thread::sleep(ACCEPT_POLL_INTERVAL);
}
}
reap_finished_handlers(handlers);
if !handlers.is_empty() {
say_err!(
"running-process-broker-v2: handler drain timed out after {timeout:?}; \
detaching {} handler(s)",
handlers.len()
);
}
}
fn accept_one(
listener: &IpcListener,
loader: Arc<ServiceDefinitionLoader>,
http: Option<Arc<HttpEndpointRegistry>>,
) -> ExitCode {
match listener.accept() {
Ok(mut stream) => {
say!("running-process-broker-v2 peer connected (--once)");
match handle_hello_with_deadline(&mut stream, &loader) {
Ok(svc) => {
if let Some(reg) = &http {
reg.track(svc.clone());
}
say!("running-process-broker-v2 Hello for service {svc:?} negotiated; exiting");
ExitCode::SUCCESS
}
Err(err) => {
say_err!("running-process-broker-v2: Hello handler failed: {err}");
ExitCode::from(1)
}
}
}
Err(err) => {
say_err!("running-process-broker-v2: accept failed: {err}");
ExitCode::from(1)
}
}
}
fn handle_hello_with_deadline(
stream: &mut IpcStream,
loader: &ServiceDefinitionLoader,
) -> Result<String, String> {
stream
.set_nonblocking(true)
.map_err(|error| format!("set Hello stream nonblocking: {error}"))?;
let read_result = {
let mut deadline_stream = DeadlineStream::new(stream, hello_read_deadline());
read_frame_with_cap(&mut deadline_stream, MAX_HELLO_BYTES)
};
stream
.set_nonblocking(false)
.map_err(|error| format!("restore Hello stream blocking mode: {error}"))?;
match read_result {
Ok(bytes) => handle_hello_bytes(stream, loader, bytes),
Err(error) => {
let reply = reply_for_framing_error(&error);
write_hello_reply_frame(stream, None, &reply)?;
Err(format!("read Hello frame: {error}"))
}
}
}
fn handle_hello_bytes<S: std::io::Write>(
stream: &mut S,
loader: &ServiceDefinitionLoader,
bytes: Vec<u8>,
) -> Result<String, String> {
let request_frame = match Frame::decode(bytes.as_slice()) {
Ok(frame) => frame,
Err(error) => {
let reply =
protocol_refused_reply(ErrorCode::ErrorPeerRejected, "malformed broker Frame", 0);
write_hello_reply_frame(stream, None, &reply)?;
return Err(format!("decode Hello Frame: {error}"));
}
};
if let Err(error) =
validate_frame_envelope(&request_frame, FrameKind::Request, CONTROL_PAYLOAD_PROTOCOL)
{
let (code, reason) = match error {
FrameValidationError::EnvelopeVersion { .. } => (
ErrorCode::ErrorVersionUnsupported,
"frame envelope_version is not v1",
),
FrameValidationError::Kind { .. } => (
ErrorCode::ErrorPeerRejected,
"Hello frame kind must be REQUEST",
),
FrameValidationError::PayloadProtocol { .. } => (
ErrorCode::ErrorPeerRejected,
"Hello frame payload_protocol must be control-plane",
),
FrameValidationError::PayloadEncoding { .. } => (
ErrorCode::ErrorPeerRejected,
"Hello payload must not be compressed",
),
};
let reply = protocol_refused_reply(code, reason, 0);
write_hello_reply_frame(stream, Some(&request_frame), &reply)?;
return Err(format!("validate Hello Frame: {error:?}"));
}
let hello = match Hello::decode(request_frame.payload.as_slice()) {
Ok(hello) => hello,
Err(error) => {
let reply =
protocol_refused_reply(ErrorCode::ErrorPeerRejected, "malformed Hello payload", 0);
write_hello_reply_frame(stream, Some(&request_frame), &reply)?;
return Err(format!("decode Hello payload: {error}"));
}
};
let backend_pipe = resolve_backend_pipe(&hello.service_name);
let reply = build_hello_reply(&hello, loader, &backend_pipe);
write_hello_reply_frame(stream, Some(&request_frame), &reply)?;
match reply.result {
Some(hello_reply::Result::Negotiated(_)) => Ok(hello.service_name),
Some(hello_reply::Result::Refused(r)) => Err(format!("refused: {}", r.reason)),
None => Err("HelloReply missing result oneof".to_string()),
}
}
fn write_hello_reply_frame<S: std::io::Write>(
stream: &mut S,
request_frame: Option<&Frame>,
reply: &HelloReply,
) -> Result<(), String> {
let response_frame = Frame {
envelope_version: PROTOCOL_VERSION,
kind: FrameKind::Response as i32,
payload_protocol: CONTROL_PAYLOAD_PROTOCOL,
payload: reply.encode_to_vec(),
request_id: request_frame.map_or(0, |frame| frame.request_id),
payload_encoding: PayloadEncoding::None as i32,
deadline_unix_ms: 0,
traceparent: request_frame
.map(|frame| frame.traceparent.clone())
.unwrap_or_default(),
tracestate: request_frame
.map(|frame| frame.tracestate.clone())
.unwrap_or_default(),
};
write_frame(stream, &response_frame.encode_to_vec())
.map_err(|error| format!("write HelloReply Frame: {error}"))?;
Ok(())
}
fn reply_for_framing_error(error: &FramingError) -> HelloReply {
match error {
FramingError::UnsupportedFramingVersion { .. } => protocol_refused_reply(
ErrorCode::ErrorVersionUnsupported,
"unsupported framing version",
0,
),
FramingError::FrameTooLarge { .. } => protocol_refused_reply(
ErrorCode::ErrorPeerRejected,
"initial Hello frame exceeds 64 KiB",
0,
),
FramingError::UnexpectedEof { .. } | FramingError::Io(_) => {
protocol_refused_reply(ErrorCode::ErrorPeerRejected, "incomplete Hello frame", 0)
}
FramingError::Decode(_) => {
protocol_refused_reply(ErrorCode::ErrorPeerRejected, "malformed Hello frame", 0)
}
}
}
fn resolve_backend_pipe(service: &str) -> String {
use running_process::broker::backend_sdk::read_daemon_identity_file;
use running_process::broker::lifecycle::names_v2::daemon_identity_path;
read_daemon_identity_file(&daemon_identity_path(service))
.map(|daemon| daemon.ipc_endpoint.path)
.unwrap_or_default()
}
fn build_hello_reply(
hello: &Hello,
loader: &ServiceDefinitionLoader,
backend_pipe: &str,
) -> HelloReply {
if hello.client_min_protocol > ENVELOPE_VERSION as u32
|| hello.client_max_protocol < ENVELOPE_VERSION as u32
{
return refused_reply(
hello,
ErrorCode::ErrorVersionUnsupported,
"client protocol range does not include the version this broker speaks",
0,
);
}
let definition = match loader.load(&hello.service_name) {
Ok(d) => d,
Err(ServiceDefinitionError::Io(err)) if err.kind() == std::io::ErrorKind::NotFound => {
return refused_reply(
hello,
ErrorCode::ErrorServiceUnknown,
"service definition was not found",
0,
);
}
Err(ServiceDefinitionError::InvalidName(_)) => {
return refused_reply(
hello,
ErrorCode::ErrorServiceUnknown,
"service name is invalid",
0,
);
}
Err(other) => {
return refused_reply(
hello,
ErrorCode::ErrorServiceUnknown,
format!("service definition could not be loaded: {other}"),
0,
);
}
};
if !definition.min_version.is_empty()
&& hello.wanted_version.as_str() < definition.min_version.as_str()
{
return refused_reply(
hello,
ErrorCode::ErrorVersionBlocked,
format!(
"wanted_version {:?} is below min_version {:?}",
hello.wanted_version, definition.min_version
),
0,
);
}
if !definition.version_allow_list.is_empty()
&& !definition
.version_allow_list
.iter()
.any(|v| v == &hello.wanted_version)
{
return refused_reply(
hello,
ErrorCode::ErrorVersionBlocked,
format!(
"wanted_version {:?} is not in version_allow_list",
hello.wanted_version
),
0,
);
}
HelloReply {
result: Some(hello_reply::Result::Negotiated(Negotiated {
negotiated_protocol: ENVELOPE_VERSION as u32,
daemon_version: env!("CARGO_PKG_VERSION").into(),
backend_pipe: backend_pipe.to_string(),
warnings: Vec::new(),
server_capabilities: 0,
keepalive_interval_secs: 0,
handle_passed_token: Vec::new(),
connection_id: hello.connection_id,
})),
}
}
fn refused_reply(
hello: &Hello,
code: ErrorCode,
reason: impl Into<String>,
retry_after_ms: u64,
) -> HelloReply {
protocol_refused_reply(code, reason, retry_after_ms).with_connection_id(hello.connection_id)
}
fn protocol_refused_reply(
code: ErrorCode,
reason: impl Into<String>,
retry_after_ms: u64,
) -> HelloReply {
HelloReply {
result: Some(hello_reply::Result::Refused(Refused {
reason: reason.into(),
daemon_min_protocol: ENVELOPE_VERSION as u32,
daemon_max_protocol: ENVELOPE_VERSION as u32,
code: code as i32,
details: std::collections::HashMap::new(),
retry_after_ms,
})),
}
}
trait HelloReplyExt {
fn with_connection_id(self, id: u64) -> Self;
}
impl HelloReplyExt for HelloReply {
fn with_connection_id(mut self, id: u64) -> Self {
if let Some(hello_reply::Result::Refused(_)) = &self.result {
} else if let Some(hello_reply::Result::Negotiated(ref mut n)) = self.result {
n.connection_id = id;
}
self
}
}
#[cfg(test)]
#[path = "running-process-broker-v2/tests.rs"]
mod tests;