use crate::common::storage::Storage;
use std::time::Duration;
use crate::common::auth::{Role, KEY_STORE};
use crate::common::{AuditEventType, AUDIT_LOGGER};
use async_stream::stream;
use once_cell::sync::Lazy;
use std::convert::Infallible;
use tokio::sync::broadcast;
pub static WATCH_CHANNEL: Lazy<broadcast::Sender<KeyChangeEvent>> = Lazy::new(|| {
let (tx, _rx) = broadcast::channel(100);
tx
});
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct KeyChangeEvent {
pub event: String, pub key: String,
pub tenant: Option<String>,
pub timestamp: i64,
}
pub async fn watch_sse(
) -> Sse<impl futures_util::Stream<Item = Result<axum::response::sse::Event, Infallible>>> {
let mut rx = WATCH_CHANNEL.subscribe();
let stream = stream! {
while let Ok(event) = rx.recv().await {
let data = serde_json::to_string(&event).unwrap();
yield Ok(axum::response::sse::Event::default().data(data));
}
};
Sse::new(stream)
}
pub async fn watch_ws(ws: WebSocketUpgrade) -> impl IntoResponse {
ws.on_upgrade(handle_ws)
}
async fn handle_ws(mut socket: WebSocket) {
let mut rx = WATCH_CHANNEL.subscribe();
while let Ok(event) = rx.recv().await {
let msg = serde_json::to_string(&event).unwrap();
if socket.send(Message::Text(msg)).await.is_err() {
break;
}
}
}
pub static STORAGE: Lazy<Storage> = Lazy::new(Storage::new_memory);
async fn admin_repair(State(_state): State<CoordState>) -> impl IntoResponse {
let res = crate::ops::repair::repair_cluster("http://localhost:5000", 3, false).await;
match res {
Ok(report) => axum::Json(json!({ "status": "ok", "report": report })),
Err(e) => axum::Json(json!({ "status": "error", "error": format!("{}", e) })),
}
}
async fn admin_compact(State(_state): State<CoordState>) -> impl IntoResponse {
let res = crate::ops::compact::compact_cluster("http://localhost:5000", None).await;
match res {
Ok(report) => axum::Json(json!({ "status": "ok", "report": report })),
Err(e) => axum::Json(json!({ "status": "error", "error": format!("{}", e) })),
}
}
async fn admin_verify(State(_state): State<CoordState>) -> impl IntoResponse {
let res = crate::ops::verify::verify_cluster("http://localhost:5000", false, 16).await;
match res {
Ok(report) => axum::Json(json!({ "status": "ok", "report": report })),
Err(e) => axum::Json(json!({ "status": "error", "error": format!("{}", e) })),
}
}
async fn admin_scale(State(_state): State<CoordState>) -> impl IntoResponse {
axum::Json(json!({ "status": "scaling triggered" }))
}
use axum::{
body::Bytes,
extract::{Path, Query, State},
http::StatusCode,
response::IntoResponse,
Router,
};
use serde::{Deserialize, Serialize};
use serde_json::json;
use std::sync::Arc;
use crate::coordinator::metadata::MetadataStore;
use crate::coordinator::placement::PlacementManager;
use crate::coordinator::raft_node::RaftNode;
use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade};
use axum::response::Sse;
#[derive(Debug, Deserialize)]
struct CreateKeyRequest {
name: String,
#[serde(default = "default_tenant")]
tenant: String,
#[serde(default)]
role: String,
expires_in_secs: Option<u64>,
}
fn default_tenant() -> String {
"default".to_string()
}
#[derive(Debug, Serialize)]
struct CreateKeyResponse {
id: String,
key: String,
tenant: String,
role: String,
warning: String,
}
async fn admin_create_key(axum::Json(req): axum::Json<CreateKeyRequest>) -> impl IntoResponse {
let role = match req.role.to_lowercase().as_str() {
"admin" => Role::Admin,
"read_write" | "readwrite" | "rw" => Role::ReadWrite,
"read_only" | "readonly" | "ro" | "" => Role::ReadOnly,
_ => {
return (
StatusCode::BAD_REQUEST,
axum::Json(json!({
"error": "Invalid role",
"valid_roles": ["admin", "read_write", "read_only"]
})),
)
.into_response();
}
};
let expires_in = req.expires_in_secs.map(Duration::from_secs);
match KEY_STORE.generate_key(&req.name, &req.tenant, role, expires_in) {
Ok((id, key)) => {
let response = CreateKeyResponse {
id: id.clone(),
key,
tenant: req.tenant.clone(),
role: format!("{:?}", role),
warning: "Store this key securely - it cannot be retrieved again!".to_string(),
};
AUDIT_LOGGER.log_event(
AuditEventType::ApiKeyCreated,
req.name.clone(),
Some(id.clone()),
format!(
"API key created for tenant {} with role {:?}",
req.tenant.clone(),
role
),
None,
);
(StatusCode::CREATED, axum::Json(json!(response))).into_response()
}
Err(e) => (
StatusCode::INTERNAL_SERVER_ERROR,
axum::Json(json!({ "error": format!("{}", e) })),
)
.into_response(),
}
}
#[derive(Debug, Deserialize)]
struct ListKeysQuery {
tenant: Option<String>,
}
async fn admin_list_keys(Query(query): Query<ListKeysQuery>) -> impl IntoResponse {
let keys = if let Some(tenant) = query.tenant {
KEY_STORE.list_keys_for_tenant(&tenant)
} else {
KEY_STORE.list_keys()
};
let safe_keys: Vec<serde_json::Value> = keys
.iter()
.map(|k| {
json!({
"id": k.id,
"name": k.name,
"tenant": k.tenant,
"role": format!("{:?}", k.role),
"active": k.active,
"created_at": k.created_at,
"expires_at": k.expires_at,
"last_used_at": k.last_used_at,
})
})
.collect();
axum::Json(json!({
"keys": safe_keys,
"total": safe_keys.len()
}))
}
async fn admin_get_key(Path(key_id): Path<String>) -> impl IntoResponse {
match KEY_STORE.get_key(&key_id) {
Some(k) => (
StatusCode::OK,
axum::Json(json!({
"id": k.id,
"name": k.name,
"tenant": k.tenant,
"role": format!("{:?}", k.role),
"active": k.active,
"created_at": k.created_at,
"expires_at": k.expires_at,
"last_used_at": k.last_used_at,
})),
)
.into_response(),
None => (
StatusCode::NOT_FOUND,
axum::Json(json!({ "error": "Key not found" })),
)
.into_response(),
}
}
async fn admin_revoke_key(Path(key_id): Path<String>) -> impl IntoResponse {
match KEY_STORE.revoke_key(&key_id) {
Ok(()) => {
AUDIT_LOGGER.log_event(
AuditEventType::ApiKeyRevoked,
"admin", Some(key_id.clone()),
"API key revoked",
None,
);
let _ = WATCH_CHANNEL.send(KeyChangeEvent {
event: "revoke".to_string(),
key: key_id.clone(),
tenant: None,
timestamp: chrono::Utc::now().timestamp(),
});
(
StatusCode::OK,
axum::Json(json!({ "status": "revoked", "id": key_id })),
)
.into_response()
}
Err(e) => (
StatusCode::NOT_FOUND,
axum::Json(json!({ "error": format!("{}", e) })),
)
.into_response(),
}
}
async fn admin_delete_key(Path(key_id): Path<String>) -> impl IntoResponse {
match KEY_STORE.delete_key(&key_id) {
Ok(()) => {
AUDIT_LOGGER.log_event(
AuditEventType::ApiKeyDeleted,
"admin", Some(key_id.clone()),
"API key deleted",
None,
);
let _ = WATCH_CHANNEL.send(KeyChangeEvent {
event: "delete".to_string(),
key: key_id.clone(),
tenant: None,
timestamp: chrono::Utc::now().timestamp(),
});
(
StatusCode::OK,
axum::Json(json!({ "status": "deleted", "id": key_id })),
)
.into_response()
}
Err(e) => (
StatusCode::NOT_FOUND,
axum::Json(json!({ "error": format!("{}", e) })),
)
.into_response(),
}
}
#[derive(Clone)]
pub struct CoordState {
pub metadata: Arc<MetadataStore>,
pub placement: Arc<std::sync::Mutex<PlacementManager>>,
pub raft: Arc<RaftNode>,
}
async fn s3_put_object(
State(state): State<CoordState>,
Path((bucket, key)): Path<(String, String)>,
headers: axum::http::HeaderMap,
body: Bytes,
) -> impl IntoResponse {
let full_key = format!("{}/{}", bucket, key);
let ttl_secs: Option<u64> = headers
.get("X-Minikv-TTL")
.and_then(|v| v.to_str().ok())
.and_then(|s| s.parse().ok());
crate::coordinator::http::STORAGE.put(&full_key, body.to_vec());
let stored_bytes = body.len();
let _ = WATCH_CHANNEL.send(KeyChangeEvent {
event: "put".to_string(),
key: full_key.clone(),
tenant: Some("default".to_string()), timestamp: chrono::Utc::now().timestamp(),
});
let placement = state.placement.lock().unwrap();
let volumes = state.metadata.get_healthy_volumes().unwrap_or_default();
let target_volumes: Vec<String> = placement
.select_volumes(&full_key, &volumes)
.unwrap_or_default();
let mut prepare_ok = true;
for _volume_id in &target_volumes {
let simulated_prepare = true;
if !simulated_prepare {
prepare_ok = false;
break;
}
}
if !prepare_ok {
return (
StatusCode::INTERNAL_SERVER_ERROR,
format!(
"PUT S3 {}/{} failed: prepare phase error (2PC)",
bucket, key
),
);
}
for _volume_id in &target_volumes {
}
let ttl_info = ttl_secs
.map(|t| format!(", TTL: {}s", t))
.unwrap_or_default();
(
StatusCode::OK,
format!(
"PUT S3 {}/{} committed via 2PC ({} bytes{})",
bucket, key, stored_bytes, ttl_info
),
)
}
async fn s3_get_object(
State(_state): State<CoordState>,
Path((bucket, key)): Path<(String, String)>,
) -> impl IntoResponse {
let full_key = format!("{}/{}", bucket, key);
if let Some(data) = crate::coordinator::http::STORAGE.get(&full_key) {
(StatusCode::OK, data)
} else {
(
StatusCode::NOT_FOUND,
format!("S3 object {}/{} not found", bucket, key).into_bytes(),
)
}
}
pub fn create_router(state: CoordState) -> Router {
Router::new()
.route("/watch/sse", axum::routing::get(watch_sse))
.route("/watch/ws", axum::routing::get(watch_ws))
.route("/s3/:bucket/:key", axum::routing::put(s3_put_object))
.route("/s3/:bucket/:key", axum::routing::get(s3_get_object))
.route("/health", axum::routing::get(health))
.route("/health/ready", axum::routing::get(health_ready))
.route("/health/live", axum::routing::get(health_live))
.route("/:key", axum::routing::post(put_key))
.route("/:key", axum::routing::get(get_key))
.route("/:key", axum::routing::delete(delete_key))
.route("/admin/repair", axum::routing::post(admin_repair))
.route("/admin/compact", axum::routing::post(admin_compact))
.route("/admin/verify", axum::routing::post(admin_verify))
.route("/admin/scale", axum::routing::post(admin_scale))
.route("/admin/status", axum::routing::get(admin_status))
.route("/admin/keys", axum::routing::post(admin_create_key))
.route("/admin/keys", axum::routing::get(admin_list_keys))
.route("/admin/keys/:key_id", axum::routing::get(admin_get_key))
.route(
"/admin/keys/:key_id/revoke",
axum::routing::post(admin_revoke_key),
)
.route(
"/admin/keys/:key_id",
axum::routing::delete(admin_delete_key),
)
.route("/admin/import", axum::routing::post(admin_import))
.route("/admin/export", axum::routing::get(admin_export))
.route("/transaction", axum::routing::post(transaction_ops))
.route("/search", axum::routing::get(search_keys))
.route("/metrics", axum::routing::get(metrics))
.route("/range", axum::routing::get(range_query))
.route("/batch", axum::routing::post(batch_ops))
.with_state(state)
}
async fn health_ready(State(state): State<CoordState>) -> impl IntoResponse {
let volumes = state.metadata.get_healthy_volumes().unwrap_or_default();
let has_leader = state.raft.is_leader() || !state.raft.get_peers().is_empty();
if !volumes.is_empty() && has_leader {
(
StatusCode::OK,
axum::Json(json!({
"ready": true,
"healthy_volumes": volumes.len(),
"is_leader": state.raft.is_leader(),
})),
)
} else {
(
StatusCode::SERVICE_UNAVAILABLE,
axum::Json(json!({
"ready": false,
"healthy_volumes": volumes.len(),
"is_leader": state.raft.is_leader(),
"reason": if volumes.is_empty() { "No healthy volumes" } else { "No Raft leader" }
})),
)
}
}
async fn health_live() -> impl IntoResponse {
(
StatusCode::OK,
axum::Json(json!({
"alive": true,
"version": env!("CARGO_PKG_VERSION"),
"timestamp": std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs(),
})),
)
}
async fn admin_status(State(state): State<CoordState>) -> impl IntoResponse {
let role = if state.raft.is_leader() {
"Leader"
} else {
"Follower"
};
let nb_peers = state.raft.get_peers().len();
let volumes = state.metadata.get_healthy_volumes().unwrap_or_default();
let nb_volumes = volumes.len();
let volume_ids: Vec<_> = volumes.iter().map(|v| v.volume_id.clone()).collect();
let nb_s3_objects = 0;
axum::Json(json!({
"role": role,
"is_leader": state.raft.is_leader(),
"nb_peers": nb_peers,
"nb_volumes": nb_volumes,
"volume_ids": volume_ids,
"nb_s3_objects": nb_s3_objects
}))
}
#[derive(Deserialize)]
struct ImportRequest {
entries: Vec<KeyValueEntry>,
}
#[derive(Deserialize)]
struct KeyValueEntry {
key: String,
value: String,
}
async fn admin_import(
State(_state): State<CoordState>,
axum::Json(req): axum::Json<ImportRequest>,
) -> impl IntoResponse {
let mut success_count = 0;
let errors: Vec<String> = Vec::new();
for entry in req.entries {
STORAGE.put(&entry.key, entry.value.into_bytes());
success_count += 1;
}
AUDIT_LOGGER.log_event(
AuditEventType::System,
"admin".to_string(),
None,
format!("Imported {} keys", success_count),
None,
);
axum::Json(json!({
"imported": success_count,
"errors": errors
}))
}
async fn admin_export(State(state): State<CoordState>) -> impl IntoResponse {
let keys = match state.metadata.list_keys() {
Ok(keys) => keys,
Err(e) => {
return (
StatusCode::INTERNAL_SERVER_ERROR,
format!("list_keys error: {}", e),
)
.into_response();
}
};
let body = stream! {
for key in keys {
if let Some(value) = STORAGE.get(&key) {
let entry = json!({
"key": key,
"value": String::from_utf8_lossy(&value)
});
yield Ok::<_, std::convert::Infallible>(axum::body::Bytes::from(format!("{}\n", entry)));
}
}
};
(
StatusCode::OK,
[("content-type", "application/x-ndjson")],
axum::body::Body::from_stream(body),
)
.into_response()
}
#[derive(Deserialize)]
struct TransactionRequest {
operations: Vec<Operation>,
}
#[derive(Deserialize)]
struct Operation {
op: String, key: String,
value: Option<String>,
}
async fn transaction_ops(
State(_state): State<CoordState>,
axum::Json(req): axum::Json<TransactionRequest>,
) -> impl IntoResponse {
let mut results = Vec::new();
let mut success_count = 0;
let total_operations = req.operations.len();
for op in &req.operations {
match op.op.as_str() {
"put" => {
if let Some(ref value) = op.value {
STORAGE.put(&op.key, value.clone().into_bytes());
success_count += 1;
results.push(TransactionResult {
op: op.op.clone(),
key: op.key.clone(),
success: true,
error: None,
});
} else {
results.push(TransactionResult {
op: op.op.clone(),
key: op.key.clone(),
success: false,
error: Some("value required for put".to_string()),
});
}
}
"delete" => {
STORAGE.delete(&op.key);
success_count += 1;
results.push(TransactionResult {
op: op.op.clone(),
key: op.key.clone(),
success: true,
error: None,
});
}
_ => {
results.push(TransactionResult {
op: op.op.clone(),
key: op.key.clone(),
success: false,
error: Some("unknown operation".to_string()),
});
}
}
}
AUDIT_LOGGER.log_event(
AuditEventType::System,
"transaction".to_string(),
None,
format!("Executed {} operations in transaction", success_count),
None,
);
axum::Json(json!({
"results": results,
"total_operations": total_operations,
"successful_operations": success_count
}))
}
#[derive(Deserialize)]
struct SearchQuery {
value: String,
}
async fn search_keys(
State(state): State<CoordState>,
Query(params): Query<SearchQuery>,
) -> impl IntoResponse {
match state.metadata.list_keys() {
Ok(keys) => {
let mut matching_keys = Vec::new();
for key in keys {
if let Some(value_bytes) = STORAGE.get(&key) {
if let Ok(value_str) = std::str::from_utf8(&value_bytes) {
if value_str.contains(¶ms.value) {
matching_keys.push(key);
}
}
}
}
axum::Json(json!({
"query": params.value,
"matching_keys": matching_keys,
"total_matches": matching_keys.len()
}))
}
Err(e) => axum::Json(json!({ "error": format!("list_keys error: {}", e) })),
}
}
#[derive(Serialize)]
struct TransactionResult {
op: String,
key: String,
success: bool,
error: Option<String>,
}
#[derive(Deserialize)]
struct RangeQuery {
start: String,
end: String,
include_values: Option<bool>,
}
async fn range_query(
State(state): State<CoordState>,
Query(params): Query<RangeQuery>,
) -> impl IntoResponse {
let keys = match state.metadata.list_keys() {
Ok(keys) => keys,
Err(e) => {
return (
StatusCode::INTERNAL_SERVER_ERROR,
format!("list_keys error: {}", e),
)
}
};
let mut filtered: Vec<String> = keys
.into_iter()
.filter(|k| k >= ¶ms.start && k <= ¶ms.end)
.collect();
filtered.sort();
if params.include_values.unwrap_or(false) {
let mut values = Vec::new();
for k in &filtered {
match state.metadata.get_key(k) {
Ok(Some(meta)) => values.push(serde_json::to_value(&meta).unwrap_or(json!(null))),
_ => values.push(json!(null)),
}
}
(
StatusCode::OK,
serde_json::to_string(&json!({ "keys": filtered, "values": values })).unwrap(),
)
} else {
(
StatusCode::OK,
serde_json::to_string(&json!({ "keys": filtered })).unwrap(),
)
}
}
#[derive(Deserialize)]
struct BatchOpReq {
op: String, key: String,
value: Option<String>,
}
#[derive(Deserialize)]
struct BatchReq {
ops: Vec<BatchOpReq>,
}
#[derive(Serialize)]
struct BatchResultResp {
ok: bool,
key: String,
value: Option<String>,
error: Option<String>,
}
async fn batch_ops(
State(state): State<CoordState>,
axum::Json(req): axum::Json<BatchReq>,
) -> impl IntoResponse {
let mut results = Vec::new();
for op in req.ops {
match op.op.as_str() {
"put" => {
if let Some(val) = op.value {
let meta = crate::coordinator::metadata::KeyMetadata {
key: op.key.clone(),
replicas: vec![],
size: val.len() as u64,
blake3: "".to_string(),
created_at: 0,
updated_at: 0,
state: crate::coordinator::metadata::KeyState::Active,
};
let r = state.metadata.put_key(&meta);
results.push(BatchResultResp {
ok: r.is_ok(),
key: op.key,
value: None,
error: r.err().map(|e| format!("{}", e)),
});
} else {
results.push(BatchResultResp {
ok: false,
key: op.key,
value: None,
error: Some("Missing value for put".to_string()),
});
}
}
"get" => {
let r = state.metadata.get_key(&op.key);
match r {
Ok(Some(meta)) => results.push(BatchResultResp {
ok: true,
key: op.key,
value: Some(serde_json::to_string(&meta).unwrap()),
error: None,
}),
Ok(None) => results.push(BatchResultResp {
ok: false,
key: op.key,
value: None,
error: Some("Not found".to_string()),
}),
Err(e) => results.push(BatchResultResp {
ok: false,
key: op.key,
value: None,
error: Some(format!("{}", e)),
}),
}
}
"delete" => {
let r = state.metadata.delete_key(&op.key);
results.push(BatchResultResp {
ok: r.is_ok(),
key: op.key,
value: None,
error: r.err().map(|e| format!("{}", e)),
});
}
_ => results.push(BatchResultResp {
ok: false,
key: op.key,
value: None,
error: Some("Unknown op".to_string()),
}),
}
}
axum::Json(json!({ "results": results }))
}
pub async fn metrics(State(state): State<CoordState>) -> impl IntoResponse {
let mut out = String::new();
let volumes: Vec<crate::coordinator::metadata::VolumeMetadata> =
state.metadata.get_healthy_volumes().unwrap_or_default();
let total_keys: u64 = volumes.iter().map(|v| v.total_keys).sum();
out += &format!("minikv_total_keys {{}} {}\n", total_keys);
out += &format!("minikv_healthy_volumes {{}} {}\n", volumes.len());
for v in &volumes {
out += &format!(
"minikv_volume_bytes {{volume_id=\"{}\"}} {}\n",
v.volume_id, v.total_bytes
);
out += &format!(
"minikv_volume_free_bytes {{volume_id=\"{}\"}} {}\n",
v.volume_id, v.free_bytes
);
out += &format!(
"minikv_volume_total_keys {{volume_id=\"{}\"}} {}\n",
v.volume_id, v.total_keys
);
}
let role = if state.raft.is_leader() {
"leader"
} else {
"follower"
};
out += &format!("minikv_raft_role {{}} \"{}\"\n", role);
out += &crate::common::METRICS.to_prometheus();
let s3_objects = 0;
let s3_objects_with_ttl = 0;
out += &format!("minikv_s3_objects_total {{}} {}\n", s3_objects);
out += &format!("minikv_s3_objects_with_ttl {{}} {}\n", s3_objects_with_ttl);
(axum::http::StatusCode::OK, out)
}
async fn health(State(state): State<CoordState>) -> impl IntoResponse {
let role = if state.raft.is_leader() {
"Leader"
} else {
"Follower"
};
axum::Json(json!({
"status": "healthy",
"role": role,
"is_leader": state.raft.is_leader(),
"version": env!("CARGO_PKG_VERSION"),
}))
}
async fn put_key(
State(state): State<CoordState>,
Path(key): Path<String>,
_body: Bytes,
) -> impl IntoResponse {
let placement = state.placement.lock().unwrap();
let volumes = state.metadata.get_healthy_volumes().unwrap_or_default();
let target_volumes: Vec<String> = placement.select_volumes(&key, &volumes).unwrap_or_default();
let mut prepare_ok = true;
for _volume_id in &target_volumes {
let simulated_prepare = true;
if !simulated_prepare {
prepare_ok = false;
break;
}
}
if !prepare_ok {
for _volume_id in &target_volumes {
}
return (
StatusCode::INTERNAL_SERVER_ERROR,
format!("PUT {} failed: prepare phase error (2PC)", key),
);
}
for _volume_id in &target_volumes {
}
(StatusCode::OK, format!("PUT {} committed via 2PC", key))
}
async fn get_key(State(_state): State<CoordState>, Path(key): Path<String>) -> impl IntoResponse {
let value = format!("Value for key {} (fetched from volume)", key);
(StatusCode::OK, value)
}
async fn delete_key(
State(_state): State<CoordState>,
Path(key): Path<String>,
) -> impl IntoResponse {
(StatusCode::OK, format!("DELETE {} succeeded", key))
}