Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 3 additions & 4 deletions services/pkg/common/db/adapters/mysql_adapter.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,15 +33,14 @@ func NewMySqlAdapter(config MySqlAdapterConfig) (SqlDataAdapter, error) {
var db *gorm.DB
var err error

retryN := 5
utils.InvokeWithRetry(utils.RetryConfig{
Count: retryN,
Count: 30,
Sleep: 1 * time.Second,
}, func(n int) error {
}, func(arg utils.RetryFuncArg) error {
db, err = gorm.Open(mysql.Open(dsn), &gorm.Config{})
if err != nil {
log.Printf("[%d/%d] Failed to connect to MySQL server: %v",
n, retryN, err)
arg.Current, arg.Total, err)
}

return err
Expand Down
6 changes: 6 additions & 0 deletions services/pkg/common/obs/tracing.go
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,12 @@ func Spanned(current context.Context, name string,
return err
}

// API to set the given attributes to various observability context
// such as Span, log etc.
func SetAttributeInContext(ctx context.Context, key string, value string) {
SetSpanAttribute(ctx, key, value)
}

func SetSpanAttribute(ctx context.Context, key string, value string) {
span := trace.SpanFromContext(ctx)
span.SetAttributes(attribute.KeyValue{
Expand Down
12 changes: 10 additions & 2 deletions services/pkg/common/utils/retry.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,15 @@ var (
errInvalidSleepDuration = errors.New("must have a valid sleep")
)

type RetriableFunc func(retryN int) error
type RetryFuncArg struct {
// Total retries to be executed
Total int

// Current try count starting with 1
Current int
}

type RetriableFunc func(arg RetryFuncArg) error

type RetryConfig struct {
Count int
Expand All @@ -29,7 +37,7 @@ func InvokeWithRetry(config RetryConfig, f RetriableFunc) error {

var err error
for i := 0; i < config.Count; i += 1 {
err = f(i + 1)
err = f(RetryFuncArg{Total: config.Count, Current: (i + 1)})
if err == nil {
break
}
Expand Down
5 changes: 3 additions & 2 deletions services/pkg/dcs/opensearch.go
Original file line number Diff line number Diff line change
Expand Up @@ -124,8 +124,9 @@ func (s *opensearchIndexer) initOpenSearchIndex(name string) error {
return utils.InvokeWithRetry(utils.RetryConfig{
Count: 30,
Sleep: time.Second * 1,
}, func(n int) error {
logger.Infof("Attempting to init opensearch index [retry=%d]", n)
}, func(arg utils.RetryFuncArg) error {
logger.Infof("Attempting to init opensearch index [retry=%d/%d]",
arg.Current, arg.Total)
return s.initOpenSearchIndexInternal(name)
})
}
Expand Down
2 changes: 1 addition & 1 deletion services/pkg/pdp/authorizer.go
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ func (s *authorizationService) checkInternal(ctx context.Context,

upstreamArtefact, upstream, err := s.resolveRequestedArtefact(httpReq)
if err != nil {
logger.Errorf("No artefact resolved: %s", err.Error())
logger.Warnf("No artefact resolved: %s", err.Error())
return &envoy_service_auth_v3.CheckResponse{}, err
}

Expand Down
6 changes: 3 additions & 3 deletions services/pkg/pdp/policy_engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ import (
)

type PolicyEngine struct {
lock sync.Mutex
lock sync.RWMutex
repository string
rego *rego.Rego
query *rego.PreparedEvalQuery
Expand All @@ -35,8 +35,8 @@ func NewPolicyEngine(path string, changeMonitor bool) (*PolicyEngine, error) {
}

func (svc *PolicyEngine) Evaluate(ctx context.Context, input PolicyInput) (PolicyResponse, error) {
svc.lock.Lock()
defer svc.lock.Unlock()
svc.lock.RLock()
defer svc.lock.RUnlock()

rs, err := svc.query.Eval(ctx, rego.EvalInput(input))
if err != nil {
Expand Down
42 changes: 30 additions & 12 deletions services/pkg/tap/tap_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,16 +2,23 @@ package tap

import (
"context"
"log"
"io"

envoy_config_core_v3 "github.com/envoyproxy/go-control-plane/envoy/config/core/v3"
envoy_v3_ext_proc_pb "github.com/envoyproxy/go-control-plane/envoy/service/ext_proc/v3"
"github.com/safedep/gateway/services/pkg/common/logger"
"github.com/safedep/gateway/services/pkg/common/messaging"
"github.com/safedep/gateway/services/pkg/common/obs"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)

const (
obsKeyTapReqType = "tap_req_type"
tapResponseTapSignatureKey = "x-gateway-tap"
tapResponseTapSignatureValue = "true"
)

type tapService struct {
handlerChain TapHandlerChain
messagingService messaging.MessagingService
Expand All @@ -29,25 +36,33 @@ func (s *tapService) RegisterHandler(handler TapHandlerRegistration) {
}

func (s *tapService) Process(srv envoy_v3_ext_proc_pb.ExternalProcessor_ProcessServer) error {
log.Printf("Tap service: Handling stream")
logger.Debugf("Tap service: Handling stream")

ctx := srv.Context()
for {
select {
case <-ctx.Done():
logger.Errorf("Context is finished: %v", ctx.Err())
logger.Debugf("Context is finished: %v", ctx.Err())
return ctx.Err()
default:
}

req, err := srv.Recv()
if err == io.EOF {
return nil
}

if err != nil {
logger.Errorf("Received error from stream: %v", err)
return status.Errorf(codes.Unknown, "Error receiving request: %v", err)
}

resp := &envoy_v3_ext_proc_pb.ProcessingResponse{}
switch req.Request.(type) {
case *envoy_v3_ext_proc_pb.ProcessingRequest_RequestHeaders:
obs.SetAttributeInContext(ctx, obsKeyTapReqType,
"ProcessingRequest_RequestHeaders")

err = s.handleRequestHeaders(ctx,
req.Request.(*envoy_v3_ext_proc_pb.ProcessingRequest_RequestHeaders))

Expand All @@ -64,21 +79,24 @@ func (s *tapService) Process(srv envoy_v3_ext_proc_pb.ExternalProcessor_ProcessS
resp.Response.(*envoy_v3_ext_proc_pb.ProcessingResponse_RequestHeaders))
break
case *envoy_v3_ext_proc_pb.ProcessingRequest_ResponseHeaders:
obs.SetAttributeInContext(ctx, obsKeyTapReqType,
"ProcessingRequest_ResponseHeaders")

err = s.handleResponseHeaders(ctx,
req.Request.(*envoy_v3_ext_proc_pb.ProcessingRequest_ResponseHeaders))
s.addTapSignature(resp)
break
default:
log.Printf("Unknown request type: %v", req.Request)
logger.Warnf("Unknown request type: %v", req.Request)
}

// TODO: How should we handle this behavior?
if err != nil {
log.Printf("Error in handling processing req: %v", err)
logger.Warnf("Error in handling processing req: %v", err)
}

if err := srv.Send(resp); err != nil {
log.Printf("Failed to send stream response: %v", err)
if err = srv.Send(resp); err != nil {
logger.Warnf("Failed to send stream response: %v", err)
}
}
}
Expand All @@ -88,7 +106,7 @@ func (s *tapService) handleRequestHeaders(ctx context.Context,
for _, registration := range s.handlerChain.Handlers {
err := registration.Handler.HandleRequestHeaders(ctx, req)
if !registration.ContinueOnError && err != nil {
log.Printf("Unable to continue on tap handler error: %v", err)
logger.Warnf("Unable to continue on tap handler error: %v", err)
return err
}
}
Expand All @@ -101,7 +119,7 @@ func (s *tapService) handleResponseHeaders(ctx context.Context,
for _, registration := range s.handlerChain.Handlers {
err := registration.Handler.HandleResponseHeaders(ctx, req)
if !registration.ContinueOnError && err != nil {
log.Printf("Unable to continue on tap handler error: %v", err)
logger.Warnf("Unable to continue on tap handler error: %v", err)
return err
}
}
Expand All @@ -115,16 +133,16 @@ func (s *tapService) addTapSignature(resp *envoy_v3_ext_proc_pb.ProcessingRespon
return
}

log.Printf("Adding tap signature to response headers")
logger.Debugf("Adding tap signature to response headers")
resp.Response = &envoy_v3_ext_proc_pb.ProcessingResponse_ResponseHeaders{
ResponseHeaders: &envoy_v3_ext_proc_pb.HeadersResponse{
Response: &envoy_v3_ext_proc_pb.CommonResponse{
HeaderMutation: &envoy_v3_ext_proc_pb.HeaderMutation{
SetHeaders: []*envoy_config_core_v3.HeaderValueOption{
{
Header: &envoy_config_core_v3.HeaderValue{
Key: "x-gateway-tap",
Value: "true",
Key: tapResponseTapSignatureKey,
Value: tapResponseTapSignatureValue,
},
},
},
Expand Down