From 4cbe6a70613153ed3636cfdf9fcab2aa0fa08b80 Mon Sep 17 00:00:00 2001 From: neha_gupta Date: Thu, 17 Nov 2022 13:13:33 +0530 Subject: [PATCH] Get logs from Pods using Kubernetes API for function log command (#2623) * add controller enablement flag * throw an error if service not found * add logs from Kubernetes in function log command * pass context in function param * add pod-namespace in function log command * pass context in function param * search for the pod in fn ns in the test --- pkg/fission-cli/cmd/function/command.go | 2 +- pkg/fission-cli/cmd/function/log.go | 45 ++++---- pkg/fission-cli/flag/flag.go | 5 +- pkg/fission-cli/flag/key/key.go | 1 + pkg/fission-cli/logdb/influxdb.go | 31 +++++- pkg/fission-cli/logdb/kubernetes_log.go | 141 ++++++++++++++++++++++++ pkg/fission-cli/logdb/logdb.go | 30 +++-- 7 files changed, 218 insertions(+), 37 deletions(-) create mode 100644 pkg/fission-cli/logdb/kubernetes_log.go diff --git a/pkg/fission-cli/cmd/function/command.go b/pkg/fission-cli/cmd/function/command.go index 89c4f043..c4a29489 100644 --- a/pkg/fission-cli/cmd/function/command.go +++ b/pkg/fission-cli/cmd/function/command.go @@ -133,7 +133,7 @@ func Commands() *cobra.Command { Required: []flag.Flag{flag.FnName}, Optional: []flag.Flag{ flag.FnLogFollow, flag.FnLogReverseQuery, flag.FnLogCount, - flag.FnLogDetail, flag.FnLogPod, flag.NamespaceFunction, flag.FnLogDBType}, + flag.FnLogDetail, flag.FnLogPod, flag.NamespaceFunction, flag.FnLogDBType, flag.NamespacePod}, }) testCmd := &cobra.Command{ diff --git a/pkg/fission-cli/cmd/function/log.go b/pkg/fission-cli/cmd/function/log.go index 7bcffa66..75808a7a 100644 --- a/pkg/fission-cli/cmd/function/log.go +++ b/pkg/fission-cli/cmd/function/log.go @@ -19,6 +19,8 @@ package function import ( "context" "fmt" + "io" + "os" "time" "github.com/pkg/errors" @@ -28,7 +30,6 @@ import ( "github.com/fission/fission/pkg/fission-cli/cmd" flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" "github.com/fission/fission/pkg/fission-cli/logdb" - "github.com/fission/fission/pkg/fission-cli/util" ) type LogSubCommand struct { @@ -60,13 +61,12 @@ func (opts *LogSubCommand) do(input cli.Input) error { return errors.Wrap(err, "error getting function") } - server, err := util.GetApplicationUrl(input.Context(), opts.Client(), "application=fission-api") - if err != nil { - return err + logDBOptions := logdb.LogDBOptions{ + Client: opts.Client(), } // request the controller to establish a proxy server to the database. - logDB, err := logdb.GetLogDB(dbType, server) + logDB, err := logdb.GetLogDB(dbType, input.Context(), logDBOptions) if err != nil { return errors.Wrapf(err, "failed to get log database") } @@ -77,31 +77,36 @@ func (opts *LogSubCommand) do(input cli.Input) error { go func(ctx context.Context, requestChan, responseChan chan struct{}) { t := time.Unix(0, 0*int64(time.Millisecond)) + detail := input.Bool(flagkey.FnLogDetail) for { select { case <-requestChan: logFilter := logdb.LogFilter{ - Pod: fnPod, - Function: f.ObjectMeta.Name, - FuncUid: string(f.ObjectMeta.UID), - Since: t, - Reverse: logReverseQuery, - RecordLimit: recordLimit, + Pod: fnPod, + PodNamespace: input.String(flagkey.NamespacePod), + Function: f.ObjectMeta.Name, + FuncUid: string(f.ObjectMeta.UID), + Since: t, + Reverse: logReverseQuery, + RecordLimit: recordLimit, + FunctionObject: f, + Details: detail, } - logEntries, err := logDB.GetLogs(logFilter) + + buf, err := logDB.GetLogs(ctx, logFilter) if err != nil { fmt.Printf("Error querying logs: %v", err) responseChan <- struct{}{} return } - for _, logEntry := range logEntries { - if input.Bool(flagkey.FnLogDetail) { - fmt.Printf("Timestamp: %s\nNamespace: %s\nFunction Name: %s\nFunction ID: %s\nPod: %s\nContainer: %s\nStream: %s\nLog: %s\n---\n", - logEntry.Timestamp, logEntry.Namespace, logEntry.FuncName, logEntry.FuncUid, logEntry.Pod, logEntry.Container, logEntry.Stream, logEntry.Message) - } else { - fmt.Printf("[%s] %s\n", logEntry.Timestamp, logEntry.Message) - } - t = logEntry.Timestamp + _, err = io.Copy(os.Stdout, buf) + if err != nil { + return + } + + t = time.Now().UTC() // next time fetch values from this time + if dbType == logdb.KUBERNETES { //in case of Kubernetes log we print pods info only once. And then print new logs + detail = false } responseChan <- struct{}{} case <-ctx.Done(): diff --git a/pkg/fission-cli/flag/flag.go b/pkg/fission-cli/flag/flag.go index 33a1aa9b..38a086df 100644 --- a/pkg/fission-cli/flag/flag.go +++ b/pkg/fission-cli/flag/flag.go @@ -117,9 +117,10 @@ var ( FnLogPod = Flag{Type: String, Name: flagkey.FnLogPod, Usage: "Function pod name (use the latest pod name if unspecified)"} FnLogFollow = Flag{Type: Bool, Name: flagkey.FnLogFollow, Short: "f", Usage: "Specify if the logs should be streamed"} FnLogDetail = Flag{Type: Bool, Name: flagkey.FnLogDetail, Short: "d", Usage: "Display detailed information"} - FnLogDBType = Flag{Type: String, Name: flagkey.FnLogDBType, Usage: "Log database type, e.g. influxdb (currently only influxdb is supported)", DefaultValue: "influxdb"} - FnLogReverseQuery = Flag{Type: Bool, Name: flagkey.FnLogReverseQuery, Short: "r", Usage: "Specify the log reverse query base on time, it will be invalid if the 'follow' flag is specified"} + FnLogDBType = Flag{Type: String, Name: flagkey.FnLogDBType, Usage: "Log database type, e.g. influxdb (currently influxdb and kubernetes logs are supported)", DefaultValue: "kubernetes"} + FnLogReverseQuery = Flag{Type: Bool, Name: flagkey.FnLogReverseQuery, Short: "r", Usage: "Specify the log reverse query base on time, it will be invalid if the 'follow' flag is specified. valid for dbtype as influxdb"} FnLogCount = Flag{Type: Int, Name: flagkey.FnLogCount, Usage: "Get N most recent log records", DefaultValue: 20} + NamespacePod = Flag{Type: String, Name: flagkey.NamespacePod, Usage: "Namespace in which function's pod are created. If not specified, function's namespace is used. Note: version <1.18 used fission-function as pod's default ns."} FnTestBody = Flag{Type: String, Name: flagkey.FnTestBody, Short: "b", Usage: "Request body"} FnTestTimeout = Flag{Type: Duration, Name: flagkey.FnTestTimeout, Short: "t", Usage: "Length of time to wait for the response. If set to zero or negative number, no timeout is set", DefaultValue: 60 * time.Second} FnTestHeader = Flag{Type: StringSlice, Name: flagkey.FnTestHeader, Short: "H", Usage: "Request headers"} diff --git a/pkg/fission-cli/flag/key/key.go b/pkg/fission-cli/flag/key/key.go index a15664e0..e1116eb6 100644 --- a/pkg/fission-cli/flag/key/key.go +++ b/pkg/fission-cli/flag/key/key.go @@ -41,6 +41,7 @@ const ( Namespace = "namespace" ForceNamespace = "force-namespace" AllNamespaces = "all-namespaces" + NamespacePod = "pod-namespace" ForceDelete = "force" RuntimeMincpu = "mincpu" diff --git a/pkg/fission-cli/logdb/influxdb.go b/pkg/fission-cli/logdb/influxdb.go index aca6201f..bff9bea9 100644 --- a/pkg/fission-cli/logdb/influxdb.go +++ b/pkg/fission-cli/logdb/influxdb.go @@ -17,6 +17,8 @@ limitations under the License. package logdb import ( + "bytes" + "context" "encoding/json" "fmt" "net/http" @@ -31,6 +33,7 @@ import ( "github.com/pkg/errors" ferror "github.com/fission/fission/pkg/error" + "github.com/fission/fission/pkg/fission-cli/util" ) const ( @@ -38,8 +41,12 @@ const ( INFLUXDB_URL = "http://influxdb:8086/query" ) -func NewInfluxDB(serverURL string) (InfluxDB, error) { - return InfluxDB{endpoint: serverURL}, nil +func NewInfluxDB(ctx context.Context, logDBOptions LogDBOptions) (InfluxDB, error) { + server, err := util.GetApplicationUrl(ctx, logDBOptions.Client, "application=fission-api") + if err != nil { + return InfluxDB{}, err + } + return InfluxDB{endpoint: server}, nil } type InfluxDB struct { @@ -55,7 +62,7 @@ func makeIndexMap(cols []string) map[string]int { return indexMap } -func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) { +func (influx InfluxDB) GetLogs(ctx context.Context, filter LogFilter) (output *bytes.Buffer, err error) { timestamp := filter.Since.UnixNano() var queryCmd string @@ -133,7 +140,23 @@ func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) { sort.Sort(ByTimestamp(logEntries, filter.Reverse)) - return logEntries, nil + output = new(bytes.Buffer) + for _, logEntry := range logEntries { + if filter.Details { + msg := fmt.Sprintf("Timestamp: %s\nNamespace: %s\nFunction Name: %s\nFunction ID: %s\nPod: %s\nContainer: %s\nStream: %s\nLog: %s\n---\n", + logEntry.Timestamp, logEntry.Namespace, logEntry.FuncName, logEntry.FuncUid, logEntry.Pod, logEntry.Container, logEntry.Stream, logEntry.Message) + if _, err := output.WriteString(msg); err != nil { + return output, errors.Wrapf(err, "error copying pod log") + } + } else { + msg := fmt.Sprintf("[%s] %s\n", logEntry.Timestamp, logEntry.Message) + if _, err := output.WriteString(msg); err != nil { + return output, errors.Wrapf(err, "error copying pod log") + } + } + } + + return output, nil } func (influx InfluxDB) query(query influxdbClient.Query) (*influxdbClient.Response, error) { diff --git a/pkg/fission-cli/logdb/kubernetes_log.go b/pkg/fission-cli/logdb/kubernetes_log.go new file mode 100644 index 00000000..107d6fd5 --- /dev/null +++ b/pkg/fission-cli/logdb/kubernetes_log.go @@ -0,0 +1,141 @@ +/* +Copyright 2016 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 logdb + +import ( + "bytes" + "context" + "fmt" + "io" + "sort" + "strconv" + "strings" + + "github.com/pkg/errors" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/client-go/kubernetes" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/fission-cli/cmd" + "github.com/fission/fission/pkg/fission-cli/console" +) + +type LogDBOptions struct { + Client cmd.Client +} + +type kubernetesLogs struct { + client cmd.Client +} + +func (k kubernetesLogs) GetLogs(ctx context.Context, logFilter LogFilter) (podLogs *bytes.Buffer, err error) { + podLogs, err = GetFunctionPodLogs(ctx, k.client, logFilter) + return podLogs, err +} + +func NewKubernetesEndpoint(logDBOptions LogDBOptions) (kubernetesLogs, error) { + return kubernetesLogs{ + client: logDBOptions.Client}, nil +} + +// FunctionPodLogs : Get logs for a function directly from pod +func GetFunctionPodLogs(ctx context.Context, client cmd.Client, logFilter LogFilter) (podLogs *bytes.Buffer, err error) { + + f := logFilter.FunctionObject + + podNs := f.Namespace + if logFilter.PodNamespace != "" { + podNs = logFilter.PodNamespace + } + // Get function Pods first + selector := map[string]string{ + fv1.FUNCTION_UID: string(f.ObjectMeta.UID), + fv1.ENVIRONMENT_NAME: f.Spec.Environment.Name, + fv1.ENVIRONMENT_NAMESPACE: f.Spec.Environment.Namespace, + } + podList, err := client.KubernetesClient.CoreV1().Pods(podNs).List(ctx, metav1.ListOptions{ + LabelSelector: labels.Set(selector).AsSelector().String(), + }) + if err != nil { + return podLogs, err + } + + // Get the logs for last Pod executed + pods := podList.Items + sort.Slice(pods, func(i, j int) bool { + rv1, _ := strconv.ParseInt(pods[i].ObjectMeta.ResourceVersion, 10, 32) + rv2, _ := strconv.ParseInt(pods[j].ObjectMeta.ResourceVersion, 10, 32) + return rv1 > rv2 + }) + + if len(pods) <= 0 { + console.Warn("version<1.18 used fission-function as pod's default namespace. Specify appropriate namespace with --pod-namespace tag.") + return podLogs, errors.New("no active pods found") + + } + + // get the pod with highest resource version + podLogs, err = streamContainerLog(ctx, client.KubernetesClient, &pods[0], logFilter) + if err != nil { + return podLogs, errors.Wrapf(err, "error getting container logs") + + } + return podLogs, err +} + +func streamContainerLog(ctx context.Context, kubernetesClient kubernetes.Interface, pod *v1.Pod, logFilter LogFilter) (output *bytes.Buffer, err error) { + + seq := strings.Repeat("=", 35) + output = new(bytes.Buffer) + + for _, container := range pod.Spec.Containers { + tailLines := int64(logFilter.RecordLimit) + sinceTime := metav1.NewTime(logFilter.Since) + podLogOpts := v1.PodLogOptions{Container: container.Name, // Only the env container, not fetcher + SinceTime: &sinceTime, + TailLines: &tailLines, + } + + podLogsReq := kubernetesClient.CoreV1().Pods(pod.Namespace).GetLogs(pod.ObjectMeta.Name, &podLogOpts) + + podLogs, err := podLogsReq.Stream(ctx) + if err != nil { + return output, errors.Wrapf(err, "error streaming pod log") + } + + if logFilter.Details { + fn := logFilter.FunctionObject + 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) + + if _, err := output.WriteString(msg); err != nil { + return output, errors.Wrapf(err, "error copying pod log") + } + } + + _, err = io.Copy(output, podLogs) + if err != nil { + return output, errors.Wrapf(err, "error copying pod log") + } + + podLogs.Close() + } + + return output, nil +} diff --git a/pkg/fission-cli/logdb/logdb.go b/pkg/fission-cli/logdb/logdb.go index feb88bc9..e2c6a65b 100644 --- a/pkg/fission-cli/logdb/logdb.go +++ b/pkg/fission-cli/logdb/logdb.go @@ -17,25 +17,33 @@ limitations under the License. package logdb import ( + "bytes" + "context" "fmt" "time" + + v1 "github.com/fission/fission/pkg/apis/core/v1" ) const ( - INFLUXDB = "influxdb" + INFLUXDB = "influxdb" + KUBERNETES = "kubernetes" ) type LogDatabase interface { - GetLogs(LogFilter) ([]LogEntry, error) + GetLogs(context.Context, LogFilter) (*bytes.Buffer, error) } type LogFilter struct { - Pod string - Function string - FuncUid string - Since time.Time - Reverse bool - RecordLimit int + Pod string + PodNamespace string + Function string + FuncUid string + Since time.Time + Reverse bool + RecordLimit int + FunctionObject *v1.Function + Details bool } type LogEntry struct { @@ -69,10 +77,12 @@ func ByTimestamp(entries []LogEntry, desc bool) ByTimestampSort { return ByTimestampSort{entries, desc} } -func GetLogDB(dbType string, serverURL string) (LogDatabase, error) { +func GetLogDB(dbType string, ctx context.Context, logDBOptions LogDBOptions) (LogDatabase, error) { switch dbType { case INFLUXDB: - return NewInfluxDB(serverURL) + return NewInfluxDB(ctx, logDBOptions) + case KUBERNETES: + return NewKubernetesEndpoint(logDBOptions) } return nil, fmt.Errorf("log database type is incorrect, now only support %s", INFLUXDB) }