Compare commits

...
21 changed files with 1249 additions and 747 deletions

No files matched your search

+2 -2
View File
@@ -102,10 +102,10 @@ func MissingURLSigningSecret(service string) error {
service, defaults.BaseConfigPath())
}
func AllComponentsDisabledError(service string) error {
func AllApiHandlersDisabledError(service string) error {
return fmt.Errorf("All request handlers and event consumers are disabled for %s; at least one component must be enabled."+
"Make sure your %s config contains the proper values "+
"(e.g. by using 'opencloud init --diff' and applying the patch or setting a value manually in "+
"the config/corresponding environment variable).",
"the config/corresponding environment variable).",
service, defaults.BaseConfigPath())
}
+65 -11
View File
@@ -21,6 +21,9 @@ import (
"github.com/opencloud-eu/opencloud/services/userlog/pkg/metrics"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/server/debug"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/server/http"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
consumerSvc "github.com/opencloud-eu/opencloud/services/userlog/pkg/service/consumer"
httpSvc "github.com/opencloud-eu/opencloud/services/userlog/pkg/service/http"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/events/stream"
@@ -84,10 +87,6 @@ func Server(cfg *config.Config) *cobra.Command {
mtrcs.BuildInfo.WithLabelValues(version.GetString()).Set(1)
connName := generators.GenerateConnectionName(cfg.Service.Name, generators.NTypeBus)
stream, err := stream.NatsFromConfig(connName, false, stream.NatsConfig(cfg.Events))
if err != nil {
return err
}
st := store.Create(
store.Store(cfg.Persistence.Store),
@@ -120,20 +119,40 @@ func Server(cfg *config.Config) *cobra.Command {
vClient := settingssvc.NewValueService("eu.opencloud.api.settings", grpcClient)
rClient := settingssvc.NewRoleService("eu.opencloud.api.settings", grpcClient)
handle, err := service.NewUserlogService(
service.Logger(logger),
service.Store(st),
service.HistoryClient(hClient),
service.TraceProvider(tracerProvider),
)
if err != nil {
return err
}
gr := runner.NewGroup()
{
if !cfg.HTTP.Disabled {
httpService, err := httpSvc.New(
handle,
httpSvc.Logger(logger),
httpSvc.Config(cfg),
httpSvc.GatewaySelector(gatewaySelector),
httpSvc.ValueClient(vClient),
httpSvc.RegisteredEvents(_registeredEvents),
httpSvc.TraceProvider(tracerProvider),
)
if err != nil {
logger.Info().Err(err).Str("transport", "http").Msg("Failed to initialize server")
return err
}
server, err := http.Server(
http.Logger(logger),
http.Context(ctx),
http.Config(cfg),
http.Metrics(mtrcs),
http.Store(st),
http.Stream(stream),
http.GatewaySelector(gatewaySelector),
http.History(hClient),
http.Value(vClient),
http.Role(rClient),
http.RegisteredEvents(_registeredEvents),
http.Service(httpService),
http.TracerProvider(tracerProvider),
)
@@ -143,6 +162,41 @@ func Server(cfg *config.Config) *cobra.Command {
}
gr.Add(runner.NewGoMicroHttpServerRunner(cfg.Service.Name+".http", server))
} else {
logger.Info().Msg("HTTP server disabled, not starting HTTP service")
}
if !cfg.Events.Disabled {
evStream, err := stream.NatsFromConfig(connName, false, stream.NatsConfig{
Endpoint: cfg.Events.Endpoint,
Cluster: cfg.Events.Cluster,
TLSInsecure: cfg.Events.TLSInsecure,
TLSRootCACertificate: cfg.Events.TLSRootCACertificate,
EnableTLS: cfg.Events.EnableTLS,
AuthUsername: cfg.Events.AuthUsername,
AuthPassword: cfg.Events.AuthPassword,
})
if err != nil {
return err
}
consumerService, err := consumerSvc.New(
handle,
evStream,
consumerSvc.Context(ctx),
consumerSvc.Logger(logger),
consumerSvc.Config(cfg),
consumerSvc.GatewaySelector(gatewaySelector),
consumerSvc.ValueClient(vClient),
consumerSvc.RegisteredEvents(_registeredEvents),
)
if err != nil {
return err
}
gr.Add(runner.New(cfg.Service.Name+".consumer", consumerService.Run, consumerService.Close))
} else {
logger.Info().Msg("event listening disabled, not starting event consumer")
}
{
+2
View File
@@ -53,6 +53,7 @@ type Persistence struct {
// Events combines the configuration options for the event bus.
type Events struct {
Disabled bool `yaml:"disabled" env:"USERLOG_EVENTS_DISABLED" desc:"Disables listening for events. Set this to true if the service should only handle HTTP requests." introductionVersion:"%NEXT%"`
Endpoint string `yaml:"endpoint" env:"OC_EVENTS_ENDPOINT;USERLOG_EVENTS_ENDPOINT" desc:"The address of the event system. The event system is the message queuing service. It is used as message broker for the microservice architecture." introductionVersion:"1.0.0"`
Cluster string `yaml:"cluster" env:"OC_EVENTS_CLUSTER;USERLOG_EVENTS_CLUSTER" desc:"The clusterID of the event system. The event system is the message queuing service. It is used as message broker for the microservice architecture. Mandatory when using NATS as event system." introductionVersion:"1.0.0"`
TLSInsecure bool `yaml:"tls_insecure" env:"OC_INSECURE;OC_EVENTS_TLS_INSECURE;USERLOG_EVENTS_TLS_INSECURE" desc:"Whether to verify the server TLS certificates." introductionVersion:"1.0.0"`
@@ -72,6 +73,7 @@ type CORS struct {
// HTTP defines the available http configuration.
type HTTP struct {
Disabled bool `yaml:"disabled" env:"USERLOG_HTTP_DISABLED" desc:"Disables the HTTP service. Set this to true if the service should only handle events." introductionVersion:"%NEXT%"`
Addr string `yaml:"addr" env:"USERLOG_HTTP_ADDR" desc:"The bind address of the HTTP service." introductionVersion:"1.0.0"`
Namespace string `yaml:"-"`
Root string `yaml:"root" env:"USERLOG_HTTP_ROOT" desc:"Subdirectory that serves as the root for this HTTP service." introductionVersion:"1.0.0"`
@@ -46,5 +46,9 @@ func Validate(cfg *config.Config) error {
return shared.MissingServiceAccountSecret(cfg.Service.Name)
}
if cfg.Events.Disabled && cfg.HTTP.Disabled {
return shared.AllApiHandlersDisabledError(cfg.Service.Name)
}
return nil
}
+21 -57
View File
@@ -2,40 +2,39 @@ package http
import (
"context"
"net/http"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
"github.com/opencloud-eu/opencloud/pkg/log"
ehsvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/eventhistory/v0"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/metrics"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/spf13/pflag"
"go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
)
// UserlogService is the http service interface the server needs
type UserlogService interface {
HandleGetEvents(w http.ResponseWriter, r *http.Request)
HandleDeleteEvents(w http.ResponseWriter, r *http.Request)
HandlePostGlobalEvent(w http.ResponseWriter, r *http.Request)
HandleDeleteGlobalEvent(w http.ResponseWriter, r *http.Request)
}
// Option defines a single option function.
type Option func(o *Options)
// Options defines the available options for this package.
type Options struct {
Logger log.Logger
Context context.Context
Config *config.Config
Metrics *metrics.Metrics
Flags []pflag.Flag
Namespace string
Store store.Store
Stream events.Stream
GatewaySelector pool.Selectable[gateway.GatewayAPIClient]
HistoryClient ehsvc.EventHistoryService
ValueClient settingssvc.ValueService
RoleClient settingssvc.RoleService
RegisteredEvents []events.Unmarshaller
TracerProvider trace.TracerProvider
Logger log.Logger
Context context.Context
Config *config.Config
Metrics *metrics.Metrics
Flags []pflag.Flag
Namespace string
RoleClient settingssvc.RoleService
UserlogService UserlogService
TracerProvider trace.TracerProvider
}
// newOptions initializes the available default options.
@@ -91,45 +90,10 @@ func Namespace(val string) Option {
}
}
// Store provides a function to configure the store
func Store(store store.Store) Option {
// Service provides a function to set the userlog http service
func Service(s UserlogService) Option {
return func(o *Options) {
o.Store = store
}
}
// Stream provides a function to configure the stream
func Stream(stream events.Stream) Option {
return func(o *Options) {
o.Stream = stream
}
}
// GatewaySelector provides a function to configure the gateway client selector
func GatewaySelector(gatewaySelector pool.Selectable[gateway.GatewayAPIClient]) Option {
return func(o *Options) {
o.GatewaySelector = gatewaySelector
}
}
// History provides a function to configure the event history client
func History(h ehsvc.EventHistoryService) Option {
return func(o *Options) {
o.HistoryClient = h
}
}
// RegisteredEvents provides a function to register events
func RegisteredEvents(evs []events.Unmarshaller) Option {
return func(o *Options) {
o.RegisteredEvents = evs
}
}
// Value provides a function to configure the value service client
func Value(vs settingssvc.ValueService) Option {
return func(o *Options) {
o.ValueClient = vs
o.UserlogService = s
}
}
+19 -20
View File
@@ -1,6 +1,7 @@
package http
import (
"errors"
"fmt"
stdhttp "net/http"
@@ -10,17 +11,15 @@ import (
"github.com/opencloud-eu/opencloud/pkg/account"
"github.com/opencloud-eu/opencloud/pkg/cors"
"github.com/opencloud-eu/opencloud/pkg/middleware"
"github.com/opencloud-eu/opencloud/pkg/roles"
"github.com/opencloud-eu/opencloud/pkg/service/http"
"github.com/opencloud-eu/opencloud/pkg/tracing"
"github.com/opencloud-eu/opencloud/pkg/version"
svc "github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
httpSvc "github.com/opencloud-eu/opencloud/services/userlog/pkg/service/http"
"github.com/riandyrn/otelchi"
"go-micro.dev/v4"
)
// Service is the service interface
type Service any
// Server initializes the http service and server.
func Server(opts ...Option) (http.Service, error) {
options := newOptions(opts...)
@@ -77,24 +76,24 @@ func Server(opts ...Option) (http.Service, error) {
),
)
handle, err := svc.NewUserlogService(
svc.Logger(options.Logger),
svc.Stream(options.Stream),
svc.Mux(mux),
svc.Store(options.Store),
svc.Config(options.Config),
svc.HistoryClient(options.HistoryClient),
svc.GatewaySelector(options.GatewaySelector),
svc.ValueClient(options.ValueClient),
svc.RoleClient(options.RoleClient),
svc.RegisteredEvents(options.RegisteredEvents),
svc.TraceProvider(options.TracerProvider),
)
if err != nil {
return http.Service{}, err
if options.UserlogService == nil {
return http.Service{}, errors.New("need non nil userlog http service to serve http requests")
}
if err := micro.RegisterHandler(service.Server(), handle); err != nil {
m := roles.NewManager(
// TODO: caching?
roles.Logger(options.Logger),
roles.RoleService(options.RoleClient),
)
mux.Route("/ocs/v2.php/apps/notifications/api/v1/notifications", func(r chi.Router) {
r.Get("/", options.UserlogService.HandleGetEvents)
r.Delete("/", options.UserlogService.HandleDeleteEvents)
r.Post("/global", httpSvc.RequireAdminOrSecret(&m, options.Config.GlobalNotificationsSecret)(options.UserlogService.HandlePostGlobalEvent))
r.Delete("/global", httpSvc.RequireAdminOrSecret(&m, options.Config.GlobalNotificationsSecret)(options.UserlogService.HandleDeleteGlobalEvent))
})
if err := micro.RegisterHandler(service.Server(), mux); err != nil {
return http.Service{}, err
}
@@ -0,0 +1,27 @@
package consumer
import (
"testing"
mRegistry "go-micro.dev/v4/registry"
"github.com/onsi/ginkgo/v2"
"github.com/onsi/gomega"
"github.com/opencloud-eu/opencloud/pkg/registry"
)
func init() {
r := registry.GetRegistry(registry.Inmemory())
service := registry.BuildGRPCService("eu.opencloud.api.gateway", "", "", "")
service.Nodes = []*mRegistry.Node{{
Address: "any",
}}
_ = r.Register(service)
}
func TestConsumer(t *testing.T) {
gomega.RegisterFailHandler(ginkgo.Fail)
ginkgo.RunSpecs(t, "Userlog consumer Suite")
}
@@ -1,4 +1,4 @@
package service
package consumer
import (
"context"
@@ -0,0 +1,178 @@
package consumer
import (
"context"
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
"github.com/opencloud-eu/opencloud/pkg/log"
settingsmsg "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/messages/settings/v0"
settings "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
settingsmocks "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0/mocks"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/pkg/errors"
"github.com/stretchr/testify/mock"
"github.com/onsi/ginkgo/v2"
"github.com/onsi/gomega"
)
var _ = ginkgo.Describe("NotificationFilter", func() {
var (
testLogger = log.NewLogger()
vs *settingsmocks.ValueService
ulf userlogFilter
)
ginkgo.BeforeEach(func() {
vs = &settingsmocks.ValueService{}
ulf = userlogFilter{
log: testLogger,
valueClient: vs,
}
})
setupMockValueService := func(inApp bool) *settingsmocks.ValueService {
vs := settingsmocks.ValueService{}
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(&settings.GetValueResponse{
Value: &settingsmsg.ValueWithIdentifier{
Value: &settingsmsg.Value{
Value: &settingsmsg.Value_CollectionValue{
CollectionValue: &settingsmsg.CollectionValue{
Values: []*settingsmsg.CollectionOption{
{
Key: "in-app",
Option: &settingsmsg.CollectionOption_BoolValue{BoolValue: inApp},
},
},
},
},
},
},
}, nil)
return &vs
}
ginkgo.Describe("execute", func() {
ginkgo.It("handles executants", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(nil, nil)
gomega.Expect(ulf.execute(context.TODO(), events.Event{}, &user.UserId{OpaqueId: "executant"}, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles connection errors", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(nil, errors.New("no connection to ValueService"))
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles no setting", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(&settings.GetValueResponse{}, nil)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles nil response", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(nil, nil)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles events that can not be disabled", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.BytesReceived{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles ShareCreated events", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles ShareRemoved events", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareRemoved{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles ShareExpired events", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareExpired{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles SpaceShared enabled", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceShared{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles SpaceUnshared enabled", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceUnshared{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles SpaceMembershipExpired enabled", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceMembershipExpired{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles SpaceDisabled enabled", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDisabled{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles SpaceDeleted enabled", func() {
ulf.valueClient = setupMockValueService(true)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDeleted{}}, nil, []string{"foo"})).To(gomega.ConsistOf("foo"))
})
ginkgo.It("handles ShareCreated disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles ShareRemoved disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareRemoved{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles ShareExpired disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareExpired{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles SpaceShared disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceShared{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles SpaceUnshared disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceUnshared{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles SpaceMembershipExpired disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceMembershipExpired{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles SpaceDisabled disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDisabled{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
ginkgo.It("handles SpaceDeleted disabled", func() {
ulf.valueClient = setupMockValueService(false)
gomega.Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDeleted{}}, nil, []string{"foo"})).To(gomega.BeEmpty())
})
})
})
@@ -0,0 +1,109 @@
package consumer
import (
"context"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
"github.com/opencloud-eu/opencloud/pkg/log"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
)
// Option for the consumer service
type Option func(*Options)
// Options for the consumer service
type Options struct {
UserlogService *service.UserlogService
Context context.Context
Logger log.Logger
Config *config.Config
Stream events.Stream
GatewaySelector pool.Selectable[gateway.GatewayAPIClient]
ValueClient settingssvc.ValueService
RegisteredEvents []events.Unmarshaller
}
// New creates a new consumer service
func New(ul *service.UserlogService, stream events.Stream, opts ...Option) (*Service, error) {
o := &Options{}
for _, opt := range opts {
opt(o)
}
if o.Context == nil {
o.Context = context.Background()
}
return &Service{
ctx: o.Context,
log: o.Logger,
cfg: o.Config,
userlog: ul,
stream: stream,
gatewaySelector: o.GatewaySelector,
valueClient: o.ValueClient,
registeredEvents: o.RegisteredEvents,
filter: newUserlogFilter(o.Logger, o.ValueClient),
stopCh: make(chan struct{}, 1),
}, nil
}
// UserlogService provides a function to set the userlog service
func UserlogService(ul *service.UserlogService) Option {
return func(o *Options) {
o.UserlogService = ul
}
}
// Context provides a function to set the context option.
func Context(val context.Context) Option {
return func(o *Options) {
o.Context = val
}
}
// Logger configures a logger for the consumer service
func Logger(log log.Logger) Option {
return func(o *Options) {
o.Logger = log
}
}
// Config adds the config for the consumer service
func Config(c *config.Config) Option {
return func(o *Options) {
o.Config = c
}
}
// Stream configures an event stream for the consumer service
func Stream(s events.Stream) Option {
return func(o *Options) {
o.Stream = s
}
}
// GatewaySelector adds a grpc client selector for the gateway service
func GatewaySelector(gatewaySelector pool.Selectable[gateway.GatewayAPIClient]) Option {
return func(o *Options) {
o.GatewaySelector = gatewaySelector
}
}
// ValueClient adds a grpc client for the value service
func ValueClient(vs settingssvc.ValueService) Option {
return func(o *Options) {
o.ValueClient = vs
}
}
// RegisteredEvents registers the events the service should listen to
func RegisteredEvents(e []events.Unmarshaller) Option {
return func(o *Options) {
o.RegisteredEvents = e
}
}
@@ -0,0 +1,230 @@
package consumer
import (
"context"
"encoding/json"
"errors"
"sync"
"sync/atomic"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
ocEvents "github.com/opencloud-eu/opencloud/pkg/events"
"github.com/opencloud-eu/opencloud/pkg/l10n"
"github.com/opencloud-eu/opencloud/pkg/log"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/opencloud-eu/reva/v2/pkg/utils"
)
// Service consumes events and stores them for the users
type Service struct {
ctx context.Context
log log.Logger
cfg *config.Config
userlog *service.UserlogService
stream events.Stream
gatewaySelector pool.Selectable[gateway.GatewayAPIClient]
valueClient settingssvc.ValueService
registeredEvents []events.Unmarshaller
filter *userlogFilter
stopCh chan struct{}
stopped atomic.Bool
}
// Run to fulfil Runner interface
func (s *Service) Run() error {
ch, err := events.Consume(s.stream, "userlog", s.registeredEvents...)
if err != nil {
return err
}
var wg sync.WaitGroup
ctx, cancel := context.WithCancel(s.ctx)
defer cancel()
for i := 0; i < s.cfg.MaxConcurrency; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for {
select {
case <-ctx.Done():
return
case event, ok := <-ch:
if !ok {
return
}
go s.processEvent(event)
}
}
}()
}
<-s.stopCh
cancel()
wg.Wait()
return nil
}
// Close will make the service to stop processing, so the `Run`
// method can finish.
func (s *Service) Close() {
if s.stopped.CompareAndSwap(false, true) {
close(s.stopCh)
}
}
func (s *Service) 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
)
gwc, err := s.gatewaySelector.Next()
if err != nil {
s.log.Error().Err(err).Msg("cannot get gateway client")
return
}
ctx, err := utils.GetServiceUserContext(s.cfg.ServiceAccount.ServiceAccountID, gwc, s.cfg.ServiceAccount.ServiceAccountSecret)
if err != nil {
s.log.Error().Err(err).Msg("cannot get service account")
return
}
gwc, err = s.gatewaySelector.Next()
if err != nil {
s.log.Error().Err(err).Msg("cannot get gateway client")
return
}
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:
return
}
// space related // TODO: how to find spaceadmins?
case events.SpaceDisabled:
executant = e.Executant
users, err = utils.GetSpaceMembers(ctx, e.ID.GetOpaqueId(), gwc, utils.ViewerRole)
case events.SpaceDeleted:
executant = e.Executant
for u := range e.FinalMembers {
users = append(users, u)
}
case events.SpaceShared:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case ocEvents.ResourceMention:
executant = e.Executant
for _, userID := range e.UserIDs {
users = append(users, userID.GetOpaqueId())
}
case events.SpaceUnshared:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case events.SpaceMembershipExpired:
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
// share related
case events.ShareCreated:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case events.ShareRemoved:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case events.ShareExpired:
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
}
if err != nil {
// TODO: Find out why this errors on ci pipeline
s.log.Debug().Err(err).Interface("event", event).Msg("error gathering members for event")
return
}
// II) filter users who want to receive the event
users = s.filter.execute(ctx, event, executant, users)
// III) store the eventID for each user
for _, id := range users {
if err := s.userlog.AddEventToUser(id, event); err != nil {
s.log.Error().Err(err).Str("userID", id).Str("eventid", event.ID).Msg("failed to store event for user")
return
}
}
// IV) send sses
if !s.cfg.DisableSSE {
if err := s.sendSSE(ctx, users, event); err != nil {
s.log.Error().Err(err).Interface("userid", users).Str("eventid", event.ID).Msg("cannot create sse event")
}
}
}
func (s *Service) sendSSE(ctx context.Context, userIDs []string, event events.Event) error {
m := make(map[string]events.SendSSE)
for _, userid := range userIDs {
loc := l10n.MustGetUserLocale(ctx, userid, "", s.valueClient)
if ev, ok := m[loc]; ok {
ev.UserIDs = append(m[loc].UserIDs, userid)
m[loc] = ev
continue
}
ev, err := service.NewConverter(ctx, loc, s.gatewaySelector, s.cfg.Service.Name, s.cfg.TranslationPath, s.cfg.DefaultLanguage).ConvertEvent(event.ID, event.Event)
if err != nil {
if utils.IsErrNotFound(err) || utils.IsErrPermissionDenied(err) {
// the resource was not found, we assume it is deleted
continue
}
return err
}
b, err := json.Marshal(ev)
if err != nil {
return err
}
m[loc] = events.SendSSE{
UserIDs: []string{userid},
Type: "userlog-notification",
Message: b,
}
}
for _, ev := range m {
if err := events.Publish(ctx, s.stream, ev); err != nil {
return err
}
}
return nil
}
@@ -0,0 +1,180 @@
package consumer_test
import (
"context"
"encoding/json"
"reflect"
"time"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
rpc "github.com/cs3org/go-cs3apis/cs3/rpc/v1beta1"
provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
"github.com/google/uuid"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/opencloud-eu/reva/v2/pkg/store"
"github.com/opencloud-eu/reva/v2/pkg/utils"
cs3mocks "github.com/opencloud-eu/reva/v2/tests/cs3mocks/mocks"
"github.com/stretchr/testify/mock"
microevents "go-micro.dev/v4/events"
microstore "go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
"google.golang.org/grpc"
"github.com/opencloud-eu/opencloud/pkg/log"
settingsmsg "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/messages/settings/v0"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
settingsmocks "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0/mocks"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
consumersvc "github.com/opencloud-eu/opencloud/services/userlog/pkg/service/consumer"
)
var _ = Describe("Userlog consumer", func() {
var (
cfg = &config.Config{
MaxConcurrency: 5,
DisableSSE: true,
Service: config.Service{
Name: "userlog",
},
}
cs *consumersvc.Service
bus testBus
sto microstore.Store
gatewayClient *cs3mocks.GatewayAPIClient
gatewaySelector pool.Selectable[gateway.GatewayAPIClient]
vc settingsmocks.ValueService
)
BeforeEach(func() {
var err error
sto = store.Create()
bus = testBus(make(chan events.Event))
pool.RemoveSelector("GatewaySelector" + "eu.opencloud.api.gateway")
gatewayClient = &cs3mocks.GatewayAPIClient{}
gatewaySelector = pool.GetSelector[gateway.GatewayAPIClient](
"GatewaySelector",
"eu.opencloud.api.gateway",
func(cc grpc.ClientConnInterface) gateway.GatewayAPIClient {
return gatewayClient
},
)
o := utils.AppendJSONToOpaque(nil, "grants", map[string]*provider.ResourcePermissions{"userid": {Stat: true}})
gatewayClient.On("ListStorageSpaces", mock.Anything, mock.Anything).Return(&provider.ListStorageSpacesResponse{StorageSpaces: []*provider.StorageSpace{
{
Opaque: o,
SpaceType: "project",
},
}, Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
gatewayClient.On("GetUser", mock.Anything, mock.Anything).Return(&user.GetUserResponse{User: &user.User{Id: &user.UserId{OpaqueId: "userid"}}, Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
gatewayClient.On("Authenticate", mock.Anything, mock.Anything).Return(&gateway.AuthenticateResponse{Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
vc.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(&settingssvc.GetValueResponse{
Value: &settingsmsg.ValueWithIdentifier{
Value: &settingsmsg.Value{
Value: &settingsmsg.Value_CollectionValue{
CollectionValue: &settingsmsg.CollectionValue{
Values: []*settingsmsg.CollectionOption{
{
Key: "in-app",
Option: &settingsmsg.CollectionOption_BoolValue{BoolValue: true},
},
},
},
},
},
},
}, nil)
ul, err := service.NewUserlogService(
service.Store(sto),
service.Logger(log.NewLogger()),
service.TraceProvider(trace.NewNoopTracerProvider()),
)
Expect(err).ToNot(HaveOccurred())
cs, err = consumersvc.New(
ul,
bus,
consumersvc.Context(context.Background()),
consumersvc.Logger(log.NewLogger()),
consumersvc.Config(cfg),
consumersvc.GatewaySelector(gatewaySelector),
consumersvc.ValueClient(&vc),
consumersvc.RegisteredEvents([]events.Unmarshaller{
events.SpaceDisabled{},
}),
)
Expect(err).ToNot(HaveOccurred())
})
It("stores events without HTTP (consumer-only mode)", func() {
go func() {
Expect(cs.Run()).To(Succeed())
}()
ids := make(map[string]struct{})
ids[bus.publish(events.SpaceDisabled{Executant: &user.UserId{OpaqueId: "executinguserid"}, ID: &provider.StorageSpaceId{OpaqueId: "spaceid"}})] = struct{}{}
time.Sleep(500 * time.Millisecond)
cs.Close()
recs, err := sto.Read("userid")
Expect(err).ToNot(HaveOccurred())
Expect(len(recs)).To(Equal(1))
var storedIDs []string
err = json.Unmarshal(recs[0].Value, &storedIDs)
Expect(err).ToNot(HaveOccurred())
Expect(len(storedIDs)).To(Equal(1))
_, exists := ids[storedIDs[0]]
Expect(exists).To(BeTrue())
})
AfterEach(func() {
close(bus)
})
})
type testBus chan events.Event
func (tb testBus) Consume(_ string, _ ...microevents.ConsumeOption) (<-chan microevents.Event, error) {
ch := make(chan microevents.Event)
go func() {
for ev := range tb {
b, _ := json.Marshal(ev.Event)
ch <- microevents.Event{
Payload: b,
Metadata: map[string]string{
events.MetadatakeyEventID: ev.ID,
events.MetadatakeyEventType: ev.Type,
},
}
}
}()
return ch, nil
}
func (tb testBus) Publish(_ string, _ any, _ ...microevents.PublishOption) error {
return nil
}
func (tb testBus) publish(e any) string {
ev := events.Event{
ID: uuid.New().String(),
Type: reflect.TypeOf(e).String(),
Event: e,
}
tb <- ev
return ev.ID
}
-184
View File
@@ -1,184 +0,0 @@
package service
import (
"context"
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
"github.com/opencloud-eu/opencloud/pkg/log"
settingsmsg "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/messages/settings/v0"
settings "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
settingsmocks "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0/mocks"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/pkg/errors"
"github.com/stretchr/testify/mock"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
)
var _ = Describe("NotificationFilter", func() {
var (
testLogger = log.NewLogger()
vs *settingsmocks.ValueService
ulf userlogFilter
)
BeforeEach(func() {
vs = &settingsmocks.ValueService{}
ulf = userlogFilter{
log: testLogger,
valueClient: vs,
}
})
setupMockValueService := func(inApp bool) *settingsmocks.ValueService {
vs := settingsmocks.ValueService{}
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(&settings.GetValueResponse{
Value: &settingsmsg.ValueWithIdentifier{
Value: &settingsmsg.Value{
Value: &settingsmsg.Value_CollectionValue{
CollectionValue: &settingsmsg.CollectionValue{
Values: []*settingsmsg.CollectionOption{
{
Key: "in-app",
Option: &settingsmsg.CollectionOption_BoolValue{BoolValue: inApp},
},
},
},
},
},
},
}, nil)
return &vs
}
Describe("execute", func() {
It("handles executants", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(nil, nil)
Expect(ulf.execute(context.TODO(), events.Event{}, &user.UserId{OpaqueId: "executant"}, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles connection errors", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(nil, errors.New("no connection to ValueService"))
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles no setting", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(&settings.GetValueResponse{}, nil)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles nil response", func() {
vs.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(nil, nil)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles events that can not be disabled", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.BytesReceived{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles ShareCreated events", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles ShareRemoved events", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareRemoved{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles ShareExpired events", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareExpired{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles SpaceShared enabled", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceShared{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles SpaceUnshared enabled", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceUnshared{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles SpaceMembershipExpired enabled", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceMembershipExpired{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles SpaceDisabled enabled", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDisabled{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles SpaceDeleted enabled", func() {
ulf.valueClient = setupMockValueService(true)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDeleted{}}, nil, []string{"foo"})).To(ConsistOf("foo"))
})
It("handles ShareCreated disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareCreated{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles ShareRemoved disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareRemoved{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles ShareExpired disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.ShareExpired{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles SpaceShared disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceShared{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles SpaceUnshared disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceUnshared{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles SpaceMembershipExpired disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceMembershipExpired{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles SpaceDisabled disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDisabled{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles SpaceDeleted disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.SpaceDeleted{}}, nil, []string{"foo"})).To(BeEmpty())
})
It("handles PostprocessingStepFinished disabled", func() {
ulf.valueClient = setupMockValueService(false)
Expect(ulf.execute(context.TODO(), events.Event{Event: events.PostprocessingStepFinished{}}, nil, []string{"foo"})).To(BeEmpty())
})
})
})
@@ -0,0 +1,27 @@
package http_test
import (
"testing"
mRegistry "go-micro.dev/v4/registry"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"github.com/opencloud-eu/opencloud/pkg/registry"
)
func init() {
r := registry.GetRegistry(registry.Inmemory())
service := registry.BuildGRPCService("eu.opencloud.api.gateway", "", "", "")
service.Nodes = []*mRegistry.Node{{
Address: "any",
}}
_ = r.Register(service)
}
func TestHttp(t *testing.T) {
RegisterFailHandler(Fail)
RunSpecs(t, "Userlog http Suite")
}
+102
View File
@@ -0,0 +1,102 @@
package http
import (
"reflect"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
"github.com/opencloud-eu/opencloud/pkg/log"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"go.opentelemetry.io/otel/trace"
)
// Option for the http service
type Option func(*Options)
// Options for the http service
type Options struct {
UserlogService *service.UserlogService
Logger log.Logger
Config *config.Config
GatewaySelector pool.Selectable[gateway.GatewayAPIClient]
ValueClient settingssvc.ValueService
RegisteredEvents []events.Unmarshaller
TraceProvider trace.TracerProvider
}
// New creates a new http service
func New(ul *service.UserlogService, opts ...Option) (*Service, error) {
o := &Options{}
for _, opt := range opts {
opt(o)
}
registeredEvents := make(map[string]events.Unmarshaller)
for _, e := range o.RegisteredEvents {
typ := reflect.TypeOf(e)
registeredEvents[typ.String()] = e
}
return &Service{
log: o.Logger,
cfg: o.Config,
userlog: ul,
gatewaySelector: o.GatewaySelector,
valueClient: o.ValueClient,
registeredEvents: registeredEvents,
tp: o.TraceProvider,
tracer: o.TraceProvider.Tracer("github.com/opencloud-eu/opencloud/services/userlog/pkg/service/http"),
}, nil
}
// UserlogService provides a function to set the userlog service
func UserlogService(ul *service.UserlogService) Option {
return func(o *Options) {
o.UserlogService = ul
}
}
// Logger configures a logger for the http service
func Logger(log log.Logger) Option {
return func(o *Options) {
o.Logger = log
}
}
// Config adds the config for the http service
func Config(c *config.Config) Option {
return func(o *Options) {
o.Config = c
}
}
// GatewaySelector adds a grpc client selector for the gateway service
func GatewaySelector(gatewaySelector pool.Selectable[gateway.GatewayAPIClient]) Option {
return func(o *Options) {
o.GatewaySelector = gatewaySelector
}
}
// ValueClient adds a grpc client for the value service
func ValueClient(vs settingssvc.ValueService) Option {
return func(o *Options) {
o.ValueClient = vs
}
}
// RegisteredEvents registers the events the service should listen to
func RegisteredEvents(e []events.Unmarshaller) Option {
return func(o *Options) {
o.RegisteredEvents = e
}
}
// TraceProvider adds a tracer provider for the http service
func TraceProvider(tp trace.TracerProvider) Option {
return func(o *Options) {
o.TraceProvider = tp
}
}
@@ -0,0 +1,28 @@
package http
import "github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
// GetEventResponseOC10 is the response from GET events endpoint in oc10 style
type GetEventResponseOC10 struct {
OCS struct {
Meta struct {
Message string `json:"message"`
Status string `json:"status"`
StatusCode int `json:"statuscode"`
} `json:"meta"`
Data []service.OC10Notification `json:"data"`
} `json:"ocs"`
}
// DeleteEventsRequest is the expected body for the delete request
type DeleteEventsRequest struct {
IDs []string `json:"ids"`
}
// PostEventsRequest is the expected body for the post request
type PostEventsRequest struct {
// the event type, e.g. "deprovision"
Type string `json:"type"`
// arbitray data for the event
Data map[string]string `json:"data"`
}
@@ -1,4 +1,4 @@
package service
package http
import (
"context"
@@ -6,37 +6,53 @@ import (
"errors"
"net/http"
"github.com/opencloud-eu/opencloud/pkg/log"
"github.com/opencloud-eu/opencloud/pkg/roles"
"github.com/opencloud-eu/opencloud/services/graph/pkg/errorcode"
settings "github.com/opencloud-eu/opencloud/services/settings/pkg/service/v0"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
"github.com/opencloud-eu/reva/v2/pkg/appctx"
revactx "github.com/opencloud-eu/reva/v2/pkg/ctx"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/opencloud-eu/reva/v2/pkg/utils"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
"github.com/opencloud-eu/reva/v2/pkg/events"
)
// Service is the http service responsible for serving userlog requests
type Service struct {
log log.Logger
cfg *config.Config
userlog *service.UserlogService
gatewaySelector pool.Selectable[gateway.GatewayAPIClient]
valueClient settingssvc.ValueService
registeredEvents map[string]events.Unmarshaller
tp trace.TracerProvider
tracer trace.Tracer
}
// HeaderAcceptLanguage is the header where the client can set the locale
var HeaderAcceptLanguage = "Accept-Language"
// ServeHTTP fulfills Handler interface
func (ul *UserlogService) ServeHTTP(w http.ResponseWriter, r *http.Request) {
ul.m.ServeHTTP(w, r)
}
// HandleGetEvents is the GET handler for events
func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request) {
ctx, span := ul.tracer.Start(r.Context(), "HandleGetEvents")
func (s *Service) HandleGetEvents(w http.ResponseWriter, r *http.Request) {
ctx, span := s.tracer.Start(r.Context(), "HandleGetEvents")
defer span.End()
u, ok := revactx.ContextGetUser(ctx)
if !ok {
ul.log.Error().Int("returned statuscode", http.StatusUnauthorized).Msg("user unauthorized")
s.log.Error().Int("returned statuscode", http.StatusUnauthorized).Msg("user unauthorized")
w.WriteHeader(http.StatusUnauthorized)
return
}
evs, err := ul.GetEvents(ctx, u.GetId().GetOpaqueId())
evs, err := s.userlog.GetEvents(ctx, u.GetId().GetOpaqueId())
if err != nil {
ul.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("get events failed")
s.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("get events failed")
w.WriteHeader(http.StatusInternalServerError)
return
}
@@ -45,33 +61,33 @@ func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request
Value: attribute.IntValue(len(evs)),
})
gwc, err := ul.gatewaySelector.Next()
gwc, err := s.gatewaySelector.Next()
if err != nil {
ul.log.Error().Err(err).Msg("cant get gateway client")
s.log.Error().Err(err).Msg("cant get gateway client")
w.WriteHeader(http.StatusInternalServerError)
return
}
ctx, err = utils.GetServiceUserContext(ul.cfg.ServiceAccount.ServiceAccountID, gwc, ul.cfg.ServiceAccount.ServiceAccountSecret)
ctx, err = utils.GetServiceUserContext(s.cfg.ServiceAccount.ServiceAccountID, gwc, s.cfg.ServiceAccount.ServiceAccountSecret)
if err != nil {
ul.log.Error().Err(err).Msg("cant get service account")
s.log.Error().Err(err).Msg("cant get service account")
w.WriteHeader(http.StatusInternalServerError)
return
}
conv := NewConverter(ctx, r.Header.Get(HeaderAcceptLanguage), ul.gatewaySelector, ul.cfg.Service.Name, ul.cfg.TranslationPath, ul.cfg.DefaultLanguage)
conv := service.NewConverter(ctx, r.Header.Get(HeaderAcceptLanguage), s.gatewaySelector, s.cfg.Service.Name, s.cfg.TranslationPath, s.cfg.DefaultLanguage)
var outdatedEvents []string
resp := GetEventResponseOC10{}
for _, e := range evs {
etype, ok := ul.registeredEvents[e.Type]
etype, ok := s.registeredEvents[e.Type]
if !ok {
ul.log.Error().Str("eventid", e.Id).Str("eventtype", e.Type).Msg("event not registered")
s.log.Error().Str("eventid", e.Id).Str("eventtype", e.Type).Msg("event not registered")
continue
}
einterface, err := etype.Unmarshal(e.Event)
if err != nil {
ul.log.Error().Str("eventid", e.Id).Str("eventtype", e.Type).Msg("failed to umarshal event")
s.log.Error().Str("eventid", e.Id).Str("eventtype", e.Type).Msg("failed to umarshal event")
continue
}
@@ -81,7 +97,7 @@ func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request
outdatedEvents = append(outdatedEvents, e.Id)
continue
}
ul.log.Error().Err(err).Str("eventid", e.Id).Str("eventtype", e.Type).Msg("failed to convert event")
s.log.Error().Err(err).Str("eventid", e.Id).Str("eventtype", e.Type).Msg("failed to convert event")
continue
}
@@ -91,16 +107,16 @@ func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request
// delete outdated events asynchronously
if len(outdatedEvents) > 0 {
go func() {
err := ul.DeleteEvents(u.GetId().GetOpaqueId(), outdatedEvents)
err := s.userlog.DeleteEvents(u.GetId().GetOpaqueId(), outdatedEvents)
if err != nil {
ul.log.Error().Err(err).Msg("failed to delete events")
s.log.Error().Err(err).Msg("failed to delete events")
}
}()
}
glevs, err := ul.GetGlobalEvents(ctx)
glevs, err := s.userlog.GetGlobalEvents(ctx)
if err != nil {
ul.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("get global events failed")
s.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("get global events failed")
w.WriteHeader(http.StatusInternalServerError)
return
}
@@ -108,7 +124,7 @@ func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request
for t, data := range glevs {
noti, err := conv.ConvertGlobalEvent(t, data)
if err != nil {
ul.log.Error().Err(err).Str("eventtype", t).Msg("failed to convert event")
s.log.Error().Err(err).Str("eventtype", t).Msg("failed to convert event")
continue
}
@@ -121,16 +137,16 @@ func (ul *UserlogService) HandleGetEvents(w http.ResponseWriter, r *http.Request
}
// HandlePostGlobaelEvent is the POST handler for global events
func (ul *UserlogService) HandlePostGlobalEvent(w http.ResponseWriter, r *http.Request) {
func (s *Service) HandlePostGlobalEvent(w http.ResponseWriter, r *http.Request) {
var req PostEventsRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
ul.log.Error().Err(err).Int("returned statuscode", http.StatusBadRequest).Msg("request body is malformed")
s.log.Error().Err(err).Int("returned statuscode", http.StatusBadRequest).Msg("request body is malformed")
w.WriteHeader(http.StatusBadRequest)
return
}
if err := ul.StoreGlobalEvent(r.Context(), req.Type, req.Data); err != nil {
ul.log.Error().Err(err).Msg("post: error storing global event")
if err := s.userlog.StoreGlobalEvent(r.Context(), req.Type, req.Data); err != nil {
s.log.Error().Err(err).Msg("post: error storing global event")
w.WriteHeader(http.StatusInternalServerError)
return
}
@@ -139,16 +155,16 @@ func (ul *UserlogService) HandlePostGlobalEvent(w http.ResponseWriter, r *http.R
}
// HandleDeleteGlobalEvent is the DELETE handler for global events
func (ul *UserlogService) HandleDeleteGlobalEvent(w http.ResponseWriter, r *http.Request) {
func (s *Service) HandleDeleteGlobalEvent(w http.ResponseWriter, r *http.Request) {
var req DeleteEventsRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
ul.log.Error().Err(err).Int("returned statuscode", http.StatusBadRequest).Msg("request body is malformed")
s.log.Error().Err(err).Int("returned statuscode", http.StatusBadRequest).Msg("request body is malformed")
w.WriteHeader(http.StatusBadRequest)
return
}
if err := ul.DeleteGlobalEvents(r.Context(), req.IDs); err != nil {
ul.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("delete events failed")
if err := s.userlog.DeleteGlobalEvents(r.Context(), req.IDs); err != nil {
s.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("delete events failed")
w.WriteHeader(http.StatusInternalServerError)
return
}
@@ -157,23 +173,23 @@ func (ul *UserlogService) HandleDeleteGlobalEvent(w http.ResponseWriter, r *http
}
// HandleDeleteEvents is the DELETE handler for events
func (ul *UserlogService) HandleDeleteEvents(w http.ResponseWriter, r *http.Request) {
func (s *Service) HandleDeleteEvents(w http.ResponseWriter, r *http.Request) {
u, ok := revactx.ContextGetUser(r.Context())
if !ok {
ul.log.Error().Int("returned statuscode", http.StatusUnauthorized).Msg("user unauthorized")
s.log.Error().Int("returned statuscode", http.StatusUnauthorized).Msg("user unauthorized")
w.WriteHeader(http.StatusUnauthorized)
return
}
var req DeleteEventsRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
ul.log.Error().Err(err).Int("returned statuscode", http.StatusBadRequest).Msg("request body is malformed")
s.log.Error().Err(err).Int("returned statuscode", http.StatusBadRequest).Msg("request body is malformed")
w.WriteHeader(http.StatusBadRequest)
return
}
if err := ul.DeleteEvents(u.GetId().GetOpaqueId(), req.IDs); err != nil {
ul.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("delete events failed")
if err := s.userlog.DeleteEvents(u.GetId().GetOpaqueId(), req.IDs); err != nil {
s.log.Error().Err(err).Int("returned statuscode", http.StatusInternalServerError).Msg("delete events failed")
w.WriteHeader(http.StatusInternalServerError)
return
}
@@ -181,31 +197,6 @@ func (ul *UserlogService) HandleDeleteEvents(w http.ResponseWriter, r *http.Requ
w.WriteHeader(http.StatusOK)
}
// GetEventResponseOC10 is the response from GET events endpoint in oc10 style
type GetEventResponseOC10 struct {
OCS struct {
Meta struct {
Message string `json:"message"`
Status string `json:"status"`
StatusCode int `json:"statuscode"`
} `json:"meta"`
Data []OC10Notification `json:"data"`
} `json:"ocs"`
}
// DeleteEventsRequest is the expected body for the delete request
type DeleteEventsRequest struct {
IDs []string `json:"ids"`
}
// PostEventsRequest is the expected body for the post request
type PostEventsRequest struct {
// the event type, e.g. "deprovision"
Type string `json:"type"`
// arbitray data for the event
Data map[string]string `json:"data"`
}
// RequireAdminOrSecret middleware allows only requests if the requesting user is an admin or knows the static secret
func RequireAdminOrSecret(rm *roles.Manager, secret string) func(http.HandlerFunc) http.HandlerFunc {
return func(next http.HandlerFunc) http.HandlerFunc {
@@ -228,7 +219,6 @@ func RequireAdminOrSecret(rm *roles.Manager, secret string) func(http.HandlerFun
}
errorcode.ItemNotFound.Render(w, r, http.StatusNotFound, "Not found")
return
}
}
}
@@ -0,0 +1,139 @@
package http_test
import (
"encoding/json"
"net/http"
"net/http/httptest"
"reflect"
"time"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
rpc "github.com/cs3org/go-cs3apis/cs3/rpc/v1beta1"
provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
revactx "github.com/opencloud-eu/reva/v2/pkg/ctx"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/opencloud-eu/reva/v2/pkg/store"
cs3mocks "github.com/opencloud-eu/reva/v2/tests/cs3mocks/mocks"
"github.com/stretchr/testify/mock"
microstore "go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
"google.golang.org/grpc"
"github.com/opencloud-eu/opencloud/pkg/log"
ehmsg "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/messages/eventhistory/v0"
ehsvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/eventhistory/v0"
"github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/eventhistory/v0/mocks"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
httpsvc "github.com/opencloud-eu/opencloud/services/userlog/pkg/service/http"
)
var _ = Describe("Userlog http service", func() {
var (
cfg = &config.Config{
Service: config.Service{
Name: "userlog",
},
ServiceAccount: config.ServiceAccount{},
DefaultLanguage: "en",
}
ul *service.UserlogService
hs *httpsvc.Service
sto microstore.Store
ehc mocks.EventHistoryService
gatewayClient *cs3mocks.GatewayAPIClient
gatewaySelector pool.Selectable[gateway.GatewayAPIClient]
)
BeforeEach(func() {
var err error
sto = store.Create()
ehc = mocks.EventHistoryService{}
pool.RemoveSelector("GatewaySelector" + "eu.opencloud.api.gateway")
gatewayClient = &cs3mocks.GatewayAPIClient{}
gatewaySelector = pool.GetSelector[gateway.GatewayAPIClient](
"GatewaySelector",
"eu.opencloud.api.gateway",
func(cc grpc.ClientConnInterface) gateway.GatewayAPIClient {
return gatewayClient
},
)
gatewayClient.On("Authenticate", mock.Anything, mock.Anything).Return(&gateway.AuthenticateResponse{Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
gatewayClient.On("GetUser", mock.Anything, mock.Anything).Return(&user.GetUserResponse{User: &user.User{Id: &user.UserId{OpaqueId: "executinguserid"}, Username: "executinguser", DisplayName: "Executing User"}, Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
gatewayClient.On("ListStorageSpaces", mock.Anything, mock.Anything).Return(&provider.ListStorageSpacesResponse{StorageSpaces: []*provider.StorageSpace{
{
Id: &provider.StorageSpaceId{OpaqueId: "spaceid"},
SpaceType: "project",
Root: &provider.ResourceId{StorageId: "storageid", SpaceId: "spaceid", OpaqueId: "spaceid"},
},
}, Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
ul, err = service.NewUserlogService(
service.Store(sto),
service.Logger(log.NewLogger()),
service.HistoryClient(&ehc),
service.TraceProvider(trace.NewNoopTracerProvider()),
)
Expect(err).ToNot(HaveOccurred())
hs, err = httpsvc.New(
ul,
httpsvc.Logger(log.NewLogger()),
httpsvc.Config(cfg),
httpsvc.GatewaySelector(gatewaySelector),
httpsvc.RegisteredEvents([]events.Unmarshaller{
events.SpaceDisabled{},
}),
httpsvc.TraceProvider(trace.NewNoopTracerProvider()),
)
Expect(err).ToNot(HaveOccurred())
})
It("returns events for the requesting user", func() {
ev := events.SpaceDisabled{
Executant: &user.UserId{OpaqueId: "executinguserid"},
ID: &provider.StorageSpaceId{OpaqueId: "spaceid"},
Timestamp: time.Now(),
}
b, err := json.Marshal(ev)
Expect(err).ToNot(HaveOccurred())
eventID := "eventid"
Expect(ul.AddEventToUser("userid", events.Event{ID: eventID})).To(Succeed())
ehc.On("GetEvents", mock.Anything, mock.Anything).Return(&ehsvc.GetEventsResponse{Events: []*ehmsg.Event{{
Id: eventID,
Type: reflect.TypeOf(ev).String(),
Event: b,
}}}, nil)
req := httptest.NewRequest(http.MethodGet, "/ocs/v2.php/apps/notifications/api/v1/notifications", nil)
req = req.WithContext(revactx.ContextSetUser(req.Context(), &user.User{Id: &user.UserId{OpaqueId: "userid"}}))
rr := httptest.NewRecorder()
hs.HandleGetEvents(rr, req)
Expect(rr.Code).To(Equal(http.StatusOK))
var resp httpsvc.GetEventResponseOC10
Expect(json.Unmarshal(rr.Body.Bytes(), &resp)).To(Succeed())
Expect(resp.OCS.Data).To(HaveLen(1))
Expect(resp.OCS.Data[0].EventID).To(Equal(eventID))
})
It("responds with unauthorized when no user is in context", func() {
req := httptest.NewRequest(http.MethodGet, "/ocs/v2.php/apps/notifications/api/v1/notifications", nil)
rr := httptest.NewRecorder()
hs.HandleGetEvents(rr, req)
Expect(rr.Code).To(Equal(http.StatusUnauthorized))
})
})
+4 -66
View File
@@ -1,14 +1,8 @@
package service
import (
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
"github.com/go-chi/chi/v5"
"github.com/opencloud-eu/opencloud/pkg/log"
ehsvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/eventhistory/v0"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
)
@@ -18,17 +12,10 @@ type Option func(*Options)
// Options for the userlog service
type Options struct {
Logger log.Logger
Stream events.Stream
Mux *chi.Mux
Store store.Store
Config *config.Config
HistoryClient ehsvc.EventHistoryService
GatewaySelector pool.Selectable[gateway.GatewayAPIClient]
ValueClient settingssvc.ValueService
RoleClient settingssvc.RoleService
RegisteredEvents []events.Unmarshaller
TraceProvider trace.TracerProvider
Logger log.Logger
Store store.Store
HistoryClient ehsvc.EventHistoryService
TraceProvider trace.TracerProvider
}
// Logger configures a logger for the userlog service
@@ -38,20 +25,6 @@ func Logger(log log.Logger) Option {
}
}
// Stream configures an event stream for the userlog service
func Stream(s events.Stream) Option {
return func(o *Options) {
o.Stream = s
}
}
// Mux defines the muxer for the userlog service
func Mux(m *chi.Mux) Option {
return func(o *Options) {
o.Mux = m
}
}
// Store defines the store for the userlog service
func Store(s store.Store) Option {
return func(o *Options) {
@@ -59,13 +32,6 @@ func Store(s store.Store) Option {
}
}
// Config adds the config for the userlog service
func Config(c *config.Config) Option {
return func(o *Options) {
o.Config = c
}
}
// HistoryClient adds a grpc client for the eventhistory service
func HistoryClient(hc ehsvc.EventHistoryService) Option {
return func(o *Options) {
@@ -73,34 +39,6 @@ func HistoryClient(hc ehsvc.EventHistoryService) Option {
}
}
// GatewaySelector adds a grpc client selector for the gateway service
func GatewaySelector(gatewaySelector pool.Selectable[gateway.GatewayAPIClient]) Option {
return func(o *Options) {
o.GatewaySelector = gatewaySelector
}
}
// RegisteredEvents registers the events the service should listen to
func RegisteredEvents(e []events.Unmarshaller) Option {
return func(o *Options) {
o.RegisteredEvents = e
}
}
// ValueClient adds a grpc client for the value service
func ValueClient(vs settingssvc.ValueService) Option {
return func(o *Options) {
o.ValueClient = vs
}
}
// RoleClient adds a grpc client for the role service
func RoleClient(rs settingssvc.RoleService) Option {
return func(o *Options) {
o.RoleClient = rs
}
}
// TraceProvider adds a tracer provider for the userlog service
func TraceProvider(tp trace.TracerProvider) Option {
return func(o *Options) {
+16 -228
View File
@@ -5,42 +5,23 @@ import (
"encoding/json"
"errors"
"fmt"
"reflect"
"time"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
"github.com/go-chi/chi/v5"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/opencloud-eu/reva/v2/pkg/utils"
"go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
ocEvents "github.com/opencloud-eu/opencloud/pkg/events"
"github.com/opencloud-eu/opencloud/pkg/l10n"
"github.com/opencloud-eu/opencloud/pkg/log"
"github.com/opencloud-eu/opencloud/pkg/roles"
ehmsg "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/messages/eventhistory/v0"
ehsvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/eventhistory/v0"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/reva/v2/pkg/events"
"go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
)
// UserlogService is the service responsible for user activities
type UserlogService struct {
log log.Logger
m *chi.Mux
store store.Store
cfg *config.Config
historyClient ehsvc.EventHistoryService
gatewaySelector pool.Selectable[gateway.GatewayAPIClient]
valueClient settingssvc.ValueService
registeredEvents map[string]events.Unmarshaller
tp trace.TracerProvider
tracer trace.Tracer
publisher events.Publisher
filter *userlogFilter
log log.Logger
store store.Store
historyClient ehsvc.EventHistoryService
tp trace.TracerProvider
tracer trace.Tracer
}
// NewUserlogService returns an EventHistory service
@@ -50,172 +31,21 @@ func NewUserlogService(opts ...Option) (*UserlogService, error) {
opt(o)
}
if o.Stream == nil || o.Store == nil {
return nil, fmt.Errorf("need non nil stream (%v) and store (%v) to work properly", o.Stream, o.Store)
}
ch, err := events.Consume(o.Stream, "userlog", o.RegisteredEvents...)
if err != nil {
return nil, err
if o.Store == nil {
return nil, fmt.Errorf("need non nil store to work properly")
}
ul := &UserlogService{
log: o.Logger,
m: o.Mux,
store: o.Store,
cfg: o.Config,
historyClient: o.HistoryClient,
gatewaySelector: o.GatewaySelector,
valueClient: o.ValueClient,
registeredEvents: make(map[string]events.Unmarshaller),
tp: o.TraceProvider,
tracer: o.TraceProvider.Tracer("github.com/opencloud-eu/opencloud/services/userlog/pkg/service"),
publisher: o.Stream,
filter: newUserlogFilter(o.Logger, o.ValueClient),
log: o.Logger,
store: o.Store,
historyClient: o.HistoryClient,
tp: o.TraceProvider,
tracer: o.TraceProvider.Tracer("github.com/opencloud-eu/opencloud/services/userlog/pkg/service"),
}
for _, e := range o.RegisteredEvents {
typ := reflect.TypeOf(e)
ul.registeredEvents[typ.String()] = e
}
m := roles.NewManager(
// TODO: caching?
roles.Logger(o.Logger),
roles.RoleService(o.RoleClient),
)
ul.m.Route("/ocs/v2.php/apps/notifications/api/v1/notifications", func(r chi.Router) {
r.Get("/", ul.HandleGetEvents)
r.Delete("/", ul.HandleDeleteEvents)
r.Post("/global", RequireAdminOrSecret(&m, o.Config.GlobalNotificationsSecret)(ul.HandlePostGlobalEvent))
r.Delete("/global", RequireAdminOrSecret(&m, o.Config.GlobalNotificationsSecret)(ul.HandleDeleteGlobalEvent))
})
go ul.MemorizeEvents(ch)
return ul, nil
}
// MemorizeEvents stores eventIDs a user wants to receive
func (ul *UserlogService) MemorizeEvents(ch <-chan events.Event) {
for i := 0; i < ul.cfg.MaxConcurrency; i++ {
go func(ch <-chan events.Event) {
for event := range ch {
go ul.processEvent(event)
}
}(ch)
}
}
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
)
gwc, err := ul.gatewaySelector.Next()
if err != nil {
ul.log.Error().Err(err).Msg("cannot get gateway client")
return
}
ctx, err := utils.GetServiceUserContext(ul.cfg.ServiceAccount.ServiceAccountID, gwc, ul.cfg.ServiceAccount.ServiceAccountSecret)
if err != nil {
ul.log.Error().Err(err).Msg("cannot get service account")
return
}
gwc, err = ul.gatewaySelector.Next()
if err != nil {
ul.log.Error().Err(err).Msg("cannot get gateway client")
return
}
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:
return
}
// space related // TODO: how to find spaceadmins?
case events.SpaceDisabled:
executant = e.Executant
users, err = utils.GetSpaceMembers(ctx, e.ID.GetOpaqueId(), gwc, utils.ViewerRole)
case events.SpaceDeleted:
executant = e.Executant
for u := range e.FinalMembers {
users = append(users, u)
}
case events.SpaceShared:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case ocEvents.ResourceMention:
executant = e.Executant
for _, userID := range e.UserIDs {
users = append(users, userID.GetOpaqueId())
}
case events.SpaceUnshared:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case events.SpaceMembershipExpired:
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
// share related
case events.ShareCreated:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case events.ShareRemoved:
executant = e.Executant
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
case events.ShareExpired:
users, err = utils.ResolveID(ctx, e.GranteeUserID, e.GranteeGroupID, gwc)
}
if err != nil {
// TODO: Find out why this errors on ci pipeline
ul.log.Debug().Err(err).Interface("event", event).Msg("error gathering members for event")
return
}
// II) filter users who want to receive the event
users = ul.filter.execute(ctx, event, executant, users)
// 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
}
}
// IV) send sses
if !ul.cfg.DisableSSE {
if err := ul.sendSSE(ctx, users, event, ul.gatewaySelector); err != nil {
ul.log.Error().Err(err).Interface("userid", users).Str("eventid", event.ID).Msg("cannot create sse event")
}
}
}
// GetEvents allows retrieving events from the eventhistory by userid
func (ul *UserlogService) GetEvents(ctx context.Context, userid string) ([]*ehmsg.Event, error) {
ctx, span := ul.tracer.Start(ctx, "GetEvents")
@@ -345,54 +175,12 @@ func (ul *UserlogService) DeleteGlobalEvents(ctx context.Context, evnames []stri
})
}
func (ul *UserlogService) addEventToUser(userid string, event events.Event) error {
func (ul *UserlogService) AddEventToUser(userid string, event events.Event) error {
return ul.alterUserEventList(userid, func(ids []string) []string {
return append(ids, event.ID)
})
}
func (ul *UserlogService) sendSSE(ctx context.Context, userIDs []string, event events.Event, gatewaySelector pool.Selectable[gateway.GatewayAPIClient]) error {
m := make(map[string]events.SendSSE)
for _, userid := range userIDs {
loc := l10n.MustGetUserLocale(ctx, userid, "", ul.valueClient)
if ev, ok := m[loc]; ok {
ev.UserIDs = append(m[loc].UserIDs, userid)
m[loc] = ev
continue
}
ev, err := NewConverter(ctx, loc, gatewaySelector, ul.cfg.Service.Name, ul.cfg.TranslationPath, ul.cfg.DefaultLanguage).ConvertEvent(event.ID, event.Event)
if err != nil {
if utils.IsErrNotFound(err) || utils.IsErrPermissionDenied(err) {
// the resource was not found, we assume it is deleted
continue
}
return err
}
b, err := json.Marshal(ev)
if err != nil {
return err
}
m[loc] = events.SendSSE{
UserIDs: []string{userid},
Type: "userlog-notification",
Message: b,
}
}
for _, ev := range m {
if err := events.Publish(ctx, ul.publisher, ev); err != nil {
return err
}
}
return nil
}
func (ul *UserlogService) removeExpiredEvents(userid string, all []string, received []*ehmsg.Event) error {
exists := make(map[string]struct{}, len(received))
for _, e := range received {
+41 -114
View File
@@ -2,125 +2,61 @@ package service_test
import (
"context"
"encoding/json"
"reflect"
"time"
settingsmsg "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/messages/settings/v0"
gateway "github.com/cs3org/go-cs3apis/cs3/gateway/v1beta1"
user "github.com/cs3org/go-cs3apis/cs3/identity/user/v1beta1"
rpc "github.com/cs3org/go-cs3apis/cs3/rpc/v1beta1"
provider "github.com/cs3org/go-cs3apis/cs3/storage/provider/v1beta1"
"github.com/go-chi/chi/v5"
"github.com/google/uuid"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"github.com/opencloud-eu/reva/v2/pkg/events"
"github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool"
"github.com/opencloud-eu/reva/v2/pkg/store"
"github.com/opencloud-eu/reva/v2/pkg/utils"
cs3mocks "github.com/opencloud-eu/reva/v2/tests/cs3mocks/mocks"
"github.com/stretchr/testify/mock"
microevents "go-micro.dev/v4/events"
microstore "go-micro.dev/v4/store"
"go.opentelemetry.io/otel/trace"
"google.golang.org/grpc"
"github.com/opencloud-eu/opencloud/pkg/log"
ehmsg "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/messages/eventhistory/v0"
ehsvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/eventhistory/v0"
"github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/eventhistory/v0/mocks"
settingssvc "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0"
settingsmocks "github.com/opencloud-eu/opencloud/protogen/gen/opencloud/services/settings/v0/mocks"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/config"
"github.com/opencloud-eu/opencloud/services/userlog/pkg/service"
)
var _ = Describe("UserlogService", func() {
var (
cfg = &config.Config{
MaxConcurrency: 5,
}
ul *service.UserlogService
bus testBus
sto microstore.Store
gatewayClient *cs3mocks.GatewayAPIClient
gatewaySelector pool.Selectable[gateway.GatewayAPIClient]
ehc mocks.EventHistoryService
vc settingsmocks.ValueService
)
BeforeEach(func() {
var err error
sto = store.Create()
bus = testBus(make(chan events.Event))
pool.RemoveSelector("GatewaySelector" + "eu.opencloud.api.gateway")
gatewayClient = &cs3mocks.GatewayAPIClient{}
gatewaySelector = pool.GetSelector[gateway.GatewayAPIClient](
"GatewaySelector",
"eu.opencloud.api.gateway",
func(cc grpc.ClientConnInterface) gateway.GatewayAPIClient {
return gatewayClient
},
)
o := utils.AppendJSONToOpaque(nil, "grants", map[string]*provider.ResourcePermissions{"userid": {Stat: true}})
gatewayClient.On("ListStorageSpaces", mock.Anything, mock.Anything).Return(&provider.ListStorageSpacesResponse{StorageSpaces: []*provider.StorageSpace{
{
Opaque: o,
SpaceType: "project",
},
}, Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
gatewayClient.On("GetUser", mock.Anything, mock.Anything).Return(&user.GetUserResponse{User: &user.User{Id: &user.UserId{OpaqueId: "userid"}}, Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
gatewayClient.On("Authenticate", mock.Anything, mock.Anything).Return(&gateway.AuthenticateResponse{Status: &rpc.Status{Code: rpc.Code_CODE_OK}}, nil)
vc.On("GetValueByUniqueIdentifiers", mock.Anything, mock.Anything).Return(&settingssvc.GetValueResponse{
Value: &settingsmsg.ValueWithIdentifier{
Value: &settingsmsg.Value{
Value: &settingsmsg.Value_CollectionValue{
CollectionValue: &settingsmsg.CollectionValue{
Values: []*settingsmsg.CollectionOption{
{
Key: "in-app",
Option: &settingsmsg.CollectionOption_BoolValue{BoolValue: true},
},
},
},
},
},
},
}, nil)
ehc = mocks.EventHistoryService{}
ul, err = service.NewUserlogService(
service.Config(cfg),
service.Stream(bus),
service.Store(sto),
service.Logger(log.NewLogger()),
service.Mux(chi.NewMux()),
service.GatewaySelector(gatewaySelector),
service.HistoryClient(&ehc),
service.ValueClient(&vc),
service.RegisteredEvents([]events.Unmarshaller{
events.SpaceDisabled{},
}),
service.TraceProvider(trace.NewNoopTracerProvider()),
)
Expect(err).ToNot(HaveOccurred())
})
newEvent := func() events.Event {
return events.Event{
ID: uuid.New().String(),
Type: reflect.TypeOf(events.SpaceDisabled{}).String(),
Event: events.SpaceDisabled{},
}
}
It("it stores, returns and deletes a couple of events", func() {
ids := make(map[string]struct{})
ids[bus.publish(events.SpaceDisabled{Executant: &user.UserId{OpaqueId: "executinguserid"}})] = struct{}{}
ids[bus.publish(events.SpaceDisabled{Executant: &user.UserId{OpaqueId: "executinguserid"}})] = struct{}{}
ids[bus.publish(events.SpaceDisabled{Executant: &user.UserId{OpaqueId: "executinguserid"}})] = struct{}{}
// ids[bus.Publish(events.SpaceMembershipExpired{SpaceOwner: &user.UserId{OpaqueId: "userid"}})] = struct{}{}
// ids[bus.Publish(events.ShareCreated{Executant: &user.UserId{OpaqueId: "userid"}})] = struct{}{}
time.Sleep(500 * time.Millisecond)
for i := 0; i < 3; i++ {
ev := newEvent()
Expect(ul.AddEventToUser("userid", ev)).To(Succeed())
ids[ev.ID] = struct{}{}
}
var events []*ehmsg.Event
for id := range ids {
@@ -150,41 +86,32 @@ var _ = Describe("UserlogService", func() {
Expect(len(evs)).To(Equal(0))
})
AfterEach(func() {
close(bus)
It("works without event consumer (HTTP-only mode)", func() {
ev := newEvent()
Expect(ul.AddEventToUser("userid", ev)).To(Succeed())
ehc.On("GetEvents", mock.Anything, mock.Anything).Return(&ehsvc.GetEventsResponse{
Events: []*ehmsg.Event{{Id: ev.ID}},
}, nil)
evs, err := ul.GetEvents(context.Background(), "userid")
Expect(err).ToNot(HaveOccurred())
Expect(len(evs)).To(Equal(1))
Expect(evs[0].Id).To(Equal(ev.ID))
})
It("stores and deletes global events", func() {
Expect(ul.StoreGlobalEvent(context.Background(), "deprovision", map[string]string{
"deprovision_date": time.Now().Add(time.Hour).Format(time.RFC3339),
})).To(Succeed())
evs, err := ul.GetGlobalEvents(context.Background())
Expect(err).ToNot(HaveOccurred())
Expect(evs).To(HaveLen(1))
Expect(ul.DeleteGlobalEvents(context.Background(), []string{"deprovision"})).To(Succeed())
evs, err = ul.GetGlobalEvents(context.Background())
Expect(err).ToNot(HaveOccurred())
Expect(evs).To(HaveLen(0))
})
})
type testBus chan events.Event
func (tb testBus) Consume(_ string, _ ...microevents.ConsumeOption) (<-chan microevents.Event, error) {
ch := make(chan microevents.Event)
go func() {
for ev := range tb {
b, _ := json.Marshal(ev.Event)
ch <- microevents.Event{
Payload: b,
Metadata: map[string]string{
events.MetadatakeyEventID: ev.ID,
events.MetadatakeyEventType: ev.Type,
},
}
}
}()
return ch, nil
}
func (tb testBus) Publish(_ string, _ any, _ ...microevents.PublishOption) error {
return nil
}
func (tb testBus) publish(e any) string {
ev := events.Event{
ID: uuid.New().String(),
Type: reflect.TypeOf(e).String(),
Event: e,
}
tb <- ev
return ev.ID
}