fix: replace client-go informer with REST loader, drop k8s deps
Informer cache sync timed out due to transient 429 from API server. REST loader does one GET + retry, polls every 30s. Zero new deps. go.mod back to 1.25.4, CI back to golang:1.25-bookworm.
This commit is contained in:
+134
-121
@@ -1,28 +1,33 @@
|
||||
package serviceadapter
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
"k8s.io/client-go/dynamic"
|
||||
"k8s.io/client-go/dynamic/dynamicinformer"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
)
|
||||
|
||||
// InformerManager watches ServiceAdapter CRs and populates the registry.
|
||||
type InformerManager struct {
|
||||
registry *Registry
|
||||
informer cache.SharedIndexInformer
|
||||
stopChan chan struct{}
|
||||
factory dynamicinformer.DynamicSharedInformerFactory
|
||||
// Loader watches ServiceAdapter CRs via the Kubernetes REST API and populates the registry.
|
||||
// No client-go dependency — uses the in-cluster service account token directly.
|
||||
type Loader struct {
|
||||
registry *Registry
|
||||
namespace string
|
||||
client *http.Client
|
||||
token string
|
||||
baseURL string
|
||||
stopChan chan struct{}
|
||||
once sync.Once
|
||||
}
|
||||
|
||||
// NewInformerManager creates a new informer that watches ServiceAdapters in the given namespace.
|
||||
func NewInformerManager(dynamicClient dynamic.Interface, registry *Registry, namespace string) (*InformerManager, error) {
|
||||
// NewLoader creates a loader that reads ServiceAdapters from the k8s API.
|
||||
// Returns an error if not running inside a cluster (missing SA token/ca).
|
||||
func NewLoader(registry *Registry, namespace string) (*Loader, error) {
|
||||
if registry == nil {
|
||||
return nil, fmt.Errorf("registry cannot be nil")
|
||||
}
|
||||
@@ -30,134 +35,142 @@ func NewInformerManager(dynamicClient dynamic.Interface, registry *Registry, nam
|
||||
namespace = "api"
|
||||
}
|
||||
|
||||
gvr := schema.GroupVersionResource{
|
||||
Group: "gateway.riotpiao.com",
|
||||
Version: "v1",
|
||||
Resource: "serviceadapters",
|
||||
}
|
||||
|
||||
factory := dynamicinformer.NewFilteredDynamicSharedInformerFactory(
|
||||
dynamicClient, 30*time.Second, namespace, nil,
|
||||
)
|
||||
|
||||
informer := factory.ForResource(gvr).Informer()
|
||||
|
||||
_, err := informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
AddFunc: func(obj interface{}) {
|
||||
u := obj.(*unstructured.Unstructured)
|
||||
if a, err := parseServiceAdapter(u); err == nil {
|
||||
if registry.Add(a) == nil {
|
||||
log.Printf("serviceadapter added: %s", a.ServiceName)
|
||||
}
|
||||
} else {
|
||||
log.Printf("serviceadapter parse error on add: %v", err)
|
||||
}
|
||||
},
|
||||
UpdateFunc: func(_, newObj interface{}) {
|
||||
u := newObj.(*unstructured.Unstructured)
|
||||
if a, err := parseServiceAdapter(u); err == nil {
|
||||
registry.Update(a)
|
||||
}
|
||||
},
|
||||
DeleteFunc: func(obj interface{}) {
|
||||
u, ok := obj.(*unstructured.Unstructured)
|
||||
if !ok {
|
||||
tombstone, ok := obj.(cache.DeletedFinalStateUnknown)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
u, ok = tombstone.Obj.(*unstructured.Unstructured)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
}
|
||||
if name, ok := u.Object["metadata"].(map[string]interface{})["name"]; ok {
|
||||
registry.Delete(name.(string))
|
||||
log.Printf("serviceadapter deleted: %s", name)
|
||||
}
|
||||
},
|
||||
})
|
||||
tokenBytes, err := os.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/token")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("adding event handler: %w", err)
|
||||
return nil, fmt.Errorf("not in cluster: %w", err)
|
||||
}
|
||||
|
||||
return &InformerManager{
|
||||
registry: registry,
|
||||
informer: informer,
|
||||
caBytes, err := os.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/ca.crt")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("missing CA cert: %w", err)
|
||||
}
|
||||
|
||||
pool := x509.NewCertPool()
|
||||
pool.AppendCertsFromPEM(caBytes)
|
||||
|
||||
return &Loader{
|
||||
registry: registry,
|
||||
namespace: namespace,
|
||||
token: string(tokenBytes),
|
||||
baseURL: "https://kubernetes.default.svc",
|
||||
client: &http.Client{
|
||||
Timeout: 10 * time.Second,
|
||||
Transport: &http.Transport{
|
||||
TLSClientConfig: &tls.Config{RootCAs: pool},
|
||||
},
|
||||
},
|
||||
stopChan: make(chan struct{}),
|
||||
factory: factory,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Start begins watching ServiceAdapter resources. Blocks until cache syncs or timeout.
|
||||
func (im *InformerManager) Start(ctx context.Context) error {
|
||||
im.factory.Start(im.stopChan)
|
||||
|
||||
done := make(chan bool, 1)
|
||||
go func() { done <- cache.WaitForCacheSync(im.stopChan, im.informer.HasSynced) }()
|
||||
|
||||
select {
|
||||
case synced := <-done:
|
||||
if !synced {
|
||||
return fmt.Errorf("cache sync failed")
|
||||
// Start loads ServiceAdapters immediately, then polls every interval.
|
||||
func (l *Loader) Start(interval time.Duration) error {
|
||||
if err := l.load(); err != nil {
|
||||
// Retry once after 2s — handles transient "storage reinitializing" (429)
|
||||
log.Printf("serviceadapter loader: first attempt failed (%v), retrying in 2s", err)
|
||||
time.Sleep(2 * time.Second)
|
||||
if err := l.load(); err != nil {
|
||||
return fmt.Errorf("serviceadapter loader: %w", err)
|
||||
}
|
||||
case <-time.After(30 * time.Second):
|
||||
return fmt.Errorf("cache sync timed out")
|
||||
}
|
||||
|
||||
// Background poll for changes
|
||||
go func() {
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
if err := l.load(); err != nil {
|
||||
log.Printf("serviceadapter loader poll error: %v", err)
|
||||
}
|
||||
case <-l.stopChan:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Stop stops the informer.
|
||||
func (im *InformerManager) Stop() {
|
||||
close(im.stopChan)
|
||||
// Stop stops the background poll.
|
||||
func (l *Loader) Stop() {
|
||||
l.once.Do(func() { close(l.stopChan) })
|
||||
}
|
||||
|
||||
// parseServiceAdapter converts an unstructured object to ServiceAdapter.
|
||||
func parseServiceAdapter(u *unstructured.Unstructured) (*ServiceAdapter, error) {
|
||||
metadata, ok := u.Object["metadata"].(map[string]interface{})
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("invalid metadata")
|
||||
// load fetches all ServiceAdapters from the k8s API and syncs the registry.
|
||||
func (l *Loader) load() error {
|
||||
url := fmt.Sprintf("%s/apis/gateway.riotpiao.com/v1/namespaces/%s/serviceadapters", l.baseURL, l.namespace)
|
||||
req, err := http.NewRequest("GET", url, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
name, _ := metadata["name"].(string)
|
||||
namespace, _ := metadata["namespace"].(string)
|
||||
if namespace == "" {
|
||||
namespace = "default"
|
||||
req.Header.Set("Authorization", "Bearer "+l.token)
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
||||
resp, err := l.client.Do(req)
|
||||
if err != nil {
|
||||
return fmt.Errorf("k8s API request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode == 429 {
|
||||
return fmt.Errorf("k8s API: storage reinitializing (429)")
|
||||
}
|
||||
if resp.StatusCode != 200 {
|
||||
body, _ := io.ReadAll(resp.Body)
|
||||
return fmt.Errorf("k8s API: %d %s", resp.StatusCode, string(body[:min(len(body), 200)]))
|
||||
}
|
||||
|
||||
spec, ok := u.Object["spec"].(map[string]interface{})
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("missing spec")
|
||||
var list crList
|
||||
if err := json.NewDecoder(resp.Body).Decode(&list); err != nil {
|
||||
return fmt.Errorf("decoding response: %w", err)
|
||||
}
|
||||
|
||||
serviceName, ok := spec["serviceName"].(string)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("missing spec.serviceName")
|
||||
}
|
||||
|
||||
adapter := &ServiceAdapter{
|
||||
Name: name,
|
||||
Namespace: namespace,
|
||||
ServiceName: serviceName,
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
|
||||
if upstreamMap, ok := spec["upstream"].(map[string]interface{}); ok {
|
||||
if url, ok := upstreamMap["url"].(string); ok {
|
||||
adapter.Spec.Upstream.URL = url
|
||||
// Sync: add/update all found, registry handles dedup
|
||||
for i := range list.Items {
|
||||
item := &list.Items[i]
|
||||
adapter := &ServiceAdapter{
|
||||
Name: item.Metadata.Name,
|
||||
Namespace: item.Metadata.Namespace,
|
||||
ServiceName: item.Spec.ServiceName,
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
if timeout, ok := upstreamMap["timeoutSeconds"].(float64); ok {
|
||||
adapter.Spec.Upstream.TimeoutSeconds = int32(timeout)
|
||||
adapter.Spec.Upstream.URL = item.Spec.Upstream.URL
|
||||
adapter.Spec.Upstream.TimeoutSeconds = item.Spec.Upstream.TimeoutSeconds
|
||||
adapter.Spec.Auth.Required = item.Spec.Auth.Required
|
||||
adapter.Spec.Auth.Capability = item.Spec.Auth.Capability
|
||||
|
||||
if existing := l.registry.Get(item.Spec.ServiceName); existing != nil {
|
||||
l.registry.Update(adapter)
|
||||
} else {
|
||||
l.registry.Add(adapter)
|
||||
log.Printf("serviceadapter loaded: %s → %s", item.Spec.ServiceName, item.Spec.Upstream.URL)
|
||||
}
|
||||
}
|
||||
|
||||
if authMap, ok := spec["auth"].(map[string]interface{}); ok {
|
||||
if required, ok := authMap["required"].(bool); ok {
|
||||
adapter.Spec.Auth.Required = required
|
||||
}
|
||||
if capability, ok := authMap["capability"].(string); ok {
|
||||
adapter.Spec.Auth.Capability = capability
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
return adapter, nil
|
||||
// Minimal JSON structs for the k8s list response — no client-go needed.
|
||||
type crList struct {
|
||||
Items []crItem `json:"items"`
|
||||
}
|
||||
|
||||
type crItem struct {
|
||||
Metadata struct {
|
||||
Name string `json:"name"`
|
||||
Namespace string `json:"namespace"`
|
||||
} `json:"metadata"`
|
||||
Spec struct {
|
||||
ServiceName string `json:"serviceName"`
|
||||
Upstream struct {
|
||||
URL string `json:"url"`
|
||||
TimeoutSeconds int32 `json:"timeoutSeconds"`
|
||||
} `json:"upstream"`
|
||||
Auth struct {
|
||||
Required bool `json:"required"`
|
||||
Capability string `json:"capability"`
|
||||
} `json:"auth"`
|
||||
} `json:"spec"`
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user