use std::io::Cursor;
use actix_web::App;
use actix_web::HttpServer;
use actix_web::web;
use actix_web::web::Data;
use openraft::BasicNode;
use openraft::ChangeMembers;
use openraft::Snapshot;
use openraft::errors::Infallible;
use openraft::errors::decompose::DecomposeResult;
use openraft::raft;
use openraft::raft::SnapshotResponse;
use serde::Deserialize;
use crate::app::EzApp;
use crate::network::SnapshotTransfer;
use crate::raft::EzRaft;
use crate::type_config::OpenRaftTypes;
type C<T> = OpenRaftTypes<T>;
pub struct EzServer<T>
where T: EzApp
{
raft: EzRaft<T>,
}
impl<T> EzServer<T>
where T: EzApp
{
pub fn new(raft: EzRaft<T>) -> Self {
Self { raft }
}
pub async fn run(self) -> std::io::Result<()> {
let addr = self.raft.addr().to_string();
let server_data = Data::new(self);
let server = HttpServer::new(move || {
App::new()
.app_data(server_data.clone())
.route("/raft/append", web::post().to(Self::handle_append))
.route("/raft/vote", web::post().to(Self::handle_vote))
.route("/raft/snapshot", web::post().to(Self::handle_snapshot))
.route("/api/write", web::post().to(Self::handle_write))
.route("/api/read", web::get().to(Self::handle_read))
.route("/api/join", web::post().to(Self::handle_join))
.route("/api/change_membership", web::post().to(Self::handle_change_membership))
.route("/api/metrics", web::get().to(Self::handle_metrics))
})
.bind(&addr)?;
server.run().await
}
async fn handle_append(
req: web::Json<raft::AppendEntriesRequest<C<T>>>,
ez: Data<Self>,
) -> Result<web::Json<Result<raft::AppendEntriesResponse<C<T>>, Infallible>>, actix_web::Error> {
let resp = ez
.raft
.inner()
.append_entries(req.into_inner())
.await
.decompose()
.map_err(|e| actix_web::error::ErrorInternalServerError(format!("append_entries failed: {}", e)))?;
Ok(web::Json(resp))
}
async fn handle_vote(
req: web::Json<raft::VoteRequest<C<T>>>,
ez: Data<Self>,
) -> Result<web::Json<Result<raft::VoteResponse<C<T>>, Infallible>>, actix_web::Error> {
let resp = ez
.raft
.inner()
.vote(req.into_inner())
.await
.decompose()
.map_err(|e| actix_web::error::ErrorInternalServerError(format!("vote failed: {}", e)))?;
Ok(web::Json(resp))
}
async fn handle_snapshot(
req: web::Json<SnapshotTransfer>,
ez: Data<Self>,
) -> Result<web::Json<Result<SnapshotResponse<C<T>>, Infallible>>, actix_web::Error> {
let SnapshotTransfer { vote, meta, data } = req.into_inner();
let snapshot = Snapshot {
meta,
snapshot: Cursor::new(data),
};
let resp = ez
.raft
.inner()
.install_full_snapshot(vote, snapshot)
.await
.map_err(|e| actix_web::error::ErrorInternalServerError(format!("install_snapshot failed: {}", e)))?;
Ok(web::Json(Ok(resp)))
}
async fn handle_write(
req: web::Json<T::Request>,
ez: Data<Self>,
) -> Result<web::Json<T::Response>, actix_web::Error> {
let resp = ez
.raft
.write(req.into_inner())
.await
.map_err(|e| actix_web::error::ErrorInternalServerError(format!("write failed: {}", e)))?;
Ok(web::Json(resp))
}
async fn handle_read(
query: web::Query<ReadQuery>,
ez: Data<Self>,
) -> Result<web::Json<serde_json::Value>, actix_web::Error> {
let key = &query.key;
let Some(value) = ez.raft.read(|app| app.read(key)).await else {
return Err(actix_web::error::ErrorNotFound(format!("no value for key {:?}", key)));
};
Ok(web::Json(value))
}
async fn handle_change_membership(
req: web::Json<ChangeMembers<u64, BasicNode>>,
ez: Data<Self>,
) -> Result<web::Json<serde_json::Value>, actix_web::Error> {
ez.raft
.change_membership(req.into_inner())
.await
.map_err(|e| actix_web::error::ErrorInternalServerError(format!("change_membership failed: {}", e)))?;
Ok(web::Json(serde_json::json!({ "status": "ok" })))
}
async fn handle_metrics(ez: Data<Self>) -> Result<web::Json<openraft::RaftMetrics<C<T>>>, actix_web::Error> {
let metrics = ez.raft.metrics().await;
Ok(web::Json(metrics))
}
async fn handle_join(
req: web::Json<JoinRequest>,
ez: Data<Self>,
) -> Result<web::Json<JoinResponse>, actix_web::Error> {
let metrics = ez.raft.metrics().await;
if metrics.current_leader != Some(metrics.id) {
let leader_addr = metrics.current_leader.and_then(|leader_id| {
metrics.membership_config.membership().get_node(&leader_id).map(|n| n.addr.clone())
});
return Ok(web::Json(Err(leader_addr)));
}
let write_result = ez
.raft
.inner()
.write_blank()
.await
.map_err(|e| actix_web::error::ErrorInternalServerError(format!("join write failed: {}", e)))?;
let node_id = write_result.log_id.index;
ez.raft
.add_learner(node_id, req.addr.clone())
.await
.map_err(|e| actix_web::error::ErrorInternalServerError(format!("add_learner failed: {}", e)))?;
Ok(web::Json(Ok(node_id)))
}
}
pub(crate) async fn run<T>(raft: EzRaft<T>) -> std::io::Result<()>
where T: EzApp {
EzServer::new(raft).run().await
}
#[derive(Debug, Deserialize)]
struct ReadQuery {
key: String,
}
#[derive(Debug, Deserialize)]
struct JoinRequest {
addr: String,
}
type JoinResponse = Result<u64, Option<String>>;