fix: restore processor.go to correct EventHandler architecture (was slack.Client corruption)

This commit is contained in:
2026-06-29 20:24:40 +00:00
parent c8b2dbe9dd
commit 72b50b88f2
+525 -170
View File
@@ -2,209 +2,564 @@ package event
import (
"context"
"fmt"
"strings"
"time"
"github.com/rs/zerolog"
"github.com/vincentc-afk/gitea-notification-hub/internal/identity"
"github.com/vincentc-afk/gitea-notification-hub/internal/slack"
"github.com/vincentc-afk/gitea-notification-hub/internal/config"
"github.com/vincentc-afk/gitea-notification-hub/internal/webhook"
)
// Processor handles processing of webhook events
// IdentityResolver resolves Gitea users to external identities (e.g., Slack)
type IdentityResolver interface {
Resolve(ctx context.Context, user User) (*ResolvedIdentity, error)
}
// ResolvedIdentity represents a resolved external identity
type ResolvedIdentity struct {
Email string `json:"email"`
SlackID string `json:"slack_id"`
SlackName string `json:"slack_name"`
}
// Notifier sends notifications
type Notifier interface {
SendDirect(ctx context.Context, userID string, msg *Notification) error
SendChannel(ctx context.Context, channel string, msg *Notification) error
}
// Processor processes webhook events and generates notifications
type Processor struct {
resolver identity.Resolver
slack slack.Client
cfg *config.Config
resolver IdentityResolver
notifier Notifier
logger zerolog.Logger
}
// NewProcessor creates a new event processor
func NewProcessor(resolver identity.Resolver, slackClient slack.Client, logger zerolog.Logger) *Processor {
func NewProcessor(cfg *config.Config, resolver IdentityResolver, notifier Notifier, logger zerolog.Logger) *Processor {
return &Processor{
cfg: cfg,
resolver: resolver,
slack: slackClient,
notifier: notifier,
logger: logger.With().Str("component", "processor").Logger(),
}
}
// ProcessPullRequest handles pull_request webhook events
func (p *Processor) ProcessPullRequest(ctx context.Context, event *webhook.PullRequestEvent) error {
logger := p.logger.With().
Int64("pr_number", event.Number).
Str("action", event.Action).
Str("repo", event.Repository.FullName).
Logger()
// HandlePullRequest processes pull request events
func (p *Processor) HandlePullRequest(e *webhook.PullRequestEvent) {
ctx := context.Background()
event := p.normalizePREvent(e)
logger.Info().Msg("processing pull request event")
p.logger.Debug().
Str("action", e.Action).
Int("assignees_count", len(event.Assignees)).
Int("reviewers_count", len(event.Reviewers)).
Str("actor", event.Actor.GiteaUsername).
Msg("processing PR event")
switch event.Action {
case "review_requested":
return p.handleReviewRequested(ctx, event, logger)
var usersToNotify []struct {
user User
reason NotificationReason
}
switch e.Action {
case "opened":
return p.handleOpened(ctx, event, logger)
event.Type = TypePROpened
// Notify assignees and reviewers
if p.cfg.Rules.PR.NotifyAssignees {
for _, u := range event.Assignees {
p.logger.Debug().Str("assignee", u.GiteaUsername).Msg("adding assignee to notify")
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{u, ReasonAssignee})
}
}
if p.cfg.Rules.PR.NotifyReviewers {
for _, u := range event.Reviewers {
p.logger.Debug().Str("reviewer", u.GiteaUsername).Msg("adding reviewer to notify")
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{u, ReasonReviewer})
}
}
case "closed":
return p.handleClosed(ctx, event, logger)
if e.PullRequest.Merged {
event.Type = TypePRMerged
} else {
event.Type = TypePRClosed
}
// Notify owner
if p.cfg.Rules.PR.NotifyOwner && event.Owner != nil {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{*event.Owner, ReasonOwner})
}
case "assigned":
event.Type = TypePRAssigned
// Notify newly assigned user
if e.Assignee != nil && p.cfg.Rules.PR.NotifyAssignees {
p.logger.Debug().Str("assignee", e.Assignee.Login).Msg("adding newly assigned user to notify")
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{giteaUserToUser(*e.Assignee), ReasonAssignee})
}
case "review_requested":
event.Type = TypePRReviewRequested
// Notify newly requested reviewers (array)
if len(e.RequestedReviewers) > 0 {
for _, reviewer := range e.RequestedReviewers {
p.logger.Debug().Str("requested_reviewer", reviewer.Login).Msg("review_requested event")
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{giteaUserToUser(reviewer), ReasonReviewer})
}
} else {
p.logger.Debug().Msg("review_requested event but no RequestedReviewers field")
}
case "synchronize":
event.Type = TypePRSynchronized
// Add commits info to event
event.Commits = make([]CommitInfo, 0, len(e.Commits))
for _, c := range e.Commits {
// Get first line of commit message
msg := c.Message
if idx := strings.Index(msg, "\n"); idx > 0 {
msg = msg[:idx]
}
event.Commits = append(event.Commits, CommitInfo{
SHA: c.ID[:7], // Short SHA
Message: msg,
URL: c.URL,
Author: c.Author.Name,
})
}
// Notify owner (will be filtered out if they pushed the commits)
if p.cfg.Rules.PR.NotifyOwner && event.Owner != nil {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{*event.Owner, ReasonOwner})
}
// Notify reviewers
if p.cfg.Rules.PR.NotifyReviewers {
for _, u := range event.Reviewers {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{u, ReasonReviewer})
}
}
// Notify assignees
if p.cfg.Rules.PR.NotifyAssignees {
for _, u := range event.Assignees {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{u, ReasonAssignee})
}
}
}
// Remove the actor from notifications (don't notify yourself)
beforeFilter := len(usersToNotify)
usersToNotify = p.filterOutActor(usersToNotify, event.Actor)
p.logger.Debug().
Int("before_filter", beforeFilter).
Int("after_filter", len(usersToNotify)).
Str("actor_filtered", event.Actor.GiteaUsername).
Msg("filtered out actor from notifications")
// Send notifications
p.sendNotifications(ctx, event, usersToNotify)
}
// HandlePullRequestReview processes PR review events
func (p *Processor) HandlePullRequestReview(e *webhook.PullRequestReviewEvent) {
ctx := context.Background()
event := p.normalizeReviewEvent(e)
var usersToNotify []struct {
user User
reason NotificationReason
}
// Notify PR owner about the review
if p.cfg.Rules.PR.NotifyOwner && event.Owner != nil {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{*event.Owner, ReasonOwner})
}
// Remove the actor
usersToNotify = p.filterOutActor(usersToNotify, event.Actor)
p.sendNotifications(ctx, event, usersToNotify)
}
// HandlePullRequestComment processes PR comment events
func (p *Processor) HandlePullRequestComment(e *webhook.PullRequestCommentEvent) {
ctx := context.Background()
event := p.normalizePRCommentEvent(e)
if e.Action != "created" {
return // Only notify on new comments
}
var usersToNotify []struct {
user User
reason NotificationReason
}
// Notify PR owner
if p.cfg.Rules.Comment.NotifyThreadOwner && event.Owner != nil {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{*event.Owner, ReasonOwner})
}
// Extract and notify mentioned users
if p.cfg.Rules.Comment.NotifyMentioned {
mentionedUsers := ExtractMentions(e.Comment.Body)
for _, username := range mentionedUsers {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{User{GiteaUsername: username}, ReasonMention})
}
}
// Notify reviewers
if p.cfg.Rules.Comment.NotifyReviewers {
for _, u := range event.Reviewers {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{u, ReasonReviewer})
}
}
// Remove the actor
usersToNotify = p.filterOutActor(usersToNotify, event.Actor)
p.sendNotifications(ctx, event, usersToNotify)
}
// HandleIssue processes issue events
func (p *Processor) HandleIssue(e *webhook.IssueEvent) {
ctx := context.Background()
event := p.normalizeIssueEvent(e)
var usersToNotify []struct {
user User
reason NotificationReason
}
switch e.Action {
case "opened":
event.Type = TypeIssueOpened
// Notify assignees
if p.cfg.Rules.Issue.NotifyAssignees {
for _, u := range event.Assignees {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{u, ReasonAssignee})
}
}
case "closed":
event.Type = TypeIssueClosed
// Notify owner
if event.Owner != nil {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{*event.Owner, ReasonOwner})
}
}
// Remove the actor
usersToNotify = p.filterOutActor(usersToNotify, event.Actor)
p.sendNotifications(ctx, event, usersToNotify)
}
// HandleIssueComment processes issue comment events
func (p *Processor) HandleIssueComment(e *webhook.IssueCommentEvent) {
ctx := context.Background()
event := p.normalizeIssueCommentEvent(e)
if e.Action != "created" {
return
}
var usersToNotify []struct {
user User
reason NotificationReason
}
// Notify issue owner
if p.cfg.Rules.Comment.NotifyThreadOwner && event.Owner != nil {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{*event.Owner, ReasonOwner})
}
// Extract and notify mentioned users
if p.cfg.Rules.Comment.NotifyMentioned {
mentionedUsers := ExtractMentions(e.Comment.Body)
for _, username := range mentionedUsers {
usersToNotify = append(usersToNotify, struct {
user User
reason NotificationReason
}{User{GiteaUsername: username}, ReasonMention})
}
}
// Remove the actor
usersToNotify = p.filterOutActor(usersToNotify, event.Actor)
p.sendNotifications(ctx, event, usersToNotify)
}
// sendNotifications resolves users and sends notifications
func (p *Processor) sendNotifications(ctx context.Context, event *Event, users []struct {
user User
reason NotificationReason
}) {
// Deduplicate users
seen := make(map[string]bool)
for _, u := range users {
key := u.user.GiteaUsername
if key == "" {
key = u.user.Email
}
if seen[key] {
continue
}
seen[key] = true
// Resolve user identity
identity, err := p.resolver.Resolve(ctx, u.user)
if err != nil {
p.logger.Warn().
Err(err).
Str("username", u.user.GiteaUsername).
Msg("failed to resolve user identity")
continue
}
if identity.SlackID == "" {
p.logger.Debug().
Str("username", u.user.GiteaUsername).
Msg("user has no Slack ID, skipping notification")
continue
}
// Create notification
notification := &Notification{
TargetUser: u.user,
Event: event,
Reason: u.reason,
Message: p.formatMessage(event, u.reason),
}
// Send DM
if err := p.notifier.SendDirect(ctx, identity.SlackID, notification); err != nil {
p.logger.Error().
Err(err).
Str("slack_id", identity.SlackID).
Msg("failed to send notification")
} else {
p.logger.Info().
Str("slack_id", identity.SlackID).
Str("reason", string(u.reason)).
Str("event_type", string(event.Type)).
Msg("notification sent")
}
}
}
// filterOutActor removes the event actor from the notification list
func (p *Processor) filterOutActor(users []struct {
user User
reason NotificationReason
}, actor User) []struct {
user User
reason NotificationReason
} {
var filtered []struct {
user User
reason NotificationReason
}
for _, u := range users {
if u.user.GiteaUsername != actor.GiteaUsername {
filtered = append(filtered, u)
}
}
return filtered
}
// formatMessage creates a human-readable notification message
func (p *Processor) formatMessage(event *Event, reason NotificationReason) string {
// This will be enhanced later with templates
switch event.Type {
case TypePROpened:
return "New PR opened"
case TypePRClosed:
return "PR closed"
case TypePRMerged:
return "PR merged"
case TypePRReviewRequested:
return "Review requested"
case TypePRReviewed:
return "PR reviewed"
case TypePRCommented:
return "New comment on PR"
case TypeIssueOpened:
return "New issue opened"
case TypeIssueClosed:
return "Issue closed"
case TypeIssueCommented:
return "New comment on issue"
default:
logger.Debug().Msg("ignoring pull_request action")
return nil
return "New notification"
}
}
// handleReviewRequested processes review_requested events
func (p *Processor) handleReviewRequested(ctx context.Context, event *webhook.PullRequestEvent, logger zerolog.Logger) error {
logger.Info().Msg("processing review_requested event")
// Normalization helpers
func (p *Processor) normalizePREvent(e *webhook.PullRequestEvent) *Event {
owner := giteaUserToUser(e.PullRequest.User)
var reviewer *webhook.GiteaUser
// Approach 1: event.RequestedReviewers array (standard)
if len(event.RequestedReviewers) > 0 {
reviewer = &event.RequestedReviewers[0]
logger.Info().
Int("count", len(event.RequestedReviewers)).
Str("reviewer", reviewer.Login).
Msg("found reviewer in RequestedReviewers array")
} else if event.RequestedReviewer != nil {
// Approach 2: singular field (edge case)
reviewer = event.RequestedReviewer
logger.Info().
Str("reviewer", reviewer.Login).
Msg("found reviewer in RequestedReviewer singular field")
} else if len(event.PullRequest.Assignees) > 0 {
// Approach 3: fall back to PR assignees
reviewer = &event.PullRequest.Assignees[0]
logger.Info().
Str("reviewer", reviewer.Login).
Msg("falling back to PR assignee as reviewer")
} else {
logger.Warn().Msg("no reviewer found in event payload")
return fmt.Errorf("no reviewer found for review_requested event (PR #%d)", event.Number)
// Get reviewers from both possible locations:
// - e.PullRequest.RequestedReviewers: present in opened action
// - e.RequestedReviewers: present in review_requested action
var reviewers []User
if len(e.PullRequest.RequestedReviewers) > 0 {
reviewers = giteaUsersToUsers(e.PullRequest.RequestedReviewers)
} else if len(e.RequestedReviewers) > 0 {
reviewers = giteaUsersToUsers(e.RequestedReviewers)
}
// Resolve the reviewer's Slack identity
resolved, err := p.resolver.Resolve(ctx, User{
GiteaUsername: reviewer.Login,
GiteaID: reviewer.ID,
Email: reviewer.Email,
FullName: reviewer.FullName,
})
if err != nil {
logger.Error().
Err(err).
Str("reviewer", reviewer.Login).
Msg("failed to resolve reviewer identity")
return fmt.Errorf("resolving reviewer %s: %w", reviewer.Login, err)
return &Event{
Timestamp: time.Now(),
Actor: giteaUserToUser(e.Sender),
RepoName: e.Repository.Name,
RepoFullName: e.Repository.FullName,
RepoURL: e.Repository.HTMLURL,
Number: e.PullRequest.Number,
Title: e.PullRequest.Title,
Body: e.PullRequest.Body,
URL: e.PullRequest.HTMLURL,
Owner: &owner,
Assignees: giteaUsersToUsers(e.PullRequest.Assignees),
Reviewers: reviewers,
}
if resolved.SlackID == "" {
logger.Warn().
Str("reviewer", reviewer.Login).
Msg("reviewer has no Slack ID, skipping DM")
return nil
}
// Send Slack DM
dmText := fmt.Sprintf(
":bell: *Review Requested*\n*Repository:* %s\n*Pull Request:* #%d - %s\n*Requested by:* %s\n\n<%s|View on Gitea>",
event.Repository.FullName,
event.Number,
event.PullRequest.Title,
event.Sender.Login,
event.PullRequest.HTMLURL,
)
if err := p.slack.SendDM(ctx, resolved.SlackID, dmText); err != nil {
logger.Error().
Err(err).
Str("slack_id", resolved.SlackID).
Msg("failed to send Slack DM")
return fmt.Errorf("sending Slack DM to %s: %w", resolved.SlackID, err)
}
logger.Info().
Str("reviewer", reviewer.Login).
Str("slack_id", resolved.SlackID).
Msg("successfully sent review_requested notification")
return nil
}
// handleOpened processes pull_request opened events
func (p *Processor) handleOpened(ctx context.Context, event *webhook.PullRequestEvent, logger zerolog.Logger) error {
logger.Info().Msg("processing PR opened event")
for _, assignee := range event.PullRequest.Assignees {
resolved, err := p.resolver.Resolve(ctx, User{
GiteaUsername: assignee.Login,
GiteaID: assignee.ID,
Email: assignee.Email,
FullName: assignee.FullName,
})
if err != nil {
logger.Error().Err(err).Str("assignee", assignee.Login).Msg("failed to resolve assignee")
continue
func (p *Processor) normalizeReviewEvent(e *webhook.PullRequestReviewEvent) *Event {
owner := giteaUserToUser(e.PullRequest.User)
return &Event{
Type: TypePRReviewed,
Timestamp: time.Now(),
Actor: giteaUserToUser(e.Sender),
RepoName: e.Repository.Name,
RepoFullName: e.Repository.FullName,
RepoURL: e.Repository.HTMLURL,
Number: e.PullRequest.Number,
Title: e.PullRequest.Title,
URL: e.PullRequest.HTMLURL,
Owner: &owner,
ReviewState: e.Review.State,
CommentBody: e.Review.Body,
}
if resolved.SlackID == "" {
continue
}
text := fmt.Sprintf(
":new: *PR Opened*\n*Repository:* %s\n*PR:* #%d - %s\n*By:* %s\n\n<%s|View on Gitea>",
event.Repository.FullName,
event.Number,
event.PullRequest.Title,
event.Sender.Login,
event.PullRequest.HTMLURL,
)
if err := p.slack.SendDM(ctx, resolved.SlackID, text); err != nil {
logger.Error().Err(err).Str("slack_id", resolved.SlackID).Msg("failed to send DM")
}
}
return nil
}
// handleClosed processes pull_request closed/merged events
func (p *Processor) handleClosed(ctx context.Context, event *webhook.PullRequestEvent, logger zerolog.Logger) error {
if event.PullRequest.Merged {
logger.Info().Msg("processing PR merged event")
} else {
logger.Info().Msg("processing PR closed (without merge) event")
func (p *Processor) normalizePRCommentEvent(e *webhook.PullRequestCommentEvent) *Event {
owner := giteaUserToUser(e.PullRequest.User)
return &Event{
Type: TypePRCommented,
Timestamp: time.Now(),
Actor: giteaUserToUser(e.Sender),
RepoName: e.Repository.Name,
RepoFullName: e.Repository.FullName,
RepoURL: e.Repository.HTMLURL,
Number: e.PullRequest.Number,
Title: e.PullRequest.Title,
URL: e.PullRequest.HTMLURL,
Owner: &owner,
Reviewers: giteaUsersToUsers(e.PullRequest.RequestedReviewers),
CommentBody: e.Comment.Body,
CommentURL: e.Comment.HTMLURL,
}
verb := "Closed"
if event.PullRequest.Merged {
verb = "Merged"
}
for _, assignee := range event.PullRequest.Assignees {
resolved, err := p.resolver.Resolve(ctx, User{
GiteaUsername: assignee.Login,
GiteaID: assignee.ID,
Email: assignee.Email,
FullName: assignee.FullName,
})
if err != nil {
logger.Error().Err(err).Str("assignee", assignee.Login).Msg("failed to resolve assignee")
continue
}
if resolved.SlackID == "" {
continue
}
var emoji string
if event.PullRequest.Merged {
emoji = ":white_check_mark:"
} else {
emoji = ":x:"
}
text := fmt.Sprintf(
"%s *PR %s*\n*Repository:* %s\n*PR:* #%d - %s\n\n<%s|View on Gitea>",
emoji,
verb,
event.Repository.FullName,
event.Number,
event.PullRequest.Title,
event.PullRequest.HTMLURL,
)
if err := p.slack.SendDM(ctx, resolved.SlackID, text); err != nil {
logger.Error().Err(err).Str("slack_id", resolved.SlackID).Msg("failed to send DM")
}
}
return nil
}
func (p *Processor) normalizeIssueEvent(e *webhook.IssueEvent) *Event {
owner := giteaUserToUser(e.Issue.User)
return &Event{
Timestamp: time.Now(),
Actor: giteaUserToUser(e.Sender),
RepoName: e.Repository.Name,
RepoFullName: e.Repository.FullName,
RepoURL: e.Repository.HTMLURL,
Number: e.Issue.Number,
Title: e.Issue.Title,
Body: e.Issue.Body,
URL: e.Issue.HTMLURL,
Owner: &owner,
Assignees: giteaUsersToUsers(e.Issue.Assignees),
}
}
func (p *Processor) normalizeIssueCommentEvent(e *webhook.IssueCommentEvent) *Event {
owner := giteaUserToUser(e.Issue.User)
return &Event{
Type: TypeIssueCommented,
Timestamp: time.Now(),
Actor: giteaUserToUser(e.Sender),
RepoName: e.Repository.Name,
RepoFullName: e.Repository.FullName,
RepoURL: e.Repository.HTMLURL,
Number: e.Issue.Number,
Title: e.Issue.Title,
URL: e.Issue.HTMLURL,
Owner: &owner,
CommentBody: e.Comment.Body,
CommentURL: e.Comment.HTMLURL,
}
}
func giteaUserToUser(u webhook.GiteaUser) User {
username := u.Login
if username == "" {
username = u.Username
}
return User{
GiteaID: u.ID,
GiteaUsername: username,
Email: u.Email,
FullName: u.FullName,
}
}
func giteaUsersToUsers(users []webhook.GiteaUser) []User {
result := make([]User, len(users))
for i, u := range users {
result[i] = giteaUserToUser(u)
}
return result
}