Skip to main content

cf_system_sdk_directory/
api.rs

1//! Directory API - contract for service discovery and instance resolution
2//!
3//! This gear defines the core traits and types for the directory service API.
4
5use anyhow::Result;
6use async_trait::async_trait;
7
8/// Represents an endpoint where a service can be reached
9#[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/// Information about a service instance
41#[derive(Debug, Clone)]
42pub struct ServiceInstanceInfo {
43    /// Gear name this instance belongs to
44    pub gear: String,
45    /// Unique instance identifier
46    pub instance_id: String,
47    /// Primary endpoint for the instance
48    pub endpoint: ServiceEndpoint,
49    /// Optional version string
50    pub version: Option<String>,
51    /// Optional REST endpoint (HTTP base URL) for this instance.
52    /// Not all gears expose a REST API.
53    pub rest_endpoint: Option<ServiceEndpoint>,
54    /// Optional `OpenAPI` spec (JSON) this instance published, if any.
55    pub openapi_spec: Option<String>,
56    /// Stable content token for the published `OpenAPI` spec, if any.
57    pub openapi_spec_hash: Option<String>,
58    /// Map of gRPC service name to endpoint published by this instance.
59    ///
60    /// Carried back by `list_instances` so a subsequent `register_instance`
61    /// (which replaces the entry wholesale) can augment — rather than clobber —
62    /// the previously-registered gRPC services when adding a REST endpoint.
63    pub grpc_services: Vec<(String, ServiceEndpoint)>,
64}
65
66/// Information for registering a new gear instance
67#[derive(Debug, Clone)]
68pub struct RegisterInstanceInfo {
69    /// Gear name
70    pub gear: String,
71    /// Unique instance identifier
72    pub instance_id: String,
73    /// Map of gRPC service name to endpoint
74    pub grpc_services: Vec<(String, ServiceEndpoint)>,
75    /// Optional version string
76    pub version: Option<String>,
77    /// Optional REST endpoint (HTTP base URL) exposed by the gear.
78    pub rest_endpoint: Option<ServiceEndpoint>,
79    /// Optional `OpenAPI` spec (JSON) published by the gear.
80    pub openapi_spec: Option<String>,
81}
82
83/// A resolved gRPC service and the endpoint it is reachable at.
84#[derive(Clone, Debug, PartialEq, Eq, Hash)]
85pub struct GrpcServiceInfo {
86    /// Fully-qualified gRPC service name (e.g. `payment.v1.PaymentApi`).
87    pub service_name: String,
88    /// Endpoint the service is reachable at.
89    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/// Sentinel error wrapped via `anyhow::Error` to signal "the requested gear or
102/// service is not registered (or has no live instance)" through the
103/// [`DirectoryClient`] trait. Consumers downcast to this type to distinguish a
104/// not-ready provider (eventual readiness) from a directory-backend failure —
105/// see `toolkit::discovery::DirectoryEndpointResolver`.
106#[derive(Debug, Clone, PartialEq, Eq)]
107pub struct DirectoryNotFound {
108    /// What was being looked up — e.g. `"gear foo"` or `"service foo.Bar"`.
109    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/// Sentinel error wrapped via `anyhow::Error` to signal "client-supplied
129/// argument is malformed" (e.g. invalid UUID) through the [`DirectoryClient`]
130/// trait. Allows the gRPC server boundary to return `Status::invalid_argument`
131/// instead of mislabeling a client bug as an internal failure.
132#[derive(Debug, Clone, PartialEq, Eq)]
133pub struct DirectoryInvalidArgument {
134    /// Human-readable description of what was invalid.
135    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/// Directory API trait for service discovery and instance management
155///
156/// This trait defines the contract for interacting with the gear directory.
157/// It can be implemented by:
158/// - A local implementation that delegates to `GearManager`
159/// - A gRPC client for out-of-process gears
160#[async_trait]
161pub trait DirectoryClient: Send + Sync {
162    /// Resolve a gRPC service by its logical name to an endpoint
163    async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint>;
164
165    /// Resolve a REST endpoint (HTTP base URL) for a gear by its name.
166    ///
167    /// Returns the base URL (e.g. `http://billing:8080`) that callers use to
168    /// make REST requests to the resolved gear.
169    async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint>;
170
171    /// Retrieve the `OpenAPI` spec (JSON) published by a gear.
172    async fn get_openapi_spec(&self, gear_name: &str) -> Result<String>;
173
174    /// List all service instances for a given gear
175    async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>>;
176
177    /// List every service instance across all registered gears.
178    ///
179    /// Used by the edge gateway to discover which gears (and their REST
180    /// endpoints) to reverse-proxy. This is a lightweight discovery snapshot:
181    /// the returned instances do **not** carry `openapi_spec` — even when the
182    /// backing store holds a stored specification. The edge fetches a gear's
183    /// document once, on first discovery, via
184    /// [`get_openapi_spec`](Self::get_openapi_spec).
185    async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>>;
186
187    /// Register a new gear instance with the directory
188    async fn register_instance(&self, info: RegisterInstanceInfo) -> Result<()>;
189
190    /// Deregister a gear instance (for graceful shutdown)
191    async fn deregister_instance(&self, gear: &str, instance_id: &str) -> Result<()>;
192
193    /// Send a heartbeat for a gear instance to indicate it's still alive
194    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}