use std::sync::Arc;
use distributed::microsvc::{self, Service};
use distributed::{AsyncAggregateBuilder, HashMapRepository, Queueable};
use serde_json::json;
use crate::handlers;
use crate::handlers::Repo;
use crate::models::counter::Counter;
fn counter_service() -> Arc<Service<Repo>> {
Arc::new(distributed::register_handlers!(
Service::with_repo(HashMapRepository::new().queued_async().async_aggregate::<Counter>()),
command handlers::counter_create,
command handlers::counter_increment,
command handlers::whoami,
))
}
async fn start_server(service: Arc<Service<Repo>>) -> String {
let app = microsvc::router(service);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
format!("http://{addr}")
}
#[tokio::test]
async fn health_check() {
let service = counter_service();
let base = start_server(service).await;
let client = reqwest::Client::new();
let resp = client.get(format!("{base}/health")).send().await.unwrap();
assert_eq!(resp.status(), 200);
let body: serde_json::Value = resp.json().await.unwrap();
assert_eq!(body["ok"], true);
let commands = body["commands"].as_array().unwrap();
assert!(commands.iter().any(|c| c == "counter.initialize"));
assert!(commands.iter().any(|c| c == "counter.increment"));
}
#[tokio::test]
async fn create_counter() {
let service = counter_service();
let base = start_server(service).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base}/counter.initialize"))
.json(&json!({ "id": "c1" }))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
let body: serde_json::Value = resp.json().await.unwrap();
assert_eq!(body, json!({ "id": "c1" }));
}
#[tokio::test]
async fn create_and_increment_counter() {
let service = counter_service();
let base = start_server(service).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base}/counter.initialize"))
.json(&json!({ "id": "c1" }))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
let resp = client
.post(format!("{base}/counter.increment"))
.json(&json!({ "id": "c1", "amount": 5 }))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
let body: serde_json::Value = resp.json().await.unwrap();
assert_eq!(body, json!({ "id": "c1", "value": 5 }));
let resp = client
.post(format!("{base}/counter.increment"))
.json(&json!({ "id": "c1", "amount": 3 }))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
let body: serde_json::Value = resp.json().await.unwrap();
assert_eq!(body, json!({ "id": "c1", "value": 8 }));
}
#[tokio::test]
async fn increment_nonexistent_returns_404() {
let service = counter_service();
let base = start_server(service).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base}/counter.increment"))
.json(&json!({ "id": "nope", "amount": 1 }))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 404);
}
#[tokio::test]
async fn unknown_command_returns_404() {
let service = counter_service();
let base = start_server(service).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base}/nonexistent"))
.json(&json!({}))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 404);
}
#[tokio::test]
async fn headers_flow_to_session() {
let service = counter_service();
let base = start_server(service).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base}/session.identify"))
.header("x-hasura-user-id", "user-42")
.json(&json!({}))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
let body: serde_json::Value = resp.json().await.unwrap();
assert_eq!(body, json!({ "user_id": "user-42" }));
}
#[tokio::test]
async fn missing_session_returns_401() {
let service = counter_service();
let base = start_server(service).await;
let client = reqwest::Client::new();
let resp = client
.post(format!("{base}/session.identify"))
.json(&json!({}))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 401);
}