Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,7 @@ type Config struct {
Password string `password:"true" info:"Password for authentication."`
}


Webhook struct {
Enable bool `public:"true" info:"Enables webhook as a contact method."`
AllowedURLs []string `public:"true" info:"If set, allows webhooks for these domains only."`
Expand Down
6 changes: 6 additions & 0 deletions engine/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
"github.com/target/goalert/engine/compatmanager"
"github.com/target/goalert/engine/escalationmanager"
"github.com/target/goalert/engine/heartbeatmanager"
"github.com/target/goalert/engine/imapmanager"
"github.com/target/goalert/engine/message"
"github.com/target/goalert/engine/metricsmanager"
"github.com/target/goalert/engine/npcyclemanager"
Expand Down Expand Up @@ -135,6 +136,10 @@ func NewEngine(ctx context.Context, db *sql.DB, c *Config) (*Engine, error) {
if err != nil {
return nil, errors.Wrap(err, "compatibility backend")
}
imapMgr, err := imapmanager.NewDB(ctx, db, c.AlertStore, c.ConfigSource, c.Logger.With(slog.String("module", "imap")))
if err != nil {
return nil, errors.Wrap(err, "imap manager backend")
}

p.modules = []processinglock.Module{
compatMgr,
Expand All @@ -147,6 +152,7 @@ func NewEngine(ctx context.Context, db *sql.DB, c *Config) (*Engine, error) {
hbMgr,
cleanMgr,
metricsMgr,
imapMgr,
}

if expflag.ContextHas(ctx, expflag.UnivKeys) {
Expand Down
53 changes: 53 additions & 0 deletions engine/imapmanager/cleanup.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
package imapmanager

import (
"context"
"database/sql"
"fmt"

"github.com/riverqueue/river"
"github.com/target/goalert/gadb"
"github.com/target/goalert/util"
)

// CleanupArgs are the arguments for the cleanup job.
type CleanupArgs struct{}

func (CleanupArgs) Kind() string { return "imap-cleanup" }

// CleanupProcessedMessages removes old processed message records.
func (db *DB) CleanupProcessedMessages(ctx context.Context, job *river.Job[CleanupArgs]) error {
db.logger.Info("IMAP cleanup: starting")

// Run cleanup in batches until no more work
var totalDeleted int64
for {
var deleted int64
err := db.lock.WithTxShared(ctx, func(ctx context.Context, tx *sql.Tx) error {
queries := gadb.New(tx)
rows, err := queries.IMAPCleanupProcessedMessages(ctx)
deleted = rows
return err
})

if err != nil {
return fmt.Errorf("cleanup processed messages: %w", err)
}

totalDeleted += deleted

if deleted == 0 {
// No more work
break
}

// Sleep briefly between batches to avoid overwhelming the database
err = util.ContextSleep(ctx, 100)
if err != nil {
return err
}
}

db.logger.Info("IMAP cleanup: completed", "deleted", totalDeleted)
return nil
}
44 changes: 44 additions & 0 deletions engine/imapmanager/db.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
package imapmanager

import (
"context"
"database/sql"
"log/slog"

"github.com/target/goalert/alert"
"github.com/target/goalert/config"
"github.com/target/goalert/engine/processinglock"
)

// DB handles IMAP email polling and alert creation.
type DB struct {
db *sql.DB
lock *processinglock.Lock

alertStore *alert.Store
cfg config.Source

logger *slog.Logger
}

// Name returns the name of the module.
func (db *DB) Name() string { return "Engine.IMAPManager" }

// NewDB creates a new DB.
func NewDB(ctx context.Context, db *sql.DB, alertStore *alert.Store, cfg config.Source, log *slog.Logger) (*DB, error) {
lock, err := processinglock.NewLock(ctx, db, processinglock.Config{
Version: 1,
Type: processinglock.TypeIMAP,
})
if err != nil {
return nil, err
}

return &DB{
db: db,
lock: lock,
logger: log,
alertStore: alertStore,
cfg: cfg,
}, nil
}
235 changes: 235 additions & 0 deletions engine/imapmanager/filter.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,235 @@
package imapmanager

import (
"bytes"
"io"
"mime"
"mime/multipart"
"net/mail"
"regexp"
"strings"

"github.com/emersion/go-imap"
"github.com/target/goalert/gadb"
)

// emailMessage represents a parsed email message.
type emailMessage struct {
MessageID string
From string
To string
Subject string
Body string
Date string // Email date
InReplyTo string // Message ID this email is replying to
}

// parseEnvelope extracts email fields from an IMAP message.
func parseEnvelope(msg *imap.Message) *emailMessage {
if msg.Envelope == nil {
return nil
}

em := &emailMessage{
MessageID: msg.Envelope.MessageId,
Subject: msg.Envelope.Subject,
Date: msg.Envelope.Date.String(),
InReplyTo: msg.Envelope.InReplyTo,
}

// Extract From address
if len(msg.Envelope.From) > 0 {
addr := msg.Envelope.From[0]
if addr.MailboxName != "" && addr.HostName != "" {
em.From = addr.MailboxName + "@" + addr.HostName
}
}

// Extract To address(es)
var toAddrs []string
for _, addr := range msg.Envelope.To {
if addr.MailboxName != "" && addr.HostName != "" {
toAddrs = append(toAddrs, addr.MailboxName+"@"+addr.HostName)
}
}
em.To = strings.Join(toAddrs, ", ")

// Extract body (RFC822) - parse to get actual content without technical headers
for _, bodyItem := range msg.Body {
if bodyItem != nil {
bodyBytes, err := io.ReadAll(bodyItem)
if err == nil {
// Parse the RFC822 message to extract just the body content
em.Body = extractEmailBody(bodyBytes)
}
}
}

if em.MessageID == "" {
return nil
}

return em
}

// extractEmailBody parses an RFC822 message and extracts just the text body content,
// stripping out all MIME headers (Delivered-To, Received, ARC-Seal, etc.).
func extractEmailBody(rfc822Data []byte) string {
// Parse the RFC822 message
msg, err := mail.ReadMessage(bytes.NewReader(rfc822Data))
if err != nil {
// If parsing fails, return empty string
return ""
}

// Get the Content-Type header to check if it's multipart
contentType := msg.Header.Get("Content-Type")
if contentType == "" {
// No Content-Type, try to read body directly
bodyBytes, err := io.ReadAll(msg.Body)
if err != nil {
return ""
}
return string(bodyBytes)
}

// Parse the media type
mediaType, params, err := mime.ParseMediaType(contentType)
if err != nil {
// If parsing fails, try to read body directly
bodyBytes, err := io.ReadAll(msg.Body)
if err != nil {
return ""
}
return string(bodyBytes)
}

// Handle multipart messages
if strings.HasPrefix(mediaType, "multipart/") {
boundary, ok := params["boundary"]
if !ok {
return ""
}

mr := multipart.NewReader(msg.Body, boundary)
var textBody, htmlBody string

// Read all parts
for {
part, err := mr.NextPart()
if err == io.EOF {
break
}
if err != nil {
break
}

partContentType := part.Header.Get("Content-Type")
partMediaType, _, _ := mime.ParseMediaType(partContentType)

partBytes, err := io.ReadAll(part)
if err != nil {
continue
}

// Prefer text/plain, fallback to text/html
if partMediaType == "text/plain" {
textBody = string(partBytes)
} else if partMediaType == "text/html" && textBody == "" {
htmlBody = string(partBytes)
}
}

// Return text/plain if available, otherwise text/html
if textBody != "" {
return textBody
}
return htmlBody
}

// Not multipart, read body directly
bodyBytes, err := io.ReadAll(msg.Body)
if err != nil {
return ""
}
return string(bodyBytes)
}

// isReply checks if an email is a reply or forward based on subject and headers.
func isReply(email *emailMessage) bool {
// Check subject line for common reply/forward prefixes
subject := strings.TrimSpace(email.Subject)
if len(subject) >= 3 {
prefix := strings.ToLower(subject[:3])
if prefix == "re:" || prefix == "fw:" {
return true
}
}
if len(subject) >= 4 {
prefix := strings.ToLower(subject[:4])
if prefix == "fwd:" {
return true
}
}

// Check In-Reply-To header
if email.InReplyTo != "" {
return true
}

return false
}

// matchesFilter checks if an email matches a filter rule.
func matchesFilter(email *emailMessage, rule gadb.IMAPFilterRulesForServiceRow) (bool, error) {
// Check exclude_replies setting first
if rule.ExcludeReplies && isReply(email) {
return false, nil
}

// All criteria in a rule must match (AND logic)
if rule.FromPattern.Valid {
matched, err := matchPattern(email.From, rule.FromPattern.String, rule.MatchMode)
if err != nil || !matched {
return false, err
}
}

if rule.SubjectPattern.Valid {
matched, err := matchPattern(email.Subject, rule.SubjectPattern.String, rule.MatchMode)
if err != nil || !matched {
return false, err
}
}

if rule.ToPattern.Valid {
matched, err := matchPattern(email.To, rule.ToPattern.String, rule.MatchMode)
if err != nil || !matched {
return false, err
}
}

return true, nil
}

// matchPattern performs pattern matching based on the match mode.
func matchPattern(value, pattern, mode string) (bool, error) {
switch mode {
case "exact":
return strings.EqualFold(value, pattern), nil

case "contains":
return strings.Contains(strings.ToLower(value), strings.ToLower(pattern)), nil

case "regex":
re, err := regexp.Compile(pattern)
if err != nil {
return false, err
}
return re.MatchString(value), nil

default:
return false, nil
}
}

Loading
Loading