package serviceadapter import ( "context" "fmt" "log" "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 resources and populates the registry. type InformerManager struct { registry *Registry informer cache.SharedIndexInformer stopChan chan struct{} factory dynamicinformer.DynamicSharedInformerFactory } // NewInformerManager creates a new informer that watches ServiceAdapters in the api namespace. func NewInformerManager(dynamicClient dynamic.Interface, registry *Registry, namespace string) (*InformerManager, error) { log.Printf("[INFORMER] NewInformerManager called with namespace=%s", namespace) if registry == nil { return nil, fmt.Errorf("registry cannot be nil") } if namespace == "" { namespace = "api" } log.Printf("[INFORMER] Using namespace: %s", namespace) // ServiceAdapter GVR gvr := schema.GroupVersionResource{ Group: "gateway.riotpiao.com", Version: "v1", Resource: "serviceadapters", } log.Printf("[INFORMER] GVR: %s/%s/%s", gvr.Group, gvr.Version, gvr.Resource) // Create informer factory scoped to namespace log.Printf("[INFORMER] Creating filtered informer factory") factory := dynamicinformer.NewFilteredDynamicSharedInformerFactory( dynamicClient, 30*time.Second, namespace, nil, ) log.Printf("[INFORMER] Factory created") // Get the informer for ServiceAdapters log.Printf("[INFORMER] Getting informer for ServiceAdapters") informer := factory.ForResource(gvr).Informer() log.Printf("[INFORMER] Informer obtained") // Add event handlers log.Printf("[INFORMER] Adding event handlers") _, err := informer.AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { log.Printf("[INFORMER] AddFunc called") unstructObj := obj.(*unstructured.Unstructured) if adapter, err := parseServiceAdapter(unstructObj); err == nil { log.Printf("[INFORMER] Parsed adapter: %s/%s", adapter.Namespace, adapter.ServiceName) if registry.Add(adapter) == nil { log.Printf("[INFORMER] ServiceAdapter added: %s/%s", adapter.Namespace, adapter.ServiceName) } } else { log.Printf("[INFORMER] ERROR parsing adapter: %v", err) } }, UpdateFunc: func(oldObj, newObj interface{}) { log.Printf("[INFORMER] UpdateFunc called") unstructObj := newObj.(*unstructured.Unstructured) if adapter, err := parseServiceAdapter(unstructObj); err == nil { log.Printf("[INFORMER] ServiceAdapter updated: %s/%s", adapter.Namespace, adapter.ServiceName) if registry.Update(adapter) == nil { log.Printf("[INFORMER] Update succeeded") } } }, DeleteFunc: func(obj interface{}) { log.Printf("[INFORMER] DeleteFunc called") unstructObj, ok := obj.(*unstructured.Unstructured) if !ok { // Handle tombstone (for objects that were deleted) tombstone, ok := obj.(cache.DeletedFinalStateUnknown) if !ok { return } unstructObj, ok = tombstone.Obj.(*unstructured.Unstructured) if !ok { return } } if name, ok := unstructObj.Object["metadata"].(map[string]interface{})["name"]; ok { registry.Delete(name.(string)) log.Printf("[INFORMER] ServiceAdapter deleted: %s", name) } }, }) log.Printf("[INFORMER] Event handlers added") if err != nil { return nil, err } return &InformerManager{ registry: registry, informer: informer, stopChan: make(chan struct{}), factory: factory, }, nil } // Start begins watching ServiceAdapter resources. func (im *InformerManager) Start(ctx context.Context) error { log.Printf("[INFORMER] Start() called") im.factory.Start(im.stopChan) log.Printf("[INFORMER] Factory started, waiting for cache sync") // Wait for informer cache to sync with 30-second timeout done := make(chan bool, 1) go func() { done <- cache.WaitForCacheSync(im.stopChan, im.informer.HasSynced) }() select { case synced := <-done: if !synced { log.Printf("[INFORMER] ERROR: Failed to sync cache") return fmt.Errorf("failed to sync ServiceAdapter informer cache") } log.Printf("[INFORMER] Cache synced successfully!") case <-time.After(30 * time.Second): log.Printf("[INFORMER] ERROR: Cache sync timed out after 30s") return fmt.Errorf("timeout waiting for cache sync") } log.Printf("[INFORMER] ServiceAdapter informer started, synced from cluster") return nil } // Stop stops watching ServiceAdapter resources. func (im *InformerManager) Stop() { close(im.stopChan) } // parseServiceAdapter converts an unstructured object to ServiceAdapter. func parseServiceAdapter(unstructured *unstructured.Unstructured) (*ServiceAdapter, error) { metadata, ok := unstructured.Object["metadata"].(map[string]interface{}) if !ok { return nil, fmt.Errorf("invalid metadata") } name, ok := metadata["name"].(string) if !ok { return nil, fmt.Errorf("missing metadata.name") } namespace, ok := metadata["namespace"].(string) if !ok { namespace = "default" } spec, ok := unstructured.Object["spec"].(map[string]interface{}) if !ok { return nil, fmt.Errorf("missing spec") } // Parse spec fields (simplified - in real implementation would use JSON unmarshaling) serviceName, ok := spec["serviceName"].(string) if !ok { return nil, fmt.Errorf("missing serviceName") } adapter := &ServiceAdapter{ Name: name, Namespace: namespace, ServiceName: serviceName, CreatedAt: time.Now(), } // Parse upstream if present if upstreamMap, ok := spec["upstream"].(map[string]interface{}); ok { if url, ok := upstreamMap["url"].(string); ok { adapter.Spec.Upstream.URL = url } if timeout, ok := upstreamMap["timeoutSeconds"].(float64); ok { adapter.Spec.Upstream.TimeoutSeconds = int32(timeout) } } // Parse auth if present 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 adapter, nil }