use base64::prelude::*; use neqo_bin::server::{HttpServer, Runner}; use neqo_common::Bytes; use neqo_common::{event::Provider, qdebug, qerror, qinfo, qtrace, Datagram, Header}; use nss_rs::{generate_ech_keys, init_db, AllowZeroRtt, AntiReplay}; use neqo_http3::{
ConnectUdpRequest, ConnectUdpServerEvent, Error, Http3OrWebTransportStream, Http3Parameters,
Http3Server, Http3ServerEvent, SessionAcceptAction, StreamId, WebTransportRequest,
WebTransportServerEvent,
}; use neqo_transport::server::ConnectionRef; use neqo_transport::{
ConnectionEvent, ConnectionParameters, OutputBatch, RandomConnectionIdGenerator, StreamType,
}; use std::env; use std::pin::Pin; use std::task::{Context, Poll}; use tokio::io::AsyncWriteExt; use tokio::io::ReadBuf; use tokio::task::LocalSet;
use std::cell::RefCell; use std::io; use std::num::NonZeroUsize; use std::path::PathBuf; use std::process::exit; use std::rc::Rc; use std::thread; use std::time::{Duration, Instant};
use cfg_if::cfg_if;
cfg_if! { if#[cfg(not(target_os = "android"))] { use std::sync::mpsc::{channel, Receiver, TryRecvError}; use http_body_util::{BodyExt, Full}; use hyper::header::{HeaderName, HeaderValue}; use hyper::Method; use hyper_util::client::legacy::Client;
}
}
use std::cmp::min; use std::collections::hash_map::DefaultHasher; use std::collections::HashSet; use std::collections::{HashMap, VecDeque}; use std::hash::{Hash, Hasher}; use std::net::SocketAddr;
const HTTP_RESPONSE_WITH_WRONG_FRAME: &[u8] = &[ 0x01, 0x06, 0x00, 0x00, 0xd9, 0x54, 0x01, 0x37, // headers 0x0, 0x3, 0x61, 0x62, 0x63, // the first data frame 0x3, 0x1, 0x5, // a cancel push frame that is not allowed
]; struct Http3TestServer {
server: Http3Server, // This a map from a post request to amount of data ithas been received on the request. // The respons will carry the amount of data received.
posts: HashMap<Http3OrWebTransportStream, usize>,
responses: HashMap<Http3OrWebTransportStream, Vec<u8>>,
connections_to_close: HashMap<Instant, Vec<ConnectionRef>>,
sessions_to_close: HashMap<Instant, Vec<WebTransportRequest>>,
sessions_to_create_stream: Vec<(WebTransportRequest, StreamType, Option<Vec<u8>>)>,
webtransport_bidi_stream: HashSet<Http3OrWebTransportStream>,
wt_unidi_conn_to_stream: HashMap<ConnectionRef, Http3OrWebTransportStream>,
wt_unidi_echo_back: HashMap<Http3OrWebTransportStream, Http3OrWebTransportStream>,
received_datagram: Option<Bytes>, // When true, server will stop processing datagrams after accepting 0-RTT, // simulating a stuck ZERORTT session that never transitions to CONNECTED.
stuck_0rtt_mode: bool,
stuck_0rtt_activated: bool,
}
fn maybe_close_session(&mutself, now: Instant) { for (expires, sessions) inself.sessions_to_close.iter_mut() { if *expires <= now { for s in sessions.iter_mut() {
drop(s.close_session(0, "", now));
}
}
} self.sessions_to_close.retain(|expires, _| *expires >= now);
}
fn maybe_close_connection(&mutself) { let now = Instant::now(); for (expires, connections) inself.connections_to_close.iter_mut() { if *expires <= now { for c in connections.iter_mut() {
c.borrow_mut().close(now, 0x0100, "");
}
}
} self.connections_to_close
.retain(|expires, _| *expires >= now);
}
impl HttpServer for Http3TestServer { fn process_multiple<'a, D: IntoIterator<Item = Datagram<&'a mut [u8]>>>(
&mutself,
dgrams: D,
now: Instant,
max_datagrams: NonZeroUsize,
) -> OutputBatch { // If stuck_0rtt_mode is enabled and we've already processed datagrams once, // stop processing to simulate a connection stuck in ZERORTT state. ifself.stuck_0rtt_mode && self.stuck_0rtt_activated {
qinfo!("Stuck 0-RTT mode active - ignoring datagrams to keep session in ZERORTT"); // Return Callback to keep the server loop running but don't process datagrams return OutputBatch::Callback(Duration::from_millis(100));
}
let output = self.server.process_multiple(dgrams, now, max_datagrams);
// If we just processed datagrams with stuck mode enabled, mark it as activated ifself.stuck_0rtt_mode && !self.stuck_0rtt_activated {
qinfo!("Stuck 0-RTT mode activated - next datagrams will be ignored"); self.stuck_0rtt_activated = true;
}
let output = ifself.sessions_to_close.is_empty() && self.connections_to_close.is_empty() {
output
} else { // In case there are pending sessions to close, use a shorter // timeout to make process_events() to be called earlier. const MIN_INTERVAL: Duration = Duration::from_millis(100);
match output {
OutputBatch::None => OutputBatch::Callback(MIN_INTERVAL),
o @ OutputBatch::DatagramBatch(_) => o,
OutputBatch::Callback(d) => OutputBatch::Callback(min(d, MIN_INTERVAL)),
}
};
// Some responses do not have content-type. This is on purpose to exercise // UnknownDecoder code. let default_ret = b"Hello World".to_vec(); let default_headers = vec![
Header::new(":status", "200"),
Header::new("cache-control", "no-cache"),
Header::new("content-length", default_ret.len().to_string()),
Header::new("x-http3-conn-hash", connection_hash.to_string()),
];
// echo unidirectional input to back to client // need to close or we hang ifself.wt_unidi_echo_back.contains_key(&stream) { let echo_back = self.wt_unidi_echo_back.remove(&stream).unwrap();
echo_back.send_data(&data, now).unwrap();
echo_back.stream_close_send(now).unwrap(); break;
}
#[cfg(not(target_os = "android"))] let output = ifself.response_to_send.is_empty() {
output
} else { // In case there are pending responses to send, make sure a reasonable // callback is returned. const MIN_INTERVAL: Duration = Duration::from_millis(100);
match output {
OutputBatch::None => OutputBatch::Callback(MIN_INTERVAL),
o @ OutputBatch::DatagramBatch(_) => o,
OutputBatch::Callback(d) => OutputBatch::Callback(min(d, MIN_INTERVAL)),
}
};
// Check if we should fallback to 127.0.0.1 before attempting connection let host_without_port = iflet Some(colon_pos) = host_str.rfind(':') {
&host_str[..colon_pos]
} else {
host_str
};
tcp_stream.set_nonblocking(true).unwrap();
qtrace!("tcp_stream to {:?} created", host_hdr);
stream
.send_headers(&[
Header::new(":status", "200"),
Header::new("cache-control", "no-cache"),
])
.unwrap(); self.tcp_streams.insert(
stream.stream_id(),
TcpStream {
send_buffer: VecDeque::new(),
recv_buffer: VecDeque::new(),
stream: tokio::net::TcpStream::from_std(tcp_stream).unwrap(),
send_fin: false,
received_fin: false,
session: stream,
},
);
}
Http3ServerEvent::Data { stream, data, fin } => {
qtrace!("tcp_stream send to server len={}", data.len()); let tcp_stream = self.tcp_streams.get_mut(&stream.stream_id()).unwrap(); // TODO: extend() effectively breaks backpressure.
tcp_stream.send_buffer.extend(data);
tcp_stream.send_fin |= fin;
}
Http3ServerEvent::DataWritable { stream } => {
qtrace!( "Http3ServerEvent::DataWritable streamid={}",
stream.stream_id()
); let tcp_stream = self.tcp_streams.get_mut(&stream.stream_id()).unwrap(); while !tcp_stream.recv_buffer.is_empty() { match stream.send_data(&tcp_stream.recv_buffer.make_contiguous(), now) {
Ok(sent) => {
qtrace!("tcp_stream send to client sent={}", sent); if sent == 0 { // no progress possible right now — stop trying to send in this loop // (could also mark for later retry) break;
}
tcp_stream.recv_buffer.drain(0..sent);
}
Err(e) => {
eprintln!("send_data failed: {:?}", e); break;
}
}
}
}
Http3ServerEvent::ConnectUdp(ConnectUdpServerEvent::NewSession {
session,
headers,
}) => {
session.response(&SessionAcceptAction::Accept, now).unwrap();
let host_hdr = headers.iter().find(|&h| h.name() == ":path").unwrap(); let path_str = host_hdr.value_utf8().unwrap(); let path_parts: Vec<&str> = path_str.split('/').collect();
// Format is /.well-known/masque/udp/{target_host}/{target_port}/ if path_parts.len() < 6 {
panic!("{}", path_str)
}
let target_host = path_parts[4]; let target_port = match path_parts[5].trim_end_matches('/').parse::<u16>() {
Ok(port) => port,
Err(_) => {
panic!("{}", path_str)
}
};
// Replace target_host with 127.0.0.1 for specific hosts let actual_host = match target_host { "foo.example.com" | "alt1.example.com" | "alt2.example.com" => "127.0.0.1",
_ => target_host,
};
let host_port = format!("{}:{}", actual_host, target_port);
qdebug!("CONNECT-UDP to {}", host_port);
let socket = { let s =
socket2::Socket::new(socket2::Domain::IPV4, socket2::Type::DGRAM, None)
.unwrap();
s.bind(&"0.0.0.0:0".parse::<SocketAddr>().unwrap().into())
.unwrap(); let s: std::net::UdpSocket = s.into();
s.connect((actual_host, target_port)).unwrap();
s.set_nonblocking(true).unwrap();
s.into()
};
// Remove failed UDP sockets from the list for stream_id in failed_udp_sockets { iflet Some(socket) = self.udp_sockets.remove(&stream_id) {
qdebug!("Removed failed UDP socket for stream {}", stream_id); // Close the session with an error code let _ = socket
.session
.close_session(0x0100, "UDP socket error", Instant::now());
}
}
let args: Vec<String> = env::args().collect(); if args.len() < 2 {
eprintln!("Wrong arguments.");
exit(1)
}
// Read data from stdin and terminate the server if EOF is detected, which // means that runxpcshelltests.py ended without shutting down the server.
thread::spawn(|| loop { letmut buffer = String::new(); match io::stdin().read_line(&mut buffer) {
Ok(n) => { if n == 0 {
exit(0);
}
}
Err(_) => {
exit(0);
}
}
});
init_db(PathBuf::from(args[1].clone())).unwrap();
let local = LocalSet::new(); letmut hosts = vec![];
let proxy_port = match env::var("MOZ_HTTP3_PROXY_PORT") {
Ok(val) => val.parse::<u16>().unwrap(),
_ => 0,
};
let anti_replay = || {
AntiReplay::new(Instant::now(), Duration::from_secs(10), 7, 14)
.expect("unable to setup anti-replay")
}; let cid_mgr = Rc::new(RefCell::new(RandomConnectionIdGenerator::new(10)));
spawn_server(
Http3TestServer::new(
Http3Server::new(
Instant::now(),
&[" HTTP2 Test Cert"],
PROTOCOLS,
anti_replay(),
cid_mgr.clone(),
Http3Parameters::default()
.max_table_size_encoder(MAX_TABLE_SIZE)
.max_table_size_decoder(MAX_TABLE_SIZE)
.max_blocked_streams(MAX_BLOCKED_STREAMS)
.webtransport(true)
.connection_parameters(ConnectionParameters::default().datagram_size(1200)),
None,
)
.expect("We cannot make a server!"),
), 0,
&local,
&mut hosts,
)?;
spawn_server(
Server(
neqo_transport::server::Server::new(
Instant::now(),
&[" HTTP2 Test Cert"],
PROTOCOLS,
anti_replay(), Box::new(AllowZeroRtt {}),
cid_mgr.clone(),
ConnectionParameters::default(),
)
.expect("We cannot make a server!"),
), 0,
&local,
&mut hosts,
)?;
let ech_config = { letmut server = Http3TestServer::new(
Http3Server::new(
Instant::now(),
&[" HTTP2 Test Cert"],
PROTOCOLS,
anti_replay(),
cid_mgr.clone(),
Http3Parameters::default()
.max_table_size_encoder(MAX_TABLE_SIZE)
.max_table_size_decoder(MAX_TABLE_SIZE)
.max_blocked_streams(MAX_BLOCKED_STREAMS),
None,
)
.expect("We cannot make a server!"),
); let (sk, pk) = generate_ech_keys().unwrap();
server
.server
.enable_ech(ECH_CONFIG_ID, ECH_PUBLIC_NAME, &sk, &pk)
.expect("unable to enable ech"); let ech_config = server.server.ech_config().to_vec();
spawn_server(server, 0, &local, &mut hosts)?;
ech_config
};
spawn_server(
{ let server_config = if env::var("MOZ_HTTP3_MOCHITEST").is_ok() {
("mochitest-cert", 8888)
} else {
(" HTTP2 Test Cert", -1)
}; let server = Http3ReverseProxyServer::new(
Http3Server::new(
Instant::now(),
&[server_config.0],
PROTOCOLS,
anti_replay(),
cid_mgr.clone(),
Http3Parameters::default()
.max_table_size_encoder(MAX_TABLE_SIZE)
.max_table_size_decoder(MAX_TABLE_SIZE)
.max_blocked_streams(MAX_BLOCKED_STREAMS)
.webtransport(true)
.connection_parameters(ConnectionParameters::default().datagram_size(1200)),
None,
)
.expect("We cannot make a server!"),
server_config.1,
);
server
},
proxy_port,
&local,
&mut hosts,
)?;
// Work around until we can use raw-dylibs. #[cfg_attr(target_os = "windows", link(name = "runtimeobject"))] extern"C" {} #[cfg_attr(target_os = "windows", link(name = "propsys"))] extern"C" {} #[cfg_attr(target_os = "windows", link(name = "iphlpapi"))] extern"C" {} #[cfg_attr(target_os = "windows", link(name = "rpcrt4"))] extern"C" {}
Messung V0.5 in Prozent
¤ Diese beiden folgenden Angebotsgruppen bietet das Unternehmen0.56Angebot
(Wie Sie bei der Firma Beratungs- und Dienstleistungen beauftragen können 2026-09-28)
¤
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.