Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,13 @@ http-body = "1"
pin-project-lite = "0.2"
# `serve_with`'s connections: HTTP/1.1 and HTTP/2 on one port, over TCP or
# TLS. Already in the tree through axum and tonic.
hyper-util = { version = "0.1", features = ["server-auto", "service", "tokio"] }
hyper-util = { version = "0.1", features = ["server-auto", "server-graceful", "service", "tokio"] }
# The executor HTTP/2 connections spawn their stream tasks on: a
# `TaskTracker` a shutdown waits on, and a `CancellationToken` that ends
# them (0.7.16: `run_until_cancelled` checks the token first). Both already
# in the tree through hyper-util and tonic.
hyper = "1"
tokio-util = { version = "0.7.16", features = ["rt"] }

# gRPC client (to upstream service). `tls-connect-info` gives a rustls server
# stream the connection record (client certificates) a tonic handler reads, so
Expand Down Expand Up @@ -180,6 +186,9 @@ tower = { version = "0.5", features = ["util", "limit"] }
http-body-util = "0.1"
# `hyper::upgrade::on`: a fallback that upgrades its connection (src/serve/tests.rs).
hyper = "1"
# A raw HTTP/2 client that sees the GOAWAY of a graceful shutdown
# (tests/shutdown.rs); already in the tree through hyper.
h2 = "0.4"
# Trailer frames for the hand-written upstream response in
# tests/upstream_controls.rs (tonic's server API cannot set success trailers).
http-body = "1"
Expand Down
67 changes: 62 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@ tonic services.
through, or translated to gRPC for an upstream that speaks only gRPC
- Built-in TLS and mTLS, and a cap on open connections; or your own TLS, with
the client's address and certificate reaching the upstream
- Graceful shutdown: calls and streams in flight finish within a bounded
drain, HTTP/2 clients get a GOAWAY ([Shutting down](#shutting-down))
- Guards you scope to the traffic they cover (transcoded calls, the proxy's
own endpoints, native gRPC, the fallback), rejecting in the protocol of the
request ([Guards and scopes](#guards-and-scopes)):
Expand Down Expand Up @@ -90,6 +92,9 @@ binary uses. Unset, `TOKIO_WORKER_THREADS` decides, else the number of CPUs
available to the process. The startup log (`RUST_LOG=info`) shows the count
and where it came from.

SIGTERM or Ctrl-C stops the binary gracefully: it stops accepting, lets the
calls in flight finish for up to `listen.drain_timeout_secs`, and exits 0.

## Configuration

```yaml
Expand All @@ -105,6 +110,10 @@ listen:
idle_timeout_secs: 60
# Seconds an HTTP/1.1 client has to send a request's headers.
header_read_timeout_secs: 30
# Seconds a shutdown (SIGTERM, Ctrl-C) waits for requests and streams in
# flight to finish; the connections still open after it are closed.
# 0 waits for all of them.
drain_timeout_secs: 25
# Optional: TLS on the listener (REST and gRPC share the port; ALPN offers
# h2 and http/1.1). With client_ca_file, client certificates are verified
# (mTLS) and reach an in-process tonic upstream as Request::peer_certs.
Expand Down Expand Up @@ -727,7 +736,7 @@ returns `application/jwk-set+json` (RFC 7517 §8.5).
command-line dependencies come with the `cli` feature. The library runs on
your tokio runtime and logs through `tracing` to the subscriber you set up.

```rust
```rust,no_run
use std::path::Path;
use structured_proxy::ProxyServer;

Expand All @@ -738,8 +747,13 @@ async fn main() -> anyhow::Result<()> {
// `ProxyConfig`.
let server = ProxyServer::from_file(Path::new("my-service.yaml"))?;

// Run the proxy on the configured listen address.
server.serve().await?;
// Run the proxy on the configured listen address until Ctrl-C, then
// drain (see "Shutting down").
server
.serve_with_shutdown(async {
tokio::signal::ctrl_c().await.ok();
})
.await?;
Ok(())
}
```
Expand Down Expand Up @@ -858,6 +872,48 @@ listener verified reaches a tonic handler in process as `Request::peer_certs`.
TLS needs a rustls crypto provider: the one a crypto backend feature brings,
or the one your process installed (see [TLS crypto](#tls-crypto)).

### Shutting down

`serve` and `serve_with` run until their future is dropped, which closes every
connection at once. For a graceful stop, hand `serve_with_shutdown` (or
`ProxyServer::serve_with_shutdown`) a future that completes when the process
should stop:

```rust
use std::time::Duration;
use structured_proxy::{ProxyServer, ServeOptions};

# async fn run(grpc: tonic::service::Routes) -> anyhow::Result<()> {
let proxy = ProxyServer::from_file(std::path::Path::new("my-service.yaml"))?.service(grpc)?;
let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await?;
let options = ServeOptions::new().drain_timeout(Some(Duration::from_secs(20)));
let stop = async {
tokio::signal::ctrl_c().await.ok();
};
structured_proxy::serve_with_shutdown(listener, proxy, options, stop).await?;
# Ok(())
# }
```

Once the future completes:

- the listening socket closes, so new connections are refused, and a
connection still in its TLS handshake or waiting for a `max_connections`
slot is dropped;
- HTTP/2 connections get a GOAWAY, so their clients open no new calls, and
HTTP/1.1 connections close after the response in progress;
- calls and streams in flight run to their end, and `serve_with_shutdown`
returns once every connection has closed;
- after `drain_timeout` (`listen.drain_timeout_secs`; `None` or 0 waits
without a bound) the connections still open are closed.

Past its grace period an orchestrator kills the process along with its calls,
so keep the drain below it. The default of 25 s fits the 30 s Kubernetes gives
a pod after SIGTERM; with a longer `terminationGracePeriodSeconds` the drain
can grow with it. A connection a fallback upgraded (a
WebSocket) belongs to the fallback's task and closes when that task lets it
go.

### Behind your own TLS

For a server of your own (another TLS stack, a Unix socket), run the service
Expand Down Expand Up @@ -941,7 +997,7 @@ a long server stream short.
`ProxyServer::router` returns the proxy's HTTP routes in front of the
configured upstream address, to serve or to merge into your own axum `Router`:

```rust
```rust,no_run
use std::path::Path;
use structured_proxy::{config::ProxyConfig, ProxyServer};

Expand Down Expand Up @@ -1033,6 +1089,7 @@ single-backend while `jsonwebtoken` sees two. Settle it once at the top of

```rust
# fn main() {
# #[cfg(feature = "builtin_jwt")]
structured_proxy::install_default_crypto_provider();
# }
```
Expand Down Expand Up @@ -1151,7 +1208,7 @@ verifies certificates with its own patched `rustls-webpki`.
At startup the proxy reads your proto descriptors and turns every
`google.api.http` rule into a REST route. Each request is then sorted once:

```
```text
REST, gRPC and gRPC-Web clients (HTTP/1.1, HTTP/2, optional TLS)
│
┌──────────────▼──────────────┐
Expand Down
4 changes: 4 additions & 0 deletions packaging/config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,10 @@
# HTTP listen address for the transcoded REST surface.
listen:
http: "0.0.0.0:8080"
# Seconds a stop (systemctl stop, SIGTERM) waits for calls in flight to
# finish before closing their connections. Keep it below TimeoutStopSec of
# the unit (30).
# drain_timeout_secs: 25

# Upstream gRPC service this proxy transcodes to. REQUIRED — replace with your
# service address.
Expand Down
3 changes: 3 additions & 0 deletions packaging/structured-proxy.service
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,9 @@ Environment=RUST_LOG=info
Restart=on-failure
RestartSec=5
LimitNOFILE=65536
# On SIGTERM the proxy drains for up to listen.drain_timeout_secs (25 by
# default), then exits; keep this above it, so systemd does not kill calls the
# drain would still finish.
TimeoutStopSec=30
StandardOutput=journal
StandardError=journal
Expand Down
5 changes: 4 additions & 1 deletion src/auth/jwks/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -207,7 +207,10 @@ async fn a_lookup_during_a_refresh_waits_for_its_keys() {
// refresh is in flight when a second lookup arrives: the second must not
// answer from the aged set in the meantime.
let (endpoint, uri) = endpoint().await;
let interval = Duration::from_millis(50);
// The second lookup checks the throttle only once the held refresh ends,
// against when that refresh started. The interval must outlast that on a
// loaded machine, or the second lookup may rightly refresh a third time.
let interval = Duration::from_millis(500);
let cache = Arc::new(
JwksCache::new(uri)
.unwrap()
Expand Down
10 changes: 10 additions & 0 deletions src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -479,6 +479,15 @@ pub struct ListenConfig {
/// 1. Default: 30.
#[serde(default = "default_header_read_timeout_secs")]
pub header_read_timeout_secs: u64,
/// Seconds a graceful shutdown waits for open connections to finish what
/// they serve; the ones still open after it are closed. 0 waits for all.
/// Default: 25, below the 30 s grace period Kubernetes gives by default.
#[serde(default = "default_drain_timeout_secs")]
pub drain_timeout_secs: u64,
}

fn default_drain_timeout_secs() -> u64 {
25
}

fn default_idle_timeout_secs() -> u64 {
Expand Down Expand Up @@ -547,6 +556,7 @@ impl Default for ListenConfig {
tls: None,
idle_timeout_secs: default_idle_timeout_secs(),
header_read_timeout_secs: default_header_read_timeout_secs(),
drain_timeout_secs: default_drain_timeout_secs(),
}
}
}
Expand Down
41 changes: 38 additions & 3 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,11 @@ compile_error!(
(or neither, and inject a verifier with `ProxyServer::with_token_verifier`)"
);

/// The README's Rust examples, compiled as doc tests.
#[cfg(doctest)]
#[doc = include_str!("../README.md")]
struct ReadmeDoctests;

pub mod auth;
pub mod config;
mod cors;
Expand All @@ -73,7 +78,7 @@ pub mod upstream;
/// [`install_default_crypto_provider`] for when a call is needed.
#[cfg(feature = "builtin_jwt")]
pub use auth::crypto::install_default_crypto_provider;
pub use serve::{serve, serve_with, ServeOptions};
pub use serve::{serve, serve_with, serve_with_shutdown, ServeOptions};
pub use service::{ConnectionInfo, ProxyService};

use axum::extract::State;
Expand Down Expand Up @@ -1152,7 +1157,11 @@ impl ProxyServer {
(listen.idle_timeout_secs > 0)
.then(|| Duration::from_secs(listen.idle_timeout_secs)),
)
.header_read_timeout(Duration::from_secs(listen.header_read_timeout_secs));
.header_read_timeout(Duration::from_secs(listen.header_read_timeout_secs))
.drain_timeout(
(listen.drain_timeout_secs > 0)
.then(|| Duration::from_secs(listen.drain_timeout_secs)),
);
if let Some(max) = listen.max_connections {
anyhow::ensure!(max > 0, "listen.max_connections must be at least 1");
options = options.max_connections(max);
Expand Down Expand Up @@ -1183,6 +1192,31 @@ impl ProxyServer {
/// [`serve_options`](Self::serve_options) reject, an invalid listen
/// address, or a listener that fails.
pub async fn serve(&self) -> anyhow::Result<()> {
self.serve_with_shutdown(std::future::pending()).await
}

/// [`serve`](Self::serve) until `signal` completes, then shut down with
/// the drain of `listen.drain_timeout_secs` (see [`serve_with_shutdown`]).
///
/// # Errors
///
/// What [`serve`](Self::serve) returns.
///
/// # Examples
///
/// ```no_run
/// # async fn run(server: structured_proxy::ProxyServer) -> anyhow::Result<()> {
/// server
/// .serve_with_shutdown(async {
/// tokio::signal::ctrl_c().await.ok();
/// })
/// .await
/// # }
/// ```
pub async fn serve_with_shutdown(
&self,
signal: impl std::future::Future<Output = ()>,
) -> anyhow::Result<()> {
let service = self.service(self.upstream()?)?;
let options = self.serve_options()?;
let addr: SocketAddr = self.config.listen.http.parse()?;
Expand All @@ -1194,7 +1228,8 @@ impl ProxyServer {
self.config.service.name,
addr
);
serve_with(listener, service, options).await?;
serve_with_shutdown(listener, service, options, signal).await?;
tracing::info!("{} stopped", self.config.service.name);
Ok(())
}
}
Expand Down
35 changes: 34 additions & 1 deletion src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,5 +52,38 @@ fn main() -> anyhow::Result<()> {
"Starting structured-proxy"
);

rt.block_on(server.serve())
rt.block_on(server.serve_with_shutdown(stop_requested()))
}

/// Completes on Ctrl-C, and on SIGTERM on Unix (what container runtimes and
/// service managers send to stop a process).
async fn stop_requested() {
#[cfg(unix)]
{
use tokio::signal::unix::{signal, SignalKind};
match signal(SignalKind::terminate()) {
Ok(mut term) => {
tokio::select! {
() = ctrl_c() => {}
_ = term.recv() => {}
}
}
Err(e) => {
tracing::warn!(error = %e, "cannot listen for SIGTERM; stopping on Ctrl-C only");
ctrl_c().await;
}
}
}
#[cfg(not(unix))]
ctrl_c().await;
tracing::info!("stop requested; draining connections");
}

/// Completes on Ctrl-C. A handler that cannot be installed never completes:
/// failing to listen for the signal is no reason to stop serving.
async fn ctrl_c() {
if let Err(e) = tokio::signal::ctrl_c().await {
tracing::warn!(error = %e, "cannot listen for Ctrl-C");
std::future::pending::<()>().await;
}
}
Loading
Loading