From 4e4c8aa14fc761b00a613bebdbaea2d5717fb2fa Mon Sep 17 00:00:00 2001 From: Soam Vasani Date: Mon, 18 Mar 2019 00:42:37 -0700 Subject: [PATCH] Switch from fluentd to fluentbit for log forwarding (#1086) This removes fluentd in favor of using fluentbit, which is lighter (in memory usage) and seems to be more actively maintained. Fluentbit's config file format is different from fluentd's. It also doesn't support the same record modification stuff that fluentd supports, so we have to change the influxdb query slightly. This means that after an upgrade, the new CLI may won't work for querying older logs. Hopefully, this slight breakage is acceptable; if users really need older logs they can use the older CLI. --- charts/README.md | 5 +- charts/fission-all/config/fluentbit.conf | 55 ++++++ charts/fission-all/config/fluentd.conf | 59 ------ charts/fission-all/config/parsers.conf | 12 ++ .../{fluentd.yaml => fluentbit.yaml} | 58 ++++-- charts/fission-all/values.yaml | 6 +- environments/fetcher/cmd/main.go | 18 +- environments/fetcher/fetcher.go | 3 +- examples/python/sourcepkg/user.py | 3 +- executor/executor.go | 5 - executor/poolmgr/gp.go | 14 +- fission-bundle/main.go | 11 ++ fission/logdb/influxdb.go | 51 +++-- hack/release-build.sh | 18 -- hack/release.sh | 9 - hack/runtests.sh | 2 +- logger/fluentd/Dockerfile | 70 ------- logger/fluentd/build.sh | 2 - logger/fluentd/fluent.conf | 55 ------ logger/fluentd/run.sh | 27 --- logger/logger.go | 181 ++++++++++++++++++ logger/types.go | 35 ---- test/build_and_test.sh | 5 +- test/test_utils.sh | 37 +--- test/upgrade/fission_upgrade_test.sh | 5 +- 25 files changed, 366 insertions(+), 380 deletions(-) create mode 100644 charts/fission-all/config/fluentbit.conf delete mode 100644 charts/fission-all/config/fluentd.conf create mode 100644 charts/fission-all/config/parsers.conf rename charts/fission-all/templates/{fluentd.yaml => fluentbit.yaml} (53%) delete mode 100644 logger/fluentd/Dockerfile delete mode 100755 logger/fluentd/build.sh delete mode 100644 logger/fluentd/fluent.conf delete mode 100755 logger/fluentd/run.sh create mode 100644 logger/logger.go delete mode 100644 logger/types.go diff --git a/charts/README.md b/charts/README.md index b4ea53ba..f35c7b50 100644 --- a/charts/README.md +++ b/charts/README.md @@ -71,8 +71,9 @@ Parameter | Description | Default --------- | ----------- | ------- `createNamespace` | If true, create `fission-function` and `fission-builder` namespaces | ` true` `logger.influxdbAdmin` | Log database admin username | `admin` -`logger.fluentdImage` | Logger fluentd image | `fission/fluentd` -`logger.fluentdImageTag` | Fission ui image tag | `1.0.0` +`logger.fluentdImageRepository` | Logger fluentbit image repository | `index.docker.io` +`logger.fluentdImage` | Logger fluentbit image | `fluent/fluent-bit` +`logger.fluentdImageTag` | Logger fluentbit image tag | `1.0.4` `nats.enabled` | Nats streaming enabled | `true` `nats.authToken` | Nats streaming auth token | `defaultFissionAuthToken` `nats.clusterID` | Nats streaming clusterID | `fissionMQTrigger` diff --git a/charts/fission-all/config/fluentbit.conf b/charts/fission-all/config/fluentbit.conf new file mode 100644 index 00000000..b3a845d1 --- /dev/null +++ b/charts/fission-all/config/fluentbit.conf @@ -0,0 +1,55 @@ +# Some stuff in this file is copied out of https://github.com/fluent/fluent-bit-kubernetes-logging/blob/master/output/kafka/fluent-bit-configmap.yaml +[SERVICE] + Flush 5 + Log_Level info + Parsers_File parsers.conf + +[INPUT] + Name tail + Tag log.* + Path ${LOG_PATH} + Mem_Buf_Limit 5MB + Parser docker + DB /var/log/fission/flb_kube.db + Skip_Long_Lines On + Refresh_Interval 10 + +[FILTER] + Name kubernetes + Match log.* + Kube_URL https://kubernetes.default.svc.cluster.local:443 + +# +# Flatten kubernetes labels so we can tag and query functionuid with influxdb +# +[FILTER] + Name nest + Match log.* + Operation lift + Nested_under kubernetes + Prefix_with kubernetes_ + +[FILTER] + Name nest + Match log.* + Operation lift + Nested_under kubernetes_labels + Prefix_with kubernetes_labels_ + +[OUTPUT] + Name influxdb + Match log.* + Host ${INFLUXDB_ADDRESS} + Port ${INFLUXDB_PORT} + Database ${INFLUXDB_DBNAME} + HTTP_User ${INFLUXDB_USERNAME} + HTTP_Passwd ${INFLUXDB_PASSWD} + Tag_Keys kubernetes_labels_functionUid + Sequence_Tag _seq + + +# Useful for testing config changes +#[OUTPUT] +# Name file +# Match log.* +# Path /flbout.txt diff --git a/charts/fission-all/config/fluentd.conf b/charts/fission-all/config/fluentd.conf deleted file mode 100644 index 49ab6bf1..00000000 --- a/charts/fission-all/config/fluentd.conf +++ /dev/null @@ -1,59 +0,0 @@ -# Hide all fluent-related logs - - type null - - -# Collect all logs from the containers in the current namespace - - type tail - format json - time_key time - path "#{ENV['FLUENTD_PATH']}" - time_format %Y-%m-%dT%H:%M:%S.%NZ - tag fission.* - read_from_head true - refresh_interval 5 - - -# Augment logs with Kubernetes metadata - - type kubernetes_metadata - - -# Simplify the nested objects to XX_YY_ZZ names - - type flatten_hash - separator _ - - -# Add `funcuid` to the record (using the functionUid label) - - type record_reformer - enable_ruby false - tag log - - funcuid ${kubernetes_labels_functionUid} - - - -# Push logs into influxdb - - @type influxdb - host "#{ENV['INFLUXDB_ADDRESS']}" - port "#{ENV['INFLUXDB_PORT']}" - dbname "#{ENV['INFLUXDB_DBNAME']}" - user "#{ENV['INFLUXDB_USERNAME']}" - password "#{ENV['INFLUXDB_PASSWD']}" - use_ssl false - time_precision ns - tag_keys ["funcuid"] - sequence_tag _seq - buffer_type file - buffer_path /var/log/fission/fluentd.buffer - buffer_chunk_limit 128m - buffer_queue_limit 256 - flush_interval 5 - retry_limit 10 - retry_wait 1.0 - num_threads 2 - diff --git a/charts/fission-all/config/parsers.conf b/charts/fission-all/config/parsers.conf new file mode 100644 index 00000000..0c53b088 --- /dev/null +++ b/charts/fission-all/config/parsers.conf @@ -0,0 +1,12 @@ +[PARSER] + Name docker + Format json + Time_Key time + Time_Format %Y-%m-%dT%H:%M:%S.%L + Time_Keep On + + # next bit is useful when the logs from the function itself are in json + # Command | Decoder | Field | Optional Action + # =============|==================|================= + #Decode_Field_As escaped log + diff --git a/charts/fission-all/templates/fluentd.yaml b/charts/fission-all/templates/fluentbit.yaml similarity index 53% rename from charts/fission-all/templates/fluentd.yaml rename to charts/fission-all/templates/fluentbit.yaml index 1eec4dc7..7d0205ba 100644 --- a/charts/fission-all/templates/fluentd.yaml +++ b/charts/fission-all/templates/fluentbit.yaml @@ -1,15 +1,21 @@ -# Fluentd deployment for Fission +# Fluentbit deployment for Fission # # Requires: # - service account: fission-svc apiVersion: v1 kind: ConfigMap metadata: - name: {{ .Release.Name }}-fission-fluentd + name: {{ .Release.Name }}-fission-fluentbit data: -{{- if .Files.Get "config/fluentd.conf" }} - td-agent.conf: | -{{ .Files.Get "config/fluentd.conf" | indent 3 }} +{{- if .Files.Get "config/fluentbit.conf" }} + fluentbit.conf: | +{{ .Files.Get "config/fluentbit.conf" | indent 3 }} +{{ else }} +{{ fail "invalid chart" }} +{{- end }} +{{- if .Files.Get "config/parsers.conf" }} + parsers.conf: | +{{ .Files.Get "config/parsers.conf" | indent 3 }} {{ else }} {{ fail "invalid chart" }} {{- end }} @@ -27,9 +33,29 @@ spec: svc: logger spec: containers: - - name: fluentd + - name: logger + image: "{{ .Values.repository }}/{{ .Values.image }}:{{ .Values.imageTag }}" + imagePullPolicy: IfNotPresent + env: + - name: NODE_NAME + valueFrom: + fieldRef: + apiVersion: v1 + fieldPath: spec.nodeName + command: ["/fission-bundle"] + args: ["--logger"] + volumeMounts: + - name: container-log + mountPath: /var/log/ + readOnly: false + - name: docker-log + mountPath: /var/lib/docker/containers + readOnly: true + - name: fluentbit image: "{{ .Values.logger.fluentdImageRepository }}/{{ .Values.logger.fluentdImage }}:{{ .Values.logger.fluentdImageTag }}" imagePullPolicy: {{ .Values.pullPolicy }} + # CMD ["/fluent-bit/bin/fluent-bit", "-c", "/fluent-bit/etc/fluent-bit.conf"] + command: ["/fluent-bit/bin/fluent-bit", "-c", "/fluent-bit/etc/fluentbit.conf"] env: - name: INFLUXDB_ADDRESS value: influxdb @@ -47,8 +73,8 @@ spec: secretKeyRef: name: influxdb key: password - - name: FLUENTD_PATH - value: /var/log/containers/*{{.Values.functionNamespace}}*.log + - name: LOG_PATH + value: /var/log/fission/*.log volumeMounts: - name: container-log mountPath: /var/log/ @@ -56,11 +82,8 @@ spec: - name: docker-log mountPath: /var/lib/docker/containers readOnly: true - - name: fission-log - mountPath: /var/log/fission - readOnly: false - - name: fluentd-config - mountPath: /etc/td-agent/ + - name: fluentbit-config + mountPath: /fluent-bit/etc/ readOnly: true serviceAccount: fission-svc volumes: @@ -70,12 +93,9 @@ spec: - name: docker-log hostPath: path: /var/lib/docker/containers - - name: fission-log - hostPath: - path: /var/log/fission - # Fluentd config location: /etc/td-agent/td-agent.conf - - name: fluentd-config + # Fluentbit config location: /fluent-bit/etc/*.conf + - name: fluentbit-config configMap: - name: {{ .Release.Name }}-fission-fluentd + name: {{ .Release.Name }}-fission-fluentbit updateStrategy: type: RollingUpdate diff --git a/charts/fission-all/values.yaml b/charts/fission-all/values.yaml index 38ba379d..ccdb1ebd 100644 --- a/charts/fission-all/values.yaml +++ b/charts/fission-all/values.yaml @@ -54,8 +54,8 @@ enableIstio: false logger: influxdbAdmin: "admin" fluentdImageRepository: index.docker.io - fluentdImage: fission/fluentd - fluentdImageTag: 1.0.0 + fluentdImage: fluent/fluent-bit + fluentdImageTag: 1.0.4 ## Message queue trigger config ### NATS Streaming, enabled by default @@ -139,4 +139,4 @@ canaryDeployment: # Use these flags to enable opentracing, the variable is endpoint of Jaeger collector in the format shown below #traceCollectorEndpoint: "http://jaeger-collector.jaeger.svc:14268/api/traces?format=jaeger.thrift" -#traceSamplingRate: 0.75 \ No newline at end of file +#traceSamplingRate: 0.75 diff --git a/environments/fetcher/cmd/main.go b/environments/fetcher/cmd/main.go index f3d8175d..c7e23985 100644 --- a/environments/fetcher/cmd/main.go +++ b/environments/fetcher/cmd/main.go @@ -8,24 +8,16 @@ import ( "log" "net/http" "os" - "os/signal" - "runtime/debug" - "syscall" - - "go.uber.org/zap" "go.opencensus.io/exporter/jaeger" "go.opencensus.io/plugin/ochttp" "go.opencensus.io/trace" + "go.uber.org/zap" "github.com/fission/fission" "github.com/fission/fission/environments/fetcher" ) -func dumpStackTrace() { - debug.PrintStack() -} - func registerTraceExporter(collectorEndpoint string) error { if collectorEndpoint == "" { return nil @@ -57,14 +49,6 @@ func main() { } defer logger.Sync() - // register signal handler for dumping stack trace. - c := make(chan os.Signal, 1) - signal.Notify(c, syscall.SIGTERM) - go func() { - <-c - logger.Fatal("received SIGTERM") - }() - flag.Usage = fetcherUsage collectorEndpoint := flag.String("jaeger-collector-endpoint", "", "") specializeOnStart := flag.Bool("specialize-on-startup", false, "Flag to activate specialize process at pod starup") diff --git a/environments/fetcher/fetcher.go b/environments/fetcher/fetcher.go index 72c7643a..012648e5 100644 --- a/environments/fetcher/fetcher.go +++ b/environments/fetcher/fetcher.go @@ -16,12 +16,11 @@ import ( "path/filepath" "time" - "go.uber.org/zap" - "github.com/mholt/archiver" "github.com/pkg/errors" uuid "github.com/satori/go.uuid" "go.opencensus.io/plugin/ochttp" + "go.uber.org/zap" "golang.org/x/net/context/ctxhttp" k8serr "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" diff --git a/examples/python/sourcepkg/user.py b/examples/python/sourcepkg/user.py index 3c198aca..aef11272 100644 --- a/examples/python/sourcepkg/user.py +++ b/examples/python/sourcepkg/user.py @@ -1,4 +1,3 @@ -import sys import yaml document = """ @@ -9,4 +8,4 @@ document = """ """ def main(): - return yaml.dump(yaml.load(document)) + return yaml.dump(yaml.load(document), default_flow_style=None) diff --git a/executor/executor.go b/executor/executor.go index 18bb40c2..c725992e 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -20,7 +20,6 @@ import ( "context" "fmt" "net/http" - "runtime/debug" "strings" "sync" "time" @@ -193,10 +192,6 @@ func (executor *Executor) isValidAddress(fsvc *fscache.FuncSvc) bool { } } -func dumpStackTrace() { - debug.PrintStack() -} - func serveMetric(logger *zap.Logger) { // Expose the registered metrics via HTTP. metricAddr := ":8080" diff --git a/executor/poolmgr/gp.go b/executor/poolmgr/gp.go index 128478b4..28de9853 100644 --- a/executor/poolmgr/gp.go +++ b/executor/poolmgr/gp.go @@ -163,6 +163,7 @@ func (gp *GenericPool) getDeployLabels() map[string]string { fission.ENVIRONMENT_NAME: gp.env.Metadata.Name, fission.ENVIRONMENT_NAMESPACE: gp.env.Metadata.Namespace, fission.ENVIRONMENT_UID: string(gp.env.Metadata.UID), + "managed": "true", // this allows us to easily find pods managed by the deployment } } @@ -261,12 +262,13 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*apiv1.Pod, erro } func (gp *GenericPool) labelsForFunction(metadata *metav1.ObjectMeta) map[string]string { - return map[string]string{ - "functionName": metadata.Name, - "functionUid": string(metadata.UID), - "unmanaged": "true", // this allows us to easily find pods not managed by the deployment - fission.EXECUTOR_INSTANCEID_LABEL: gp.instanceId, - } + label := gp.getDeployLabels() + label[fission.FUNCTION_NAME] = metadata.Name + label[fission.FUNCTION_UID] = string(metadata.UID) + label[fission.FUNCTION_NAMESPACE] = metadata.Namespace // function CRD must stay within same namespace of environment CRD + label["managed"] = "false" // this allows us to easily find pods not managed by the deployment + return label + } func (gp *GenericPool) scheduleDeletePod(name string) { diff --git a/fission-bundle/main.go b/fission-bundle/main.go index 7ea36de4..81d2410b 100644 --- a/fission-bundle/main.go +++ b/fission-bundle/main.go @@ -16,6 +16,7 @@ import ( "github.com/fission/fission/controller" "github.com/fission/fission/executor" "github.com/fission/fission/kubewatcher" + functionLogger "github.com/fission/fission/logger" "github.com/fission/fission/mqtrigger" "github.com/fission/fission/router" "github.com/fission/fission/storagesvc" @@ -77,6 +78,11 @@ func runBuilderMgr(logger *zap.Logger, storageSvcUrl string, envBuilderNamespace } } +func runLogger() { + functionLogger.Start() + log.Fatalf("Error: Logger exited.") +} + func getPort(logger *zap.Logger, portArg interface{}) int { portArgStr := portArg.(string) port, err := strconv.Atoi(portArgStr) @@ -174,6 +180,7 @@ Usage: fission-bundle --builderMgr [--storageSvcUrl=] [--envbuilder-namespace=] [--collectorEndpoint=] fission-bundle --timer [--routerUrl=] [--collectorEndpoint=] fission-bundle --mqt [--routerUrl=] [--collectorEndpoint=] + fission-bundle --logger fission-bundle --version Options: --collectorEndpoint= Jaeger HTTP Thrift collector URL. @@ -249,6 +256,10 @@ Options: runBuilderMgr(logger, storageSvcUrl, envBuilderNs) } + if arguments["--logger"] == true { + runLogger() + } + if arguments["--storageServicePort"] != nil { port := getPort(logger, arguments["--storageServicePort"]) filePath := arguments["--filePath"].(string) diff --git a/fission/logdb/influxdb.go b/fission/logdb/influxdb.go index 07cbeabf..a5667aab 100644 --- a/fission/logdb/influxdb.go +++ b/fission/logdb/influxdb.go @@ -66,10 +66,12 @@ func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) { //the parameters above are only for the where clause and do not work with LIMIT if filter.Pod != "" { - queryCmd = "select * from \"log\" where \"funcuid\" = $funcuid AND \"pod\" = $pod AND \"time\" > $time LIMIT " + strconv.Itoa(filter.RecordLimit) + // wait for bug fix for fluent-bit influxdb plugin + queryCmd = "select * from /^log*/ where (\"funcuid\" = $funcuid OR \"kubernetes_labels_functionUid\" = $funcuid) AND \"pod\" = $pod AND \"time\" > $time LIMIT " + strconv.Itoa(filter.RecordLimit) parameters["pod"] = filter.Pod } else { - queryCmd = "select * from \"log\" where \"funcuid\" = $funcuid AND \"time\" > $time LIMIT " + strconv.Itoa(filter.RecordLimit) + // wait for bug fix for fluent-bit influxdb plugin + queryCmd = "select * from /^log*/ where (\"funcuid\" = $funcuid OR \"kubernetes_labels_functionUid\" = $funcuid) AND \"time\" > $time LIMIT " + strconv.Itoa(filter.RecordLimit) } query := influxdbClient.NewQueryWithParameters(queryCmd, INFLUXDB_DATABASE, "", parameters) @@ -84,9 +86,13 @@ func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) { //create map of columns to row indeces indexMap := makeIndexMap(series.Columns) - container := indexMap["docker_container_id"] + // TODO: Remove fallback indexes. Some of index's name changed in fluent-bit, here we add extra fallbackIndexes to address compatibility problem. + container := indexMap["kubernetes_docker_id"] + container_1 := indexMap["docker_container_id"] // for backward compatibility functionName := indexMap["kubernetes_labels_functionName"] funcuid := indexMap["kubernetes_labels_functionUid"] + funcuid_1 := indexMap["funcuid"] // for backward compatibility + funcuid_2 := indexMap["kubernetes_labels_functionUid_1"] // for backward compatibility logMessage := indexMap["log"] nameSpace := indexMap["kubernetes_namespace_name"] podName := indexMap["kubernetes_pod_name"] @@ -102,18 +108,19 @@ func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) { if err != nil { return logEntries, err } - logEntries = append(logEntries, LogEntry{ + entry := LogEntry{ //The attributes of the LogEntry are selected as relative to their position in InfluxDB's line protocol response Timestamp: t, - Container: row[container].(string), //docker_container_id - FuncName: row[functionName].(string), //kubernetes_labels_functionName - FuncUid: row[funcuid].(string), //funcuid - Message: strings.TrimSuffix(row[logMessage].(string), "\n"), //log field - Namespace: row[nameSpace].(string), //kubernetes_namespace_name - Pod: row[podName].(string), //kubernetes_pod_name - Stream: row[stream].(string), //stream - Sequence: seqNum, //sequence tag - }) + Container: getEntryValue(row, container, container_1), + FuncName: getEntryValue(row, functionName, -1), + FuncUid: getEntryValue(row, funcuid, funcuid_1, funcuid_2), + Message: strings.TrimSuffix(getEntryValue(row, logMessage, -1), "\n"), //log field + Namespace: getEntryValue(row, nameSpace, -1), + Pod: getEntryValue(row, podName, -1), + Stream: getEntryValue(row, stream, -1), + Sequence: seqNum, //sequence tag + } + logEntries = append(logEntries, entry) } } } @@ -175,3 +182,21 @@ func (influx InfluxDB) query(query influxdbClient.Query) (*influxdbClient.Respon } return &response, nil } + +// getEntryValue returns a field value in string type of log entry by providing index of log entry. +// Since we switch from fluentd to fluent-bit, there are some field names' changed which will break +// CLI due to empty value. For backward compatibility, getEntryValue also supports to get value from +// fallbackIndex if exists, otherwise an empty string returned instead. +func getEntryValue(list []interface{}, index int, fallbackIndex ...int) string { + if index < len(list) && list[index] != nil { + return list[index].(string) + } + + for _, i := range fallbackIndex { + if i >= 0 && i < len(list) && list[i] != nil { + return list[i].(string) + } + } + + return "" +} diff --git a/hack/release-build.sh b/hack/release-build.sh index 60c040de..b33d8710 100755 --- a/hack/release-build.sh +++ b/hack/release-build.sh @@ -96,23 +96,6 @@ build_builder_image() { popd } -build_logger_image() { - local version=$1 - local tag=fission/fluentd:$version - - pushd $DIR/logger/fluentd - - docker build -t $tag . - docker tag $tag fission/fluentd:latest - - popd -} - -push_logger_image() { - local version=$1 - local tag=fission/fluentd:$version -} - build_env_image() { local version=$1 envdir=$2 @@ -294,7 +277,6 @@ build_all() { build_fission_bundle_image $version $date $gitcommit build_fetcher_image $version $date $gitcommit build_builder_image $version $date $gitcommit - build_logger_image $version build_all_cli $version $date $gitcommit build_pre_upgrade_checks_image $version $date $gitcommit } diff --git a/hack/release.sh b/hack/release.sh index 3e6e4627..8fdc35ea 100755 --- a/hack/release.sh +++ b/hack/release.sh @@ -46,12 +46,6 @@ push_builder_image() { docker push $tag } -push_logger_image() { - local version=$1 - local tag=fission/fluentd:$version - docker push $tag -} - push_env_image() { local version=$1 envdir=$2 @@ -141,9 +135,6 @@ push_all() { push_builder_image $version push_builder_image latest - push_logger_image $version - push_logger_image latest - push_pre_upgrade_checks_image $version push_pre_upgrade_checks_image latest } diff --git a/hack/runtests.sh b/hack/runtests.sh index d61e63f3..9c38243b 100755 --- a/hack/runtests.sh +++ b/hack/runtests.sh @@ -1,5 +1,5 @@ #!/bin/bash - +set -euxo pipefail if [ ! -f ${KUBECONFIG} ] then diff --git a/logger/fluentd/Dockerfile b/logger/fluentd/Dockerfile deleted file mode 100644 index 404b69d8..00000000 --- a/logger/fluentd/Dockerfile +++ /dev/null @@ -1,70 +0,0 @@ -# This file originally came from official Kubernetes GitHub repository. -# You can reach original file with the following link: -# https://github.com/kubernetes/kubernetes/tree/42fbf93fb0bb48d0592e2aa08c5ce6d28ab6d4b0/cluster/addons/fluentd-gcp/fluentd-gcp-image - -# Modification: -# 1. add plugin "fluent-plugin-influxdb" for influxdb support - -# Copyright 2016 The Kubernetes 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. - -# This Dockerfile will build an image that is configured -# to use Fluentd to collect all Docker container log files -# and then cause them to be ingested using the Google Cloud -# Logging API. This configuration assumes that the host performning -# the collection is a VM that has been created with a logging.write -# scope and that the Logging API has been enabled for the project -# in the Google Developer Console. - -FROM gcr.io/google_containers/ubuntu-slim:0.6 - - -# Disable prompts from apt -ENV DEBIAN_FRONTEND noninteractive - -# Install build tools -RUN apt-get -qq update && \ - apt-get install -y -qq curl ca-certificates gcc g++ make bash sudo && \ - apt-get install -y -qq --reinstall lsb-base lsb-release && \ - # Install logging agent and required gems - /usr/bin/curl -sSL https://toolbelt.treasuredata.com/sh/install-ubuntu-xenial-td-agent2.sh | sh && \ - sed -i -e "s/USER=td-agent/USER=root/" -e "s/GROUP=td-agent/GROUP=root/" /etc/init.d/td-agent && \ - td-agent-gem install --no-document fluent-plugin-record-reformer -v 0.8.2 && \ - td-agent-gem install --no-document fluent-plugin-systemd -v 0.0.5 && \ - td-agent-gem install --no-document fluent-plugin-google-cloud -v 0.5.2 && \ - td-agent-gem install --no-document fluent-plugin-detect-exceptions -v 0.0.4 && \ - td-agent-gem install --no-document fluent-plugin-influxdb -v 1.1.0 && \ - td-agent-gem install --no-document fluent-plugin-kubernetes_metadata_filter -v 2.1.5 && \ - td-agent-gem install --no-document fluent-plugin-flatten-hash -v 0.5.1 && \ - # Remove build tools - apt-get remove -y -qq gcc make && \ - apt-get autoremove -y -qq && \ - apt-get clean -qq && \ - # Remove unnecessary files - rm -rf /var/lib/apt/lists/* /tmp/* /var/tmp/* \ - /opt/td-agent/embedded/share/doc \ - /opt/td-agent/embedded/share/gtk-doc \ - /opt/td-agent/embedded/lib/postgresql \ - /opt/td-agent/embedded/bin/postgres \ - /opt/td-agent/embedded/share/postgresql \ - /etc/td-agent/td-agent.conf - -# Copy the Fluentd configuration file for logging Docker container logs. -COPY fluent.conf /etc/td-agent/td-agent.conf - -# Copy the entrypoint for the container -COPY run.sh /run.sh - -# Start Fluentd to pick up our config that watches Docker container logs. -CMD /run.sh $FLUENTD_ARGS diff --git a/logger/fluentd/build.sh b/logger/fluentd/build.sh deleted file mode 100755 index 06aa0b98..00000000 --- a/logger/fluentd/build.sh +++ /dev/null @@ -1,2 +0,0 @@ -#!/bin/sh -docker build -t fission-daemonset-fluentd:latest . diff --git a/logger/fluentd/fluent.conf b/logger/fluentd/fluent.conf deleted file mode 100644 index efa64a44..00000000 --- a/logger/fluentd/fluent.conf +++ /dev/null @@ -1,55 +0,0 @@ -# Default logger configfile - which is generally replaced by a more specific logger config at runtime (see charts/) - - type null - - - - type tail - format json - time_key time - path "#{ENV['FLUENTD_PATH']}" - time_format %Y-%m-%dT%H:%M:%S.%NZ - tag fission.* - read_from_head true - refresh_interval 5 - - - - type kubernetes_metadata - - - - - type flatten_hash - separator _ - - - - type record_reformer - enable_ruby false - tag log - - funcuid ${kubernetes_labels_functionUid} - - - - - @type influxdb - host "#{ENV['INFLUXDB_ADDRESS']}" - port "#{ENV['INFLUXDB_PORT']}" - dbname "#{ENV['INFLUXDB_DBNAME']}" - user "#{ENV['INFLUXDB_USERNAME']}" - password "#{ENV['INFLUXDB_PASSWD']}" - use_ssl false - time_precision ns - tag_keys ["funcuid"] - sequence_tag _seq - buffer_type file - buffer_path /var/log/fission/fluentd.buffer - buffer_chunk_limit 128m - buffer_queue_limit 256 - flush_interval 5 - retry_limit 10 - retry_wait 1.0 - num_threads 2 - diff --git a/logger/fluentd/run.sh b/logger/fluentd/run.sh deleted file mode 100755 index d5e3036b..00000000 --- a/logger/fluentd/run.sh +++ /dev/null @@ -1,27 +0,0 @@ -# This file originally came from official Kubernetes GitHub repository. -# You can reach original file with the following link: -# https://github.com/kubernetes/kubernetes/tree/42fbf93fb0bb48d0592e2aa08c5ce6d28ab6d4b0/cluster/addons/fluentd-gcp/fluentd-gcp-image - -#!/bin/sh - -# Copyright 2016 The Kubernetes 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. - -# For systems without journald -mkdir -p /var/log/journal - -LD_PRELOAD=/opt/td-agent/embedded/lib/libjemalloc.so -RUBY_GC_HEAP_OLDOBJECT_LIMIT_FACTOR=0.9 - -/usr/sbin/td-agent $@ diff --git a/logger/logger.go b/logger/logger.go new file mode 100644 index 00000000..d9ee148a --- /dev/null +++ b/logger/logger.go @@ -0,0 +1,181 @@ +/* +Copyright 2018 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 logger + +import ( + "fmt" + "os" + "path/filepath" + "strings" + "time" + + log "github.com/sirupsen/logrus" + "go.uber.org/zap" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/fields" + "k8s.io/client-go/kubernetes" + k8sCache "k8s.io/client-go/tools/cache" + + "github.com/fission/fission" + "github.com/fission/fission/crd" +) + +var nodeName = os.Getenv("NODE_NAME") + +const ( + originalContainerLogPath = "/var/log/containers" + fissionSymlinkPath = "/var/log/fission" +) + +func makePodLoggerController(zapLogger *zap.Logger, k8sClientSet *kubernetes.Clientset) k8sCache.Controller { + resyncPeriod := 30 * time.Second + lw := k8sCache.NewListWatchFromClient(k8sClientSet.CoreV1().RESTClient(), "pods", metav1.NamespaceAll, fields.Everything()) + _, controller := k8sCache.NewInformer(lw, &corev1.Pod{}, resyncPeriod, + k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + pod := obj.(*corev1.Pod) + if !isValidFunctionPodOnNode(pod) || !fission.IsReadyPod(pod) { + return + } + err := createLogSymlinks(zapLogger, pod) + if err != nil { + funcName := pod.Labels[fission.FUNCTION_NAME] + zapLogger.Error("error creating symlink", + zap.String("function", funcName), zap.Error(err)) + } + }, + UpdateFunc: func(_, obj interface{}) { + pod := obj.(*corev1.Pod) + if !isValidFunctionPodOnNode(pod) || !fission.IsReadyPod(pod) { + return + } + err := createLogSymlinks(zapLogger, pod) + if err != nil { + funcName := pod.Labels[fission.FUNCTION_NAME] + zapLogger.Error("error creating symlink", + zap.String("function", funcName), zap.Error(err)) + } + }, + DeleteFunc: func(obj interface{}) { + // Do nothing here, let symlink reaper to recycle orphan symlink file + }, + }) + return controller +} + +func createLogSymlinks(zapLogger *zap.Logger, pod *corev1.Pod) error { + for _, container := range pod.Status.ContainerStatuses { + containerUID, err := parseContainerString(container.ContainerID) + if err != nil { + zapLogger.Error("error parsing container uid", + zap.String("container", container.Name), + zap.String("pod", pod.Name), + zap.String("namespace", pod.Namespace), + zap.Error(err)) + continue + } + containerLogPath := getLogPath(originalContainerLogPath, pod.Name, pod.Namespace, container.Name, containerUID) + symlinkLogPath := getLogPath(fissionSymlinkPath, pod.Name, pod.Namespace, container.Name, containerUID) + + // check whether a symlink exists, if yes then ignore it + if _, err := os.Stat(symlinkLogPath); os.IsNotExist(err) { + err := os.Symlink(containerLogPath, symlinkLogPath) + if err != nil { + zapLogger.Error("error creating symlink", + zap.String("container", container.Name), + zap.String("pod", pod.Name), + zap.String("namespace", pod.Namespace), + zap.Error(err)) + } + } + } + + return nil +} + +// isValidFunctionPodOnNode checks whether a pod is scheduled to the node the logger runs on +// and examines it's metadata labels to ensure it's a qualified function pod. +func isValidFunctionPodOnNode(pod *corev1.Pod) bool { + if pod.Spec.NodeName != nodeName { + return false + } + labels := []string{fission.ENVIRONMENT_NAMESPACE, fission.ENVIRONMENT_NAME, fission.ENVIRONMENT_UID, + fission.FUNCTION_NAMESPACE, fission.FUNCTION_NAME, fission.FUNCTION_UID, fission.EXECUTOR_TYPE} + for _, l := range labels { + if len(pod.Labels[l]) == 0 { + return false + } + } + return true +} + +// The ContainerID is consist of container engine type (docker://) and uuid of container. +// (e.g., docker://f4ca66baaa715030e20273aaf5232635a144165f1cd8e34ca5175064c245b679) +// This function tries to extract container uuid from ContainerID. +func parseContainerString(containerID string) (string, error) { + // Trim the quotes and split the type and ID. + parts := strings.Split(strings.Trim(containerID, "\""), "://") + if len(parts) != 2 { + return "", fmt.Errorf("invalid container ID: %q", containerID) + } + _, ID := parts[0], parts[1] + return ID, nil +} + +func getLogPath(pathPrefix, podName, podNamespace, containerName, containerID string) string { + logName := fmt.Sprintf("%s_%s_%s-%s.log", podName, podNamespace, containerName, containerID) + return filepath.Join(pathPrefix, logName) +} + +// symlinkReaper periodically checks and removes symlink file if it's target container log file is no longer exists. +func symlinkReaper(zapLogger *zap.Logger) { + for { + select { + case <-time.After(5 * time.Minute): + err := filepath.Walk(fissionSymlinkPath, func(path string, info os.FileInfo, err error) error { + if target, e := os.Readlink(path); e == nil { + if _, pathErr := os.Stat(target); os.IsNotExist(pathErr) { + zapLogger.Debug("remove symlink file", zap.String("filepath", path)) + os.Remove(path) + } + } + return nil + }) + if err != nil { + zapLogger.Error("error reaping symlink", zap.Error(err)) + } + } + } +} + +func Start() { + zapLogger, err := zap.NewProduction() + if err != nil { + log.Fatalf("can't initialize zap logger: %v", err) + } + defer zapLogger.Sync() + + go symlinkReaper(zapLogger) + _, kubernetesClient, _, err := crd.MakeFissionClient() + if err != nil { + log.Fatalf("Error starting pod watcher: %v", err) + } + controller := makePodLoggerController(zapLogger, kubernetesClient) + controller.Run(make(chan struct{})) + zapLogger.Fatal("Stop watching pod changes") +} diff --git a/logger/types.go b/logger/types.go deleted file mode 100644 index 51bbd070..00000000 --- a/logger/types.go +++ /dev/null @@ -1,35 +0,0 @@ -/* -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 logger - -import "sync" - -type ( - LogRequest struct { - Namespace string `json:"namespace"` - Pod string `json:"pod"` - Container string `json:"container"` - FuncName string `json:"funcname"` - FuncUid string `json:"funcuid"` - ContainerID string `json:"-"` - } - - logRequestTracker struct { - sync.RWMutex - logMap map[string]LogRequest - } -) diff --git a/test/build_and_test.sh b/test/build_and_test.sh index 75f8bc40..de681c33 100755 --- a/test/build_and_test.sh +++ b/test/build_and_test.sh @@ -19,7 +19,6 @@ source $(dirname $0)/test_utils.sh REPO=gcr.io/fission-ci IMAGE=fission-bundle -FLUENTD_IMAGE=fluentd FETCHER_IMAGE=$REPO/fetcher BUILDER_IMAGE=$REPO/builder TAG=test @@ -46,8 +45,6 @@ build_and_push_env_builder python $REPO/python-env-builder:$TAG $BUILDER_IMAGE:$ build_and_push_env_builder jvm $REPO/jvm-env-builder:$TAG $BUILDER_IMAGE:$TAG build_and_push_env_builder go $REPO/go-env-builder:$TAG $BUILDER_IMAGE:$TAG -build_and_push_fluentd $REPO/$FLUENTD_IMAGE:$TAG - build_fission_cli -install_and_test $REPO $IMAGE $TAG $FETCHER_IMAGE $TAG $FLUENTD_IMAGE $TAG $PRUNE_INTERVAL $ROUTER_SERVICE_TYPE $SERVICE_TYPE $PRE_UPGRADE_CHECK_IMAGE +install_and_test $REPO $IMAGE $TAG $FETCHER_IMAGE $TAG $PRUNE_INTERVAL $ROUTER_SERVICE_TYPE $SERVICE_TYPE $PRE_UPGRADE_CHECK_IMAGE diff --git a/test/test_utils.sh b/test/test_utils.sh index fd8749ed..97dbc45c 100755 --- a/test/test_utils.sh +++ b/test/test_utils.sh @@ -107,19 +107,6 @@ build_and_push_builder() { popd } -build_and_push_fluentd(){ - image_tag=$1 - - pushd $ROOT/logger/fluentd - docker build -q -t $image_tag . - - gcloud_login - - gcloud docker -- push $image_tag - popd - -} - build_and_push_env_runtime() { env=$1 image_tag=$2 @@ -183,17 +170,15 @@ helm_install_fission() { fetcherImageTag=$6 controllerNodeport=$7 routerNodeport=$8 - fluentdImage=$9 - fluentdImageTag=${10} - pruneInterval="${11}" - routerServiceType=${12} - serviceType=${13} - preUpgradeCheckImage=${14} + pruneInterval=$9 + routerServiceType=${10} + serviceType=${11} + preUpgradeCheckImage=${12} ns=f-$id fns=f-func-$id - helmVars=repository=$repo,image=$image,imageTag=$imageTag,fetcherImage=$fetcherImage,fetcherImageTag=$fetcherImageTag,functionNamespace=$fns,controllerPort=$controllerNodeport,routerPort=$routerNodeport,pullPolicy=Always,analytics=false,logger.fluentdImageRepository=$repo,logger.fluentdImage=$fluentdImage,logger.fluentdImageTag=$fluentdImageTag,pruneInterval=$pruneInterval,routerServiceType=$routerServiceType,serviceType=$serviceType,preUpgradeChecksImage=$preUpgradeCheckImage,prometheus.server.persistentVolume.enabled=false,prometheus.alertmanager.enabled=false,prometheus.kubeStateMetrics.enabled=false,prometheus.nodeExporter.enabled=false + helmVars=repository=$repo,image=$image,imageTag=$imageTag,fetcherImage=$fetcherImage,fetcherImageTag=$fetcherImageTag,functionNamespace=$fns,controllerPort=$controllerNodeport,routerPort=$routerNodeport,pullPolicy=Always,analytics=false,pruneInterval=$pruneInterval,routerServiceType=$routerServiceType,serviceType=$serviceType,preUpgradeChecksImage=$preUpgradeCheckImage,prometheus.server.persistentVolume.enabled=false,prometheus.alertmanager.enabled=false,prometheus.kubeStateMetrics.enabled=false,prometheus.nodeExporter.enabled=false timeout 30 bash -c "helm_setup" @@ -495,12 +480,10 @@ install_and_test() { imageTag=$3 fetcherImage=$4 fetcherImageTag=$5 - fluentdImage=$6 - fluentdImageTag=$7 - pruneInterval=$8 - routerServiceType=$9 - serviceType=${10} - preUpgradeCheckImage=${11} + pruneInterval=$6 + routerServiceType=$7 + serviceType=$8 + preUpgradeCheckImage=$9 controllerPort=31234 @@ -510,7 +493,7 @@ install_and_test() { id=$(generate_test_id) trap "helm_uninstall_fission $id" EXIT - helm_install_fission $id $repo $image $imageTag $fetcherImage $fetcherImageTag $controllerPort $routerPort $fluentdImage $fluentdImageTag $pruneInterval $routerServiceType $serviceType $preUpgradeCheckImage + helm_install_fission $id $repo $image $imageTag $fetcherImage $fetcherImageTag $controllerPort $routerPort $pruneInterval $routerServiceType $serviceType $preUpgradeCheckImage helm status $id | grep STATUS | grep -i deployed if [ $? -ne 0 ]; then describe_all_pods $id diff --git a/test/upgrade/fission_upgrade_test.sh b/test/upgrade/fission_upgrade_test.sh index 480aea81..479cae6b 100755 --- a/test/upgrade/fission_upgrade_test.sh +++ b/test/upgrade/fission_upgrade_test.sh @@ -75,7 +75,6 @@ upgrade_tests REPO=gcr.io/fission-ci IMAGE=fission-bundle FETCHER_IMAGE=$REPO/fetcher -FLUENTD_IMAGE=gcr.io/fission-ci/fluentd BUILDER_IMAGE=$REPO/builder TAG=upgrade-test PRUNE_INTERVAL=1 # Unit - Minutes; Controls the interval to run archivePruner. @@ -85,15 +84,13 @@ build_and_push_fission_bundle $IMAGE:$TAG build_and_push_fetcher $FETCHER_IMAGE:$TAG -build_and_push_fluentd $FLUENTD_IMAGE:$TAG - build_fission_cli sudo mv $ROOT/fission/fission /usr/local/bin/ ## Upgrade -helmVars=repository=$repo,image=$IMAGE,imageTag=$TAG,fetcherImage=$FETCHER_IMAGE,fetcherImageTag=$TAG,logger.fluentdImageRepository=$repo,logger.fluentdImage=$FLUENTD_IMAGE,logger.fluentdImageTag=$TAG,functionNamespace=$fns,controllerPort=$controllerNodeport,pullPolicy=Always,analytics=false,pruneInterval=$pruneInterval,routerServiceType=$routerServiceType +helmVars=repository=$repo,image=$IMAGE,imageTag=$TAG,fetcherImage=$FETCHER_IMAGE,fetcherImageTag=$TAG,functionNamespace=$fns,controllerPort=$controllerNodeport,pullPolicy=Always,analytics=false,pruneInterval=$pruneInterval,routerServiceType=$routerServiceType echo "Upgrading fission" helm upgrade \