mirror of
https://github.com/opencloud-eu/opencloud.git
synced 2026-09-08 11:53:07 -04:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b964da7daf | ||
|
|
9c628b697c | ||
|
|
a9cede5438 |
No files matched your search
@@ -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())
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
{
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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
-1
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
@@ -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"`
|
||||
}
|
||||
+54
-64
@@ -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))
|
||||
})
|
||||
})
|
||||
@@ -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) {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
Reference in new issue
Block a user