diff --git a/services/userlog/pkg/command/server.go b/services/userlog/pkg/command/server.go index 32f08c4cc8..e461b4ba3b 100644 --- a/services/userlog/pkg/command/server.go +++ b/services/userlog/pkg/command/server.go @@ -21,7 +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" + "github.com/go-chi/chi/v5" "github.com/opencloud-eu/reva/v2/pkg/events" "github.com/opencloud-eu/reva/v2/pkg/events/stream" "github.com/opencloud-eu/reva/v2/pkg/rgrpc/todo/pool" @@ -84,9 +86,12 @@ 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 + var evStream events.Stream + if !cfg.EventsDisabled { + evStream, err = stream.NatsFromConfig(connName, false, stream.NatsConfig(cfg.Events)) + if err != nil { + return err + } } st := store.Create( @@ -121,14 +126,15 @@ func Server(cfg *config.Config) *cobra.Command { rClient := settingssvc.NewRoleService("eu.opencloud.api.settings", grpcClient) gr := runner.NewGroup() - { + + if !cfg.HTTPDisabled { server, err := http.Server( http.Logger(logger), http.Context(ctx), http.Config(cfg), http.Metrics(mtrcs), http.Store(st), - http.Stream(stream), + http.Stream(evStream), http.GatewaySelector(gatewaySelector), http.History(hClient), http.Value(vClient), @@ -143,6 +149,36 @@ 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.EventsDisabled { + _, err := service.NewUserlogService( + service.Logger(logger), + service.Stream(evStream), + service.Mux(chi.NewMux()), + service.Store(st), + service.Config(cfg), + service.HistoryClient(hClient), + service.GatewaySelector(gatewaySelector), + service.ValueClient(vClient), + service.RoleClient(rClient), + service.RegisteredEvents(_registeredEvents), + service.TraceProvider(tracerProvider), + ) + if err != nil { + return err + } + + gr.Add(runner.New(cfg.Service.Name+".consumer", func() error { + <-ctx.Done() + return nil + }, func() {})) + } + } + + if cfg.EventsDisabled { + logger.Info().Msg("event listening disabled, not starting event consumer") } { diff --git a/services/userlog/pkg/config/config.go b/services/userlog/pkg/config/config.go index d2c3a5baec..d168185e42 100644 --- a/services/userlog/pkg/config/config.go +++ b/services/userlog/pkg/config/config.go @@ -30,6 +30,9 @@ type Config struct { DisableSSE bool `yaml:"disable_sse" env:"OC_DISABLE_SSE,USERLOG_DISABLE_SSE" desc:"Disables server-sent events (sse). When disabled, clients will no longer receive sse notifications." introductionVersion:"1.0.0"` + EventsDisabled bool `yaml:"events_disabled" env:"USERLOG_EVENTS_DISABLED" desc:"Disables listening for events. Set this to true if the service should only handle HTTP requests." introductionVersion:"%NEXT%"` + HTTPDisabled bool `yaml:"http_disabled" env:"USERLOG_HTTP_DISABLED" desc:"Disables the HTTP service. Set this to true if the service should only handle events." introductionVersion:"%NEXT%"` + GlobalNotificationsSecret string `yaml:"global_notifications_secret" env:"USERLOG_GLOBAL_NOTIFICATIONS_SECRET" desc:"The secret to secure the global notifications endpoint. Only system admins and users knowing that secret can call the global notifications POST/DELETE endpoints." introductionVersion:"1.0.0"` ServiceAccount ServiceAccount `yaml:"service_account"` diff --git a/services/userlog/pkg/config/parser/parse.go b/services/userlog/pkg/config/parser/parse.go index a034d0a920..20bea22d36 100644 --- a/services/userlog/pkg/config/parser/parse.go +++ b/services/userlog/pkg/config/parser/parse.go @@ -46,5 +46,9 @@ func Validate(cfg *config.Config) error { return shared.MissingServiceAccountSecret(cfg.Service.Name) } + if cfg.EventsDisabled && cfg.HTTPDisabled { + return shared.AllComponentsDisabledError(cfg.Service.Name) + } + return nil } diff --git a/services/userlog/pkg/service/http.go b/services/userlog/pkg/service/http.go index 97e87c961c..ba67452e83 100644 --- a/services/userlog/pkg/service/http.go +++ b/services/userlog/pkg/service/http.go @@ -228,7 +228,6 @@ func RequireAdminOrSecret(rm *roles.Manager, secret string) func(http.HandlerFun } errorcode.ItemNotFound.Render(w, r, http.StatusNotFound, "Not found") - return } } } diff --git a/services/userlog/pkg/service/service.go b/services/userlog/pkg/service/service.go index ee59e7b40c..ef0234ad31 100644 --- a/services/userlog/pkg/service/service.go +++ b/services/userlog/pkg/service/service.go @@ -50,13 +50,8 @@ 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{ @@ -92,7 +87,13 @@ func NewUserlogService(opts ...Option) (*UserlogService, error) { r.Delete("/global", RequireAdminOrSecret(&m, o.Config.GlobalNotificationsSecret)(ul.HandleDeleteGlobalEvent)) }) - go ul.MemorizeEvents(ch) + if !o.Config.EventsDisabled { + ch, err := events.Consume(o.Stream, "userlog", o.RegisteredEvents...) + if err != nil { + return nil, err + } + go ul.MemorizeEvents(ch) + } return ul, nil } diff --git a/services/userlog/pkg/service/service_test.go b/services/userlog/pkg/service/service_test.go index 8138044bbf..40795295b1 100644 --- a/services/userlog/pkg/service/service_test.go +++ b/services/userlog/pkg/service/service_test.go @@ -150,6 +150,60 @@ var _ = Describe("UserlogService", func() { Expect(len(evs)).To(Equal(0)) }) + It("verifies events are stored in store without using HTTP (consumer-only mode)", func() { + ids := make(map[string]struct{}) + ids[bus.publish(events.SpaceDisabled{Executant: &user.UserId{OpaqueId: "executinguserid"}})] = struct{}{} + + time.Sleep(500 * time.Millisecond) + + 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()) + }) + + It("works without event consumer (HTTP-only mode)", func() { + localEhc := mocks.EventHistoryService{} + cfgNoEvents := &config.Config{ + MaxConcurrency: 5, + EventsDisabled: true, + } + ulNoConsumer, err := service.NewUserlogService( + service.Config(cfgNoEvents), + service.Stream(bus), + service.Store(sto), + service.Logger(log.NewLogger()), + service.Mux(chi.NewMux()), + service.GatewaySelector(gatewaySelector), + service.HistoryClient(&localEhc), + service.ValueClient(&vc), + service.RegisteredEvents([]events.Unmarshaller{ + events.SpaceDisabled{}, + }), + service.TraceProvider(trace.NewNoopTracerProvider()), + ) + Expect(err).ToNot(HaveOccurred()) + Expect(ulNoConsumer).ToNot(BeNil()) + + id := bus.publish(events.SpaceDisabled{Executant: &user.UserId{OpaqueId: "executinguserid"}}) + time.Sleep(500 * time.Millisecond) + + localEhc.On("GetEvents", mock.Anything, mock.Anything).Return(&ehsvc.GetEventsResponse{ + Events: []*ehmsg.Event{{Id: id}}, + }, nil) + + evs, err := ulNoConsumer.GetEvents(context.Background(), "userid") + Expect(err).ToNot(HaveOccurred()) + Expect(len(evs)).To(Equal(1)) + Expect(evs[0].Id).To(Equal(id)) + }) + AfterEach(func() { close(bus) })