Fix components crash before crds creation (#602)
* Wait for CRDs creation for 30 sec when component start * Fix ensureCRD return nil while the error is not empty
This commit is contained in:
@@ -31,6 +31,11 @@ func Start(storageSvcUrl string, envBuilderNamespace string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
err = fissionClient.WaitForCRDs()
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Error waiting for CRDs: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
envWatcher := makeEnvironmentWatcher(fissionClient, kubernetesClient, envBuilderNamespace)
|
envWatcher := makeEnvironmentWatcher(fissionClient, kubernetesClient, envBuilderNamespace)
|
||||||
go envWatcher.watchEnvironments()
|
go envWatcher.watchEnvironments()
|
||||||
|
|
||||||
|
|||||||
@@ -37,7 +37,10 @@ func Start(port int) {
|
|||||||
log.Fatalf("Failed to create fission CRDs: %v", err)
|
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()
|
api, err := MakeAPI()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
+2
-2
@@ -212,8 +212,8 @@ func (fc *FissionClient) Packages(ns string) PackageInterface {
|
|||||||
return MakePackageInterface(fc.crdClient, ns)
|
return MakePackageInterface(fc.crdClient, ns)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (fc *FissionClient) WaitForCRDs() {
|
func (fc *FissionClient) WaitForCRDs() error {
|
||||||
waitForCRDs(fc.crdClient)
|
return waitForCRDs(fc.crdClient)
|
||||||
}
|
}
|
||||||
func (fc *FissionClient) GetCrdClient() *rest.RESTClient {
|
func (fc *FissionClient) GetCrdClient() *rest.RESTClient {
|
||||||
return fc.crdClient
|
return fc.crdClient
|
||||||
|
|||||||
+20
-20
@@ -34,31 +34,31 @@ const (
|
|||||||
// ensureCRD checks if the given CRD type exists, and creates it if
|
// 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
|
// needed. (Note that this creates the CRD type; it doesn't create any
|
||||||
// _instances_ of that type.)
|
// _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
|
maxRetries := 5
|
||||||
for i := 0; i < maxRetries; i++ {
|
|
||||||
_, err := clientset.ApiextensionsV1beta1().CustomResourceDefinitions().Get(crd.ObjectMeta.Name, metav1.GetOptions{})
|
|
||||||
|
|
||||||
if err != nil {
|
for i := 0; i < maxRetries; i++ {
|
||||||
if errors.IsNotFound(err) {
|
_, err = clientset.ApiextensionsV1beta1().CustomResourceDefinitions().Get(crd.ObjectMeta.Name, metav1.GetOptions{})
|
||||||
// crd resource not found error
|
if err == nil {
|
||||||
_, err := clientset.ApiextensionsV1beta1().CustomResourceDefinitions().Create(crd)
|
return nil
|
||||||
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
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// resource already exists
|
if errors.IsNotFound(err) {
|
||||||
break
|
// 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
|
// Ensure CRDs
|
||||||
|
|||||||
@@ -193,6 +193,12 @@ func StartExecutor(fissionNamespace string, functionNamespace string, port int)
|
|||||||
fission.SetupStackTraceHandler()
|
fission.SetupStackTraceHandler()
|
||||||
|
|
||||||
fissionClient, kubernetesClient, _, err := crd.MakeFissionClient()
|
fissionClient, kubernetesClient, _, err := crd.MakeFissionClient()
|
||||||
|
|
||||||
|
err = fissionClient.WaitForCRDs()
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Error waiting for CRDs: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
restClient := fissionClient.GetCrdClient()
|
restClient := fissionClient.GetCrdClient()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("Failed to get kubernetes client: %v", err)
|
log.Printf("Failed to get kubernetes client: %v", err)
|
||||||
|
|||||||
@@ -138,7 +138,11 @@ func TestExecutor(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
log.Panicf("failed to ensure crds: %v", err)
|
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
|
// create an env on the cluster
|
||||||
env, err := fissionClient.Environments(fissionNs).Create(&crd.Environment{
|
env, err := fissionClient.Environments(fissionNs).Create(&crd.Environment{
|
||||||
|
|||||||
@@ -17,6 +17,8 @@ limitations under the License.
|
|||||||
package kubewatcher
|
package kubewatcher
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"log"
|
||||||
|
|
||||||
"github.com/fission/fission/crd"
|
"github.com/fission/fission/crd"
|
||||||
"github.com/fission/fission/publisher"
|
"github.com/fission/fission/publisher"
|
||||||
)
|
)
|
||||||
@@ -26,6 +28,12 @@ func Start(routerUrl string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
err = fissionClient.WaitForCRDs()
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Error waiting for CRDs: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
poster := publisher.MakeWebhookPublisher(routerUrl)
|
poster := publisher.MakeWebhookPublisher(routerUrl)
|
||||||
kubeWatch := MakeKubeWatcher(kubeClient, poster)
|
kubeWatch := MakeKubeWatcher(kubeClient, poster)
|
||||||
MakeWatchSync(fissionClient, kubeWatch)
|
MakeWatchSync(fissionClient, kubeWatch)
|
||||||
|
|||||||
@@ -30,6 +30,11 @@ func Start(routerUrl string) error {
|
|||||||
log.Fatalf("Failed to get fission client: %v", err)
|
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
|
// Message queue type: nats is the only supported one for now
|
||||||
mqType := os.Getenv("MESSAGE_QUEUE_TYPE")
|
mqType := os.Getenv("MESSAGE_QUEUE_TYPE")
|
||||||
mqUrl := os.Getenv("MESSAGE_QUEUE_URL")
|
mqUrl := os.Getenv("MESSAGE_QUEUE_URL")
|
||||||
|
|||||||
@@ -82,6 +82,11 @@ func Start(port int, executorUrl string) {
|
|||||||
log.Fatalf("Error connecting to kubernetes API: %v", err)
|
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()
|
restClient := fissionClient.GetCrdClient()
|
||||||
|
|
||||||
executor := executorClient.MakeClient(executorUrl)
|
executor := executorClient.MakeClient(executorUrl)
|
||||||
|
|||||||
@@ -17,6 +17,8 @@ limitations under the License.
|
|||||||
package timer
|
package timer
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"log"
|
||||||
|
|
||||||
"github.com/fission/fission/crd"
|
"github.com/fission/fission/crd"
|
||||||
"github.com/fission/fission/publisher"
|
"github.com/fission/fission/publisher"
|
||||||
)
|
)
|
||||||
@@ -27,6 +29,11 @@ func Start(routerUrl string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
err = fissionClient.WaitForCRDs()
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("Error waiting for CRDs: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
poster := publisher.MakeWebhookPublisher(routerUrl)
|
poster := publisher.MakeWebhookPublisher(routerUrl)
|
||||||
MakeTimerSync(fissionClient, MakeTimer(poster))
|
MakeTimerSync(fissionClient, MakeTimer(poster))
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user