Skip to content

Commit 707beb4

Browse files
committed
feat: feature-gated Kubernetes service discovery
Move Kubernetes discovery from the example into the library behind a `kubernetes` feature flag. When disabled (default), zero k8s dependencies are pulled in. Usage: groupcache = { version = "0.3", features = ["kubernetes"] } let discovery = KubernetesDiscovery::builder() .client(client) .label_selector("app=my-service") .port(8080) .build()?; Groupcache::builder(me, loader) .service_discovery(discovery) .build(); Improvements over the example's implementation: - Configurable label selector (was hardcoded) - Configurable port (was hardcoded to 3000) - Configurable namespace (defaults to pod's own) - Configurable poll interval - Builder pattern with validation Also bumps kube from 2.0.1 to 3.x (absorbs dependabot PR Petroniuss#52).
1 parent 290a918 commit 707beb4

7 files changed

Lines changed: 178 additions & 66 deletions

File tree

examples/kubernetes-service-discovery/Cargo.toml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ publish = false
77
[dependencies]
88
# This pulls version from main branch so that docker build works (docker was confused by paths)
99
#groupcache = { git = "https://github.com/Petroniuss/groupcache.git" }
10-
groupcache = { path = "../../groupcache" }
10+
groupcache = { path = "../../groupcache", features = ["kubernetes"] }
1111
tonic = "0.14.2"
1212
axum = "0.8.0"
1313

@@ -25,5 +25,5 @@ axum-prometheus = "0.9.0"
2525
anyhow = "1"
2626
async-trait = "0.1"
2727

28-
kube = { version = "2.0.1", features = ["runtime", "derive"] }
29-
k8s-openapi = { version = "0.26.0", features = ["latest"] }
28+
# kube and k8s-openapi come through groupcache's "kubernetes" feature
29+
kube = { version = "3", features = ["runtime"] }
Lines changed: 2 additions & 61 deletions
Original file line numberDiff line numberDiff line change
@@ -1,61 +1,2 @@
1-
use async_trait::async_trait;
2-
use groupcache::{GroupcachePeer, ServiceDiscovery};
3-
use k8s_openapi::api::core::v1::Pod;
4-
use kube::api::ListParams;
5-
use kube::{Api, Client};
6-
use std::collections::HashSet;
7-
use std::error::Error;
8-
use std::net::SocketAddr;
9-
10-
pub struct Kubernetes {
11-
api: Api<Pod>,
12-
}
13-
14-
pub struct KubernetesBuilder {
15-
client: Option<Client>,
16-
}
17-
18-
impl KubernetesBuilder {
19-
pub fn build(self) -> Kubernetes {
20-
Kubernetes {
21-
api: Api::default_namespaced(self.client.unwrap()),
22-
}
23-
}
24-
pub fn client(mut self, client: Client) -> Self {
25-
self.client = Some(client);
26-
self
27-
}
28-
}
29-
30-
impl Kubernetes {
31-
pub fn builder() -> KubernetesBuilder {
32-
KubernetesBuilder { client: None }
33-
}
34-
}
35-
36-
#[async_trait]
37-
impl ServiceDiscovery for Kubernetes {
38-
async fn pull_instances(
39-
&self,
40-
) -> Result<HashSet<GroupcachePeer>, Box<dyn Error + Send + Sync + 'static>> {
41-
let pods_with_label_query = ListParams::default().labels("app=groupcache-powered-backend");
42-
Ok(self
43-
.api
44-
.list(&pods_with_label_query)
45-
.await
46-
.unwrap()
47-
.into_iter()
48-
.filter_map(|pod| {
49-
let status = pod.status?;
50-
let pod_ip = status.pod_ip?;
51-
52-
let Ok(ip) = pod_ip.parse() else {
53-
return None;
54-
};
55-
56-
let addr = SocketAddr::new(ip, 3000);
57-
Some(GroupcachePeer::from_socket(addr))
58-
})
59-
.collect())
60-
}
61-
}
1+
// Re-export from the library's built-in kubernetes discovery.
2+
pub use groupcache::discovery::kubernetes::KubernetesDiscovery;

examples/kubernetes-service-discovery/src/main.rs

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ mod cache;
33
mod k8s;
44

55
use crate::cache::CachedValue;
6-
use crate::k8s::Kubernetes;
6+
use crate::k8s::KubernetesDiscovery;
77
use anyhow::Context;
88
use anyhow::Result;
99
use axum::extract::{Path, State};
@@ -59,8 +59,15 @@ async fn main() -> Result<()> {
5959
let client = Client::try_default().await?;
6060

6161
// Configuring groupcache to use Kubernetes API server for peer auto-discovery.
62+
let discovery = KubernetesDiscovery::builder()
63+
.client(client)
64+
.label_selector("app=groupcache-powered-backend")
65+
.port(pod_port.parse()?)
66+
.build()
67+
.map_err(|e| anyhow::anyhow!("{}", e))?;
68+
6269
let groupcache = Groupcache::builder(addr.into(), loader)
63-
.service_discovery(Kubernetes::builder().client(client).build())
70+
.service_discovery(discovery)
6471
.build();
6572

6673
// Example axum app with endpoint to retrieve value from groupcache.

groupcache/Cargo.toml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ repository = "https://github.com/Petroniuss/groupcache"
1717

1818
[features]
1919
bincode = ["dep:bincode"]
20+
kubernetes = ["dep:kube", "dep:k8s-openapi"]
2021

2122
[dependencies]
2223
groupcache-pb = { path = "../groupcache-pb", version = "0.3.0" }
@@ -35,6 +36,8 @@ singleflight-async = "0.2.0"
3536
moka = { version = "0.12.1", features = ["future"] }
3637
log = "0.4.20"
3738
metrics = "0.24.0"
39+
kube = { version = "3", features = ["runtime", "derive"], optional = true }
40+
k8s-openapi = { version = "0.27", features = ["latest"], optional = true }
3841

3942
[dev-dependencies]
4043
cargo-husky = { workspace = true }
Lines changed: 153 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,153 @@
1+
//! Kubernetes service discovery for groupcache.
2+
//!
3+
//! Discovers peers by listing pods matching a label selector via the
4+
//! Kubernetes API server. Requires the `kubernetes` feature.
5+
//!
6+
//! # Example
7+
//!
8+
//! ```no_run
9+
//! use groupcache::discovery::kubernetes::KubernetesDiscovery;
10+
//!
11+
//! # async fn example() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
12+
//! let client = kube::Client::try_default().await?;
13+
//! let discovery = KubernetesDiscovery::builder()
14+
//! .client(client)
15+
//! .label_selector("app=my-cache")
16+
//! .port(8080)
17+
//! .build()?;
18+
//! # Ok(())
19+
//! # }
20+
//! ```
21+
22+
use crate::groupcache::GroupcachePeer;
23+
use crate::service_discovery::ServiceDiscovery;
24+
use async_trait::async_trait;
25+
use k8s_openapi::api::core::v1::Pod;
26+
use kube::api::ListParams;
27+
use kube::{Api, Client};
28+
use std::collections::HashSet;
29+
use std::error::Error;
30+
use std::net::SocketAddr;
31+
use std::time::Duration;
32+
33+
/// Kubernetes-based service discovery for groupcache peers.
34+
///
35+
/// Lists pods matching a label selector and extracts pod IPs to build
36+
/// the peer set. Uses the Kubernetes API server (via the `kube` crate).
37+
pub struct KubernetesDiscovery {
38+
api: Api<Pod>,
39+
label_selector: String,
40+
port: u16,
41+
poll_interval: Duration,
42+
}
43+
44+
/// Builder for [`KubernetesDiscovery`].
45+
pub struct KubernetesDiscoveryBuilder {
46+
client: Option<Client>,
47+
namespace: Option<String>,
48+
label_selector: Option<String>,
49+
port: Option<u16>,
50+
poll_interval: Option<Duration>,
51+
}
52+
53+
impl KubernetesDiscoveryBuilder {
54+
/// Set the Kubernetes client. Required.
55+
pub fn client(mut self, client: Client) -> Self {
56+
self.client = Some(client);
57+
self
58+
}
59+
60+
/// Set a specific namespace. Defaults to the pod's own namespace.
61+
pub fn namespace(mut self, namespace: impl Into<String>) -> Self {
62+
self.namespace = Some(namespace.into());
63+
self
64+
}
65+
66+
/// Set the label selector for discovering peer pods. Required.
67+
///
68+
/// Example: `"app=my-groupcache-service"`
69+
pub fn label_selector(mut self, selector: impl Into<String>) -> Self {
70+
self.label_selector = Some(selector.into());
71+
self
72+
}
73+
74+
/// Set the port that groupcache peers listen on. Required.
75+
pub fn port(mut self, port: u16) -> Self {
76+
self.port = Some(port);
77+
self
78+
}
79+
80+
/// Set the polling interval for service discovery. Default: 10 seconds.
81+
pub fn poll_interval(mut self, interval: Duration) -> Self {
82+
self.poll_interval = Some(interval);
83+
self
84+
}
85+
86+
/// Build the [`KubernetesDiscovery`] instance.
87+
///
88+
/// # Errors
89+
///
90+
/// Returns an error if `client`, `label_selector`, or `port` were not set.
91+
pub fn build(self) -> Result<KubernetesDiscovery, Box<dyn Error + Send + Sync>> {
92+
let client = self
93+
.client
94+
.ok_or("KubernetesDiscovery requires a kube::Client")?;
95+
let label_selector = self
96+
.label_selector
97+
.ok_or("KubernetesDiscovery requires a label_selector")?;
98+
let port = self
99+
.port
100+
.ok_or("KubernetesDiscovery requires a port")?;
101+
102+
let api = match self.namespace {
103+
Some(ns) => Api::namespaced(client, &ns),
104+
None => Api::default_namespaced(client),
105+
};
106+
107+
Ok(KubernetesDiscovery {
108+
api,
109+
label_selector,
110+
port,
111+
poll_interval: self.poll_interval.unwrap_or(Duration::from_secs(10)),
112+
})
113+
}
114+
}
115+
116+
impl KubernetesDiscovery {
117+
/// Create a new builder.
118+
pub fn builder() -> KubernetesDiscoveryBuilder {
119+
KubernetesDiscoveryBuilder {
120+
client: None,
121+
namespace: None,
122+
label_selector: None,
123+
port: None,
124+
poll_interval: None,
125+
}
126+
}
127+
}
128+
129+
#[async_trait]
130+
impl ServiceDiscovery for KubernetesDiscovery {
131+
async fn pull_instances(
132+
&self,
133+
) -> Result<HashSet<GroupcachePeer>, Box<dyn Error + Send + Sync + 'static>> {
134+
let params = ListParams::default().labels(&self.label_selector);
135+
let pods = self.api.list(&params).await?;
136+
137+
let peers = pods
138+
.into_iter()
139+
.filter_map(|pod| {
140+
let pod_ip = pod.status?.pod_ip?;
141+
let ip = pod_ip.parse().ok()?;
142+
let addr = SocketAddr::new(ip, self.port);
143+
Some(GroupcachePeer::from_socket(addr))
144+
})
145+
.collect();
146+
147+
Ok(peers)
148+
}
149+
150+
fn interval(&self) -> Duration {
151+
self.poll_interval
152+
}
153+
}

groupcache/src/discovery/mod.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
//! Built-in service discovery implementations.
2+
//!
3+
//! Enable via feature flags:
4+
//! - `kubernetes` — discover peers via Kubernetes API server
5+
6+
#[cfg(feature = "kubernetes")]
7+
pub mod kubernetes;

groupcache/src/lib.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
#![doc = include_str!("../readme.md")]
22

33
mod codec;
4+
pub mod discovery;
45
mod errors;
56
mod groupcache;
67
mod groupcache_builder;

0 commit comments

Comments
 (0)