Refactor kubewatch command (#1365)

This commit is contained in:
Ta-Ching Chen
2019-10-31 10:48:27 +08:00
committed by GitHub
parent 698d591788
commit e22cdee5f5
5 changed files with 195 additions and 72 deletions
+122
View File
@@ -0,0 +1,122 @@
/*
Copyright 2019 The Fission Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package kubewatch
import (
"fmt"
"github.com/pkg/errors"
uuid "github.com/satori/go.uuid"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
fv1 "github.com/fission/fission/pkg/apis/fission.io/v1"
"github.com/fission/fission/pkg/controller/client"
"github.com/fission/fission/pkg/fission-cli/cliwrapper/cli"
"github.com/fission/fission/pkg/fission-cli/cmd"
"github.com/fission/fission/pkg/fission-cli/cmd/spec"
"github.com/fission/fission/pkg/fission-cli/log"
)
type CreateSubCommand struct {
client *client.Client
watcher *fv1.KubernetesWatchTrigger
}
func Create(flags cli.Input) error {
opts := CreateSubCommand{
client: cmd.GetServer(flags),
}
return opts.do(flags)
}
func (opts *CreateSubCommand) do(flags cli.Input) error {
err := opts.complete(flags)
if err != nil {
return err
}
return opts.run(flags)
}
func (opts *CreateSubCommand) complete(flags cli.Input) error {
fnName := flags.String("function")
if len(fnName) == 0 {
log.Fatal("Need a function name to create a watch, use --function")
}
fnNamespace := flags.String("fnNamespace")
namespace := flags.String("ns")
if len(namespace) == 0 {
fmt.Println("Watch 'default' namespace. Use --ns <namespace> to override.")
namespace = "default"
}
objType := flags.String("type")
if len(objType) == 0 {
fmt.Println("Object type unspecified, will watch pods. Use --type <type> to override.")
objType = "pod"
}
labels := flags.String("labels")
// empty 'labels' selects everything
if len(labels) == 0 {
fmt.Printf("Watching all objects of type '%v', use --labels to refine selection.\n", objType)
} else {
// TODO
fmt.Printf("Label selector not implemented, watching all objects")
}
// automatically name watches
watchName := uuid.NewV4().String()
opts.watcher = &fv1.KubernetesWatchTrigger{
Metadata: metav1.ObjectMeta{
Name: watchName,
Namespace: fnNamespace,
},
Spec: fv1.KubernetesWatchTriggerSpec{
Namespace: namespace,
Type: objType,
//LabelSelector: labels,
FunctionReference: fv1.FunctionReference{
Name: fnName,
Type: fv1.FunctionReferenceTypeFunctionName,
},
},
}
return nil
}
func (opts *CreateSubCommand) run(flags cli.Input) error {
// if we're writing a spec, don't call the API
if flags.Bool("spec") {
specFile := fmt.Sprintf("kubewatch-%v.yaml", opts.watcher.Metadata.Name)
err := spec.SpecSave(*opts.watcher, specFile)
if err != nil {
return errors.Wrap(err, "error creating kubewatch spec")
}
return nil
}
_, err := opts.client.WatchCreate(opts.watcher)
if err != nil {
return errors.Wrap(err, "error creating kubewatch")
}
fmt.Printf("kubewatch '%v' created\n", opts.watcher.Metadata.Name)
return nil
}
+71
View File
@@ -0,0 +1,71 @@
/*
Copyright 2019 The Fission Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package kubewatch
import (
"fmt"
"github.com/pkg/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"github.com/fission/fission/pkg/controller/client"
"github.com/fission/fission/pkg/fission-cli/cliwrapper/cli"
"github.com/fission/fission/pkg/fission-cli/cmd"
)
type DeleteSubCommand struct {
client *client.Client
name string
namespace string
}
func Delete(flags cli.Input) error {
opts := DeleteSubCommand{
client: cmd.GetServer(flags),
}
return opts.do(flags)
}
func (opts *DeleteSubCommand) do(flags cli.Input) error {
err := opts.complete(flags)
if err != nil {
return err
}
return opts.run(flags)
}
func (opts *DeleteSubCommand) complete(flags cli.Input) error {
opts.name = flags.String("name")
if len(opts.name) == 0 {
return errors.New("need name of watch to delete, use --name")
}
opts.namespace = flags.String("triggerns")
return nil
}
func (opts *DeleteSubCommand) run(flags cli.Input) error {
err := opts.client.WatchDelete(&metav1.ObjectMeta{
Name: opts.name,
Namespace: opts.namespace,
})
if err != nil {
return errors.Wrap(err, "error deleting kubewatch")
}
fmt.Printf("watch '%v' deleted\n", opts.name)
return nil
}
+73
View File
@@ -0,0 +1,73 @@
/*
Copyright 2019 The Fission Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package kubewatch
import (
"fmt"
"os"
"text/tabwriter"
"github.com/pkg/errors"
"github.com/fission/fission/pkg/controller/client"
"github.com/fission/fission/pkg/fission-cli/cliwrapper/cli"
"github.com/fission/fission/pkg/fission-cli/cmd"
)
type ListSubCommand struct {
client *client.Client
namespace string
}
func List(flags cli.Input) error {
opts := ListSubCommand{
client: cmd.GetServer(flags),
}
return opts.do(flags)
}
func (opts *ListSubCommand) do(flags cli.Input) error {
err := opts.complete(flags)
if err != nil {
return err
}
return opts.run(flags)
}
func (opts *ListSubCommand) complete(flags cli.Input) error {
opts.namespace = flags.String("triggerns")
return nil
}
func (opts *ListSubCommand) run(flags cli.Input) error {
ws, err := opts.client.WatchList(opts.namespace)
if err != nil {
return errors.Wrap(err, "error listing kubewatches")
}
w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0)
fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\n",
"NAME", "NAMESPACE", "OBJTYPE", "LABELS", "FUNCTION_NAME")
for _, wa := range ws {
fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\n",
wa.Metadata.Name, wa.Spec.Namespace, wa.Spec.Type, wa.Spec.LabelSelector, wa.Spec.FunctionReference.Name)
}
w.Flush()
return nil
}
-1
View File
@@ -56,7 +56,6 @@ func (opts *CreateSubCommand) do(flags cli.Input) error {
return nil
}
// complete creates a environment objects and populates it with default value and CLI inputs.
func (opts *CreateSubCommand) complete(flags cli.Input) error {
pkgNamespace := flags.String("pkgNamespace")
envName := flags.String("env")