From 1d156899d7d8a83301aba91e0b9013c952502752 Mon Sep 17 00:00:00 2001 From: Peter Grlica Date: Thu, 17 Mar 2022 15:54:23 +0100 Subject: [PATCH] Moved filter, sorter to system profiler service --- automation/automation/http_request_handler.go | 16 +- pkg/apigw/profiler/filter.go | 79 ---- pkg/apigw/profiler/profiler.go | 8 - pkg/apigw/profiler/sort.go | 83 ----- pkg/apigw/route.go | 4 - system/rest/apigw_profiler.go | 2 +- system/service/apigw_profiler.go | 338 ++++++++++++++++++ system/service/apigw_route.go | 163 --------- system/service/service.go | 2 + 9 files changed, 342 insertions(+), 353 deletions(-) delete mode 100644 pkg/apigw/profiler/filter.go create mode 100644 system/service/apigw_profiler.go diff --git a/automation/automation/http_request_handler.go b/automation/automation/http_request_handler.go index 79762d8bd..6f39f638d 100644 --- a/automation/automation/http_request_handler.go +++ b/automation/automation/http_request_handler.go @@ -9,7 +9,6 @@ import ( "net/http" "net/url" "strings" - "time" "github.com/cortezaproject/corteza-server/pkg/version" ) @@ -35,8 +34,6 @@ func (h httpRequestHandler) send(ctx context.Context, args *httpRequestSendArgs) rsp *http.Response ) - // spew.Dump("TIMEOUT", args.Timeout) - r = &httpRequestSendResults{} req, err = h.makeRequest(ctx, args) @@ -44,16 +41,7 @@ func (h httpRequestHandler) send(ctx context.Context, args *httpRequestSendArgs) return nil, err } - // spew.Dump("ERR, REQ", err) - // spew.Dump(httputil.DumpRequestOut(req, false)) - - client := http.DefaultClient - client.Timeout = 5 * time.Minute - - rsp, err = client.Do(req) - // spew.Dump("ERR", err) - // spew.Dump("RESP") - // spew.Dump(httputil.DumpResponse(rsp, true)) + rsp, err = http.DefaultClient.Do(req) if err != nil { return } @@ -141,8 +129,6 @@ func (h httpRequestHandler) makeRequest(ctx context.Context, args *httpRequestSe args.Url = purl.String() } - // spew.Dump("CTX", ctx) - req, err = http.NewRequestWithContext(ctx, args.Method, args.Url, args.bodyStream) if err != nil { return nil, err diff --git a/pkg/apigw/profiler/filter.go b/pkg/apigw/profiler/filter.go deleted file mode 100644 index a5eaf23d8..000000000 --- a/pkg/apigw/profiler/filter.go +++ /dev/null @@ -1,79 +0,0 @@ -package profiler - -import ( - "encoding/base64" - - "github.com/cortezaproject/corteza-server/system/types" -) - -func FilterAggregation(list *types.ApigwProfilerAggregationSet, filter *types.ApigwProfilerFilter) { - var ( - dec string = "" - i uint = 0 - b = filter.Before == "" - ) - - if filter.Limit == 0 { - filter.Limit = FILTER_NUM_AGG_ITEMS - } - - dec, _ = decodeRoutePath(filter.Before) - - *list, _ = list.Filter(func(apa *types.ApigwProfilerAggregation) (bool, error) { - // after a specific hit and inside the limits - if b && i < filter.Limit { - i++ - filter.Next = encodeRoutePath(apa.Path) - return true, nil - } - - // after the specific hit check - if dec != "" && b == false { - b = apa.Path == dec - } - - return false, nil - }) - - return -} - -func FilterHits(list *types.ApigwProfilerHitSet, filter *types.ApigwProfilerFilter) { - var ( - i uint = 0 - b = filter.Before == "" - ) - - if filter.Limit == 0 { - filter.Limit = FILTER_NUM_ITEMS - } - - *list, _ = list.Filter(func(aph *types.ApigwProfilerHit) (bool, error) { - // after a specific hit and inside the limits - if b && i < filter.Limit { - i++ - filter.Next = aph.ID - return true, nil - } - - // after the specific hit check - if filter.Before != "" && b == false { - b = aph.ID == filter.Before - } - - return false, nil - }) - - return -} - -func encodeRoutePath(p string) string { - return base64.URLEncoding.EncodeToString([]byte(p)) -} - -func decodeRoutePath(p string) (s string, err error) { - b, err := base64.URLEncoding.DecodeString(p) - s = string(b) - - return -} diff --git a/pkg/apigw/profiler/profiler.go b/pkg/apigw/profiler/profiler.go index e198b8658..c3a5a3219 100644 --- a/pkg/apigw/profiler/profiler.go +++ b/pkg/apigw/profiler/profiler.go @@ -9,14 +9,6 @@ import ( h "github.com/cortezaproject/corteza-server/pkg/http" ) -const ( - // default fallback on amount of items - FILTER_NUM_ITEMS = 20 - - // default fallback on amount of aggregated items - FILTER_NUM_AGG_ITEMS = 10 -) - type ( Hits map[string][]*Hit diff --git a/pkg/apigw/profiler/sort.go b/pkg/apigw/profiler/sort.go index 9de144f1f..e908a8e7f 100644 --- a/pkg/apigw/profiler/sort.go +++ b/pkg/apigw/profiler/sort.go @@ -1,16 +1,5 @@ package profiler -import ( - "sort" - - "github.com/cortezaproject/corteza-server/system/types" -) - -var ( - sortAggFields = []string{"path", "count", "size_min", "size_max", "size_avg", "time_min", "time_max", "time_avg"} - sortRouteFields = []string{"time_start", "time_finish", "time_duration"} -) - type ( Sort struct { Hit string @@ -19,75 +8,3 @@ type ( Before string } ) - -func SortAggregation(list *types.ApigwProfilerAggregationSet, filter *types.ApigwProfilerFilter) { - for _, ff := range sortAggFields { - fe := filter.Sort.Get(ff) - - if fe == nil { - continue - } - - if filter.Sort.Get(ff).Descending { - sort.Sort(sort.Reverse(getSortType(ff, list))) - break - } - - sort.Sort(getSortType(ff, list)) - break - } -} - -func SortHits(list *types.ApigwProfilerHitSet, filter *types.ApigwProfilerFilter) { - for _, ff := range sortRouteFields { - fe := filter.Sort.Get(ff) - - if fe == nil { - continue - } - - if filter.Sort.Get(ff).Descending { - sort.Sort(sort.Reverse(getSortTypeHit(ff, list))) - break - } - - sort.Sort(getSortTypeHit(ff, list)) - break - } -} - -func getSortType(s string, list *types.ApigwProfilerAggregationSet) sort.Interface { - switch s { - case "path": - return types.ByPath(*list) - case "count": - return types.ByCount(*list) - case "size_min": - return types.BySizeMin(*list) - case "size_max": - return types.BySizeMax(*list) - case "size_avg": - return types.BySizeAvg(*list) - case "time_min": - return types.ByTimeMin(*list) - case "time_max": - return types.ByTimeMax(*list) - case "time_avg": - return types.ByTimeAvg(*list) - default: - return types.ByCount(*list) - } -} - -func getSortTypeHit(s string, list *types.ApigwProfilerHitSet) sort.Interface { - switch s { - case "time_start": - return types.BySTime(*list) - case "time_finish": - return types.ByFTime(*list) - case "time_duration": - return types.ByDuration(*list) - default: - return types.BySTime(*list) - } -} diff --git a/pkg/apigw/route.go b/pkg/apigw/route.go index 8aaae1ccc..aeb18a9de 100644 --- a/pkg/apigw/route.go +++ b/pkg/apigw/route.go @@ -11,7 +11,6 @@ import ( "github.com/cortezaproject/corteza-server/pkg/auth" h "github.com/cortezaproject/corteza-server/pkg/http" "github.com/cortezaproject/corteza-server/pkg/options" - "github.com/davecgh/go-spew/spew" "go.uber.org/zap" ) @@ -55,11 +54,8 @@ func (r route) ServeHTTP(w http.ResponseWriter, req *http.Request) { scope.Set("opts", r.opts) scope.Set("request", ar) - spew.Dump("OPTS", r.opts) - // use profiler, override any profiling prefilter if r.opts.ProfilerEnabled && r.opts.ProfilerGlobal { - spew.Dump("adding to profiler") // add request to profiler hit = r.pr.Hit(ar) hit.Route = r.ID diff --git a/system/rest/apigw_profiler.go b/system/rest/apigw_profiler.go index e8de767f0..b260ab651 100644 --- a/system/rest/apigw_profiler.go +++ b/system/rest/apigw_profiler.go @@ -41,7 +41,7 @@ type ( func (ApigwProfiler) New() *ApigwProfiler { return &ApigwProfiler{ - svc: service.DefaultApigwRoute, + svc: service.DefaultApigwProfiler, ac: service.DefaultAccessControl, } } diff --git a/system/service/apigw_profiler.go b/system/service/apigw_profiler.go new file mode 100644 index 000000000..cac98dd9d --- /dev/null +++ b/system/service/apigw_profiler.go @@ -0,0 +1,338 @@ +package service + +import ( + "context" + "encoding/base64" + "errors" + "io/ioutil" + "math" + "sort" + "time" + + "github.com/cortezaproject/corteza-server/pkg/apigw" + "github.com/cortezaproject/corteza-server/pkg/apigw/profiler" + + "github.com/cortezaproject/corteza-server/system/types" +) + +var ( + sortAggFields = []string{"path", "count", "size_min", "size_max", "size_avg", "time_min", "time_max", "time_avg"} + sortRouteFields = []string{"time_start", "time_finish", "time_duration"} +) + +const ( + // default fallback on amount of items + FILTER_NUM_ITEMS = 20 + + // default fallback on amount of aggregated items + FILTER_NUM_AGG_ITEMS = 10 +) + +type ( + apigwProfiler struct{} +) + +func Profiler() *apigwProfiler { + return &apigwProfiler{} +} + +// HitsAggregated fetches a list of hits from integration gateway profiler +func (svc *apigwProfiler) Hits(ctx context.Context, filter types.ApigwProfilerFilter) (r types.ApigwProfilerHitSet, f types.ApigwProfilerFilter, err error) { + + f = filter + r = make(types.ApigwProfilerHitSet, 0) + + uDec, err := base64.URLEncoding.DecodeString(filter.Path) + + if err != nil { + return + } + + filter.Path = string(uDec) + + if filter.Path == "" && filter.Hit == "" { + err = errors.New("fetching all hits (no route and hit specified) not supported") + return + } + + var sorting = profiler.Sort{ + Hit: filter.Hit, + Path: filter.Path, + Before: filter.Before, + } + + var ( + list = apigw.Service().Profiler().Hits(sorting) + ) + + var pp = "" + + for k, _ := range list { + if filter.Hit != "" { + pp = k + break + } + + if filter.Path != "" && k == filter.Path { + pp = k + break + } + } + + if pp == "" { + return + } + + for _, h := range list[pp] { + hh := &types.ApigwProfilerHit{ + ID: h.ID, + + Route: h.Route, + Status: h.Status, + Request: *h.R, + + Ts: h.Ts, + Tf: h.Tf, + D: h.D, + Dr: float64(h.D.Microseconds()) / 1000, + } + + // fetch body only on hit details + if filter.Hit != "" { + hh.Body, _ = ioutil.ReadAll(hh.Request.Body) + } + + r = append(r, hh) + } + + // sort + sortHits(&r, &f) + + // filter sorted + filterHits(&r, &f) + + return +} + +// HitsAggregated fetches a list of hits from integration gateway profiler +// and aggregates them with assigned filters +func (svc *apigwProfiler) HitsAggregated(ctx context.Context, filter types.ApigwProfilerFilter) (r types.ApigwProfilerAggregationSet, f types.ApigwProfilerFilter, err error) { + f = filter + r = make(types.ApigwProfilerAggregationSet, 0) + + uDec, err := base64.URLEncoding.DecodeString(filter.Path) + + if err != nil { + return + } + + filter.Path = string(uDec) + + var ( + list = apigw.Service().Profiler().Hits(profiler.Sort{ + Path: filter.Path, + Before: filter.Before, + }) + + tsum, tmin, tmax time.Duration + ssum, smin, smax int64 + i uint64 = 1 + ) + + for p, v := range list { + tmin, tmax, tsum = time.Hour, 0, 0 + smin, smax, ssum = math.MaxInt64, 0, 0 + + i = 0 + + for _, vv := range v { + var ( + d = *vv.D + s = vv.R.ContentLength + ) + + if d < tmin { + tmin = d + } + + if d > tmax { + tmax = d + } + + if s < smin { + smin = s + } + + if s > smax { + smax = s + } + + tsum += d + ssum += s + i++ + } + + r = append(r, &types.ApigwProfilerAggregation{ + Path: p, + Count: i, + Tmin: float64(tmin.Microseconds()) / 1000, + Tmax: float64(tmax.Microseconds()) / 1000, + Tavg: float64(tsum.Microseconds()) / float64(i) / 1000, + Smin: smin, + Smax: smax, + Savg: float64(ssum) / float64(i), + }) + } + + // sort + sortAggregation(&r, &f) + + // filter + filterAggregation(&r, &f) + + return +} + +func sortAggregation(list *types.ApigwProfilerAggregationSet, filter *types.ApigwProfilerFilter) { + for _, ff := range sortAggFields { + fe := filter.Sort.Get(ff) + + if fe == nil { + continue + } + + if filter.Sort.Get(ff).Descending { + sort.Sort(sort.Reverse(getSortType(ff, list))) + break + } + + sort.Sort(getSortType(ff, list)) + break + } +} + +func sortHits(list *types.ApigwProfilerHitSet, filter *types.ApigwProfilerFilter) { + for _, ff := range sortRouteFields { + fe := filter.Sort.Get(ff) + + if fe == nil { + continue + } + + if filter.Sort.Get(ff).Descending { + sort.Sort(sort.Reverse(getSortTypeHit(ff, list))) + break + } + + sort.Sort(getSortTypeHit(ff, list)) + break + } +} + +func filterAggregation(list *types.ApigwProfilerAggregationSet, filter *types.ApigwProfilerFilter) { + var ( + dec string = "" + i uint = 0 + b = filter.Before == "" + ) + + if filter.Limit == 0 { + filter.Limit = FILTER_NUM_AGG_ITEMS + } + + dec, _ = decodeRoutePath(filter.Before) + + *list, _ = list.Filter(func(apa *types.ApigwProfilerAggregation) (bool, error) { + // after a specific hit and inside the limits + if b && i < filter.Limit { + i++ + filter.Next = encodeRoutePath(apa.Path) + return true, nil + } + + // after the specific hit check + if dec != "" && b == false { + b = apa.Path == dec + } + + return false, nil + }) + + return +} + +func filterHits(list *types.ApigwProfilerHitSet, filter *types.ApigwProfilerFilter) { + var ( + i uint = 0 + b = filter.Before == "" + ) + + if filter.Limit == 0 { + filter.Limit = FILTER_NUM_ITEMS + } + + *list, _ = list.Filter(func(aph *types.ApigwProfilerHit) (bool, error) { + // after a specific hit and inside the limits + if b && i < filter.Limit { + i++ + filter.Next = aph.ID + return true, nil + } + + // after the specific hit check + if filter.Before != "" && b == false { + b = aph.ID == filter.Before + } + + return false, nil + }) + + return +} + +func encodeRoutePath(p string) string { + return base64.URLEncoding.EncodeToString([]byte(p)) +} + +func decodeRoutePath(p string) (s string, err error) { + b, err := base64.URLEncoding.DecodeString(p) + s = string(b) + + return +} + +func getSortType(s string, list *types.ApigwProfilerAggregationSet) sort.Interface { + switch s { + case "path": + return types.ByPath(*list) + case "count": + return types.ByCount(*list) + case "size_min": + return types.BySizeMin(*list) + case "size_max": + return types.BySizeMax(*list) + case "size_avg": + return types.BySizeAvg(*list) + case "time_min": + return types.ByTimeMin(*list) + case "time_max": + return types.ByTimeMax(*list) + case "time_avg": + return types.ByTimeAvg(*list) + default: + return types.ByCount(*list) + } +} + +func getSortTypeHit(s string, list *types.ApigwProfilerHitSet) sort.Interface { + switch s { + case "time_start": + return types.BySTime(*list) + case "time_finish": + return types.ByFTime(*list) + case "time_duration": + return types.ByDuration(*list) + default: + return types.BySTime(*list) + } +} diff --git a/system/service/apigw_route.go b/system/service/apigw_route.go index 47efef18e..552fced59 100644 --- a/system/service/apigw_route.go +++ b/system/service/apigw_route.go @@ -2,15 +2,9 @@ package service import ( "context" - "encoding/base64" - "errors" - "io/ioutil" - "math" - "time" "github.com/cortezaproject/corteza-server/pkg/actionlog" "github.com/cortezaproject/corteza-server/pkg/apigw" - "github.com/cortezaproject/corteza-server/pkg/apigw/profiler" a "github.com/cortezaproject/corteza-server/pkg/auth" "github.com/cortezaproject/corteza-server/store" @@ -259,160 +253,3 @@ func (svc *apigwRoute) Search(ctx context.Context, filter types.ApigwRouteFilter return r, f, svc.recordAction(ctx, aProps, ApigwRouteActionSearch, err) } - -// HitsAggregated fetches a list of hits from integration gateway profiler -func (svc *apigwRoute) Hits(ctx context.Context, filter types.ApigwProfilerFilter) (r types.ApigwProfilerHitSet, f types.ApigwProfilerFilter, err error) { - - f = filter - r = make(types.ApigwProfilerHitSet, 0) - - uDec, err := base64.URLEncoding.DecodeString(filter.Path) - - if err != nil { - return - } - - filter.Path = string(uDec) - - if filter.Path == "" && filter.Hit == "" { - err = errors.New("fetching all hits (no route and hit specified) not supported") - return - } - - var sorting = profiler.Sort{ - Hit: filter.Hit, - Path: filter.Path, - Before: filter.Before, - } - - var ( - list = apigw.Service().Profiler().Hits(sorting) - ) - - var pp = "" - - for k, _ := range list { - if filter.Hit != "" { - pp = k - break - } - - if filter.Path != "" && k == filter.Path { - pp = k - break - } - } - - if pp == "" { - return - } - - for _, h := range list[pp] { - hh := &types.ApigwProfilerHit{ - ID: h.ID, - - Route: h.Route, - Status: h.Status, - Request: *h.R, - - Ts: h.Ts, - Tf: h.Tf, - D: h.D, - Dr: float64(h.D.Microseconds()) / 1000, - } - - // fetch body only on hit details - if filter.Hit != "" { - hh.Body, _ = ioutil.ReadAll(hh.Request.Body) - } - - r = append(r, hh) - } - - // sort - profiler.SortHits(&r, &f) - - // filter sorted - profiler.FilterHits(&r, &f) - - return -} - -// HitsAggregated fetches a list of hits from integration gateway profiler -// and aggregates them with assigned filters -func (svc *apigwRoute) HitsAggregated(ctx context.Context, filter types.ApigwProfilerFilter) (r types.ApigwProfilerAggregationSet, f types.ApigwProfilerFilter, err error) { - f = filter - r = make(types.ApigwProfilerAggregationSet, 0) - - uDec, err := base64.URLEncoding.DecodeString(filter.Path) - - if err != nil { - return - } - - filter.Path = string(uDec) - - var ( - list = apigw.Service().Profiler().Hits(profiler.Sort{ - Path: filter.Path, - Before: filter.Before, - }) - - tsum, tmin, tmax time.Duration - ssum, smin, smax int64 - i uint64 = 1 - ) - - for p, v := range list { - tmin, tmax, tsum = time.Hour, 0, 0 - smin, smax, ssum = math.MaxInt64, 0, 0 - - i = 0 - - for _, vv := range v { - var ( - d = *vv.D - s = vv.R.ContentLength - ) - - if d < tmin { - tmin = d - } - - if d > tmax { - tmax = d - } - - if s < smin { - smin = s - } - - if s > smax { - smax = s - } - - tsum += d - ssum += s - i++ - } - - r = append(r, &types.ApigwProfilerAggregation{ - Path: p, - Count: i, - Tmin: float64(tmin.Microseconds()) / 1000, - Tmax: float64(tmax.Microseconds()) / 1000, - Tavg: float64(tsum.Microseconds()) / float64(i) / 1000, - Smin: smin, - Smax: smax, - Savg: float64(ssum) / float64(i), - }) - } - - // sort - profiler.SortAggregation(&r, &f) - - // filter - profiler.FilterAggregation(&r, &f) - - return -} diff --git a/system/service/service.go b/system/service/service.go index 4a966ab1c..f4e20fab2 100644 --- a/system/service/service.go +++ b/system/service/service.go @@ -82,6 +82,7 @@ var ( DefaultQueue *queue DefaultApigwRoute *apigwRoute DefaultApigwFilter *apigwFilter + DefaultApigwProfiler *apigwProfiler DefaultReport *report DefaultStatistics *statistics @@ -197,6 +198,7 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, ws websock DefaultAttachment = Attachment(DefaultObjectStore) DefaultQueue = Queue() DefaultApigwRoute = Route() + DefaultApigwProfiler = Profiler() DefaultApigwFilter = Filter() if err = initRoles(ctx, log.Named("rbac.roles"), c.RBAC, eventbus.Service(), rbac.Global()); err != nil {