diff --git a/crates/dbx-core/src/nacos/http.rs b/crates/dbx-core/src/nacos/http.rs index 306927234..7ba8d79a9 100644 --- a/crates/dbx-core/src/nacos/http.rs +++ b/crates/dbx-core/src/nacos/http.rs @@ -1,4 +1,4 @@ -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::time::{Duration, Instant}; use async_trait::async_trait; @@ -340,6 +340,68 @@ impl NacosOpenApiAdmin { } list } + + async fn list_v1_catalog_instances( + &self, + query: &NacosInstanceQuery, + namespace: &str, + ) -> Result, String> { + // Nacos catalog controllers derive the group from serviceName and ignore a separate groupName parameter. + let catalog_service_name = qualified_nacos_service_name(&query.service_name, query.group_name.as_deref()); + let mut cluster_names = split_nacos_cluster_names(query.clusters.as_deref()); + if cluster_names.is_empty() { + let detail_params = vec![ + ("serviceName".to_string(), catalog_service_name.clone()), + ("namespaceId".to_string(), namespace.to_string()), + ]; + let detail = self.get_json("/v1/ns/catalog/service", detail_params).await?; + cluster_names = parse_catalog_cluster_names(&detail); + } + + let page_size = self.cfg.page_size.max(100).clamp(1, 500); + let mut instances = Vec::new(); + for cluster_name in cluster_names { + let mut page_no = 1u32; + let mut loaded = 0u64; + loop { + let params = vec![ + ("serviceName".to_string(), catalog_service_name.clone()), + ("namespaceId".to_string(), namespace.to_string()), + ("clusterName".to_string(), cluster_name.clone()), + ("pageNo".to_string(), page_no.to_string()), + ("pageSize".to_string(), page_size.to_string()), + ]; + let value = self.get_json("/v1/ns/catalog/instances", params).await?; + let total_count = catalog_instance_count(&value); + let page = parse_instances(value); + let page_len = page.len(); + loaded = loaded.saturating_add(page_len as u64); + instances.extend(page); + + let has_more = total_count + .filter(|total| *total > 0) + .map(|total| loaded < total) + .unwrap_or(page_len == page_size as usize); + if !has_more || page_len == 0 { + break; + } + page_no = page_no + .checked_add(1) + .ok_or_else(|| "Nacos instance pagination exceeded the supported page range".to_string())?; + } + } + + let mut seen = HashSet::new(); + instances.retain(|instance| seen.insert((instance.ip.clone(), instance.port, instance.cluster_name.clone()))); + Ok(instances) + } +} + +fn qualified_nacos_service_name(service_name: &str, group_name: Option<&str>) -> String { + match group_name.map(str::trim).filter(|group| !group.is_empty()) { + Some(group) => format!("{group}@@{service_name}"), + None => service_name.to_string(), + } } #[async_trait] @@ -795,20 +857,33 @@ impl NacosAdmin for NacosOpenApiAdmin { async fn list_instances(&self, query: NacosInstanceQuery) -> Result, String> { let namespace = self.namespace(query.namespace.as_deref()); - let mut params = vec![("serviceName".to_string(), query.service_name), ("namespaceId".to_string(), namespace)]; - push_optional(&mut params, "groupName", query.group_name); - push_optional(&mut params, "clusters", query.clusters); - let value = self - .get_json_from_candidates( - "list Nacos instances", - vec![ - ("/v3/console/ns/instance/list", params.clone()), - ("/v3/console/ns/instance", params.clone()), - ("/v1/ns/instance/list", params), - ], - ) - .await?; - Ok(parse_instances(value)) + let mut params = vec![ + ("serviceName".to_string(), query.service_name.clone()), + ("namespaceId".to_string(), namespace.clone()), + ]; + push_optional(&mut params, "groupName", query.group_name.clone()); + push_optional(&mut params, "clusters", query.clusters.clone()); + + let mut errors = Vec::new(); + for path in ["/v3/console/ns/instance/list", "/v3/console/ns/instance"] { + match self.get_json(path, params.clone()).await { + Ok(value) => return Ok(parse_instances(value)), + Err(err) => errors.push(err), + } + } + + match self.list_v1_catalog_instances(&query, &namespace).await { + Ok(instances) => return Ok(instances), + Err(err) => errors.push(err), + } + + match self.get_json("/v1/ns/instance/list", params).await { + Ok(value) => Ok(parse_instances(value)), + Err(err) => { + errors.push(err); + Err(format!("Failed to list Nacos instances: {}", errors.join("; "))) + } + } } async fn update_instance(&self, req: NacosInstanceUpdate) -> Result<(), String> { @@ -1279,14 +1354,51 @@ fn split_nacos_service_name(value: &str) -> (Option, String) { (None, trimmed.to_string()) } +fn split_nacos_cluster_names(value: Option<&str>) -> Vec { + let mut seen = HashSet::new(); + value + .unwrap_or_default() + .split(',') + .map(str::trim) + .filter(|name| !name.is_empty()) + .filter(|name| seen.insert((*name).to_string())) + .map(str::to_string) + .collect() +} + +fn parse_catalog_cluster_names(value: &Value) -> Vec { + let data = value.get("data").unwrap_or(value); + let mut seen = HashSet::new(); + data.get("clusters") + .or_else(|| value.get("clusters")) + .and_then(Value::as_array) + .into_iter() + .flatten() + .filter_map(|cluster| optional_string_field(cluster, &["name", "clusterName"])) + .filter(|name| seen.insert(name.clone())) + .collect() +} + +fn catalog_instance_count(value: &Value) -> Option { + let data = value.get("data").unwrap_or(value); + data.get("count") + .or_else(|| data.get("totalCount")) + .or_else(|| data.get("total")) + .or_else(|| value.get("count")) + .or_else(|| value.get("totalCount")) + .and_then(Value::as_u64) +} + fn parse_instances(value: Value) -> Vec { let data = value.get("data").unwrap_or(&value); data.get("hosts") .or_else(|| data.get("instances")) + .or_else(|| data.get("list")) .or_else(|| data.get("pageItems")) .or_else(|| data.get("items")) .or_else(|| value.get("hosts")) .or_else(|| value.get("instances")) + .or_else(|| value.get("list")) .and_then(Value::as_array) .cloned() .unwrap_or_default() @@ -1405,6 +1517,46 @@ impl EmptyFallback for String { #[cfg(test)] mod tests { use super::*; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + async fn read_request_target(socket: &mut tokio::net::TcpStream) -> String { + let mut request = Vec::new(); + let mut buffer = [0u8; 1024]; + loop { + let read = socket.read(&mut buffer).await.unwrap(); + if read == 0 { + break; + } + request.extend_from_slice(&buffer[..read]); + if request.windows(4).any(|window| window == b"\r\n\r\n") { + break; + } + } + let request = String::from_utf8(request).unwrap(); + request.split_whitespace().nth(1).unwrap().to_string() + } + + async fn write_json_response(socket: &mut tokio::net::TcpStream, body: &str) { + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ); + socket.write_all(response.as_bytes()).await.unwrap(); + } + + fn test_admin_config(server_addr: String) -> NacosAdminConfig { + NacosAdminConfig { + server_addr: server_addr.clone(), + display_server_addr: server_addr, + namespace: String::new(), + context_path: String::new(), + auth: NacosAuthConfig::None, + tls_skip_verify: false, + page_size: 100, + connect_override: None, + } + } #[test] fn parses_config_list_shapes() { @@ -1711,6 +1863,86 @@ mod tests { assert_eq!(parsed[0].healthy, Some(true)); } + #[test] + fn parses_v1_catalog_instance_list_including_disabled_instances() { + let parsed = parse_instances(serde_json::json!({ + "list": [{ + "ip": "192.0.2.59", + "port": 3259, + "clusterName": "DEFAULT", + "healthy": false, + "enabled": false, + "ephemeral": false + }], + "count": 1 + })); + + assert_eq!(catalog_instance_count(&serde_json::json!({ "list": [], "count": 1 })), Some(1)); + assert_eq!(parsed.len(), 1); + assert_eq!(parsed[0].ip, "192.0.2.59"); + assert_eq!(parsed[0].cluster_name.as_deref(), Some("DEFAULT")); + assert_eq!(parsed[0].healthy, Some(false)); + assert_eq!(parsed[0].enabled, Some(false)); + } + + #[test] + fn parses_v1_catalog_service_clusters_and_requested_cluster_filter() { + let clusters = parse_catalog_cluster_names(&serde_json::json!({ + "service": { "name": "svc" }, + "clusters": [ + { "name": "DEFAULT" }, + { "clusterName": "GRAY" }, + { "name": "DEFAULT" } + ] + })); + + assert_eq!(clusters, vec!["DEFAULT", "GRAY"]); + assert_eq!(split_nacos_cluster_names(Some(" DEFAULT,GRAY, DEFAULT ,")), vec!["DEFAULT", "GRAY"]); + assert!(split_nacos_cluster_names(None).is_empty()); + } + + #[tokio::test] + async fn qualifies_group_in_v1_catalog_service_requests() { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + let (mut detail_socket, _) = listener.accept().await.unwrap(); + let detail_target = read_request_target(&mut detail_socket).await; + let detail_url = reqwest::Url::parse(&format!("http://localhost{detail_target}")).unwrap(); + let detail_params = detail_url.query_pairs().collect::>(); + assert_eq!(detail_url.path(), "/v1/ns/catalog/service"); + assert_eq!(detail_params.get("serviceName").map(|value| value.as_ref()), Some("GRAY_GROUP@@orders")); + assert!(!detail_params.contains_key("groupName")); + write_json_response(&mut detail_socket, r#"{"clusters":[{"name":"DEFAULT"}]}"#).await; + + let (mut instances_socket, _) = listener.accept().await.unwrap(); + let instances_target = read_request_target(&mut instances_socket).await; + let instances_url = reqwest::Url::parse(&format!("http://localhost{instances_target}")).unwrap(); + let instances_params = instances_url.query_pairs().collect::>(); + assert_eq!(instances_url.path(), "/v1/ns/catalog/instances"); + assert_eq!(instances_params.get("serviceName").map(|value| value.as_ref()), Some("GRAY_GROUP@@orders")); + assert!(!instances_params.contains_key("groupName")); + write_json_response(&mut instances_socket, r#"{"list":[],"count":0}"#).await; + }); + + let admin = NacosOpenApiAdmin::new(test_admin_config(format!("http://{address}"))).unwrap(); + let instances = admin + .list_v1_catalog_instances( + &NacosInstanceQuery { + namespace: Some("public".to_string()), + service_name: "orders".to_string(), + group_name: Some("GRAY_GROUP".to_string()), + clusters: None, + }, + "public", + ) + .await + .unwrap(); + + assert!(instances.is_empty()); + server.await.unwrap(); + } + #[test] fn parses_namespace_list_shape() { let parsed = parse_namespaces(serde_json::json!({