mirror of
https://github.com/opencloud-eu/opencloud.git
synced 2026-09-21 03:25:24 -04:00
userlog service, posibility to disable/enable event listener and http server
This commit is contained in:
1 parent
d6c6b6fd0d
commit
a9cede5438
6 files changed
+111
-14
No files matched your search
@@ -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")
|
||||
}
|
||||
|
||||
{
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -228,7 +228,6 @@ func RequireAdminOrSecret(rm *roles.Manager, secret string) func(http.HandlerFun
|
||||
}
|
||||
|
||||
errorcode.ItemNotFound.Render(w, r, http.StatusNotFound, "Not found")
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
|
||||
Reference in new issue
Block a user