3
0

Moved filter, sorter to system profiler service

This commit is contained in:
Peter Grlica
2022-03-17 15:54:23 +01:00
parent f33336b21d
commit 1d156899d7
9 changed files with 342 additions and 353 deletions
+1 -15
View File
@@ -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
-79
View File
@@ -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
}
-8
View File
@@ -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
-83
View File
@@ -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)
}
}
-4
View File
@@ -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
+1 -1
View File
@@ -41,7 +41,7 @@ type (
func (ApigwProfiler) New() *ApigwProfiler {
return &ApigwProfiler{
svc: service.DefaultApigwRoute,
svc: service.DefaultApigwProfiler,
ac: service.DefaultAccessControl,
}
}
+338
View File
@@ -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)
}
}
-163
View File
@@ -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
}
+2
View File
@@ -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 {