From 95eef49cf8032dd9f9a2edcb985f91a7b1978e17 Mon Sep 17 00:00:00 2001 From: Nafisa Shazia Date: Fri, 17 Aug 2018 13:17:54 -0700 Subject: [PATCH] Replay recorded requests by ReqUID (#864) --- controller/api.go | 2 + controller/client/replayer.go | 47 ++++++++++++ controller/replayAPI.go | 40 +++++++++++ fission/main.go | 4 ++ fission/replay.go | 50 +++++++++++++ redis/redis.go | 2 + redis/redisApi.go | 130 +++++++++++++++++++++++++++------- 7 files changed, 251 insertions(+), 24 deletions(-) create mode 100644 controller/client/replayer.go create mode 100644 controller/replayAPI.go create mode 100644 fission/replay.go diff --git a/controller/api.go b/controller/api.go index f25b32a2..208e4271 100644 --- a/controller/api.go +++ b/controller/api.go @@ -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") diff --git a/controller/client/replayer.go b/controller/client/replayer.go new file mode 100644 index 00000000..b81b2738 --- /dev/null +++ b/controller/client/replayer.go @@ -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 +} diff --git a/controller/replayAPI.go b/controller/replayAPI.go new file mode 100644 index 00000000..28e44b10 --- /dev/null +++ b/controller/replayAPI.go @@ -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) +} diff --git a/fission/main.go b/fission/main.go index 23a5dff2..5fd29237 100644 --- a/fission/main.go +++ b/fission/main.go @@ -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}, diff --git a/fission/replay.go b/fission/replay.go new file mode 100644 index 00000000..3513e627 --- /dev/null +++ b/fission/replay.go @@ -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 +} diff --git a/redis/redis.go b/redis/redis.go index 79ed5722..cf2c0935 100644 --- a/redis/redis.go +++ b/redis/redis.go @@ -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), diff --git a/redis/redisApi.go b/redis/redisApi.go index 7df08425..c404f88a 100644 --- a/redis/redisApi.go +++ b/redis/redisApi.go @@ -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 +}