Add tracing to more places in userlog.

This passes the treaceprovider to more places where it belongs
and also adds tracing to the event processing.
This commit is contained in:
Daniël Franke
2023-08-25 10:39:47 +02:00
parent 0eb4b5dc56
commit 7a1e8f1f11
4 changed files with 96 additions and 89 deletions
+2 -2
View File
@@ -19,8 +19,7 @@ import (
) )
// Service is the service interface // Service is the service interface
type Service interface { type Service interface{}
}
// Server initializes the http service and server. // Server initializes the http service and server.
func Server(opts ...Option) (http.Service, error) { func Server(opts ...Option) (http.Service, error) {
@@ -90,6 +89,7 @@ func Server(opts ...Option) (http.Service, error) {
svc.ValueClient(options.ValueClient), svc.ValueClient(options.ValueClient),
svc.RoleClient(options.RoleClient), svc.RoleClient(options.RoleClient),
svc.RegisteredEvents(options.RegisteredEvents), svc.RegisteredEvents(options.RegisteredEvents),
svc.TraceProvider(options.TracerProvider),
) )
if err != nil { if err != nil {
return http.Service{}, err return http.Service{}, err
+1 -1
View File
@@ -24,7 +24,7 @@ func (ul *UserlogService) ServeHTTP(w http.ResponseWriter, r *http.Request) {
// HandleGetEvents is the GET handler for events // HandleGetEvents is the GET handler for events
func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request) { func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request) {
ctx, span := tracer.Start(r.Context(), "HandleGetEvents") ctx, span := ul.tracer.Start(r.Context(), "HandleGetEvents")
defer span.End() defer span.End()
u, ok := revactx.ContextGetUser(ctx) u, ok := revactx.ContextGetUser(ctx)
if !ok { if !ok {
+9
View File
@@ -10,6 +10,7 @@ import (
settingssvc "github.com/owncloud/ocis/v2/protogen/gen/ocis/services/settings/v0" settingssvc "github.com/owncloud/ocis/v2/protogen/gen/ocis/services/settings/v0"
"github.com/owncloud/ocis/v2/services/userlog/pkg/config" "github.com/owncloud/ocis/v2/services/userlog/pkg/config"
"go-micro.dev/v4/store" "go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
) )
// Option for the userlog service // Option for the userlog service
@@ -27,6 +28,7 @@ type Options struct {
ValueClient settingssvc.ValueService ValueClient settingssvc.ValueService
RoleClient settingssvc.RoleService RoleClient settingssvc.RoleService
RegisteredEvents []events.Unmarshaller RegisteredEvents []events.Unmarshaller
TraceProvider trace.TracerProvider
} }
// Logger configures a logger for the userlog service // Logger configures a logger for the userlog service
@@ -98,3 +100,10 @@ func RoleClient(rs settingssvc.RoleService) Option {
o.RoleClient = rs o.RoleClient = rs
} }
} }
// TraceProvider adds a tracer provider for the userlog service
func TraceProvider(tp trace.TracerProvider) Option {
return func(o *Options) {
o.TraceProvider = tp
}
}
+84 -86
View File
@@ -29,17 +29,10 @@ import (
"github.com/r3labs/sse/v2" "github.com/r3labs/sse/v2"
micrometadata "go-micro.dev/v4/metadata" micrometadata "go-micro.dev/v4/metadata"
"go-micro.dev/v4/store" "go-micro.dev/v4/store"
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/trace" "go.opentelemetry.io/otel/trace"
"google.golang.org/grpc/metadata" "google.golang.org/grpc/metadata"
) )
var tracer trace.Tracer
func init() {
tracer = otel.Tracer("github.com/owncloud/ocis/services/userlog/pkg/service")
}
// UserlogService is the service responsible for user activities // UserlogService is the service responsible for user activities
type UserlogService struct { type UserlogService struct {
log log.Logger log log.Logger
@@ -51,6 +44,8 @@ type UserlogService struct {
valueClient settingssvc.ValueService valueClient settingssvc.ValueService
sse *sse.Server sse *sse.Server
registeredEvents map[string]events.Unmarshaller registeredEvents map[string]events.Unmarshaller
tp trace.TracerProvider
tracer trace.Tracer
} }
// NewUserlogService returns an EventHistory service // NewUserlogService returns an EventHistory service
@@ -78,6 +73,8 @@ func NewUserlogService(opts ...Option) (*UserlogService, error) {
gatewaySelector: o.GatewaySelector, gatewaySelector: o.GatewaySelector,
valueClient: o.ValueClient, valueClient: o.ValueClient,
registeredEvents: make(map[string]events.Unmarshaller), registeredEvents: make(map[string]events.Unmarshaller),
tp: o.TraceProvider,
tracer: o.TraceProvider.Tracer("github.com/owncloud/ocis/services/userlog/pkg/service"),
} }
if !ul.cfg.DisableSSE { if !ul.cfg.DisableSSE {
@@ -114,91 +111,95 @@ func NewUserlogService(opts ...Option) (*UserlogService, error) {
// MemorizeEvents stores eventIDs a user wants to receive // MemorizeEvents stores eventIDs a user wants to receive
func (ul *UserlogService) MemorizeEvents(ch <-chan events.Event) { func (ul *UserlogService) MemorizeEvents(ch <-chan events.Event) {
for event := range ch { for event := range ch {
// for each event we need to: ul.processEvent(event)
// I) find users eligible to receive the event }
var ( }
users []string
executant *user.UserId
err error
)
switch e := event.Event.(type) { func (ul *UserlogService) processEvent(event events.Event) {
// for each event we need to:
// I) find users eligible to receive the event
var (
users []string
executant *user.UserId
err error
)
switch e := event.Event.(type) {
default:
err = errors.New("unhandled event")
// file related
case events.PostprocessingStepFinished:
switch e.FinishedStep {
case events.PPStepAntivirus:
result := e.Result.(events.VirusscanResult)
if !result.Infected {
return
}
// TODO: should space mangers also be informed?
users = append(users, e.ExecutingUser.GetId().GetOpaqueId())
case events.PPStepPolicies:
if e.Outcome == events.PPOutcomeContinue {
return
}
users = append(users, e.ExecutingUser.GetId().GetOpaqueId())
default: default:
err = errors.New("unhandled event") return
// file related
case events.PostprocessingStepFinished:
switch e.FinishedStep {
case events.PPStepAntivirus:
result := e.Result.(events.VirusscanResult)
if !result.Infected {
continue
}
// TODO: should space mangers also be informed?
users = append(users, e.ExecutingUser.GetId().GetOpaqueId())
case events.PPStepPolicies:
if e.Outcome == events.PPOutcomeContinue {
continue
}
users = append(users, e.ExecutingUser.GetId().GetOpaqueId())
default:
continue
}
// space related // TODO: how to find spaceadmins?
case events.SpaceDisabled:
executant = e.Executant
users, err = ul.findSpaceMembers(ul.impersonate(e.Executant), e.ID.GetOpaqueId(), viewer)
case events.SpaceDeleted:
executant = e.Executant
for u := range e.FinalMembers {
users = append(users, u)
}
case events.SpaceShared:
executant = e.Executant
users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
case events.SpaceUnshared:
executant = e.Executant
users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
case events.SpaceMembershipExpired:
users, err = ul.resolveID(ul.impersonate(e.SpaceOwner), e.GranteeUserID, e.GranteeGroupID)
// share related
case events.ShareCreated:
executant = e.Executant
users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
case events.ShareRemoved:
executant = e.Executant
users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
case events.ShareExpired:
users, err = ul.resolveID(ul.impersonate(e.ShareOwner), e.GranteeUserID, e.GranteeGroupID)
} }
// space related // TODO: how to find spaceadmins?
if err != nil { case events.SpaceDisabled:
// TODO: Find out why this errors on ci pipeline executant = e.Executant
ul.log.Debug().Err(err).Interface("event", event).Msg("error gathering members for event") users, err = ul.findSpaceMembers(ul.impersonate(e.Executant), e.ID.GetOpaqueId(), viewer)
continue case events.SpaceDeleted:
executant = e.Executant
for u := range e.FinalMembers {
users = append(users, u)
} }
case events.SpaceShared:
executant = e.Executant
users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
case events.SpaceUnshared:
executant = e.Executant
users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
case events.SpaceMembershipExpired:
users, err = ul.resolveID(ul.impersonate(e.SpaceOwner), e.GranteeUserID, e.GranteeGroupID)
// II) filter users who want to receive the event // share related
// This step is postponed for later. case events.ShareCreated:
// For now each user should get all events she is eligible to receive executant = e.Executant
// ...except notifications for their own actions users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
users = removeExecutant(users, executant) case events.ShareRemoved:
executant = e.Executant
users, err = ul.resolveID(ul.impersonate(e.Executant), e.GranteeUserID, e.GranteeGroupID)
case events.ShareExpired:
users, err = ul.resolveID(ul.impersonate(e.ShareOwner), e.GranteeUserID, e.GranteeGroupID)
}
// III) store the eventID for each user if err != nil {
for _, id := range users { // TODO: Find out why this errors on ci pipeline
if err := ul.addEventToUser(id, event); err != nil { ul.log.Debug().Err(err).Interface("event", event).Msg("error gathering members for event")
ul.log.Error().Err(err).Str("userID", id).Str("eventid", event.ID).Msg("failed to store event for user") return
continue }
}
// II) filter users who want to receive the event
// This step is postponed for later.
// For now each user should get all events she is eligible to receive
// ...except notifications for their own actions
users = removeExecutant(users, executant)
// III) store the eventID for each user
for _, id := range users {
if err := ul.addEventToUser(id, event); err != nil {
ul.log.Error().Err(err).Str("userID", id).Str("eventid", event.ID).Msg("failed to store event for user")
return
} }
} }
} }
// GetEvents allows retrieving events from the eventhistory by userid // GetEvents allows retrieving events from the eventhistory by userid
func (ul *UserlogService) GetEvents(ctx context.Context, userid string) ([]*ehmsg.Event, error) { func (ul *UserlogService) GetEvents(ctx context.Context, userid string) ([]*ehmsg.Event, error) {
ctx, span := tracer.Start(ctx, "GetEvents") ctx, span := ul.tracer.Start(ctx, "GetEvents")
defer span.End() defer span.End()
rec, err := ul.store.Read(userid) rec, err := ul.store.Read(userid)
if err != nil && err != store.ErrNotFound { if err != nil && err != store.ErrNotFound {
@@ -227,11 +228,9 @@ func (ul *UserlogService) GetEvents(ctx context.Context, userid string) ([]*ehms
if err := ul.removeExpiredEvents(userid, eventIDs, resp.Events); err != nil { if err := ul.removeExpiredEvents(userid, eventIDs, resp.Events); err != nil {
ul.log.Error().Err(err).Str("userid", userid).Msg("could not remove expired events from user") ul.log.Error().Err(err).Str("userid", userid).Msg("could not remove expired events from user")
} }
}() }()
return resp.Events, nil return resp.Events, nil
} }
// DeleteEvents will delete the specified events // DeleteEvents will delete the specified events
@@ -256,7 +255,7 @@ func (ul *UserlogService) DeleteEvents(userid string, evids []string) error {
// StoreGlobalEvent will store a global event that will be returned with each `GetEvents` request // StoreGlobalEvent will store a global event that will be returned with each `GetEvents` request
func (ul *UserlogService) StoreGlobalEvent(ctx context.Context, typ string, data map[string]string) error { func (ul *UserlogService) StoreGlobalEvent(ctx context.Context, typ string, data map[string]string) error {
ctx, span := tracer.Start(ctx, "StoreGlobalEvent") ctx, span := ul.tracer.Start(ctx, "StoreGlobalEvent")
defer span.End() defer span.End()
switch typ { switch typ {
default: default:
@@ -293,12 +292,11 @@ func (ul *UserlogService) StoreGlobalEvent(ctx context.Context, typ string, data
return nil return nil
}) })
} }
} }
// GetGlobalEvents will return all global events // GetGlobalEvents will return all global events
func (ul *UserlogService) GetGlobalEvents(ctx context.Context) (map[string]json.RawMessage, error) { func (ul *UserlogService) GetGlobalEvents(ctx context.Context) (map[string]json.RawMessage, error) {
_, span := tracer.Start(ctx, "GetGlobalEvents") _, span := ul.tracer.Start(ctx, "GetGlobalEvents")
defer span.End() defer span.End()
out := make(map[string]json.RawMessage) out := make(map[string]json.RawMessage)
@@ -318,7 +316,7 @@ func (ul *UserlogService) GetGlobalEvents(ctx context.Context) (map[string]json.
// DeleteGlobalEvents will delete the specified event // DeleteGlobalEvents will delete the specified event
func (ul *UserlogService) DeleteGlobalEvents(ctx context.Context, evnames []string) error { func (ul *UserlogService) DeleteGlobalEvents(ctx context.Context, evnames []string) error {
_, span := tracer.Start(ctx, "DeleteGlobalEvents") _, span := ul.tracer.Start(ctx, "DeleteGlobalEvents")
defer span.End() defer span.End()
return ul.alterGlobalEvents(ctx, func(evs map[string]json.RawMessage) error { return ul.alterGlobalEvents(ctx, func(evs map[string]json.RawMessage) error {
for _, name := range evnames { for _, name := range evnames {
@@ -406,7 +404,7 @@ func (ul *UserlogService) alterUserEventList(userid string, alter func([]string)
} }
func (ul *UserlogService) alterGlobalEvents(ctx context.Context, alter func(map[string]json.RawMessage) error) error { func (ul *UserlogService) alterGlobalEvents(ctx context.Context, alter func(map[string]json.RawMessage) error) error {
_, span := tracer.Start(ctx, "alterGlobalEvents") _, span := ul.tracer.Start(ctx, "alterGlobalEvents")
defer span.End() defer span.End()
evs, err := ul.GetGlobalEvents(ctx) evs, err := ul.GetGlobalEvents(ctx)
if err != nil && err != store.ErrNotFound { if err != nil && err != store.ErrNotFound {