From ebb486d121204059a351b164f7175443bef78719 Mon Sep 17 00:00:00 2001 From: Soam Vasani Date: Thu, 5 Jan 2017 15:44:49 -0800 Subject: [PATCH 1/4] Delete pods that failed to load user function --- poolmgr/gp.go | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/poolmgr/gp.go b/poolmgr/gp.go index 23a91051..123b7aeb 100644 --- a/poolmgr/gp.go +++ b/poolmgr/gp.go @@ -215,6 +215,12 @@ func (gp *GenericPool) labelsForMetadata(metadata *fission.Metadata) map[string] } } +func (gp *GenericPool) tryDeletePod(name string) { + go func() { + gp.kubernetesClient.Core().Pods(gp.namespace).Delete(name, nil) + }() +} + // specializePod chooses a pod, copies the required user-defined function to that pod // (via fetcher), and calls the function-run container to load it, resulting in a // specialized pod. @@ -230,6 +236,7 @@ func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error // for fetcher we don't need to create a service, just talk to the pod directly podIP := pod.Status.PodIP if len(podIP) == 0 { + gp.tryDeletePod(pod.ObjectMeta.Name) return nil, errors.New("Pod has no IP") } @@ -242,11 +249,13 @@ func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error log.Printf("[%v] calling fetcher to copy function", metadata) resp, err := http.Post(fetcherUrl, "application/json", bytes.NewReader([]byte(fetcherRequest))) if err != nil { + gp.tryDeletePod(pod.ObjectMeta.Name) // TODO we should retry this call in case fetcher hasn't come up yet return nil, err } defer resp.Body.Close() if resp.StatusCode != 200 { + gp.tryDeletePod(pod.ObjectMeta.Name) return nil, errors.New(fmt.Sprintf("Error from fetcher: %v", resp.Status)) } @@ -280,6 +289,7 @@ func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error err = fission.MakeErrorFromHTTP(resp2) } log.Printf("Failed to specialize pod: %v", err) + gp.tryDeletePod(pod.ObjectMeta.Name) return nil, err } From e4f7027a6984198edcc1bd3c7b2d6228cf4f7947 Mon Sep 17 00:00:00 2001 From: Soam Vasani Date: Fri, 6 Jan 2017 21:54:15 -0800 Subject: [PATCH 2/4] Reorganize a bit to avoid duplicated error handling --- poolmgr/gp.go | 46 ++++++++++++++++++++++++---------------------- 1 file changed, 24 insertions(+), 22 deletions(-) diff --git a/poolmgr/gp.go b/poolmgr/gp.go index 123b7aeb..820cc941 100644 --- a/poolmgr/gp.go +++ b/poolmgr/gp.go @@ -215,8 +215,12 @@ func (gp *GenericPool) labelsForMetadata(metadata *fission.Metadata) map[string] } } -func (gp *GenericPool) tryDeletePod(name string) { +func (gp *GenericPool) scheduleDeletePod(name string) { go func() { + // The sleep allows debugging or collecting logs from the pod before it's + // cleaned up. (We need a better solutions for both those things; log + // aggregation and storage will help.) + time.Sleep(5 * time.Minute) gp.kubernetesClient.Core().Pods(gp.namespace).Delete(name, nil) }() } @@ -224,23 +228,14 @@ func (gp *GenericPool) tryDeletePod(name string) { // specializePod chooses a pod, copies the required user-defined function to that pod // (via fetcher), and calls the function-run container to load it, resulting in a // specialized pod. -func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error) { - newLabels := gp.labelsForMetadata(metadata) - - log.Printf("[%v] Choosing pod from pool", metadata) - pod, err := gp.choosePod(newLabels) - if err != nil { - return nil, err - } - +func (gp *GenericPool) specializePod(pod *v1.Pod, metadata *fission.Metadata) error { // for fetcher we don't need to create a service, just talk to the pod directly podIP := pod.Status.PodIP if len(podIP) == 0 { - gp.tryDeletePod(pod.ObjectMeta.Name) - return nil, errors.New("Pod has no IP") + return errors.New("Pod has no IP") } - // tell fetcher to get the function + // tell fetcher to get the function. fetcherUrl := fmt.Sprintf("http://%v:8000/", podIP) functionUrl := fmt.Sprintf("%v/v1/functions/%v?uid=%v&raw=1", gp.controllerUrl, metadata.Name, metadata.Uid) @@ -249,14 +244,12 @@ func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error log.Printf("[%v] calling fetcher to copy function", metadata) resp, err := http.Post(fetcherUrl, "application/json", bytes.NewReader([]byte(fetcherRequest))) if err != nil { - gp.tryDeletePod(pod.ObjectMeta.Name) // TODO we should retry this call in case fetcher hasn't come up yet - return nil, err + return err } defer resp.Body.Close() if resp.StatusCode != 200 { - gp.tryDeletePod(pod.ObjectMeta.Name) - return nil, errors.New(fmt.Sprintf("Error from fetcher: %v", resp.Status)) + return errors.New(fmt.Sprintf("Error from fetcher: %v", resp.Status)) } // get function run container to specialize @@ -268,8 +261,9 @@ func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error for i := 0; i < maxRetries; i++ { resp2, err := http.Post(specializeUrl, "text/plain", bytes.NewReader([]byte{})) if err == nil && resp2.StatusCode < 300 { + // Success resp2.Body.Close() - return pod, nil + return nil } // Only retry for the specific case of a connection error. @@ -289,11 +283,10 @@ func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error err = fission.MakeErrorFromHTTP(resp2) } log.Printf("Failed to specialize pod: %v", err) - gp.tryDeletePod(pod.ObjectMeta.Name) - return nil, err + return err } - return pod, nil + return nil } // A pool is a deployment of generic containers for an env. This @@ -417,10 +410,19 @@ func (gp *GenericPool) createSvc(name string, labels map[string]string) (*v1.Ser } func (gp *GenericPool) GetFuncSvc(m *fission.Metadata) (*funcSvc, error) { - pod, err := gp.specializePod(m) + + log.Printf("[%v] Choosing pod from pool", m) + newLabels := gp.labelsForMetadata(m) + pod, err := gp.choosePod(newLabels) if err != nil { return nil, err } + + err = gp.specializePod(pod, m) + if err != nil { + gp.scheduleDeletePod(pod.ObjectMeta.Name) + return nil, err + } log.Printf("Specialized pod: %v", pod.ObjectMeta.Name) var svcHost string From 3e386fb41d4b901fece19f387332b057844eeb29 Mon Sep 17 00:00:00 2001 From: Soam Vasani Date: Fri, 6 Jan 2017 22:09:17 -0800 Subject: [PATCH 3/4] Handle pod cleanup on svc creation error --- poolmgr/gp.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/poolmgr/gp.go b/poolmgr/gp.go index 820cc941..8264f7aa 100644 --- a/poolmgr/gp.go +++ b/poolmgr/gp.go @@ -435,9 +435,11 @@ func (gp *GenericPool) GetFuncSvc(m *fission.Metadata) (*funcSvc, error) { labels := gp.labelsForMetadata(m) svc, err := gp.createSvc(svcName, labels) if err != nil { + gp.scheduleDeletePod(pod.ObjectMeta.Name) return nil, err } if svc.ObjectMeta.Name != svcName { + gp.scheduleDeletePod(pod.ObjectMeta.Name) return nil, errors.New(fmt.Sprintf("sanity check failed for svc %v", svc.ObjectMeta.Name)) } From 0fb12af4e551700a0c2df938802a3ba753d4c6b9 Mon Sep 17 00:00:00 2001 From: Soam Vasani Date: Fri, 6 Jan 2017 22:16:49 -0800 Subject: [PATCH 4/4] add log on specializePod error/clean up --- poolmgr/gp.go | 1 + 1 file changed, 1 insertion(+) diff --git a/poolmgr/gp.go b/poolmgr/gp.go index 8264f7aa..be290066 100644 --- a/poolmgr/gp.go +++ b/poolmgr/gp.go @@ -220,6 +220,7 @@ func (gp *GenericPool) scheduleDeletePod(name string) { // The sleep allows debugging or collecting logs from the pod before it's // cleaned up. (We need a better solutions for both those things; log // aggregation and storage will help.) + log.Printf("Error in pod '%v', scheduling cleanup", name) time.Sleep(5 * time.Minute) gp.kubernetesClient.Core().Pods(gp.namespace).Delete(name, nil) }()