karawaci.kode

← Semua snippet

Rust Menengah Utility

Tokio graceful shutdown dengan signal handling

Graceful shutdown Rust Tokio: tangkap SIGTERM, broadcast cancel ke semua task, tunggu drain. Pattern wajib untuk production service.

Dipublikasikan 10 Juli 2026

Rust production service tanpa graceful shutdown bikin in-flight request kepotong dan transaksi inconsistent. Tokio kasih primitives bagus tapi pattern-nya harus disusun sendiri. Snippet ini lengkap: signal handler + cancellation token + worker drain dengan timeout.

Kode

// Cargo.toml
// [dependencies]
// tokio = { version = "1.40", features = ["full"] }
// tokio-util = { version = "0.7", features = ["rt"] }
// axum = "0.7"
// tracing = "0.1"
// tracing-subscriber = "0.3"
// futures-util = "0.3"

// main.rs
use std::sync::Arc;
use std::time::Duration;

use axum::{routing::get, Router};
use tokio::net::TcpListener;
use tokio::signal;
use tokio::sync::broadcast;
use tokio::task::JoinSet;
use tokio_util::sync::CancellationToken;
use tracing::{error, info, warn};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    tracing_subscriber::fmt::init();

    // Token utama — propagate ke semua subsystem
    let shutdown_token = CancellationToken::new();

    // Watcher untuk OS signal
    let signal_token = shutdown_token.clone();
    tokio::spawn(async move {
        wait_for_shutdown_signal().await;
        info!("Sinyal shutdown diterima, broadcast cancel");
        signal_token.cancel();
    });

    // Group untuk track semua long-running task
    let mut tasks = JoinSet::new();

    // 1. HTTP server (axum)
    let http_token = shutdown_token.clone();
    tasks.spawn(async move {
        if let Err(e) = run_http_server(http_token).await {
            error!("HTTP server error: {e}");
        }
        info!("HTTP server task selesai");
    });

    // 2. Background worker (consume queue, process job)
    let worker_token = shutdown_token.clone();
    tasks.spawn(async move {
        run_worker(worker_token).await;
        info!("Worker task selesai");
    });

    // 3. Scheduled task (heartbeat)
    let heartbeat_token = shutdown_token.clone();
    tasks.spawn(async move {
        run_heartbeat(heartbeat_token).await;
        info!("Heartbeat task selesai");
    });

    // Tunggu semua task drain — dengan timeout safety net
    let drain_result = tokio::time::timeout(
        Duration::from_secs(30),
        async {
            while let Some(res) = tasks.join_next().await {
                if let Err(e) = res {
                    error!("Task panic: {e}");
                }
            }
        },
    ).await;

    match drain_result {
        Ok(_) => info!("Semua task selesai dengan bersih"),
        Err(_) => {
            warn!("Timeout drain 30s — force abort sisa task");
            tasks.abort_all();
        }
    }

    info!("Shutdown selesai");
    Ok(())
}

/// Wait untuk SIGTERM atau SIGINT (Ctrl+C)
async fn wait_for_shutdown_signal() {
    #[cfg(unix)]
    {
        use signal::unix::{signal, SignalKind};
        let mut sigterm = signal(SignalKind::terminate()).expect("install SIGTERM handler");
        let mut sigint = signal(SignalKind::interrupt()).expect("install SIGINT handler");

        tokio::select! {
            _ = sigterm.recv() => info!("SIGTERM diterima"),
            _ = sigint.recv() => info!("SIGINT diterima"),
        }
    }

    #[cfg(not(unix))]
    {
        signal::ctrl_c().await.expect("install Ctrl+C handler");
        info!("Ctrl+C diterima");
    }
}

/// HTTP server dengan graceful shutdown via axum::serve
async fn run_http_server(token: CancellationToken) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
    let app = Router::new()
        .route("/", get(|| async { "OK" }))
        .route("/healthz", get(|| async { "OK" }));

    let listener = TcpListener::bind("0.0.0.0:8080").await?;
    info!("HTTP server listening :8080");

    axum::serve(listener, app)
        .with_graceful_shutdown(async move {
            token.cancelled().await;
            info!("HTTP: stop terima koneksi baru, tunggu inflight selesai");
        })
        .await?;

    Ok(())
}

/// Background worker — process job loop
async fn run_worker(token: CancellationToken) {
    info!("Worker started");

    loop {
        tokio::select! {
            // Polling job dengan timeout supaya cek cancel berkala
            _ = token.cancelled() => {
                info!("Worker dapat cancel signal, drain in-flight");
                drain_inflight_jobs().await;
                return;
            }
            _ = tokio::time::sleep(Duration::from_secs(1)) => {
                // Process satu batch job
                if let Err(e) = process_next_batch().await {
                    error!("Worker batch error: {e}");
                }
            }
        }
    }
}

/// Heartbeat — kirim metric atau check-in
async fn run_heartbeat(token: CancellationToken) {
    let mut interval = tokio::time::interval(Duration::from_secs(10));

    loop {
        tokio::select! {
            _ = token.cancelled() => {
                info!("Heartbeat stopped");
                return;
            }
            _ = interval.tick() => {
                info!("Heartbeat tick");
            }
        }
    }
}

async fn process_next_batch() -> Result<(), Box<dyn std::error::Error>> {
    // Dummy work
    tokio::time::sleep(Duration::from_millis(100)).await;
    Ok(())
}

async fn drain_inflight_jobs() {
    // Cleanup in-flight: ack pending message, commit DB transaction
    tokio::time::sleep(Duration::from_secs(2)).await;
    info!("In-flight job drain complete");
}

Pemakaian

# Build & run
cargo run --release

# Trigger graceful shutdown
# Terminal 1: cargo run
# Terminal 2: kill -SIGTERM $(pgrep app)
# Atau Ctrl+C di terminal 1

# Log output:
# [INFO] HTTP server listening :8080
# [INFO] Worker started
# [INFO] Heartbeat tick
# ^C
# [INFO] SIGINT diterima
# [INFO] Sinyal shutdown diterima, broadcast cancel
# [INFO] HTTP: stop terima koneksi baru, tunggu inflight selesai
# [INFO] Worker dapat cancel signal, drain in-flight
# [INFO] Heartbeat stopped
# [INFO] In-flight job drain complete
# [INFO] HTTP server task selesai
# [INFO] Worker task selesai
# [INFO] Semua task selesai dengan bersih
# [INFO] Shutdown selesai
// Pattern: child token untuk subset cancellation
// Misal mau cancel hanya worker (bukan HTTP)
let worker_only_token = shutdown_token.child_token();

// Pakai worker_only_token di worker
// Parent cancel akan propagate, tapi worker bisa di-cancel sendiri
worker_only_token.cancel();  // Hanya worker yang cancel
# k8s — penting set terminationGracePeriodSeconds
apiVersion: apps/v1
kind: Deployment
spec:
  template:
    spec:
      terminationGracePeriodSeconds: 35   # > drain timeout di code (30s)
      containers:
        - name: app
          lifecycle:
            preStop:
              exec:
                # Wait Service endpoint propagate "not ready"
                command: ["sleep", "5"]

Kapan dipakai

  • Production HTTP / gRPC server (axum, tonic).
  • Background worker yang process message queue.
  • Long-running ETL pipeline.
  • WebSocket server dengan banyak connection terbuka.

Catatan

  • CancellationToken vs broadcast — Token lebih ergonomis untuk pure signal. Broadcast cocok kalau perlu kirim data dengan cancel (rare).
  • JoinSet untuk track task — better dari Vec karena bisa join_next iteratively dan abort_all sekaligus.
  • timeout pada drain — safety net. Kalau ada task yang infinite loop, abort paksa setelah 30s supaya pod tidak stuck di k8s.
  • axum graceful shutdown built-in via with_graceful_shutdown. Tunggu future complete sebelum stop accept connection.
  • Signal Unix only — di Windows pakai ctrl_c. Snippet ini handle keduanya via cfg.
  • child_token untuk scoped cancellation — subset task bisa di-cancel sendiri tanpa cancel parent.

tokio::spawn task yang gak punya cancel path bakal jalan terus walaupun shutdown. Setiap loop background HARUS punya cek cancel via select! atau token.is_cancelled().

# tags

rusttokiograceful-shutdownsignalasync

Ditulis oleh Asti Larasati · 10 Juli 2026