Urgensi Idempotensi pada Mutasi Data

Kegagalan jaringan pada komunikasi HTTP bersifat asimetris. Klien dapat mengirim mutasi (POST/PATCH), server berhasil memproses transaksi ke database, namun response koneksi terputus (TCP reset atau timeout) sebelum tiba di klien. Ketika klien mengeksekusi mekanisme exponential backoff retry, server berisiko mengeksekusi operasi finansial atau mutasi resource berulang kali.

Metode HTTP POST tidak bersifat idempoten menurut spesifikasi RFC 9110. Untuk menjamin keselamatan eksekusi ulang, integrasi API modern mengadopsi draft spesifikasi IETF Idempotency-Key. Klien melampirkan token unik pada header HTTP. Server memanfaatkan token ini untuk memastikan logika bisnis hanya dijalankan tepat satu kali.

Arsitektur State Machine dan Verifikasi Payload

Middleware idempotensi mengelola siklus hidup permintaan menggunakan state machine berbasis kunci:

  • Header Missing: Jika request tidak memuat Idempotency-Key, request langsung dialirkan ke inner service tanpa inspeksi idempotensi.
  • State IN_PROGRESS: Request dengan kunci yang sedang diproses di-lock. Jika ada request konkuren dengan kunci sama datang bersamaan, server merespons 409 Conflict untuk memutus race condition.
  • Payload Verification (SHA-256): Hash dari body request disimpan. Jika request berikutnya datang dengan kunci identik namun hash payload berbeda, server menggagalkan eksekusi dengan status 422 Unprocessable Entity.
  • State COMPLETED: Setelah handler controller selesai, status, header, dan body response disimpan di cache storage dengan TTL (Time-To-Live). Request replay dengan kunci dan hash identik langsung disajikan dari cache (HTTP Short-circuit).

Setup Dependensi Cargo.toml

Tambahkan dependensi berikut pada berkas Cargo.toml:

[dependencies]
actix-web = "4.4"
actix-http = "4.4"
futures-util = "0.3"
sha2 = "0.10"
tokio = { version = "1", features = ["sync", "time"] }
serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"

Implementasi State Store dan Trait Storage

Buat struktur data untuk melacak status transaksi serta membedakan fase eksekusi request.

use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};

#[derive(Clone, Debug)]
pub struct CachedResponse {
    pub status: u16,
    pub headers: Vec<(String, String)>,
    pub body: Vec<u8>,
}

#[derive(Clone, Debug)]
pub enum IdempotencyRecord {
    InProgress { payload_hash: String, created_at: Instant },
    Completed { payload_hash: String, response: CachedResponse, created_at: Instant },
}

pub trait IdempotencyStore: Send + Sync {
    fn get(&self, key: &str) -> Option<IdempotencyRecord>;
    fn lock_key(&self, key: &str, payload_hash: String) -> Result<(), ()>;
    fn complete(&self, key: &str, payload_hash: String, response: CachedResponse);
    fn remove(&self, key: &str);
}

pub struct InMemoryStore {
    ttl: Duration,
    entries: Arc<Mutex<HashMap<String, IdempotencyRecord>>>,
}

impl InMemoryStore {
    pub fn new(ttl: Duration) -> Self {
        Self {
            ttl,
            entries: Arc::new(Mutex::new(HashMap::new())),
        }
    }
}

impl IdempotencyStore for InMemoryStore {
    fn get(&self, key: &str) -> Option<IdempotencyRecord> {
        let mut map = self.entries.lock().unwrap();
        if let Some(record) = map.get(key) {
            let created_at = match record {
                IdempotencyRecord::InProgress { created_at, .. } => *created_at,
                IdempotencyRecord::Completed { created_at, .. } => *created_at,
            };
            if created_at.elapsed() > self.ttl {
                map.remove(key);
                return None;
            }
            return Some(record.clone());
        }
        None
    }

    fn lock_key(&self, key: &str, payload_hash: String) -> Result<(), ()> {
        let mut map = self.entries.lock().unwrap();
        if map.contains_key(key) {
            return Err(());
        }
        map.insert(key.to_string(), IdempotencyRecord::InProgress {
            payload_hash,
            created_at: Instant::now(),
        });
        Ok(())
    }

    fn complete(&self, key: &str, payload_hash: String, response: CachedResponse) {
        let mut map = self.entries.lock().unwrap();
        map.insert(key.to_string(), IdempotencyRecord::Completed {
            payload_hash,
            response,
            created_at: Instant::now(),
        });
    }

    fn remove(&self, key: &str) {
        let mut map = self.entries.lock().unwrap();
        map.remove(key);
    }
}

Membangun Middleware Actix Web Idiomatik

Middleware Actix Web membutuhkan implementasi trait Transform dan Service. Tantangan teknis utama: membaca payload stream untuk hashing tanpa menghabiskannya dari pipeline controller. Solusinya adalah menyalin bytes stream, menghitung hash, lalu merekonstruksi payload ke dalam ServiceRequest via actix_http::h1::Payload.

use actix_web::dev::{forward_ready, Service, ServiceRequest, ServiceResponse, Transform};
use actix_web::{body::to_bytes, http::header::HeaderName, http::StatusCode, Error, HttpResponse};
use futures_util::future::{ready, LocalBoxFuture, Ready};
use futures_util::StreamExt;
use sha2::{Digest, Sha256};
use std::rc::Rc;
use std::sync::Arc;

pub struct IdempotencyMiddleware {
    pub store: Arc<dyn IdempotencyStore>,
}

impl<S, B> Transform<S, ServiceRequest> for IdempotencyMiddleware
where
    S: Service<ServiceRequest, Response = ServiceResponse<B>, Error = Error> + 'static,
    B: actix_web::body::MessageBody + 'static,
{
    type Response = ServiceResponse<actix_web::body::BoxBody>;
    type Error = Error;
    type InitError = ();
    type Transform = IdempotencyService<S>;
    type Future = Ready<Result<Self::Transform, Self::InitError>>;

    fn new_transform(&self, service: S) -> Self::Future {
        ready(Ok(IdempotencyService {
            service: Rc::new(service),
            store: self.store.clone(),
        }))
    }
}

pub struct IdempotencyService<S> {
    service: Rc<S>,
    store: Arc<dyn IdempotencyStore>,
}

impl<S, B> Service<ServiceRequest> for IdempotencyService<S>
where
    S: Service<ServiceRequest, Response = ServiceResponse<B>, Error = Error> + 'static,
    B: actix_web::body::MessageBody + 'static,
{
    type Response = ServiceResponse<actix_web::body::BoxBody>;
    type Error = Error;
    type Future = LocalBoxFuture<'static, Result<Self::Response, Self::Error>>;

    forward_ready!(service);

    fn call(&self, mut req: ServiceRequest) -> Self::Future {
        let svc = self.service.clone();
        let store = self.store.clone();

        Box::pin(async move {
            let header_key = HeaderName::from_static("idempotency-key");
            let key = match req.headers().get(&header_key) {
                Some(val) => match val.to_str() {
                    Ok(v) => v.to_string(),
                    Err(_) => {
                        let res = HttpResponse::BadRequest().body("Invalid Idempotency-Key encoding");
                        return Ok(req.into_response(res.map_into_boxed_body()));
                    }
                },
                None => {
                    return svc.call(req).await.map(|res| res.map_into_boxed_body());
                }
            };

            // Baca body stream dan hitung SHA-256
            let mut body_bytes = actix_web::web::BytesMut::new();
            let mut stream = req.take_payload();
            while let Some(chunk) = stream.next().await {
                let chunk = chunk.map_err(|e| actix_web::error::ErrorBadRequest(e.to_string()))?;
                body_bytes.extend_from_slice(&chunk);
            }
            let payload = body_bytes.freeze();

            let mut hasher = Sha256::new();
            hasher.update(&payload);
            let payload_hash = format!("{:x}", hasher.finalize());

            // Kembalikan payload stream ke request agar controller dapat mengekstrak data
            let (_, mut new_payload) = actix_http::h1::Payload::create(true);
            new_payload.unread_data(payload.clone());
            req.set_payload(new_payload.into());

            // Evaluasi State Cache
            if let Some(record) = store.get(&key) {
                match record {
                    IdempotencyRecord::InProgress { payload_hash: existing_hash, .. } => {
                        if existing_hash != payload_hash {
                            let res = HttpResponse::UnprocessableEntity()
                                .body("Payload hash mismatch on active key");
                            return Ok(req.into_response(res.map_into_boxed_body()));
                        }
                        let res = HttpResponse::Conflict()
                            .body("Request with this key is already in progress");
                        return Ok(req.into_response(res.map_into_boxed_body()));
                    }
                    IdempotencyRecord::Completed { payload_hash: existing_hash, response, .. } => {
                        if existing_hash != payload_hash {
                            let res = HttpResponse::UnprocessableEntity()
                                .body("Payload hash mismatch with completed key");
                            return Ok(req.into_response(res.map_into_boxed_body()));
                        }

                        // Replay cached response tanpa memanggil downstream service
                        let mut resp_builder = HttpResponse::build(
                            StatusCode::from_u16(response.status).unwrap_or(StatusCode::OK)
                        );
                        for (hk, hv) in response.headers {
                            resp_builder.append_header((hk, hv));
                        }
                        let final_res = resp_builder.body(response.body);
                        return Ok(req.into_response(final_res.map_into_boxed_body()));
                    }
                }
            }

            // Akuisisi lock state
            if store.lock_key(&key, payload_hash.clone()).is_err() {
                let res = HttpResponse::Conflict().body("Concurrent request detected");
                return Ok(req.into_response(res.map_into_boxed_body()));
            }

            // Eksekusi downstream service
            let res = match svc.call(req).await {
                Ok(resp) => resp,
                Err(err) => {
                    store.remove(&key);
                    return Err(err);
                }
            };

            let (http_req, http_resp) = res.into_parts();
            let status = http_resp.status().as_u16();

            // Tangani kegagalan server: hapus key jika status 5xx agar klien dapat mencoba lagi
            if http_resp.status().is_server_error() {
                store.remove(&key);
                return Ok(ServiceResponse::new(http_req, http_resp.map_into_boxed_body()));
            }

            let mut headers_vec = Vec::new();
            for (name, value) in http_resp.headers().iter() {
                if let Ok(v) = value.to_str() {
                    headers_vec.push((name.as_str().to_string(), v.to_string()));
                }
            }

            let body_bytes = match to_bytes(http_resp.into_body()).await {
                Ok(bytes) => bytes.to_vec(),
                Err(e) => {
                    store.remove(&key);
                    return Err(actix_web::error::ErrorInternalServerError(e.to_string()));
                }
            };

            store.complete(&key, payload_hash, CachedResponse {
                status,
                headers: headers_vec.clone(),
                body: body_bytes.clone(),
            });

            let mut final_builder = HttpResponse::build(StatusCode::from_u16(status).unwrap());
            for (hk, hv) in headers_vec {
                final_builder.append_header((hk, hv));
            }
            let reconstructed = final_builder.body(body_bytes);
            Ok(ServiceResponse::new(http_req, reconstructed.map_into_boxed_body()))
        })
    }
}

Pengujian Otomatis Menggunakan actix_web::test

Gunakan suite tes berikut untuk memvalidasi keempat skenario kritis: mutasi normal, response replay, payload mismatch, dan penanganan konkurensi.

#[cfg(test)]
mod tests {
    use super::*;
    use actix_web::{test, web, App, HttpResponse, Responder};
    use serde::{Deserialize, Serialize};
    use std::sync::atomic::{AtomicUsize, Ordering};

    #[derive(Serialize, Deserialize, Clone)]
    struct OrderRequest {
        item_id: u32,
        quantity: u32,
    }

    static EXECUTION_COUNT: AtomicUsize = AtomicUsize::new(0);

    async fn handle_checkout(data: web::Json<OrderRequest>) -> impl Responder {
        EXECUTION_COUNT.fetch_add(1, Ordering::SeqCst);
        HttpResponse::Created().json(serde_json::json!({
            "order_id": 101,
            "item": data.item_id,
            "status": "CONFIRMED"
        }))
    }

    #[actix_web::test]
    async fn test_idempotency_lifecycle() {
        let store = Arc::new(InMemoryStore::new(Duration::from_secs(60)));
        let app = test::init_service(
            App::new()
                .wrap(IdempotencyMiddleware { store: store.clone() })
                .route("/checkout", web::post().to(handle_checkout)),
        )
        .await;

        let payload_a = serde_json::to_string(&OrderRequest { item_id: 1, quantity: 2 }).unwrap();
        let payload_b = serde_json::to_string(&OrderRequest { item_id: 2, quantity: 5 }).unwrap();

        // 1. Eksekusi pertama: Sukses 201 Created
        let req1 = test::TestRequest::post()
            .uri("/checkout")
            .insert_header(("Idempotency-Key", "key-xyz-123"))
            .insert_header(("Content-Type", "application/json"))
            .set_payload(payload_a.clone())
            .to_request();

        let resp1 = test::call_service(&app, req1).await;
        assert_eq!(resp1.status(), StatusCode::CREATED);
        assert_eq!(EXECUTION_COUNT.load(Ordering::SeqCst), 1);

        // 2. Replay request dengan key & payload identik: Return 201 via cache, handler tidak dieksekusi
        let req2 = test::TestRequest::post()
            .uri("/checkout")
            .insert_header(("Idempotency-Key", "key-xyz-123"))
            .insert_header(("Content-Type", "application/json"))
            .set_payload(payload_a)
            .to_request();

        let resp2 = test::call_service(&app, req2).await;
        assert_eq!(resp2.status(), StatusCode::CREATED);
        assert_eq!(EXECUTION_COUNT.load(Ordering::SeqCst), 1); // Counter tetap 1

        // 3. Request dengan key sama tetapi payload beda: Return 422 Unprocessable Entity
        let req3 = test::TestRequest::post()
            .uri("/checkout")
            .insert_header(("Idempotency-Key", "key-xyz-123"))
            .insert_header(("Content-Type", "application/json"))
            .set_payload(payload_b)
            .to_request();

        let resp3 = test::call_service(&app, req3).await;
        assert_eq!(resp3.status(), StatusCode::UNPROCESSABLE_ENTITY);

        // 4. Request paralel saat status IN_PROGRESS: Return 409 Conflict
        store.lock_key("lock-concurrent", "dummy-hash".to_string()).unwrap();

        let req4 = test::TestRequest::post()
            .uri("/checkout")
            .insert_header(("Idempotency-Key", "lock-concurrent"))
            .insert_header(("Content-Type", "application/json"))
            .set_payload("dummy-hash")
            .to_request();

        let resp4 = test::call_service(&app, req4).await;
        assert_eq!(resp4.status(), StatusCode::CONFLICT);
    }
}

Pertimbangan Produksi

Peringatan Lingkungan Multi-Instance: Implementasi InMemoryStore di atas menggunakan mutex thread lokal dan hanya cocok untuk single process binary. Pada deployment multi-pod (seperti Kubernetes), ganti storage layer menggunakan Redis.

Untuk migrasi ke Redis Cluster, gunakan atomic command:

  • Locking: Jalankan SET idempotency:{key} {hash}:IN_PROGRESS NX EX 120. Return 409 jika command menghasilkan nil.
  • Completion: Simpan hasil JSON serialized via SET idempotency:{key} {serialized_data} XX EX 86400.
  • Recovery: Pastikan durasi TTL IN_PROGRESS pendek (misal 30-120 detik) agar worker crash tidak memicu dead-lock permanen pada key tersebut.