Replay recorded requests by ReqUID (#864)

This commit is contained in:
Nafisa Shazia
2018-08-17 13:17:54 -07:00
committed by smruthi2187
parent 2d01382a82
commit 95eef49cf8
7 changed files with 251 additions and 24 deletions
+2
View File
@@ -248,6 +248,8 @@ func (api *API) Serve(port int) {
r.HandleFunc("/v2/records/trigger/{trigger}", api.RecordsApiFilterByTrigger).Methods("GET")
r.HandleFunc("/v2/records/time", api.RecordsApiFilterByTime).Methods("GET")
r.HandleFunc("/v2/replay/{reqUID}", api.ReplayByReqUID).Methods("GET")
r.HandleFunc("/v2/secrets/{secret}", api.SecretGet).Methods("GET")
r.HandleFunc("/v2/configmaps/{configmap}", api.ConfigMapGet).Methods("GET")
+47
View File
@@ -0,0 +1,47 @@
/*
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 client
import (
"encoding/json"
"fmt"
"net/http"
)
func (c *Client) ReplayByReqUID(reqUID string) ([]string, error) {
relativeUrl := fmt.Sprintf("replay/%v", reqUID)
resp, err := http.Get(c.url(relativeUrl))
if err != nil {
return nil, err
}
defer resp.Body.Close()
body, err := c.handleResponse(resp)
if err != nil {
return nil, err
}
replayed := make([]string, 0)
err = json.Unmarshal(body, &replayed)
if err != nil {
return nil, err
}
return replayed, nil
}
+40
View File
@@ -0,0 +1,40 @@
/*
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 controller
import (
"fmt"
"net/http"
"github.com/gorilla/mux"
"github.com/fission/fission/redis"
)
func (a *API) ReplayByReqUID(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
queriedID := vars["reqUID"]
routerUrl := fmt.Sprintf("http://router.%v", podNamespace)
resp, err := redis.ReplayByReqUID(routerUrl, queriedID)
if err != nil {
a.respondWithError(w, err)
return
}
a.respondWithSuccess(w, resp)
}
+4
View File
@@ -227,6 +227,9 @@ func main() {
{Name: "view", Usage: "View existing records", Flags: []cli.Flag{filterTimeTo, filterTimeFrom, filterFunction, filterTrigger, verbosityFlag, vvFlag}, Action: recordsView},
}
// Replay records
reqIDFlag := cli.StringFlag{Name: "reqUID", Usage: "Replay a particular request by providing the reqUID (to view reqUIDs, do 'fission records view')"}
// environments
envNameFlag := cli.StringFlag{Name: "name", Usage: "Environment name"}
envPoolsizeFlag := cli.IntFlag{Name: "poolsize", Value: 3, Usage: "Size of the pool"}
@@ -307,6 +310,7 @@ func main() {
{Name: "mqtrigger", Aliases: []string{"mqt", "messagequeue"}, Usage: "Manage message queue triggers for functions", Subcommands: mqtSubcommands},
{Name: "recorder", Usage: "Manage recorders for functions", Subcommands: recSubcommands, Hidden: true},
{Name: "records", Usage: "View records with optional filters", Subcommands: recViewSubcommands, Hidden: true},
{Name: "replay", Usage: "Replay records", Flags: []cli.Flag{reqIDFlag}, Action: replay},
{Name: "environment", Aliases: []string{"env"}, Usage: "Manage environments", Subcommands: envSubcommands},
{Name: "watch", Aliases: []string{"w"}, Usage: "Manage watches", Subcommands: wSubCommands},
{Name: "package", Aliases: []string{"pkg"}, Usage: "Manage packages", Subcommands: pkgSubCommands},
+50
View File
@@ -0,0 +1,50 @@
/*
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 main
import (
"fmt"
"log"
"os"
"text/tabwriter"
"github.com/urfave/cli"
)
func replay(c *cli.Context) error {
fc := getClient(c.GlobalString("server"))
reqUID := c.String("reqUID")
if len(reqUID) == 0 {
log.Fatal("Need a reqUID, use --reqUID flag to specify")
}
responses, err := fc.ReplayByReqUID(reqUID)
checkErr(err, "replay records")
w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0)
for _, resp := range responses {
fmt.Fprintf(w, "%v",
resp,
)
}
w.Flush()
return nil
}
+2
View File
@@ -92,6 +92,8 @@ func Record(triggerName string, recorderName string, reqUID string, request *htt
PostForm: postForm,
}
log.Info("Storing request > ", req)
resp := &redisCache.Response{
Status: response.Status,
StatusCode: int32(response.StatusCode),
+106 -24
View File
@@ -17,8 +17,12 @@ limitations under the License.
package redis
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"io/ioutil"
"net/http"
"strconv"
"strings"
"time"
@@ -34,7 +38,7 @@ import (
func RecordsListAll() ([]byte, error) {
client := NewClient()
if client == nil {
return []byte{}, errors.New("failed to create redis client")
return nil, errors.New("failed to create redis client")
}
iter := 0
@@ -45,7 +49,7 @@ func RecordsListAll() ([]byte, error) {
// Redis tells us there are no keys left to traverse.
arr, err := redis.Values(client.Do("SCAN", iter))
if err != nil {
return []byte{}, err
return nil, err
}
// SCAN return value is an array of two values: the first value is the new cursor to use in the next call,
// the second value is an array of elements.
@@ -56,12 +60,12 @@ func RecordsListAll() ([]byte, error) {
val, err := redis.Bytes(client.Do("HGET", key, "ReqResponse"))
if err != nil {
log.Error("Error retrieving request from Redis: ", err)
return []byte{}, err
return nil, err
}
entry, err := deserializeReqResponse(val, key)
if err != nil {
log.Error("Error deserializing request: ", err)
return []byte{}, err
return nil, err
}
filtered = append(filtered, entry)
}
@@ -73,7 +77,7 @@ func RecordsListAll() ([]byte, error) {
resp, err := json.Marshal(filtered)
if err != nil {
return []byte{}, err
return nil, err
}
return resp, nil
}
@@ -86,12 +90,12 @@ func RecordsFilterByTime(from string, to string) ([]byte, error) {
if rangeStart >= rangeEnd {
log.Error("Invalid chronology")
return []byte{}, err
return nil, err
}
client := NewClient()
if client == nil {
return []byte{}, errors.New("failed to create redis client")
return nil, errors.New("failed to create redis client")
}
iter := 0
@@ -100,7 +104,7 @@ func RecordsFilterByTime(from string, to string) ([]byte, error) {
for {
arr, err := redis.Values(client.Do("SCAN", iter))
if err != nil {
return []byte{}, err
return nil, err
}
// SCAN return value is an array of two values: the first value is the new cursor to use in the next call,
// the second value is an array of elements.
@@ -111,24 +115,24 @@ func RecordsFilterByTime(from string, to string) ([]byte, error) {
val, err := redis.Strings(client.Do("HMGET", key, "Timestamp"))
if err != nil {
log.Error("Error retrieving timestamp from Redis: ", err)
return []byte{}, err
return nil, err
}
tsO, err := strconv.Atoi(val[0])
if err != nil {
log.Error("Error converting timestamp to int: ", err)
return []byte{}, err
return nil, err
}
ts := int64(tsO)
if ts >= rangeStart && ts <= rangeEnd {
val2, err := redis.Bytes(client.Do("HGET", key, "ReqResponse"))
if err != nil {
log.Error("Error retrieving request from Redis: ", err)
return []byte{}, err
return nil, err
}
entry, err := deserializeReqResponse(val2, key)
if err != nil {
log.Error("Error deserializing request: ", err)
return []byte{}, err
return nil, err
}
filtered = append(filtered, entry)
}
@@ -142,7 +146,7 @@ func RecordsFilterByTime(from string, to string) ([]byte, error) {
resp, err := json.Marshal(filtered)
if err != nil {
return []byte{}, err
return nil, err
}
return resp, nil
}
@@ -176,7 +180,7 @@ func RecordsFilterByTrigger(queriedTriggerName string, recorders *crd.RecorderLi
client := NewClient()
if client == nil {
return []byte{}, errors.New("failed to create redis client")
return nil, errors.New("failed to create redis client")
}
var filtered []*redisCache.RecordedEntry
@@ -186,25 +190,25 @@ func RecordsFilterByTrigger(queriedTriggerName string, recorders *crd.RecorderLi
val, err := redis.Strings(client.Do("LRANGE", key, "0", "-1")) // TODO: Prefix that distinguishes recorder lists
if err != nil {
// TODO: Handle deleted recorder? Or is this a non-issue because our list of recorders is up to date?
return []byte{}, err
return nil, err
}
for _, reqUID := range val {
val, err := redis.Strings(client.Do("HMGET", reqUID, "Trigger")) // 1-to-1 reqUID - trigger?
if err != nil {
log.Error("Error retrieving trigger for a request from Redis: ", err)
return []byte{}, err
return nil, err
}
if val[0] == queriedTriggerName {
// TODO: Reconsider multiple commands
val, err := redis.Bytes(client.Do("HGET", reqUID, "ReqResponse"))
if err != nil {
log.Error("Error retrieving request from Redis: ", err)
return []byte{}, err
return nil, err
}
entry, err := deserializeReqResponse(val, reqUID)
if err != nil {
log.Error("Error deserializing request: ", err)
return []byte{}, err
return nil, err
}
filtered = append(filtered, entry)
}
@@ -213,7 +217,7 @@ func RecordsFilterByTrigger(queriedTriggerName string, recorders *crd.RecorderLi
resp, err := json.Marshal(filtered)
if err != nil {
return []byte{}, err
return nil, err
}
return resp, nil
}
@@ -249,7 +253,7 @@ func RecordsFilterByFunction(queriedFunctionName string, recorders *crd.Recorder
client := NewClient()
if client == nil {
return []byte{}, errors.New("failed to create redis client")
return nil, errors.New("failed to create redis client")
}
var filtered []*redisCache.RecordedEntry
@@ -257,7 +261,7 @@ func RecordsFilterByFunction(queriedFunctionName string, recorders *crd.Recorder
for key := range matchingRecorders {
val, err := redis.Strings(client.Do("LRANGE", key, "0", "-1")) // TODO: Prefix that distinguishes recorder lists
if err != nil {
return []byte{}, err
return nil, err
}
for _, reqUID := range val {
@@ -270,12 +274,12 @@ func RecordsFilterByFunction(queriedFunctionName string, recorders *crd.Recorder
val, err := redis.Bytes(client.Do("HGET", reqUID, "ReqResponse"))
if err != nil {
log.Error("Error retrieving request from Redis: ", err)
return []byte{}, err
return nil, err
}
entry, err := deserializeReqResponse(val, reqUID)
if err != nil {
log.Error("Error deserializing request: ", err)
return []byte{}, err
return nil, err
}
filtered = append(filtered, entry)
}
@@ -284,7 +288,7 @@ func RecordsFilterByFunction(queriedFunctionName string, recorders *crd.Recorder
resp, err := json.Marshal(filtered)
if err != nil {
return []byte{}, err
return nil, err
}
return resp, nil
}
@@ -336,3 +340,81 @@ func includesTrigger(triggers []string, query string) bool {
}
return false
}
func ReplayByReqUID(routerUrl string, queriedID string) ([]byte, error) {
client := NewClient()
if client == nil {
return nil, errors.New("failed to create redis client")
}
exists, err := redis.Int(client.Do("EXISTS", queriedID))
if exists != 1 || err != nil {
log.Error("couldn't find request to replay")
return nil, err
}
val, err := redis.Bytes(client.Do("HGET", queriedID, "ReqResponse"))
if err != nil {
log.Error("couldn't obtain ReqResponse for this ID")
return nil, err
}
entry, err := deserializeReqResponse(val, queriedID)
if err != nil {
log.Error("couldn't deserialize ReqResponse")
return nil, err
}
replayed, err := ReplayRequest(routerUrl, entry.Req)
if err != nil {
log.Error("couldn't replay request")
return nil, err
}
resp, err := json.Marshal(replayed)
if err != nil {
log.Error("couldn't marshall replayed request response")
return nil, err
}
return resp, nil
}
func ReplayRequest(routerUrl string, request *redisCache.Request) ([]string, error) {
path := request.URL["Path"] // Includes slash prefix
payload := request.URL["Payload"]
targetUrl := fmt.Sprintf("%v%v", routerUrl, path)
var req *http.Request
var err error
client := http.DefaultClient
if request.Method == http.MethodGet {
req, err = http.NewRequest("GET", targetUrl, nil)
if err != nil {
return nil, err
}
} else {
req, err = http.NewRequest(request.Method, targetUrl, bytes.NewReader([]byte(payload)))
if err != nil {
return nil, err
}
}
req.Header.Add("X-Fission-Replayed", "true")
resp, err := client.Do(req)
if err != nil {
return nil, errors.New(fmt.Sprintf("failed to make request: %v", err))
}
defer resp.Body.Close()
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, errors.New(fmt.Sprintf("failed to read response: %v", err))
}
bodyStr := string(body)
return []string{bodyStr}, nil
}