cloudflare / cloudflare/pingora
Allow mutable reference gathering for custom ShutdownSignalWatch implementation
- Dominant language
- Rust
- Stars
- 27.4k
- Forks
- 1.7k
- Avg merge
- 6h 22m
- Merged PRs (30d)
- 3
Description
## What is the problem your feature solves, or the need it fulfills?
These changes allow you to create a custom implementation of 'ShutdownSignalWatch' to control the server lifecycle with more than just unix signals (e.g. tokio mspc).
## Describe the solution you'd like
My PR: https://github.com/cloudflare/pingora/pull/562
## Describe alternatives you've considered
NaN
## Additional context
Sample example for understand (unix signals omitted for shortness):
```rust
use async_trait::async_trait;
use log::info;
use pingora::{
Result,
http::RequestHeader,
lb::{LoadBalancer, health_check, selection::RoundRobin},
prelude::{HttpPeer, Opt, background_service},
proxy::{ProxyHttp, Session, http_proxy_service},
server::{RunArgs, Server, ShutdownSignal, ShutdownSignalWatch},
};
use std::{sync::Arc, time::Duration};
use tokio::sync::mpsc::Receiver;
pub struct LB(Arc>);
#[async_trait]
impl ProxyHttp for LB {
type CTX = ();
fn new_ctx(&self) -> Self::CTX {}
async fn upstream_peer(&self, _session: &mut Session, _ctx: &mut ()) -> Result> {
let upstream = self
.0
.select(b"", 256) // hash doesn't matter
.unwrap();
info!("upstream peer is: {:?}", upstream);
let peer = Box::new(HttpPeer::new(upstream, true, "one.one.one.one".to_string()));
Ok(peer)
}
async fn upstream_request_filter(
&self,
_session: &mut Session,
upstream_request: &mut RequestHeader,
_ctx: &mut Self::CTX,
) -> Result<()> {
upstream_request
.insert_header("Host", "one.one.one.one")
.unwrap();
Ok(())
}
}
fn construct() -> Server {
let opt = Opt::parse_args();
let mut my_server = Server::new(Some(opt)).unwrap();
my_server.bootstrap();
let mut upstreams =
LoadBalancer::try_from_iter(["1.1.1.1:443", "1.0.0.1:443", "127.0.0.1:343"]).unwrap();
let hc = health_check::TcpHealthCheck::new();
upstreams.set_health_check(hc);
upstreams.health_check_frequency = Some(Duration::from_secs(1));
let background = background_service("health check", upstreams);
let upstreams = background.task();
let mut lb = http_proxy_service(&my_server.configuration, LB(upstreams));
lb.add_tcp("0.0.0.0:6188");
my_server.add_service(lb);
my_server.add_service(background);
return my_server;
}
fn run(srv: Server, args: RunArgs) {
srv.run(args);
}
pub struct CustomShutdown {
rx: Receiver,
}
impl CustomShutdown {
pub fn new(rx: Receiver) -> Self {
return CustomShutdown { rx: rx };
}
}
#[async_trait]
impl ShutdownSignalWatch for CustomShutdown {
async fn recv(&mut self) -> ShutdownSignal {
while let Some(x) = self.rx.recv().await {
return x;
}
return ShutdownSignal::FastShutdown;
}
}
fn main() {
env_logger::init();
let (tx, rx) = tokio::sync::mpsc::channel::(1);
let run_args = RunArgs {
shutdown_signal: Box::new(CustomShutdown::new(rx)),
};
let custom_runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.unwrap();
custom_runtime.spawn(async move {
tokio::time::sleep(Duration::from_millis(1000)).await;
info!("send shutdown signal");
tx.send(ShutdownSignal::GracefulTerminate).await.unwrap();
});
let my_server = construct();
run(my_server, run_args);
}
```
Contributor guide
Assessment
This issue has not been assessed yet.