#![cfg(feature = "real-nacos")]
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
#[derive(Debug, thiserror::Error)]
pub enum NacosError {
#[error("HTTP error: {0}")]
Http(#[from] reqwest::Error),
#[error("Nacos API error: status={status}, body={body}")]
Api { status: u16, body: String },
#[error("Config not found: {0}")]
NotFound(String),
#[error("Auth error: {0}")]
Auth(String),
#[error("Invalid response: {0}")]
InvalidResponse(String),
}
#[derive(Debug, Clone)]
pub struct NacosConfig {
pub endpoint: String,
pub username: Option<String>,
pub password: Option<String>,
pub namespace: Option<String>,
pub timeout_secs: u64,
}
impl Default for NacosConfig {
fn default() -> Self {
Self {
endpoint: "http://127.0.0.1:8848".to_string(),
username: None,
password: None,
namespace: None,
timeout_secs: 10,
}
}
}
impl NacosConfig {
pub fn new(endpoint: impl Into<String>) -> Self {
Self {
endpoint: endpoint.into(),
..Default::default()
}
}
pub fn with_auth(mut self, username: impl Into<String>, password: impl Into<String>) -> Self {
self.username = Some(username.into());
self.password = Some(password.into());
self
}
pub fn with_namespace(mut self, ns: impl Into<String>) -> Self {
self.namespace = Some(ns.into());
self
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct NacosLoginResponse {
#[serde(rename = "accessToken")]
access_token: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NacosServiceRegistration {
pub name: String,
pub ip: String,
pub port: u16,
#[serde(skip_serializing_if = "Option::is_none")]
pub weight: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub metadata: Option<HashMap<String, String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cluster_name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub group_name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub enabled: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ephemeral: Option<bool>,
}
pub struct NacosClient {
config: NacosConfig,
http: reqwest::Client,
access_token: Option<String>,
}
impl NacosClient {
pub fn new(config: NacosConfig) -> Result<Self, NacosError> {
let http = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(config.timeout_secs))
.build()?;
Ok(Self {
config,
http,
access_token: None,
})
}
fn build_url(&self, path: &str) -> String {
format!("{}{}", self.config.endpoint, path)
}
fn add_auth(&self, req: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
if let Some(token) = &self.access_token {
req.query(&[("accessToken", token)])
} else {
req
}
}
fn add_namespace(&self, req: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
if let Some(ns) = &self.config.namespace {
req.query(&[("tenant", ns)])
} else {
req
}
}
pub async fn login(&mut self) -> Result<(), NacosError> {
let (username, password) = match (&self.config.username, &self.config.password) {
(Some(u), Some(p)) => (u.clone(), p.clone()),
_ => return Err(NacosError::Auth("username and password required".into())),
};
let url = self.build_url("/nacos/v1/auth/login");
let resp = self
.http
.post(&url)
.form(&[
("username", username.as_str()),
("password", password.as_str()),
])
.send()
.await?;
if !resp.status().is_success() {
let status = resp.status().as_u16();
let body = resp.text().await.unwrap_or_default();
return Err(NacosError::Api { status, body });
}
let login_resp: NacosLoginResponse = resp.json().await?;
self.access_token = Some(login_resp.access_token);
Ok(())
}
pub async fn get_config(&self, data_id: &str, group: &str) -> Result<String, NacosError> {
let url = self.build_url("/nacos/v1/cs/configs");
let req = self
.http
.get(&url)
.query(&[("dataId", data_id), ("group", group)]);
let req = self.add_auth(req);
let req = self.add_namespace(req);
let resp = req.send().await?;
if resp.status() == reqwest::StatusCode::NOT_FOUND
|| resp.status() == reqwest::StatusCode::CONFLICT
{
return Err(NacosError::NotFound(format!("{}/{}", group, data_id)));
}
if !resp.status().is_success() {
let status = resp.status().as_u16();
let body = resp.text().await.unwrap_or_default();
return Err(NacosError::Api { status, body });
}
let config = resp.text().await?;
if config.is_empty() {
return Err(NacosError::NotFound(format!("{}/{}", group, data_id)));
}
Ok(config)
}
pub async fn set_config(
&self,
data_id: &str,
group: &str,
content: &str,
) -> Result<(), NacosError> {
let url = self.build_url("/nacos/v1/cs/configs");
let mut form: Vec<(&str, &str)> =
vec![("dataId", data_id), ("group", group), ("content", content)];
let ns_val: String;
if let Some(ns) = &self.config.namespace {
ns_val = ns.clone();
form.push(("tenant", ns_val.as_str()));
}
let req = self.http.post(&url).form(&form);
let req = self.add_auth(req);
let resp = req.send().await?;
if !resp.status().is_success() {
let status = resp.status().as_u16();
let body = resp.text().await.unwrap_or_default();
return Err(NacosError::Api { status, body });
}
Ok(())
}
pub async fn delete_config(&self, data_id: &str, group: &str) -> Result<(), NacosError> {
let url = self.build_url("/nacos/v1/cs/configs");
let req = self
.http
.delete(&url)
.query(&[("dataId", data_id), ("group", group)]);
let req = self.add_auth(req);
let req = self.add_namespace(req);
let resp = req.send().await?;
if !resp.status().is_success() {
let status = resp.status().as_u16();
let body = resp.text().await.unwrap_or_default();
return Err(NacosError::Api { status, body });
}
Ok(())
}
pub async fn watch(
&self,
data_id: &str,
group: &str,
timeout_millis: u64,
) -> Result<String, NacosError> {
let url = self.build_url("/nacos/v1/cs/configs/listener");
let listening_str = format!(
"{}^2{}^1{}",
data_id,
group,
self.config.namespace.as_deref().unwrap_or("")
);
let req = self.http.post(&url).query(&[
("Listening-Configs", listening_str.as_str()),
("timeout", &timeout_millis.to_string()),
]);
let req = self.add_auth(req);
let resp = req.send().await?;
if !resp.status().is_success() {
let status = resp.status().as_u16();
let body = resp.text().await.unwrap_or_default();
return Err(NacosError::Api { status, body });
}
let content = resp.text().await?;
if content.is_empty() {
return Err(NacosError::NotFound(format!("{}/{}", group, data_id)));
}
Ok(content)
}
pub async fn register_service(
&self,
service: NacosServiceRegistration,
) -> Result<(), NacosError> {
let url = self.build_url("/nacos/v1/ns/instance");
let mut form: Vec<(&str, String)> = vec![
("serviceName", service.name.clone()),
("ip", service.ip.clone()),
("port", service.port.to_string()),
];
if let Some(w) = service.weight {
form.push(("weight", w.to_string()));
}
if let Some(cluster) = &service.cluster_name {
form.push(("clusterName", cluster.clone()));
}
if let Some(group) = &service.group_name {
form.push(("groupName", group.clone()));
}
if let Some(enabled) = service.enabled {
form.push(("enabled", enabled.to_string()));
}
if let Some(ephemeral) = service.ephemeral {
form.push(("ephemeral", ephemeral.to_string()));
}
if let Some(metadata) = &service.metadata {
let meta_str = metadata
.iter()
.map(|(k, v)| format!("{}={}", k, v))
.collect::<Vec<_>>()
.join(",");
form.push(("metadata", meta_str));
}
let form_refs: Vec<(&str, &str)> = form.iter().map(|(k, v)| (*k, v.as_str())).collect();
let req = self.http.post(&url).form(&form_refs);
let req = self.add_auth(req);
let resp = req.send().await?;
if !resp.status().is_success() {
let status = resp.status().as_u16();
let body = resp.text().await.unwrap_or_default();
return Err(NacosError::Api { status, body });
}
Ok(())
}
pub async fn deregister_service(
&self,
name: &str,
ip: &str,
port: u16,
) -> Result<(), NacosError> {
let url = self.build_url("/nacos/v1/ns/instance");
let req = self.http.delete(&url).query(&[
("serviceName", name),
("ip", ip),
("port", &port.to_string()),
]);
let req = self.add_auth(req);
let resp = req.send().await?;
if !resp.status().is_success() {
let status = resp.status().as_u16();
let body = resp.text().await.unwrap_or_default();
return Err(NacosError::Api { status, body });
}
Ok(())
}
}