Fix fission function logs (#448)
Remove hard coded field indices in influx query, instead searching for them by name.
This commit is contained in:
committed by
Soam Vasani
parent
bed189c715
commit
c03bb01bca
@@ -67,6 +67,15 @@ func (influx InfluxDB) GetPods(filter LogFilter) ([]string, error) {
|
|||||||
return pods, nil
|
return pods, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func makeIndexMap(cols []string) map[string]int {
|
||||||
|
indexMap := make(map[string]int, len(cols))
|
||||||
|
for i := range cols {
|
||||||
|
indexMap[cols[i]] = i
|
||||||
|
}
|
||||||
|
|
||||||
|
return indexMap
|
||||||
|
}
|
||||||
|
|
||||||
func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) {
|
func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) {
|
||||||
timestamp := filter.Since.UnixNano()
|
timestamp := filter.Since.UnixNano()
|
||||||
var queryCmd string
|
var queryCmd string
|
||||||
@@ -93,26 +102,39 @@ func (influx InfluxDB) GetLogs(filter LogFilter) ([]LogEntry, error) {
|
|||||||
}
|
}
|
||||||
for _, r := range response.Results {
|
for _, r := range response.Results {
|
||||||
for _, series := range r.Series {
|
for _, series := range r.Series {
|
||||||
|
|
||||||
|
//create map of columns to row indeces
|
||||||
|
indexMap := makeIndexMap(series.Columns)
|
||||||
|
|
||||||
|
container := indexMap["docker_container_id"]
|
||||||
|
functionName := indexMap["kubernetes_labels_functionName"]
|
||||||
|
funcuid := indexMap["kubernetes_labels_functionUid"]
|
||||||
|
logMessage := indexMap["log"]
|
||||||
|
nameSpace := indexMap["kubernetes_namespace_name"]
|
||||||
|
podName := indexMap["kubernetes_pod_name"]
|
||||||
|
stream := indexMap["stream"]
|
||||||
|
seq := indexMap["_seq"]
|
||||||
|
|
||||||
for _, row := range series.Values {
|
for _, row := range series.Values {
|
||||||
t, err := time.Parse(time.RFC3339, row[0].(string))
|
t, err := time.Parse(time.RFC3339, row[0].(string))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Fatal(err)
|
log.Fatal(err)
|
||||||
}
|
}
|
||||||
seqNum, err := strconv.Atoi(row[1].(string))
|
seqNum, err := strconv.Atoi(row[seq].(string))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return logEntries, err
|
return logEntries, err
|
||||||
}
|
}
|
||||||
logEntries = append(logEntries, LogEntry{
|
logEntries = append(logEntries, LogEntry{
|
||||||
//The attributes of the LogEntry are selected as relative to their position in InfluxDB's line protocol response
|
//The attributes of the LogEntry are selected as relative to their position in InfluxDB's line protocol response
|
||||||
Timestamp: t,
|
Timestamp: t,
|
||||||
Container: row[2].(string), //docker_container_id
|
Container: row[container].(string), //docker_container_id
|
||||||
FuncName: row[8].(string), //kubernetes_labels_functionName
|
FuncName: row[functionName].(string), //kubernetes_labels_functionName
|
||||||
FuncUid: row[3].(string), //funcuid
|
FuncUid: row[funcuid].(string), //funcuid
|
||||||
Message: strings.TrimSuffix(row[17].(string), "\n"), //log field
|
Message: strings.TrimSuffix(row[logMessage].(string), "\n"), //log field
|
||||||
Namespace: row[14].(string), //kubernetes_namespace_name
|
Namespace: row[nameSpace].(string), //kubernetes_namespace_name
|
||||||
Pod: row[15].(string), //kubernetes_pod_name
|
Pod: row[podName].(string), //kubernetes_pod_name
|
||||||
Stream: row[18].(string), //stream
|
Stream: row[stream].(string), //stream
|
||||||
Sequence: seqNum, //sequence tag
|
Sequence: seqNum, //sequence tag
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,6 +1,5 @@
|
|||||||
#!/bin/bash
|
#!/bin/bash
|
||||||
|
|
||||||
#test:disabled
|
|
||||||
|
|
||||||
set -euo pipefail
|
set -euo pipefail
|
||||||
|
|
||||||
@@ -12,6 +11,8 @@ function cleanup {
|
|||||||
echo "Cleanup route"
|
echo "Cleanup route"
|
||||||
var=$(fission route list | grep $fn | awk '{print $1;}')
|
var=$(fission route list | grep $fn | awk '{print $1;}')
|
||||||
fission route delete --name $var
|
fission route delete --name $var
|
||||||
|
echo "delete logfile"
|
||||||
|
rm "/tmp/logfile"
|
||||||
}
|
}
|
||||||
|
|
||||||
# Create a hello world function in nodejs, test it with an http trigger
|
# Create a hello world function in nodejs, test it with an http trigger
|
||||||
@@ -31,7 +32,7 @@ fission route create --function $fn --url /$fn --method GET
|
|||||||
trap cleanup EXIT
|
trap cleanup EXIT
|
||||||
|
|
||||||
echo "Waiting for router to catch up"
|
echo "Waiting for router to catch up"
|
||||||
sleep 3
|
sleep 15
|
||||||
|
|
||||||
echo "Doing 4 HTTP GETs on the function's route"
|
echo "Doing 4 HTTP GETs on the function's route"
|
||||||
for i in 1 2 3 4
|
for i in 1 2 3 4
|
||||||
@@ -43,11 +44,18 @@ echo "Grabbing logs, should have 4 calls in logs"
|
|||||||
|
|
||||||
sleep 15
|
sleep 15
|
||||||
|
|
||||||
logs=$(fission function logs --name $fn --detail)
|
fission function logs --name $fn --detail > /tmp/logfile
|
||||||
|
|
||||||
|
size=$(wc -c </tmp/logfile)
|
||||||
|
if [ $size = 0 ]
|
||||||
|
then
|
||||||
|
fission function logs --name $fn --detail > /tmp/logfile
|
||||||
|
fi
|
||||||
|
|
||||||
echo "---function logs---"
|
echo "---function logs---"
|
||||||
echo $logs
|
cat /tmp/logfile
|
||||||
echo "------"
|
echo "------"
|
||||||
num=(cat "$logs" | grep 'log test' | wc -l)
|
num=$(grep 'log test' /tmp/logfile | wc -l)
|
||||||
echo $num logs found
|
echo $num logs found
|
||||||
|
|
||||||
if [ $num -ne 4 ]
|
if [ $num -ne 4 ]
|
||||||
|
|||||||
Reference in New Issue
Block a user