Skip to content

Commit a8e8b99

Browse files
committed
feat[proxier server]: support /pods interface for http server
1 parent 271f1f4 commit a8e8b99

5 files changed

Lines changed: 714 additions & 18 deletions

File tree

cmd/kubeocean-proxier/main.go

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@ import (
4444

4545
cloudv1beta1 "github.com/gocrane/kubeocean/api/v1beta1"
4646
"github.com/gocrane/kubeocean/pkg/proxier"
47+
"github.com/gocrane/kubeocean/pkg/proxier/podcache"
4748
"github.com/gocrane/kubeocean/pkg/version"
4849
)
4950

@@ -112,7 +113,7 @@ func main() {
112113
}
113114

114115
// Setup proxier services
115-
kubeletProxy, httpServer, vnodeProxierAgent, err := setupProxierServices(ctx, config, tlsConfig, virtualClient, physicalClient, physicalConfig, clusterBinding, tokenManager)
116+
kubeletProxy, httpServer, vnodeProxierAgent, err := setupProxierServices(ctx, config, tlsConfig, virtualClient, physicalClient, physicalConfig, clusterBinding, tokenManager, virtualManager)
116117
if err != nil {
117118
setupLog.Error(err, "failed to setup proxier services")
118119
os.Exit(1)
@@ -550,7 +551,7 @@ func setupTLSConfiguration(ctx context.Context, config *ProxierConfig, clusterBi
550551
}
551552

552553
// setupProxierServices sets up kubelet proxy, HTTP server, and optionally VNode proxier agent
553-
func setupProxierServices(ctx context.Context, config *ProxierConfig, tlsConfig *TLSConfiguration, virtualClient client.Client, physicalClient kubernetes.Interface, physicalConfig *rest.Config, clusterBinding *cloudv1beta1.ClusterBinding, tokenManager *proxier.TokenManager) (proxier.KubeletProxy, proxier.HTTPServer, *proxier.VNodeProxierAgent, error) {
554+
func setupProxierServices(ctx context.Context, config *ProxierConfig, tlsConfig *TLSConfiguration, virtualClient client.Client, physicalClient kubernetes.Interface, physicalConfig *rest.Config, clusterBinding *cloudv1beta1.ClusterBinding, tokenManager *proxier.TokenManager, virtualManager ctrl.Manager) (proxier.KubeletProxy, proxier.HTTPServer, *proxier.VNodeProxierAgent, error) {
554555
virtualClientset, err := kubernetes.NewForConfig(ctrl.GetConfigOrDie())
555556
if err != nil {
556557
return nil, nil, nil, fmt.Errorf("unable to create virtual cluster clientset: %w", err)
@@ -590,6 +591,13 @@ func setupProxierServices(ctx context.Context, config *ProxierConfig, tlsConfig
590591
ctrl.Log.WithName("http-server"),
591592
)
592593

594+
// Create pod cache manager using virtualManager
595+
podCacheManager := podcache.NewPodCacheManager(virtualManager)
596+
if err := podCacheManager.Setup(ctx); err != nil {
597+
return nil, nil, nil, fmt.Errorf("failed to setup pod cache manager: %w", err)
598+
}
599+
setupLog.Info("Pod cache manager setup completed")
600+
593601
// Create VNode proxier agent if metrics enabled
594602
var vnodeProxierAgent *proxier.VNodeProxierAgent
595603
if config.MetricsEnabled {
@@ -614,6 +622,7 @@ func setupProxierServices(ctx context.Context, config *ProxierConfig, tlsConfig
614622
virtualClientset,
615623
clusterBinding.Spec.ClusterID,
616624
ctrl.Log.WithName("vnode-proxier-agent"),
625+
podCacheManager,
617626
)
618627

619628
if err := vnodeProxierAgent.Start(ctx); err != nil {

pkg/proxier/metrics_collector.go

Lines changed: 79 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -32,11 +32,15 @@ import (
3232
"github.com/gorilla/mux"
3333
corev1 "k8s.io/api/core/v1"
3434
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
35+
"k8s.io/apimachinery/pkg/runtime"
36+
"k8s.io/apimachinery/pkg/runtime/schema"
3537
"k8s.io/apimachinery/pkg/types"
3638
"k8s.io/apiserver/pkg/util/flushwriter"
3739
"k8s.io/client-go/kubernetes"
40+
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
3841
clientremotecommand "k8s.io/client-go/tools/remotecommand"
3942

43+
"github.com/gocrane/kubeocean/pkg/proxier/podcache"
4044
localremotecommand "github.com/gocrane/kubeocean/pkg/proxier/remotecommand"
4145
)
4246

@@ -60,12 +64,13 @@ type NodeEventHandler interface {
6064

6165
// NodeInfo node information (consistent with NodeController)
6266
type NodeInfo struct {
67+
NodeName string
6368
InternalIP string
6469
ProxierPort string
6570
}
6671

6772
func (n NodeInfo) String() string {
68-
return fmt.Sprintf("%s %s", n.InternalIP, n.ProxierPort)
73+
return fmt.Sprintf("%s %s %s", n.NodeName, n.InternalIP, n.ProxierPort)
6974
}
7075

7176
// ServerEntry HTTP server entry
@@ -120,6 +125,9 @@ type VNodeProxierAgent struct {
120125
// Cluster identification for VNode name generation
121126
clusterID string // ClusterBinding.Spec.ClusterID for generating VNode names
122127

128+
// Pod cache manager
129+
podCacheManager podcache.Manager
130+
123131
// Node state management
124132
nodeStates map[string]NodeInfo // key: nodeName, value: NodeInfo
125133
httpServers map[string]*ServerEntry // key: port
@@ -135,26 +143,27 @@ type VNodeProxierAgent struct {
135143
}
136144

137145
// NewVNodeProxierAgent creates a new VNode proxier agent
138-
func NewVNodeProxierAgent(config *MetricsConfig, tokenManager *TokenManager, kubeletProxy KubeletProxy, kubeClient kubernetes.Interface, clusterID string, log logr.Logger) *VNodeProxierAgent {
146+
func NewVNodeProxierAgent(config *MetricsConfig, tokenManager *TokenManager, kubeletProxy KubeletProxy, kubeClient kubernetes.Interface, clusterID string, log logr.Logger, podCacheManager podcache.Manager) *VNodeProxierAgent {
139147
// Initialize VictoriaMetrics unmarshal workers (referring to vnode_metrics)
140148
log.Info("Starting VictoriaMetrics unmarshal workers")
141149
common.StartUnmarshalWorkers()
142150

143151
mc := &VNodeProxierAgent{
144-
config: config,
145-
tokenManager: tokenManager,
146-
kubeletClient: NewKubeletClient(log.WithName("kubelet-client"), tokenManager),
147-
metricsParser: NewMetricsParser(),
148-
kubeletProxy: kubeletProxy,
149-
log: log,
150-
kubeClient: kubeClient,
151-
clusterID: clusterID,
152-
nodeStates: make(map[string]NodeInfo),
153-
httpServers: make(map[string]*ServerEntry),
154-
metricsCache: make(map[string][]byte),
155-
summaryCache: make(map[string]*Summary),
156-
lastUpdate: make(map[string]time.Time),
157-
stopChan: make(chan struct{}),
152+
config: config,
153+
tokenManager: tokenManager,
154+
kubeletClient: NewKubeletClient(log.WithName("kubelet-client"), tokenManager),
155+
metricsParser: NewMetricsParser(),
156+
kubeletProxy: kubeletProxy,
157+
log: log,
158+
kubeClient: kubeClient,
159+
clusterID: clusterID,
160+
podCacheManager: podCacheManager,
161+
nodeStates: make(map[string]NodeInfo),
162+
httpServers: make(map[string]*ServerEntry),
163+
metricsCache: make(map[string][]byte),
164+
summaryCache: make(map[string]*Summary),
165+
lastUpdate: make(map[string]time.Time),
166+
stopChan: make(chan struct{}),
158167
}
159168

160169
// Load TLS configuration if provided
@@ -646,6 +655,9 @@ func (va *VNodeProxierAgent) setupRoutes(port string, nodeInfo NodeInfo) *mux.Ro
646655
// Container exec endpoint (same as main server)
647656
router.HandleFunc("/exec/{namespace}/{pod}/{container}", va.handleContainerExec).Methods("POST", "GET")
648657

658+
// Pods endpoint
659+
router.HandleFunc("/pods", va.handlePods(nodeInfo.NodeName)).Methods("GET")
660+
649661
// Health check endpoint
650662
router.HandleFunc("/healthz", va.handleHealthz).Methods("GET")
651663

@@ -703,6 +715,57 @@ func (va *VNodeProxierAgent) handleSummary(port string) http.HandlerFunc {
703715
}
704716
}
705717

718+
// handlePods handles /pods endpoint to return pods by node name
719+
func (va *VNodeProxierAgent) handlePods(nodeName string) http.HandlerFunc {
720+
return func(w http.ResponseWriter, r *http.Request) {
721+
if va.podCacheManager == nil {
722+
va.log.Error(nil, "Pod cache manager is not available")
723+
http.Error(w, "Pod cache manager is not available", http.StatusServiceUnavailable)
724+
return
725+
}
726+
727+
va.log.V(1).Info("Processing pods request",
728+
"nodeName", nodeName,
729+
"remoteAddr", r.RemoteAddr,
730+
"userAgent", r.UserAgent(),
731+
)
732+
733+
// Get pods by node name
734+
ctx := r.Context()
735+
podList, err := va.podCacheManager.GetPodsByNode(ctx, nodeName)
736+
if err != nil {
737+
va.log.Error(err, "Failed to get pods by node name", "nodeName", nodeName)
738+
http.Error(w, fmt.Sprintf("Failed to get pods: %v", err), http.StatusInternalServerError)
739+
return
740+
}
741+
742+
// Encode PodList to JSON using k8s client-go codec (following kubelet's encodePods pattern)
743+
// Use LegacyCodec like kubelet does for v1 API version
744+
codec := clientgoscheme.Codecs.LegacyCodec(schema.GroupVersion{Group: corev1.GroupName, Version: "v1"})
745+
jsonBytes, err := runtime.Encode(codec, podList)
746+
if err != nil {
747+
va.log.Error(err, "Failed to encode PodList to JSON", "nodeName", nodeName)
748+
http.Error(w, fmt.Sprintf("Failed to encode pods: %v", err), http.StatusInternalServerError)
749+
return
750+
}
751+
752+
// Set response headers
753+
w.Header().Set("Content-Type", "application/json")
754+
w.WriteHeader(http.StatusOK)
755+
756+
// Write JSON response
757+
if _, err := w.Write(jsonBytes); err != nil {
758+
va.log.Error(err, "Failed to write response", "nodeName", nodeName)
759+
return
760+
}
761+
762+
va.log.V(2).Info("Successfully served pods data",
763+
"nodeName", nodeName,
764+
"podCount", len(podList.Items),
765+
"dataSize", len(jsonBytes))
766+
}
767+
}
768+
706769
// handleContainerLogs handles container logs requests (same as main server)
707770
func (va *VNodeProxierAgent) handleContainerLogs(w http.ResponseWriter, r *http.Request) {
708771
if va.kubeletProxy == nil {

pkg/proxier/node_controller.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,7 @@ func (r *NodeController) extractNodeInfo(node *corev1.Node) (*NodeInfo, error) {
139139
}
140140

141141
return &NodeInfo{
142+
NodeName: node.Name,
142143
InternalIP: internalIP,
143144
ProxierPort: proxierPort,
144145
}, nil
Lines changed: 146 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,146 @@
1+
// Copyright 2025 The Kubeocean Authors
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package podcache
16+
17+
import (
18+
"context"
19+
"fmt"
20+
"sync"
21+
22+
"github.com/go-logr/logr"
23+
corev1 "k8s.io/api/core/v1"
24+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
25+
"k8s.io/client-go/tools/cache"
26+
ctrl "sigs.k8s.io/controller-runtime"
27+
)
28+
29+
const (
30+
// IndexNameNodeName is the index name for pods by node name
31+
IndexNameNodeName = "spec.nodeName"
32+
)
33+
34+
// Manager defines the interface for pod cache management
35+
type Manager interface {
36+
// GetPodsByNode returns all pods running on the specified node
37+
GetPodsByNode(ctx context.Context, nodeName string) (*corev1.PodList, error)
38+
}
39+
40+
// PodCacheManager manages a cache of pods using controller-runtime manager
41+
// PodCacheManager implements the Manager interface
42+
type PodCacheManager struct {
43+
manager ctrl.Manager
44+
indexer cache.Indexer
45+
logger logr.Logger
46+
mu sync.RWMutex
47+
}
48+
49+
// NewPodCacheManager creates a new PodCacheManager
50+
func NewPodCacheManager(manager ctrl.Manager) *PodCacheManager {
51+
return &PodCacheManager{
52+
manager: manager,
53+
logger: ctrl.Log.WithName("pod-cache"),
54+
}
55+
}
56+
57+
// Setup sets up the index for pods by node name
58+
// Note: This should be called before starting the manager
59+
func (pcm *PodCacheManager) Setup(ctx context.Context) error {
60+
pcm.mu.Lock()
61+
defer pcm.mu.Unlock()
62+
63+
// Get pod informer from controller-runtime cache
64+
podInformer, err := pcm.manager.GetCache().GetInformer(ctx, &corev1.Pod{})
65+
if err != nil {
66+
return fmt.Errorf("failed to get pod informer: %w", err)
67+
}
68+
69+
// Type assert to get the underlying SharedIndexInformer
70+
// controller-runtime's cache wraps client-go's SharedIndexInformer
71+
sharedIndexInformer, ok := podInformer.(cache.SharedIndexInformer)
72+
if !ok {
73+
return fmt.Errorf("failed to get indexer from informer: informer is not a SharedIndexInformer")
74+
}
75+
76+
// Add indexer using client-go's AddIndexers method
77+
// This must be called before the informer starts
78+
indexers := cache.Indexers{
79+
IndexNameNodeName: func(obj interface{}) ([]string, error) {
80+
pod, ok := obj.(*corev1.Pod)
81+
if !ok {
82+
return nil, fmt.Errorf("object is not a Pod")
83+
}
84+
if pod.Spec.NodeName == "" {
85+
return nil, nil
86+
}
87+
return []string{pod.Spec.NodeName}, nil
88+
},
89+
}
90+
91+
// Add the indexers to the informer's indexer
92+
// Note: AddIndexers will return an error if the index already exists
93+
// This can happen if the informer has already started or if the index was added elsewhere
94+
if err := sharedIndexInformer.GetIndexer().AddIndexers(indexers); err != nil {
95+
// Check if it's an "indexer conflict" error, which means the index already exists
96+
// In that case, we can safely proceed as the index is already set up
97+
if err.Error() != "indexer conflict: "+fmt.Sprintf("map[%s:{}]", IndexNameNodeName) {
98+
return fmt.Errorf("failed to add indexers: %w", err)
99+
}
100+
pcm.logger.V(1).Info("Index already exists, reusing existing index", "indexName", IndexNameNodeName)
101+
}
102+
103+
// Store the indexer for later use
104+
pcm.indexer = sharedIndexInformer.GetIndexer()
105+
106+
pcm.logger.Info("Pod cache manager index setup completed")
107+
108+
return nil
109+
}
110+
111+
// GetPodsByNode returns all pods running on the specified node
112+
func (pcm *PodCacheManager) GetPodsByNode(ctx context.Context, nodeName string) (*corev1.PodList, error) {
113+
pcm.mu.RLock()
114+
defer pcm.mu.RUnlock()
115+
116+
if pcm.indexer == nil {
117+
return nil, fmt.Errorf("pod indexer is not initialized")
118+
}
119+
120+
// Directly use indexer.ByIndex for optimal performance
121+
// This bypasses the cache.List overhead and directly uses the index
122+
pods, err := pcm.indexer.ByIndex(IndexNameNodeName, nodeName)
123+
if err != nil {
124+
return nil, fmt.Errorf("failed to get pods by node name %s: %w", nodeName, err)
125+
}
126+
127+
// Convert to PodList
128+
podList := &corev1.PodList{
129+
TypeMeta: metav1.TypeMeta{
130+
Kind: "PodList",
131+
APIVersion: "v1",
132+
},
133+
Items: make([]corev1.Pod, 0, len(pods)),
134+
}
135+
136+
for _, obj := range pods {
137+
pod, ok := obj.(*corev1.Pod)
138+
if !ok {
139+
pcm.logger.V(1).Info("Object in indexer is not a pod, skipping", "nodeName", nodeName)
140+
continue
141+
}
142+
podList.Items = append(podList.Items, *pod)
143+
}
144+
145+
return podList, nil
146+
}

0 commit comments

Comments
 (0)