Fix archive reindex and add target dedup

This commit is contained in:
DonVoo 2026-07-14 20:57:09 +02:00
parent e612896e06
commit c83225f5ad
9 changed files with 796 additions and 6 deletions

View file

@ -132,6 +132,12 @@ func ConnectDB(initDB bool) error {
inner_len INTEGER NOT NULL,
UNIQUE(account_id, folder, message_id)
)`,
`CREATE TABLE IF NOT EXISTS mbox_index_state(
account_id INTEGER NOT NULL,
folder TEXT NOT NULL,
indexed_bytes INTEGER NOT NULL DEFAULT 0,
UNIQUE(account_id, folder)
)`,
`CREATE TABLE IF NOT EXISTS jobs(
id INTEGER PRIMARY KEY AUTOINCREMENT,
account_id INTEGER NOT NULL REFERENCES accounts(id) ON DELETE CASCADE,
@ -442,6 +448,23 @@ func DeleteAccount(name string) error {
return err
}
func FolderMap(accountID int64) (map[string]string, error) {
rows, err := DB.Query(`SELECT src_folder, dst_folder FROM folder_map WHERE account_id=?`, accountID)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]string{}
for rows.Next() {
var src, dst string
if err := rows.Scan(&src, &dst); err != nil {
return nil, err
}
out[src] = dst
}
return out, rows.Err()
}
func ListArchiveMailboxes() ([]string, error) {
rows, err := DB.Query(`SELECT name FROM archive_mailboxes ORDER BY lower(name), name`)
if err != nil {
@ -544,6 +567,64 @@ func SaveMboxIndex(e MboxIndexEntry) error {
return err
}
func ReplaceMboxIndex(accountID int64, folder string, entries []MboxIndexEntry, indexedBytes int64) error {
if DB == nil || accountID == 0 {
return nil
}
tx, err := DB.Begin()
if err != nil {
return err
}
defer tx.Rollback()
if _, err := tx.Exec(`DELETE FROM mbox_index WHERE account_id=? AND folder=?`, accountID, folder); err != nil {
return err
}
for i, entry := range entries {
entry.AccountID = accountID
entry.Folder = folder
entry.Seq = i
if entry.MessageID == "" {
continue
}
if _, err := tx.Exec(`INSERT INTO mbox_index(account_id, folder, seq, message_id, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len)
VALUES(?,?,?,?,?,?,?,?,?,?,?)`,
entry.AccountID, entry.Folder, entry.Seq, entry.MessageID, entry.Subject, entry.From, entry.Date,
entry.FileOffset, entry.FrameLen, entry.InnerOffset, entry.InnerLen); err != nil {
return err
}
}
if _, err := tx.Exec(`INSERT INTO mbox_index_state(account_id, folder, indexed_bytes)
VALUES(?,?,?)
ON CONFLICT(account_id, folder) DO UPDATE SET indexed_bytes=excluded.indexed_bytes`,
accountID, folder, indexedBytes); err != nil {
return err
}
return tx.Commit()
}
func UpdateMboxIndexState(accountID int64, folder string, indexedBytes int64) error {
if DB == nil || accountID == 0 {
return nil
}
_, err := DB.Exec(`INSERT INTO mbox_index_state(account_id, folder, indexed_bytes)
VALUES(?,?,?)
ON CONFLICT(account_id, folder) DO UPDATE SET indexed_bytes=excluded.indexed_bytes`,
accountID, folder, indexedBytes)
return err
}
func MboxIndexedBytes(accountID int64, folder string) (int64, error) {
if DB == nil || accountID == 0 {
return 0, sql.ErrNoRows
}
var n int64
err := DB.QueryRow(`SELECT indexed_bytes FROM mbox_index_state WHERE account_id=? AND folder=?`, accountID, folder).Scan(&n)
if errors.Is(err, sql.ErrNoRows) {
return 0, nil
}
return n, err
}
func ListMboxIndex(accountID int64, folder string) ([]MboxIndexEntry, error) {
if DB == nil || accountID == 0 {
return nil, sql.ErrNoRows

View file

@ -18,6 +18,7 @@ type TargetMailbox interface {
Append(folder string, m RawMessage) error // mit m.Flags und m.InternalDate
Headers(folder string, limit, offset int) ([]MessageHeader, error)
FetchOne(folder string, uid uint32) (RawMessage, error)
DeleteUIDs(folder string, uids []uint32) error
Close() error
}
@ -129,6 +130,32 @@ func (t *imapTarget) FetchOne(folder string, uid uint32) (RawMessage, error) {
return t.mailbox.fetchOne(folder, uid)
}
func (t *imapTarget) DeleteUIDs(folder string, uids []uint32) error {
if len(uids) == 0 {
return nil
}
if _, err := t.mailbox.c.Select(folder, nil).Wait(); err != nil {
return err
}
uidSet := imap.UIDSetNum(uint32sToUIDs(uids)...)
flags := &imap.StoreFlags{Op: imap.StoreFlagsAdd, Silent: true, Flags: []imap.Flag{imap.FlagDeleted}}
if err := t.mailbox.c.Store(uidSet, flags, nil).Close(); err != nil {
return err
}
_, err := t.mailbox.c.UIDExpunge(uidSet).Collect()
return err
}
func (t *imapTarget) Close() error {
return t.mailbox.Close()
}
func uint32sToUIDs(values []uint32) []imap.UID {
out := make([]imap.UID, 0, len(values))
for _, value := range values {
if value != 0 {
out = append(out, imap.UID(value))
}
}
return out
}

View file

@ -114,6 +114,11 @@ func mboxRecord(m RawMessage) []byte {
func mboxAppendInfoFromMessage(m RawMessage, record []byte) MboxAppendInfo {
info := MboxAppendInfo{InnerLen: int64(len(record))}
parts := splitMboxRecordsWithOffsets(record)
if len(parts) > 0 {
info.InnerOffset = parts[0].Offset
info.InnerLen = int64(len(parts[0].Message))
}
msg, err := mail.ReadMessage(bytes.NewReader(m.Body))
if err != nil {
return info
@ -137,6 +142,9 @@ type MboxEntry struct {
// ReadMboxList parst die Kopfzeilen aller Mails einer mbox-Datei (fuer die
// Nachrichtenliste). ReadMboxMessage liefert eine einzelne Mail als Rohtext.
func ReadMboxList(path string) ([]MboxEntry, error) {
if err := EnsureMboxIndexed(path); err != nil {
return nil, err
}
if entries, ok := readMboxListFromIndex(path); ok {
return entries, nil
}
@ -162,6 +170,9 @@ func ReadMboxList(path string) ([]MboxEntry, error) {
}
func ReadMboxMessage(path string, index int) ([]byte, error) {
if err := EnsureMboxIndexed(path); err != nil {
return nil, err
}
if raw, ok, err := readMboxMessageFromIndex(path, index); ok || err != nil {
return raw, err
}
@ -194,6 +205,229 @@ func readMboxMessages(path string) ([][]byte, error) {
return readMboxMessagesBytes(b), nil
}
func EnsureMboxIndexed(path string) error {
accountID, folder, ok := mboxIndexContext(path)
if !ok {
return nil
}
info, err := os.Stat(path)
if err != nil {
return err
}
indexed, err := MboxIndexedBytes(accountID, folder)
if err != nil {
return err
}
if indexed == info.Size() {
return nil
}
return ReindexMboxFile(accountID, folder, path)
}
func ReindexMboxFile(accountID int64, folder, path string) error {
info, err := os.Stat(path)
if err != nil {
return err
}
var entries []MboxIndexEntry
if strings.HasSuffix(strings.ToLower(path), ".zst") {
entries, err = reindexZstdMbox(path)
} else {
entries, err = reindexPlainMbox(path)
}
if err != nil {
return err
}
return ReplaceMboxIndex(accountID, folder, entries, info.Size())
}
func ReindexArchives(name string) error {
name = strings.TrimSpace(name)
if name == "" {
return fmt.Errorf("reindex ziel fehlt")
}
accounts, err := ListAccounts()
if err != nil {
return err
}
matched := 0
for _, account := range accounts {
archive := accountMboxDir(account)
if archive == "" {
continue
}
if name != "all" && !strings.EqualFold(name, account.Name) && !strings.EqualFold(name, archive) {
continue
}
if err := reindexAccountArchive(account); err != nil {
return err
}
matched++
}
if matched == 0 {
return fmt.Errorf("kein Konto/Archiv fuer reindex %q gefunden", name)
}
return nil
}
func reindexAccountArchive(account Account) error {
dir := filepath.Join(Cfg.MboxRoot, accountMboxDir(account))
for _, path := range archiveMboxPaths(dir) {
folder := folderNameFromMboxPath(path)
if folder == "" {
continue
}
if err := ReindexMboxFile(account.ID, folder, path); err != nil {
return fmt.Errorf("%s %s: %w", account.Name, folder, err)
}
}
return nil
}
func reindexPlainMbox(path string) ([]MboxIndexEntry, error) {
b, err := os.ReadFile(path)
if err != nil {
return nil, err
}
var out []MboxIndexEntry
for _, part := range splitMboxRecordsWithOffsets(b) {
entry := indexEntryFromMboxRecord(part.Message, part.Offset, int64(part.Length), 0, int64(part.Length))
if entry.MessageID != "" {
out = append(out, entry)
}
}
return out, nil
}
func reindexZstdMbox(path string) ([]MboxIndexEntry, error) {
b, err := os.ReadFile(path)
if err != nil {
return nil, err
}
dec, err := zstd.NewReader(nil)
if err != nil {
return nil, err
}
defer dec.Close()
var out []MboxIndexEntry
for _, frame := range splitZstdFrames(b) {
record, err := dec.DecodeAll(frame.Data, nil)
if err != nil {
return nil, err
}
for _, part := range splitMboxRecordsWithOffsets(record) {
entry := indexEntryFromMboxRecord(part.Message, frame.Offset, int64(frame.Length), part.Offset, int64(len(part.Message)))
if entry.MessageID != "" {
out = append(out, entry)
}
}
}
return out, nil
}
type mboxRecordPart struct {
Offset int64
Length int
Message []byte
}
func splitMboxRecordsWithOffsets(b []byte) []mboxRecordPart {
normalized := bytes.ReplaceAll(b, []byte("\r\n"), []byte("\n"))
lines := bytes.SplitAfter(normalized, []byte("\n"))
var out []mboxRecordPart
var cur bytes.Buffer
var curStart int64
var pos int64
inMsg := false
for _, lineWithNL := range lines {
line := bytes.TrimSuffix(lineWithNL, []byte("\n"))
if bytes.HasPrefix(line, []byte("From ")) {
if inMsg && cur.Len() > 0 {
msg := bytes.TrimRight(cur.Bytes(), "\n")
out = append(out, mboxRecordPart{Offset: curStart, Length: int(pos - curStart), Message: append([]byte(nil), msg...)})
cur.Reset()
}
inMsg = true
pos += int64(len(lineWithNL))
curStart = pos
continue
}
if inMsg {
if bytes.HasPrefix(line, []byte(">From ")) {
line = line[1:]
}
_, _ = cur.Write(line)
_ = cur.WriteByte('\n')
}
pos += int64(len(lineWithNL))
}
if inMsg && cur.Len() > 0 {
msg := bytes.TrimRight(cur.Bytes(), "\n")
out = append(out, mboxRecordPart{Offset: curStart, Length: int(pos - curStart), Message: append([]byte(nil), msg...)})
}
return out
}
type zstdFramePart struct {
Offset int64
Length int
Data []byte
}
func splitZstdFrames(b []byte) []zstdFramePart {
magic := []byte{0x28, 0xb5, 0x2f, 0xfd}
var starts []int
for i := 0; i <= len(b)-len(magic); i++ {
if bytes.Equal(b[i:i+len(magic)], magic) {
starts = append(starts, i)
}
}
if len(starts) == 0 {
return nil
}
out := make([]zstdFramePart, 0, len(starts))
for i, start := range starts {
end := len(b)
if i+1 < len(starts) {
end = starts[i+1]
}
out = append(out, zstdFramePart{Offset: int64(start), Length: end - start, Data: b[start:end]})
}
return out
}
func indexEntryFromMboxRecord(raw []byte, fileOffset, frameLen, innerOffset, innerLen int64) MboxIndexEntry {
entry := MboxIndexEntry{FileOffset: fileOffset, FrameLen: frameLen, InnerOffset: innerOffset, InnerLen: innerLen}
msg, err := mail.ReadMessage(bytes.NewReader(raw))
if err != nil {
return entry
}
entry.MessageID = normalizeMessageID(msg.Header.Get("Message-ID"))
if entry.MessageID == "" {
entry.MessageID = messageID(raw)
}
entry.Subject = decodeHeader(msg.Header.Get("Subject"))
entry.From = decodeHeader(msg.Header.Get("From"))
entry.Date = decodeHeader(msg.Header.Get("Date"))
return entry
}
func unescapeMboxMessage(raw []byte) []byte {
raw = bytes.ReplaceAll(raw, []byte("\r\n"), []byte("\n"))
lines := bytes.Split(raw, []byte("\n"))
var b bytes.Buffer
for i, line := range lines {
if i > 0 {
_ = b.WriteByte('\n')
}
if bytes.HasPrefix(line, []byte(">From ")) {
line = line[1:]
}
_, _ = b.Write(line)
}
return bytes.TrimRight(b.Bytes(), "\n")
}
func readMboxMessagesBytes(b []byte) [][]byte {
b = bytes.ReplaceAll(b, []byte("\r\n"), []byte("\n"))
lines := bytes.Split(b, []byte("\n"))
@ -278,11 +512,12 @@ func readMboxMessageFromIndex(path string, index int) ([]byte, bool, error) {
if err != nil {
return nil, true, err
}
msgs := readMboxMessagesBytes(record)
if len(msgs) == 0 {
start := entry.InnerOffset
end := entry.InnerOffset + entry.InnerLen
if start < 0 || end > int64(len(record)) || start >= end {
return nil, true, fmt.Errorf("message index out of range")
}
return msgs[0], true, nil
return unescapeMboxMessage(record[start:end]), true, nil
}
func mboxIndexContext(path string) (int64, string, bool) {

View file

@ -118,3 +118,100 @@ func TestZstdMboxRoundTripAndIndexRead(t *testing.T) {
t.Fatalf("unexpected message body: %q", raw)
}
}
func TestPlainMboxPartialIndexIsRebuiltBeforeList(t *testing.T) {
oldCfg, oldDB := Cfg, DB
root := t.TempDir()
t.Cleanup(func() {
if DB != nil {
_ = DB.Close()
}
Cfg, DB = oldCfg, oldDB
})
Cfg = Config{
DBPath: filepath.Join(root, "mail-graveyard.db"),
MboxRoot: filepath.Join(root, "backup"),
MboxCompression: "none",
}
if err := ConnectDB(true); err != nil {
t.Fatal(err)
}
if err := SaveArchiveMailbox("archive"); err != nil {
t.Fatal(err)
}
if err := SaveAccount(Account{
Name: "test-account",
SrcHost: "source.example",
SrcPort: 993,
SrcSecurity: "tls",
SrcUser: "source@example.com",
SrcPass: "x",
DstHost: "target.example",
DstPort: 993,
DstSecurity: "tls",
DstUser: "target@example.com",
DstPass: "x",
MboxDir: "archive",
Active: true,
}); err != nil {
t.Fatal(err)
}
account, err := GetAccount("test-account")
if err != nil {
t.Fatal(err)
}
writer, err := NewMboxWriter(filepath.Join(Cfg.MboxRoot, "archive"))
if err != nil {
t.Fatal(err)
}
first := testRawMessage("first@example.com", "First")
firstInfo, err := writer.Append("INBOX", first)
if err != nil {
t.Fatal(err)
}
if err := SaveMboxIndex(MboxIndexEntry{
AccountID: account.ID,
Folder: "INBOX",
MessageID: first.MessageID,
Subject: firstInfo.Subject,
From: firstInfo.From,
Date: firstInfo.Date,
FileOffset: firstInfo.FileOffset,
FrameLen: firstInfo.FrameLen,
InnerOffset: firstInfo.InnerOffset,
InnerLen: firstInfo.InnerLen,
}); err != nil {
t.Fatal(err)
}
if err := UpdateMboxIndexState(account.ID, "INBOX", firstInfo.FileOffset+firstInfo.FrameLen); err != nil {
t.Fatal(err)
}
if _, err := writer.Append("INBOX", testRawMessage("second@example.com", "Second")); err != nil {
t.Fatal(err)
}
entries, err := ReadMboxList(firstInfo.Path)
if err != nil {
t.Fatal(err)
}
if len(entries) != 2 {
t.Fatalf("expected rebuilt full list with 2 entries, got %#v", entries)
}
if entries[0].Subject != "First" || entries[1].Subject != "Second" {
t.Fatalf("unexpected entries after reindex: %#v", entries)
}
}
func testRawMessage(id, subject string) RawMessage {
return RawMessage{
MessageID: id,
InternalDate: time.Date(2026, 7, 14, 10, 11, 12, 0, time.UTC),
Body: []byte("Message-ID: <" + id + ">\r\n" +
"From: Sender <sender@example.com>\r\n" +
"Subject: " + subject + "\r\n" +
"Date: Tue, 14 Jul 2026 10:11:12 +0000\r\n" +
"\r\n" +
"Hello\r\n"),
}
}

View file

@ -114,12 +114,17 @@ func migrateAccount(a Account, selectedFolders map[string]bool) error {
_ = finishJob(jobID, 0, 0, 1, "error")
return err
}
folderMap, err := FolderMap(a.ID)
if err != nil {
_ = finishJob(jobID, 0, 0, 1, "error")
return err
}
total, done, errs := 0, 0, 0
for _, folder := range folders {
if len(selectedFolders) > 0 && !selectedFolders[folder.Name] {
continue
}
dstFolder := MapSourceToTarget(folder.Name, folder.Attrs, folder.Delim, dst.Delim(), targets, "")
dstFolder := MapSourceToTarget(folder.Name, folder.Attrs, folder.Delim, dst.Delim(), targets, folderMap[folder.Name])
if dstFolder == "" {
dstFolder = folder.Name
}
@ -173,6 +178,12 @@ func migrateAccount(a Account, selectedFolders map[string]bool) error {
log.Printf("%s %s mbox index error %s: %v", a.Name, folder.Name, m.MessageID, err)
return nil
}
if err := UpdateMboxIndexState(a.ID, folder.Name, appendInfo.FileOffset+appendInfo.FrameLen); err != nil {
errs++
folderErrs++
log.Printf("%s %s mbox index state error %s: %v", a.Name, folder.Name, m.MessageID, err)
return nil
}
if err := MarkCopied(a.ID, folder.Name, m.MessageID); err != nil {
errs++
folderErrs++

140
backend/13-target-dedup.go Normal file
View file

@ -0,0 +1,140 @@
package backend
import (
"fmt"
"log"
"sort"
"strings"
)
type TargetDedupReport struct {
Account string
Apply bool
Folders int
Groups int
ExtraCopies int
WithoutID int
Deleted int
FolderReports []TargetDedupFolderReport
}
type TargetDedupFolderReport struct {
Folder string
Groups int
ExtraCopies int
WithoutID int
Deletes []TargetDedupDelete
}
type TargetDedupDelete struct {
UID uint32
MessageID string
}
func DedupTarget(name string, folders []string, apply bool) (TargetDedupReport, error) {
a, err := GetAccount(name)
if err != nil {
return TargetDedupReport{}, err
}
dst, err := OpenIMAPTarget(a)
if err != nil {
return TargetDedupReport{}, err
}
defer dst.Close()
allFolders, err := dst.Folders()
if err != nil {
return TargetDedupReport{}, err
}
selected := selectedFolderSet(folders)
report := TargetDedupReport{Account: name, Apply: apply}
for _, folder := range allFolders {
if len(selected) > 0 && !selected[folder.Name] {
continue
}
folderReport, err := dedupTargetFolder(dst, folder.Name, apply)
if err != nil {
return report, err
}
report.Folders++
report.Groups += folderReport.Groups
report.ExtraCopies += folderReport.ExtraCopies
report.WithoutID += folderReport.WithoutID
report.Deleted += len(folderReport.Deletes)
report.FolderReports = append(report.FolderReports, folderReport)
}
log.Printf("target dedup %s apply=%v folders=%d groups=%d extra=%d deleted=%d without_id=%d",
report.Account, report.Apply, report.Folders, report.Groups, report.ExtraCopies, report.Deleted, report.WithoutID)
return report, nil
}
func dedupTargetFolder(dst TargetMailbox, folder string, apply bool) (TargetDedupFolderReport, error) {
report := TargetDedupFolderReport{Folder: folder}
byID := map[string][]uint32{}
for offset := 0; ; offset += mailboxListLimit {
headers, err := dst.Headers(folder, mailboxListLimit, offset)
if err != nil {
return report, fmt.Errorf("%s headers: %w", folder, err)
}
if len(headers) == 0 {
break
}
for _, header := range headers {
id := normalizeMessageID(header.MessageID)
if id == "" {
report.WithoutID++
continue
}
byID[id] = append(byID[id], header.UID)
}
}
var deleteUIDs []uint32
for id, uids := range byID {
if len(uids) < 2 {
continue
}
sort.Slice(uids, func(i, j int) bool { return uids[i] < uids[j] })
report.Groups++
report.ExtraCopies += len(uids) - 1
for _, uid := range uids[1:] {
report.Deletes = append(report.Deletes, TargetDedupDelete{UID: uid, MessageID: id})
deleteUIDs = append(deleteUIDs, uid)
}
}
sort.Slice(report.Deletes, func(i, j int) bool {
if report.Deletes[i].MessageID == report.Deletes[j].MessageID {
return report.Deletes[i].UID < report.Deletes[j].UID
}
return strings.Compare(report.Deletes[i].MessageID, report.Deletes[j].MessageID) < 0
})
if apply && len(deleteUIDs) > 0 {
sort.Slice(deleteUIDs, func(i, j int) bool { return deleteUIDs[i] < deleteUIDs[j] })
for _, del := range report.Deletes {
log.Printf("target dedup delete folder=%s uid=%d message_id=%s", folder, del.UID, del.MessageID)
}
if err := dst.DeleteUIDs(folder, deleteUIDs); err != nil {
return report, fmt.Errorf("%s delete uids: %w", folder, err)
}
}
if report.Groups > 0 || report.WithoutID > 0 {
log.Printf("target dedup folder=%s groups=%d extra=%d without_id=%d apply=%v", folder, report.Groups, report.ExtraCopies, report.WithoutID, apply)
}
return report, nil
}
func LogTargetDedupReport(report TargetDedupReport) {
mode := "DRY-RUN"
if report.Apply {
mode = "APPLY"
}
log.Printf("target dedup report %s account=%s folders=%d groups=%d extra=%d deleted=%d without_id=%d",
mode, report.Account, report.Folders, report.Groups, report.ExtraCopies, report.Deleted, report.WithoutID)
for _, folder := range report.FolderReports {
if folder.Groups == 0 && folder.WithoutID == 0 {
continue
}
log.Printf("target dedup report folder=%s groups=%d extra=%d without_id=%d", folder.Folder, folder.Groups, folder.ExtraCopies, folder.WithoutID)
for _, del := range folder.Deletes {
log.Printf("target dedup report candidate folder=%s uid=%d message_id=%s", folder.Folder, del.UID, del.MessageID)
}
}
}