use std::future::Future;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use hyper::body::Bytes;
use hyper::body::Incoming;
use hyper::server::conn::http1;
use hyper::{Method, Request, Response, StatusCode};
use hyper::service::service_fn;
use hyper_util::rt::TokioIo;
use hyper_util::rt::TokioTimer;
use http_body_util::combinators::BoxBody;
use http_body_util::{BodyExt, Empty, Full};
use tokio::net::{TcpListener, TcpStream};
#[cfg(feature = "tls")]
use tokio_rustls::TlsAcceptor;
use crate::body::{MaxBodySize, DEFAULT_MAX_BODY_SIZE};
use crate::cors::CorsConfig;
use crate::error::ServeError;
use crate::handler::{Handler, ResponseBody};
use crate::router::{QueryParams, Router};
use crate::state::State;
const MAX_PATH_LEN: usize = 8_192;
const MAX_QUERY_LEN: usize = 4_096;
const DEFAULT_HEADER_READ_TIMEOUT: Duration = Duration::from_secs(30);
const DEFAULT_MAX_CONNECTIONS: usize = 1024;
#[cfg(feature = "tls")]
const DEFAULT_TLS_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
const ACCEPT_BACKOFF_INITIAL: Duration = Duration::from_millis(10);
const ACCEPT_BACKOFF_MAX: Duration = Duration::from_secs(1);
#[cfg(test)]
thread_local! {
static ERROR_LOG: std::cell::RefCell<Vec<(u16, String)>> = std::cell::RefCell::new(Vec::new());
}
#[cfg(test)]
fn capture_error(code: u16, message: String) {
ERROR_LOG.with(|log| log.borrow_mut().push((code, message)));
}
#[cfg(test)]
fn take_error_log() -> Vec<(u16, String)> {
ERROR_LOG.with(|log| log.borrow_mut().drain(..).collect())
}
trait TcpAccept {
async fn accept(&self) -> std::io::Result<(TcpStream, SocketAddr)>;
}
impl TcpAccept for TcpListener {
async fn accept(&self) -> std::io::Result<(TcpStream, SocketAddr)> {
TcpListener::accept(self).await
}
}
struct Backoff {
delay: Duration,
}
impl Backoff {
fn new() -> Self {
Backoff { delay: ACCEPT_BACKOFF_INITIAL }
}
fn next_delay(&mut self) -> Duration {
let delay = self.delay;
self.delay = (self.delay * 2).min(ACCEPT_BACKOFF_MAX);
delay
}
fn reset(&mut self) {
self.delay = ACCEPT_BACKOFF_INITIAL;
}
}
async fn accept_with_backoff<L: TcpAccept>(
listener: &L,
backoff: &mut Backoff,
) -> (TcpStream, SocketAddr) {
loop {
match listener.accept().await {
Ok(conn) => {
backoff.reset();
return conn;
}
Err(_) => {
tokio::time::sleep(backoff.next_delay()).await;
}
}
}
}
pub type ErrorHandler =
Arc<dyn Fn(StatusCode, &str) -> Response<ResponseBody> + Send + Sync>;
pub struct App<S> {
state: Arc<S>,
router: Arc<Router<S>>,
max_body_size: usize,
pub(crate) header_read_timeout: Duration,
pub(crate) max_connections: usize,
#[cfg(feature = "tls")]
pub(crate) tls_handshake_timeout: Duration,
error_handler: ErrorHandler,
cors_config: Option<CorsConfig>,
}
fn parse_query(query: Option<&str>) -> QueryParams {
let mut map = std::collections::HashMap::new();
if let Some(query) = query {
for pair in query.split('&').filter(|s| !s.is_empty()) {
if let Some((key, value)) = pair.split_once('=') {
let key = decode_query_component(key);
let value = decode_query_component(value);
map.insert(key, value);
} else {
let pair = decode_query_component(pair);
map.insert(pair, String::new());
}
}
}
QueryParams(map)
}
fn decode_query_component(s: &str) -> String {
let with_spaces = s.replace('+', " ");
percent_encoding::percent_decode_str(&with_spaces)
.decode_utf8_lossy()
.into_owned()
}
fn error_response(status: StatusCode, message: &str) -> Response<ResponseBody> {
#[cfg(test)]
capture_error(status.as_u16(), message.to_string());
let client_message = if status.is_server_error() {
"internal server error"
} else {
message
};
let body = serde_json::json!({ "message": client_message });
let json = serde_json::to_string(&body)
.unwrap_or_else(|_| r#"{"message":"internal server error"}"#.to_string());
Response::builder()
.status(status)
.header("content-type", "application/json")
.body(BoxBody::new(Full::new(Bytes::from(json)).map_err(|never: std::convert::Infallible| match never {})))
.expect("status is valid and headers are static ASCII")
}
fn default_error_handler() -> ErrorHandler {
Arc::new(error_response)
}
fn ephemeral_bind_addr() -> SocketAddr {
(std::net::Ipv4Addr::LOCALHOST, 0).into()
}
impl<S: Send + Sync + 'static> App<S> {
pub fn new(state: S) -> Self {
App {
state: Arc::new(state),
router: Arc::new(Router::new()),
max_body_size: DEFAULT_MAX_BODY_SIZE,
header_read_timeout: DEFAULT_HEADER_READ_TIMEOUT,
max_connections: DEFAULT_MAX_CONNECTIONS,
#[cfg(feature = "tls")]
tls_handshake_timeout: DEFAULT_TLS_HANDSHAKE_TIMEOUT,
error_handler: default_error_handler(),
cors_config: None,
}
}
pub fn state_arc(&self) -> Arc<S> {
Arc::clone(&self.state)
}
pub async fn route(&self, req: Request<Incoming>) -> Response<ResponseBody> {
let method = req.method().clone();
let resp = self.route_inner(req).await;
if method == Method::HEAD {
let (parts, _) = resp.into_parts();
Response::from_parts(parts, BoxBody::new(Empty::new().map_err(|never: std::convert::Infallible| match never {})))
} else {
resp
}
}
async fn route_inner(&self, req: Request<Incoming>) -> Response<ResponseBody> {
if req.uri().path().len() > MAX_PATH_LEN {
return (self.error_handler)(StatusCode::BAD_REQUEST, "path too long");
}
if req.uri().query().map(|q| q.len()).unwrap_or(0) > MAX_QUERY_LEN {
return (self.error_handler)(StatusCode::BAD_REQUEST, "query string too long");
}
let method = req.method().clone();
let path = req.uri().path().to_string();
let state = State::from_arc(Arc::clone(&self.state));
let query_params = parse_query(req.uri().query());
let req_origin = req
.headers()
.get("origin")
.and_then(|v| v.to_str().ok())
.map(|s| s.to_string());
if method == Method::OPTIONS && req_origin.is_some() {
if let Some(cfg) = &self.cors_config {
if self.router.path_exists(&path) {
let requested_headers = req
.headers()
.get("access-control-request-headers")
.and_then(|v| v.to_str().ok());
let allowed = self.allowed_methods_with_head(&path);
return cfg.preflight_response(req_origin.as_deref(), requested_headers, &allowed);
}
}
}
let method_to_match = if method == Method::HEAD {
Method::GET
} else {
method.clone()
};
match self.router.match_route(&method_to_match, &path) {
Some((handler, params)) => {
let mut req = req;
req.extensions_mut().insert(query_params);
req.extensions_mut().insert(params);
req.extensions_mut().insert(MaxBodySize(self.max_body_size));
match handler(req, state).await {
Ok(mut resp) => {
if let Some(cfg) = &self.cors_config {
cfg.apply_to_response(&mut resp, req_origin.as_deref());
}
resp
}
Err(e) => (self.error_handler)(
StatusCode::from_u16(e.code)
.unwrap_or(StatusCode::INTERNAL_SERVER_ERROR),
&e.message,
),
}
}
None => {
let allowed = self.allowed_methods_with_head(&path);
if !allowed.is_empty() {
let mut method_strs: Vec<&str> = allowed.iter().map(|m| m.as_str()).collect();
method_strs.sort();
method_strs.dedup();
let allow_header = method_strs.join(", ");
let mut resp = (self.error_handler)(StatusCode::METHOD_NOT_ALLOWED, "method not allowed");
if let Ok(val) = allow_header.parse() {
resp.headers_mut().insert("allow", val);
}
resp
} else {
(self.error_handler)(StatusCode::NOT_FOUND, "not found")
}
}
}
}
fn allowed_methods_with_head(&self, path: &str) -> Vec<Method> {
let mut allowed = self.router.allowed_methods(path);
if allowed.contains(&Method::GET) {
allowed.push(Method::HEAD);
}
allowed
}
pub async fn bind_ephemeral(self) -> Result<u16, ServeError> {
let listener = TcpListener::bind(ephemeral_bind_addr())
.await
.map_err(|e| ServeError::new(500, format!("failed to bind to ephemeral port: {e}")))?;
let port = listener
.local_addr()
.map_err(|e| ServeError::new(500, format!("failed to get assigned port: {e}")))?
.port();
let app = Arc::new(self);
tokio::spawn(async move {
serve_inner(listener, app).await;
});
Ok(port)
}
#[cfg(feature = "tls")]
pub async fn bind_tls_ephemeral(
self,
config: Arc<rustls::ServerConfig>,
) -> Result<u16, ServeError> {
let listener = TcpListener::bind(ephemeral_bind_addr())
.await
.map_err(|e| ServeError::new(500, format!("failed to bind to ephemeral port: {e}")))?;
let port = listener
.local_addr()
.map_err(|e| ServeError::new(500, format!("failed to get assigned port: {e}")))?
.port();
let acceptor = TlsAcceptor::from(config);
let app = Arc::new(self);
tokio::spawn(async move {
serve_tls_inner(listener, app, acceptor).await;
});
Ok(port)
}
pub async fn run<F>(self, listener: TcpListener, shutdown: F) -> Result<(), ServeError>
where
F: Future<Output = ()> + Send + 'static,
{
let app = Arc::new(self);
serve_with_shutdown(listener, app, shutdown).await;
Ok(())
}
pub async fn bind(self, addr: SocketAddr) -> Result<(), ServeError> {
let listener = TcpListener::bind(addr)
.await
.map_err(|e| ServeError::new(500, format!("failed to bind to {addr}: {e}")))?;
self.run(listener, signal_shutdown()).await
}
#[cfg(feature = "tls")]
pub async fn run_tls<F>(
self,
listener: TcpListener,
config: Arc<rustls::ServerConfig>,
shutdown: F,
) -> Result<(), ServeError>
where
F: Future<Output = ()> + Send + 'static,
{
let acceptor = TlsAcceptor::from(config);
let app = Arc::new(self);
serve_tls_with_shutdown(listener, app, acceptor, shutdown).await;
Ok(())
}
#[cfg(feature = "tls")]
pub async fn bind_tls(
self,
addr: SocketAddr,
config: Arc<rustls::ServerConfig>,
) -> Result<(), ServeError> {
let listener = TcpListener::bind(addr)
.await
.map_err(|e| ServeError::new(500, format!("failed to bind to {addr}: {e}")))?;
self.run_tls(listener, config, signal_shutdown()).await
}
}
impl App<()> {
pub fn stateless() -> Self {
App::new(())
}
}
async fn serve_connection<S, IO>(io: IO, app: Arc<App<S>>, header_read_timeout: Duration)
where
S: Send + Sync + 'static,
IO: hyper::rt::Read + hyper::rt::Write + Unpin + 'static,
{
let svc = service_fn(move |req: Request<Incoming>| {
let app = app.clone();
async move {
Ok::<_, hyper::Error>(app.route(req).await)
}
});
let mut builder = http1::Builder::new();
builder.timer(TokioTimer::new());
builder.header_read_timeout(header_read_timeout);
let conn = builder.serve_connection(io, svc);
let _ = conn.await;
}
async fn serve_inner<S: Send + Sync + 'static>(
listener: TcpListener,
app: Arc<App<S>>,
) {
let header_read_timeout = app.header_read_timeout;
let semaphore = Arc::new(tokio::sync::Semaphore::new(app.max_connections));
let mut backoff = Backoff::new();
loop {
let (stream, _) = accept_with_backoff(&listener, &mut backoff).await;
let sem = semaphore.clone();
let permit = match sem.acquire_owned().await {
Ok(p) => p,
Err(_) => continue,
};
let app = app.clone();
tokio::spawn(async move {
let _permit = permit;
serve_connection(TokioIo::new(stream), app, header_read_timeout).await;
});
}
}
#[cfg(feature = "tls")]
async fn serve_tls_inner<S: Send + Sync + 'static>(
listener: TcpListener,
app: Arc<App<S>>,
acceptor: TlsAcceptor,
) {
let header_read_timeout = app.header_read_timeout;
let handshake_timeout = app.tls_handshake_timeout;
let semaphore = Arc::new(tokio::sync::Semaphore::new(app.max_connections));
let mut backoff = Backoff::new();
loop {
let (stream, _) = accept_with_backoff(&listener, &mut backoff).await;
let sem = semaphore.clone();
let permit = match sem.acquire_owned().await {
Ok(p) => p,
Err(_) => continue,
};
let app = app.clone();
let acceptor = acceptor.clone();
tokio::spawn(async move {
let _permit = permit;
let tls_stream = match tokio::time::timeout(handshake_timeout, acceptor.accept(stream)).await {
Ok(Ok(s)) => s,
Ok(Err(_)) | Err(_) => return,
};
serve_connection(TokioIo::new(tls_stream), app, header_read_timeout).await;
});
}
}
async fn accept_and_permit<L: TcpAccept>(
listener: &L,
backoff: &mut Backoff,
semaphore: &Arc<tokio::sync::Semaphore>,
) -> Option<(TcpStream, tokio::sync::OwnedSemaphorePermit)> {
loop {
let (stream, _) = match listener.accept().await {
Ok(conn) => {
backoff.reset();
conn
}
Err(_) => {
tokio::time::sleep(backoff.next_delay()).await;
continue;
}
};
return match semaphore.clone().acquire_owned().await {
Ok(permit) => Some((stream, permit)),
Err(_) => None,
};
}
}
async fn signal_shutdown() {
let mut sigint = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt())
.expect("failed to install SIGINT handler");
let mut sigterm = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("failed to install SIGTERM handler");
tokio::select! {
_ = sigint.recv() => {}
_ = sigterm.recv() => {}
}
}
async fn serve_with_shutdown<S, F>(listener: TcpListener, app: Arc<App<S>>, shutdown: F)
where
S: Send + Sync + 'static,
F: Future<Output = ()> + Send + 'static,
{
let header_read_timeout = app.header_read_timeout;
let semaphore = Arc::new(tokio::sync::Semaphore::new(app.max_connections));
let mut backoff = Backoff::new();
let mut join_set: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
let mut shutdown_pin = std::pin::pin!(shutdown);
let mut shutting_down = false;
loop {
if !shutting_down {
tokio::select! {
accepted = accept_and_permit(&listener, &mut backoff, &semaphore) => {
match accepted {
Some((stream, permit)) => {
let app = app.clone();
join_set.spawn(async move {
let _permit = permit;
serve_connection(TokioIo::new(stream), app, header_read_timeout).await;
});
}
None => shutting_down = true,
}
}
_ = shutdown_pin.as_mut() => {
shutting_down = true;
}
}
continue;
}
match join_set.join_next().await {
Some(_) => continue,
None => break,
}
}
}
#[cfg(feature = "tls")]
async fn serve_tls_with_shutdown<S, F>(
listener: TcpListener,
app: Arc<App<S>>,
acceptor: TlsAcceptor,
shutdown: F,
) where
S: Send + Sync + 'static,
F: Future<Output = ()> + Send + 'static,
{
let header_read_timeout = app.header_read_timeout;
let handshake_timeout = app.tls_handshake_timeout;
let semaphore = Arc::new(tokio::sync::Semaphore::new(app.max_connections));
let mut backoff = Backoff::new();
let mut join_set: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
let mut shutdown_pin = std::pin::pin!(shutdown);
let mut shutting_down = false;
loop {
if !shutting_down {
tokio::select! {
accepted = accept_and_permit(&listener, &mut backoff, &semaphore) => {
match accepted {
Some((stream, permit)) => {
let app = app.clone();
let acceptor = acceptor.clone();
join_set.spawn(async move {
let _permit = permit;
let tls_stream = match tokio::time::timeout(handshake_timeout, acceptor.accept(stream)).await {
Ok(Ok(s)) => s,
Ok(Err(_)) | Err(_) => return,
};
serve_connection(TokioIo::new(tls_stream), app, header_read_timeout).await;
});
}
None => shutting_down = true,
}
}
_ = shutdown_pin.as_mut() => {
shutting_down = true;
}
}
continue;
}
match join_set.join_next().await {
Some(_) => continue,
None => break,
}
}
}
#[must_use = "RouteBuilder does nothing until .seal() is called"]
pub struct RouteBuilder<S> {
state: Arc<S>,
router: Router<S>,
max_body_size: usize,
header_read_timeout: Duration,
max_connections: usize,
#[cfg(feature = "tls")]
tls_handshake_timeout: Duration,
error_handler: ErrorHandler,
cors_config: Option<CorsConfig>,
}
impl<S: Send + Sync + 'static> RouteBuilder<S> {
pub fn new(state: S) -> Self {
RouteBuilder {
state: Arc::new(state),
router: Router::new(),
max_body_size: DEFAULT_MAX_BODY_SIZE,
header_read_timeout: DEFAULT_HEADER_READ_TIMEOUT,
max_connections: DEFAULT_MAX_CONNECTIONS,
#[cfg(feature = "tls")]
tls_handshake_timeout: DEFAULT_TLS_HANDSHAKE_TIMEOUT,
error_handler: default_error_handler(),
cors_config: None,
}
}
pub fn with_max_body_size(mut self, max: usize) -> Self {
self.max_body_size = max;
self
}
pub fn with_header_read_timeout(mut self, d: Duration) -> Self {
self.header_read_timeout = d;
self
}
#[cfg(feature = "tls")]
pub fn with_tls_handshake_timeout(mut self, d: Duration) -> Self {
self.tls_handshake_timeout = d;
self
}
pub fn with_max_connections(mut self, max: usize) -> Self {
self.max_connections = max;
self
}
pub fn with_error_handler(
mut self,
f: impl Fn(StatusCode, &str) -> Response<ResponseBody> + Send + Sync + 'static,
) -> Self {
self.error_handler = Arc::new(f);
self
}
pub fn with_cors(mut self, config: CorsConfig) -> Self {
self.cors_config = Some(config);
self
}
pub fn get(mut self, path: &str, handler: Handler<S>) -> Self {
self.router.insert(Method::GET, path, handler);
self
}
pub fn post(mut self, path: &str, handler: Handler<S>) -> Self {
self.router.insert(Method::POST, path, handler);
self
}
pub fn put(mut self, path: &str, handler: Handler<S>) -> Self {
self.router.insert(Method::PUT, path, handler);
self
}
pub fn delete(mut self, path: &str, handler: Handler<S>) -> Self {
self.router.insert(Method::DELETE, path, handler);
self
}
pub fn seal(self) -> App<S> {
App {
state: self.state,
router: Arc::new(self.router),
max_body_size: self.max_body_size,
header_read_timeout: self.header_read_timeout,
max_connections: self.max_connections,
#[cfg(feature = "tls")]
tls_handshake_timeout: self.tls_handshake_timeout,
error_handler: self.error_handler,
cors_config: self.cors_config,
}
}
}
impl RouteBuilder<()> {
pub fn stateless() -> Self {
RouteBuilder::new(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use http_body_util::BodyExt;
#[test]
fn ephemeral_bind_addr_is_loopback_only() {
assert_eq!(
ephemeral_bind_addr().ip(),
std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)
);
}
#[test]
fn backoff_delays_double_up_to_a_cap() {
let mut backoff = Backoff::new();
assert_eq!(backoff.next_delay(), ACCEPT_BACKOFF_INITIAL);
assert_eq!(backoff.next_delay(), ACCEPT_BACKOFF_INITIAL * 2);
assert_eq!(backoff.next_delay(), ACCEPT_BACKOFF_INITIAL * 4);
let mut last = Duration::ZERO;
for _ in 0..20 {
last = backoff.next_delay();
}
assert_eq!(last, ACCEPT_BACKOFF_MAX);
}
#[test]
fn backoff_reset_returns_to_initial_delay() {
let mut backoff = Backoff::new();
backoff.next_delay();
backoff.next_delay();
backoff.reset();
assert_eq!(backoff.next_delay(), ACCEPT_BACKOFF_INITIAL);
}
struct FlakyListener {
inner: TcpListener,
remaining_failures: AtomicUsize,
attempts: Mutex<Vec<tokio::time::Instant>>,
}
impl TcpAccept for FlakyListener {
async fn accept(&self) -> std::io::Result<(TcpStream, SocketAddr)> {
self.attempts.lock().unwrap().push(tokio::time::Instant::now());
if self.remaining_failures.fetch_sub(1, Ordering::SeqCst) > 0 {
Err(std::io::Error::other("simulated accept error"))
} else {
TcpAccept::accept(&self.inner).await
}
}
}
#[tokio::test(start_paused = true)]
async fn accept_loop_backs_off_between_repeated_errors_instead_of_busy_spinning() {
let inner = TcpListener::bind(("127.0.0.1", 0)).await.unwrap();
let addr = inner.local_addr().unwrap();
let flaky = FlakyListener {
inner,
remaining_failures: AtomicUsize::new(5),
attempts: Mutex::new(Vec::new()),
};
tokio::spawn(async move {
let _ = TcpStream::connect(addr).await;
});
let mut backoff = Backoff::new();
accept_with_backoff(&flaky, &mut backoff).await;
let recorded = flaky.attempts.lock().unwrap();
assert_eq!(recorded.len(), 6, "5 failures then 1 success");
let expected_gaps = [
ACCEPT_BACKOFF_INITIAL,
ACCEPT_BACKOFF_INITIAL * 2,
ACCEPT_BACKOFF_INITIAL * 4,
ACCEPT_BACKOFF_INITIAL * 8,
ACCEPT_BACKOFF_INITIAL * 16,
];
for (i, expected) in expected_gaps.iter().enumerate() {
let gap = recorded[i + 1] - recorded[i];
assert_eq!(
gap, *expected,
"gap between attempt {i} and {} should reflect the backoff delay, not a busy spin",
i + 1
);
}
}
#[tokio::test]
async fn error_handler_sanitizes_5xx_in_response_body() {
take_error_log(); let resp = error_response(StatusCode::INTERNAL_SERVER_ERROR, "raw db connection string leaked");
let (parts, body) = resp.into_parts();
assert_eq!(parts.status, StatusCode::INTERNAL_SERVER_ERROR);
let collected = body.collect().await.unwrap();
let bytes = collected.to_bytes();
let json: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
let msg = json.get("message").and_then(|v| v.as_str()).unwrap();
assert_eq!(msg, "internal server error", "5xx message should be sanitized");
}
#[test]
fn error_handler_captures_5xx_message_in_log() {
take_error_log(); let sensitive_msg = "raw db connection string leaked";
error_response(StatusCode::INTERNAL_SERVER_ERROR, sensitive_msg);
let log = take_error_log();
assert_eq!(log.len(), 1);
assert_eq!(log[0].0, 500);
assert_eq!(log[0].1, sensitive_msg);
}
#[tokio::test]
async fn error_handler_passes_through_4xx_in_response_body() {
take_error_log(); let msg = "bad request";
let resp = error_response(StatusCode::BAD_REQUEST, msg);
let (parts, body) = resp.into_parts();
assert_eq!(parts.status, StatusCode::BAD_REQUEST);
let collected = body.collect().await.unwrap();
let bytes = collected.to_bytes();
let json: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
let response_msg = json.get("message").and_then(|v| v.as_str()).unwrap();
assert_eq!(response_msg, msg, "4xx message should pass through");
}
#[test]
fn error_handler_logs_4xx_messages() {
take_error_log(); let msg = "bad request";
error_response(StatusCode::BAD_REQUEST, msg);
let log = take_error_log();
assert_eq!(log.len(), 1);
assert_eq!(log[0].0, 400);
assert_eq!(log[0].1, msg);
}
}