1use anyhow::Result;
6use async_trait::async_trait;
7
8#[derive(Clone, Debug, PartialEq, Eq, Hash)]
10pub struct ServiceEndpoint {
11 pub uri: String,
12}
13
14impl ServiceEndpoint {
15 pub fn new(uri: impl Into<String>) -> Self {
16 Self { uri: uri.into() }
17 }
18
19 #[must_use]
20 pub fn http(host: &str, port: u16) -> Self {
21 Self {
22 uri: format!("{}://{}:{}", "http", host, port),
23 }
24 }
25
26 #[must_use]
27 pub fn https(host: &str, port: u16) -> Self {
28 Self {
29 uri: format!("https://{host}:{port}"),
30 }
31 }
32
33 pub fn uds(path: impl AsRef<std::path::Path>) -> Self {
34 Self {
35 uri: format!("unix://{}", path.as_ref().display()),
36 }
37 }
38}
39
40#[derive(Debug, Clone)]
42pub struct ServiceInstanceInfo {
43 pub gear: String,
45 pub instance_id: String,
47 pub endpoint: ServiceEndpoint,
49 pub version: Option<String>,
51 pub rest_endpoint: Option<ServiceEndpoint>,
54 pub openapi_spec: Option<String>,
56 pub openapi_spec_hash: Option<String>,
58 pub grpc_services: Vec<(String, ServiceEndpoint)>,
64}
65
66#[derive(Debug, Clone)]
68pub struct RegisterInstanceInfo {
69 pub gear: String,
71 pub instance_id: String,
73 pub grpc_services: Vec<(String, ServiceEndpoint)>,
75 pub version: Option<String>,
77 pub rest_endpoint: Option<ServiceEndpoint>,
79 pub openapi_spec: Option<String>,
81}
82
83#[derive(Clone, Debug, PartialEq, Eq, Hash)]
85pub struct GrpcServiceInfo {
86 pub service_name: String,
88 pub endpoint: ServiceEndpoint,
90}
91
92impl GrpcServiceInfo {
93 pub fn new(service_name: impl Into<String>, endpoint: ServiceEndpoint) -> Self {
94 Self {
95 service_name: service_name.into(),
96 endpoint,
97 }
98 }
99}
100
101#[derive(Debug, Clone, PartialEq, Eq)]
107pub struct DirectoryNotFound {
108 pub resource: String,
110}
111
112impl DirectoryNotFound {
113 pub fn new(resource: impl Into<String>) -> Self {
114 Self {
115 resource: resource.into(),
116 }
117 }
118}
119
120impl std::fmt::Display for DirectoryNotFound {
121 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
122 write!(f, "directory: not found: {}", self.resource)
123 }
124}
125
126impl std::error::Error for DirectoryNotFound {}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
133pub struct DirectoryInvalidArgument {
134 pub message: String,
136}
137
138impl DirectoryInvalidArgument {
139 pub fn new(message: impl Into<String>) -> Self {
140 Self {
141 message: message.into(),
142 }
143 }
144}
145
146impl std::fmt::Display for DirectoryInvalidArgument {
147 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
148 write!(f, "directory: invalid argument: {}", self.message)
149 }
150}
151
152impl std::error::Error for DirectoryInvalidArgument {}
153
154#[async_trait]
161pub trait DirectoryClient: Send + Sync {
162 async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint>;
164
165 async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint>;
170
171 async fn get_openapi_spec(&self, gear_name: &str) -> Result<String>;
173
174 async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>>;
176
177 async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>>;
186
187 async fn register_instance(&self, info: RegisterInstanceInfo) -> Result<()>;
189
190 async fn deregister_instance(&self, gear: &str, instance_id: &str) -> Result<()>;
192
193 async fn send_heartbeat(&self, gear: &str, instance_id: &str) -> Result<()>;
195}
196
197#[cfg(test)]
198#[cfg_attr(coverage_nightly, coverage(off))]
199mod tests {
200 use super::*;
201
202 #[test]
203 fn test_service_endpoint_creation() {
204 let http_ep = ServiceEndpoint::http("localhost", 8080);
205 assert_eq!(http_ep.uri, concat!("http", "://localhost:8080"));
206
207 let https_endpoint = ServiceEndpoint::https("localhost", 8443);
208 assert_eq!(https_endpoint.uri, "https://localhost:8443");
209
210 let uds_ep = ServiceEndpoint::uds("/tmp/socket.sock");
211 assert!(uds_ep.uri.starts_with("unix://"));
212 assert!(uds_ep.uri.contains("socket.sock"));
213
214 let custom_ep = ServiceEndpoint::new(concat!("http", "://example.com"));
215 assert_eq!(custom_ep.uri, concat!("http", "://example.com"));
216 }
217
218 #[test]
219 fn test_register_instance_info() {
220 let info = RegisterInstanceInfo {
221 gear: "test_gear".to_owned(),
222 instance_id: "instance1".to_owned(),
223 grpc_services: vec![(
224 "test.Service".to_owned(),
225 ServiceEndpoint::http("127.0.0.1", 8001),
226 )],
227 version: Some("1.0.0".to_owned()),
228 rest_endpoint: None,
229 openapi_spec: None,
230 };
231
232 assert_eq!(info.gear, "test_gear");
233 assert_eq!(info.instance_id, "instance1");
234 assert_eq!(info.grpc_services.len(), 1);
235 assert!(info.rest_endpoint.is_none());
236 assert!(info.openapi_spec.is_none());
237 }
238
239 #[test]
240 fn test_register_instance_info_with_rest() {
241 let info = RegisterInstanceInfo {
242 gear: "billing".to_owned(),
243 instance_id: "instance1".to_owned(),
244 grpc_services: vec![],
245 version: Some("2.0.0".to_owned()),
246 rest_endpoint: Some(ServiceEndpoint::http("billing", 8080)),
247 openapi_spec: Some("{\"openapi\":\"3.1.0\"}".to_owned()),
248 };
249
250 assert_eq!(info.gear, "billing");
251 assert_eq!(
252 info.rest_endpoint.as_ref().unwrap().uri,
253 concat!("http", "://billing:8080")
254 );
255 assert!(info.openapi_spec.is_some());
256 }
257}