Skip to content

[FR] Starting Event Listeners should be optional #1312

Description

@butonic

In #1124 I analyzed which services start what. In order to scale request handlers and event listeners independently we need to be able to run only the handler or the listener part of a service.

The affected services are:

We cannot set the service address to emptystring, eg. ACTIVITYLOG_HTTP_ADDR="", because that would prevent services from using a random port. So, we should introduce a new env var eg. ACTIVITYLOG_HTTP_DISABLE="true" or POLICIES_GRPC_DISABLE="true".

To disable event listeners we cannot just set the service specific OC_EVENTS_ENDPOINT to emptystring because we may need to differentiate between publishing and listening. So, we need a new env var, eg. GRAPH_EVENTS_DISABLE_CONSUMER to not listen for events.

Using grep we can find services that Consume
vendor/github.com/opencloud-eu/reva/v2/pkg/share/manager/jsoncs3/jsoncs3.go:            ch, err := events.Consume(m.eventStream, "jsoncs3sharemanager", _registeredEvents...)
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/decomposedfs.go:            ch, err := events.Consume(fs.stream, "dcfs", _registeredEvents...)
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/decomposedfs.go:          ch, err := events.Consume(fs.stream, "dcfs", _registeredEvents...)
services/userlog/pkg/service/service.go:        ch, err := events.Consume(o.Stream, "userlog", o.RegisteredEvents...)
services/sse/pkg/server/http/server.go: ch, err := events.Consume(options.Consumer, "sse-"+uuid.New().String(), options.RegisteredEvents...)
services/graph/pkg/service/v0/service.go:       evChannel, err := events.Consume(g.eventsConsumer, "graph", _registeredEvents...)
services/antivirus/pkg/service/service.go:      ch, err := events.Consume(natsStream, "antivirus", events.StartPostprocessingStep{})
services/policies/pkg/service/event/service.go: ch, err := events.Consume(s.stream, "policies", events.StartPostprocessingStep{})
services/storage-users/pkg/event/service.go:    ch, err := events.Consume(s.eventStream, consumerGroup, PurgeTrashBin{})
services/audit/pkg/command/server.go:                   evts, err := events.Consume(client, "audit", types.RegisteredEvents()...)
services/frontend/pkg/command/events.go:        evChannel, err := events.Consume(bus, "frontend", _registeredEvents...)
services/clientlog/pkg/service/service.go:      ch, err := events.Consume(o.Stream, "clientlog", o.RegisteredEvents...)
services/notifications/pkg/command/server.go:                   evts, err := events.Consume(client, "notifications", evs...)
services/activitylog/pkg/service/service.go:    ch, err := events.Consume(o.Stream, o.Config.Service.Name, o.RegisteredEvents...)

services/eventhistory/pkg/service/service.go:   ch, err := events.ConsumeAll(consumer, "evhistory")
Using grep we can find services that Publish
vendor/github.com/opencloud-eu/reva/v2/internal/http/services/sciencemesh/token.go:             if err := events.Publish(ctx, h.eventStream, events.ScienceMeshInviteTokenGenerated{
vendor/github.com/opencloud-eu/reva/v2/internal/grpc/interceptors/eventsmiddleware/events.go:                   if err := events.Publish(ctx, publisher, ev); err != nil {
vendor/github.com/opencloud-eu/reva/v2/pkg/share/manager/jsoncs3/jsoncs3.go:            if err := events.Publish(ctx, m.eventStream, events.ShareExpired{
vendor/github.com/opencloud-eu/reva/v2/pkg/share/manager/jsoncs3/jsoncs3.go:                                    if err := events.Publish(ctx, m.eventStream, events.ShareExpired{
vendor/github.com/opencloud-eu/reva/v2/pkg/share/manager/jsoncs3/jsoncs3.go:                                            if err := events.Publish(ctx, m.eventStream, events.ShareExpired{
vendor/github.com/opencloud-eu/reva/v2/pkg/share/manager/jsoncs3/jsoncs3.go:                                            if err := events.Publish(ctx, m.eventStream, events.ShareExpired{
vendor/github.com/opencloud-eu/reva/v2/pkg/share/manager/jsoncs3/jsoncs3.go:            if err := events.Publish(ctx, m.eventStream, events.ShareExpired{
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/events.go:  if err := events.Publish(ctx, fs.stream, ev); err != nil {
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/decomposedfs.go:                    if err := events.Publish(
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/decomposedfs.go:                    if err := events.Publish(ctx, fs.stream, events.BytesReceived{
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/pkg/decomposedfs/upload/upload.go:           if err := events.Publish(ctx, session.store.pub, events.BytesReceived{
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/fs/posix/tree/tree.go:       if err := events.Publish(context.Background(), t.es, ev); err != nil {
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/events.go:        if err := events.Publish(ctx, fs.stream, ev); err != nil {
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/decomposedfs.go:                  if err := events.Publish(
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/decomposedfs.go:                  if err := events.Publish(ctx, fs.stream, events.BytesReceived{
vendor/github.com/opencloud-eu/reva/v2/pkg/storage/utils/decomposedfs/upload/upload.go:         if err := events.Publish(ctx, session.store.pub, events.BytesReceived{
vendor/github.com/opencloud-eu/reva/v2/pkg/rhttp/datatx/datatx.go:      return events.Publish(context.Background(), publisher, uploadedEv)
services/userlog/pkg/service/service.go:                if err := events.Publish(ctx, ul.publisher, ev); err != nil {
services/graph/pkg/service/v0/graph.go:         if err := events.Publish(ctx, g.eventsPublisher, ev); err != nil {
services/graph/pkg/service/v0/tags.go:          if err := events.Publish(r.Context(), g.eventsPublisher, ev); err != nil {
services/graph/pkg/service/v0/tags.go:          if err := events.Publish(ctx, g.eventsPublisher, ev); err != nil {
services/graph/pkg/service/v0/personaldata.go:  if err := events.Publish(ctx, g.eventsPublisher, events.PersonalDataExtracted{
services/antivirus/pkg/service/service.go:              if err := events.Publish(ctx, s, events.PostprocessingStepFinished{
services/antivirus/pkg/service/service.go:      if err := events.Publish(ctx, s, events.PostprocessingStepFinished{
services/policies/pkg/service/event/service.go:         if err := events.Publish(ctx, s.stream, events.PostprocessingStepFinished{
services/storage-users/pkg/command/trash_bin.go:                        if err := events.Publish(c.Context, stream, event.PurgeTrashBin{ExecutionTime: time.Now()}); err != nil {
services/storage-users/pkg/command/uploads.go:                                  if err := events.Publish(context.Background(), stream, events.RestartPostprocessing{
services/storage-users/pkg/command/uploads.go:                                  if err := events.Publish(context.Background(), stream, events.ResumePostprocessing{
services/postprocessing/pkg/command/postprocessing.go:                  return events.Publish(context.Background(), stream, ev)
services/postprocessing/pkg/service/service.go:                         err := events.Publish(ctx, pps.pub, retryEvent)
services/postprocessing/pkg/service/service.go:                 if err := events.Publish(ctx, pps.pub, next); err != nil {
services/postprocessing/pkg/service/service.go:         if err := events.Publish(ctx, pps.pub, next); err != nil {
services/postprocessing/pkg/service/service.go:                 if err := events.Publish(ctx, pps.pub, events.RestartPostprocessing{
services/postprocessing/pkg/service/service.go: return events.Publish(ctx, pps.pub, pp.CurrentStep())
services/proxy/pkg/middleware/account_resolver.go:                      if err := events.Publish(req.Context(), m.eventsPublisher, event); err != nil {
services/proxy/pkg/staticroutes/backchannellogout.go:   if err := events.Publish(ctx, s.EventsPublisher, e); err != nil {
services/clientlog/pkg/service/service.go:      return events.Publish(context.Background(), cl.publisher, events.SendSSE{
services/notifications/pkg/command/send_email.go:                               err = events.Publish(c.Context, s, events.SendEmailsEvent{
services/notifications/pkg/command/send_email.go:                               err = events.Publish(c.Context, s, events.SendEmailsEvent{

Which gives us this table. I added the handlers manually. If a service has no handler or only produces events we can ignore it for this issue.

For services that have a handler and only consume events we can disable either by disabling the handler or the consumer:

service handlers consumer publisher
activitylog HTTP y -
eventhistory GRPC y -
frontend HTTP y -
sse HTTP y -

For services that have a handler and consume as well as publish events we need to be able to

  1. disable the handlers by setting the address to an empty string
  2. disable only the event consumer part so requests can still publish events.
service handlers consumer publisher
graph HTTP y y
policies GRPC y y
storage-users HTTP, GRPC y y
sharing GRPC y y
userlog HTTP y y

storage-users is a special case, because I think we should split it, which is why it is left out of the initial list of services.

Metadata

Metadata

Assignees

Projects

Status
In Progress

Milestone

No milestone

Relationships

None yet

Development

No branches or pull requests

Issue actions