diff --git a/cmd/fission-bundle/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go index 50f65b29..0ea06c66 100644 --- a/cmd/fission-bundle/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -41,7 +41,7 @@ func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { return errors.Wrap(err, "failed to get fission or kubernetes client") } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { return errors.Wrap(err, "error waiting for CRDs") } diff --git a/cmd/reporter/app/cmd_event.go b/cmd/reporter/app/cmd_event.go index 21170616..436f2346 100644 --- a/cmd/reporter/app/cmd_event.go +++ b/cmd/reporter/app/cmd_event.go @@ -16,7 +16,6 @@ limitations under the License. package app import ( - "context" "log" "github.com/spf13/cobra" @@ -50,13 +49,11 @@ func eventCommandHandler(cmd *cobra.Command, args []string) error { return err } - ctx := context.Background() - t, err := tracker.NewTracker() if err != nil { return err } - return t.SendEvent(ctx, event) + return t.SendEvent(cmd.Context(), event) } // EventCommand reports an event to analytics diff --git a/pkg/buildermgr/buildermgr.go b/pkg/buildermgr/buildermgr.go index b86ab5c2..c5642b29 100644 --- a/pkg/buildermgr/buildermgr.go +++ b/pkg/buildermgr/buildermgr.go @@ -42,7 +42,7 @@ func Start(ctx context.Context, logger *zap.Logger, storageSvcUrl string, envBui return errors.Wrap(err, "failed to get fission or kubernetes client") } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { return errors.Wrap(err, "error waiting for CRDs") } diff --git a/pkg/controller/client/rest/client.go b/pkg/controller/client/rest/client.go index 849d0776..a74f9417 100644 --- a/pkg/controller/client/rest/client.go +++ b/pkg/controller/client/rest/client.go @@ -24,7 +24,7 @@ import ( "strings" "github.com/pkg/errors" - "golang.org/x/net/context/ctxhttp" + "go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp" ) type ( @@ -39,13 +39,17 @@ type ( } RESTClient struct { - url string + url string + HTTPClient *http.Client } ) func NewRESTClient(serverUrl string) Interface { return &RESTClient{ url: strings.TrimSuffix(serverUrl, "/"), + HTTPClient: &http.Client{ + Transport: otelhttp.NewTransport(http.DefaultTransport), + }, } } @@ -54,7 +58,7 @@ func (c *RESTClient) Create(relativeUrl string, contentType string, payload []by if len(payload) > 0 { reader = bytes.NewReader(payload) } - return c.sendRequest(http.MethodPost, c.v2CrdUrl(relativeUrl), map[string]string{"Content-type": contentType}, reader) + return c.sendRequest(context.TODO(), http.MethodPost, c.v2CrdUrl(relativeUrl), map[string]string{"Content-type": contentType}, reader) } func (c *RESTClient) Put(relativeUrl string, contentType string, payload []byte) (*http.Response, error) { @@ -62,15 +66,15 @@ func (c *RESTClient) Put(relativeUrl string, contentType string, payload []byte) if len(payload) > 0 { reader = bytes.NewReader(payload) } - return c.sendRequest(http.MethodPut, c.v2CrdUrl(relativeUrl), map[string]string{"Content-type": contentType}, reader) + return c.sendRequest(context.TODO(), http.MethodPut, c.v2CrdUrl(relativeUrl), map[string]string{"Content-type": contentType}, reader) } func (c *RESTClient) Get(relativeUrl string) (*http.Response, error) { - return c.sendRequest(http.MethodGet, c.v2CrdUrl(relativeUrl), nil, nil) + return c.sendRequest(context.TODO(), http.MethodGet, c.v2CrdUrl(relativeUrl), nil, nil) } func (c *RESTClient) Delete(relativeUrl string) error { - resp, err := c.sendRequest(http.MethodDelete, c.v2CrdUrl(relativeUrl), nil, nil) + resp, err := c.sendRequest(context.TODO(), http.MethodDelete, c.v2CrdUrl(relativeUrl), nil, nil) if err != nil { return err } @@ -93,27 +97,26 @@ func (c *RESTClient) Proxy(method string, relativeUrl string, payload []byte) (* if len(payload) > 0 { reader = bytes.NewReader(payload) } - return c.sendRequest(method, c.proxyUrl(relativeUrl), nil, reader) + return c.sendRequest(context.TODO(), method, c.proxyUrl(relativeUrl), nil, reader) } func (c *RESTClient) ServerInfo() (*http.Response, error) { - return c.sendRequest(http.MethodGet, c.url, nil, nil) + return c.sendRequest(context.TODO(), http.MethodGet, c.url, nil, nil) } func (c *RESTClient) ServerURL() string { return c.url } -func (c *RESTClient) sendRequest(method string, relativeUrl string, headers map[string]string, reader io.Reader) (*http.Response, error) { - req, err := http.NewRequest(method, relativeUrl, reader) +func (c *RESTClient) sendRequest(ctx context.Context, method string, relativeUrl string, headers map[string]string, reader io.Reader) (*http.Response, error) { + req, err := http.NewRequestWithContext(ctx, method, relativeUrl, reader) if err != nil { return nil, err } for k, v := range headers { req.Header.Set(k, v) } - // TODO: accept context - return ctxhttp.Do(context.Background(), &http.Client{}, req) + return c.HTTPClient.Do(req) } func (c *RESTClient) v2CrdUrl(relativeUrl string) string { diff --git a/pkg/controller/controller.go b/pkg/controller/controller.go index 0bf82387..9dd3ac0c 100644 --- a/pkg/controller/controller.go +++ b/pkg/controller/controller.go @@ -32,12 +32,12 @@ func Start(ctx context.Context, logger *zap.Logger, port int, unitTestFlag bool) cLogger.Fatal("failed to connect to k8s API", zap.Error(err)) } - err = crd.EnsureFissionCRDs(cLogger, apiExtClient) + err = crd.EnsureFissionCRDs(ctx, cLogger, apiExtClient) if err != nil { cLogger.Fatal("failed to find fission CRDs", zap.Error(err)) } - err = crd.WaitForCRDs(fc) + err = crd.WaitForCRDs(ctx, fc) if err != nil { cLogger.Fatal("error waiting for CRDs", zap.Error(err)) } diff --git a/pkg/crd/client.go b/pkg/crd/client.go index 49b38ecd..76a2890c 100644 --- a/pkg/crd/client.go +++ b/pkg/crd/client.go @@ -88,11 +88,11 @@ func MakeFissionClient() (versioned.Interface, kubernetes.Interface, apiextensio } // WaitForCRDs does a timeout to check if CRDs have been installed -func WaitForCRDs(fissionClient versioned.Interface) error { +func WaitForCRDs(ctx context.Context, fissionClient versioned.Interface) error { start := time.Now() for { fi := fissionClient.CoreV1().Functions(metav1.NamespaceDefault) - _, err := fi.List(context.TODO(), metav1.ListOptions{}) + _, err := fi.List(ctx, metav1.ListOptions{}) if err != nil { time.Sleep(100 * time.Millisecond) } else { diff --git a/pkg/crd/crd.go b/pkg/crd/crd.go index ac9b1829..9c712062 100644 --- a/pkg/crd/crd.go +++ b/pkg/crd/crd.go @@ -27,7 +27,7 @@ import ( ) // EnsureFissionCRDs checks if all Fission CRDs are present -func EnsureFissionCRDs(logger *zap.Logger, clientset apiextensionsclient.Interface) error { +func EnsureFissionCRDs(ctx context.Context, logger *zap.Logger, clientset apiextensionsclient.Interface) error { crdsExpected := []string{ "canaryconfigs.fission.io", "environments.fission.io", @@ -40,7 +40,7 @@ func EnsureFissionCRDs(logger *zap.Logger, clientset apiextensionsclient.Interfa } errs := &multierror.Error{} for _, crdName := range crdsExpected { - crd, err := clientset.ApiextensionsV1().CustomResourceDefinitions().Get(context.TODO(), crdName, metav1.GetOptions{}) + crd, err := clientset.ApiextensionsV1().CustomResourceDefinitions().Get(ctx, crdName, metav1.GetOptions{}) if err != nil { errs = multierror.Append(errs, fmt.Errorf("CRD %s not found: %s", crdName, err)) } diff --git a/pkg/executor/executor.go b/pkg/executor/executor.go index b4d3a946..d17bd0df 100644 --- a/pkg/executor/executor.go +++ b/pkg/executor/executor.go @@ -259,7 +259,7 @@ func StartExecutor(ctx context.Context, logger *zap.Logger, functionNamespace st return errors.Wrap(err, "failed to get kubernetes client") } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { return errors.Wrap(err, "error waiting for CRDs") } diff --git a/pkg/executor/executor_test.go b/pkg/executor/executor_test.go index 1ea60938..cd74296e 100644 --- a/pkg/executor/executor_test.go +++ b/pkg/executor/executor_test.go @@ -143,12 +143,12 @@ func TestExecutor(t *testing.T) { panicIf(err) // make sure CRD types exist on cluster - err = crd.EnsureFissionCRDs(logger, apiExtClient) + err = crd.EnsureFissionCRDs(context.TODO(), logger, apiExtClient) if err != nil { log.Panicf("failed to ensure crds: %v", err) } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(context.TODO(), fissionClient) if err != nil { log.Panicf("failed to wait crds: %v", err) } diff --git a/pkg/fission-cli/cliwrapper/cli/cli.go b/pkg/fission-cli/cliwrapper/cli/cli.go index 4152f024..bd4009ae 100644 --- a/pkg/fission-cli/cliwrapper/cli/cli.go +++ b/pkg/fission-cli/cliwrapper/cli/cli.go @@ -17,11 +17,14 @@ limitations under the License. package cli import ( + "context" "time" ) type ( Input interface { + Context() context.Context + //Parse(input interface{}) error // IsSet checks whether a flag has been set by the user diff --git a/pkg/fission-cli/cliwrapper/driver/cobra/cobra.go b/pkg/fission-cli/cliwrapper/driver/cobra/cobra.go index 474cadbf..09a1094f 100644 --- a/pkg/fission-cli/cliwrapper/driver/cobra/cobra.go +++ b/pkg/fission-cli/cliwrapper/driver/cobra/cobra.go @@ -17,6 +17,7 @@ limitations under the License. package cobra import ( + "context" "fmt" "strings" "time" @@ -217,6 +218,10 @@ func WrapperChain(actions ...cmd.CommandAction) func(*cobra.Command, []string) e } } +func (u Cli) Context() context.Context { + return u.c.Context() +} + func (u Cli) IsSet(key string) bool { return u.c.Flags().Changed(key) } diff --git a/pkg/fission-cli/cliwrapper/driver/dummy/dummy.go b/pkg/fission-cli/cliwrapper/driver/dummy/dummy.go index e3b0f6df..c9e3807f 100644 --- a/pkg/fission-cli/cliwrapper/driver/dummy/dummy.go +++ b/pkg/fission-cli/cliwrapper/driver/dummy/dummy.go @@ -17,6 +17,7 @@ limitations under the License. package dummy import ( + "context" "time" fCli "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" @@ -33,6 +34,10 @@ func TestFlagSet() Cli { return Cli{c: make(map[string]interface{})} } +func (u Cli) Context() context.Context { + return context.TODO() +} + // Set allows to set any kinds of value with given key. // The type of set value should be matched with the returned // type of GetXXX function. diff --git a/pkg/fission-cli/cmd/archive/delete.go b/pkg/fission-cli/cmd/archive/delete.go index d6d28bd5..dc195a5e 100644 --- a/pkg/fission-cli/cmd/archive/delete.go +++ b/pkg/fission-cli/cmd/archive/delete.go @@ -17,7 +17,6 @@ limitations under the License. package archive import ( - "context" "fmt" "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" @@ -40,14 +39,14 @@ func (opts *DeleteSubCommand) do(input cli.Input) error { kubeContext := input.String(flagkey.KubeContext) archiveID := input.String(flagkey.ArchiveID) - storagesvcURL, err := util.GetStorageURL(kubeContext) + storagesvcURL, err := util.GetStorageURL(input.Context(), kubeContext) if err != nil { return err } client := storagesvcClient.MakeClient(storagesvcURL.String()) - err = client.Delete(context.Background(), archiveID) + err = client.Delete(input.Context(), archiveID) if err != nil { return err } diff --git a/pkg/fission-cli/cmd/archive/download.go b/pkg/fission-cli/cmd/archive/download.go index 974f52cd..e64e3abf 100644 --- a/pkg/fission-cli/cmd/archive/download.go +++ b/pkg/fission-cli/cmd/archive/download.go @@ -17,7 +17,6 @@ limitations under the License. package archive import ( - "context" "fmt" "strings" @@ -46,13 +45,13 @@ func (opts *DownloadSubCommand) do(input cli.Input) error { archiveOutput = strings.TrimPrefix(archiveID, "/fission/fission-functions/") } - storageAccessURL, err := util.GetStorageURL(kubeContext) + storageAccessURL, err := util.GetStorageURL(input.Context(), kubeContext) if err != nil { return err } client := storagesvcClient.MakeClient(storageAccessURL.String()) - err = client.Download(context.Background(), archiveID, archiveOutput) + err = client.Download(input.Context(), archiveID, archiveOutput) if err != nil { return err } diff --git a/pkg/fission-cli/cmd/archive/geturl.go b/pkg/fission-cli/cmd/archive/geturl.go index c993ab13..75c9acf8 100644 --- a/pkg/fission-cli/cmd/archive/geturl.go +++ b/pkg/fission-cli/cmd/archive/geturl.go @@ -41,7 +41,7 @@ func (opts *GetURLSubCommand) do(input cli.Input) error { kubeContext := input.String(flagkey.KubeContext) archiveID := input.String(flagkey.ArchiveID) - serverURL, err := util.GetStorageURL(kubeContext) + serverURL, err := util.GetStorageURL(input.Context(), kubeContext) if err != nil { return err } @@ -61,7 +61,7 @@ func (opts *GetURLSubCommand) do(input cli.Input) error { defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - return fmt.Errorf("Error getting URL. Exited with Status: %s", resp.Status) + return fmt.Errorf("error getting URL. Exited with Status: %s", resp.Status) } storageType := resp.Header.Get("X-FISSION-STORAGETYPE") diff --git a/pkg/fission-cli/cmd/archive/list.go b/pkg/fission-cli/cmd/archive/list.go index 109cce7e..4c1e18dc 100644 --- a/pkg/fission-cli/cmd/archive/list.go +++ b/pkg/fission-cli/cmd/archive/list.go @@ -17,7 +17,6 @@ limitations under the License. package archive import ( - "context" "fmt" "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" @@ -39,13 +38,13 @@ func (opts *ListSubCommand) do(input cli.Input) error { kubeContext := input.String(flagkey.KubeContext) - storageAccessURL, err := util.GetStorageURL(kubeContext) + storageAccessURL, err := util.GetStorageURL(input.Context(), kubeContext) if err != nil { return err } client := storagesvcClient.MakeClient(storageAccessURL.String()) - files, err := client.List(context.Background()) + files, err := client.List(input.Context()) if err != nil { return err } diff --git a/pkg/fission-cli/cmd/archive/upload.go b/pkg/fission-cli/cmd/archive/upload.go index 782085b7..db231b57 100644 --- a/pkg/fission-cli/cmd/archive/upload.go +++ b/pkg/fission-cli/cmd/archive/upload.go @@ -17,7 +17,6 @@ limitations under the License. package archive import ( - "context" "fmt" "github.com/fission/fission/pkg/fission-cli/cliwrapper/cli" @@ -40,13 +39,13 @@ func (opts *UploadSubCommand) do(input cli.Input) error { kubeContext := input.String(flagkey.KubeContext) archiveName := input.String(flagkey.ArchiveName) - storagesvcURL, err := util.GetStorageURL(kubeContext) + storagesvcURL, err := util.GetStorageURL(input.Context(), kubeContext) if err != nil { return err } client := storagesvcClient.MakeClient(storagesvcURL.String()) - archiveID, err := client.Upload(context.Background(), archiveName, nil) + archiveID, err := client.Upload(input.Context(), archiveName, nil) if err != nil { return err } diff --git a/pkg/fission-cli/cmd/function/log.go b/pkg/fission-cli/cmd/function/log.go index 73754761..0c41cd74 100644 --- a/pkg/fission-cli/cmd/function/log.go +++ b/pkg/fission-cli/cmd/function/log.go @@ -59,7 +59,7 @@ func (opts *LogSubCommand) do(input cli.Input) error { return errors.Wrap(err, "error getting function") } - server, err := util.GetApplicationUrl("application=fission-api", kubeContext) + server, err := util.GetApplicationUrl(input.Context(), "application=fission-api", kubeContext) if err != nil { return err } @@ -72,7 +72,7 @@ func (opts *LogSubCommand) do(input cli.Input) error { requestChan := make(chan struct{}) responseChan := make(chan struct{}) - ctx := context.Background() + ctx := input.Context() go func(ctx context.Context, requestChan, responseChan chan struct{}) { t := time.Unix(0, 0*int64(time.Millisecond)) diff --git a/pkg/fission-cli/cmd/function/test.go b/pkg/fission-cli/cmd/function/test.go index 614cb2df..ed7bb8e6 100644 --- a/pkg/fission-cli/cmd/function/test.go +++ b/pkg/fission-cli/cmd/function/test.go @@ -62,7 +62,7 @@ func (opts *TestSubCommand) do(input cli.Input) error { } // Portforward to the fission router - localRouterPort, err := util.SetupPortForward(util.GetFissionNamespace(), "application=fission-router", kubeContext) + localRouterPort, err := util.SetupPortForward(input.Context(), util.GetFissionNamespace(), "application=fission-router", kubeContext) if err != nil { return err } @@ -107,10 +107,10 @@ func (opts *TestSubCommand) do(input cli.Input) error { testTimeout := input.Duration(flagkey.FnTestTimeout) if testTimeout <= 0*time.Second { - ctx = context.Background() + ctx = input.Context() } else { var closeCtx context.CancelFunc - ctx, closeCtx = context.WithTimeout(context.Background(), input.Duration(flagkey.FnTestTimeout)) + ctx, closeCtx = context.WithTimeout(input.Context(), input.Duration(flagkey.FnTestTimeout)) defer closeCtx() } diff --git a/pkg/fission-cli/cmd/package/package.go b/pkg/fission-cli/cmd/package/package.go index f402fd5c..83f2a1aa 100644 --- a/pkg/fission-cli/cmd/package/package.go +++ b/pkg/fission-cli/cmd/package/package.go @@ -17,7 +17,6 @@ limitations under the License. package _package import ( - "context" "fmt" "net/http" "os" @@ -128,7 +127,7 @@ func CreateArchive(client client.Interface, input cli.Input, includeFiles []stri return nil, err } file := filepath.Join(tmpDir, id.String()) - err = utils.DownloadUrl(context.Background(), http.DefaultClient, fileURL, file) + err = utils.DownloadUrl(input.Context(), http.DefaultClient, fileURL, file) if err != nil { return nil, errors.Wrap(err, "error downloading file from the given URL") } @@ -193,8 +192,7 @@ func CreateArchive(client client.Interface, input cli.Input, includeFiles []stri return nil, err } - ctx := context.Background() - return pkgutil.UploadArchiveFile(ctx, client, archivePath) + return pkgutil.UploadArchiveFile(input.Context(), client, archivePath) } // makeArchiveFile creates a zip file from the given list of input files, diff --git a/pkg/fission-cli/cmd/spec/apply.go b/pkg/fission-cli/cmd/spec/apply.go index 9427c1fd..fb926e9d 100644 --- a/pkg/fission-cli/cmd/spec/apply.go +++ b/pkg/fission-cli/cmd/spec/apply.go @@ -130,7 +130,7 @@ func (opts *ApplySubCommand) run(input cli.Input) error { } // make changes to the cluster based on the specs - pkgMetas, as, err := applyResources(opts.Client(), specDir, fr, deleteResources, input.Bool(flagkey.SpecAllowConflicts)) + pkgMetas, as, err := applyResources(input.Context(), opts.Client(), specDir, fr, deleteResources, input.Bool(flagkey.SpecAllowConflicts)) if err != nil { return errors.Wrap(err, "error applying specs") } @@ -141,7 +141,7 @@ func (opts *ApplySubCommand) run(input cli.Input) error { pbw.addPackages(pkgMetas) } - ctx, pkgWatchCancel := context.WithCancel(context.Background()) + ctx, pkgWatchCancel := context.WithCancel(input.Context()) if watchResources { // if we're watching for files, we don't need to wait for builds to complete @@ -283,7 +283,7 @@ func pluralize(num int, word string) string { } // applyArchives figures out the set of archives that need to be uploaded, and uploads them. -func applyArchives(fclient client.Interface, specDir string, fr *FissionResources) error { +func applyArchives(ctx context.Context, fclient client.Interface, specDir string, fr *FissionResources) error { // archive:// URL -> archive map. archiveFiles := make(map[string]fv1.Archive) @@ -329,7 +329,6 @@ func applyArchives(fclient client.Interface, specDir string, fr *FissionResource // doesn't exist, upload fmt.Printf("uploading archive %v\n", name) // ar.URL is actually a local filename at this stage - ctx := context.Background() uploadedAr, err := pkgutil.UploadArchiveFile(ctx, fclient, ar.URL) if err != nil { return err @@ -357,12 +356,12 @@ func applyArchives(fclient client.Interface, specDir string, fr *FissionResource } // applyResources applies the given set of fission resources. -func applyResources(fclient client.Interface, specDir string, fr *FissionResources, delete bool, specAllowConflicts bool) (map[string]metav1.ObjectMeta, map[string]ResourceApplyStatus, error) { +func applyResources(ctx context.Context, fclient client.Interface, specDir string, fr *FissionResources, delete bool, specAllowConflicts bool) (map[string]metav1.ObjectMeta, map[string]ResourceApplyStatus, error) { applyStatus := make(map[string]ResourceApplyStatus) // upload archives that need to be uploaded. Changes archive references in fr.Packages. - err := applyArchives(fclient, specDir, fr) + err := applyArchives(ctx, fclient, specDir, fr) if err != nil { return nil, nil, err } diff --git a/pkg/fission-cli/cmd/support/dump.go b/pkg/fission-cli/cmd/support/dump.go index 9b10a78d..8366244d 100644 --- a/pkg/fission-cli/cmd/support/dump.go +++ b/pkg/fission-cli/cmd/support/dump.go @@ -137,7 +137,7 @@ func (opts *DumpSubCommand) do(input cli.Input) error { wg.Add(1) go func(res resources.Resource, dir string) { defer wg.Done() - res.Dump(dir) + res.Dump(input.Context(), dir) }(res, dir) } diff --git a/pkg/fission-cli/cmd/support/resources/crd.go b/pkg/fission-cli/cmd/support/resources/crd.go index bbf611fc..18e9ee34 100644 --- a/pkg/fission-cli/cmd/support/resources/crd.go +++ b/pkg/fission-cli/cmd/support/resources/crd.go @@ -17,6 +17,7 @@ limitations under the License. package resources import ( + "context" "fmt" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -46,7 +47,7 @@ func NewCrdDumper(client client.Interface, crdType string) Resource { return CrdDumper{client: client, crdType: crdType} } -func (res CrdDumper) Dump(dumpDir string) { +func (res CrdDumper) Dump(ctx context.Context, dumpDir string) { switch res.crdType { case CrdEnvironment: diff --git a/pkg/fission-cli/cmd/support/resources/fissionversion.go b/pkg/fission-cli/cmd/support/resources/fissionversion.go index ede1145b..24c4a844 100644 --- a/pkg/fission-cli/cmd/support/resources/fissionversion.go +++ b/pkg/fission-cli/cmd/support/resources/fissionversion.go @@ -17,6 +17,7 @@ limitations under the License. package resources import ( + "context" "fmt" "path/filepath" @@ -32,7 +33,7 @@ func NewFissionVersion(client client.Interface) Resource { return FissionVersion{client: client} } -func (res FissionVersion) Dump(dumpDir string) { +func (res FissionVersion) Dump(ctx context.Context, dumpDir string) { ver := util.GetVersion(res.client) file := filepath.Clean(fmt.Sprintf("%v/%v", dumpDir, "fission-version.txt")) writeToFile(file, ver) diff --git a/pkg/fission-cli/cmd/support/resources/kubernetes.go b/pkg/fission-cli/cmd/support/resources/kubernetes.go index 381b31bc..eb392449 100644 --- a/pkg/fission-cli/cmd/support/resources/kubernetes.go +++ b/pkg/fission-cli/cmd/support/resources/kubernetes.go @@ -50,7 +50,7 @@ func NewKubernetesVersion(clientset kubernetes.Interface) Resource { return KubernetesVersion{client: clientset} } -func (res KubernetesVersion) Dump(dumpDir string) { +func (res KubernetesVersion) Dump(ctx context.Context, dumpDir string) { serverVer, err := res.client.Discovery().ServerVersion() if err != nil { console.Error(fmt.Sprintf("Error setting up kubernetes client: %v", err)) @@ -76,10 +76,10 @@ func NewKubernetesObjectDumper(clientset kubernetes.Interface, objType string, s } } -func (res KubernetesObjectDumper) Dump(dumpDir string) { +func (res KubernetesObjectDumper) Dump(ctx context.Context, dumpDir string) { switch res.objType { case KubernetesService: - objs, err := res.client.CoreV1().Services(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector}) + objs, err := res.client.CoreV1().Services(metav1.NamespaceAll).List(ctx, metav1.ListOptions{LabelSelector: res.selector}) if err != nil { console.Error(fmt.Sprintf("Error getting %v list with selector %v: %v", res.objType, res.selector, err)) return @@ -92,7 +92,7 @@ func (res KubernetesObjectDumper) Dump(dumpDir string) { } case KubernetesDeployment: - objs, err := res.client.AppsV1().Deployments(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector}) + objs, err := res.client.AppsV1().Deployments(metav1.NamespaceAll).List(ctx, metav1.ListOptions{LabelSelector: res.selector}) if err != nil { console.Error(fmt.Sprintf("Error getting %v list with selector %v: %v", res.objType, res.selector, err)) return @@ -104,7 +104,7 @@ func (res KubernetesObjectDumper) Dump(dumpDir string) { } case KubernetesPod: - objs, err := res.client.CoreV1().Pods(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector}) + objs, err := res.client.CoreV1().Pods(metav1.NamespaceAll).List(ctx, metav1.ListOptions{LabelSelector: res.selector}) if err != nil { console.Error(fmt.Sprintf("Error getting %v list with selector %v: %v", res.objType, res.selector, err)) return @@ -116,7 +116,7 @@ func (res KubernetesObjectDumper) Dump(dumpDir string) { } case KubernetesHPA: - objs, err := res.client.AutoscalingV2beta2().HorizontalPodAutoscalers(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector}) + objs, err := res.client.AutoscalingV2beta2().HorizontalPodAutoscalers(metav1.NamespaceAll).List(ctx, metav1.ListOptions{LabelSelector: res.selector}) if err != nil { console.Error(fmt.Sprintf("Error getting %v list with selector %v: %v", res.objType, res.selector, err)) return @@ -128,7 +128,7 @@ func (res KubernetesObjectDumper) Dump(dumpDir string) { } case KubernetesDaemonSet: - objs, err := res.client.AppsV1().DaemonSets(metav1.NamespaceAll).List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector}) + objs, err := res.client.AppsV1().DaemonSets(metav1.NamespaceAll).List(ctx, metav1.ListOptions{LabelSelector: res.selector}) if err != nil { console.Error(fmt.Sprintf("Error getting %v list with selector %v: %v", res.objType, res.selector, err)) return @@ -140,7 +140,7 @@ func (res KubernetesObjectDumper) Dump(dumpDir string) { } case KubernetesNode: - objs, err := res.client.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{LabelSelector: res.selector}) + objs, err := res.client.CoreV1().Nodes().List(ctx, metav1.ListOptions{LabelSelector: res.selector}) if err != nil { console.Error(fmt.Sprintf("Error getting %v list with selector %v: %v", res.objType, res.selector, err)) return @@ -195,10 +195,10 @@ func NewKubernetesPodLogDumper(clientset kubernetes.Interface, selector string) } } -func (res KubernetesPodLogDumper) Dump(dumpDir string) { +func (res KubernetesPodLogDumper) Dump(ctx context.Context, dumpDir string) { l, err := res.client.CoreV1(). Pods(metav1.NamespaceAll). - List(context.TODO(), metav1.ListOptions{LabelSelector: res.labelSelector}) + List(ctx, metav1.ListOptions{LabelSelector: res.labelSelector}) if err != nil { console.Error(fmt.Sprintf("Error getting controller list: %v", err)) return @@ -217,7 +217,7 @@ func (res KubernetesPodLogDumper) Dump(dumpDir string) { req := res.client.CoreV1().Pods(pod.Namespace). GetLogs(pod.Name, &corev1.PodLogOptions{Container: container.Name}) - stream, err := req.Stream(context.Background()) + stream, err := req.Stream(ctx) if err != nil { console.Error(fmt.Sprintf("Error streaming logs for pod %v: %v", pod.Name, err)) return diff --git a/pkg/fission-cli/cmd/support/resources/resource.go b/pkg/fission-cli/cmd/support/resources/resource.go index 621f94c9..09117860 100644 --- a/pkg/fission-cli/cmd/support/resources/resource.go +++ b/pkg/fission-cli/cmd/support/resources/resource.go @@ -17,6 +17,7 @@ limitations under the License. package resources import ( + "context" "fmt" "os" "path/filepath" @@ -29,7 +30,7 @@ import ( ) type Resource interface { - Dump(string) + Dump(context.Context, string) } func getFileName(dumpdir string, meta metav1.ObjectMeta) string { diff --git a/pkg/fission-cli/cmd/token/create.go b/pkg/fission-cli/cmd/token/create.go index 64d6d2b9..669cb77f 100644 --- a/pkg/fission-cli/cmd/token/create.go +++ b/pkg/fission-cli/cmd/token/create.go @@ -66,7 +66,7 @@ func (opts *CreateSubCommand) run(input cli.Input) error { kubeContext := input.String(flagkey.KubeContext) // Portforward to the fission router - localRouterPort, err := util.SetupPortForward(util.GetFissionNamespace(), "application=fission-router", kubeContext) + localRouterPort, err := util.SetupPortForward(input.Context(), util.GetFissionNamespace(), "application=fission-router", kubeContext) if err != nil { return err } diff --git a/pkg/fission-cli/util/portforward.go b/pkg/fission-cli/util/portforward.go index cccfdbaf..b87029c2 100644 --- a/pkg/fission-cli/util/portforward.go +++ b/pkg/fission-cli/util/portforward.go @@ -43,7 +43,7 @@ const maxDuration time.Duration = 2000 // is found by looking for a service in the same namespace and using // its targetPort. Once the port forward is started, wait for it to // start accepting connections before returning. -func SetupPortForward(namespace, labelSelector string, kubeContext string) (string, error) { +func SetupPortForward(ctx context.Context, namespace, labelSelector string, kubeContext string) (string, error) { console.Verbose(2, "Setting up port forward to %s in namespace %s", labelSelector, namespace) @@ -71,7 +71,7 @@ func SetupPortForward(namespace, labelSelector string, kubeContext string) (stri console.Verbose(2, "Starting port forward from local port %v", localPort) - readyC, _, err := runPortForward(context.Background(), labelSelector, localPort, namespace, kubeContext) + readyC, _, err := runPortForward(ctx, labelSelector, localPort, namespace, kubeContext) if err != nil { fmt.Printf("Error forwarding to port %v: %s", localPort, err.Error()) return "", err diff --git a/pkg/fission-cli/util/util.go b/pkg/fission-cli/util/util.go index ed512d98..686afaae 100644 --- a/pkg/fission-cli/util/util.go +++ b/pkg/fission-cli/util/util.go @@ -17,6 +17,7 @@ limitations under the License. package util import ( + "context" "fmt" "net/url" "os" @@ -53,13 +54,13 @@ func GetFissionNamespace() string { return fissionNamespace } -func GetApplicationUrl(selector string, kubeContext string) (string, error) { +func GetApplicationUrl(ctx context.Context, selector string, kubeContext string) (string, error) { var serverUrl string // Use FISSION_URL env variable if set; otherwise, port-forward to controller. fissionUrl := os.Getenv("FISSION_URL") if len(fissionUrl) == 0 { fissionNamespace := GetFissionNamespace() - localPort, err := SetupPortForward(fissionNamespace, selector, kubeContext) + localPort, err := SetupPortForward(ctx, fissionNamespace, selector, kubeContext) if err != nil { return "", err } @@ -213,7 +214,7 @@ func GetServerURL(input cli.Input) (serverUrl string, err error) { kubeContext := input.String(flagkey.KubeContext) if len(serverUrl) == 0 { // starts local portforwarder etc. - serverUrl, err = GetApplicationUrl("application=fission-api", kubeContext) + serverUrl, err = GetApplicationUrl(input.Context(), "application=fission-api", kubeContext) if err != nil { return "", err } @@ -445,8 +446,8 @@ func ApplyLabelsAndAnnotations(input cli.Input, objectMeta *metav1.ObjectMeta) e return nil } -func GetStorageURL(kubeContext string) (*url.URL, error) { - storageLocalPort, err := SetupPortForward(GetFissionNamespace(), "application=fission-storage", kubeContext) +func GetStorageURL(ctx context.Context, kubeContext string) (*url.URL, error) { + storageLocalPort, err := SetupPortForward(ctx, GetFissionNamespace(), "application=fission-storage", kubeContext) if err != nil { return nil, err } diff --git a/pkg/kubewatcher/main.go b/pkg/kubewatcher/main.go index a76544a7..41540293 100644 --- a/pkg/kubewatcher/main.go +++ b/pkg/kubewatcher/main.go @@ -32,7 +32,7 @@ func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { return errors.Wrap(err, "failed to get fission or kubernetes client") } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { return errors.Wrap(err, "error waiting for CRDs") } diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index 38c28981..6a3bfbe8 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -154,7 +154,7 @@ func StartScalerManager(ctx context.Context, logger *zap.Logger, routerURL strin if err != nil { return err } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { return errors.Wrap(err, "error waiting for CRDs") } diff --git a/pkg/router/router.go b/pkg/router/router.go index f126470c..2f9bd6be 100644 --- a/pkg/router/router.go +++ b/pkg/router/router.go @@ -95,7 +95,7 @@ func Start(ctx context.Context, logger *zap.Logger, port int, executorURL string logger.Fatal("error connecting to kubernetes API", zap.Error(err)) } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { logger.Fatal("error waiting for CRDs", zap.Error(err)) } diff --git a/pkg/timer/main.go b/pkg/timer/main.go index 91de9c23..ecae3784 100644 --- a/pkg/timer/main.go +++ b/pkg/timer/main.go @@ -32,7 +32,7 @@ func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { return errors.Wrap(err, "failed to get fission or kubernetes client") } - err = crd.WaitForCRDs(fissionClient) + err = crd.WaitForCRDs(ctx, fissionClient) if err != nil { return errors.Wrap(err, "error waiting for CRDs") }