use super::*;
pub struct TlsConfig {
pub cert: Vec<u8>,
pub key: Vec<u8>,
}
pub(super) fn serve_http(
port: u16,
handler_name: String,
program: Arc<Program>,
policy: Policy,
tls: Option<TlsConfig>,
opts: ServeOpts,
) -> Result<Value, String> {
match tls {
None => serve_http_plain(port, handler_name, program, policy, opts),
Some(cfg) => serve_http_tls_legacy(port, handler_name, program, policy, cfg),
}
}
pub(super) fn serve_http_plain(
port: u16,
handler_name: String,
program: Arc<Program>,
policy: Policy,
opts: ServeOpts,
) -> Result<Value, String> {
use http_body_util::BodyExt as _;
use hyper::server::conn::http1;
use hyper::service::service_fn;
use hyper_util::rt::{TokioExecutor, TokioIo};
use hyper_util::server::conn::auto;
use tokio::net::TcpListener as TokioTcpListener;
let inline_vm = opts.inline_vm;
let http2 = opts.http2;
let host = opts.host.clone();
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.map_err(|e| format!("net.serve: tokio runtime: {e}"))?;
rt.block_on(async move {
let listener = TokioTcpListener::bind((host.as_str(), port))
.await
.map_err(|e| format!("net.serve bind {host}:{port}: {e}"))?;
eprintln!(
"net.serve: listening on http://{host}:{port}{}{}",
if inline_vm { " (inline-vm)" } else { "" },
if http2 { " (http1+http2)" } else { "" }
);
loop {
let (stream, _) = listener
.accept()
.await
.map_err(|e| format!("net.serve accept: {e}"))?;
let io = TokioIo::new(stream);
let program = Arc::clone(&program);
let policy = policy.clone();
let handler_name = handler_name.clone();
tokio::spawn(async move {
let program2 = Arc::clone(&program);
let policy2 = policy.clone();
let handler_name2 = handler_name.clone();
let svc = service_fn(move |req: hyper::Request<hyper::body::Incoming>| {
let program = Arc::clone(&program2);
let policy = policy2.clone();
let handler_name = handler_name2.clone();
async move {
let (parts, body) = req.into_parts();
let body_bytes = body
.collect()
.await
.map(|c| c.to_bytes())
.unwrap_or_default();
let result = if inline_vm {
let lex_req = build_request_value_parts(&parts, &body_bytes);
let handler = DefaultHandler::new(policy)
.with_program(Arc::clone(&program));
let mut vm = Vm::with_handler(&program, Box::new(handler));
let r = vm.call(&handler_name, vec![lex_req]);
Ok(r.map(|v| unpack_response(&mut vm, &v)))
} else {
tokio::task::spawn_blocking(move || {
let lex_req = build_request_value_parts(&parts, &body_bytes);
let handler = DefaultHandler::new(policy)
.with_program(Arc::clone(&program));
let mut vm = Vm::with_handler(&program, Box::new(handler));
let r = vm.call(&handler_name, vec![lex_req]);
r.map(|v| unpack_response(&mut vm, &v))
})
.await
};
Ok::<_, std::convert::Infallible>(match result {
Ok(Ok(unpacked)) => build_hyper_response(unpacked),
Ok(Err(e)) => error_response(500, &format!("internal error: {e}")),
Err(e) => error_response(500, &format!("task panicked: {e}")),
})
}
});
let result = if http2 {
auto::Builder::new(TokioExecutor::new())
.serve_connection(io, svc)
.await
.map_err(|e| e.to_string())
} else {
http1::Builder::new()
.serve_connection(io, svc)
.await
.map_err(|e| e.to_string())
};
if let Err(e) = result {
eprintln!("net.serve: connection error: {e}");
}
});
}
})
}
pub(super) fn serve_http_tls_legacy(
port: u16,
handler_name: String,
program: Arc<Program>,
policy: Policy,
cfg: TlsConfig,
) -> Result<Value, String> {
let ssl = tiny_http::SslConfig {
certificate: cfg.cert,
private_key: cfg.key,
};
let server = tiny_http::Server::https(("0.0.0.0", port), ssl)
.map_err(|e| format!("net.serve_tls bind {port}: {e}"))?;
eprintln!("net.serve: listening on https://0.0.0.0:{port}");
for req in server.incoming_requests() {
let program = Arc::clone(&program);
let policy = policy.clone();
let handler_name = handler_name.clone();
std::thread::spawn(move || handle_request_tls(req, program, policy, handler_name));
}
Ok(Value::Unit)
}
pub(super) fn handle_request_tls(
mut req: tiny_http::Request,
program: Arc<Program>,
policy: Policy,
handler_name: String,
) {
let lex_req = build_request_value_tiny(&mut req);
let handler = DefaultHandler::new(policy).with_program(Arc::clone(&program));
let mut vm = Vm::with_handler(&program, Box::new(handler));
match vm.call(&handler_name, vec![lex_req]) {
Ok(resp) => {
let (status, body, headers) = unpack_response(&mut vm, &resp);
respond_with_body_tls(req, status, body, headers);
}
Err(e) => {
let response = tiny_http::Response::from_string(format!("internal error: {e}"))
.with_status_code(500);
let _ = req.respond(response);
}
}
}
pub(super) fn serve_http_fn(
port: u16,
closure: Value,
program: Arc<Program>,
policy: Policy,
opts: ServeOpts,
) -> Result<Value, String> {
use http_body_util::BodyExt as _;
use hyper::server::conn::http1;
use hyper::service::service_fn;
use hyper_util::rt::{TokioExecutor, TokioIo};
use hyper_util::server::conn::auto;
use tokio::net::TcpListener as TokioTcpListener;
let inline_vm = opts.inline_vm;
let http2 = opts.http2;
let host = opts.host.clone();
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.map_err(|e| format!("net.serve_fn: tokio runtime: {e}"))?;
rt.block_on(async move {
let listener = TokioTcpListener::bind((host.as_str(), port))
.await
.map_err(|e| format!("net.serve_fn bind {host}:{port}: {e}"))?;
eprintln!(
"net.serve_fn: listening on http://{host}:{port}{}{}",
if inline_vm { " (inline-vm)" } else { "" },
if http2 { " (http1+http2)" } else { "" }
);
loop {
let (stream, _) = listener
.accept()
.await
.map_err(|e| format!("net.serve_fn accept: {e}"))?;
let io = TokioIo::new(stream);
let program = Arc::clone(&program);
let policy = policy.clone();
let closure = closure.clone();
tokio::spawn(async move {
let program2 = Arc::clone(&program);
let policy2 = policy.clone();
let closure2 = closure.clone();
let svc = service_fn(move |req: hyper::Request<hyper::body::Incoming>| {
let program = Arc::clone(&program2);
let policy = policy2.clone();
let closure = closure2.clone();
async move {
let (parts, body) = req.into_parts();
let body_bytes = body
.collect()
.await
.map(|c| c.to_bytes())
.unwrap_or_default();
let result = if inline_vm {
let lex_req = build_request_value_parts(&parts, &body_bytes);
let handler = DefaultHandler::new(policy)
.with_program(Arc::clone(&program));
let mut vm = Vm::with_handler(&program, Box::new(handler));
let scope = vm.enter_request_scope();
let r = vm.invoke_closure_value(closure, vec![lex_req]);
let r = r.map(|v| unpack_response(&mut vm, &v));
vm.exit_request_scope(scope);
Ok(r)
} else {
tokio::task::spawn_blocking(move || {
let lex_req = build_request_value_parts(&parts, &body_bytes);
let handler = DefaultHandler::new(policy)
.with_program(Arc::clone(&program));
let mut vm = Vm::with_handler(&program, Box::new(handler));
let scope = vm.enter_request_scope();
let r = vm.invoke_closure_value(closure, vec![lex_req]);
let r = r.map(|v| unpack_response(&mut vm, &v));
vm.exit_request_scope(scope);
r
})
.await
};
Ok::<_, std::convert::Infallible>(match result {
Ok(Ok(unpacked)) => build_hyper_response(unpacked),
Ok(Err(e)) => error_response(500, &format!("internal error: {e}")),
Err(e) => error_response(500, &format!("task panicked: {e}")),
})
}
});
let result = if http2 {
auto::Builder::new(TokioExecutor::new())
.serve_connection(io, svc)
.await
.map_err(|e| e.to_string())
} else {
http1::Builder::new()
.serve_connection(io, svc)
.await
.map_err(|e| e.to_string())
};
if let Err(e) = result {
eprintln!("net.serve_fn: connection error: {e}");
}
});
}
})
}
#[derive(Clone, Debug)]
pub(crate) enum RouteSeg {
Literal(String),
Param(String),
}
pub(super) fn compile_path_pattern(pat: &str) -> Result<Vec<RouteSeg>, String> {
if pat.is_empty() {
return Err("path pattern must be non-empty (use \"/\" for the root)".into());
}
if !pat.starts_with('/') {
return Err(format!("path pattern must start with '/' (got {pat:?})"));
}
let mut segs = Vec::new();
for raw in pat.split('/') {
if let Some(name) = raw.strip_prefix(':') {
if name.is_empty() {
return Err(format!(
":-segment in pattern {pat:?} must have a name (e.g. :id)"
));
}
segs.push(RouteSeg::Param(name.to_string()));
} else {
segs.push(RouteSeg::Literal(raw.to_string()));
}
}
Ok(segs)
}
pub(super) fn match_path_pattern(
segs: &[RouteSeg],
path: &str,
) -> Option<std::collections::BTreeMap<lex_bytecode::MapKey, Value>> {
let path_segs: Vec<&str> = path.split('/').collect();
if path_segs.len() != segs.len() {
return None;
}
let mut params = std::collections::BTreeMap::new();
for (pat, p) in segs.iter().zip(path_segs.iter()) {
match pat {
RouteSeg::Literal(lit) => {
if lit != p {
return None;
}
}
RouteSeg::Param(name) => {
params.insert(
lex_bytecode::MapKey::Str(name.clone()),
Value::Str((*p).into()),
);
}
}
}
Some(params)
}
pub(super) fn decode_routes_arg(
v: Value,
) -> Result<Vec<(String, Vec<RouteSeg>, Value)>, String> {
let list = match v {
Value::List(xs) => xs,
_ => return Err("net.serve_routed: routes must be a List".into()),
};
let mut out = Vec::with_capacity(list.len());
for (i, item) in list.into_iter().enumerate() {
let tup = match item {
Value::Tuple(xs) if xs.len() == 3 => xs,
other => return Err(format!(
"net.serve_routed: route #{i} must be a (method, pattern, handler) 3-tuple, got {other:?}"
)),
};
let mut it = tup.into_iter();
let method_raw = match it.next() {
Some(Value::Str(s)) => s.to_string(),
_ => return Err(format!("net.serve_routed: route #{i} method must be Str")),
};
let method = if method_raw == "*" { method_raw } else { method_raw.to_uppercase() };
let pattern = match it.next() {
Some(Value::Str(s)) => s.to_string(),
_ => return Err(format!("net.serve_routed: route #{i} path-pattern must be Str")),
};
let segs = compile_path_pattern(&pattern)
.map_err(|e| format!("net.serve_routed: route #{i} ({pattern:?}): {e}"))?;
let closure = match it.next() {
Some(c @ Value::Closure { .. }) => c,
_ => return Err(format!("net.serve_routed: route #{i} handler must be a closure")),
};
out.push((method, segs, closure));
}
Ok(out)
}
pub(crate) fn dispatch_route<'a>(
routes: &'a [(String, Vec<RouteSeg>, Value)],
req_method: &str,
req_path: &str,
) -> Option<(&'a Value, std::collections::BTreeMap<lex_bytecode::MapKey, Value>)> {
let req_method_upper = req_method.to_ascii_uppercase();
for (m, segs, closure) in routes {
if m != "*" && m != &req_method_upper {
continue;
}
if let Some(params) = match_path_pattern(segs, req_path) {
return Some((closure, params));
}
}
None
}
pub(crate) fn stamp_path_params(
req: &mut Value,
params: std::collections::BTreeMap<lex_bytecode::MapKey, Value>,
) {
if let Value::Record { fields: rec, .. } = req {
rec.insert("path_params".into(), Value::Map(params));
}
}
pub(super) fn serve_http_routed(
port: u16,
routes: Vec<(String, Vec<RouteSeg>, Value)>,
fallback: Value,
program: Arc<Program>,
policy: Policy,
opts: ServeOpts,
) -> Result<Value, String> {
use http_body_util::BodyExt as _;
use hyper::server::conn::http1;
use hyper::service::service_fn;
use hyper_util::rt::{TokioExecutor, TokioIo};
use hyper_util::server::conn::auto;
use tokio::net::TcpListener as TokioTcpListener;
let inline_vm = opts.inline_vm;
let http2 = opts.http2;
let host = opts.host.clone();
let routes = Arc::new(routes);
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.map_err(|e| format!("net.serve_routed: tokio runtime: {e}"))?;
rt.block_on(async move {
let listener = TokioTcpListener::bind((host.as_str(), port))
.await
.map_err(|e| format!("net.serve_routed bind {host}:{port}: {e}"))?;
eprintln!(
"net.serve_routed: listening on http://{host}:{port} ({} routes{}{})",
routes.len(),
if inline_vm { ", inline-vm" } else { "" },
if http2 { ", http1+http2" } else { "" }
);
loop {
let (stream, _) = listener
.accept()
.await
.map_err(|e| format!("net.serve_routed accept: {e}"))?;
let io = TokioIo::new(stream);
let program = Arc::clone(&program);
let policy = policy.clone();
let routes = Arc::clone(&routes);
let fallback = fallback.clone();
tokio::spawn(async move {
let program2 = Arc::clone(&program);
let policy2 = policy.clone();
let routes2 = Arc::clone(&routes);
let fallback2 = fallback.clone();
let svc = service_fn(move |req: hyper::Request<hyper::body::Incoming>| {
let program = Arc::clone(&program2);
let policy = policy2.clone();
let routes = Arc::clone(&routes2);
let fallback = fallback2.clone();
async move {
let (parts, body) = req.into_parts();
let body_bytes = body
.collect()
.await
.map(|c| c.to_bytes())
.unwrap_or_default();
let method = parts.method.as_str().to_string();
let path = match parts.uri.path() {
"" => "/".to_string(),
p => p.to_string(),
};
let result = if inline_vm {
let mut lex_req = build_request_value_parts(&parts, &body_bytes);
let (closure, params) = match dispatch_route(&routes, &method, &path) {
Some((c, p)) => (c.clone(), p),
None => (fallback.clone(), std::collections::BTreeMap::new()),
};
stamp_path_params(&mut lex_req, params);
let handler = DefaultHandler::new(policy)
.with_program(Arc::clone(&program));
let mut vm = Vm::with_handler(&program, Box::new(handler));
let r = vm.invoke_closure_value(closure, vec![lex_req]);
Ok(r.map(|v| unpack_response(&mut vm, &v)))
} else {
tokio::task::spawn_blocking(move || {
let mut lex_req = build_request_value_parts(&parts, &body_bytes);
let (closure, params) = match dispatch_route(&routes, &method, &path) {
Some((c, p)) => (c.clone(), p),
None => (fallback.clone(), std::collections::BTreeMap::new()),
};
stamp_path_params(&mut lex_req, params);
let handler = DefaultHandler::new(policy)
.with_program(Arc::clone(&program));
let mut vm = Vm::with_handler(&program, Box::new(handler));
let r = vm.invoke_closure_value(closure, vec![lex_req]);
r.map(|v| unpack_response(&mut vm, &v))
})
.await
};
Ok::<_, std::convert::Infallible>(match result {
Ok(Ok(unpacked)) => build_hyper_response(unpacked),
Ok(Err(e)) => error_response(500, &format!("internal error: {e}")),
Err(e) => error_response(500, &format!("task panicked: {e}")),
})
}
});
let result = if http2 {
auto::Builder::new(TokioExecutor::new())
.serve_connection(io, svc)
.await
.map_err(|e| e.to_string())
} else {
http1::Builder::new()
.serve_connection(io, svc)
.await
.map_err(|e| e.to_string())
};
if let Err(e) = result {
eprintln!("net.serve_routed: connection error: {e}");
}
});
}
})
}
pub(super) fn env_inline_vm() -> bool {
match std::env::var("LEX_NET_INLINE_VM") {
Ok(v) => {
let s = v.trim().to_ascii_lowercase();
s == "1" || s == "true"
}
Err(_) => false,
}
}
#[derive(Debug, Clone)]
pub(crate) struct ServeOpts {
pub(crate) http2: bool,
pub(crate) inline_vm: bool,
pub(crate) host: String,
}
impl ServeOpts {
pub(super) fn from_env() -> Self {
Self {
http2: env_http2(),
inline_vm: env_inline_vm(),
host: "0.0.0.0".to_string(),
}
}
pub(super) fn lex_defaults() -> Self {
Self {
http2: false,
inline_vm: false,
host: "0.0.0.0".to_string(),
}
}
pub(super) fn to_value(&self) -> Value {
let mut rec = indexmap::IndexMap::new();
rec.insert("http2".to_string(), Value::Bool(self.http2));
rec.insert("inline_vm".to_string(), Value::Bool(self.inline_vm));
rec.insert("host".to_string(), Value::Str(self.host.clone().into()));
Value::record_dynamic(rec)
}
}
pub(super) fn decode_serve_opts(v: &Value) -> Result<ServeOpts, String> {
let rec = match v {
Value::Record { fields: r, .. } => r,
other => return Err(format!("opts must be a Record, got {other:?}")),
};
let http2 = match rec.get("http2") {
Some(Value::Bool(b)) => *b,
_ => return Err("opts.http2 must be Bool".into()),
};
let inline_vm = match rec.get("inline_vm") {
Some(Value::Bool(b)) => *b,
_ => return Err("opts.inline_vm must be Bool".into()),
};
let host = match rec.get("host") {
Some(Value::Str(s)) => s.to_string(),
_ => return Err("opts.host must be Str".into()),
};
Ok(ServeOpts { http2, inline_vm, host })
}
pub(super) fn make_tls_config_value(cert_pem: Vec<u8>, key_pem: Vec<u8>) -> Value {
let mut rec = indexmap::IndexMap::new();
rec.insert("cert".into(), Value::Bytes(cert_pem));
rec.insert("key".into(), Value::Bytes(key_pem));
Value::record_dynamic(rec)
}
#[cfg(feature = "quic")]
pub(super) fn decode_tls_config(v: &Value) -> Result<crate::quic::QuicTls, String> {
let rec = match v {
Value::Record { fields: r, .. } => r,
other => return Err(format!("TlsConfig: expected Record, got {other:?}")),
};
let cert = match rec.get("cert") {
Some(Value::Bytes(b)) => b.to_vec(),
_ => return Err("TlsConfig.cert: must be Bytes".into()),
};
let key = match rec.get("key") {
Some(Value::Bytes(b)) => b.to_vec(),
_ => return Err("TlsConfig.key: must be Bytes".into()),
};
Ok(crate::quic::QuicTls { cert_pem: cert, key_pem: key })
}
pub(super) fn dispatch_tls_from_pem_files(
handler: &DefaultHandler,
args: Vec<Value>,
) -> Result<Value, String> {
let cert_path = expect_str(args.first())?.to_string();
let key_path = expect_str(args.get(1))?.to_string();
let cert_resolved = handler.resolve_read_path(&cert_path);
let key_resolved = handler.resolve_read_path(&key_path);
if !handler.policy.allow_fs_read.is_empty() {
let allowed = |p: &std::path::Path| -> bool {
handler.policy.allow_fs_read.iter().any(|a| p.starts_with(a))
};
if !allowed(&cert_resolved) {
return Ok(err(Value::Str(
format!("tls.from_pem_files: cert `{cert_path}` outside --allow-fs-read").into(),
)));
}
if !allowed(&key_resolved) {
return Ok(err(Value::Str(
format!("tls.from_pem_files: key `{key_path}` outside --allow-fs-read").into(),
)));
}
}
let cert = match std::fs::read(&cert_resolved) {
Ok(b) => b,
Err(e) => return Ok(err(Value::Str(format!("read cert {cert_path}: {e}").into()))),
};
let key = match std::fs::read(&key_resolved) {
Ok(b) => b,
Err(e) => return Ok(err(Value::Str(format!("read key {key_path}: {e}").into()))),
};
Ok(ok(make_tls_config_value(cert, key)))
}
#[cfg(feature = "quic")]
pub(super) fn dispatch_tls_self_signed(args: Vec<Value>) -> Result<Value, String> {
let hostname = expect_str(args.first())?.to_string();
match crate::quic::self_signed_pem(&hostname) {
Ok((cert, key)) => Ok(ok(make_tls_config_value(cert, key))),
Err(e) => Ok(err(Value::Str(format!("tls.self_signed: {e}").into()))),
}
}
#[cfg(not(feature = "quic"))]
pub(super) fn dispatch_tls_self_signed(_args: Vec<Value>) -> Result<Value, String> {
Ok(err(Value::Str(
"tls.self_signed: lex-runtime was compiled without the `quic` feature (needed for rcgen)".into(),
)))
}
impl DefaultHandler {
#[cfg(feature = "quic")]
pub(super) fn dispatch_serve_quic_named(&self, args: Vec<Value>) -> Result<Value, String> {
let port = match args.first() {
Some(Value::Int(n)) if (0..=65535).contains(n) => *n as u16,
_ => return Err("net.serve_quic(port, tls, handler): port must be Int 0..=65535".into()),
};
let tls = decode_tls_config(args.get(1)
.ok_or_else(|| "net.serve_quic(port, tls, handler): missing tls".to_string())?)?;
let handler_name = expect_str(args.get(2))?.to_string();
let program = self.program.clone()
.ok_or_else(|| "net.serve_quic requires a Program reference; use DefaultHandler::with_program".to_string())?;
let policy = self.policy.clone();
crate::quic::serve_http3_named(port, handler_name, tls, program, policy, ServeOpts::from_env())
}
#[cfg(feature = "quic")]
pub(super) fn dispatch_serve_quic_fn(&self, args: Vec<Value>) -> Result<Value, String> {
let port = match args.first() {
Some(Value::Int(n)) if (0..=65535).contains(n) => *n as u16,
_ => return Err("net.serve_quic_fn(port, tls, handler): port must be Int 0..=65535".into()),
};
let tls = decode_tls_config(args.get(1)
.ok_or_else(|| "net.serve_quic_fn(port, tls, handler): missing tls".to_string())?)?;
let closure = match args.into_iter().nth(2) {
Some(c @ Value::Closure { .. }) => c,
_ => return Err("net.serve_quic_fn(port, tls, handler): handler must be a closure".into()),
};
let program = self.program.clone()
.ok_or_else(|| "net.serve_quic_fn requires a Program reference; use DefaultHandler::with_program".to_string())?;
let policy = self.policy.clone();
crate::quic::serve_http3_fn(port, closure, tls, program, policy, ServeOpts::from_env())
}
#[cfg(feature = "quic")]
pub(super) fn dispatch_serve_quic_routed(&self, args: Vec<Value>) -> Result<Value, String> {
let port = match args.first() {
Some(Value::Int(n)) if (0..=65535).contains(n) => *n as u16,
_ => return Err("net.serve_quic_routed(port, tls, routes, fallback): port must be Int 0..=65535".into()),
};
let tls = decode_tls_config(args.get(1)
.ok_or_else(|| "net.serve_quic_routed(port, tls, routes, fallback): missing tls".to_string())?)?;
let routes_val = args.get(2).cloned()
.ok_or_else(|| "net.serve_quic_routed(port, tls, routes, fallback): missing routes".to_string())?;
let fallback = match args.into_iter().nth(3) {
Some(c @ Value::Closure { .. }) => c,
_ => return Err("net.serve_quic_routed(port, tls, routes, fallback): fallback must be a closure".into()),
};
let routes = decode_routes_arg(routes_val)?;
let program = self.program.clone()
.ok_or_else(|| "net.serve_quic_routed requires a Program reference; use DefaultHandler::with_program".to_string())?;
let policy = self.policy.clone();
crate::quic::serve_http3_routed(port, routes, fallback, tls, program, policy, ServeOpts::from_env())
}
#[cfg(not(feature = "quic"))]
pub(super) fn dispatch_serve_quic_named(&self, _args: Vec<Value>) -> Result<Value, String> {
Err("net.serve_quic: lex-runtime was compiled without the `quic` feature (needed for quinn + h3)".into())
}
#[cfg(not(feature = "quic"))]
pub(super) fn dispatch_serve_quic_fn(&self, _args: Vec<Value>) -> Result<Value, String> {
Err("net.serve_quic_fn: lex-runtime was compiled without the `quic` feature (needed for quinn + h3)".into())
}
#[cfg(not(feature = "quic"))]
pub(super) fn dispatch_serve_quic_routed(&self, _args: Vec<Value>) -> Result<Value, String> {
Err("net.serve_quic_routed: lex-runtime was compiled without the `quic` feature (needed for quinn + h3)".into())
}
}
pub(super) fn env_http2() -> bool {
match std::env::var("LEX_NET_HTTP2") {
Ok(v) => {
let s = v.trim().to_ascii_lowercase();
s == "1" || s == "true"
}
Err(_) => false,
}
}
pub(crate) fn build_request_value_parts(
parts: &hyper::http::request::Parts,
body: &bytes::Bytes,
) -> Value {
let method = parts.method.as_str().to_string();
let path = parts.uri.path().to_string();
let query = parts.uri.query().map(str::to_string).unwrap_or_default();
let mut headers_map = std::collections::BTreeMap::new();
for (name, val) in &parts.headers {
if let Ok(v) = val.to_str() {
headers_map.insert(
lex_bytecode::MapKey::Str(name.as_str().to_ascii_lowercase()),
Value::Str(v.to_string().into()),
);
}
}
let body_str = String::from_utf8_lossy(body).into_owned();
let mut rec = indexmap::IndexMap::new();
rec.insert("method".into(), Value::Str(method.into()));
rec.insert("path".into(), Value::Str(path.into()));
rec.insert("query".into(), Value::Str(query.into()));
rec.insert("body".into(), Value::Str(body_str.into()));
rec.insert("headers".into(), Value::Map(headers_map));
rec.insert("path_params".into(), Value::Map(std::collections::BTreeMap::new()));
Value::record_dynamic(rec)
}
pub(super) fn build_request_value_tiny(req: &mut tiny_http::Request) -> Value {
let method = format!("{:?}", req.method()).to_uppercase();
let url = req.url().to_string();
let (path, query) = match url.split_once('?') {
Some((p, q)) => (p.to_string(), q.to_string()),
None => (url, String::new()),
};
let mut headers_map = std::collections::BTreeMap::new();
for h in req.headers() {
headers_map.insert(
lex_bytecode::MapKey::Str(h.field.as_str().as_str().to_ascii_lowercase()),
Value::Str(h.value.as_str().to_string().into()),
);
}
let mut body = String::new();
let _ = req.as_reader().read_to_string(&mut body);
let mut rec = indexmap::IndexMap::new();
rec.insert("method".into(), Value::Str(method.into()));
rec.insert("path".into(), Value::Str(path.into()));
rec.insert("query".into(), Value::Str(query.into()));
rec.insert("body".into(), Value::Str(body.into()));
rec.insert("headers".into(), Value::Map(headers_map));
rec.insert("path_params".into(), Value::Map(std::collections::BTreeMap::new()));
Value::record_dynamic(rec)
}
pub(crate) fn unpack_response(vm: &mut Vm, v: &Value) -> UnpackedResponse {
if !matches!(v, Value::Record { .. } | Value::ArenaRecord { .. }) {
return (
500,
ResponseBodyOut::Str(format!("handler returned non-record: {v:?}")),
vec![],
);
}
let status = vm.get_record_field(v, "status").and_then(|s| match s {
Value::Int(n) => Some(n as u16),
_ => None,
}).unwrap_or(200);
let body = match vm.get_record_field(v, "body") {
Some(Value::Variant { name, mut args }) if args.len() == 1 => {
let inner = args.pop().unwrap();
match (name.as_str(), inner) {
("BodyStr", Value::Str(s)) => ResponseBodyOut::Str(s.to_string()),
("BodyStream", iter_v) => {
let drained = materialize_lazy_iter(vm, iter_v);
ResponseBodyOut::TextChunks(drain_iter_str(&drained))
}
("BodyBytes", iter_v) => {
let drained = materialize_lazy_iter(vm, iter_v);
ResponseBodyOut::BytesChunks(drain_iter_bytes(&drained))
}
_ => ResponseBodyOut::Str(String::new()),
}
}
Some(Value::Str(s)) => ResponseBodyOut::Str(s.to_string()),
_ => ResponseBodyOut::Str(String::new()),
};
let headers: Vec<(String, String)> = match vm.get_record_field(v, "headers") {
Some(Value::Map(hmap)) => hmap.iter().filter_map(|(k, val)| {
if let (lex_bytecode::MapKey::Str(name), Value::Str(s)) = (k, val) {
Some((name.clone(), s.to_string()))
} else {
None
}
}).collect(),
_ => vec![],
};
(status, body, headers)
}
pub(super) type HyperRespBody =
http_body_util::combinators::BoxBody<bytes::Bytes, std::convert::Infallible>;
pub(super) fn build_hyper_response(
(status, body, headers): UnpackedResponse,
) -> hyper::Response<HyperRespBody> {
use http_body_util::BodyExt as _;
let boxed_body: HyperRespBody = match body {
ResponseBodyOut::Str(s) => {
http_body_util::Full::new(bytes::Bytes::from(s.into_bytes())).boxed()
}
ResponseBodyOut::TextChunks(chunks) | ResponseBodyOut::BytesChunks(chunks) => {
HyperChunkedBody::from(chunks).boxed()
}
};
let mut builder = hyper::Response::builder().status(status);
for (name, val) in headers {
builder = builder.header(name, val);
}
builder
.body(boxed_body)
.unwrap_or_else(|_| error_response(500, "response build error"))
}
pub(super) fn error_response(status: u16, msg: &str) -> hyper::Response<HyperRespBody> {
use http_body_util::BodyExt as _;
hyper::Response::builder()
.status(status)
.body(
http_body_util::Full::new(bytes::Bytes::from(msg.to_owned()))
.boxed(),
)
.unwrap_or_else(|_| {
use http_body_util::BodyExt as _;
hyper::Response::new(http_body_util::Empty::new().map_err(|e| match e {}).boxed())
})
}
pub(super) struct HyperChunkedBody {
pub(super) chunks: std::collections::VecDeque<Vec<u8>>,
}
impl From<Vec<Vec<u8>>> for HyperChunkedBody {
fn from(chunks: Vec<Vec<u8>>) -> Self {
Self {
chunks: chunks.into_iter().filter(|c| !c.is_empty()).collect(),
}
}
}
impl hyper::body::Body for HyperChunkedBody {
type Data = bytes::Bytes;
type Error = std::convert::Infallible;
fn poll_frame(
mut self: std::pin::Pin<&mut Self>,
_cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Result<hyper::body::Frame<Self::Data>, Self::Error>>> {
match self.chunks.pop_front() {
Some(chunk) => std::task::Poll::Ready(Some(Ok(hyper::body::Frame::data(
bytes::Bytes::from(chunk),
)))),
None => std::task::Poll::Ready(None),
}
}
}
pub(super) fn respond_with_body_tls(
req: tiny_http::Request,
status: u16,
body: ResponseBodyOut,
headers: Vec<(String, String)>,
) {
let tiny_headers: Vec<tiny_http::Header> = headers
.into_iter()
.filter_map(|(name, val)| format!("{name}: {val}").parse::<tiny_http::Header>().ok())
.collect();
match body {
ResponseBodyOut::Str(s) => {
let mut response = tiny_http::Response::from_string(s).with_status_code(status);
for h in tiny_headers {
response.add_header(h);
}
let _ = req.respond(response);
}
ResponseBodyOut::TextChunks(chunks) | ResponseBodyOut::BytesChunks(chunks) => {
let reader = ChunkReader::new(chunks);
let response = tiny_http::Response::new(
tiny_http::StatusCode(status),
tiny_headers,
reader,
None,
None,
);
let _ = req.respond(response);
}
}
}
pub(crate) type UnpackedResponse = (u16, ResponseBodyOut, Vec<(String, String)>);
pub(crate) enum ResponseBodyOut {
Str(String),
TextChunks(Vec<Vec<u8>>),
BytesChunks(Vec<Vec<u8>>),
}
pub(super) fn drain_iter_str(v: &Value) -> Vec<Vec<u8>> {
match v {
Value::Variant { name, args }
if name == "__IterEager" && args.len() == 2 =>
{
if let (Value::List(items), Value::Int(idx)) = (&args[0], &args[1]) {
items.iter().skip(*idx as usize).filter_map(|item| {
if let Value::Str(s) = item { Some(s.as_bytes().to_vec()) } else { None }
}).collect()
} else {
Vec::new()
}
}
_ => Vec::new(),
}
}
pub(super) fn drain_iter_bytes(v: &Value) -> Vec<Vec<u8>> {
match v {
Value::Variant { name, args }
if name == "__IterEager" && args.len() == 2 =>
{
if let (Value::List(items), Value::Int(idx)) = (&args[0], &args[1]) {
items.iter().skip(*idx as usize).filter_map(|item| {
if let Value::List(ints) = item {
Some(ints.iter().filter_map(|i| match i {
Value::Int(n) => Some((*n & 0xff) as u8),
_ => None,
}).collect::<Vec<u8>>())
} else {
None
}
}).collect()
} else {
Vec::new()
}
}
_ => Vec::new(),
}
}
pub(super) fn materialize_lazy_iter(vm: &mut Vm, v: Value) -> Value {
let mut current = v;
let mut items: Vec<Value> = Vec::new();
loop {
match current {
Value::Variant { name, args } if name == "__IterLazy" && args.len() == 2 => {
let seed = args[0].clone();
let step = args[1].clone();
match vm.invoke_closure_value(step.clone(), vec![seed]) {
Ok(Value::Variant { name: opt, args: opt_args })
if opt == "None" =>
{
let _ = opt_args;
break;
}
Ok(Value::Variant { name: opt, args: opt_args })
if opt == "Some" && opt_args.len() == 1 =>
{
if let Value::Tuple(pair) = &opt_args[0] {
if pair.len() == 2 {
items.push(pair[0].clone());
current = Value::Variant {
name: "__IterLazy".to_string(),
args: vec![pair[1].clone(), step],
};
continue;
}
}
break;
}
_ => break,
}
}
other => {
if items.is_empty() {
return other;
}
let _ = other;
break;
}
}
}
Value::Variant {
name: "__IterEager".to_string(),
args: vec![
Value::List(items.into_iter().collect()),
Value::Int(0),
],
}
}
pub(super) struct ChunkReader {
pub(super) chunks: std::collections::VecDeque<Vec<u8>>,
pub(super) cursor: usize,
}
impl ChunkReader {
pub(super) fn new(chunks: Vec<Vec<u8>>) -> Self {
Self {
chunks: chunks.into_iter().filter(|c| !c.is_empty()).collect(),
cursor: 0,
}
}
}
impl std::io::Read for ChunkReader {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
loop {
let Some(front) = self.chunks.front() else {
return Ok(0);
};
let remaining = &front[self.cursor..];
if remaining.is_empty() {
self.chunks.pop_front();
self.cursor = 0;
continue;
}
let n = remaining.len().min(buf.len());
buf[..n].copy_from_slice(&remaining[..n]);
self.cursor += n;
if self.cursor >= front.len() {
self.chunks.pop_front();
self.cursor = 0;
}
return Ok(n);
}
}
}
#[cfg(test)]
mod unpack_response_tests {
use super::*;
use std::sync::Arc;
use indexmap::IndexMap;
use lex_bytecode::{Const, Op, Program, Value};
use lex_bytecode::program::{Function, ZERO_BODY_HASH};
use lex_bytecode::vm::Vm;
fn build_arena_response_program() -> Arc<Program> {
let constants = vec![
Const::FieldName("status".into()), Const::FieldName("body".into()), Const::Int(200), Const::VariantName("BodyStr".into()), Const::Str("hello".into()), ];
let mut function_names = IndexMap::new();
function_names.insert("handler".to_string(), 0);
Arc::new(Program {
constants,
functions: vec![Function {
name: "handler".into(),
arity: 0,
locals_count: 0,
code: vec![
Op::PushConst(2), Op::PushConst(4), Op::MakeVariant { name_idx: 3, arity: 1 }, Op::AllocArenaRecord { shape_idx: 0, field_count: 2 }, Op::Return,
],
effects: vec![],
body_hash: ZERO_BODY_HASH,
refinements: vec![],
field_ic_sites: 0,
}],
function_names,
module_aliases: IndexMap::new(),
entry: Some(0),
record_shapes: vec![vec![0, 1]], })
}
#[test]
fn unpack_response_reads_arena_record_via_slab() {
let p = build_arena_response_program();
let mut vm = Vm::new(&p);
let scope = vm.enter_request_scope();
let resp = vm.invoke(p.function_names["handler"], vec![]).unwrap();
assert!(matches!(resp, Value::ArenaRecord { .. }),
"expected ArenaRecord (slab path), got {resp:?}");
let (status, body, headers) = unpack_response(&mut vm, &resp);
vm.exit_request_scope(scope);
assert_eq!(status, 200);
assert!(headers.is_empty());
match body {
ResponseBodyOut::Str(s) => assert_eq!(s, "hello"),
_ => panic!("expected BodyStr"),
}
}
#[test]
fn unpack_response_reads_heap_record() {
let p = build_arena_response_program();
let mut vm = Vm::new(&p);
let resp = vm.invoke(p.function_names["handler"], vec![]).unwrap();
assert!(matches!(resp, Value::Record { .. }),
"expected heap Record (fallback path), got {resp:?}");
let (status, body, headers) = unpack_response(&mut vm, &resp);
assert_eq!(status, 200);
assert!(headers.is_empty());
match body {
ResponseBodyOut::Str(s) => assert_eq!(s, "hello"),
_ => panic!("expected BodyStr"),
}
}
#[test]
fn unpack_response_falls_back_to_500_on_non_record() {
let p = build_arena_response_program();
let mut vm = Vm::new(&p);
let v = Value::Int(7);
let (status, _body, _headers) = unpack_response(&mut vm, &v);
assert_eq!(status, 500);
}
}