diff --git a/backend/modules/eventprocessing/connectors/repository.go b/backend/modules/eventprocessing/connectors/repository.go index de45bd82d..24051ffdd 100644 --- a/backend/modules/eventprocessing/connectors/repository.go +++ b/backend/modules/eventprocessing/connectors/repository.go @@ -71,7 +71,7 @@ type PipelineRepository interface { List(tenantId string) []domain.Pipeline GetByRelPath(relPath string) *domain.Pipeline Create(relPath string, content []byte, tenantId string) (*domain.Pipeline, error) - Update(relPath string, content []byte) (*domain.Pipeline, error) + Update(relPath string, content []byte, tenantId string) (*domain.Pipeline, error) Delete(relPath string) error SetEnabled(tenantId, relPath string, active bool) error } diff --git a/backend/modules/eventprocessing/repository/pipeline_store.go b/backend/modules/eventprocessing/repository/pipeline_store.go index 2f229c60c..ce3a3fdb8 100644 --- a/backend/modules/eventprocessing/repository/pipeline_store.go +++ b/backend/modules/eventprocessing/repository/pipeline_store.go @@ -196,7 +196,21 @@ func (s *PipelineStore) Create(relPath string, content []byte, tenantId string) return &cp, nil } -func (s *PipelineStore) Update(relPath string, content []byte) (*domain.Pipeline, error) { +func (s *PipelineStore) Update(relPath string, content []byte,tenantId string) (*domain.Pipeline, error) { + if tenantId != "" { + injected, err := withTenantID(content, tenantId) + if err != nil { + return nil, fmt.Errorf("invalid filter content: %w", err) + } + content = injected + + base := filepath.Base(relPath) + ext := filepath.Ext(base) + name := strings.TrimSuffix(base, ext) + relPath = filepath.ToSlash(filepath.Join(tenantId, name+"-"+tenantId+ext)) + } + + s.mu.Lock() defer s.mu.Unlock() existing, ok := s.filters[relPath] diff --git a/backend/modules/eventprocessing/usecase/pipeline.go b/backend/modules/eventprocessing/usecase/pipeline.go index c505a3909..9cef24efd 100644 --- a/backend/modules/eventprocessing/usecase/pipeline.go +++ b/backend/modules/eventprocessing/usecase/pipeline.go @@ -241,7 +241,7 @@ func (u *pipelineUsecase) Update(ctx context.Context, req dto.UpdatePipelineRequ if err != nil { return nil, fmt.Errorf("%w: %v", domain.ErrPipelineInvalidContent, err) } - entry, err := u.store.Update(req.RelPath, []byte(content)) + entry, err := u.store.Update(req.RelPath, []byte(content),authz.TenantIDFromContext(ctx)) if err != nil { return nil, err }