diff --git a/buildermgr/buildermgr.go b/buildermgr/buildermgr.go index 9dac1495..a1c2fe9c 100644 --- a/buildermgr/buildermgr.go +++ b/buildermgr/buildermgr.go @@ -31,6 +31,11 @@ func Start(storageSvcUrl string, envBuilderNamespace string) error { return err } + err = fissionClient.WaitForCRDs() + if err != nil { + log.Fatalf("Error waiting for CRDs: %v", err) + } + envWatcher := makeEnvironmentWatcher(fissionClient, kubernetesClient, envBuilderNamespace) go envWatcher.watchEnvironments() diff --git a/controller/controller.go b/controller/controller.go index ba9afd95..5855737b 100644 --- a/controller/controller.go +++ b/controller/controller.go @@ -37,7 +37,10 @@ func Start(port int) { log.Fatalf("Failed to create fission CRDs: %v", err) } - fc.WaitForCRDs() + err = fc.WaitForCRDs() + if err != nil { + log.Fatalf("Error waiting for CRDs: %v", err) + } api, err := MakeAPI() if err != nil { diff --git a/crd/client.go b/crd/client.go index 5b8fff29..57389f42 100644 --- a/crd/client.go +++ b/crd/client.go @@ -212,8 +212,8 @@ func (fc *FissionClient) Packages(ns string) PackageInterface { return MakePackageInterface(fc.crdClient, ns) } -func (fc *FissionClient) WaitForCRDs() { - waitForCRDs(fc.crdClient) +func (fc *FissionClient) WaitForCRDs() error { + return waitForCRDs(fc.crdClient) } func (fc *FissionClient) GetCrdClient() *rest.RESTClient { return fc.crdClient diff --git a/crd/crd.go b/crd/crd.go index 0e2c54be..8dd2acb1 100644 --- a/crd/crd.go +++ b/crd/crd.go @@ -34,31 +34,31 @@ const ( // ensureCRD checks if the given CRD type exists, and creates it if // needed. (Note that this creates the CRD type; it doesn't create any // _instances_ of that type.) -func ensureCRD(clientset *apiextensionsclient.Clientset, crd *apiextensionsv1beta1.CustomResourceDefinition) error { +func ensureCRD(clientset *apiextensionsclient.Clientset, crd *apiextensionsv1beta1.CustomResourceDefinition) (err error) { maxRetries := 5 - for i := 0; i < maxRetries; i++ { - _, err := clientset.ApiextensionsV1beta1().CustomResourceDefinitions().Get(crd.ObjectMeta.Name, metav1.GetOptions{}) - if err != nil { - if errors.IsNotFound(err) { - // crd resource not found error - _, err := clientset.ApiextensionsV1beta1().CustomResourceDefinitions().Create(crd) - if err != nil { - return err - } - } else { - // The requests fail to connect to k8s api server before - // istio-prxoy is ready to serve traffic. Retry again. - log.Printf("Error connecting to kubernetes api service (%v), retrying", err) - time.Sleep(500 * time.Duration(2*i) * time.Millisecond) - continue - } + for i := 0; i < maxRetries; i++ { + _, err = clientset.ApiextensionsV1beta1().CustomResourceDefinitions().Get(crd.ObjectMeta.Name, metav1.GetOptions{}) + if err == nil { + return nil } - // resource already exists - break + if errors.IsNotFound(err) { + // crd resource not found error + _, err = clientset.ApiextensionsV1beta1().CustomResourceDefinitions().Create(crd) + if err != nil { + return err + } + } else { + // The requests fail to connect to k8s api server before + // istio-prxoy is ready to serve traffic. Retry again. + log.Printf("Error connecting to kubernetes api service (%v), retrying", err) + time.Sleep(500 * time.Duration(2*i) * time.Millisecond) + continue + } } - return nil + + return err } // Ensure CRDs diff --git a/executor/executor.go b/executor/executor.go index d6cb7a58..fced4e8d 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -193,6 +193,12 @@ func StartExecutor(fissionNamespace string, functionNamespace string, port int) fission.SetupStackTraceHandler() fissionClient, kubernetesClient, _, err := crd.MakeFissionClient() + + err = fissionClient.WaitForCRDs() + if err != nil { + log.Fatalf("Error waiting for CRDs: %v", err) + } + restClient := fissionClient.GetCrdClient() if err != nil { log.Printf("Failed to get kubernetes client: %v", err) diff --git a/executor/executor_test.go b/executor/executor_test.go index 39da2e40..19153277 100644 --- a/executor/executor_test.go +++ b/executor/executor_test.go @@ -138,7 +138,11 @@ func TestExecutor(t *testing.T) { if err != nil { log.Panicf("failed to ensure crds: %v", err) } - fissionClient.WaitForCRDs() + + err = fissionClient.WaitForCRDs() + if err != nil { + log.Panicf("failed to wait crds: %v", err) + } // create an env on the cluster env, err := fissionClient.Environments(fissionNs).Create(&crd.Environment{ diff --git a/kubewatcher/main.go b/kubewatcher/main.go index b99e898f..acad8dcd 100644 --- a/kubewatcher/main.go +++ b/kubewatcher/main.go @@ -17,6 +17,8 @@ limitations under the License. package kubewatcher import ( + "log" + "github.com/fission/fission/crd" "github.com/fission/fission/publisher" ) @@ -26,6 +28,12 @@ func Start(routerUrl string) error { if err != nil { return err } + + err = fissionClient.WaitForCRDs() + if err != nil { + log.Fatalf("Error waiting for CRDs: %v", err) + } + poster := publisher.MakeWebhookPublisher(routerUrl) kubeWatch := MakeKubeWatcher(kubeClient, poster) MakeWatchSync(fissionClient, kubeWatch) diff --git a/mqtrigger/main.go b/mqtrigger/main.go index 81fea523..3f888ecb 100644 --- a/mqtrigger/main.go +++ b/mqtrigger/main.go @@ -30,6 +30,11 @@ func Start(routerUrl string) error { log.Fatalf("Failed to get fission client: %v", err) } + err = fissionClient.WaitForCRDs() + if err != nil { + log.Fatalf("Error waiting for CRDs: %v", err) + } + // Message queue type: nats is the only supported one for now mqType := os.Getenv("MESSAGE_QUEUE_TYPE") mqUrl := os.Getenv("MESSAGE_QUEUE_URL") diff --git a/router/router.go b/router/router.go index 4497b7ae..50fce5a7 100644 --- a/router/router.go +++ b/router/router.go @@ -82,6 +82,11 @@ func Start(port int, executorUrl string) { log.Fatalf("Error connecting to kubernetes API: %v", err) } + err = fissionClient.WaitForCRDs() + if err != nil { + log.Fatalf("Error waiting for CRDs: %v", err) + } + restClient := fissionClient.GetCrdClient() executor := executorClient.MakeClient(executorUrl) diff --git a/timer/main.go b/timer/main.go index 6234efb5..0186c211 100644 --- a/timer/main.go +++ b/timer/main.go @@ -17,6 +17,8 @@ limitations under the License. package timer import ( + "log" + "github.com/fission/fission/crd" "github.com/fission/fission/publisher" ) @@ -27,6 +29,11 @@ func Start(routerUrl string) error { return err } + err = fissionClient.WaitForCRDs() + if err != nil { + log.Fatalf("Error waiting for CRDs: %v", err) + } + poster := publisher.MakeWebhookPublisher(routerUrl) MakeTimerSync(fissionClient, MakeTimer(poster))