package imapx import ( "bytes" "context" "errors" "fmt" "io" "log/slog" "net" "os" "strings" "time" "github.com/emersion/go-imap/v2" "github.com/emersion/go-imap/v2/imapclient" ) // CopyDeps injects the dedup/progress hooks used by CopyFolder. APPEND to // dst always happens before MarkMigrated is called, so a crash between the // two only ever causes a message to be re-copied (never lost) on the next // run. type CopyDeps struct { IsMigrated func(key string) (bool, error) MarkMigrated func(folder, key string) error OnProgress func(copied, skipped int) // OnFolder is called once per folder right after EXAMINE, before the // (potentially long) envelope fetch, with the message count in the source // folder — for progress visibility. OnFolder func(srcFolder, dstFolder string, total int64) // OnScan is called during the streaming metadata pass with how many of the // folder's messages have been examined so far — so the UI shows movement // while dedup decisions are made, before bodies start copying. OnScan func(scanned, total int64) // OnError is called for each message-level error, with a message reference // ("UID N: subject", empty when the envelope is unavailable) and the error // text — so the orchestrator can persist individual errors for the // per-account error modal. Folder-level errors are reported by the caller. OnError func(ref, msg string) // OnActivity is called repeatedly WHILE a single message body is streamed // (each FETCH read chunk and each APPEND write chunk). Copying one large // message can take longer than the orchestrator's stall timeout; without an // in-body signal the watchdog can't tell a slow-but-live transfer from a // wedged connection and cancels a healthy copy. May be nil. OnActivity func() // ReconnectSrc dials and logs in a FRESH source client, returning it. It is // called when a body read times out (server under-delivered a literal), // which leaves the current src connection desynced and unusable. CopyFolder // swaps to the returned client, re-EXAMINEs the folder, and resumes. The // implementation is expected to also update any external reference to the // live src client (e.g. so a cancel path closes the right connection). If // nil, a body-read timeout aborts the folder instead of recovering. ReconnectSrc func() (*Client, error) } // ErrBodyTimeout means a message body did not finish transferring within the // idle deadline — the server stopped sending mid-literal. It is almost always a // server announcing a BODY[] literal larger than the bytes it actually sends, // which makes go-imap wait forever for bytes that never come. Distinct from a // closed connection so the caller can skip just this one message and resume. var ErrBodyTimeout = errors.New("message body read timed out") // bodyIdleTimeout bounds how long a body read may go with NO bytes arriving // before it is abandoned. It is generous enough for legitimately slow servers // (even ~15 KB/s links keep bytes flowing far more often than this) yet well // under the orchestrator's multi-minute stall watchdog, so an under-delivered // literal is caught quickly and locally instead of stalling the whole account. // A var (not const) so tests can shorten it. var bodyIdleTimeout = 30 * time.Second // CopyResult summarizes the outcome of one CopyFolder run. type CopyResult struct { Copied int Skipped int Errors int } // metaScanBatch bounds how many messages one Pass-1 metadata FETCH covers. A // single unbounded FETCH 1:* over a large mailbox keeps one command open for // the entire scan; under parallel load the server can stop responding and, // since go-imap has no per-command deadline, the worker wedges forever. Short // windows keep each command brief so the server stays responsive and ctx is // checked between windows. const metaScanBatch = 1000 // msgRef builds a human-readable reference for a message error: its UID plus // subject when known, e.g. "UID 42: Invoice". Falls back to just the UID. func msgRef(uid imap.UID, subject string) string { if subject == "" { return fmt.Sprintf("UID %d", uid) } return fmt.Sprintf("UID %d: %s", uid, subject) } // metaBatches tiles 1..total into contiguous, non-overlapping windows of at // most batchSize, covering every sequence number exactly once. func metaBatches(total, batchSize uint32) []imap.SeqRange { if total == 0 || batchSize == 0 { return nil } var out []imap.SeqRange for start := uint32(1); start <= total; start += batchSize { stop := start + batchSize - 1 if stop > total { stop = total } out = append(out, imap.SeqRange{Start: start, Stop: stop}) if stop == total { break // guard against uint32 overflow when total is near max } } return out } // CopyFolder streams messages from srcFolder on src to dstFolder on dst. // // The source folder is opened read-only (EXAMINE) and is never mutated: // no \Deleted flags are set and no EXPUNGE is issued. Each message body is // held in memory only for the duration of a single FETCH->APPEND and is // never written to disk. Messages already migrated (per deps.IsMigrated) // are skipped without re-fetching their bodies. func CopyFolder(ctx context.Context, src, dst *Client, srcFolder, dstFolder string, deps CopyDeps) (CopyResult, error) { var res CopyResult sel, err := src.Select(srcFolder, &imap.SelectOptions{ReadOnly: true}).Wait() if err != nil { return res, fmt.Errorf("examine src %q: %w", srcFolder, err) } total := int64(sel.NumMessages) if deps.OnFolder != nil { deps.OnFolder(srcFolder, dstFolder, total) } if total == 0 { return res, nil } // dst folder must exist (idempotent create; ignore "already exists"). _ = dst.Create(dstFolder, nil).Wait() // Pass 1: STREAM metadata (no bodies) via Next(), dedup as we go, and queue // only the new messages. Streaming (not Collect) means progress shows during // the scan and memory stays flat — we hold small meta for new messages only. type queued struct { uid imap.UID key string subject string flags []imap.Flag internalDate time.Time } var todo []queued var scanned int64 reportErr := func(ref, msg string) { if deps.OnError != nil { deps.OnError(ref, msg) } } // Scan metadata in bounded windows instead of one FETCH 1:*, so each // command is short (the server stays responsive) and ctx is checked on // every window boundary — not just between messages of one giant command. for _, win := range metaBatches(sel.NumMessages, metaScanBatch) { if err := ctx.Err(); err != nil { return res, err } fc := src.Fetch(imap.SeqSet{win}, &imap.FetchOptions{ UID: true, Envelope: true, RFC822Size: true, Flags: true, InternalDate: true, }) for { if err := ctx.Err(); err != nil { _ = fc.Close() return res, err } msg := fc.Next() if msg == nil { break } buf, err := msg.Collect() if err != nil { res.Errors++ reportErr("", "read message metadata: "+err.Error()) continue } scanned++ key := MessageKey(buf.Envelope, buf.RFC822Size) subject := "" if buf.Envelope != nil { subject = buf.Envelope.Subject } already, err := deps.IsMigrated(key) if err != nil { res.Errors++ reportErr(msgRef(buf.UID, subject), "dedup lookup: "+err.Error()) } else if already { res.Skipped++ if deps.OnProgress != nil { deps.OnProgress(res.Copied, res.Skipped) } } else { todo = append(todo, queued{uid: buf.UID, key: key, subject: subject, flags: buf.Flags, internalDate: buf.InternalDate}) } if deps.OnScan != nil { deps.OnScan(scanned, total) } } if err := fc.Close(); err != nil { return res, fmt.Errorf("fetch meta %q: %w", srcFolder, err) } } // Pass 2: fetch bodies for the queued (new) messages, one at a time. for _, q := range todo { if err := ctx.Err(); err != nil { return res, err } if err := streamOne(src, dst, dstFolder, q.uid, q.flags, q.internalDate, deps.OnActivity); err != nil { res.Errors++ reportErr(msgRef(q.uid, q.subject), "copy message: "+err.Error()) // A body-read timeout means the server under-delivered this message's // literal; the src connection is now desynced. Mark the message // migrated so this and future runs skip it (it is un-fetchable via a // conforming client), then reconnect src and resume the folder with // the remaining queued messages. if errors.Is(err, ErrBodyTimeout) && deps.ReconnectSrc != nil { if merr := deps.MarkMigrated(dstFolder, q.key); merr != nil { reportErr(msgRef(q.uid, q.subject), "mark skipped: "+merr.Error()) } newSrc, rerr := deps.ReconnectSrc() if rerr != nil { return res, fmt.Errorf("reconnect src after body timeout in %q: %w", srcFolder, rerr) } src = newSrc if _, serr := src.Select(srcFolder, &imap.SelectOptions{ReadOnly: true}).Wait(); serr != nil { return res, fmt.Errorf("re-examine %q after reconnect: %w", srcFolder, serr) } continue } // A closed/broken connection won't recover: every remaining APPEND // would fail identically. Abort the folder instead of logging // thousands of the same error; a re-run resumes via dedup. if errors.Is(err, net.ErrClosed) || strings.Contains(err.Error(), "use of closed network connection") { return res, fmt.Errorf("dst connection lost in %q: %w", dstFolder, err) } continue } if err := deps.MarkMigrated(dstFolder, q.key); err != nil { res.Errors++ reportErr(msgRef(q.uid, q.subject), "mark migrated: "+err.Error()) continue } res.Copied++ if deps.OnProgress != nil { deps.OnProgress(res.Copied, res.Skipped) } } return res, nil } // slowMessage marks how long one message's FETCH or APPEND phase may take // before it is logged as anomalous. Well below the orchestrator's 3-minute // stall timeout, so a message that trips the watchdog always leaves a log line // naming the phase (FETCH vs APPEND) and size — turning a silent stall into // evidence of which side and which message is the culprit. const slowMessage = 20 * time.Second // deadlineReader arms an idle read deadline on the source socket before every // read of a message body, so a server that stops sending mid-literal — having // announced a larger BODY[] size than it actually delivers — trips the deadline // instead of blocking go-imap forever waiting for bytes that never arrive. It // also pings onActivity as bytes arrive, feeding the orchestrator's stall // watchdog. The deadline is refreshed on each read, so it bounds IDLE time // (no bytes) rather than total transfer time — a legitimately slow but steady // download never trips it. type deadlineReader struct { c *Client r io.Reader idle time.Duration on func() } func (d deadlineReader) Read(p []byte) (int, error) { _ = d.c.SetReadDeadline(time.Now().Add(d.idle)) n, err := d.r.Read(p) if n > 0 && d.on != nil { d.on() } return n, err } // touchWriter wraps the APPEND write stream and pings onActivity on every // non-empty write, so a long upload keeps the stall watchdog fed byte-by-byte. type touchWriter struct { w io.Writer on func() } func (t touchWriter) Write(p []byte) (int, error) { n, err := t.w.Write(p) if n > 0 && t.on != nil { t.on() } return n, err } // streamOne FETCHes BODY[] for one message and APPENDs it into dst without // spooling to disk. The body is buffered in RAM only for the duration of // this single FETCH->APPEND round trip. onActivity (may be nil) fires as bytes // move in either direction, feeding the orchestrator's stall watchdog. // // The body read is guarded by an idle deadline on src: if the server goes // silent mid-literal, streamOne returns ErrBodyTimeout rather than hanging, so // CopyFolder can skip the message and reconnect. func streamOne(src, dst *Client, dstFolder string, uid imap.UID, flags []imap.Flag, internalDate time.Time, onActivity func()) error { bodySection := &imap.FetchItemBodySection{} fetchStart := time.Now() fetchCmd := src.Fetch(imap.UIDSetNum(uid), &imap.FetchOptions{ BodySection: []*imap.FetchItemBodySection{bodySection}, }) defer fetchCmd.Close() msg := fetchCmd.Next() if msg == nil { return fmt.Errorf("no message for uid %v", uid) } var body []byte var readErr error for { item := msg.Next() if item == nil { break } if d, ok := item.(imapclient.FetchItemDataBodySection); ok { body, readErr = io.ReadAll(deadlineReader{c: src, r: d.Literal, idle: bodyIdleTimeout, on: onActivity}) } } // Clear the deadline before any further I/O on src (fetchCmd.Close reads the // command's completion off the same socket). _ = src.SetReadDeadline(time.Time{}) if readErr != nil { if errors.Is(readErr, os.ErrDeadlineExceeded) { return fmt.Errorf("%w: uid %v (server sent fewer bytes than the announced literal)", ErrBodyTimeout, uid) } return readErr } if err := fetchCmd.Close(); err != nil { return err } if body == nil { return fmt.Errorf("empty body uid %v", uid) } fetchDur := time.Since(fetchStart) appendStart := time.Now() appendCmd := dst.Append(dstFolder, int64(len(body)), &imap.AppendOptions{Flags: keepFlags(flags), Time: internalDate}) // Append acquires go-imap's per-client encoder mutex and holds it until // Close() calls enc.end(). Close() MUST run on every path: if io.Copy // fails mid-write (server stall, idle timeout), returning without Close() // leaks the mutex and the NEXT Append on this client deadlocks forever on // beginCommand. Close() is idempotent and always releases the lock. _, copyErr := io.Copy(touchWriter{w: appendCmd, on: onActivity}, bytes.NewReader(body)) closeErr := appendCmd.Close() if copyErr != nil { return fmt.Errorf("append body uid %v: %w", uid, copyErr) } if closeErr != nil { return closeErr } if _, err := appendCmd.Wait(); err != nil { return err } appendDur := time.Since(appendStart) // One message that individually eats a large slice of the stall budget is // the prime suspect behind a "no progress" cancel; name it, its size, and // which phase was slow so the culprit is visible in the logs. if fetchDur > slowMessage || appendDur > slowMessage { slog.Warn("slow message copy", "uid", uid, "bytes", len(body), "fetch", fetchDur.Round(time.Millisecond), "append", appendDur.Round(time.Millisecond)) } return nil } // keepFlags drops \Recent: it cannot be set via APPEND. go-imap v2 beta.8 // no longer defines an imap.FlagRecent constant (RFC 9051 dropped \Recent // from IMAP4rev2), so match it by its literal wire form instead. func keepFlags(flags []imap.Flag) []imap.Flag { out := make([]imap.Flag, 0, len(flags)) for _, f := range flags { if f == "\\Recent" { continue } out = append(out, f) } return out }