use crate::host::{invoke, with_host, IoTask, JsObj};
use fusevm::Value;
use hpack::{Decoder, Encoder};
use indexmap::IndexMap;
use rustls::pki_types::{CertificateDer, PrivateKeyDer};
use rustls::{ServerConfig, ServerConnection, StreamOwned};
use std::collections::HashMap;
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::mpsc::Sender;
use std::sync::Arc;
pub const METHODS: &[&str] = &[
"createSecureServer",
"createServer",
"connect",
"getDefaultSettings",
"getPackedSettings",
"getUnpackedSettings",
];
pub const SERVER_METHODS: &[&str] = &["listen", "close", "address", "setTimeout"];
pub const STREAM_METHODS: &[&str] = &[
"respond",
"write",
"end",
"close",
"setEncoding",
"setTimeout",
"pause",
"resume",
"writeHead",
"setHeader",
"getHeader",
"removeHeader",
];
pub const SESSION_METHODS: &[&str] = &[
"settings",
"ping",
"goaway",
"close",
"destroy",
"ref",
"unref",
"setTimeout",
];
const PREFACE: &[u8] = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n";
const FT_DATA: u8 = 0x0;
const FT_HEADERS: u8 = 0x1;
const FT_PRIORITY: u8 = 0x2;
const FT_RST_STREAM: u8 = 0x3;
const FT_SETTINGS: u8 = 0x4;
const FT_PING: u8 = 0x6;
const FT_GOAWAY: u8 = 0x7;
const FT_WINDOW_UPDATE: u8 = 0x8;
const FT_CONTINUATION: u8 = 0x9;
const FL_END_STREAM: u8 = 0x1;
const FL_ACK: u8 = 0x1; const FL_END_HEADERS: u8 = 0x4;
const FL_PADDED: u8 = 0x8;
const FL_PRIORITY: u8 = 0x20;
const MAX_FRAME_SIZE: usize = 16384;
static NEXT_SERVER_ID: AtomicU64 = AtomicU64::new(1);
static NEXT_STREAM_KEY: AtomicU64 = AtomicU64::new(1);
static NEXT_SESSION_KEY: AtomicU64 = AtomicU64::new(1);
fn next_server_id() -> u64 {
NEXT_SERVER_ID.fetch_add(1, Ordering::Relaxed)
}
fn next_stream_key() -> u64 {
NEXT_STREAM_KEY.fetch_add(1, Ordering::Relaxed)
}
fn next_session_key() -> u64 {
NEXT_SESSION_KEY.fetch_add(1, Ordering::Relaxed)
}
enum H2Cmd {
Respond {
stream_id: u32,
headers: Vec<(String, String)>,
end: bool,
},
Data {
stream_id: u32,
data: Vec<u8>,
end: bool,
},
Close { stream_id: u32 },
Goaway,
}
struct H2ServerRec {
emitter: Value,
stop: Arc<AtomicBool>,
}
struct H2StreamRec {
emitter: Value,
tx: Sender<H2Cmd>,
stream_id: u32,
responded: bool,
}
struct H2SessionRec {
#[allow(dead_code)]
emitter: Value,
tx: Sender<H2Cmd>,
}
#[derive(Default)]
struct H2State {
servers: HashMap<u64, H2ServerRec>,
streams: HashMap<u64, H2StreamRec>,
sessions: HashMap<u64, H2SessionRec>,
}
thread_local! {
static H2: std::cell::RefCell<H2State> = std::cell::RefCell::new(H2State::default());
static PENDING_CONFIGS: std::cell::RefCell<Vec<(Value, Arc<ServerConfig>)>> =
const { std::cell::RefCell::new(Vec::new()) };
}
fn get_prop(recv: &Value, key: &str) -> Option<Value> {
with_host(|h| match h.get(recv) {
Some(JsObj::Object(p)) => p.get(key).cloned(),
_ => None,
})
}
fn set_prop(recv: &Value, key: &str, val: Value) {
with_host(|h| {
if let Some(JsObj::Object(p)) = h.get_mut(recv) {
p.insert(key.to_string(), val);
}
});
}
fn u64_prop(recv: &Value, key: &str) -> Option<u64> {
get_prop(recv, key).map(|v| with_host(|h| h.to_number(&v)) as u64)
}
fn emitter_dispatch(recv: &Value, method: &str, args: &[Value]) -> Option<Result<Value, String>> {
super::events::METHODS
.contains(&method)
.then(|| super::events::instance_call(recv, method, args.to_vec()))
}
fn value_bytes(v: Option<&Value>) -> Vec<u8> {
let Some(v) = v else { return Vec::new() };
let is_buffer =
with_host(|h| matches!(h.get(v), Some(JsObj::Object(p)) if p.contains_key("@@bytes")));
if is_buffer {
return with_host(|h| match h.get(v) {
Some(JsObj::Object(p)) => match p.get("@@bytes").and_then(|b| h.get(b)) {
Some(JsObj::Array(items)) => items.iter().map(|x| h.to_number(x) as u8).collect(),
_ => Vec::new(),
},
_ => Vec::new(),
});
}
with_host(|h| h.str_of(v)).into_bytes()
}
fn object_pairs(obj: &Value) -> Vec<(String, String)> {
with_host(|h| match h.get(obj) {
Some(JsObj::Object(p)) => p
.iter()
.filter(|(k, _)| !k.starts_with("@@") && !k.starts_with('#'))
.map(|(k, v)| (k.clone(), h.str_of(v)))
.collect(),
_ => Vec::new(),
})
}
pub fn call(method: &str, args: &[Value]) -> Option<Result<Value, String>> {
Some(match method {
"createSecureServer" => create_secure_server(args),
"createServer" => Err(
"Error: http2.createServer (cleartext h2c / prior-knowledge) \
is not implemented in node-js; use http2.createSecureServer (h2 over TLS)"
.to_string(),
),
"connect" => Err(
"Error: http2.connect (HTTP/2 client) is not implemented in node-js; \
only the HTTP/2 server (http2.createSecureServer) is implemented"
.to_string(),
),
"getDefaultSettings" => Ok(default_settings_object()),
"getPackedSettings" => Ok(pack_settings(args)),
"getUnpackedSettings" => unpack_settings(args),
_ => return None,
})
}
pub fn constant(name: &str) -> Option<Value> {
match name {
"constants" => Some(constants_object()),
"Http2ServerRequest" => Some(with_host(|h| {
h.alloc(JsObj::Builtin("Http2ServerRequest".into()))
})),
"Http2ServerResponse" => Some(with_host(|h| {
h.alloc(JsObj::Builtin("Http2ServerResponse".into()))
})),
_ => None,
}
}
fn constants_object() -> Value {
with_host(|h| {
let mut m = IndexMap::new();
let put_i = |m: &mut IndexMap<String, Value>, k: &str, v: i64| {
m.insert(k.to_string(), Value::Float(v as f64));
};
put_i(&mut m, "HTTP_STATUS_OK", 200);
put_i(&mut m, "HTTP_STATUS_NO_CONTENT", 204);
put_i(&mut m, "HTTP_STATUS_MOVED_PERMANENTLY", 301);
put_i(&mut m, "HTTP_STATUS_FOUND", 302);
put_i(&mut m, "HTTP_STATUS_NOT_MODIFIED", 304);
put_i(&mut m, "HTTP_STATUS_BAD_REQUEST", 400);
put_i(&mut m, "HTTP_STATUS_UNAUTHORIZED", 401);
put_i(&mut m, "HTTP_STATUS_FORBIDDEN", 403);
put_i(&mut m, "HTTP_STATUS_NOT_FOUND", 404);
put_i(&mut m, "HTTP_STATUS_INTERNAL_SERVER_ERROR", 500);
put_i(&mut m, "NGHTTP2_NO_ERROR", 0x0);
put_i(&mut m, "NGHTTP2_PROTOCOL_ERROR", 0x1);
put_i(&mut m, "NGHTTP2_INTERNAL_ERROR", 0x2);
put_i(&mut m, "NGHTTP2_FLOW_CONTROL_ERROR", 0x3);
put_i(&mut m, "NGHTTP2_SETTINGS_TIMEOUT", 0x4);
put_i(&mut m, "NGHTTP2_STREAM_CLOSED", 0x5);
put_i(&mut m, "NGHTTP2_FRAME_SIZE_ERROR", 0x6);
put_i(&mut m, "NGHTTP2_REFUSED_STREAM", 0x7);
put_i(&mut m, "NGHTTP2_CANCEL", 0x8);
put_i(&mut m, "NGHTTP2_COMPRESSION_ERROR", 0x9);
put_i(&mut m, "NGHTTP2_ENHANCE_YOUR_CALM", 0xb);
put_i(&mut m, "NGHTTP2_SETTINGS_HEADER_TABLE_SIZE", 0x1);
put_i(&mut m, "NGHTTP2_SETTINGS_ENABLE_PUSH", 0x2);
put_i(&mut m, "NGHTTP2_SETTINGS_MAX_CONCURRENT_STREAMS", 0x3);
put_i(&mut m, "NGHTTP2_SETTINGS_INITIAL_WINDOW_SIZE", 0x4);
put_i(&mut m, "NGHTTP2_SETTINGS_MAX_FRAME_SIZE", 0x5);
put_i(&mut m, "NGHTTP2_SETTINGS_MAX_HEADER_LIST_SIZE", 0x6);
let hdr =
|m: &mut IndexMap<String, Value>, k: &str, v: &str, h: &mut crate::host::JsHost| {
let s = h.new_str(v);
m.insert(k.to_string(), s);
};
hdr(&mut m, "HTTP2_HEADER_STATUS", ":status", h);
hdr(&mut m, "HTTP2_HEADER_METHOD", ":method", h);
hdr(&mut m, "HTTP2_HEADER_AUTHORITY", ":authority", h);
hdr(&mut m, "HTTP2_HEADER_SCHEME", ":scheme", h);
hdr(&mut m, "HTTP2_HEADER_PATH", ":path", h);
hdr(&mut m, "HTTP2_HEADER_CONTENT_TYPE", "content-type", h);
hdr(&mut m, "HTTP2_HEADER_CONTENT_LENGTH", "content-length", h);
hdr(&mut m, "HTTP2_METHOD_GET", "GET", h);
hdr(&mut m, "HTTP2_METHOD_POST", "POST", h);
h.new_object(m)
})
}
fn default_settings_object() -> Value {
with_host(|h| {
let mut m = IndexMap::new();
m.insert("headerTableSize".into(), Value::Float(4096.0));
m.insert("enablePush".into(), Value::Bool(false));
m.insert("initialWindowSize".into(), Value::Float(65535.0));
m.insert("maxFrameSize".into(), Value::Float(MAX_FRAME_SIZE as f64));
m.insert("maxConcurrentStreams".into(), Value::Float(100.0));
h.new_object(m)
})
}
fn push_setting(out: &mut Vec<u8>, id: u16, val: u32) {
out.extend_from_slice(&id.to_be_bytes());
out.extend_from_slice(&val.to_be_bytes());
}
fn pack_settings(args: &[Value]) -> Value {
let settings = args.first().cloned().unwrap_or(Value::Undef);
let num = |key: &str| -> Option<u32> {
get_prop(&settings, key)
.filter(|v| !matches!(v, Value::Undef))
.map(|v| with_host(|h| h.to_number(&v)) as u32)
};
let mut out: Vec<u8> = Vec::new();
if let Some(v) = num("headerTableSize") {
push_setting(&mut out, 0x1, v);
}
if let Some(p) = get_prop(&settings, "enablePush").filter(|v| !matches!(v, Value::Undef)) {
let on = with_host(|h| h.truthy(&p));
push_setting(&mut out, 0x2, u32::from(on));
}
if let Some(v) = num("maxConcurrentStreams") {
push_setting(&mut out, 0x3, v);
}
if let Some(v) = num("initialWindowSize") {
push_setting(&mut out, 0x4, v);
}
if let Some(v) = num("maxFrameSize") {
push_setting(&mut out, 0x5, v);
}
if let Some(v) = num("maxHeaderListSize").or_else(|| num("maxHeaderSize")) {
push_setting(&mut out, 0x6, v);
}
super::buffer::from_bytes(&out)
}
fn unpack_settings(args: &[Value]) -> Result<Value, String> {
let bytes = value_bytes(args.first());
if bytes.len() % 6 != 0 {
return Err("RangeError [ERR_HTTP2_INVALID_PACKED_SETTINGS_LENGTH]: \
Packed settings length must be a multiple of six"
.to_string());
}
let mut m = IndexMap::new();
for chunk in bytes.chunks_exact(6) {
let id = u16::from_be_bytes([chunk[0], chunk[1]]);
let val = u32::from_be_bytes([chunk[2], chunk[3], chunk[4], chunk[5]]);
match id {
0x1 => {
m.insert("headerTableSize".to_string(), Value::Float(val as f64));
}
0x2 => {
m.insert("enablePush".to_string(), Value::Bool(val != 0));
}
0x3 => {
m.insert("maxConcurrentStreams".to_string(), Value::Float(val as f64));
}
0x4 => {
m.insert("initialWindowSize".to_string(), Value::Float(val as f64));
}
0x5 => {
m.insert("maxFrameSize".to_string(), Value::Float(val as f64));
}
0x6 => {
m.insert("maxHeaderSize".to_string(), Value::Float(val as f64));
m.insert("maxHeaderListSize".to_string(), Value::Float(val as f64));
}
_ => {}
}
}
Ok(with_host(|h| h.new_object(m)))
}
fn create_secure_server(args: &[Value]) -> Result<Value, String> {
let mut options: Option<Value> = None;
let mut handler: Option<Value> = None;
for a in args {
if with_host(|h| crate::host::is_callable(h, a)) {
handler = Some(a.clone());
} else if matches!(a, Value::Obj(_)) {
options = Some(a.clone());
}
}
let opts = options.ok_or_else(|| {
crate::host::type_error("http2.createSecureServer requires options with `key` and `cert`")
})?;
let cert = value_bytes(get_prop(&opts, "cert").as_ref());
let key = value_bytes(get_prop(&opts, "key").as_ref());
if cert.is_empty() || key.is_empty() {
return Err(crate::host::type_error(
"http2.createSecureServer requires `key` and `cert`",
));
}
let config = build_h2_server_config(&cert, &key)?;
let server = new_emitter_object("Http2Server", IndexMap::new());
if let Some(cb) = handler {
super::events::instance_call(&server, "on", vec![with_host(|h| h.new_str("request")), cb])?;
}
PENDING_CONFIGS.with(|p| p.borrow_mut().push((server.clone(), config)));
Ok(server)
}
fn build_h2_server_config(cert_pem: &[u8], key_pem: &[u8]) -> Result<Arc<ServerConfig>, String> {
let certs: Vec<CertificateDer<'static>> = rustls_pemfile::certs(&mut &cert_pem[..])
.collect::<Result<_, _>>()
.map_err(|e| format!("Error: http2: bad certificate PEM: {e}"))?;
if certs.is_empty() {
return Err("Error: http2: no certificates found in `cert`".to_string());
}
let key: PrivateKeyDer<'static> = rustls_pemfile::private_key(&mut &key_pem[..])
.map_err(|e| format!("Error: http2: bad private key PEM: {e}"))?
.ok_or_else(|| "Error: http2: no private key found in `key`".to_string())?;
let mut cfg = ServerConfig::builder()
.with_no_client_auth()
.with_single_cert(certs, key)
.map_err(|e| format!("Error: http2: invalid key/cert: {e}"))?;
cfg.alpn_protocols = vec![b"h2".to_vec()];
Ok(Arc::new(cfg))
}
fn take_pending_config(server: &Value) -> Option<Arc<ServerConfig>> {
PENDING_CONFIGS.with(|p| {
let mut p = p.borrow_mut();
p.iter()
.position(|(s, _)| s == server)
.map(|pos| p.remove(pos).1)
})
}
pub fn instance_call(
tag: &str,
recv: &Value,
method: &str,
args: Vec<Value>,
) -> Result<Value, String> {
match tag {
"Http2Server" => server_call(recv, method, args),
"Http2Stream" => stream_call(recv, method, args),
"Http2Session" => session_call(recv, method, args),
_ => Err(crate::host::type_error(&format!(
"{method} is not a function"
))),
}
}
fn server_call(recv: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
if let Some(r) = emitter_dispatch(recv, method, &args) {
return r;
}
match method {
"listen" => server_listen(recv, &args),
"close" => server_close(recv, &args),
"address" => Ok(get_prop(recv, "@@address").unwrap_or(Value::Undef)),
"setTimeout" => Ok(recv.clone()),
_ => Err(crate::host::type_error(&format!(
"server.{method} is not a function"
))),
}
}
fn server_listen(recv: &Value, args: &[Value]) -> Result<Value, String> {
let port = with_host(|h| args.first().map(|v| h.to_number(v)).unwrap_or(0.0)) as u16;
let mut host = "0.0.0.0".to_string();
let mut cb: Option<Value> = None;
for a in &args[1.min(args.len())..] {
if with_host(|h| h.as_str(a)).is_some() {
host = with_host(|h| h.str_of(a));
} else if with_host(|h| crate::host::is_callable(h, a)) {
cb = Some(a.clone());
}
}
let config = take_pending_config(recv)
.ok_or_else(|| crate::host::type_error("http2 server has no secure context"))?;
let listener = TcpListener::bind((host.as_str(), port))
.map_err(|e| format!("Error: listen EADDRINUSE: {e}"))?;
let local = listener.local_addr().ok();
let id = next_server_id();
set_prop(recv, "@@serverid", Value::Float(id as f64));
if let Some(addr) = local {
let mut a = IndexMap::new();
a.insert("port".into(), Value::Float(addr.port() as f64));
a.insert(
"address".into(),
with_host(|h| h.new_str(addr.ip().to_string())),
);
a.insert(
"family".into(),
with_host(|h| h.new_str(if addr.is_ipv6() { "IPv6" } else { "IPv4" })),
);
let addr_obj = with_host(|h| h.new_object(a));
set_prop(recv, "@@address", addr_obj);
}
let stop = Arc::new(AtomicBool::new(false));
H2.with(|s| {
s.borrow_mut().servers.insert(
id,
H2ServerRec {
emitter: recv.clone(),
stop: stop.clone(),
},
);
});
with_host(|h| h.incr_handle());
let io_tx = with_host(|h| h.io_sender());
listener.set_nonblocking(true).ok();
std::thread::spawn(move || loop {
if stop.load(Ordering::Acquire) {
break;
}
match listener.accept() {
Ok((stream, _addr)) => {
let cfg = config.clone();
let tx = io_tx.clone();
std::thread::spawn(move || serve_connection(id, stream, cfg, tx));
}
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
std::thread::sleep(std::time::Duration::from_millis(5));
}
Err(_) => break,
}
});
let server = recv.clone();
let _ = with_host(|h| h.io_sender()).send(Box::new(move || {
super::events::instance_call(&server, "emit", vec![with_host(|h| h.new_str("listening"))])?;
if let Some(cb) = cb {
invoke(&cb, Vec::new(), None)?;
}
Ok(())
}));
Ok(recv.clone())
}
fn server_close(recv: &Value, args: &[Value]) -> Result<Value, String> {
if let Some(id) = u64_prop(recv, "@@serverid") {
let rec = H2.with(|s| s.borrow_mut().servers.remove(&id));
if let Some(rec) = rec {
rec.stop.store(true, Ordering::Release);
with_host(|h| h.decr_handle());
let _ = with_host(|h| h.io_sender()).send(Box::new(|| Ok(())));
}
}
if let Some(cb) = args
.first()
.filter(|v| with_host(|h| crate::host::is_callable(h, v)))
{
invoke(cb, Vec::new(), None)?;
}
super::events::instance_call(recv, "emit", vec![with_host(|h| h.new_str("close"))])?;
Ok(recv.clone())
}
fn serve_connection(
server_id: u64,
mut sock: TcpStream,
config: Arc<ServerConfig>,
io_tx: Sender<IoTask>,
) {
sock.set_nonblocking(false).ok();
let mut conn = match ServerConnection::new(config) {
Ok(c) => c,
Err(_) => return,
};
if conn.complete_io(&mut sock).is_err() {
return;
}
let is_h2 = conn.alpn_protocol().map(|p| p == b"h2").unwrap_or(false);
let mut stream = StreamOwned::new(conn, sock);
if !is_h2 {
stream.conn.send_close_notify();
let _ = stream.flush();
let _ = stream.sock.shutdown(std::net::Shutdown::Both);
return;
}
let (tx, rx) = std::sync::mpsc::channel::<H2Cmd>();
let session_key = next_session_key();
{
let tx_sess = tx.clone();
let _ = io_tx.send(Box::new(move || {
on_session(server_id, session_key, tx_sess)
}));
}
if h2_debug() {
eprintln!("[http2] connection {session_key}: ALPN h2 negotiated, starting framing loop");
}
let mut ok = true;
ok &= write_frame(&mut stream, FT_SETTINGS, 0, 0, &[]).is_ok();
ok &= stream.flush().is_ok();
stream
.sock
.set_read_timeout(Some(std::time::Duration::from_millis(20)))
.ok();
let mut decoder = Decoder::new();
let mut encoder = Encoder::new();
let mut inbuf: Vec<u8> = Vec::new();
let mut got_preface = false;
let mut max_stream_id: u32 = 0;
let mut id_to_key: HashMap<u32, u64> = HashMap::new();
let mut buf = [0u8; MAX_FRAME_SIZE];
'conn: loop {
if !ok {
break;
}
loop {
match rx.try_recv() {
Ok(cmd) => {
if h2_debug() {
eprintln!(
"[http2] connection {session_key}: draining {}",
cmd_name(&cmd)
);
}
if !apply_cmd(&mut stream, &mut encoder, max_stream_id, cmd) {
if h2_debug() {
eprintln!("[http2] connection {session_key}: write failed, closing");
}
break 'conn;
}
}
Err(std::sync::mpsc::TryRecvError::Empty) => break,
Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
}
}
match stream.read(&mut buf) {
Ok(0) => {
if h2_debug() {
eprintln!("[http2] connection {session_key}: read EOF (Ok 0)");
}
break;
}
Ok(n) => inbuf.extend_from_slice(&buf[..n]),
Err(ref e)
if matches!(
e.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
) =>
{
continue;
}
Err(e) => {
if h2_debug() {
eprintln!(
"[http2] connection {session_key}: read error {:?} ({e})",
e.kind()
);
}
break;
}
}
if !got_preface {
if inbuf.len() < PREFACE.len() {
continue;
}
if &inbuf[..PREFACE.len()] != PREFACE {
break; }
inbuf.drain(..PREFACE.len());
got_preface = true;
}
loop {
if inbuf.len() < 9 {
break;
}
let len =
((inbuf[0] as usize) << 16) | ((inbuf[1] as usize) << 8) | (inbuf[2] as usize);
if inbuf.len() < 9 + len {
break; }
let ftype = inbuf[3];
let flags = inbuf[4];
let stream_id =
u32::from_be_bytes([inbuf[5], inbuf[6], inbuf[7], inbuf[8]]) & 0x7fff_ffff;
let payload: Vec<u8> = inbuf[9..9 + len].to_vec();
inbuf.drain(..9 + len);
if h2_debug() {
eprintln!(
"[http2] connection {session_key}: recv frame type={ftype} flags={flags:#04x} \
stream={stream_id} len={len}"
);
}
match ftype {
FT_SETTINGS => {
if flags & FL_ACK == 0 {
if write_frame(&mut stream, FT_SETTINGS, FL_ACK, 0, &[])
.and_then(|_| stream.flush())
.is_err()
{
break 'conn;
}
}
}
FT_PING => {
if flags & FL_ACK == 0 {
if write_frame(&mut stream, FT_PING, FL_ACK, 0, &payload)
.and_then(|_| stream.flush())
.is_err()
{
break 'conn;
}
}
}
FT_HEADERS => {
if stream_id > max_stream_id {
max_stream_id = stream_id;
}
if flags & FL_END_HEADERS == 0 {
continue;
}
let block = strip_headers_padding_priority(&payload, flags);
let decoded = match decoder.decode(&block) {
Ok(d) => d,
Err(_) => continue,
};
let headers: Vec<(String, String)> = decoded
.into_iter()
.map(|(k, v)| {
(
String::from_utf8_lossy(&k).into_owned(),
String::from_utf8_lossy(&v).into_owned(),
)
})
.collect();
let end_stream = flags & FL_END_STREAM != 0;
let key = next_stream_key();
id_to_key.insert(stream_id, key);
let tx_stream = tx.clone();
let _ = io_tx.send(Box::new(move || {
on_headers(server_id, key, stream_id, headers, end_stream, tx_stream)
}));
}
FT_DATA => {
if let Some(&key) = id_to_key.get(&stream_id) {
let data = strip_data_padding(&payload, flags);
let end_stream = flags & FL_END_STREAM != 0;
let _ = io_tx.send(Box::new(move || on_data(key, data, end_stream)));
}
}
FT_GOAWAY => break 'conn,
FT_WINDOW_UPDATE | FT_PRIORITY | FT_RST_STREAM | FT_CONTINUATION => {}
_ => {}
}
}
}
if h2_debug() {
eprintln!(
"[http2] connection {session_key}: framing loop exited, sending GOAWAY + closing"
);
}
let mut goaway = Vec::with_capacity(8);
goaway.extend_from_slice(&(max_stream_id & 0x7fff_ffff).to_be_bytes());
goaway.extend_from_slice(&0u32.to_be_bytes()); let _ = write_frame(&mut stream, FT_GOAWAY, 0, 0, &goaway);
let _ = stream.flush();
stream.conn.send_close_notify();
let _ = stream.flush();
let _ = stream.sock.shutdown(std::net::Shutdown::Both);
let stream_keys: Vec<u64> = id_to_key.values().copied().collect();
let _ = io_tx.send(Box::new(move || on_session_close(session_key, stream_keys)));
}
fn cmd_name(cmd: &H2Cmd) -> &'static str {
match cmd {
H2Cmd::Respond { .. } => "respond(HEADERS)",
H2Cmd::Data { .. } => "data(DATA)",
H2Cmd::Close { .. } => "close(RST_STREAM)",
H2Cmd::Goaway => "goaway(GOAWAY)",
}
}
fn apply_cmd(
stream: &mut StreamOwned<ServerConnection, TcpStream>,
encoder: &mut Encoder<'_>,
max_stream_id: u32,
cmd: H2Cmd,
) -> bool {
match cmd {
H2Cmd::Respond {
stream_id,
headers,
end,
} => {
let block = encode_header_block(encoder, &headers);
if h2_debug() {
eprintln!(
"[http2] write HEADERS stream={stream_id} end_stream={end} \
hpack_len={} headers={headers:?}",
block.len()
);
}
let flags = FL_END_HEADERS | if end { FL_END_STREAM } else { 0 };
write_frame(stream, FT_HEADERS, flags, stream_id, &block)
.and_then(|_| stream.flush())
.is_ok()
}
H2Cmd::Data {
stream_id,
data,
end,
} => {
if h2_debug() {
eprintln!(
"[http2] write DATA stream={stream_id} len={} end_stream={end}",
data.len()
);
}
send_data(stream, stream_id, &data, end)
}
H2Cmd::Close { stream_id } => {
write_frame(stream, FT_RST_STREAM, 0, stream_id, &0u32.to_be_bytes())
.and_then(|_| stream.flush())
.is_ok()
}
H2Cmd::Goaway => {
let mut g = Vec::with_capacity(8);
g.extend_from_slice(&(max_stream_id & 0x7fff_ffff).to_be_bytes());
g.extend_from_slice(&0u32.to_be_bytes());
let _ = write_frame(stream, FT_GOAWAY, 0, 0, &g);
let _ = stream.flush();
false
}
}
}
fn send_data(
stream: &mut StreamOwned<ServerConnection, TcpStream>,
stream_id: u32,
data: &[u8],
end: bool,
) -> bool {
if data.is_empty() {
let flags = if end { FL_END_STREAM } else { 0 };
return write_frame(stream, FT_DATA, flags, stream_id, &[])
.and_then(|_| stream.flush())
.is_ok();
}
let chunks: Vec<&[u8]> = data.chunks(MAX_FRAME_SIZE).collect();
let last = chunks.len() - 1;
for (i, chunk) in chunks.iter().enumerate() {
let flags = if end && i == last { FL_END_STREAM } else { 0 };
if write_frame(stream, FT_DATA, flags, stream_id, chunk).is_err() {
return false;
}
}
stream.flush().is_ok()
}
fn write_frame<W: Write>(
w: &mut W,
ftype: u8,
flags: u8,
stream_id: u32,
payload: &[u8],
) -> std::io::Result<()> {
let len = payload.len();
let mut hdr = [0u8; 9];
hdr[0] = (len >> 16) as u8;
hdr[1] = (len >> 8) as u8;
hdr[2] = len as u8;
hdr[3] = ftype;
hdr[4] = flags;
hdr[5..9].copy_from_slice(&(stream_id & 0x7fff_ffff).to_be_bytes());
w.write_all(&hdr)?;
w.write_all(payload)?;
Ok(())
}
fn encode_header_block(encoder: &mut Encoder<'_>, headers: &[(String, String)]) -> Vec<u8> {
let owned: Vec<(Vec<u8>, Vec<u8>)> = headers
.iter()
.map(|(k, v)| (k.as_bytes().to_vec(), v.as_bytes().to_vec()))
.collect();
encoder.encode(owned.iter().map(|(k, v)| (k.as_slice(), v.as_slice())))
}
fn strip_headers_padding_priority(payload: &[u8], flags: u8) -> Vec<u8> {
let mut start = 0usize;
let mut pad_len = 0usize;
if flags & FL_PADDED != 0 && !payload.is_empty() {
pad_len = payload[0] as usize;
start = 1;
}
if flags & FL_PRIORITY != 0 {
start += 5; }
let end = payload.len().saturating_sub(pad_len);
if start > end {
return Vec::new();
}
payload[start..end].to_vec()
}
fn strip_data_padding(payload: &[u8], flags: u8) -> Vec<u8> {
if flags & FL_PADDED != 0 && !payload.is_empty() {
let pad_len = payload[0] as usize;
let end = payload.len().saturating_sub(pad_len);
if 1 <= end {
return payload[1..end].to_vec();
}
return Vec::new();
}
payload.to_vec()
}
fn on_session(server_id: u64, session_key: u64, tx: Sender<H2Cmd>) -> Result<(), String> {
with_host(|h| h.incr_handle());
let server = H2.with(|s| {
s.borrow()
.servers
.get(&server_id)
.map(|r| r.emitter.clone())
});
let Some(server) = server else { return Ok(()) };
let mut extra = IndexMap::new();
extra.insert("@@h2session".into(), Value::Float(session_key as f64));
let session = new_emitter_object("Http2Session", extra);
H2.with(|s| {
s.borrow_mut().sessions.insert(
session_key,
H2SessionRec {
emitter: session.clone(),
tx,
},
);
});
if let Err(e) = super::events::instance_call(
&server,
"emit",
vec![with_host(|h| h.new_str("session")), session],
) {
report_handler_error("session", &e);
}
Ok(())
}
fn on_session_close(session_key: u64, stream_keys: Vec<u64>) -> Result<(), String> {
H2.with(|s| {
let mut st = s.borrow_mut();
st.sessions.remove(&session_key);
for k in &stream_keys {
st.streams.remove(k);
}
});
with_host(|h| h.decr_handle());
let _ = with_host(|h| h.io_sender()).send(Box::new(|| Ok(())));
Ok(())
}
fn on_headers(
server_id: u64,
stream_key: u64,
stream_id: u32,
headers: Vec<(String, String)>,
end_stream: bool,
tx: Sender<H2Cmd>,
) -> Result<(), String> {
let server = H2.with(|s| {
s.borrow()
.servers
.get(&server_id)
.map(|r| r.emitter.clone())
});
let Some(server) = server else { return Ok(()) };
let headers_obj = with_host(|h| {
let mut m = IndexMap::new();
for (k, v) in &headers {
m.insert(k.clone(), h.new_str(v.clone()));
}
h.new_object(m)
});
let mut extra = IndexMap::new();
extra.insert("@@h2key".into(), Value::Float(stream_key as f64));
extra.insert("id".into(), Value::Float(stream_id as f64));
let stream_obj = new_emitter_object("Http2Stream", extra);
H2.with(|s| {
s.borrow_mut().streams.insert(
stream_key,
H2StreamRec {
emitter: stream_obj.clone(),
tx,
stream_id,
responded: false,
},
);
});
let method = header_value(&headers, ":method").unwrap_or_else(|| "GET".to_string());
let path = header_value(&headers, ":path").unwrap_or_else(|| "/".to_string());
if h2_debug() {
eprintln!("[http2] dispatch stream={stream_id} {method} {path} (end_stream={end_stream})");
}
if let Err(e) = super::events::instance_call(
&server,
"emit",
vec![
with_host(|h| h.new_str("stream")),
stream_obj.clone(),
headers_obj.clone(),
],
) {
report_handler_error("stream", &e);
return Ok(());
}
let req = super::events::new_emitter();
set_prop(&req, "method", with_host(|h| h.new_str(method)));
set_prop(&req, "url", with_host(|h| h.new_str(path)));
set_prop(&req, "headers", headers_obj);
if let Err(e) = super::events::instance_call(
&server,
"emit",
vec![with_host(|h| h.new_str("request")), req, stream_obj.clone()],
) {
report_handler_error("request", &e);
return Ok(());
}
if end_stream {
if let Err(e) =
super::events::instance_call(&stream_obj, "emit", vec![with_host(|h| h.new_str("end"))])
{
report_handler_error("stream.end", &e);
}
}
Ok(())
}
fn report_handler_error(event: &str, err: &str) {
eprintln!("http2: uncaught error in '{event}' handler: {err}");
}
fn h2_debug() -> bool {
std::env::var_os("HTTP2_DEBUG").is_some()
}
fn on_data(stream_key: u64, data: Vec<u8>, end_stream: bool) -> Result<(), String> {
let stream = H2.with(|s| {
s.borrow()
.streams
.get(&stream_key)
.map(|r| r.emitter.clone())
});
let Some(stream) = stream else { return Ok(()) };
if !data.is_empty() {
let chunk = super::buffer::from_bytes(&data);
if let Err(e) = super::events::instance_call(
&stream,
"emit",
vec![with_host(|h| h.new_str("data")), chunk],
) {
report_handler_error("data", &e);
return Ok(());
}
}
if end_stream {
if let Err(e) =
super::events::instance_call(&stream, "emit", vec![with_host(|h| h.new_str("end"))])
{
report_handler_error("end", &e);
}
}
Ok(())
}
fn header_value(headers: &[(String, String)], name: &str) -> Option<String> {
headers
.iter()
.find(|(k, _)| k == name)
.map(|(_, v)| v.clone())
}
fn stream_key_of(recv: &Value) -> Option<u64> {
u64_prop(recv, "@@h2key")
}
fn stream_call(recv: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
if let Some(r) = emitter_dispatch(recv, method, &args) {
return r;
}
match method {
"respond" => {
let hdrs = args.first().map(object_pairs).unwrap_or_default();
let end = args
.get(1)
.and_then(|o| get_prop(o, "endStream"))
.map(|v| with_host(|h| h.truthy(&v)))
.unwrap_or(false);
do_respond(recv, hdrs, end)?;
Ok(recv.clone())
}
"writeHead" => {
let status =
with_host(|h| args.first().map(|v| h.to_number(v)).unwrap_or(200.0)) as u32;
let mut hdrs: Vec<(String, String)> = Vec::new();
for a in args.iter().skip(1) {
if matches!(a, Value::Obj(_)) {
hdrs = object_pairs(a);
break;
}
}
hdrs.insert(0, (":status".to_string(), status.to_string()));
do_respond(recv, hdrs, false)?;
Ok(recv.clone())
}
"setHeader" => {
let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)));
let v = with_host(|h| h.str_of(&args.get(1).cloned().unwrap_or(Value::Undef)));
let bag = pending_headers_obj(recv);
set_prop(&bag, &k, with_host(|h| h.new_str(v)));
Ok(Value::Undef)
}
"getHeader" => {
let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)));
Ok(get_prop(recv, "@@pendingHeaders")
.and_then(|bag| get_prop(&bag, &k))
.unwrap_or(Value::Undef))
}
"removeHeader" => {
let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)));
if let Some(bag) = get_prop(recv, "@@pendingHeaders") {
with_host(|h| {
if let Some(JsObj::Object(p)) = h.get_mut(&bag) {
p.shift_remove(&k);
}
});
}
Ok(Value::Undef)
}
"write" => {
ensure_responded(recv)?;
let bytes = value_bytes(args.first());
send_stream_data(recv, bytes, false);
Ok(Value::Bool(true))
}
"end" => {
ensure_responded(recv)?;
let bytes = args
.first()
.filter(|v| !matches!(v, Value::Undef))
.map(|v| value_bytes(Some(v)))
.unwrap_or_default();
send_stream_data(recv, bytes, true);
super::events::instance_call(recv, "emit", vec![with_host(|h| h.new_str("finish"))])?;
Ok(recv.clone())
}
"close" => {
if let Some(key) = stream_key_of(recv) {
let sent = H2.with(|s| {
s.borrow().streams.get(&key).map(|r| {
let _ = r.tx.send(H2Cmd::Close {
stream_id: r.stream_id,
});
})
});
let _ = sent;
}
Ok(recv.clone())
}
"setEncoding" | "setTimeout" | "pause" | "resume" => Ok(recv.clone()),
_ => Err(crate::host::type_error(&format!(
"stream.{method} is not a function"
))),
}
}
fn pending_headers_obj(recv: &Value) -> Value {
if let Some(bag) = get_prop(recv, "@@pendingHeaders") {
return bag;
}
let bag = with_host(|h| h.new_object(IndexMap::new()));
set_prop(recv, "@@pendingHeaders", bag.clone());
bag
}
fn do_respond(recv: &Value, mut headers: Vec<(String, String)>, end: bool) -> Result<(), String> {
if let Some(bag) = get_prop(recv, "@@pendingHeaders") {
for (k, v) in object_pairs(&bag) {
if !headers.iter().any(|(hk, _)| hk.eq_ignore_ascii_case(&k)) {
headers.push((k, v));
}
}
}
let status = headers
.iter()
.find(|(k, _)| k == ":status")
.map(|(_, v)| v.clone())
.unwrap_or_else(|| "200".to_string());
let mut ordered: Vec<(String, String)> = vec![(":status".to_string(), status)];
for (k, v) in headers.into_iter() {
if k == ":status" {
continue;
}
ordered.push((k.to_ascii_lowercase(), v));
}
let Some(key) = stream_key_of(recv) else {
return Ok(());
};
H2.with(|s| {
if let Some(r) = s.borrow_mut().streams.get_mut(&key) {
if !r.responded {
r.responded = true;
let _ = r.tx.send(H2Cmd::Respond {
stream_id: r.stream_id,
headers: ordered,
end,
});
}
}
});
Ok(())
}
fn ensure_responded(recv: &Value) -> Result<(), String> {
let Some(key) = stream_key_of(recv) else {
return Ok(());
};
let responded = H2.with(|s| {
s.borrow()
.streams
.get(&key)
.map(|r| r.responded)
.unwrap_or(true)
});
if !responded {
do_respond(recv, Vec::new(), false)?;
}
Ok(())
}
fn send_stream_data(recv: &Value, data: Vec<u8>, end: bool) {
if let Some(key) = stream_key_of(recv) {
H2.with(|s| {
if let Some(r) = s.borrow().streams.get(&key) {
let _ = r.tx.send(H2Cmd::Data {
stream_id: r.stream_id,
data,
end,
});
}
});
}
}
fn session_call(recv: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
if let Some(r) = emitter_dispatch(recv, method, &args) {
return r;
}
match method {
"close" | "destroy" | "goaway" => {
if let Some(key) = u64_prop(recv, "@@h2session") {
H2.with(|s| {
if let Some(r) = s.borrow().sessions.get(&key) {
let _ = r.tx.send(H2Cmd::Goaway);
}
});
}
super::events::instance_call(recv, "emit", vec![with_host(|h| h.new_str("close"))])?;
Ok(recv.clone())
}
"settings" | "ping" | "ref" | "unref" | "setTimeout" => Ok(recv.clone()),
_ => Err(crate::host::type_error(&format!(
"session.{method} is not a function"
))),
}
}
pub fn new_emitter_object(tag: &str, mut extra: IndexMap<String, Value>) -> Value {
with_host(|h| {
let on = h.new_object(IndexMap::new());
let once = h.new_object(IndexMap::new());
let mut m = IndexMap::new();
m.insert("@@native".into(), h.new_str(tag));
m.insert("@@on".into(), on);
m.insert("@@once".into(), once);
for (k, v) in extra.drain(..) {
m.insert(k, v);
}
h.new_object(m)
})
}