diff --git a/.claude/skills/project-service-dev/references/nats-messaging.md b/.claude/skills/project-service-dev/references/nats-messaging.md index 69ad6a7e..d5c78f52 100644 --- a/.claude/skills/project-service-dev/references/nats-messaging.md +++ b/.claude/skills/project-service-dev/references/nats-messaging.md @@ -36,6 +36,22 @@ All six live as `Project*Subject` constants in `pkg/constants/nats.go`. `lfx.index.*` envelopes are owned by `lfx-v2-indexer-service`. `lfx.fga-sync.*` envelopes are owned by `lfx-v2-fga-sync`. For per-resource fields, see `docs/indexer-contract.md` and `docs/fga-contract.md`. The `lfx.projects-api.*.created` payloads are `events.ProjectDocumentCreatedMessage` and `events.ProjectLinkCreatedMessage` in `pkg/events/`; this service also subscribes to them itself (`internal/service/document_subscriber.go`) to send upload-notification emails. +## Inbound event subjects (handled by this service) + +```go +"lfx.projects-api.project_settings.updated" // handled by HandleProjectSettingsUpdated (role-change notifications) +"lfx.invite-service.invite_accepted" // handled by HandleInviteAccepted (promote email-only settings entries to LFID) +"lfx.v1-sync-helper.user.deleted" // handled by HandleUserDeleted (scrub deleted user's username from project settings) +"lfx.projects-api.project_document.created" // handled by HandleProjectDocumentCreated (upload notification emails) +"lfx.projects-api.project_link.created" // handled by HandleProjectLinkCreated +``` + +### Inbound event delivery (core NATS) + +Lifecycle handlers (`invite_accepted`, `user.deleted`, document/link created) subscribe via core NATS queue groups (`QueueSubscribe`), the same transport used by `lfx-v2-committee-service` for `user.deleted` (LFXV2-2645). v1-sync-helper publishes deletions with core NATS `Publish`. Messages emitted while this service is disconnected are not automatically replayed; durable JetStream delivery or a periodic reconciliation job would be a cross-service follow-up. + +Full project-settings scans inside those handlers use `settingsScanTimeout` (2 minutes), not `notificationTimeout` (5 seconds). + ## Owned KV buckets | Bucket | Purpose | History | diff --git a/cmd/project-api/main.go b/cmd/project-api/main.go index ff19d68b..a5957faa 100644 --- a/cmd/project-api/main.go +++ b/cmd/project-api/main.go @@ -497,9 +497,13 @@ func createNatsSubcriptions(ctx context.Context, svc *ProjectsAPI, natsConn *nat for _, eh := range []eventHandler{ {constants.ProjectSettingsUpdatedSubject, svc.service.HandleProjectSettingsUpdated}, {inviteapi.InviteServiceAcceptedSubject, svc.service.HandleInviteAccepted}, + {constants.V1SyncHelperUserDeletedSubject, svc.service.HandleUserDeleted}, {constants.ProjectDocumentCreatedSubject, svc.service.HandleProjectDocumentCreated}, {constants.ProjectLinkCreatedSubject, svc.service.HandleProjectLinkCreated}, } { + // Inbound lifecycle events use core NATS queue subscriptions, matching + // lfx-v2-committee-service (LFXV2-2645). v1-sync-helper publishes with core NATS + // Publish; missed messages during disconnect are not replayed automatically. slog.With("subject", eh.subject, "queue", queueName).Debug("subscribing to NATS subject") _, err := natsConn.QueueSubscribe(eh.subject, queueName, func(msg *nats.Msg) { msgCtx, end := internalnats.ExtractMsgContext(ctx, msg, eh.subject) diff --git a/internal/service/project_subscriber.go b/internal/service/project_subscriber.go index 47484007..9f91e0cf 100644 --- a/internal/service/project_subscriber.go +++ b/internal/service/project_subscriber.go @@ -13,6 +13,7 @@ import ( "time" emailapi "github.com/linuxfoundation/lfx-v2-email-service/pkg/api" + fgaconstants "github.com/linuxfoundation/lfx-v2-fga-sync/pkg/constants" indexerConstants "github.com/linuxfoundation/lfx-v2-indexer-service/pkg/constants" indexerTypes "github.com/linuxfoundation/lfx-v2-indexer-service/pkg/types" inviteapi "github.com/linuxfoundation/lfx-v2-invite-service/pkg/api" @@ -29,6 +30,21 @@ import ( // auth-service actor name lookup all run under this deadline. const notificationTimeout = 5 * time.Second +// settingsScanTimeout caps ListAllProjectsSettings and per-project reconciliation work. +// Unlike notificationTimeout (single RPC), a full project-settings bucket scan can require +// many sequential KV reads under load. +const settingsScanTimeout = 2 * time.Minute + +// scrubMaxRetries is the number of attempts for settings KV writes and indexer/FGA publishes +// after a successful scrub. Retries are independent of the username match so a transient +// conflict or NATS failure does not leave access tuples stale after settings were scrubbed. +const scrubMaxRetries = 4 + +// serviceAuthBearer is the static JWT audience token used for background NATS handler +// side effects (indexer/FGA) that have no originating HTTP request context. +// Pattern matches lfx-v2-committee-service message_handler.go. +const serviceAuthBearer = "Bearer lfx-v2-project-service" + const ( roleWriter = "Writer" roleAuditor = "Auditor" @@ -402,7 +418,7 @@ func (s *ProjectsService) HandleInviteAccepted(ctx context.Context, msg domain.M } // Scan all project settings for email-only entries that match the recipient. - listCtx, listCancel := context.WithTimeout(ctx, notificationTimeout) + listCtx, listCancel := context.WithTimeout(ctx, settingsScanTimeout) allSettings, listErr := s.ProjectRepository.ListAllProjectsSettings(listCtx) listCancel() if listErr != nil { @@ -416,7 +432,7 @@ func (s *ProjectsService) HandleInviteAccepted(ctx context.Context, msg domain.M continue } projectUID := candidate.UID - promoteCtx, promoteCancel := context.WithTimeout(ctx, notificationTimeout) + promoteCtx, promoteCancel := context.WithTimeout(ctx, settingsScanTimeout) s.promoteInvitedUserInProjectSettings(promoteCtx, projectUID, normalizedEmail, event.AcceptedBy, event.UID, event.Role) promoteCancel() } @@ -504,7 +520,7 @@ func (s *ProjectsService) promoteInvitedUserInProjectSettings(ctx context.Contex Data: *settings, IndexingConfig: settings.IndexingConfig(projectUID), } - if indexErr := s.MessageBuilder.SendIndexerMessage(ctx, constants.IndexProjectSettingsSubject, indexMsg, false); indexErr != nil { + if indexErr := s.MessageBuilder.SendIndexerMessage(ctxWithServiceAuth(ctx), constants.IndexProjectSettingsSubject, indexMsg, false); indexErr != nil { slog.WarnContext(ctx, "project_subscriber: failed to reindex project settings after invite acceptance", constants.ErrKey, indexErr, "project_uid", projectUID) } @@ -520,6 +536,271 @@ func (s *ProjectsService) promoteInvitedUserInProjectSettings(ctx context.Contex } } +// HandleUserDeleted scrubs the deleted user's username from project settings. +// Best-effort: partial failures are logged but do not block the overall scrub. +// Unlike committee data, project settings have no separate member records — only +// the settings writers/auditors/meeting coordinators and named role fields are updated. +func (s *ProjectsService) HandleUserDeleted(ctx context.Context, msg domain.Message) error { + var event events.V1UserDeletedEvent + if err := json.Unmarshal(msg.Data(), &event); err != nil { + slog.WarnContext(ctx, "project_subscriber: failed to unmarshal user.deleted event", constants.ErrKey, err) + return nil + } + if strings.TrimSpace(event.Username) == "" { + slog.WarnContext(ctx, "project_subscriber: user.deleted event missing username — nothing to scrub") + return nil + } + + slog.InfoContext(ctx, "project_subscriber: scrubbing deleted user's username from project settings") + + listCtx, listCancel := context.WithTimeout(ctx, settingsScanTimeout) + allSettings, listErr := s.ProjectRepository.ListAllProjectsSettings(listCtx) + listCancel() + if listErr != nil { + slog.WarnContext(ctx, "project_subscriber: failed to list project settings for username scrub", + constants.ErrKey, listErr) + return nil + } + + for _, candidate := range allSettings { + if !projectSettingsHasUsername(candidate, event.Username) { + continue + } + scrubCtx, scrubCancel := context.WithTimeout(ctx, settingsScanTimeout) + s.scrubProjectSettingsUsername(scrubCtx, candidate.UID, event.Username, event.Email) + scrubCancel() + } + + return nil +} + +// usernameMatches reports whether storedUsername represents the same LFID as deletedUsername. +func usernameMatches(deletedUsername, storedUsername string) bool { + deleted := strings.TrimSpace(deletedUsername) + stored := strings.TrimSpace(storedUsername) + if deleted == "" || stored == "" { + return false + } + return strings.EqualFold(deleted, stored) +} + +// projectSettingsHasUsername reports whether any role-bearing field in settings carries username. +func projectSettingsHasUsername(s *models.ProjectSettings, username string) bool { + if s == nil { + return false + } + for _, u := range s.Writers { + if usernameMatches(username, u.Username) { + return true + } + } + for _, u := range s.Auditors { + if usernameMatches(username, u.Username) { + return true + } + } + for _, u := range s.MeetingCoordinators { + if usernameMatches(username, u.Username) { + return true + } + } + if s.ExecutiveDirector != nil && usernameMatches(username, s.ExecutiveDirector.Username) { + return true + } + if s.ProgramManager != nil && usernameMatches(username, s.ProgramManager.Username) { + return true + } + if s.OpportunityOwner != nil && usernameMatches(username, s.OpportunityOwner.Username) { + return true + } + return false +} + +// clearUsernameInSettings clears username on every matching entry in settings +// that still represents the deleted account. Returns true when at least one field was changed. +func (s *ProjectsService) clearUsernameInSettings(ctx context.Context, settings *models.ProjectSettings, username, deletedEmail string) bool { + changed := false + clearIfMatch := func(u *models.UserInfo) { + if u == nil || !usernameMatches(username, u.Username) { + return + } + if !s.shouldScrubSettingsUsername(ctx, *u, username, deletedEmail) { + slog.DebugContext(ctx, "project_subscriber: skipping username scrub — email still maps to active LFID", + "username", username, "email", u.Email) + return + } + u.Username = "" + changed = true + } + + for i := range settings.Writers { + clearIfMatch(&settings.Writers[i]) + } + for i := range settings.Auditors { + clearIfMatch(&settings.Auditors[i]) + } + for i := range settings.MeetingCoordinators { + clearIfMatch(&settings.MeetingCoordinators[i]) + } + clearIfMatch(settings.ExecutiveDirector) + clearIfMatch(settings.ProgramManager) + clearIfMatch(settings.OpportunityOwner) + return changed +} + +// shouldScrubSettingsUsername reports whether a settings entry carrying deletedUsername should +// be cleared. When the deletion event carries an email, entries with a different non-empty email +// are treated as reuse and skipped; a matching entry email is a definitive identification and +// scrubs without auth lookup. When the event omits email, auth is consulted for entries that +// have an email so a reassigned LFID reused by a new account is not scrubbed. +// +// Entries without email scrub when the username matches only when the deletion event +// also omits email (M2M/legacy). When the event carries email, username-only entries are +// skipped because they cannot be verified as the deleted account. +func (s *ProjectsService) shouldScrubSettingsUsername(ctx context.Context, u models.UserInfo, deletedUsername, deletedEmail string) bool { + entryEmail := strings.ToLower(strings.TrimSpace(u.Email)) + deletedEmailNorm := strings.ToLower(strings.TrimSpace(deletedEmail)) + if deletedEmailNorm != "" { + if entryEmail == "" { + return false + } + if entryEmail != deletedEmailNorm { + return false + } + return true // definitive match — scrub without auth lookup + } + if entryEmail == "" { + return true + } + if s.UserReader == nil { + return true + } + + lookupCtx, cancel := context.WithTimeout(ctx, notificationTimeout) + defer cancel() + + resolved, err := s.UserReader.UsernameByEmail(lookupCtx, entryEmail) + if err != nil { + if errors.Is(err, domain.ErrUserNotFound) { + return true + } + slog.WarnContext(ctx, "project_subscriber: auth lookup failed during username scrub — skipping entry", + constants.ErrKey, err) + return false + } + return !usernameMatches(resolved, deletedUsername) +} + +// scrubProjectSettingsUsername fetches settings for a single project, clears the +// username on any matching entry, persists, and reindexes. Retries on revision conflicts. +func (s *ProjectsService) scrubProjectSettingsUsername(ctx context.Context, projectUID, username, deletedEmail string) { + for attempt := 0; attempt < scrubMaxRetries; attempt++ { + settings, revision, err := s.ProjectRepository.GetProjectSettingsWithRevision(ctx, projectUID) + if err != nil { + slog.DebugContext(ctx, "project_subscriber: failed to get settings for username scrub — skipping", + constants.ErrKey, err, "project_uid", projectUID) + return + } + if settings == nil { + return + } + + if !s.clearUsernameInSettings(ctx, settings, username, deletedEmail) { + return + } + + updateErr := s.ProjectRepository.UpdateProjectSettings(ctx, settings, revision) + if updateErr == nil { + slog.InfoContext(ctx, "project_subscriber: cleared username from project settings", + "project_uid", projectUID) + s.publishProjectSettingsScrubSideEffects(ctx, projectUID) + return + } + if !errors.Is(updateErr, domain.ErrRevisionMismatch) || attempt == scrubMaxRetries-1 { + slog.WarnContext(ctx, "project_subscriber: failed to clear username from project settings", + constants.ErrKey, updateErr, "project_uid", projectUID) + return + } + slog.DebugContext(ctx, "project_subscriber: revision mismatch scrubbing username — retrying", + "attempt", attempt+1, "project_uid", projectUID) + } +} + +// publishProjectSettingsScrubSideEffects reindexes scrubbed settings and refreshes OpenFGA +// access tuples. Each attempt reloads the current KV record before publishing; indexer and FGA +// projections are full-state and idempotent. ProjectSettingsUpdatedSubject is intentionally +// omitted to avoid role-change emails. +func (s *ProjectsService) publishProjectSettingsScrubSideEffects(ctx context.Context, projectUID string) { + ctx = ctxWithServiceAuth(ctx) + for attempt := 0; attempt < scrubMaxRetries; attempt++ { + settings, _, err := s.ProjectRepository.GetProjectSettingsWithRevision(ctx, projectUID) + if err != nil { + if attempt == scrubMaxRetries-1 { + slog.WarnContext(ctx, "project_subscriber: failed to reload settings for scrub side effects", + constants.ErrKey, err, "project_uid", projectUID, "attempts", scrubMaxRetries) + return + } + slog.DebugContext(ctx, "project_subscriber: retrying scrub side effects after settings reload failure", + "attempt", attempt+1, "project_uid", projectUID) + continue + } + if settings == nil { + return + } + + projectBase, baseErr := s.ProjectRepository.GetProjectBase(ctx, projectUID) + if baseErr != nil { + slog.WarnContext(ctx, "project_subscriber: failed to load project for FGA refresh after username scrub", + constants.ErrKey, baseErr, "project_uid", projectUID) + if attempt == scrubMaxRetries-1 { + return + } + slog.DebugContext(ctx, "project_subscriber: retrying scrub side effects after project load failure", + "attempt", attempt+1, "project_uid", projectUID) + continue + } + + indexMsg := indexerTypes.IndexerMessageEnvelope{ + Action: indexerConstants.ActionUpdated, + Data: *settings, + IndexingConfig: settings.IndexingConfig(projectUID), + } + indexErr := s.MessageBuilder.SendIndexerMessage(ctx, constants.IndexProjectSettingsSubject, indexMsg, false) + + fgaMsg := buildFGAUpdateAccessMessage(projectBase, settings) + accessErr := s.MessageBuilder.PublishAccessMessage(ctx, fgaconstants.GenericUpdateAccessSubject, fgaMsg) + + if indexErr == nil && accessErr == nil { + return + } + + if attempt == scrubMaxRetries-1 { + if indexErr != nil { + slog.WarnContext(ctx, "project_subscriber: failed to reindex project settings after username scrub", + constants.ErrKey, indexErr, "project_uid", projectUID, "attempts", scrubMaxRetries) + } + if accessErr != nil { + slog.WarnContext(ctx, "project_subscriber: failed to publish FGA update after username scrub", + constants.ErrKey, accessErr, "project_uid", projectUID, "attempts", scrubMaxRetries) + } + return + } + + slog.DebugContext(ctx, "project_subscriber: retrying scrub side effects after reload", + "attempt", attempt+1, "project_uid", projectUID) + } +} + +// ctxWithServiceAuth returns ctx with a service-identity bearer when no inbound JWT is present. +// NATS queue subscribers have no HTTP middleware, but indexer V2 transactions require an +// authorization header in the message envelope. +func ctxWithServiceAuth(ctx context.Context) context.Context { + if _, ok := ctx.Value(constants.AuthorizationContextID).(string); ok { + return ctx + } + return context.WithValue(ctx, constants.AuthorizationContextID, serviceAuthBearer) +} + // buildProjectURL constructs the deep-link URL for a project's overview page. func buildProjectURL(baseURL, slug string) string { base := strings.TrimRight(baseURL, "/") + "/project/overview" diff --git a/internal/service/project_subscriber_test.go b/internal/service/project_subscriber_test.go index 6ff2f37e..6c192d14 100644 --- a/internal/service/project_subscriber_test.go +++ b/internal/service/project_subscriber_test.go @@ -12,6 +12,7 @@ import ( "time" emailapi "github.com/linuxfoundation/lfx-v2-email-service/pkg/api" + fgaconstants "github.com/linuxfoundation/lfx-v2-fga-sync/pkg/constants" indexerTypes "github.com/linuxfoundation/lfx-v2-indexer-service/pkg/types" inviteapi "github.com/linuxfoundation/lfx-v2-invite-service/pkg/api" "github.com/stretchr/testify/assert" @@ -20,6 +21,7 @@ import ( "github.com/linuxfoundation/lfx-v2-project-service/internal/domain" "github.com/linuxfoundation/lfx-v2-project-service/internal/domain/models" + "github.com/linuxfoundation/lfx-v2-project-service/pkg/constants" "github.com/linuxfoundation/lfx-v2-project-service/pkg/events" ) @@ -1148,6 +1150,437 @@ func TestHandleInviteAccepted(t *testing.T) { } } +func TestCtxWithServiceAuth(t *testing.T) { + t.Run("injects service bearer when absent", func(t *testing.T) { + ctx := ctxWithServiceAuth(context.Background()) + auth, ok := ctx.Value(constants.AuthorizationContextID).(string) + require.True(t, ok) + assert.Equal(t, serviceAuthBearer, auth) + }) + + t.Run("preserves existing bearer", func(t *testing.T) { + existing := "Bearer caller-jwt" + ctx := context.WithValue(context.Background(), constants.AuthorizationContextID, existing) + ctx = ctxWithServiceAuth(ctx) + auth, ok := ctx.Value(constants.AuthorizationContextID).(string) + require.True(t, ok) + assert.Equal(t, existing, auth) + }) +} + +func TestHandleUserDeleted(t *testing.T) { + const deletedUsername = "deleted.user" + const projectUID = "proj-1" + const project2UID = "proj-2" + + makeEvent := func(username string) events.V1UserDeletedEvent { + return events.V1UserDeletedEvent{Username: username} + } + + tests := []struct { + name string + payload any + setupRepo func(*domain.MockProjectRepository) + setupMsg func(*domain.MockMessageBuilder) + setupUserReader func(*domain.MockUserReader) + }{ + { + name: "malformed payload — returns nil without crashing", + payload: []byte("not json"), + }, + { + name: "empty username — discarded", + payload: makeEvent(""), + }, + { + name: "case-insensitive prefilter — stored Alice scrubbed for event alice", + payload: makeEvent("alice"), + setupRepo: func(r *domain.MockProjectRepository) { + settings := &models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: "Alice", Email: "deleted@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{settings}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(settings, uint64(1), nil).Times(2) + r.On("UpdateProjectSettings", mock.Anything, mock.MatchedBy(func(s *models.ProjectSettings) bool { + return len(s.Writers) == 1 && s.Writers[0].Username == "" + }), uint64(1)).Return(nil) + r.On("GetProjectBase", mock.Anything, projectUID).Return(&models.ProjectBase{UID: projectUID}, nil) + }, + setupMsg: func(m *domain.MockMessageBuilder) { + m.On("SendIndexerMessage", mock.Anything, constants.IndexProjectSettingsSubject, mock.Anything, false).Return(nil) + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")).Return(nil) + }, + }, + { + name: "settings match — writer username cleared and reindexed", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + settings := &models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: deletedUsername, Email: "deleted@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{settings}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(settings, uint64(1), nil).Times(2) + r.On("UpdateProjectSettings", mock.Anything, mock.MatchedBy(func(s *models.ProjectSettings) bool { + return len(s.Writers) == 1 && s.Writers[0].Username == "" && s.Writers[0].Email == "deleted@example.com" + }), uint64(1)).Return(nil) + r.On("GetProjectBase", mock.Anything, projectUID).Return(&models.ProjectBase{UID: projectUID}, nil) + }, + setupMsg: func(m *domain.MockMessageBuilder) { + m.On("SendIndexerMessage", mock.Anything, constants.IndexProjectSettingsSubject, mock.Anything, false).Return(nil) + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")).Return(nil) + }, + }, + { + name: "settings no match — no update", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + settings := &models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: "other", Email: "other@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{settings}, nil) + }, + }, + { + name: "conflict retry — succeeds on second attempt", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + firstRead := &models.ProjectSettings{ + UID: projectUID, + Auditors: []models.UserInfo{{Username: deletedUsername, Email: "deleted@example.com"}}, + } + secondRead := &models.ProjectSettings{ + UID: projectUID, + Auditors: []models.UserInfo{{Username: deletedUsername, Email: "deleted@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{firstRead}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(firstRead, uint64(1), nil).Once() + r.On("UpdateProjectSettings", mock.Anything, mock.Anything, uint64(1)). + Return(domain.ErrRevisionMismatch).Once() + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(secondRead, uint64(2), nil).Once() + r.On("UpdateProjectSettings", mock.Anything, mock.MatchedBy(func(s *models.ProjectSettings) bool { + return len(s.Auditors) == 1 && s.Auditors[0].Username == "" + }), uint64(2)).Return(nil).Once() + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(secondRead, uint64(2), nil).Once() + r.On("GetProjectBase", mock.Anything, projectUID).Return(&models.ProjectBase{UID: projectUID}, nil) + }, + setupMsg: func(m *domain.MockMessageBuilder) { + m.On("SendIndexerMessage", mock.Anything, constants.IndexProjectSettingsSubject, mock.Anything, false).Return(nil) + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")).Return(nil) + }, + }, + { + name: "reuse guard — skips scrub when email still maps to username", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + settings := &models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: deletedUsername, Email: "reassigned@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{settings}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(settings, uint64(1), nil).Once() + }, + setupUserReader: func(u *domain.MockUserReader) { + u.On("UsernameByEmail", mock.Anything, "reassigned@example.com").Return(deletedUsername, nil) + }, + }, + { + name: "ListAllProjectsSettings error — early return", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + r.On("ListAllProjectsSettings", mock.Anything).Return(nil, errors.New("kv unavailable")) + }, + }, + { + name: "meeting coordinator and named roles cleared", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + settings := &models.ProjectSettings{ + UID: projectUID, + MeetingCoordinators: []models.UserInfo{{Username: deletedUsername, Email: "mc@example.com"}}, + ExecutiveDirector: &models.UserInfo{Username: deletedUsername, Email: "ed@example.com"}, + ProgramManager: &models.UserInfo{Username: deletedUsername, Email: "pm@example.com"}, + OpportunityOwner: &models.UserInfo{Username: deletedUsername, Email: "oo@example.com"}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{settings}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(settings, uint64(1), nil).Times(2) + r.On("UpdateProjectSettings", mock.Anything, mock.MatchedBy(func(s *models.ProjectSettings) bool { + return len(s.MeetingCoordinators) == 1 && s.MeetingCoordinators[0].Username == "" && + s.ExecutiveDirector != nil && s.ExecutiveDirector.Username == "" && + s.ProgramManager != nil && s.ProgramManager.Username == "" && + s.OpportunityOwner != nil && s.OpportunityOwner.Username == "" + }), uint64(1)).Return(nil) + r.On("GetProjectBase", mock.Anything, projectUID).Return(&models.ProjectBase{UID: projectUID}, nil) + }, + setupMsg: func(m *domain.MockMessageBuilder) { + m.On("SendIndexerMessage", mock.Anything, constants.IndexProjectSettingsSubject, mock.Anything, false).Return(nil) + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")).Return(nil) + }, + }, + { + name: "GetProjectBase retry — succeeds on second attempt", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + settings := &models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: deletedUsername, Email: "deleted@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{settings}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(settings, uint64(1), nil).Times(3) + r.On("UpdateProjectSettings", mock.Anything, mock.Anything, uint64(1)).Return(nil) + r.On("GetProjectBase", mock.Anything, projectUID). + Return((*models.ProjectBase)(nil), errors.New("transient read failure")).Once() + r.On("GetProjectBase", mock.Anything, projectUID). + Return(&models.ProjectBase{UID: projectUID}, nil).Once() + }, + setupMsg: func(m *domain.MockMessageBuilder) { + m.On("SendIndexerMessage", mock.Anything, constants.IndexProjectSettingsSubject, mock.Anything, false).Return(nil).Maybe() + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")).Return(nil).Once() + }, + }, + { + name: "FGA publish retry — succeeds on second attempt", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + settings := &models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: deletedUsername, Email: "deleted@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{settings}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(settings, uint64(1), nil).Times(3) + r.On("UpdateProjectSettings", mock.Anything, mock.Anything, uint64(1)).Return(nil) + r.On("GetProjectBase", mock.Anything, projectUID).Return(&models.ProjectBase{UID: projectUID}, nil).Times(2) + }, + setupMsg: func(m *domain.MockMessageBuilder) { + m.On("SendIndexerMessage", mock.Anything, constants.IndexProjectSettingsSubject, mock.Anything, false).Return(nil).Maybe() + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")). + Return(errors.New("transient nats failure")).Once() + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")). + Return(nil).Once() + }, + }, + { + name: "multiple projects — only matching project updated", + payload: makeEvent(deletedUsername), + setupRepo: func(r *domain.MockProjectRepository) { + match := &models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: deletedUsername, Email: "deleted@example.com"}}, + } + other := &models.ProjectSettings{ + UID: project2UID, + Writers: []models.UserInfo{{Username: "other", Email: "other@example.com"}}, + } + r.On("ListAllProjectsSettings", mock.Anything).Return([]*models.ProjectSettings{match, other}, nil) + r.On("GetProjectSettingsWithRevision", mock.Anything, projectUID).Return(match, uint64(1), nil).Times(2) + r.On("UpdateProjectSettings", mock.Anything, mock.Anything, uint64(1)).Return(nil) + r.On("GetProjectBase", mock.Anything, projectUID).Return(&models.ProjectBase{UID: projectUID}, nil) + }, + setupMsg: func(m *domain.MockMessageBuilder) { + m.On("SendIndexerMessage", mock.Anything, constants.IndexProjectSettingsSubject, mock.Anything, false).Return(nil) + m.On("PublishAccessMessage", mock.Anything, fgaconstants.GenericUpdateAccessSubject, mock.AnythingOfType("types.GenericFGAMessage")).Return(nil) + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + mockRepo := &domain.MockProjectRepository{} + mockMsg := &domain.MockMessageBuilder{} + var mockUserReader *domain.MockUserReader + if tt.setupRepo != nil { + tt.setupRepo(mockRepo) + } + if tt.setupMsg != nil { + tt.setupMsg(mockMsg) + } + + svc := &ProjectsService{ + ProjectRepository: mockRepo, + MessageBuilder: mockMsg, + } + if tt.setupUserReader != nil { + mockUserReader = &domain.MockUserReader{} + svc.UserReader = mockUserReader + tt.setupUserReader(mockUserReader) + } + + var data []byte + if raw, ok := tt.payload.([]byte); ok { + data = raw + } else { + data = marshalEvent(t, tt.payload) + } + msg := domain.NewMockMessage(data, constants.V1SyncHelperUserDeletedSubject) + err := svc.HandleUserDeleted(context.Background(), msg) + assert.NoError(t, err) + + mockRepo.AssertExpectations(t) + mockMsg.AssertExpectations(t) + if mockUserReader != nil { + mockUserReader.AssertExpectations(t) + } + }) + } +} + +func TestShouldScrubSettingsUsername(t *testing.T) { + const deletedUsername = "deleted.user" + ctx := context.Background() + + tests := []struct { + name string + entry models.UserInfo + deletedEmail string + setupUser func(*domain.MockUserReader) + userReader bool + skipAuthLookup bool + want bool + }{ + { + name: "no email — always scrub", + entry: models.UserInfo{Username: deletedUsername}, + want: true, + }, + { + name: "nil user reader with email — scrub", + entry: models.UserInfo{Username: deletedUsername, Email: "a@example.com"}, + want: true, + }, + { + name: "email maps to same username — do not scrub", + entry: models.UserInfo{Username: deletedUsername, Email: "active@example.com"}, + setupUser: func(u *domain.MockUserReader) { + u.On("UsernameByEmail", mock.Anything, "active@example.com").Return(deletedUsername, nil) + }, + userReader: true, + want: false, + }, + { + name: "email maps to different username — scrub", + entry: models.UserInfo{Username: deletedUsername, Email: "reassigned@example.com"}, + setupUser: func(u *domain.MockUserReader) { + u.On("UsernameByEmail", mock.Anything, "reassigned@example.com").Return("new.user", nil) + }, + userReader: true, + want: true, + }, + { + name: "email not found in auth — scrub", + entry: models.UserInfo{Username: deletedUsername, Email: "gone@example.com"}, + setupUser: func(u *domain.MockUserReader) { + u.On("UsernameByEmail", mock.Anything, "gone@example.com").Return("", domain.ErrUserNotFound) + }, + userReader: true, + want: true, + }, + { + name: "auth lookup error — skip entry", + entry: models.UserInfo{Username: deletedUsername, Email: "err@example.com"}, + setupUser: func(u *domain.MockUserReader) { + u.On("UsernameByEmail", mock.Anything, "err@example.com").Return("", errors.New("auth unavailable")) + }, + userReader: true, + want: false, + }, + { + name: "event email differs from entry email — do not scrub", + entry: models.UserInfo{Username: deletedUsername, Email: "new@example.com"}, + deletedEmail: "old@example.com", + want: false, + }, + { + name: "event email with username-only entry — do not scrub", + entry: models.UserInfo{Username: deletedUsername}, + deletedEmail: "deleted@example.com", + want: false, + }, + { + name: "event email matches entry email — scrub without auth lookup", + entry: models.UserInfo{Username: deletedUsername, Email: "deleted@example.com"}, + deletedEmail: "deleted@example.com", + setupUser: func(u *domain.MockUserReader) { + u.On("UsernameByEmail", mock.Anything, "deleted@example.com").Return(deletedUsername, nil) + }, + userReader: true, + skipAuthLookup: true, + want: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + svc := &ProjectsService{} + if tt.userReader { + mockUser := &domain.MockUserReader{} + if tt.setupUser != nil { + tt.setupUser(mockUser) + } + svc.UserReader = mockUser + got := svc.shouldScrubSettingsUsername(ctx, tt.entry, deletedUsername, tt.deletedEmail) + assert.Equal(t, tt.want, got) + if tt.skipAuthLookup { + mockUser.AssertNumberOfCalls(t, "UsernameByEmail", 0) + } else { + mockUser.AssertExpectations(t) + } + return + } + got := svc.shouldScrubSettingsUsername(ctx, tt.entry, deletedUsername, tt.deletedEmail) + assert.Equal(t, tt.want, got) + }) + } +} + +func TestScrubProjectSettingsUsernameRetryExhaustion(t *testing.T) { + const ( + deletedUsername = "deleted.user" + projectUID = "proj-1" + ) + mockRepo := &domain.MockProjectRepository{} + for range scrubMaxRetries { + mockRepo.On("GetProjectSettingsWithRevision", mock.Anything, projectUID). + Return(&models.ProjectSettings{ + UID: projectUID, + Writers: []models.UserInfo{{Username: deletedUsername, Email: "deleted@example.com"}}, + }, uint64(1), nil).Once() + } + mockRepo.On("UpdateProjectSettings", mock.Anything, mock.Anything, uint64(1)). + Return(domain.ErrRevisionMismatch) + + svc := &ProjectsService{ProjectRepository: mockRepo} + svc.scrubProjectSettingsUsername(context.Background(), projectUID, deletedUsername, "") + + mockRepo.AssertNumberOfCalls(t, "UpdateProjectSettings", scrubMaxRetries) + mockRepo.AssertExpectations(t) +} + +func TestProjectSettingsHasUsername(t *testing.T) { + const username = "alice" + tests := []struct { + name string + settings *models.ProjectSettings + want bool + }{ + {name: "nil settings", settings: nil, want: false}, + {name: "writer match", settings: &models.ProjectSettings{Writers: []models.UserInfo{{Username: username}}}, want: true}, + {name: "writer case-insensitive match", settings: &models.ProjectSettings{Writers: []models.UserInfo{{Username: "Alice"}}}, want: true}, + {name: "auditor match", settings: &models.ProjectSettings{Auditors: []models.UserInfo{{Username: username}}}, want: true}, + {name: "meeting coordinator match", settings: &models.ProjectSettings{MeetingCoordinators: []models.UserInfo{{Username: username}}}, want: true}, + {name: "executive director match", settings: &models.ProjectSettings{ExecutiveDirector: &models.UserInfo{Username: username}}, want: true}, + {name: "program manager match", settings: &models.ProjectSettings{ProgramManager: &models.UserInfo{Username: username}}, want: true}, + {name: "opportunity owner match", settings: &models.ProjectSettings{OpportunityOwner: &models.UserInfo{Username: username}}, want: true}, + {name: "no match", settings: &models.ProjectSettings{Writers: []models.UserInfo{{Username: "other"}}}, want: false}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, projectSettingsHasUsername(tt.settings, username)) + }) + } +} + // Compile-time checks for imported API types. var ( _ emailapi.SendEmailRequest diff --git a/pkg/constants/nats.go b/pkg/constants/nats.go index 46e91e24..8d53c39d 100644 --- a/pkg/constants/nats.go +++ b/pkg/constants/nats.go @@ -109,3 +109,12 @@ const ( // Request: plain-text email. Reply: plain-text username on success, JSON error envelope on miss. AuthEmailToUsernameSubject = "lfx.auth-service.email_to_username" ) + +// NATS subjects consumed from other services. +const ( + // V1SyncHelperUserDeletedSubject is emitted by v1-sync-helper when a merged user record is + // soft-deleted. The project service subscribes to scrub the deleted user's username from + // project settings writers/auditors/meeting coordinators and named role fields. + // Payload: JSON-encoded events.V1UserDeletedEvent. + V1SyncHelperUserDeletedSubject = "lfx.v1-sync-helper.user.deleted" +) diff --git a/pkg/events/project.go b/pkg/events/project.go index 068a1021..28fc9437 100644 --- a/pkg/events/project.go +++ b/pkg/events/project.go @@ -90,3 +90,12 @@ type InviteAccepted struct { InviteUID string `json:"invite_uid"` Username string `json:"username"` } + +// V1UserDeletedEvent is the payload published by v1-sync-helper on +// lfx.v1-sync-helper.user.deleted when a merged user record is soft-deleted. +// Username is the normalized LFID; Email is the deleted account's primary email when +// available so scrubbers can distinguish LFID reuse. +type V1UserDeletedEvent struct { + Username string `json:"username"` + Email string `json:"email,omitempty"` +}