package mail import ( "context" "crypto/tls" "fmt" "io" "log/slog" "net" "net/url" "strings" "time" gomessage "github.com/emersion/go-message/mail" "github.com/emersion/go-imap/v2" "github.com/emersion/go-imap/v2/imapclient" "github.com/google/uuid" "golang.org/x/net/proxy" "gorm.io/gorm" "thuanle.me/claw-email-bridge/internal/config" "thuanle.me/claw-email-bridge/internal/database" "thuanle.me/claw-email-bridge/internal/logging" ) // IMAPWatcher monitors an IMAP mailbox via IDLE and processes new emails. type IMAPWatcher struct { cfg *config.Config db *gorm.DB client *imapclient.Client dialIMAP func(addr string, options *imapclient.Options) (*imapclient.Client, error) dialIMAPViaProxy func(addr, proxyURL string, options *imapclient.Options) (*imapclient.Client, error) // onReceived is called after a task is saved as RECEIVED. // This will be wired to the OpenClaw dispatch in a future step. OnReceived func(task *database.Task) } const imapDialTimeout = 30 * time.Second type timeoutContextDialer struct { timeout time.Duration } func (d timeoutContextDialer) Dial(network, address string) (net.Conn, error) { ctx, cancel := context.WithTimeout(context.Background(), d.timeout) defer cancel() return d.DialContext(ctx, network, address) } func (d timeoutContextDialer) DialContext(ctx context.Context, network, address string) (net.Conn, error) { return (&net.Dialer{Timeout: d.timeout}).DialContext(ctx, network, address) } var socks5DialerFactory = proxy.SOCKS5 // NewIMAPWatcher creates a new IMAP watcher. func NewIMAPWatcher(cfg *config.Config, db *gorm.DB) *IMAPWatcher { return &IMAPWatcher{ cfg: cfg, db: db, dialIMAP: imapclient.DialTLS, dialIMAPViaProxy: dialTLSViaSOCKS5, } } // Run connects to IMAP, selects INBOX, and enters the IDLE loop. // It blocks until ctx is cancelled. Reconnects automatically on error. func (w *IMAPWatcher) Run(ctx context.Context) { for { if err := w.connectAndWatch(ctx); err != nil { slog.Error("imap watcher error, reconnecting in 10s", "error", err) } select { case <-ctx.Done(): slog.Info("imap watcher stopped") return case <-time.After(10 * time.Second): // Reconnect after delay. } } } // connectAndWatch handles a single IMAP session lifecycle. func (w *IMAPWatcher) connectAndWatch(ctx context.Context) error { numMessages := uint32(0) options := &imapclient.Options{ UnilateralDataHandler: &imapclient.UnilateralDataHandler{ Mailbox: func(data *imapclient.UnilateralDataMailbox) { if data.NumMessages != nil && *data.NumMessages > numMessages { slog.Info("imap: new message notification", "num_messages", *data.NumMessages) numMessages = *data.NumMessages } }, }, } addr := fmt.Sprintf("%s:%s", w.cfg.IMAPHost, w.cfg.IMAPPort) slog.Info("imap: connecting", "addr", addr) var c *imapclient.Client var err error if strings.TrimSpace(w.cfg.IMAPProxyURL) == "" { c, err = w.dialIMAP(addr, options) } else { c, err = w.dialIMAPViaProxy(addr, w.cfg.IMAPProxyURL, options) } if err != nil { return fmt.Errorf("dial IMAP: %w", err) } w.client = c defer func() { _ = c.Close() w.client = nil }() if err := c.Login(w.cfg.IMAPUser, w.cfg.IMAPPass).Wait(); err != nil { return fmt.Errorf("login: %w", err) } slog.Info("imap: logged in", "user", w.cfg.IMAPUser) selectedMbox, err := c.Select("INBOX", nil).Wait() if err != nil { return fmt.Errorf("select INBOX: %w", err) } numMessages = selectedMbox.NumMessages slog.Info("imap: INBOX selected", "messages", numMessages) // Fetch any unseen messages on startup. if err := w.fetchUnseen(c); err != nil { slog.Error("imap: initial fetch unseen failed", "error", err) } // IDLE loop: wait for new messages, then fetch unseen. for { select { case <-ctx.Done(): return nil default: } idleCmd, err := c.Idle() if err != nil { return fmt.Errorf("idle: %w", err) } done := make(chan error, 1) go func() { done <- idleCmd.Wait() }() // Re-check every 28 minutes (RFC recommends <29 min IDLE timeout). idleTimer := time.NewTimer(28 * time.Minute) select { case <-ctx.Done(): _ = idleCmd.Close() idleTimer.Stop() return nil case err := <-done: idleTimer.Stop() if err != nil { return fmt.Errorf("idle wait: %w", err) } case <-idleTimer.C: // Stop IDLE and restart it. if err := idleCmd.Close(); err != nil { return fmt.Errorf("idle close: %w", err) } <-done } // After exiting IDLE, fetch unseen messages. if err := w.fetchUnseen(c); err != nil { slog.Error("imap: fetch unseen failed", "error", err) } } } func dialTLSViaSOCKS5(addr, proxyURL string, options *imapclient.Options) (*imapclient.Client, error) { u, err := url.Parse(proxyURL) if err != nil { return nil, fmt.Errorf("parse proxy url: %w", err) } if u.Scheme != "socks5" { return nil, fmt.Errorf("unsupported proxy scheme: %s", u.Scheme) } var auth *proxy.Auth if u.User != nil { pw, _ := u.User.Password() auth = &proxy.Auth{User: u.User.Username(), Password: pw} } dialer, err := socks5DialerFactory("tcp", u.Host, auth, timeoutContextDialer{timeout: imapDialTimeout}) if err != nil { return nil, fmt.Errorf("create socks5 dialer: %w", err) } dialCtx, dialCancel := context.WithTimeout(context.Background(), imapDialTimeout) defer dialCancel() var conn net.Conn if cd, ok := dialer.(proxy.ContextDialer); ok { conn, err = cd.DialContext(dialCtx, "tcp", addr) } else { conn, err = dialer.Dial("tcp", addr) } if err != nil { return nil, fmt.Errorf("proxy dial: %w", err) } host, _, splitErr := net.SplitHostPort(addr) if splitErr != nil || host == "" { host = addr } tlsConn := tls.Client(conn, &tls.Config{ServerName: host}) _ = tlsConn.SetDeadline(time.Now().Add(imapDialTimeout)) handshakeCtx, handshakeCancel := context.WithTimeout(context.Background(), imapDialTimeout) defer handshakeCancel() if err := tlsConn.HandshakeContext(handshakeCtx); err != nil { _ = tlsConn.Close() return nil, fmt.Errorf("tls handshake: %w", err) } _ = tlsConn.SetDeadline(time.Time{}) return imapclient.New(tlsConn, options), nil } // fetchUnseen searches for UNSEEN messages and processes each one. func (w *IMAPWatcher) fetchUnseen(c *imapclient.Client) error { criteria := &imap.SearchCriteria{ NotFlag: []imap.Flag{imap.FlagSeen}, } searchData, err := c.UIDSearch(criteria, nil).Wait() if err != nil { return fmt.Errorf("uid search: %w", err) } uids := searchData.AllUIDs() if len(uids) == 0 { return nil } slog.Info("imap: unseen messages found", "count", len(uids)) uidSet := imap.UIDSet{} for _, uid := range uids { uidSet.AddNum(uid) } bodySection := &imap.FetchItemBodySection{} fetchOptions := &imap.FetchOptions{ Envelope: true, UID: true, BodySection: []*imap.FetchItemBodySection{bodySection}, } fetchCmd := c.Fetch(uidSet, fetchOptions) defer fetchCmd.Close() for { msg := fetchCmd.Next() if msg == nil { break } buf, err := msg.Collect() if err != nil { slog.Error("imap: collect message failed", "error", err) continue } w.processMessage(c, buf, bodySection) } if err := fetchCmd.Close(); err != nil { return fmt.Errorf("fetch close: %w", err) } return nil } // resolveThreadID determines the thread_id for a message. // It tries In-Reply-To first, then falls back to References (last-to-first). // Returns the message's own ID if no parent is found. func resolveThreadID(messageID, inReplyTo, references string, lookup func(string) (*database.Task, error)) string { if inReplyTo != "" { if parent, err := lookup(inReplyTo); err == nil && parent != nil { return parent.ThreadID } } for _, refID := range parseReferences(references) { if parent, err := lookup(refID); err == nil && parent != nil { return parent.ThreadID } } return messageID } // parseReferences splits a References header value into individual message IDs, // returned in reverse order (most recent first). func parseReferences(header string) []string { if header == "" { return nil } raw := strings.Fields(header) for i := range raw { raw[i] = strings.TrimSpace(raw[i]) } // Reverse: last (most recent) first. for i, j := 0, len(raw)-1; i < j; i, j = i+1, j-1 { raw[i], raw[j] = raw[j], raw[i] } return raw } // processMessage handles a single fetched email message. func (w *IMAPWatcher) processMessage(c *imapclient.Client, buf *imapclient.FetchMessageBuffer, bodySection *imap.FetchItemBodySection) { env := buf.Envelope if env == nil { slog.Warn("imap: message has no envelope, skipping") return } messageID := env.MessageID subject := env.Subject // Extract sender email. sender := "" if len(env.From) > 0 { sender = env.From[0].Addr() } log := logging.TaskLogger("", "", messageID) // 1. Anti-loop: skip if sender is the system email. if strings.EqualFold(sender, w.cfg.SystemEmail) { log.Info("imap: anti-loop, ignoring own email", "sender", sender) w.saveIgnored(messageID, sender, subject, "anti-loop: system email") w.markSeen(c, buf.UID) return } // 2. Whitelist: skip if sender not in allowed list. if !w.isWhitelisted(sender) { log.Info("imap: sender not whitelisted, ignoring", "sender", sender) w.saveIgnored(messageID, sender, subject, "not whitelisted") w.markSeen(c, buf.UID) return } // 3. Idempotency: skip if message_id already exists. existing, err := database.FindByMessageID(w.db, messageID) if err != nil { log.Error("imap: db lookup failed", "error", err) return } if existing != nil { log.Info("imap: duplicate message_id, skipping") w.markSeen(c, buf.UID) return } // 4. Threading: determine thread_id. threadID := messageID inReplyTo := "" // Parse body to get In-Reply-To / References headers and plain text. bodyPlain := "" rawBody := buf.FindBodySection(bodySection) if rawBody != nil { mr, err := gomessage.CreateReader(io.NopCloser(strings.NewReader(string(rawBody)))) if err == nil { // Extract In-Reply-To header. if irt, err := mr.Header.Text("In-Reply-To"); err == nil && irt != "" { inReplyTo = strings.TrimSpace(irt) } // Read body parts for plain text. for { p, err := mr.NextPart() if err != nil { break } if _, ok := p.Header.(*gomessage.InlineHeader); ok { b, _ := io.ReadAll(p.Body) if bodyPlain == "" { bodyPlain = string(b) } } } } } if inReplyTo != "" { parent, err := database.FindByMessageID(w.db, inReplyTo) if err == nil && parent != nil { threadID = parent.ThreadID } } // 5. Save task with status RECEIVED. taskUUID := uuid.New().String() task := &database.Task{ TaskUUID: taskUUID, ThreadID: threadID, MessageID: messageID, Sender: sender, Subject: subject, BodyPlain: bodyPlain, Status: database.StatusReceived, } taskLog := logging.TaskLogger(taskUUID, threadID, messageID) if err := database.CreateTask(w.db, task); err != nil { taskLog.Error("imap: failed to save task", "error", err) return } taskLog.Info("imap: task created", "sender", sender, "subject", subject, "status", database.StatusReceived, ) // Mark as seen in IMAP so we don't re-process. w.markSeen(c, buf.UID) // Trigger dispatch (will be wired later). if w.OnReceived != nil { w.OnReceived(task) } } // saveIgnored records an ignored email in the DB for auditability. func (w *IMAPWatcher) saveIgnored(messageID, sender, subject, reason string) { taskUUID := uuid.New().String() errMsg := reason task := &database.Task{ TaskUUID: taskUUID, ThreadID: messageID, MessageID: messageID, Sender: sender, Subject: subject, Status: database.StatusIgnored, LastError: &errMsg, } if err := database.CreateTask(w.db, task); err != nil { slog.Error("imap: failed to save ignored task", "error", err, "message_id", messageID) } } // isWhitelisted checks if the sender is in the allowed list. // Empty whitelist means all senders are allowed. func (w *IMAPWatcher) isWhitelisted(sender string) bool { if len(w.cfg.WhitelistEmails) == 0 { return true } lower := strings.ToLower(sender) for _, email := range w.cfg.WhitelistEmails { if strings.ToLower(email) == lower { return true } } return false } // markSeen flags a message as \Seen in IMAP. func (w *IMAPWatcher) markSeen(c *imapclient.Client, uid imap.UID) { storeFlags := imap.StoreFlags{ Op: imap.StoreFlagsAdd, Flags: []imap.Flag{imap.FlagSeen}, Silent: true, } if err := c.Store(imap.UIDSetNum(uid), &storeFlags, nil).Close(); err != nil { slog.Error("imap: failed to mark message as seen", "error", err, "uid", uid) } }