Skip to content

Commit a113543

Browse files
committed
feat: implemented shutdown
1 parent 24fc621 commit a113543

1 file changed

Lines changed: 54 additions & 23 deletions

File tree

adapter/rest/src/main.rs

Lines changed: 54 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,10 @@ mod route;
2929

3030
#[tokio::main]
3131
async fn main() {
32-
let server = HttpServer { addr: None };
32+
let server = HttpServer {
33+
shutdown_tx: None,
34+
addr: None,
35+
};
3336
let runner = match ServerRunner::new(server).await {
3437
Ok(runner) => runner,
3538
Err(err) => panic!("Failed to create server runner: {:?}", err),
@@ -42,6 +45,7 @@ async fn main() {
4245
}
4346

4447
struct HttpServer {
48+
shutdown_tx: Option<tokio::sync::broadcast::Sender<()>>,
4549
addr: Option<SocketAddr>,
4650
}
4751

@@ -265,6 +269,10 @@ pub async fn handle_request(
265269
impl ServerTrait<config::HttpServerConfig> for HttpServer {
266270
async fn init(&mut self, ctx: &ServerContext<config::HttpServerConfig>) -> anyhow::Result<()> {
267271
log::info!("Initializing http server");
272+
273+
let (shutdown_tx, _) = tokio::sync::broadcast::channel(1);
274+
self.shutdown_tx = Some(shutdown_tx);
275+
268276
let bind = format!("{}:{}", ctx.server_config.host, ctx.server_config.port);
269277

270278
self.addr = Some(
@@ -277,49 +285,72 @@ impl ServerTrait<config::HttpServerConfig> for HttpServer {
277285
}
278286

279287
async fn run(&mut self, ctx: &ServerContext<config::HttpServerConfig>) -> anyhow::Result<()> {
280-
let addr = match self.addr {
281-
Some(addr) => addr,
282-
None => panic!("cannot start tcp listener with empty address"),
283-
};
284-
285-
let listener = match TcpListener::bind(addr).await {
286-
Ok(listener) => listener,
287-
Err(err) => {
288-
panic!("failed to register tcp listener on address: {:?}", err);
289-
}
290-
};
288+
let addr = self
289+
.addr
290+
.expect("cannot start tcp listener with empty address");
291+
292+
let listener = TcpListener::bind(addr)
293+
.await
294+
.map_err(|e| anyhow::anyhow!("failed to bind {addr}: {e}"))?;
295+
296+
// Create a receiver for this run loop
297+
let shutdown_tx = self
298+
.shutdown_tx
299+
.as_ref()
300+
.expect("shutdown_tx not initialized; init() must run first")
301+
.clone();
302+
let mut shutdown_rx = shutdown_tx.subscribe();
291303

292304
loop {
293-
let (stream, _) = match listener.accept().await {
294-
Ok(res) => res,
295-
Err(e) => {
296-
panic!("listener failed to accept requests: {:?}", e);
305+
let (stream, _) = tokio::select! {
306+
_ = shutdown_rx.recv() => {
307+
log::info!("HTTP server: shutdown received, stopping accept loop");
308+
break;
309+
}
310+
res = listener.accept() => {
311+
res.map_err(|e| anyhow::anyhow!("accept failed: {e}"))?
297312
}
298313
};
299314

300315
let io = TokioIo::new(stream);
301-
302316
let store = Arc::clone(&ctx.adapter_store);
303317

304-
tokio::task::spawn(async move {
305-
let store = Arc::clone(&store);
318+
let mut conn_shutdown_rx = shutdown_tx.subscribe();
319+
320+
tokio::spawn(async move {
306321
let svc = hyper::service::service_fn(move |req| {
307322
let store = Arc::clone(&store);
308323
async move { handle_request(req, store).await }
309324
});
310325

311-
if let Err(err) = http1::Builder::new().serve_connection(io, svc).await {
312-
log::error!("Error serving connection: {:?}", err);
326+
let conn = http1::Builder::new().serve_connection(io, svc);
327+
328+
tokio::pin!(conn);
329+
330+
tokio::select! {
331+
res = conn.as_mut() => {
332+
if let Err(err) = res {
333+
log::error!("Error serving connection: {:?}", err);
334+
}
335+
}
336+
_ = conn_shutdown_rx.recv() => {
337+
conn.as_mut().graceful_shutdown();
338+
}
313339
}
314340
});
315341
}
316-
}
317342

343+
Ok(())
344+
}
318345
async fn shutdown(
319346
&mut self,
320347
_ctx: &ServerContext<config::HttpServerConfig>,
321348
) -> anyhow::Result<()> {
322-
todo!("Implement shutdown!");
349+
if let Some(ref tx) = self.shutdown_tx {
350+
log::info!("Recieved a shutdown signal for Adapter Server");
351+
let _ = tx.send(());
352+
}
353+
323354
Ok(())
324355
}
325356
}

0 commit comments

Comments
 (0)