Add verbosity flag and verbose logs for portforwarder (#575)

Also, parse flags before connecting to server.
This commit is contained in:
Soam Vasani
2018-03-26 12:23:45 -07:00
committed by GitHub
parent 6b188c053a
commit 8708fdec93
3 changed files with 53 additions and 15 deletions
+13 -2
View File
@@ -40,6 +40,11 @@ import (
storageSvcClient "github.com/fission/fission/storagesvc/client" storageSvcClient "github.com/fission/fission/storagesvc/client"
) )
var (
// global verbosity of our CLI
verbosity int
)
func fatal(msg string) { func fatal(msg string) {
os.Stderr.WriteString(msg + "\n") os.Stderr.WriteString(msg + "\n")
os.Exit(1) os.Exit(1)
@@ -49,10 +54,16 @@ func warn(msg string) {
os.Stderr.WriteString(msg + "\n") os.Stderr.WriteString(msg + "\n")
} }
func getClient(serverUrl string) *client.Client { func verbose(msglevel int, format string, args ...interface{}) {
if verbosity >= msglevel {
fmt.Printf(format+"\n", args...)
}
}
func getClient(serverUrl string) *client.Client {
if len(serverUrl) == 0 { if len(serverUrl) == 0 {
fatal("Need --server or FISSION_URL set to your fission server.") // starts local portforwarder etc.
serverUrl = getServerUrl()
} }
isHTTPS := strings.Index(serverUrl, "https://") == 0 isHTTPS := strings.Index(serverUrl, "https://") == 0
+23 -12
View File
@@ -60,29 +60,38 @@ func getFissionAPIVersion(apiUrl string) (string, error) {
return strings.TrimRight(string(body), "\n"), nil return strings.TrimRight(string(body), "\n"), nil
} }
func main() { func getServerUrl() string {
app := cli.NewApp() var serverUrl string
app.Name = "fission" // Use FISSION_URL env variable if set; otherwise, port-forward to controller.
app.Usage = "Serverless functions for Kubernetes"
app.Version = version.Version
// fetch the FISSION_URL env variable. If not set, port-forward to controller.
var value string
fissionUrl := os.Getenv("FISSION_URL") fissionUrl := os.Getenv("FISSION_URL")
if len(fissionUrl) == 0 { if len(fissionUrl) == 0 {
fissionNamespace := getFissionNamespace() fissionNamespace := getFissionNamespace()
kubeConfig := getKubeConfigPath() kubeConfig := getKubeConfigPath()
localPort := setupPortForward( localPort := setupPortForward(
kubeConfig, fissionNamespace, "application=fission-api") kubeConfig, fissionNamespace, "application=fission-api")
value = "http://127.0.0.1:" + localPort serverUrl = "http://127.0.0.1:" + localPort
} else { } else {
value = fissionUrl serverUrl = fissionUrl
} }
return serverUrl
}
func cliHook(c *cli.Context) error {
verbosity = c.Int("verbosity")
verbose(2, "Verbosity = 2")
return nil
}
func main() {
app := cli.NewApp()
app.Name = "fission"
app.Usage = "Serverless functions for Kubernetes"
app.Version = version.Version
cli.VersionPrinter = func(c *cli.Context) { cli.VersionPrinter = func(c *cli.Context) {
clientVer := version.VersionInfo().String() clientVer := version.VersionInfo().String()
fmt.Printf("Client Version: %v\n", clientVer) fmt.Printf("Client Version: %v\n", clientVer)
serverVer, err := getFissionAPIVersion(value) serverVer, err := getFissionAPIVersion(getServerUrl())
if err != nil { if err != nil {
fmt.Printf("Error getting Fission API version: %v", err) fmt.Printf("Error getting Fission API version: %v", err)
} else { } else {
@@ -91,7 +100,8 @@ func main() {
} }
app.Flags = []cli.Flag{ app.Flags = []cli.Flag{
cli.StringFlag{Name: "server", Value: value, Usage: "Fission server URL"}, cli.StringFlag{Name: "server", Value: "", Usage: "Fission server URL"},
cli.IntFlag{Name: "verbosity", Value: 1, Usage: "CLI verbosity (0 is quiet, 1 is the default, 2 is verbose.)"},
} }
// trigger method and url flags (used in function and route CLIs) // trigger method and url flags (used in function and route CLIs)
@@ -273,5 +283,6 @@ func main() {
{Name: "tpr2crd", Aliases: []string{}, Usage: "Migrate tool for TPR to CRD", Subcommands: migrateSubCommands}, {Name: "tpr2crd", Aliases: []string{}, Usage: "Migrate tool for TPR to CRD", Subcommands: migrateSubCommands},
} }
app.Before = cliHook
app.Run(os.Args) app.Run(os.Args)
} }
+17 -1
View File
@@ -52,6 +52,8 @@ func runPortForward(kubeConfig string, labelSelector string, localPort string, f
fatal(fmt.Sprintf("Failed to connect to Kubernetes: %s", err)) fatal(fmt.Sprintf("Failed to connect to Kubernetes: %s", err))
} }
verbose(2, "Connected to Kubernetes API")
// if fission namespace is unset, try to find a fission pod in any namespace // if fission namespace is unset, try to find a fission pod in any namespace
if len(fissionNamespace) == 0 { if len(fissionNamespace) == 0 {
fissionNamespace = meta_v1.NamespaceAll fissionNamespace = meta_v1.NamespaceAll
@@ -93,6 +95,7 @@ func runPortForward(kubeConfig string, labelSelector string, localPort string, f
for _, servicePort := range service.Spec.Ports { for _, servicePort := range service.Spec.Ports {
targetPort = servicePort.TargetPort.String() targetPort = servicePort.TargetPort.String()
} }
verbose(2, "Connecting to port %v on pod %v/%v", targetPort, podNameSpace, podNameSpace)
stopChannel := make(chan struct{}, 1) stopChannel := make(chan struct{}, 1)
readyChannel := make(chan struct{}) readyChannel := make(chan struct{})
@@ -113,12 +116,17 @@ func runPortForward(kubeConfig string, labelSelector string, localPort string, f
fatal(msg) fatal(msg)
} }
fw, err := portforward.New(dialer, ports, stopChannel, readyChannel, nil, os.Stderr) outStream := os.Stdout
if verbosity < 2 {
outStream = nil
}
fw, err := portforward.New(dialer, ports, stopChannel, readyChannel, outStream, os.Stderr)
if err != nil { if err != nil {
msg := fmt.Sprintf("portforward.new errored out :%v", err.Error()) msg := fmt.Sprintf("portforward.new errored out :%v", err.Error())
fatal(msg) fatal(msg)
} }
verbose(2, "Starting port forwarder")
return fw.ForwardPorts() return fw.ForwardPorts()
} }
@@ -128,11 +136,15 @@ func runPortForward(kubeConfig string, labelSelector string, localPort string, f
// its targetPort. Once the port forward is started, wait for it to // its targetPort. Once the port forward is started, wait for it to
// start accepting connections before returning. // start accepting connections before returning.
func setupPortForward(kubeConfig, namespace, labelSelector string) string { func setupPortForward(kubeConfig, namespace, labelSelector string) string {
verbose(2, "Setting up port forward to %s in namespace %s using the kubeconfig at %s",
labelSelector, namespace, kubeConfig)
localPort, err := findFreePort() localPort, err := findFreePort()
if err != nil { if err != nil {
fatal(fmt.Sprintf("Error finding unused port :%v", err.Error())) fatal(fmt.Sprintf("Error finding unused port :%v", err.Error()))
} }
verbose(2, "Waiting for local port %v", localPort)
for { for {
conn, _ := net.DialTimeout("tcp", conn, _ := net.DialTimeout("tcp",
net.JoinHostPort("", localPort), time.Millisecond) net.JoinHostPort("", localPort), time.Millisecond)
@@ -144,6 +156,7 @@ func setupPortForward(kubeConfig, namespace, labelSelector string) string {
time.Sleep(time.Millisecond * 50) time.Sleep(time.Millisecond * 50)
} }
verbose(2, "Starting port forward from local port %v", localPort)
go func() { go func() {
err := runPortForward(kubeConfig, labelSelector, localPort, namespace) err := runPortForward(kubeConfig, labelSelector, localPort, namespace)
if err != nil { if err != nil {
@@ -151,6 +164,7 @@ func setupPortForward(kubeConfig, namespace, labelSelector string) string {
} }
}() }()
verbose(2, "Waiting for port forward %v to start...", localPort)
for { for {
conn, _ := net.DialTimeout("tcp", conn, _ := net.DialTimeout("tcp",
net.JoinHostPort("", localPort), time.Millisecond) net.JoinHostPort("", localPort), time.Millisecond)
@@ -161,5 +175,7 @@ func setupPortForward(kubeConfig, namespace, labelSelector string) string {
time.Sleep(time.Millisecond * 50) time.Sleep(time.Millisecond * 50)
} }
verbose(2, "Port forward from local port %v started", localPort)
return localPort return localPort
} }