diff --git a/go.mod b/go.mod index 81a245e863..d57708d43f 100644 --- a/go.mod +++ b/go.mod @@ -24,6 +24,8 @@ require ( github.com/cloudwego/eino v0.9.14 github.com/denisenkom/go-mssqldb v0.12.3 github.com/elastic/go-elasticsearch/v8 v8.19.1 + github.com/emersion/go-imap/v2 v2.0.0-beta.8 + github.com/emersion/go-message v0.18.2 github.com/eric642/e2b-go-sdk v0.1.3 github.com/gin-gonic/gin v1.12.0 github.com/glebarez/sqlite v1.11.0 @@ -78,6 +80,8 @@ require ( gorm.io/gorm v1.25.7 ) +require github.com/emersion/go-sasl v0.0.0-20241020182733-b788ff22d5a6 // indirect + require ( buf.build/gen/go/bufbuild/protovalidate/protocolbuffers/go v1.36.11-20260415201107-50325440f8f2.1 // indirect cel.dev/expr v0.25.1 // indirect diff --git a/go.sum b/go.sum index 10b93787e2..9dd5340ba8 100644 --- a/go.sum +++ b/go.sum @@ -193,6 +193,12 @@ github.com/elastic/elastic-transport-go/v8 v8.8.0 h1:7k1Ua+qluFr6p1jfJjGDl97ssJS github.com/elastic/elastic-transport-go/v8 v8.8.0/go.mod h1:YLHer5cj0csTzNFXoNQ8qhtGY1GTvSqPnKWKaqQE3Hk= github.com/elastic/go-elasticsearch/v8 v8.19.1 h1:0iEGt5/Ds9MNVxEp3hqLsXdbe6SjleaVHONg/FuR09Q= github.com/elastic/go-elasticsearch/v8 v8.19.1/go.mod h1:tHJQdInFa6abmDbDCEH2LJja07l/SIpaGpJcm13nt7s= +github.com/emersion/go-imap/v2 v2.0.0-beta.8 h1:5IXZK1E33DyeP526320J3RS7eFlCYGFgtbrfapqDPug= +github.com/emersion/go-imap/v2 v2.0.0-beta.8/go.mod h1:dhoFe2Q0PwLrMD7oZw8ODuaD0vLYPe5uj2wcOMnvh48= +github.com/emersion/go-message v0.18.2 h1:rl55SQdjd9oJcIoQNhubD2Acs1E6IzlZISRTK7x/Lpg= +github.com/emersion/go-message v0.18.2/go.mod h1:XpJyL70LwRvq2a8rVbHXikPgKj8+aI0kGdHlg16ibYA= +github.com/emersion/go-sasl v0.0.0-20241020182733-b788ff22d5a6 h1:oP4q0fw+fOSWn3DfFi4EXdT+B+gTtzx8GC9xsc26Znk= +github.com/emersion/go-sasl v0.0.0-20241020182733-b788ff22d5a6/go.mod h1:iL2twTeMvZnrg54ZoPDNfJaJaqy0xIQFuBdrLsmspwQ= github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= github.com/envoyproxy/go-control-plane v0.9.4/go.mod h1:6rpuAdCZL397s3pYoYcLgu1mIlRU8Am5FuJP05cCM98= github.com/envoyproxy/go-control-plane v0.14.0 h1:hbG2kr4RuFj222B6+7T83thSPqLjwBIfQawTkC++2HA= diff --git a/internal/syncer/README.md b/internal/syncer/README.md index 0325fb138b..dd1a9026cc 100644 --- a/internal/syncer/README.md +++ b/internal/syncer/README.md @@ -289,7 +289,7 @@ Supplementary details (beyond the diagram): - `registry.go`: registry mapping source names to connector factories. - `builtin.go`: registration entry for connectors built into the current binary. - `fingerprint.go`: stable fingerprint and file name normalization utilities. -- `.go`: per data-source implementations, e.g. `rss.go`, `github.go`, `gmail.go`, `google_drive.go`, `outlook.go`, `rest_api.go`, `mysql.go`, `postgresql.go`, `discord.go`. +- `.go`: per data-source implementations, e.g. `rss.go`, `github.go`, `gmail.go`, `imap.go`, `google_drive.go`, `outlook.go`, `rest_api.go`, `mysql.go`, `postgresql.go`, `discord.go`. - `_test.go`: unit tests for each data source. - `mock/mock.go`: mock connector for syncer testing. @@ -568,4 +568,3 @@ Tests that need real MySQL, MinIO, Elasticsearch, Infinity, LLMs, or external Sa > + [feat[Go]: monitoring NATs and refactoring concurrency logic - #18049](https://github.com/infiniflow/ragflow/pull/18049) > + [feat[Go]: resuming transmission from the point of interruption during data source synchronisation - #18176](https://github.com/infiniflow/ragflow/pull/18176) - diff --git a/internal/syncer/connector/builtin.go b/internal/syncer/connector/builtin.go index 9150ee4925..84f4066482 100644 --- a/internal/syncer/connector/builtin.go +++ b/internal/syncer/connector/builtin.go @@ -36,6 +36,7 @@ func RegisterBuiltIns(registry *Registry) { registerBuiltIn(registry, "google-drive", NewGoogleDriveConnector) registerBuiltIn(registry, "google_drive", NewGoogleDriveConnector) registerBuiltIn(registry, "google_cloud_storage", NewGoogleCloudStorageConnector) + registerBuiltIn(registry, "imap", NewIMAPConnector) registerBuiltIn(registry, "outlook", NewOutlookConnector) registerBuiltIn(registry, "notion", NewNotionConnector) registerBuiltIn(registry, "rest_api", NewRestAPIConnector) diff --git a/internal/syncer/connector/imap.go b/internal/syncer/connector/imap.go new file mode 100644 index 0000000000..6402f97889 --- /dev/null +++ b/internal/syncer/connector/imap.go @@ -0,0 +1,851 @@ +// +// Copyright 2026 The InfiniFlow Authors. All Rights Reserved. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// + +package connector + +import ( + "bytes" + "context" + "crypto/sha256" + "crypto/tls" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "log" + "net" + netmail "net/mail" + "os" + "strconv" + "strings" + "time" + "unicode/utf8" + + "github.com/emersion/go-imap/v2" + "github.com/emersion/go-imap/v2/imapclient" + "github.com/emersion/go-message" + // Register the charset decoder so text/* parts in non-UTF-8 charsets + // (e.g. ISO-8859-1, Windows-1252) are decoded to UTF-8 on read. + _ "github.com/emersion/go-message/charset" + xhtml "golang.org/x/net/html" + + "ragflow/internal/utility" +) + +const ( + defaultIMAPBatchSize = 32 + defaultIMAPPort = 993 + defaultIMAPSizeThreshold = 10 * 1024 * 1024 + imapDialTimeout = 30 * time.Second + imapCommandTimeout = 30 * time.Second +) + +// IMAPConnector reads email messages and attachments from an IMAP server. +type IMAPConnector struct { + host string + port int + mailboxes []string + username string + password string + batchSize int + sizeThreshold int64 + + dial func(ctx context.Context, host string, port int, username, password string) (imapClient, error) +} + +// imapClient is the IMAP surface used by the connector. It is injected so unit +// tests can exercise the full connector without a live server. +type imapClient interface { + List(ctx context.Context) ([]string, error) + SelectMailbox(ctx context.Context, mailbox string) error + Search(ctx context.Context, since, before time.Time) ([]uint32, error) + Fetch(ctx context.Context, seqNum uint32) ([]byte, error) + Close() error +} + +// NewIMAPConnector creates an IMAP connector from the given config. +func NewIMAPConnector(config map[string]any) (*IMAPConnector, error) { + credentials, _ := config["credentials"].(map[string]any) + threshold := int64(defaultIMAPSizeThreshold) + if rawThreshold, err := strconv.ParseInt(os.Getenv("IMAP_CONNECTOR_SIZE_THRESHOLD"), 10, 64); err == nil && rawThreshold > 0 { + threshold = rawThreshold + } + return &IMAPConnector{ + host: strings.TrimSpace(stringConfig(config["imap_host"])), + port: configInt(config["imap_port"], defaultIMAPPort), + mailboxes: imapMailboxList(config["imap_mailbox"]), + username: strings.TrimSpace(stringConfig(credentials["imap_username"])), + password: stringConfig(credentials["imap_password"]), + batchSize: configInt(firstNonEmpty(stringConfig(config["sync_batch_size"]), stringConfig(config["batch_size"])), defaultIMAPBatchSize), + sizeThreshold: threshold, + dial: dialRealIMAPClient, + }, nil +} + +// Validate validates IMAP connector settings and credentials. +func (c *IMAPConnector) Validate(ctx context.Context) error { + if c == nil { + return fmt.Errorf("imap connector is nil") + } + if c.host == "" { + return fmt.Errorf("Invalid connector settings: 'imap_host' must be provided") + } + if c.port <= 0 { + return fmt.Errorf("Invalid connector settings: 'imap_port' must be a positive integer") + } + if c.username == "" { + return fmt.Errorf("Missing imap_username in credentials") + } + if c.password == "" { + return fmt.Errorf("Missing imap_password in credentials") + } + if c.batchSize <= 0 { + return fmt.Errorf("batch_size must be a positive integer") + } + client, err := c.dial(ctx, c.host, c.port, c.username, c.password) + if err != nil { + return err + } + return client.Close() +} + +// ValidateConnectorSetting validates IMAP settings from an unsaved config. +func (c *IMAPConnector) ValidateConnectorSetting(ctx context.Context, request map[string]any) error { + ctx, cancel := context.WithTimeout(ctx, connectorSettingValidationTimeout) + defer cancel() + return c.Validate(ctx) +} + +// OpenSync opens one IMAP sync session. +func (c *IMAPConnector) OpenSync(ctx context.Context, request SyncRequest) (SyncSession, error) { + client, err := c.dial(ctx, c.host, c.port, c.username, c.password) + if err != nil { + return nil, err + } + session := &imapSyncSession{ + connector: c, + client: client, + batchSize: c.batchSize, + windowStart: request.WindowStart, + windowEnd: request.WindowEnd, + hasMore: true, + } + session.applyResume(request.Resume) + return session, nil +} + +// OpenPrune opens one complete IMAP prune snapshot session. +func (c *IMAPConnector) OpenPrune(ctx context.Context, request PruneRequest) (PruneSession, error) { + client, err := c.dial(ctx, c.host, c.port, c.username, c.password) + if err != nil { + return nil, err + } + return &imapPruneSession{connector: c, client: client, batchSize: c.batchSize, hasMore: true}, nil +} + +type imapAttachment struct { + filename string + contentType string + content []byte +} + +type imapSyncSession struct { + connector *IMAPConnector + client imapClient + batchSize int + windowStart *time.Time + windowEnd time.Time + + todoMailboxes []string + currentMailbox string + todoEmailIDs []string + selected string + hasMore bool +} + +// NextBatch returns the next IMAP document batch. +func (s *imapSyncSession) NextBatch(ctx context.Context) (SyncBatch, error) { + documents := make([]SourceDocument, 0, s.batchSize) + var lastDoc *SourceDocument + for len(documents) < s.batchSize && s.hasMore { + if err := s.ensureCurrentEmail(ctx); err != nil { + return SyncBatch{}, err + } + if !s.hasMore { + break + } + + emailID := s.todoEmailIDs[0] + s.todoEmailIDs = s.todoEmailIDs[1:] + seq, err := strconv.ParseUint(emailID, 10, 32) + if err != nil { + continue + } + raw, err := s.client.Fetch(ctx, uint32(seq)) + if err != nil { + return SyncBatch{}, err + } + emailDoc, attachments, err := parseIMAPMessage(raw, s.connector.sizeThreshold) + if err != nil { + log.Printf("imap: skip message seq %d in mailbox %q: %v", seq, s.currentMailbox, err) + continue + } + if !s.inWindow(emailDoc.UpdatedAt) { + continue + } + documents = append(documents, emailDoc) + lastDoc = &emailDoc + documents = append(documents, attachments...) + } + if len(documents) == 0 { + return SyncBatch{}, io.EOF + } + return SyncBatch{Documents: documents, Checkpoint: s.checkpoint(lastDoc)}, nil +} + +// Close closes the IMAP sync session. +func (s *imapSyncSession) Close() error { + return s.client.Close() +} + +// ensureCurrentEmail makes sure the session has a list of mailboxes and a +// current mailbox with remaining email IDs, advancing as needed. +func (s *imapSyncSession) ensureCurrentEmail(ctx context.Context) error { + if s.todoMailboxes == nil { + mailboxes, err := s.listMailboxes(ctx) + if err != nil { + return err + } + s.todoMailboxes = mailboxes + if len(mailboxes) == 0 { + s.hasMore = false + return nil + } + } + for { + for s.currentMailbox == "" || len(s.todoEmailIDs) == 0 { + if len(s.todoMailboxes) == 0 { + s.hasMore = false + return nil + } + mailbox := s.todoMailboxes[0] + s.todoMailboxes = s.todoMailboxes[1:] + emailIDs, err := s.searchMailbox(ctx, mailbox) + if err != nil { + return err + } + s.currentMailbox = mailbox + s.todoEmailIDs = emailIDs + } + if s.selected == s.currentMailbox { + return nil + } + if err := s.client.SelectMailbox(ctx, s.currentMailbox); err != nil { + s.currentMailbox = "" + s.todoEmailIDs = nil + continue + } + s.selected = s.currentMailbox + return nil + } +} + +// listMailboxes returns configured mailboxes or discovers all mailboxes. +func (s *imapSyncSession) listMailboxes(ctx context.Context) ([]string, error) { + if len(s.connector.mailboxes) > 0 { + return s.connector.mailboxes, nil + } + mailboxes, err := s.client.List(ctx) + if err != nil { + return nil, err + } + if len(mailboxes) == 0 { + return []string{"INBOX"}, nil + } + return mailboxes, nil +} + +// searchMailbox selects a mailbox and returns the email IDs in its window. +func (s *imapSyncSession) searchMailbox(ctx context.Context, mailbox string) ([]string, error) { + if err := s.client.SelectMailbox(ctx, mailbox); err != nil { + return nil, err + } + s.selected = mailbox + start := time.Time{} + if s.windowStart != nil { + start = *s.windowStart + } + before := s.windowEnd.AddDate(0, 0, 1) + nums, err := s.client.Search(ctx, start, before) + if err != nil { + return nil, err + } + emailIDs := make([]string, 0, len(nums)) + for _, num := range nums { + emailIDs = append(emailIDs, strconv.FormatUint(uint64(num), 10)) + } + return emailIDs, nil +} + +func (s *imapSyncSession) inWindow(t time.Time) bool { + start := time.Time{} + if s.windowStart != nil { + start = *s.windowStart + } + return t.After(start) && !t.After(s.windowEnd) +} + +func (s *imapSyncSession) checkpoint(lastDoc *SourceDocument) *SyncCheckpoint { + cursor := imapCursor{ + TodoMailboxes: s.todoMailboxes, + HasMore: s.hasMore, + } + if s.currentMailbox != "" { + cursor.CurrentMailbox = &imapMailboxCursor{ + Mailbox: s.currentMailbox, + TodoEmailIDs: s.todoEmailIDs, + } + } + data, _ := json.Marshal(cursor) + checkpoint := &SyncCheckpoint{Cursor: string(data)} + if lastDoc != nil { + checkpoint.SourceID = lastDoc.SourceID + updatedAt := lastDoc.UpdatedAt + checkpoint.UpdatedAt = &updatedAt + } + return checkpoint +} + +// applyResume restores a sync session from a previously committed checkpoint. +func (s *imapSyncSession) applyResume(checkpoint *SyncCheckpoint) { + if checkpoint == nil || checkpoint.Cursor == "" { + return + } + var cursor imapCursor + if err := json.Unmarshal([]byte(checkpoint.Cursor), &cursor); err != nil { + return + } + s.todoMailboxes = cursor.TodoMailboxes + s.hasMore = cursor.HasMore + if cursor.CurrentMailbox != nil { + s.currentMailbox = cursor.CurrentMailbox.Mailbox + s.todoEmailIDs = cursor.CurrentMailbox.TodoEmailIDs + } +} + +type imapCursor struct { + TodoMailboxes []string `json:"todo_mailboxes"` + CurrentMailbox *imapMailboxCursor `json:"current_mailbox,omitempty"` + HasMore bool `json:"has_more"` +} + +type imapMailboxCursor struct { + Mailbox string `json:"mailbox"` + TodoEmailIDs []string `json:"todo_email_ids"` +} + +type imapPruneSession struct { + connector *IMAPConnector + client imapClient + batchSize int + + todoMailboxes []string + currentMailbox string + todoEmailIDs []string + hasMore bool + buffer []SlimDocument +} + +// NextBatch returns the next IMAP prune snapshot batch. +func (s *imapPruneSession) NextBatch(ctx context.Context) (PruneBatch, error) { + documents := make([]SlimDocument, 0, s.batchSize) + for len(documents) < s.batchSize { + if len(s.buffer) > 0 { + n := min(s.batchSize-len(documents), len(s.buffer)) + documents = append(documents, s.buffer[:n]...) + s.buffer = s.buffer[n:] + continue + } + if !s.hasMore { + break + } + if err := s.ensureCurrentEmail(ctx); err != nil { + return PruneBatch{}, err + } + if !s.hasMore { + break + } + emailID := s.todoEmailIDs[0] + s.todoEmailIDs = s.todoEmailIDs[1:] + seq, err := strconv.ParseUint(emailID, 10, 32) + if err != nil { + continue + } + raw, err := s.client.Fetch(ctx, uint32(seq)) + if err != nil { + return PruneBatch{}, err + } + emailDoc, attachments, err := parseIMAPMessage(raw, s.connector.sizeThreshold) + if err != nil { + return PruneBatch{}, err + } + s.buffer = append(s.buffer, SlimDocument{SourceID: emailDoc.SourceID}) + for _, attachment := range attachments { + s.buffer = append(s.buffer, SlimDocument{SourceID: attachment.SourceID}) + } + } + if len(documents) == 0 { + return PruneBatch{}, io.EOF + } + return PruneBatch{Documents: documents}, nil +} + +// Close closes the IMAP prune session. +func (s *imapPruneSession) Close() error { + return s.client.Close() +} + +// Ensure the prune session has a current mailbox, listing mailboxes if needed. +func (s *imapPruneSession) ensureCurrentEmail(ctx context.Context) error { + if s.todoMailboxes == nil { + mailboxes, err := s.listMailboxes(ctx) + if err != nil { + return err + } + s.todoMailboxes = mailboxes + if len(mailboxes) == 0 { + s.hasMore = false + return nil + } + } + for s.currentMailbox == "" || len(s.todoEmailIDs) == 0 { + if len(s.todoMailboxes) == 0 { + s.currentMailbox = "" + s.hasMore = false + return nil + } + mailbox := s.todoMailboxes[0] + s.todoMailboxes = s.todoMailboxes[1:] + if err := s.client.SelectMailbox(ctx, mailbox); err != nil { + return err + } + nums, err := s.client.Search(ctx, time.Time{}, time.Time{}) + if err != nil { + return err + } + emailIDs := make([]string, 0, len(nums)) + for _, num := range nums { + emailIDs = append(emailIDs, strconv.FormatUint(uint64(num), 10)) + } + s.currentMailbox = mailbox + s.todoEmailIDs = emailIDs + s.hasMore = true + } + return nil +} + +func (s *imapPruneSession) listMailboxes(ctx context.Context) ([]string, error) { + if len(s.connector.mailboxes) > 0 { + return s.connector.mailboxes, nil + } + mailboxes, err := s.client.List(ctx) + if err != nil { + return nil, err + } + if len(mailboxes) == 0 { + return []string{"INBOX"}, nil + } + return mailboxes, nil +} + +// parseIMAPMessage parses a raw RFC 822 message into an email document and its +// attachment documents. +func parseIMAPMessage(raw []byte, sizeThreshold int64) (SourceDocument, []SourceDocument, error) { + entity, err := message.Read(bytes.NewReader(raw)) + if err != nil { + if !message.IsUnknownCharset(err) && !message.IsUnknownEncoding(err) { + return SourceDocument{}, nil, err + } + } + if entity == nil { + return SourceDocument{}, nil, fmt.Errorf("imap: failed to read message") + } + + body, attachments, err := walkIMAPParts(entity, sizeThreshold) + if err != nil { + return SourceDocument{}, nil, fmt.Errorf("imap: walk message parts: %w", err) + } + header := entity.Header + subject := decodedIMAPHeader(header, "Subject") + if subject == "" { + subject = "Unknown Subject" + } + from := decodedIMAPHeader(header, "From") + to := decodedIMAPHeader(header, "To") + if to == "" { + to = decodedIMAPHeader(header, "Delivered-To") + } + cc := decodedIMAPHeader(header, "Cc") + dateString := decodedIMAPHeader(header, "Date") + parsedDate := parseIMAPDate(dateString) + date := parsedDate + if date.IsZero() { + date = time.Now().UTC() + } + + messageID := decodedIMAPHeader(header, "Message-ID") + if messageID == "" { + messageID = buildGeneratedIMAPID(raw, subject, dateString, parsedDate, from, to, cc, body) + } + + emailDoc := SourceDocument{ + SourceID: messageID, + SemanticIdentifier: subject, + Extension: ".txt", + Blob: []byte(body), + UpdatedAt: date, + SizeBytes: int64(len(body)), + Fingerprint: contentFingerprint([]byte(body)), + Metadata: map[string]any{}, + } + + attachmentDocs := make([]SourceDocument, 0, len(attachments)) + for index, attachment := range attachments { + if utility.FilenameType(attachment.filename) == utility.FileTypeOTHER { + continue + } + attachmentDocs = append(attachmentDocs, SourceDocument{ + SourceID: messageID + "#att:" + strconv.Itoa(index) + ":" + attachment.filename, + SemanticIdentifier: attachment.filename, + Extension: imapAttachmentExtension(attachment.filename), + Blob: attachment.content, + UpdatedAt: date, + SizeBytes: int64(len(attachment.content)), + Fingerprint: contentFingerprint(attachment.content), + Metadata: map[string]any{ + "parent_email_id": messageID, + "parent_subject": subject, + "attachment_filename": attachment.filename, + "attachment_content_type": attachment.contentType, + }, + }) + } + return emailDoc, attachmentDocs, nil +} + +// walkIMAPParts collects the first decodable text body and the attachments. +func walkIMAPParts(entity *message.Entity, sizeThreshold int64) (string, []imapAttachment, error) { + var body string + var htmlBody string + var attachments []imapAttachment + err := entity.Walk(func(path []int, part *message.Entity, partErr error) error { + if partErr != nil { + return partErr + } + if part == nil { + return nil + } + if part.MultipartReader() != nil { + return nil + } + disposition, dispositionParams, _ := part.Header.ContentDisposition() + contentType, contentTypeParams, _ := part.Header.ContentType() + dispositionLower := strings.ToLower(disposition) + filename := firstNonEmpty(dispositionParams["filename"], contentTypeParams["name"]) + isAttachment := strings.HasPrefix(dispositionLower, "attachment") || + (strings.HasPrefix(dispositionLower, "inline") && filename != "") + + var payload []byte + var err error + if isAttachment { + payload, err = io.ReadAll(io.LimitReader(part.Body, sizeThreshold+1)) + } else { + payload, err = io.ReadAll(part.Body) + } + if err != nil { + return err + } + + if isAttachment { + if len(payload) > 0 && int64(len(payload)) <= sizeThreshold { + name := strings.TrimSpace(filename) + if name == "" { + name = "attachment.bin" + } + attachments = append(attachments, imapAttachment{ + filename: name, + contentType: contentType, + content: payload, + }) + } + // Walk only advances after the current part body is fully consumed; + // drain anything beyond the capped read so the next part is reached. + if _, err := io.Copy(io.Discard, part.Body); err != nil { + return err + } + return nil + } + + if !utf8.Valid(payload) { + return nil + } + switch strings.ToLower(contentType) { + case "text/plain": + if body == "" { + body = string(payload) + } + case "text/html": + if htmlBody == "" { + htmlBody = imapHTMLToText(string(payload)) + } + } + return nil + }) + if err != nil { + return "", nil, err + } + if body == "" { + body = htmlBody + } + return body, attachments, nil +} + +func decodedIMAPHeader(header message.Header, key string) string { + value, err := header.Text(key) + if err != nil { + value = header.Get(key) + } + return strings.TrimSpace(value) +} + +func parseIMAPDate(value string) time.Time { + if value == "" { + return time.Time{} + } + if parsed, err := netmail.ParseDate(value); err == nil { + return parsed + } + for _, layout := range []string{time.RFC3339, "Mon, 2 Jan 2006 15:04:05 -0700"} { + if parsed, err := time.Parse(layout, value); err == nil { + return parsed + } + } + return time.Time{} +} + +func buildGeneratedIMAPID(raw []byte, subject, dateString string, parsedDate time.Time, from, to, cc, body string) string { + if parsedDate.IsZero() { + dateString = strings.TrimSpace(dateString) + } else { + dateString = parsedDate.UTC().Format("2006-01-02T15:04:05") + "+00:00" + } + rawDigest := sha256.Sum256(raw) + bodyValue := []byte(body) + bodyDigest := sha256.Sum256(bodyValue) + material := strings.Join([]string{subject, dateString, from, to, cc, + hex.EncodeToString(bodyDigest[:]), hex.EncodeToString(rawDigest[:])}, "\n") + digest := sha256.Sum256([]byte(material)) + return "generated:" + hex.EncodeToString(digest[:]) +} + +func imapAttachmentExtension(filename string) string { + if index := strings.LastIndex(filename, "."); index >= 0 { + return filename[index:] + } + return "" +} + +// imapHTMLToText flattens HTML into space-separated text. +func imapHTMLToText(value string) string { + document, err := xhtml.Parse(strings.NewReader(value)) + if err != nil { + return htmlToText(value) + } + parts := []string{} + var walk func(*xhtml.Node) + walk = func(node *xhtml.Node) { + if node.Type == xhtml.ElementNode { + switch strings.ToLower(node.Data) { + case "script", "style": + return + } + } + if node.Type == xhtml.TextNode { + if text := strings.TrimSpace(node.Data); text != "" { + parts = append(parts, text) + } + } + for child := node.FirstChild; child != nil; child = child.NextSibling { + walk(child) + } + } + walk(document) + return strings.Join(parts, " ") +} + +func imapMailboxList(value any) []string { + switch typed := value.(type) { + case []any: + names := []string{} + for _, item := range typed { + if name := strings.TrimSpace(stringConfig(item)); name != "" { + names = append(names, name) + } + } + return names + default: + return splitCommaList(stringConfig(value)) + } +} + +type realIMAPClient struct { + client *imapclient.Client +} + +type imapCommandResult[T any] struct { + value T + err error +} + +func runIMAPCommand[T any](ctx context.Context, client *imapclient.Client, run func() (T, error)) (T, error) { + ctx, cancel := context.WithTimeout(ctx, imapCommandTimeout) + defer cancel() + if err := ctx.Err(); err != nil { + var zero T + return zero, err + } + done := make(chan imapCommandResult[T], 1) + go func() { + value, err := run() + done <- imapCommandResult[T]{value: value, err: err} + }() + select { + case <-ctx.Done(): + _ = client.Close() + var zero T + return zero, ctx.Err() + case result := <-done: + return result.value, result.err + } +} + +func dialRealIMAPClient(ctx context.Context, host string, port int, username, password string) (imapClient, error) { + address := net.JoinHostPort(host, strconv.Itoa(port)) + dialCtx, cancelDial := context.WithTimeout(ctx, imapDialTimeout) + defer cancelDial() + rawConn, err := (&net.Dialer{}).DialContext(dialCtx, "tcp", address) + if err != nil { + return nil, err + } + tlsConn := tls.Client(rawConn, &tls.Config{ + ServerName: host, + NextProtos: []string{"imap"}, + }) + if err := tlsConn.HandshakeContext(dialCtx); err != nil { + _ = rawConn.Close() + return nil, err + } + client := imapclient.New(tlsConn, nil) + if _, err := runIMAPCommand(ctx, client, func() (struct{}, error) { + return struct{}{}, client.Login(username, password).Wait() + }); err != nil { + if errors.Is(err, context.DeadlineExceeded) { + // runIMAPCommand already closed the client on timeout; skip LOGOUT. + _ = client.Close() + } else { + _ = client.Logout().Wait() + _ = client.Close() + } + return nil, err + } + return &realIMAPClient{client: client}, nil +} + +func (c *realIMAPClient) List(ctx context.Context) ([]string, error) { + listed, err := runIMAPCommand(ctx, c.client, func() ([]*imap.ListData, error) { + return c.client.List("", "*", nil).Collect() + }) + if err != nil { + return nil, err + } + mailboxes := []string{} + for _, data := range listed { + if data == nil || data.Mailbox == "" { + continue + } + mailboxes = append(mailboxes, data.Mailbox) + } + return mailboxes, nil +} + +func (c *realIMAPClient) SelectMailbox(ctx context.Context, mailbox string) error { + // Send a real SELECT, not EXAMINE. go-imap emits EXAMINE when ReadOnly + // is set, which some servers acknowledge without entering the selected + // state. Fetches already use Peek, so a writable SELECT is read-safe. + _, err := runIMAPCommand(ctx, c.client, func() (*imap.SelectData, error) { + return c.client.Select(mailbox, nil).Wait() + }) + return err +} + +func (c *realIMAPClient) Search(ctx context.Context, since, before time.Time) ([]uint32, error) { + criteria := &imap.SearchCriteria{} + if !since.IsZero() { + criteria.Since = since + } + if !before.IsZero() { + criteria.Before = before + } + data, err := runIMAPCommand(ctx, c.client, func() (*imap.SearchData, error) { + return c.client.Search(criteria, nil).Wait() + }) + if err != nil { + return nil, err + } + return data.AllSeqNums(), nil +} + +func (c *realIMAPClient) Fetch(ctx context.Context, seqNum uint32) ([]byte, error) { + buffers, err := runIMAPCommand(ctx, c.client, func() ([]*imapclient.FetchMessageBuffer, error) { + return c.client.Fetch(imap.SeqSetNum(seqNum), &imap.FetchOptions{ + BodySection: []*imap.FetchItemBodySection{{Peek: true}}, + }).Collect() + }) + if err != nil { + return nil, err + } + for _, buffer := range buffers { + if buffer == nil { + continue + } + for _, section := range buffer.BodySection { + if len(section.Bytes) > 0 { + return section.Bytes, nil + } + } + } + return nil, nil +} + +func (c *realIMAPClient) Close() error { + _, logoutErr := runIMAPCommand(context.Background(), c.client, func() (struct{}, error) { + return struct{}{}, c.client.Logout().Wait() + }) + // Always close the connection even when LOGOUT times out or fails. + _ = c.client.Close() + return logoutErr +} diff --git a/internal/syncer/connector/imap_test.go b/internal/syncer/connector/imap_test.go new file mode 100644 index 0000000000..7b8a2e394e --- /dev/null +++ b/internal/syncer/connector/imap_test.go @@ -0,0 +1,567 @@ +// +// Copyright 2026 The InfiniFlow Authors. All Rights Reserved. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// + +package connector + +import ( + "context" + "encoding/base64" + "encoding/json" + "io" + "strings" + "testing" + "time" +) + +type imapSearchCall struct { + since time.Time + before time.Time +} + +type fakeIMAPClient struct { + mailboxes []string + current string + selected []string + searchByMailbox map[string][]uint32 + searchCalls []imapSearchCall + rawBySeq map[uint32][]byte + fetched []uint32 + closed bool + listErr error + searchErr error + selectErrs map[string]error +} + +func (f *fakeIMAPClient) List(ctx context.Context) ([]string, error) { + return f.mailboxes, f.listErr +} + +func (f *fakeIMAPClient) SelectMailbox(ctx context.Context, mailbox string) error { + f.selected = append(f.selected, mailbox) + f.current = mailbox + if f.selectErrs != nil { + if err := f.selectErrs[mailbox]; err != nil { + return err + } + } + return nil +} + +func (f *fakeIMAPClient) Search(ctx context.Context, since, before time.Time) ([]uint32, error) { + f.searchCalls = append(f.searchCalls, imapSearchCall{since: since, before: before}) + if f.searchErr != nil { + return nil, f.searchErr + } + return f.searchByMailbox[f.current], nil +} + +func (f *fakeIMAPClient) Fetch(ctx context.Context, seqNum uint32) ([]byte, error) { + f.fetched = append(f.fetched, seqNum) + return f.rawBySeq[seqNum], nil +} + +func (f *fakeIMAPClient) Close() error { + f.closed = true + return nil +} + +func rawIMAPEmail(id, subject, date, body string) []byte { + var b strings.Builder + b.WriteString("From: Sender \r\n") + b.WriteString("To: Recipient \r\n") + b.WriteString("Subject: " + subject + "\r\n") + b.WriteString("Date: " + date + "\r\n") + b.WriteString("Message-ID: <" + id + "@example.com>\r\n") + b.WriteString("MIME-Version: 1.0\r\n") + b.WriteString("Content-Type: text/plain; charset=utf-8\r\n") + b.WriteString("Content-Transfer-Encoding: 8bit\r\n\r\n") + b.WriteString(body + "\r\n") + return []byte(b.String()) +} + +func multipartIMAPEmail(id, subject, date, body, attachmentName string, attachmentContent []byte) []byte { + var b strings.Builder + b.WriteString("From: Sender \r\n") + b.WriteString("To: Recipient \r\n") + b.WriteString("Subject: " + subject + "\r\n") + b.WriteString("Date: " + date + "\r\n") + b.WriteString("Message-ID: <" + id + "@example.com>\r\n") + b.WriteString("MIME-Version: 1.0\r\n") + b.WriteString("Content-Type: multipart/mixed; boundary=\"BOUNDARY\"\r\n\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: text/plain; charset=utf-8\r\n") + b.WriteString("Content-Transfer-Encoding: 8bit\r\n\r\n") + b.WriteString(body + "\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: application/pdf\r\n") + b.WriteString("Content-Disposition: attachment; filename=\"" + attachmentName + "\"\r\n") + b.WriteString("Content-Transfer-Encoding: base64\r\n\r\n") + b.WriteString(base64.StdEncoding.EncodeToString(attachmentContent) + "\r\n") + b.WriteString("--BOUNDARY--\r\n") + return []byte(b.String()) +} + +func multipartIMAPEmailWithCharset(id, subject, date, charset string, body []byte) []byte { + var b strings.Builder + b.WriteString("From: Sender \r\n") + b.WriteString("To: Recipient \r\n") + b.WriteString("Subject: " + subject + "\r\n") + b.WriteString("Date: " + date + "\r\n") + b.WriteString("Message-ID: <" + id + "@example.com>\r\n") + b.WriteString("MIME-Version: 1.0\r\n") + b.WriteString("Content-Type: multipart/mixed; boundary=\"BOUNDARY\"\r\n\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: text/plain; charset=" + charset + "\r\n") + b.WriteString("Content-Transfer-Encoding: 8bit\r\n\r\n") + b.Write(body) + b.WriteString("\r\n--BOUNDARY--\r\n") + return []byte(b.String()) +} + +func newTestIMAPConnector(client *fakeIMAPClient, batchSize int) *IMAPConnector { + return &IMAPConnector{ + host: "imap.example.com", + port: 993, + mailboxes: []string{"INBOX"}, + username: "user", + password: "pass", + batchSize: batchSize, + sizeThreshold: defaultIMAPSizeThreshold, + dial: func(ctx context.Context, host string, port int, username, password string) (imapClient, error) { + return client, nil + }, + } +} + +func TestNewIMAPConnectorParsesConfig(t *testing.T) { + connector, err := NewIMAPConnector(map[string]any{ + "imap_host": " imap.example.com ", + "imap_port": 993, + "imap_mailbox": []any{"INBOX", ""}, + "sync_batch_size": 7, + "credentials": map[string]any{ + "imap_username": "user", + "imap_password": "pass", + }, + }) + if err != nil { + t.Fatalf("NewIMAPConnector failed: %v", err) + } + if connector.host != "imap.example.com" { + t.Fatalf("host = %q", connector.host) + } + if connector.port != 993 { + t.Fatalf("port = %d", connector.port) + } + if len(connector.mailboxes) != 1 || connector.mailboxes[0] != "INBOX" { + t.Fatalf("mailboxes = %v", connector.mailboxes) + } + if connector.batchSize != 7 { + t.Fatalf("batch_size = %d", connector.batchSize) + } + if connector.username != "user" || connector.password != "pass" { + t.Fatalf("credentials = %q %q", connector.username, connector.password) + } +} + +func TestIMAPValidate(t *testing.T) { + client := &fakeIMAPClient{} + connector := newTestIMAPConnector(client, 32) + if err := connector.Validate(context.Background()); err != nil { + t.Fatalf("Validate failed: %v", err) + } + if !client.closed { + t.Fatalf("expected the IMAP client to be closed after validation") + } + + missing := newTestIMAPConnector(&fakeIMAPClient{}, 32) + missing.host = "" + if err := missing.Validate(context.Background()); err == nil || !strings.Contains(err.Error(), "imap_host") { + t.Fatalf("Validate error = %v, want missing imap_host", err) + } +} + +func TestIMAPSyncInitialBatchAndResume(t *testing.T) { + start := mustTime(t, "2026-01-02T12:00:00Z") + end := mustTime(t, "2026-01-03T00:00:00Z") + client := &fakeIMAPClient{ + mailboxes: []string{"INBOX"}, + searchByMailbox: map[string][]uint32{ + "INBOX": {1, 2}, + }, + rawBySeq: map[uint32][]byte{ + 1: rawIMAPEmail("msg1", "Hello one", "Mon, 2 Jan 2026 12:30:00 +0000", "body one"), + 2: rawIMAPEmail("msg2", "Hello two", "Mon, 2 Jan 2026 22:00:00 +0000", "body two"), + }, + } + connector := newTestIMAPConnector(client, 1) + + session, err := connector.OpenSync(context.Background(), SyncRequest{ + WindowStart: &start, + WindowEnd: end, + }) + if err != nil { + t.Fatalf("OpenSync failed: %v", err) + } + batch1, err := session.NextBatch(context.Background()) + if err != nil { + t.Fatalf("first NextBatch failed: %v", err) + } + if len(batch1.Documents) != 1 || batch1.Documents[0].SourceID != "" { + t.Fatalf("batch1 documents = %v", batch1.Documents) + } + if batch1.Checkpoint == nil || batch1.Checkpoint.Cursor == "" { + t.Fatalf("batch1 checkpoint is nil") + } + if len(client.searchCalls) != 1 { + t.Fatalf("searchCalls = %d, want 1", len(client.searchCalls)) + } + if !client.searchCalls[0].since.Equal(start) { + t.Fatalf("search since = %v, want %v", client.searchCalls[0].since, start) + } + wantBefore := end.AddDate(0, 0, 1) + if !client.searchCalls[0].before.Equal(wantBefore) { + t.Fatalf("search before = %v, want %v", client.searchCalls[0].before, wantBefore) + } + + var cursor imapCursor + if err := json.Unmarshal([]byte(batch1.Checkpoint.Cursor), &cursor); err != nil { + t.Fatalf("unmarshal cursor: %v", err) + } + if cursor.CurrentMailbox == nil || cursor.CurrentMailbox.Mailbox != "INBOX" || len(cursor.CurrentMailbox.TodoEmailIDs) != 1 || cursor.CurrentMailbox.TodoEmailIDs[0] != "2" { + t.Fatalf("cursor = %+v", cursor) + } + + client.fetched = nil + session2, err := connector.OpenSync(context.Background(), SyncRequest{ + WindowStart: &start, + WindowEnd: end, + Resume: batch1.Checkpoint, + }) + if err != nil { + t.Fatalf("OpenSync on resume failed: %v", err) + } + batch2, err := session2.NextBatch(context.Background()) + if err != nil { + t.Fatalf("resumed NextBatch failed: %v", err) + } + if len(batch2.Documents) != 1 || batch2.Documents[0].SourceID != "" { + t.Fatalf("resumed documents = %v", batch2.Documents) + } + if len(client.fetched) != 1 || client.fetched[0] != 2 { + t.Fatalf("resumed fetch = %v, want [2]", client.fetched) + } +} + +func TestIMAPSyncResumeSelectsMailboxBeforeFetch(t *testing.T) { + start := mustTime(t, "2026-01-02T12:00:00Z") + end := mustTime(t, "2026-01-03T00:00:00Z") + cursorData, err := json.Marshal(imapCursor{ + TodoMailboxes: []string{}, + HasMore: true, + CurrentMailbox: &imapMailboxCursor{ + Mailbox: "INBOX", + TodoEmailIDs: []string{"2"}, + }, + }) + if err != nil { + t.Fatalf("marshal cursor: %v", err) + } + + // Use a fresh connection so the resumed session can only rely on its own + // select call before fetching, matching a real resumed IMAP connection. + client := &fakeIMAPClient{ + mailboxes: []string{"INBOX"}, + rawBySeq: map[uint32][]byte{ + 2: rawIMAPEmail("msg2", "Hello two", "Mon, 2 Jan 2026 22:00:00 +0000", "body two"), + }, + } + connector := newTestIMAPConnector(client, 32) + session, err := connector.OpenSync(context.Background(), SyncRequest{ + WindowStart: &start, + WindowEnd: end, + Resume: &SyncCheckpoint{Cursor: string(cursorData)}, + }) + if err != nil { + t.Fatalf("OpenSync on resume failed: %v", err) + } + batch, err := session.NextBatch(context.Background()) + if err != nil { + t.Fatalf("resumed NextBatch failed: %v", err) + } + if len(client.selected) == 0 || client.selected[0] != "INBOX" { + t.Fatalf("resumed select = %v, want INBOX selected before fetch", client.selected) + } + if len(client.fetched) != 1 || client.fetched[0] != 2 { + t.Fatalf("resumed fetch = %v, want [2]", client.fetched) + } + if len(batch.Documents) != 1 || batch.Documents[0].SourceID != "" { + t.Fatalf("resumed documents = %v", batch.Documents) + } +} + +func TestIMAPSyncFiltersOutsideWindow(t *testing.T) { + start := mustTime(t, "2026-01-02T00:00:00Z") + end := mustTime(t, "2026-01-03T00:00:00Z") + client := &fakeIMAPClient{ + mailboxes: []string{"INBOX"}, + searchByMailbox: map[string][]uint32{ + "INBOX": {1}, + }, + rawBySeq: map[uint32][]byte{ + 1: rawIMAPEmail("old", "Old mail", "Thu, 1 Jan 2026 09:00:00 +0000", "old"), + }, + } + connector := newTestIMAPConnector(client, 32) + session, err := connector.OpenSync(context.Background(), SyncRequest{ + WindowStart: &start, + WindowEnd: end, + }) + if err != nil { + t.Fatalf("OpenSync failed: %v", err) + } + if _, err := session.NextBatch(context.Background()); err != io.EOF { + t.Fatalf("NextBatch error = %v, want io.EOF", err) + } +} + +func TestIMAPParseMessageWithAttachment(t *testing.T) { + attachment := []byte("%PDF-1.4 fake") + raw := multipartIMAPEmail("msg3", "With att", "Mon, 2 Jan 2026 09:00:00 +0000", "Hello body", "report.pdf", attachment) + emailDoc, attachments, err := parseIMAPMessage(raw, defaultIMAPSizeThreshold) + if err != nil { + t.Fatalf("parseIMAPMessage failed: %v", err) + } + if emailDoc.SourceID != "" { + t.Fatalf("email SourceID = %q", emailDoc.SourceID) + } + if string(emailDoc.Blob) != "Hello body" { + t.Fatalf("email body = %q", emailDoc.Blob) + } + if len(attachments) != 1 { + t.Fatalf("attachments = %d, want 1", len(attachments)) + } + attr := attachments[0] + if attr.SourceID != "#att:0:report.pdf" { + t.Fatalf("attachment SourceID = %q", attr.SourceID) + } + if attr.Extension != ".pdf" { + t.Fatalf("attachment extension = %q", attr.Extension) + } + if string(attr.Blob) != string(attachment) { + t.Fatalf("attachment content mismatch") + } +} + +func TestIMAPParseMessageSkipsUnsupportedAttachment(t *testing.T) { + var b strings.Builder + b.WriteString("From: Sender \r\n") + b.WriteString("To: Recipient \r\n") + b.WriteString("Subject: Mixed att\r\n") + b.WriteString("Date: Mon, 2 Jan 2026 09:00:00 +0000\r\n") + b.WriteString("Message-ID: \r\n") + b.WriteString("MIME-Version: 1.0\r\n") + b.WriteString("Content-Type: multipart/mixed; boundary=\"BOUNDARY\"\r\n\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: text/plain; charset=utf-8\r\n") + b.WriteString("Content-Transfer-Encoding: 8bit\r\n\r\n") + b.WriteString("Hello body\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: application/pdf\r\n") + b.WriteString("Content-Disposition: attachment; filename=\"report.pdf\"\r\n") + b.WriteString("Content-Transfer-Encoding: base64\r\n\r\n") + b.WriteString(base64.StdEncoding.EncodeToString([]byte("%PDF-1.4 fake")) + "\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: application/zip\r\n") + b.WriteString("Content-Disposition: attachment; filename=\"archive.zip\"\r\n") + b.WriteString("Content-Transfer-Encoding: base64\r\n\r\n") + b.WriteString(base64.StdEncoding.EncodeToString([]byte("PK\x03\x04fake")) + "\r\n") + b.WriteString("--BOUNDARY--\r\n") + + emailDoc, attachments, err := parseIMAPMessage([]byte(b.String()), defaultIMAPSizeThreshold) + if err != nil { + t.Fatalf("parseIMAPMessage failed: %v", err) + } + if emailDoc.SourceID != "" { + t.Fatalf("email SourceID = %q", emailDoc.SourceID) + } + if len(attachments) != 1 { + t.Fatalf("attachments = %d, want 1", len(attachments)) + } + attr := attachments[0] + if attr.SourceID != "#att:0:report.pdf" { + t.Fatalf("attachment SourceID = %q", attr.SourceID) + } + if attr.Extension != ".pdf" { + t.Fatalf("attachment extension = %q", attr.Extension) + } +} + +func TestIMAPParseMessageAttachmentSourceIDsUnique(t *testing.T) { + var b strings.Builder + b.WriteString("From: Sender \r\n") + b.WriteString("To: Recipient \r\n") + b.WriteString("Subject: Duplicate att names\r\n") + b.WriteString("Date: Mon, 2 Jan 2026 09:00:00 +0000\r\n") + b.WriteString("Message-ID: \r\n") + b.WriteString("MIME-Version: 1.0\r\n") + b.WriteString("Content-Type: multipart/mixed; boundary=\"BOUNDARY\"\r\n\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: text/plain; charset=utf-8\r\n") + b.WriteString("Content-Transfer-Encoding: 8bit\r\n\r\n") + b.WriteString("Hello body\r\n") + for _, content := range []string{"%PDF-1.4 first", "%PDF-1.4 second"} { + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: application/pdf\r\n") + b.WriteString("Content-Disposition: attachment; filename=\"report.pdf\"\r\n") + b.WriteString("Content-Transfer-Encoding: base64\r\n\r\n") + b.WriteString(base64.StdEncoding.EncodeToString([]byte(content)) + "\r\n") + } + b.WriteString("--BOUNDARY--\r\n") + + emailDoc, attachments, err := parseIMAPMessage([]byte(b.String()), defaultIMAPSizeThreshold) + if err != nil { + t.Fatalf("parseIMAPMessage failed: %v", err) + } + if emailDoc.SourceID != "" { + t.Fatalf("email SourceID = %q", emailDoc.SourceID) + } + if len(attachments) != 2 { + t.Fatalf("attachments = %d, want 2", len(attachments)) + } + want := []string{"#att:0:report.pdf", "#att:1:report.pdf"} + for i, attr := range attachments { + if attr.SourceID != want[i] { + t.Fatalf("attachment %d SourceID = %q, want %q", i, attr.SourceID, want[i]) + } + } +} + +func TestIMAPParseMessageIgnoresOversizedTextAttachment(t *testing.T) { + sizeThreshold := int64(16) + var b strings.Builder + b.WriteString("From: Sender \r\n") + b.WriteString("To: Recipient \r\n") + b.WriteString("Subject: Oversized text attachment\r\n") + b.WriteString("Date: Mon, 2 Jan 2026 09:00:00 +0000\r\n") + b.WriteString("Message-ID: \r\n") + b.WriteString("MIME-Version: 1.0\r\n") + b.WriteString("Content-Type: multipart/mixed; boundary=\"BOUNDARY\"\r\n\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: text/plain; charset=utf-8\r\n") + b.WriteString("Content-Disposition: attachment; filename=\"oversized.txt\"\r\n") + b.WriteString("Content-Transfer-Encoding: 8bit\r\n\r\n") + b.WriteString(strings.Repeat("x", 1024) + "\r\n") + b.WriteString("--BOUNDARY\r\n") + b.WriteString("Content-Type: text/plain; charset=utf-8\r\n") + b.WriteString("Content-Transfer-Encoding: 8bit\r\n\r\n") + b.WriteString("Hello body\r\n") + b.WriteString("--BOUNDARY--\r\n") + + emailDoc, attachments, err := parseIMAPMessage([]byte(b.String()), sizeThreshold) + if err != nil { + t.Fatalf("parseIMAPMessage failed: %v", err) + } + if string(emailDoc.Blob) != "Hello body" { + t.Fatalf("email body = %q, want %q", emailDoc.Blob, "Hello body") + } + if len(attachments) != 0 { + t.Fatalf("attachments = %d, want 0 oversized text attachments", len(attachments)) + } +} + +func TestIMAPParseMessageDecodesISO88591Body(t *testing.T) { + raw := multipartIMAPEmailWithCharset("msg6", "Latin body", "Mon, 2 Jan 2026 09:00:00 +0000", "iso-8859-1", []byte("caf\xe9")) + emailDoc, attachments, err := parseIMAPMessage(raw, defaultIMAPSizeThreshold) + if err != nil { + t.Fatalf("parseIMAPMessage failed: %v", err) + } + if string(emailDoc.Blob) != "caf\u00e9" { + t.Fatalf("email body = %q, want %q", emailDoc.Blob, "caf\u00e9") + } + if len(attachments) != 0 { + t.Fatalf("attachments = %d, want 0", len(attachments)) + } +} + +func TestIMAPParseMessageDecodesWindows1252Body(t *testing.T) { + raw := multipartIMAPEmailWithCharset("msg7", "Euro body", "Mon, 2 Jan 2026 09:00:00 +0000", "windows-1252", []byte("price: 5\x80")) + emailDoc, attachments, err := parseIMAPMessage(raw, defaultIMAPSizeThreshold) + if err != nil { + t.Fatalf("parseIMAPMessage failed: %v", err) + } + if string(emailDoc.Blob) != "price: 5\u20ac" { + t.Fatalf("email body = %q, want %q", emailDoc.Blob, "price: 5\u20ac") + } + if len(attachments) != 0 { + t.Fatalf("attachments = %d, want 0", len(attachments)) + } +} + +func TestIMAPParseMessageKeepsBodyOverAttachmentThreshold(t *testing.T) { + sizeThreshold := int64(16) + body := strings.Repeat("y", 64) + raw := multipartIMAPEmailWithCharset("msg8", "Long body", "Mon, 2 Jan 2026 09:00:00 +0000", "utf-8", []byte(body)) + emailDoc, attachments, err := parseIMAPMessage(raw, sizeThreshold) + if err != nil { + t.Fatalf("parseIMAPMessage failed: %v", err) + } + if string(emailDoc.Blob) != body { + t.Fatalf("email body = %q, want %q", emailDoc.Blob, body) + } + if len(attachments) != 0 { + t.Fatalf("attachments = %d, want 0", len(attachments)) + } +} + +func TestIMAPPruneReturnsSlimDocs(t *testing.T) { + client := &fakeIMAPClient{ + mailboxes: []string{"INBOX"}, + searchByMailbox: map[string][]uint32{ + "INBOX": {1, 2}, + }, + rawBySeq: map[uint32][]byte{ + 1: rawIMAPEmail("p1", "Prune one", "Mon, 2 Jan 2026 10:00:00 +0000", "one"), + 2: rawIMAPEmail("p2", "Prune two", "Mon, 2 Jan 2026 11:00:00 +0000", "two"), + }, + } + connector := newTestIMAPConnector(client, 1) + prune, err := connector.OpenPrune(context.Background(), PruneRequest{}) + if err != nil { + t.Fatalf("OpenPrune failed: %v", err) + } + + batch1, err := prune.NextBatch(context.Background()) + if err != nil { + t.Fatalf("first prune batch failed: %v", err) + } + if len(batch1.Documents) != 1 || batch1.Documents[0].SourceID != "" { + t.Fatalf("prune batch1 = %v", batch1.Documents) + } + batch2, err := prune.NextBatch(context.Background()) + if err != nil { + t.Fatalf("second prune batch failed: %v", err) + } + if len(batch2.Documents) != 1 || batch2.Documents[0].SourceID != "" { + t.Fatalf("prune batch2 = %v", batch2.Documents) + } + if _, err := prune.NextBatch(context.Background()); err != io.EOF { + t.Fatalf("final prune batch error = %v, want io.EOF", err) + } + if len(client.fetched) != 2 || client.fetched[0] != 1 || client.fetched[1] != 2 { + t.Fatalf("prune fetch = %v, want [1 2]", client.fetched) + } +}