Fixed golangci-lint issues: /fission/pkg (#1902)

This commit is contained in:
Gaurav Gahlot
2021-01-19 23:29:20 +05:30
committed by GitHub
parent 6dbdf4d88b
commit 695b0759e4
21 changed files with 178 additions and 66 deletions
+14 -2
View File
@@ -74,7 +74,13 @@ func MakeBuilder(logger *zap.Logger, sharedVolumePath string) *Builder {
func (builder *Builder) VersionHandler(w http.ResponseWriter, r *http.Request) { func (builder *Builder) VersionHandler(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json; charset=utf-8") w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.Write([]byte(info.BuildInfo().String())) _, err := w.Write([]byte(info.BuildInfo().String()))
if err != nil {
builder.logger.Error(
"error writing response",
zap.Error(err),
)
}
} }
func (builder *Builder) Handler(w http.ResponseWriter, r *http.Request) { func (builder *Builder) Handler(w http.ResponseWriter, r *http.Request) {
@@ -149,7 +155,13 @@ func (builder *Builder) reply(w http.ResponseWriter, pkgFilename string, buildLo
// should write header before writing the body, // should write header before writing the body,
// or client will receive HTTP 200 regardless the real status code // or client will receive HTTP 200 regardless the real status code
w.WriteHeader(statusCode) w.WriteHeader(statusCode)
w.Write(rBody) _, err = w.Write(rBody)
if err != nil {
builder.logger.Error(
"error writing response",
zap.Error(err),
)
}
} }
func (builder *Builder) build(command string, srcPkgPath string, deployPkgPath string) (string, error) { func (builder *Builder) build(command string, srcPkgPath string, deployPkgPath string) (string, error) {
+60 -7
View File
@@ -81,7 +81,12 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
if err != nil { if err != nil {
return return
} }
defer buildCache.Delete(key) defer func() {
err := buildCache.Delete(key)
if err != nil {
pkgw.logger.Error("error deleting key from cache", zap.String("key", key), zap.Error(err))
}
}()
pkgw.logger.Info("starting build for package", zap.String("package_name", srcpkg.ObjectMeta.Name), zap.String("resource_version", srcpkg.ObjectMeta.ResourceVersion)) pkgw.logger.Info("starting build for package", zap.String("package_name", srcpkg.ObjectMeta.Name), zap.String("resource_version", srcpkg.ObjectMeta.ResourceVersion))
@@ -95,8 +100,16 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
if k8serrors.IsNotFound(err) { if k8serrors.IsNotFound(err) {
e := "environment does not exist" e := "environment does not exist"
pkgw.logger.Error(e, zap.String("environment", pkg.Spec.Environment.Name)) pkgw.logger.Error(e, zap.String("environment", pkg.Spec.Environment.Name))
updatePackage(pkgw.logger, pkgw.fissionClient, pkg, _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg,
fv1.BuildStatusFailed, fmt.Sprintf("%s: %q", e, pkg.Spec.Environment.Name), nil) fv1.BuildStatusFailed, fmt.Sprintf("%s: %q", e, pkg.Spec.Environment.Name), nil)
if er != nil {
pkgw.logger.Error(
"error updating package",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
zap.Error(er),
)
}
return return
} }
@@ -173,7 +186,15 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
uploadResp, buildLogs, err := buildPackage(ctx, pkgw.logger, pkgw.fissionClient, builderNs, pkgw.storageSvcUrl, pkg) uploadResp, buildLogs, err := buildPackage(ctx, pkgw.logger, pkgw.fissionClient, builderNs, pkgw.storageSvcUrl, pkg)
if err != nil { if err != nil {
pkgw.logger.Error("error building package", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name)) pkgw.logger.Error("error building package", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name))
updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
if er != nil {
pkgw.logger.Error(
"error updating package",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
zap.Error(er),
)
}
return return
} }
@@ -185,7 +206,15 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
e := "error getting function list" e := "error getting function list"
pkgw.logger.Error(e, zap.Error(err)) pkgw.logger.Error(e, zap.Error(err))
buildLogs += fmt.Sprintf("%s: %v\n", e, err) buildLogs += fmt.Sprintf("%s: %v\n", e, err)
updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
if er != nil {
pkgw.logger.Error(
"error updating package",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
zap.Error(er),
)
}
} }
// A package may be used by multiple functions. Update // A package may be used by multiple functions. Update
@@ -201,7 +230,15 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
e := "error updating function package resource version" e := "error updating function package resource version"
pkgw.logger.Error(e, zap.Error(err)) pkgw.logger.Error(e, zap.Error(err))
buildLogs += fmt.Sprintf("%s: %v\n", e, err) buildLogs += fmt.Sprintf("%s: %v\n", e, err)
updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
if er != nil {
pkgw.logger.Error(
"error updating package",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
zap.Error(er),
)
}
return return
} }
} }
@@ -211,7 +248,15 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
fv1.BuildStatusSucceeded, buildLogs, uploadResp) fv1.BuildStatusSucceeded, buildLogs, uploadResp)
if err != nil { if err != nil {
pkgw.logger.Error("error updating package info", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name)) pkgw.logger.Error("error updating package info", zap.Error(err), zap.String("package_name", pkg.ObjectMeta.Name))
updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil) _, er := updatePackage(pkgw.logger, pkgw.fissionClient, pkg, fv1.BuildStatusFailed, buildLogs, nil)
if er != nil {
pkgw.logger.Error(
"error updating package",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
zap.Error(er),
)
}
return return
} }
@@ -221,8 +266,16 @@ func (pkgw *packageWatcher) build(buildCache *cache.Cache, srcpkg *fv1.Package)
time.Sleep(healthCheckBackOff.GetNext()) time.Sleep(healthCheckBackOff.GetNext())
} }
// build timeout // build timeout
updatePackage(pkgw.logger, pkgw.fissionClient, pkg, _, err = updatePackage(pkgw.logger, pkgw.fissionClient, pkg,
fv1.BuildStatusFailed, "Build timeout due to environment builder not ready", nil) fv1.BuildStatusFailed, "Build timeout due to environment builder not ready", nil)
if err != nil {
pkgw.logger.Error(
"error updating package",
zap.String("package_name", pkg.ObjectMeta.Name),
zap.String("resource_version", pkg.ObjectMeta.ResourceVersion),
zap.Error(err),
)
}
pkgw.logger.Error("max retries exceeded in building source package, timeout due to environment builder not ready", pkgw.logger.Error("max retries exceeded in building source package, timeout due to environment builder not ready",
zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace))) zap.String("package", fmt.Sprintf("%s.%s", pkg.ObjectMeta.Name, pkg.ObjectMeta.Namespace)))
+4 -1
View File
@@ -166,7 +166,10 @@ func (api *API) getLogDBConfig(dbType string) logDBConfig {
func (api *API) HomeHandler(w http.ResponseWriter, r *http.Request) { func (api *API) HomeHandler(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json; charset=utf-8") w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.Write([]byte(info.ApiInfo().String())) _, err := w.Write([]byte(info.ApiInfo().String()))
if err != nil {
api.respondWithError(w, err)
}
} }
func (api *API) ApiVersionMismatchHandler(w http.ResponseWriter, r *http.Request) { func (api *API) ApiVersionMismatchHandler(w http.ResponseWriter, r *http.Request) {
+11 -10
View File
@@ -129,7 +129,7 @@ func TestFunctionApi(t *testing.T) {
testFunc.ObjectMeta.Name = "bar" testFunc.ObjectMeta.Name = "bar"
m2, err := g.Client().V1().Function().Create(testFunc) m2, err := g.Client().V1().Function().Create(testFunc)
panicIf(err) panicIf(err)
defer g.Client().V1().Function().Delete(m2) defer panicIf(g.Client().V1().Function().Delete(m2))
funcs, err := g.Client().V1().Function().List(testNS) funcs, err := g.Client().V1().Function().List(testNS)
panicIf(err) panicIf(err)
@@ -173,7 +173,7 @@ func TestHTTPTriggerApi(t *testing.T) {
m, err := g.Client().V1().HTTPTrigger().Create(testTrigger) m, err := g.Client().V1().HTTPTrigger().Create(testTrigger)
panicIf(err) panicIf(err)
defer g.Client().V1().HTTPTrigger().Delete(m) defer panicIf(g.Client().V1().HTTPTrigger().Delete(m))
_, err = g.Client().V1().HTTPTrigger().Create(testTrigger) _, err = g.Client().V1().HTTPTrigger().Create(testTrigger)
assertNameReuseFailure(err, "httptrigger") assertNameReuseFailure(err, "httptrigger")
@@ -198,7 +198,7 @@ func TestHTTPTriggerApi(t *testing.T) {
testTrigger.Spec.RelativeURL = "/hi2" testTrigger.Spec.RelativeURL = "/hi2"
m2, err := g.Client().V1().HTTPTrigger().Create(testTrigger) m2, err := g.Client().V1().HTTPTrigger().Create(testTrigger)
panicIf(err) panicIf(err)
defer g.Client().V1().HTTPTrigger().Delete(m2) defer panicIf(g.Client().V1().HTTPTrigger().Delete(m2))
ts, err := g.Client().V1().HTTPTrigger().List(testNS) ts, err := g.Client().V1().HTTPTrigger().List(testNS)
panicIf(err) panicIf(err)
@@ -227,7 +227,7 @@ func TestEnvironmentApi(t *testing.T) {
m, err := g.Client().V1().Environment().Create(testEnv) m, err := g.Client().V1().Environment().Create(testEnv)
panicIf(err) panicIf(err)
defer g.Client().V1().Environment().Delete(m) defer panicIf(g.Client().V1().Environment().Delete(m))
_, err = g.Client().V1().Environment().Create(testEnv) _, err = g.Client().V1().Environment().Create(testEnv)
assertNameReuseFailure(err, "environment") assertNameReuseFailure(err, "environment")
@@ -246,7 +246,7 @@ func TestEnvironmentApi(t *testing.T) {
m2, err := g.Client().V1().Environment().Create(testEnv) m2, err := g.Client().V1().Environment().Create(testEnv)
panicIf(err) panicIf(err)
defer g.Client().V1().Environment().Delete(m2) defer panicIf(g.Client().V1().Environment().Delete(m2))
ts, err := g.Client().V1().Environment().List(testNS) ts, err := g.Client().V1().Environment().List(testNS)
panicIf(err) panicIf(err)
@@ -276,7 +276,7 @@ func TestWatchApi(t *testing.T) {
m, err := g.Client().V1().KubeWatcher().Create(testWatch) m, err := g.Client().V1().KubeWatcher().Create(testWatch)
panicIf(err) panicIf(err)
defer g.Client().V1().KubeWatcher().Delete(m) defer panicIf(g.Client().V1().KubeWatcher().Delete(m))
_, err = g.Client().V1().KubeWatcher().Create(testWatch) _, err = g.Client().V1().KubeWatcher().Create(testWatch)
assertNameReuseFailure(err, "watch") assertNameReuseFailure(err, "watch")
@@ -291,7 +291,7 @@ func TestWatchApi(t *testing.T) {
testWatch.ObjectMeta.Name = "yyy" testWatch.ObjectMeta.Name = "yyy"
m2, err := g.Client().V1().KubeWatcher().Create(testWatch) m2, err := g.Client().V1().KubeWatcher().Create(testWatch)
panicIf(err) panicIf(err)
defer g.Client().V1().KubeWatcher().Delete(m2) defer panicIf(g.Client().V1().KubeWatcher().Delete(m2))
ws, err := g.Client().V1().KubeWatcher().List(testNS) ws, err := g.Client().V1().KubeWatcher().List(testNS)
panicIf(err) panicIf(err)
@@ -317,7 +317,7 @@ func TestTimeTriggerApi(t *testing.T) {
m, err := g.Client().V1().TimeTrigger().Create(testTrigger) m, err := g.Client().V1().TimeTrigger().Create(testTrigger)
panicIf(err) panicIf(err)
defer g.Client().V1().TimeTrigger().Delete(m) defer panicIf(g.Client().V1().TimeTrigger().Delete(m))
_, err = g.Client().V1().TimeTrigger().Create(testTrigger) _, err = g.Client().V1().TimeTrigger().Create(testTrigger)
assertNameReuseFailure(err, "trigger") assertNameReuseFailure(err, "trigger")
@@ -359,12 +359,13 @@ func TestMain(m *testing.M) {
// testNS isolation for running multiple CI builds concurrently. // testNS isolation for running multiple CI builds concurrently.
testNS = uuid.NewV4().String() testNS = uuid.NewV4().String()
kubeClient.CoreV1().Namespaces().Create(&v1.Namespace{ _, err = kubeClient.CoreV1().Namespaces().Create(&v1.Namespace{
ObjectMeta: metav1.ObjectMeta{ ObjectMeta: metav1.ObjectMeta{
Name: testNS, Name: testNS,
}, },
}) })
defer kubeClient.CoreV1().Namespaces().Delete(testNS, nil) panicIf(err)
defer panicIf(kubeClient.CoreV1().Namespaces().Delete(testNS, nil))
config := zap.NewDevelopmentConfig() config := zap.NewDevelopmentConfig()
config.EncoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder config.EncoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder
+1 -1
View File
@@ -107,7 +107,7 @@ func (c *Function) Get(m *metav1.ObjectMeta) (*fv1.Function, error) {
func (c *Function) GetRawDeployment(m *metav1.ObjectMeta) ([]byte, error) { func (c *Function) GetRawDeployment(m *metav1.ObjectMeta) ([]byte, error) {
relativeUrl := fmt.Sprintf("functions/%v", m.Name) relativeUrl := fmt.Sprintf("functions/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace) relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
relativeUrl += fmt.Sprintf("&deploymentraw=1") relativeUrl += "&deploymentraw=1"
resp, err := c.client.Get(relativeUrl) resp, err := c.client.Get(relativeUrl)
if err != nil { if err != nil {
+4 -1
View File
@@ -353,7 +353,10 @@ func getContainerLog(kubernetesClient *kubernetes.Clientset, w http.ResponseWrit
msg := fmt.Sprintf("\n%v\nFunction: %v\nEnvironment: %v\nNamespace: %v\nPod: %v\nContainer: %v\nNode: %v\n%v\n", seq, msg := fmt.Sprintf("\n%v\nFunction: %v\nEnvironment: %v\nNamespace: %v\nPod: %v\nContainer: %v\nNode: %v\n%v\n", seq,
fn.ObjectMeta.Name, fn.Spec.Environment.Name, pod.Namespace, pod.Name, container.Name, pod.Spec.NodeName, seq) fn.ObjectMeta.Name, fn.Spec.Environment.Name, pod.Namespace, pod.Name, container.Name, pod.Spec.NodeName, seq)
w.Write([]byte(msg)) _, err = w.Write([]byte(msg))
if err != nil {
return errors.Wrap(err, "error writing response")
}
_, err = io.Copy(w, podLogs) _, err = io.Copy(w, podLogs)
if err != nil { if err != nil {
+26 -18
View File
@@ -69,7 +69,7 @@ func functionTests(crdClient genInformerCoreV1.CoreV1Interface) {
fi := crdClient.Functions(testNS) fi := crdClient.Functions(testNS)
// cleanup from old crashed tests, ignore errors // cleanup from old crashed tests, ignore errors
fi.Delete(function.ObjectMeta.Name, nil) fi.Delete(function.ObjectMeta.Name, nil) //nolint: errcheck
// create // create
f, err := fi.Create(function) f, err := fi.Create(function)
@@ -117,7 +117,10 @@ func functionTests(crdClient genInformerCoreV1.CoreV1Interface) {
function.ObjectMeta.ResourceVersion = "" function.ObjectMeta.ResourceVersion = ""
f, err = fi.Create(function) f, err = fi.Create(function)
panicIf(err) panicIf(err)
defer fi.Delete(f.ObjectMeta.Name, nil) defer func() {
err := fi.Delete(f.ObjectMeta.Name, nil)
panicIf(err)
}()
// assert that we get a watch event for the new function // assert that we get a watch event for the new function
recvd := false recvd := false
@@ -135,9 +138,7 @@ func functionTests(crdClient genInformerCoreV1.CoreV1Interface) {
log.Panicf("Bad object from watch: %#v", wf) log.Panicf("Bad object from watch: %#v", wf)
} }
log.Printf("watch event took %v", time.Since(start)) log.Printf("watch event took %v", time.Since(start))
recvd = true
} }
} }
func environmentTests(crdClient genInformerCoreV1.CoreV1Interface) { func environmentTests(crdClient genInformerCoreV1.CoreV1Interface) {
@@ -167,7 +168,7 @@ func environmentTests(crdClient genInformerCoreV1.CoreV1Interface) {
ei := crdClient.Environments(testNS) ei := crdClient.Environments(testNS)
// cleanup from old crashed tests, ignore errors // cleanup from old crashed tests, ignore errors
ei.Delete(environment.ObjectMeta.Name, nil) ei.Delete(environment.ObjectMeta.Name, nil) //nolint: errCheck
// create // create
e, err := ei.Create(environment) e, err := ei.Create(environment)
@@ -211,7 +212,10 @@ func environmentTests(crdClient genInformerCoreV1.CoreV1Interface) {
environment.ObjectMeta.ResourceVersion = "" environment.ObjectMeta.ResourceVersion = ""
e, err = ei.Create(environment) e, err = ei.Create(environment)
panicIf(err) panicIf(err)
defer ei.Delete(e.ObjectMeta.Name, nil) defer func() {
err := ei.Delete(e.ObjectMeta.Name, nil)
panicIf(err)
}()
// assert that we get a watch event for the new environment // assert that we get a watch event for the new environment
recvd := false recvd := false
@@ -229,9 +233,7 @@ func environmentTests(crdClient genInformerCoreV1.CoreV1Interface) {
log.Panicf("Bad object from watch: %#v", obj) log.Panicf("Bad object from watch: %#v", obj)
} }
log.Printf("watch event took %v", time.Since(start)) log.Printf("watch event took %v", time.Since(start))
recvd = true
} }
} }
func httpTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) { func httpTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
@@ -259,7 +261,7 @@ func httpTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
ei := crdClient.HTTPTriggers(testNS) ei := crdClient.HTTPTriggers(testNS)
// cleanup from old crashed tests, ignore errors // cleanup from old crashed tests, ignore errors
ei.Delete(httpTrigger.ObjectMeta.Name, nil) ei.Delete(httpTrigger.ObjectMeta.Name, nil) //nolint: errCheck
// create // create
e, err := ei.Create(httpTrigger) e, err := ei.Create(httpTrigger)
@@ -303,7 +305,10 @@ func httpTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
httpTrigger.ObjectMeta.ResourceVersion = "" httpTrigger.ObjectMeta.ResourceVersion = ""
e, err = ei.Create(httpTrigger) e, err = ei.Create(httpTrigger)
panicIf(err) panicIf(err)
defer ei.Delete(e.ObjectMeta.Name, nil) defer func() {
err := ei.Delete(e.ObjectMeta.Name, nil)
panicIf(err)
}()
// assert that we get a watch event for the new httpTrigger // assert that we get a watch event for the new httpTrigger
recvd := false recvd := false
@@ -321,9 +326,7 @@ func httpTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
log.Panicf("Bad object from watch: %#v", obj) log.Panicf("Bad object from watch: %#v", obj)
} }
log.Printf("watch event took %v", time.Since(start)) log.Printf("watch event took %v", time.Since(start))
recvd = true
} }
} }
func kubernetesWatchTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) { func kubernetesWatchTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
@@ -354,7 +357,7 @@ func kubernetesWatchTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
ei := crdClient.KubernetesWatchTriggers(testNS) ei := crdClient.KubernetesWatchTriggers(testNS)
// cleanup from old crashed tests, ignore errors // cleanup from old crashed tests, ignore errors
ei.Delete(kubernetesWatchTrigger.ObjectMeta.Name, nil) ei.Delete(kubernetesWatchTrigger.ObjectMeta.Name, nil) //nolint: errCheck
// create // create
e, err := ei.Create(kubernetesWatchTrigger) e, err := ei.Create(kubernetesWatchTrigger)
@@ -398,7 +401,10 @@ func kubernetesWatchTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
kubernetesWatchTrigger.ObjectMeta.ResourceVersion = "" kubernetesWatchTrigger.ObjectMeta.ResourceVersion = ""
e, err = ei.Create(kubernetesWatchTrigger) e, err = ei.Create(kubernetesWatchTrigger)
panicIf(err) panicIf(err)
defer ei.Delete(e.ObjectMeta.Name, nil) defer func() {
err := ei.Delete(e.ObjectMeta.Name, nil)
panicIf(err)
}()
// assert that we get a watch event for the new kubernetesWatchTrigger // assert that we get a watch event for the new kubernetesWatchTrigger
recvd := false recvd := false
@@ -416,9 +422,7 @@ func kubernetesWatchTriggerTests(crdClient genInformerCoreV1.CoreV1Interface) {
log.Panicf("Bad object from watch: %#v", obj) log.Panicf("Bad object from watch: %#v", obj)
} }
log.Printf("watch event took %v", time.Since(start)) log.Printf("watch event took %v", time.Since(start))
recvd = true
} }
} }
func TestCrd(t *testing.T) { func TestCrd(t *testing.T) {
@@ -442,12 +446,16 @@ func TestCrd(t *testing.T) {
// testNS isolation for running multiple CI builds concurrently. // testNS isolation for running multiple CI builds concurrently.
testNS = uuid.NewV4().String() testNS = uuid.NewV4().String()
kubeClient.CoreV1().Namespaces().Create(&v1.Namespace{ _, err = kubeClient.CoreV1().Namespaces().Create(&v1.Namespace{
ObjectMeta: metav1.ObjectMeta{ ObjectMeta: metav1.ObjectMeta{
Name: testNS, Name: testNS,
}, },
}) })
defer kubeClient.CoreV1().Namespaces().Delete(testNS, nil) panicIf(err)
defer func() {
err := kubeClient.CoreV1().Namespaces().Delete(testNS, nil)
panicIf(err)
}()
// init our types // init our types
err = EnsureFissionCRDs(logger, apiExtClient) err = EnsureFissionCRDs(logger, apiExtClient)
+10 -2
View File
@@ -116,7 +116,10 @@ func writeSecretOrConfigMap(dataMap map[string][]byte, dirPath string) error {
func (fetcher *Fetcher) VersionHandler(w http.ResponseWriter, r *http.Request) { func (fetcher *Fetcher) VersionHandler(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json; charset=utf-8") w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.Write([]byte(info.BuildInfo().String())) _, err := w.Write([]byte(info.BuildInfo().String()))
if err != nil {
fetcher.logger.Error("error writing response", zap.Error(err))
}
} }
func (fetcher *Fetcher) FetchHandler(w http.ResponseWriter, r *http.Request) { func (fetcher *Fetcher) FetchHandler(w http.ResponseWriter, r *http.Request) {
@@ -506,7 +509,12 @@ func (fetcher *Fetcher) UploadHandler(w http.ResponseWriter, r *http.Request) {
fetcher.logger.Info("completed upload request") fetcher.logger.Info("completed upload request")
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
w.Write(rBody) _, err = w.Write(rBody)
if err != nil {
e := "error writing response"
fetcher.logger.Error(e, zap.Error(err))
http.Error(w, fmt.Sprintf("%s: %v", e, err), http.StatusInternalServerError)
}
} }
func (fetcher *Fetcher) rename(src string, dst string) error { func (fetcher *Fetcher) rename(src string, dst string) error {
@@ -95,9 +95,9 @@ func optionalFlags(cmd *cobra.Command, flags ...flag.Flag) {
toCobraFlag(cmd, f, false) toCobraFlag(cmd, f, false)
if f.Deprecated { if f.Deprecated {
usage := fmt.Sprintf("Use --%v instead. The flag still works for now and will be removed in future", f.Substitute) usage := fmt.Sprintf("Use --%v instead. The flag still works for now and will be removed in future", f.Substitute)
cmd.Flags().MarkDeprecated(f.Name, usage) cmd.Flags().MarkDeprecated(f.Name, usage) //nolint: errCheck
} else if f.Hidden { } else if f.Hidden {
cmd.Flags().MarkHidden(f.Name) cmd.Flags().MarkHidden(f.Name) //nolint: errCheck
} }
} }
} }
@@ -105,7 +105,7 @@ func optionalFlags(cmd *cobra.Command, flags ...flag.Flag) {
func requiredFlags(cmd *cobra.Command, flags ...flag.Flag) { func requiredFlags(cmd *cobra.Command, flags ...flag.Flag) {
for _, f := range flags { for _, f := range flags {
toCobraFlag(cmd, f, false) toCobraFlag(cmd, f, false)
cmd.MarkFlagRequired(f.Name) cmd.MarkFlagRequired(f.Name) //nolint: errCheck
} }
} }
+5 -5
View File
@@ -441,7 +441,7 @@ func (fr *FissionResources) Validate(input cli.Input) ([]string, error) {
} }
if len(t.Spec.Host) > 0 { if len(t.Spec.Host) > 0 {
warnings = append(warnings, (fmt.Sprintf("Host in HTTPTrigger spec.Host is now marked as deprecated, see 'help' for details"))) warnings = append(warnings, "Host in HTTPTrigger spec.Host is now marked as deprecated, see 'help' for details")
} }
result = multierror.Append(result, t.Validate()) result = multierror.Append(result, t.Validate())
} }
@@ -475,25 +475,25 @@ func (fr *FissionResources) Validate(input cli.Input) ([]string, error) {
for _, e := range fr.Environments { for _, e := range fr.Environments {
environments[fmt.Sprintf("%s:%s", e.ObjectMeta.Name, e.ObjectMeta.Namespace)] = struct{}{} environments[fmt.Sprintf("%s:%s", e.ObjectMeta.Name, e.ObjectMeta.Namespace)] = struct{}{}
if ((e.Spec.Runtime.Container != nil) && (e.Spec.Runtime.PodSpec != nil)) || ((e.Spec.Builder.Container != nil) && (e.Spec.Builder.PodSpec != nil)) { if ((e.Spec.Runtime.Container != nil) && (e.Spec.Runtime.PodSpec != nil)) || ((e.Spec.Builder.Container != nil) && (e.Spec.Builder.PodSpec != nil)) {
warnings = append(warnings, fmt.Sprintf("You have provided both - container spec and pod spec and while merging the pod spec will take precedence.")) warnings = append(warnings, "You have provided both - container spec and pod spec and while merging the pod spec will take precedence.")
} }
// Unlike CLI can change the environment version silently, // Unlike CLI can change the environment version silently,
// we have to warn the user to modify spec file when this takes place. // we have to warn the user to modify spec file when this takes place.
if e.Spec.Poolsize != 3 && e.Spec.Version < 3 { if e.Spec.Poolsize != 3 && e.Spec.Version < 3 {
warnings = append(warnings, fmt.Sprintf("Poolsize can only be configured when environment version equals to 3, default poolsize 3 will be used for creating environment pool.")) warnings = append(warnings, "Poolsize can only be configured when environment version equals to 3, default poolsize 3 will be used for creating environment pool.")
} }
} }
for _, f := range fr.Functions { for _, f := range fr.Functions {
if _, ok := environments[fmt.Sprintf("%s:%s", f.Spec.Environment.Name, f.Spec.Environment.Namespace)]; !ok { if _, ok := environments[fmt.Sprintf("%s:%s", f.Spec.Environment.Name, f.Spec.Environment.Namespace)]; !ok {
warnings = append(warnings, fmt.Sprintf("Environment %s is referenced in function %s but not declared in specs", f.Spec.Environment.Name, f.ObjectMeta.Name)) warnings = append(warnings, "Environment %s is referenced in function %s but not declared in specs", f.Spec.Environment.Name, f.ObjectMeta.Name)
} }
strategy := f.Spec.InvokeStrategy.ExecutionStrategy strategy := f.Spec.InvokeStrategy.ExecutionStrategy
if strategy.ExecutorType == fv1.ExecutorTypeNewdeploy && strategy.SpecializationTimeout < fv1.DefaultSpecializationTimeOut { if strategy.ExecutorType == fv1.ExecutorTypeNewdeploy && strategy.SpecializationTimeout < fv1.DefaultSpecializationTimeOut {
warnings = append(warnings, fmt.Sprintf("SpecializationTimeout in function spec.InvokeStrategy.ExecutionStrategy should be a value equal to or greater than %v", fv1.DefaultSpecializationTimeOut)) warnings = append(warnings, fmt.Sprintf("SpecializationTimeout in function spec.InvokeStrategy.ExecutionStrategy should be a value equal to or greater than %v", fv1.DefaultSpecializationTimeOut))
} }
if f.Spec.FunctionTimeout <= 0 { if f.Spec.FunctionTimeout <= 0 {
warnings = append(warnings, fmt.Sprintf("FunctionTimeout in function spec should be a field which should have a value greater than 0")) warnings = append(warnings, "FunctionTimeout in function spec should be a field which should have a value greater than 0")
} }
} }
// (ErrorOrNil returns nil if there were no errors appended.) // (ErrorOrNil returns nil if there were no errors appended.)
+2 -2
View File
@@ -111,13 +111,13 @@ func (kw *KubeWatcher) svc() {
// Remove old watches // Remove old watches
for uid, ws := range kw.watches { for uid, ws := range kw.watches {
if _, ok := newWatchUids[uid]; !ok { if _, ok := newWatchUids[uid]; !ok {
kw.removeWatch(&ws.watch) kw.removeWatch(&ws.watch) //nolint: errCheck
} }
} }
// Add new watches // Add new watches
for _, w := range req.watches { for _, w := range req.watches {
if _, ok := kw.watches[w.ObjectMeta.UID]; !ok { if _, ok := kw.watches[w.ObjectMeta.UID]; !ok {
kw.addWatch(&w) kw.addWatch(&w) //nolint: errCheck
} }
} }
req.responseChannel <- &kubeWatcherResponse{error: nil} req.responseChannel <- &kubeWatcherResponse{error: nil}
+4 -1
View File
@@ -51,7 +51,10 @@ func (ws *WatchSync) syncSvc() {
ws.logger.Fatal("failed to get Kubernetes watch trigger list", zap.Error(err)) ws.logger.Fatal("failed to get Kubernetes watch trigger list", zap.Error(err))
} }
ws.kubeWatcher.Sync(watches.Items) err = ws.kubeWatcher.Sync(watches.Items)
if err != nil {
ws.logger.Fatal("failed to sync watches", zap.Error(err))
}
time.Sleep(3 * time.Second) time.Sleep(3 * time.Second)
} }
} }
+7 -1
View File
@@ -170,7 +170,13 @@ func Start() {
if err != nil { if err != nil {
log.Fatalf("can't initialize zap logger: %v", err) log.Fatalf("can't initialize zap logger: %v", err)
} }
defer zapLogger.Sync()
defer func() {
err := zapLogger.Sync()
if err != nil {
log.Fatalf("failed to sync zap logger: %v", err)
}
}()
if _, err := os.Stat(fissionSymlinkPath); os.IsNotExist(err) { if _, err := os.Stat(fissionSymlinkPath); os.IsNotExist(err) {
zapLogger.Info("symlink path not exist, create it", zapLogger.Info("symlink path not exist, create it",
@@ -319,7 +319,7 @@ func pollAzureQueueSubscription(conn AzureStorageConnection, sub *AzureQueueSubs
} }
func invokeTriggeredFunction(conn AzureStorageConnection, sub *AzureQueueSubscription, message AzureMessage) { func invokeTriggeredFunction(conn AzureStorageConnection, sub *AzureQueueSubscription, message AzureMessage) {
defer message.Delete(nil) defer message.Delete(nil) //nolint: errCheck
conn.logger.Info("making HTTP request to invoke function", zap.String("function_url", sub.functionURL)) conn.logger.Info("making HTTP request to invoke function", zap.String("function_url", sub.functionURL))
@@ -329,7 +329,7 @@ func TestAzureStorageQueuePoisonMessage(t *testing.T) {
require.NoError(t, err) require.NoError(t, err)
require.NotNil(t, subscription) require.NotNil(t, subscription)
connection.Unsubscribe(subscription) panicIf(connection.Unsubscribe(subscription))
mock.AssertExpectationsForObjects(t, httpClient, message, poisonMessage, queue, poisonQueue, service) mock.AssertExpectationsForObjects(t, httpClient, message, poisonMessage, queue, poisonQueue, service)
} }
@@ -480,7 +480,7 @@ func runAzureStorageQueueTest(t *testing.T, count int, output bool) {
require.NoError(t, err) require.NoError(t, err)
require.NotNil(t, subscription) require.NotNil(t, subscription)
connection.Unsubscribe(subscription) panicIf(connection.Unsubscribe(subscription))
mock.AssertExpectationsForObjects(t, httpClient, message, outputMessage, queue, outputQueue, service) mock.AssertExpectationsForObjects(t, httpClient, message, outputMessage, queue, outputQueue, service)
} }
+2 -2
View File
@@ -30,7 +30,7 @@ func TestPoolCache(t *testing.T) {
log.Panicf("expected 2 items") log.Panicf("expected 2 items")
} }
c.DeleteValue("func2", "ip2") checkErr(c.DeleteValue("func2", "ip2"))
c.MarkAvailable("func", "ip") c.MarkAvailable("func", "ip")
cc = c.ListAvailableValue() cc = c.ListAvailableValue()
@@ -41,7 +41,7 @@ func TestPoolCache(t *testing.T) {
_, err := c.GetValue("func") _, err := c.GetValue("func")
checkErr(err) checkErr(err)
c.DeleteValue("func", "ip") checkErr(c.DeleteValue("func", "ip"))
_, err = c.GetValue("func") _, err = c.GetValue("func")
if err == nil { if err == nil {
+6 -1
View File
@@ -126,7 +126,12 @@ func TestS3StorageService(t *testing.T) {
log.Fatalf("Could not connect to docker: %s", err) log.Fatalf("Could not connect to docker: %s", err)
} }
defer pool.Purge(resource) defer func() {
err := pool.Purge(resource)
if err != nil {
log.Fatal(err)
}
}()
// Start storagesvc // Start storagesvc
bucketName := "test-s3-service" bucketName := "test-s3-service"
+13 -2
View File
@@ -65,7 +65,10 @@ func getStorageLocation(config *storageConfig) (stow.Location, error) {
// Handle multipart file uploads. // Handle multipart file uploads.
func (ss *StorageService) uploadHandler(w http.ResponseWriter, r *http.Request) { func (ss *StorageService) uploadHandler(w http.ResponseWriter, r *http.Request) {
// handle upload // handle upload
r.ParseMultipartForm(0) err := r.ParseMultipartForm(0)
if err != nil {
http.Error(w, "failed to parse request", http.StatusBadRequest)
}
file, handler, err := r.FormFile("uploadfile") file, handler, err := r.FormFile("uploadfile")
if err != nil { if err != nil {
http.Error(w, "missing upload file", http.StatusBadRequest) http.Error(w, "missing upload file", http.StatusBadRequest)
@@ -121,7 +124,15 @@ func (ss *StorageService) uploadHandler(w http.ResponseWriter, r *http.Request)
http.Error(w, "Error marshaling response", http.StatusInternalServerError) http.Error(w, "Error marshaling response", http.StatusInternalServerError)
return return
} }
w.Write(resp) _, err = w.Write(resp)
if err != nil {
ss.logger.Error(
"error writing HTTP response",
zap.Error(err),
zap.String("filename", handler.Filename),
)
}
} }
func (ss *StorageService) getIdFromRequest(r *http.Request) (string, error) { func (ss *StorageService) getIdFromRequest(r *http.Request) (string, error) {
+1 -2
View File
@@ -171,8 +171,7 @@ func (client *StowClient) getItemIDsWithFilter(filterFunc filter, filterFuncPara
for { for {
items, cursor, err = client.container.Items(stow.NoPrefix, cursor, PaginationSize) items, cursor, err = client.container.Items(stow.NoPrefix, cursor, PaginationSize)
if err != nil { if err != nil {
errors.Wrap(err, "error getting items from container") return nil, errors.Wrap(err, "error getting items from container")
return nil, err
} }
for _, item := range items { for _, item := range items {
+1 -1
View File
@@ -127,7 +127,7 @@ func (timer *Timer) syncCron(triggers []fv1.TimeTrigger) error {
func (timer *Timer) newCron(t fv1.TimeTrigger) *cron.Cron { func (timer *Timer) newCron(t fv1.TimeTrigger) *cron.Cron {
c := cron.New() c := cron.New()
c.AddFunc(t.Spec.Cron, func() { c.AddFunc(t.Spec.Cron, func() { //nolint: errCheck
headers := map[string]string{ headers := map[string]string{
"X-Fission-Timer-Name": t.ObjectMeta.Name, "X-Fission-Timer-Name": t.ObjectMeta.Name,
} }
+1 -1
View File
@@ -55,7 +55,7 @@ func (ws *TimerSync) syncSvc() {
} }
ws.logger.Fatal("failed to get time trigger list", zap.Error(err)) ws.logger.Fatal("failed to get time trigger list", zap.Error(err))
} }
ws.timer.Sync(triggers.Items) ws.timer.Sync(triggers.Items) //nolint: errCheck
// TODO switch to watches // TODO switch to watches
time.Sleep(3 * time.Second) time.Sleep(3 * time.Second)